diff --git a/.gitignore b/.gitignore index c9429baa..ea5806c8 100644 --- a/.gitignore +++ b/.gitignore @@ -6,7 +6,7 @@ .pnpm-store/ # logs -logs +/logs/ *.log # config diff --git a/CLAUDE.md b/CLAUDE.md new file mode 100644 index 00000000..af30abe1 --- /dev/null +++ b/CLAUDE.md @@ -0,0 +1,511 @@ +# 项目开发规范 + +> 本文档面向 AI 代理(Agent)和开发者,描述项目的目录结构、模块职责与开发规范。 + +--- + +## 一、技术栈 + +### 后端 + +| 技术 | 用途 | +|------|------| +| Go (1.25+) | 主语言 | +| Gin | HTTP 框架 | +| GORM | ORM,主库 PostgreSQL,可选 ClickHouse | +| Redis | 缓存 / Session / 队列 | +| Asynq | 异步任务队列(基于 Redis) | +| Cobra + Viper | CLI 入口 + 配置加载 | +| Swaggo | Swagger 文档生成 | +| OpenTelemetry | 链路追踪 | +| Zap | 结构化日志 | +| AWS SDK v2 | S3 兼容文件存储 | +| Snowflake | 分布式 ID 生成 | + +### 前端 + +| 技术 | 用途 | +|------|------| +| Next.js (App Router) | 前端框架 | +| TypeScript | 主语言 | +| Tailwind CSS | 样式 | +| pnpm | 包管理 | +| shadcn/ui | 组件库 | + +--- + +## 二、顶层目录结构 + +以下是项目的顶层目录结构及其职责, 如果有新增目录或文件,请务必在此处同步更新: + +``` +Refreshing/ # 项目根目录(模块名: github.com/linux-do/credit) +├── main.go # 程序入口,调用 internal/cmd +├── go.mod / go.sum # Go 模块依赖 +├── config.yaml # 运行时配置(不提交到 Git) +├── config.example.yaml # 配置模板(需提交) +├── DEPLOYMENT_zh.md # 部署说明文档(中文版) +├── Makefile # 常用命令(swagger/tidy/license) +├── Dockerfile # 后端容器镜像构建 +├── docker-compose.yml # 本地依赖服务(PostgreSQL / Redis / ClickHouse) +├── .editorconfig # 编辑器格式规范 +├── .gitignore +├── docs/ # Swagger 自动生成文档(不要手动编辑) +├── frontend/ # Next.js 前端项目 +├── internal/ # 后端核心代码(Go private,不对外暴露) +├── scripts/ # CI/本地工具脚本 +└── support-files/ # 辅助文件(如 nginx 配置等) +``` + +--- + +## 三、后端 `internal/` 目录结构 + +以下是 `internal/` 目录的结构及其职责, 如果有新增目录或文件,请务必在此处同步更新: + +``` +internal/ +├── cmd/ # CLI 命令入口(Cobra) +│ ├── root.go # 根命令,加载配置、初始化依赖 +│ ├── api.go # 启动 HTTP API 服务器子命令 +│ ├── scheduler.go # 启动定时任务调度器子命令 +│ └── worker.go # 启动 Asynq Worker 子命令 +│ +├── config/ # 配置加载与结构定义 +│ ├── model.go # 所有配置结构体(AppConfig / DB / Redis 等) +│ └── config.go # Viper 加载逻辑,暴露全局 config.Config +│ +├── router/ # HTTP 路由注册(唯一路由注册点) +│ ├── router.go # 路由总入口,注册所有分组路由、中间件、启动 HTTP Server +│ └── middlewares.go # 全局中间件(如请求日志) +│ +├── apps/ # 业务功能模块(按功能域划分) +│ ├── oauth/ # OAuth / OIDC 登录、会话、用户信息 +│ ├── user/ # 用户密码登录、注册、登出 +│ ├── upload/ # 文件上传、文件服务、清理任务 +│ ├── health/ # 健康检查端点 +│ ├── config/ # 公开配置接口(前端读取) +│ └── admin/ # 管理后台功能(需 Admin 权限) +│ ├── middlewares.go # Admin 鉴权中间件 +│ ├── errs.go # Admin 错误常量 +│ ├── auth_source/ # 认证源管理(CRUD) +│ ├── system_config/ # 系统配置管理(CRUD) +│ ├── task/ # 任务手动调度接口 +│ └── user/ # 用户管理(列表、状态) +│ +├── model/ # 数据模型(GORM 实体 + 业务方法) +│ ├── users.go # User 实体、OAuthUserInfo、查询/更新方法 +│ ├── auth_source.go # AuthSource 实体(OAuth 接入源) +│ ├── system_configs.go # SystemConfig 实体(KV 系统配置) +│ ├── uploads.go # Upload 实体(上传文件记录) +│ └── task_execution.go # TaskExecution 实体(异步任务执行记录 + CRUD) +│ +├── db/ # 数据库连接与基础设施 +│ ├── postgres.go # PostgreSQL 初始化、读写分离、GORM 配置 +│ ├── redis.go # Redis 初始化(单机/哨兵/集群) +│ ├── clickhouse.go # ClickHouse 初始化(可选) +│ ├── postgres_logger.go # 自定义 GORM 日志(对接 Zap) +│ ├── idgen/ # Snowflake 分布式 ID 生成器 +│ └── migrator/ # 数据库迁移(AutoMigrate) +│ +├── storage/ # 文件存储抽象层 +│ ├── s3.go # S3 兼容存储(上传/下载/URL 生成) +│ ├── cache.go # 本地磁盘缓存(S3 内容缓存) +│ └── errs.go # 存储层错误常量 +│ +├── task/ # 异步任务定义与调度 +│ ├── constants.go # 任务类型名称常量(TaskType)、队列名、TaskMeta(含 Retryable) +│ ├── handler.go # TaskHandler 接口定义 + TaskResult 结构体 +│ ├── executor.go # 核心运行机制:RegisterHandler / DispatchTask / ProcessTask / RetryTask / AppendLog +│ ├── utils.go # 任务工具函数(RedisOpt、AsynqClient) +│ ├── scheduler/ # Asynq 定时任务调度器(Cron 注册) +│ └── worker/ # Asynq Worker 服务端(任务处理器注册) +│ ├── worker.go # StartWorker 入口,注册 Handler +│ └── middlewares.go # Worker 中间件 +│ +├── service/ # 复杂业务逻辑服务层(当前占位,待填充) +│ +├── common/ # 跨模块共享代码 +│ ├── constants.go # 全局常量(错误消息字符串等) +│ ├── errs.go # 通用错误定义 +│ ├── bind/ # 请求参数绑定封装(统一处理错误响应) +│ └── response/ # 统一 HTTP 响应格式封装 +│ +├── util/ # 无业务依赖的纯工具函数 +│ ├── crypto.go # 加密/签名工具 +│ ├── password.go # 密码 Hash(bcrypt) +│ ├── http_clients.go # HTTP 客户端封装 +│ ├── context.go # Context 存取工具 +│ ├── response.go # ResponseAny 等响应结构体 +│ ├── session.go # Session 选项构建 +│ ├── uuid.go # UUID / 唯一 ID 生成 +│ ├── strings.go # 字符串工具 +│ ├── validate.go # 参数校验工具 +│ └── custom_types.go # 自定义类型 +│ +├── logger/ # 日志封装(基于 Zap + OTel) +│ ├── logger.go # 全局 Logger 初始化 +│ └── utils.go # InfoF / WarnF / ErrorF 快捷函数 +│ +├── listener/ # 事件监听器(Webhook / 消息消费) +│ +└── otel_trace/ # OpenTelemetry 链路追踪封装 + └── ... # Span 创建、Exporter 配置 +``` + +--- + +## 四、`apps/` 模块内部文件规范 + +每个业务模块(`apps//`)内部按照以下约定组织文件: + +| 文件名 | 职责 | +|--------|------| +| `routers.go` | **HTTP Handler 函数**(业务逻辑入口,对应 Controller 层)| +| `controllers.go` | 可选,当 Handler 较多时拆分(同 `routers.go` 职责)| +| `middlewares.go` | 本模块专属中间件(如 `LoginRequired`、`LoginAdminRequired`)| +| `errs.go` | 本模块专属错误消息字符串常量(`const`)| +| `constants.go` | 本模块专属业务常量(非错误)| + +> **规则**: +> - 路由 **不在** 模块内部注册,统一在 `internal/router/router.go` 中注册。 +> - `errs.go` 只定义字符串常量,不定义 `error` 类型值,错误通过 `response.RespondFailure(c, errMsg)` 输出。 + +### `admin/` 子模块结构示例 + +``` +apps/admin/ +├── middlewares.go # LoginAdminRequired 中间件 +├── errs.go # admin 级别错误常量 +├── auth_source/ # 认证源 CRUD +│ └── routers.go +├── system_config/ # 系统 KV 配置 CRUD +│ └── routers.go +├── task/ # 任务调度接口 +│ └── routers.go +├── user/ # 用户管理 +│ ├── routers.go +│ └── errs.go +└── user_pay_config/ # 用户支付配置 + └── routers.go +``` + +--- + +## 五、前端 `frontend/` 目录结构 + +以下是前端 `frontend/` 目录的结构, 如果需要调整请在此处同步更改: + +``` +frontend/ +├── app/ # Next.js App Router 页面目录 +│ ├── layout.tsx # 根布局(全局 Provider、字体、meta) +│ ├── globals.css # 全局样式 +│ ├── page.tsx # 首页重定向 +│ ├── (auth)/ # 认证相关页面组(登录/注册/OAuth 回调) +│ ├── (main)/ # 主应用页面组(用户界面) +│ └── (docs)/ # 文档类页面组 +│ +├── components/ # 可复用 React 组件 +│ ├── ui/ # shadcn/ui 基础组件(Button/Input/Dialog 等) +│ ├── common/ # 通用业务组件(跨页面复用),详见下方说明 +│ ├── layout/ # 布局组件(Header / Sidebar / Footer) +│ ├── auth/ # 认证相关组件 +│ ├── home/ # 首页专属组件 +│ ├── animate-ui/ # 动画 UI 组件 +│ └── providers/ # Context Provider 组件 +│ +├── contexts/ # React Context(全局状态) +├── hooks/ # 自定义 React Hooks +├── lib/ # 前端工具函数、API 客户端封装 +├── types/ # TypeScript 类型定义 +├── public/ # 静态资源 +├── proxy.ts # 开发环境代理配置 +├── next.config.ts # Next.js 配置 +├── package.json +├── tsconfig.json +├── .env # 环境变量(不提交) +└── .env.example # 环境变量模板(需提交) +``` + +--- + +## 5.1 前端 `components/common/` 通用业务组件详解 + +`common/` 目录存放跨页面复用的业务组件,按功能域分为五个子目录。以下是每个文件的职责说明: + +``` +components/common/ +├── admin/ # 管理员后台组件 +│ ├── tasks.tsx # TaskManager — 异步任务调度管理页面,展示所有可用任务类型, +│ │ # 支持通过弹窗配置参数后立即下发任务到后台队列执行 +│ ├── task-executions.tsx # TaskExecutionsManager — 任务日志页面,展示异步任务执行记录, +│ │ # 支持状态/类型筛选、分页、详情抽屉查看完整日志与失败任务重试 +│ ├── system.tsx # SystemConfigs — 系统 KV 配置管理页面,以表格展示系统/业务两类 +│ │ # 配置项,支持在线编辑(布尔类型自动渲染为 Switch)并保存/删除 +│ └── users.tsx # UsersManager — 用户管理页面,提供分页、搜索、筛选的用户列表表格, +│ # 支持在侧边抽屉查看用户详情,以及启用/禁用(封禁/解封)切换 +│ +├── docs/ # 文档页面组件,包括法律文档(隐私政策/服务条款)和接口文档 +│ +├── general/ # 通用框架组件 +│ ├── manage-pannel.tsx # ManagePage(泛型)— 通用管理页面框架,封装"列表 + 详情面板"布局, +│ │ # 包含数据加载/错误/空状态处理、表格渲染、选中/悬停交互、 +│ │ # 编辑/保存/删除逻辑;ManageDetailPanel 为带保存按钮的详情面板; +│ │ # ManageTable 为配置驱动型表格组件 +│ └── password-dialog.tsx # PasswordDialog — 密码确认弹窗,用于敏感操作前的二次身份验证, +│ # 包含 6 位 OTP 输入框,支持 Enter 快捷确认,带加载状态显示 +│ +├── home/ # 首页组件 +│ └── home-main.tsx # HomeMain — 系统首页主内容,展示当前用户的快捷导航卡片 +│ # (个人资料、开发接口文档、使用文档),管理员额外显示后台管理入口 +│ +└── settings/ # 设置页面组件 + ├── access-token.tsx # AccessTokenMain — 个人访问令牌管理页面,展示用户 API 密钥列表, + │ # 支持创建(仅展示一次明文)、轮换、撤销/删除令牌 + ├── appearance.tsx # AppearanceMain — 外观设置页面,分为主题模式选择 + │ # (明亮/黑暗/自动)和界面配色方案(可视化色卡网格切换) + ├── auth-source-modal.tsx # AuthSourceModal — OIDC 认证源新增/编辑弹窗,包含标识符、 + │ # Client ID/Secret、Discovery URL、Scopes、图标等表单字段 + ├── notifications.tsx # NotificationsMain — 通知设置页面,控制顶部导航栏 + │ # 是否显示通知铃铛图标,通过 Context 持久化偏好 + ├── profile.tsx # ProfileMain — 个人资料页面,展示用户基本信息,提供第三方 + │ # 账号绑定管理(查看已绑定 OIDC 账号、解除绑定、绑定新认证源) + └── security.tsx # SecurityMain — 系统安全设置页面(管理员专属),包含系统登录与 + # 注册控制(密码登录/注册/密码注册/OIDC 登录四个开关)、 + # 认证源管理(新增、编辑、启用/禁用、删除 OIDC 认证源) +``` + +--- + +## 六、开发规范 + +### 6.1 命名规范 + +| 对象 | 规范 | 示例 | +|------|------|------| +| Go 包名 | 小写,下划线分词(单词) | `auth_source`、`system_config` | +| Go 文件名 | 小写,下划线分词 | `routers.go`、`postgres_logger.go` | +| Go 导出函数 | PascalCase | `ListUsers`、`StartWorker` | +| Go 未导出函数 | camelCase | `buildQueuesFromConfig` | +| Go 结构体请求/响应 | camelCase + 后缀 | `listUsersRequest`、`listUsersResponse` | +| 错误常量 | camelCase 字符串 `const` | `const userNotFound = "用户不存在"` | +| 任务类型常量 | 全大写蛇形 | `CleanupUnusedUploadsTask` | +| 配置 Key | 全小写蛇形(YAML) | `session_cookie_name`、`max_idle_conn` | + +### 6.2 HTTP Handler 规范 + +```go +// Handler 函数命名:动词 + 名词(PascalCase) +func ListUsers(c *gin.Context) { + // 1. 参数绑定(使用 ShouldBindQuery / ShouldBindJSON) + var req listUsersRequest + if err := c.ShouldBindQuery(&req); err != nil { + c.JSON(http.StatusBadRequest, util.Err(err.Error())) + return + } + + // 2. 业务逻辑 + + // 3. 统一响应 + c.JSON(http.StatusOK, util.OK(data)) +} +``` + +**响应格式约定**: +- 成功:`util.OK(data)` 或 `util.OKNil()` +- 失败:`util.Err(msg)` + 对应 HTTP 状态码 +- 通过 `response.RespondSuccess / RespondFailure` 也可(两套工具共存) + +### 6.4 错误处理规范 + +- **模块内错误消息**:定义在本模块 `errs.go` 中,使用 `const` 字符串。 +- **跨模块错误消息**:定义在 `internal/common/errs.go` 或 `common/constants.go`。 +- **数据库错误**:直接 `err.Error()` 返回给响应(开发阶段),生产环境应屏蔽详情。 +- **gorm.ErrRecordNotFound**:显式判断,返回 404。 + +### 6.5 中间件使用规范 + +| 中间件 | 位置 | 作用 | +|--------|------|------| +| `gin.Recovery()` | 全局 | Panic 恢复 | +| `otelgin.Middleware()` | 全局 | OTel 链路追踪 | +| `loggerMiddleware()` | 全局 | 请求日志 | +| `sessions.Sessions()` | 全局 | Session 注入 | +| `oauth.LoginRequired()` | 路由组 | 登录校验 | +| `admin.LoginAdminRequired()` | Admin 路由组 | 管理员校验 | + +### 6.6 配置访问规范 + +- 所有配置通过 `config.Config.
.` 访问(全局单例)。 +- 不允许在业务代码中使用 `os.Getenv()` 读取配置,统一通过 Viper 加载。 +- 新增配置项:先在 `config.example.yaml` 添加注释模板,再在 `internal/config/model.go` 添加结构体字段。 + +### 6.7 数据库访问规范 + +- 直接使用 GORM:`model.DB.Where(...).Find(&result)`(适合简单查询)。 +- 通过 `db.DB(ctx)` 获取带链路追踪的 DB 实例(Admin 模块推荐)。 +- 禁止在 Handler 层直接写复杂 SQL,应封装到 `model/` 层方法或 `service/` 层。 +- 数据库迁移使用 `db/migrator/` 中的 AutoMigrate,不允许手动执行 DDL。 + +### 6.8 异步任务规范 + +**定义任务**: +1. 在 `internal/task/constants.go` 中定义任务类型常量。 +2. 实现 Handler 函数(放在对应 `apps/` 模块的 `tasks.go` 文件中)。 +3. 在 `internal/task/worker/worker.go` 中注册 Handler:`mux.HandleFunc(task.XxxTask, handler)`。 +4. 调度:在 `internal/task/scheduler/` 中按 Cron 表达式调度,或通过 Admin API 手动触发。 + +**队列优先级**(从高到低):`webhook` > `whitelist_only` > `default` + +### 6.9 前端组件样式规范 + +**基础组件必须遵循系统的色彩主题系统。** 所有基于 shadcn/ui 的基础组件(Button、Dialog、Input 等)应使用组件内置的 `variant` 属性来控制样式,禁止通过 `className` 手写颜色或背景等样式。 + +**错误示例(禁止)**: + +```tsx +// ❌ 禁止通过 className 手写颜色、背景、阴影等样式 + +``` + +**正确示例**: + +```tsx +// ✅ 使用 variant 属性,让组件遵循系统主题 + +``` + +> **原则**:组件的视觉表现由 shadcn/ui 的 variant 系统和全局 CSS 变量统一控制,保持应用内所有页面风格一致。如现有 variant 无法满足需求,应扩展 shadcn/ui 组件的 variant 定义,而非在业务代码中硬编码颜色值。 + +### 6.10 严格禁止事项 + +| 禁止行为 | 说明 | +|----------|------| +| **禁止删除 `node_modules` 目录** | `node_modules` 为前端依赖安装目录,删除会导致项目无法运行。如需重新安装依赖,使用 `pnpm install` 覆盖更新即可,严禁执行 `rm -rf node_modules`。 | +| **`internal/util/` 下禁止引用框架包** | `util/` 及其子包(如 `util/cap`)定位为**纯工具层**,不得 `import` 任何 HTTP / ORM / 框架包,包括但不限于 `github.com/gin-gonic/gin`、`gorm.io/gorm`、`github.com/gin-contrib/sessions`。违反此约束会导致工具层与框架产生耦合,无法独立测试。详见 **6.11** 的建议方案。 | + +### 6.11 `util/` 包依赖约束与建议方案 + +#### 约束范围 + +`internal/util/` 及其全部子包(如 `util/cap`、`util/crypto` 等)只允许引用: + +- Go 标准库(`context`、`crypto`、`encoding`、`net/http` 原生包等) +- 项目内同级别的纯工具包(`internal/config`、`internal/db`、`internal/model` 等无框架依赖的包) +- 与框架无关的第三方库(如 `github.com/redis/go-redis`、`github.com/shopspring/decimal` 等) + +**严禁引用**:`github.com/gin-gonic/gin`、`gorm.io/gorm`、`github.com/gin-contrib/sessions` 及任何 HTTP 框架 / Web 中间件相关包。 + +#### 常见误区与建议方案 + +| 误区 | 建议方案 | +|------|----------| +| 在 `util/` 中写 `gin.HandlerFunc` 形式的中间件 | 将中间件移至对应的 `apps//middleware.go`,通过**函数参数**接收 `util/` 层的核心对象(如 `*cap.Manager`) | +| 在 `util/` 中通过 `*gin.Context` 写响应 | 只在 `util/` 中计算/校验逻辑并返回 `(result, error)`,由 `apps/` 层的 Handler 负责调用 `c.AbortWithStatusJSON` 写响应 | +| 在 `util/` 中使用 `gorm.DB` 直接查询 | 将数据库查询封装在 `internal/model/` 层方法中,`util/` 只接收已查出的数据结构 | + +#### 正确示例 + +```go +// ✅ internal/util/cap/manager.go — 纯逻辑,无框架依赖 +func (m *Manager) VerifyToken(ctx context.Context, token, scope string) (bool, error) { + // 只依赖 context、标准库、redis client + ... +} + +// ✅ internal/apps/cap/middleware.go — 框架胶水层,持有 gin 依赖 +func VerifyMiddleware(mgr *caputil.Manager, scope string, enabledFunc func() bool) gin.HandlerFunc { + return func(c *gin.Context) { + valid, err := mgr.VerifyToken(c.Request.Context(), token, scope) // 调用纯逻辑 + if err != nil || !valid { + c.AbortWithStatusJSON(http.StatusUnauthorized, util.Err("验证码校验失败")) + return + } + c.Next() + } +} +``` + +#### 错误示例(禁止) + +```go +// ❌ internal/util/cap/middleware.go — util/ 层不应出现 gin +import "github.com/gin-gonic/gin" + +func (m *Manager) VerifyMiddleware(...) gin.HandlerFunc { ... } +``` + +--- + +## 七、新增功能开发流程 + +新增 **异步任务** :使用项目专属 SKILL: new-async-task 进行开发 + +以新增 **管理员功能模块** 为例: + +``` +1. 在 internal/model/ 中定义/扩展数据模型 +2. 在 db/migrator/ 中注册 AutoMigrate +3. 在 internal/apps/admin// 中创建: + - routers.go (Handler 实现 + Swagger 注释) + - errs.go (错误常量,按需) +4. 在 internal/router/router.go 中注册路由 +5. 执行 make swagger 更新文档 +``` + +**Handler 文件拆分规则**: + +逻辑简单的 CRUD 可以全部放在 `routers.go` 中。但当文件代码行数增长时,必须按以下规则拆分: + +| 条件 | 拆分方式 | +|---------------------------|----------| +| 文件超过 **600 行** | 必须拆分 | +| 包含复杂业务逻辑(如外部调用、多步校验、事务处理) | 将业务逻辑拆到 `logic.go` 或 `logics.go` | +| 同一模块有多个独立功能域 | 按功能域拆分多个文件,如 `user_routers.go`、`role_routers.go` | + +拆分后的模块文件结构示例: + +``` +apps/admin// +├── routers.go # 路由注册入口 + 简单 Handler(参数绑定 → 调用逻辑 → 响应) +├── logics.go # 复杂业务逻辑(外部调用、事务、多步处理) +├── errs.go # 错误常量 +└── constants.go # 业务常量(按需) +``` + +**职责边界**: + +- `routers.go` 只做三件事:参数绑定、调用 logic 函数、返回响应。不包含任何业务判断逻辑。 +- `logics.go` 负责所有业务逻辑,接收已校验的参数,返回处理结果和错误。函数以 `PascalCase` 导出,供 `routers.go` 调用。 + + +--- + +## 八、前端任务管理页面 + +任务管理 API 路由(Admin): + +| 方法 | 路径 | 说明 | +|------|------|------| +| GET | `/api/v1/admin/tasks/types` | 获取可调度任务类型列表 | +| POST | `/api/v1/admin/tasks/dispatch` | 手动下发任务 | +| GET | `/api/v1/admin/tasks/executions` | 分页查询任务执行记录(支持 status / task_type 筛选) | +| GET | `/api/v1/admin/tasks/executions/:id` | 查询单条任务执行详情(含完整 Log) | +| POST | `/api/v1/admin/tasks/executions/:id/retry` | 重试失败任务(校验 Retryable && RetryCount < MaxRetry) | + +--- diff --git a/frontend/app/(main)/admin/logs/page.tsx b/frontend/app/(main)/admin/logs/page.tsx new file mode 100644 index 00000000..317a750d --- /dev/null +++ b/frontend/app/(main)/admin/logs/page.tsx @@ -0,0 +1,8 @@ +"use client" + +import {SystemLogs} from "@/components/common/admin/system-logs" + +/* 系统日志页面 */ +export default function LogsPage() { + return +} diff --git a/frontend/components/common/admin/system-logs.tsx b/frontend/components/common/admin/system-logs.tsx new file mode 100644 index 00000000..c464d753 --- /dev/null +++ b/frontend/components/common/admin/system-logs.tsx @@ -0,0 +1,334 @@ +"use client" + +import {useCallback, useEffect, useRef, useState} from "react" +import {toast} from "sonner" +import {ArrowDown, ChevronUp, Loader2, Pause, Play, Terminal} from "lucide-react" + +import {AdminService} from "@/lib/services" +import {ErrorInline} from "@/components/layout/error" +import {LoadingStateWithBorder} from "@/components/layout/loading" +import {Badge} from "@/components/ui/badge" +import {Button} from "@/components/ui/button" + +interface LogEntry { + index: number + data: string +} + +function getApiBaseUrl(): string { + return process.env.NEXT_PUBLIC_LINUX_DO_CREDIT_BACKEND_URL || "" +} + +function buildWsUrl(): string { + const base = getApiBaseUrl() + const wsBase = base.replace(/^http/, "ws") + return `${wsBase}/api/v1/admin/logs/ws` +} + +function parseLogLevel(line: string): "debug" | "info" | "warn" | "error" | "unknown" { + const lower = line.toLowerCase() + if (lower.includes("\"level\":\"error\"") || lower.includes("level=error")) return "error" + if (lower.includes("\"level\":\"warn\"") || lower.includes("level=warn")) return "warn" + if (lower.includes("\"level\":\"debug\"") || lower.includes("level=debug")) return "debug" + if (lower.includes("\"level\":\"info\"") || lower.includes("level=info")) return "info" + return "unknown" +} + +// Distance (px) from bottom to treat as "at bottom" — keeps auto-scroll +// active even when the user's last row is a few pixels shy of the edge. +const BOTTOM_THRESHOLD = 40 + +export function SystemLogs() { + const [loading, setLoading] = useState(true) + const [error, setError] = useState(null) + + const [logs, setLogs] = useState([]) + const [hasMore, setHasMore] = useState(false) + const [nextCursor, setNextCursor] = useState(0) + const [loadingMore, setLoadingMore] = useState(false) + + const [connected, setConnected] = useState(false) + const [paused, setPaused] = useState(false) + + // autoScroll = true → new logs auto-scroll to bottom + // autoScroll = false → user is browsing history, lock scroll position + const [autoScroll, setAutoScroll] = useState(true) + + const containerRef = useRef(null) + const wsRef = useRef(null) + const pausedRef = useRef(paused) + const autoScrollRef = useRef(autoScroll) + const isUserScrolling = useRef(false) + + // Keep refs in sync for use inside callbacks without stale closures + useEffect(() => { pausedRef.current = paused }, [paused]) + useEffect(() => { autoScrollRef.current = autoScroll }, [autoScroll]) + + // ---- Scroll detection ------------------------------------------------ + + // We distinguish programmatic scrolls (triggered by our own auto-scroll or + // load-history offset fixup) from user-initiated scrolls by gating with + // the `isUserScrolling` flag. Only a user scroll can toggle `autoScroll`. + + const handleScroll = useCallback(() => { + const el = containerRef.current + if (!el) return + + // Ignore programmatic scrolls + if (!isUserScrolling.current) return + + const atBottom = el.scrollHeight - el.scrollTop - el.clientHeight < BOTTOM_THRESHOLD + if (atBottom && !autoScrollRef.current) { + setAutoScroll(true) + } else if (!atBottom && autoScrollRef.current) { + setAutoScroll(false) + } + }, []) + + // Mark user-initiated scrolls + const handleWheel = useCallback(() => { isUserScrolling.current = true }, []) + + // Touch devices + const handleTouchStart = useCallback(() => { isUserScrolling.current = true }, []) + + // After the user stops scrolling, reset the flag so our programmatic + // scrolls won't accidentally toggle autoScroll. + useEffect(() => { + const el = containerRef.current + if (!el) return + + let timer: ReturnType + const onScrollEnd = () => { + clearTimeout(timer) + timer = setTimeout(() => { isUserScrolling.current = false }, 150) + } + el.addEventListener("scroll", onScrollEnd, { passive: true }) + return () => { + el.removeEventListener("scroll", onScrollEnd) + clearTimeout(timer) + } + }, []) + + // ---- Auto-scroll to bottom when new logs arrive ---------------------- + + useEffect(() => { + if (!autoScroll || !containerRef.current) return + isUserScrolling.current = false + const el = containerRef.current + requestAnimationFrame(() => { + el.scrollTop = el.scrollHeight + }) + }, [logs, autoScroll]) + + // ---- Data fetching --------------------------------------------------- + + const fetchLogs = useCallback(async (cursor: number = 0) => { + try { + return await AdminService.getLogs(cursor) + } catch (err) { + throw err instanceof Error ? err : new Error("获取日志失败") + } + }, []) + + const loadHistory = useCallback(async (cursor: number = 0) => { + const isInitial = cursor === 0 + if (isInitial) { + setLoading(true) + setError(null) + } else { + setLoadingMore(true) + } + + try { + const data = await fetchLogs(cursor) + if (isInitial) { + setLogs(data.lines || []) + } else { + // Prepend older logs; keep the viewport showing the same content + // by restoring the scroll offset after the DOM update. + setLogs(prev => [...(data.lines || []), ...prev]) + + // Wait for React to render, then compensate scroll position + requestAnimationFrame(() => { + const el = containerRef.current + if (!el) return + // The number of new rows prepended × approximate line height + const newCount = (data.lines || []).length + const lineH = 20 // matches leading-5 ≈ 20px + isUserScrolling.current = false + el.scrollTop = el.scrollTop + newCount * lineH + }) + } + setHasMore(data.has_more) + setNextCursor(data.next_cursor) + } catch (err) { + if (isInitial) { + setError(err instanceof Error ? err : new Error("获取日志失败")) + } else { + toast.error("加载更早日志失败") + } + } finally { + if (isInitial) setLoading(false) + else setLoadingMore(false) + } + }, [fetchLogs]) + + // ---- WebSocket ------------------------------------------------------- + + const connectWs = useCallback(() => { + if (wsRef.current) wsRef.current.close() + + const ws = new WebSocket(buildWsUrl()) + wsRef.current = ws + + ws.onopen = () => { setConnected(true) } + + ws.onmessage = (event) => { + if (pausedRef.current) return + try { + const msg = JSON.parse(event.data) + if (msg.type === "log" && msg.data) { + const entry: LogEntry = msg.data + setLogs(prev => { + const next = [...prev, entry] + return next.length > 2000 ? next.slice(-2000) : next + }) + } + } catch { /* ignore */ } + } + + ws.onclose = () => { setConnected(false); wsRef.current = null } + ws.onerror = () => { setConnected(false) } + }, []) + + // ---- Initialize ------------------------------------------------------ + + useEffect(() => { + loadHistory(0).then(() => connectWs()) + return () => { + wsRef.current?.close() + wsRef.current = null + } + // eslint-disable-next-line react-hooks/exhaustive-deps + }, []) + + // ---- Actions --------------------------------------------------------- + + const scrollToBottom = useCallback(() => { + setAutoScroll(true) + requestAnimationFrame(() => { + if (containerRef.current) { + containerRef.current.scrollTop = containerRef.current.scrollHeight + } + }) + }, []) + + const togglePause = useCallback(() => setPaused(p => !p), []) + const reconnect = useCallback(() => connectWs(), [connectWs]) + const handleLoadMore = useCallback(() => { + if (nextCursor > 0) loadHistory(nextCursor) + }, [nextCursor, loadHistory]) + + // ---- Render ---------------------------------------------------------- + + if (loading) return + if (error) return loadHistory(0)} /> + + return ( +
+ {/* Header */} +
+
+ +

系统日志

+
+
+ + {connected ? "已连接" : "未连接"} + + {connected && ( + + )} + {!connected && ( + + )} +
+
+ + {/* Log viewer — fixed height, scrollable */} +
+ {/* Load older logs */} + {hasMore && ( +
+ +
+ )} + + {/* Log lines */} +
+ {logs.length === 0 ? ( +
暂无日志
+ ) : ( + logs.map((entry) => { + const level = parseLogLevel(entry.data) + const color = level === "error" + ? "text-red-400" + : level === "warn" + ? "text-yellow-400" + : level === "debug" + ? "text-gray-500" + : "text-gray-300" + return ( +
+ {entry.data} +
+ ) + }) + )} +
+
+ + {/* Floating "back to latest" button */} + {!autoScroll && ( +
+ +
+ )} +
+ ) +} diff --git a/frontend/components/layout/sidebar.tsx b/frontend/components/layout/sidebar.tsx index fa079528..6fb7c2ec 100644 --- a/frontend/components/layout/sidebar.tsx +++ b/frontend/components/layout/sidebar.tsx @@ -56,6 +56,7 @@ import { Palette, Settings, ShieldCheck, + Terminal, UserRound, } from "lucide-react" @@ -71,6 +72,7 @@ const data = { { title: "文件管理", url: "/admin/files", icon: FolderOpen }, { title: "任务管理", url: "/admin/tasks", icon: Layers }, { title: "系统配置", url: "/admin/system", icon: ShieldCheck }, + { title: "系统日志", url: "/admin/logs", icon: Terminal }, { title: "系统设置", url: "/admin/settings", icon: Settings }, ], diff --git a/frontend/lib/services/admin/admin.service.ts b/frontend/lib/services/admin/admin.service.ts index dafed3d7..22e30210 100644 --- a/frontend/lib/services/admin/admin.service.ts +++ b/frontend/lib/services/admin/admin.service.ts @@ -347,4 +347,19 @@ export class AdminService extends BaseService { static async getSystemStatus(): Promise { return this.get('/status'); } + + // ==================== 系统日志 ==================== + + /** + * 获取系统历史日志 + * @param cursor - 日志游标,0=获取最新,>0=获取更早 + * @param limit - 每页条数,默认 200 + */ + static async getLogs(cursor: number = 0, limit: number = 200): Promise<{ + lines: Array<{ index: number; data: string }>; + has_more: boolean; + next_cursor: number; + }> { + return this.get('/logs', { cursor, limit }); + } } diff --git a/go.mod b/go.mod index 51a0dcb9..a31a2501 100644 --- a/go.mod +++ b/go.mod @@ -16,6 +16,7 @@ require ( github.com/glebarez/sqlite v1.11.0 github.com/go-jose/go-jose/v4 v4.1.3 github.com/google/uuid v1.6.0 + github.com/gorilla/websocket v1.5.3 github.com/hibiken/asynq v0.25.1 github.com/redis/go-redis/extra/redisotel/v9 v9.16.0 github.com/redis/go-redis/v9 v9.16.0 diff --git a/go.sum b/go.sum index 9c9389b0..9b58df74 100644 --- a/go.sum +++ b/go.sum @@ -172,6 +172,8 @@ github.com/gorilla/securecookie v1.1.2 h1:YCIWL56dvtr73r6715mJs5ZvhtnY73hBvEF8kX github.com/gorilla/securecookie v1.1.2/go.mod h1:NfCASbcHqRSY+3a8tlWJwsQap2VX5pwzwo4h3eOamfo= github.com/gorilla/sessions v1.4.0 h1:kpIYOp/oi6MG/p5PgxApU8srsSw9tuFbt46Lt7auzqQ= github.com/gorilla/sessions v1.4.0/go.mod h1:FLWm50oby91+hl7p/wRxDth9bWSuk0qVL2emc7lT5ik= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= github.com/grpc-ecosystem/grpc-gateway/v2 v2.26.3 h1:5ZPtiqj0JL5oKWmcsq4VMaAW5ukBEgSGXEN89zeH1Jo= github.com/grpc-ecosystem/grpc-gateway/v2 v2.26.3/go.mod h1:ndYquD05frm2vACXE1nsccT4oJzjhw2arTS2cpUD1PI= github.com/hashicorp/go-version v1.7.0 h1:5tqGy27NaOTB8yJKUZELlFAS/LTKJkrmONwQKeRZfjY= diff --git a/internal/apps/admin/logs/routers.go b/internal/apps/admin/logs/routers.go new file mode 100644 index 00000000..e7148e79 --- /dev/null +++ b/internal/apps/admin/logs/routers.go @@ -0,0 +1,131 @@ +/* +Copyright 2025-2026 linux.do + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package logs + +import ( + "encoding/json" + "net/http" + + "github.com/gin-gonic/gin" + "github.com/linux-do/credit/internal/logger" + "github.com/linux-do/credit/internal/util" +) + +const defaultLimit = 200 + +// logsResponse 历史日志查询响应 +type logsResponse struct { + Lines []logger.LogEntry `json:"lines"` + HasMore bool `json:"has_more"` + NextCursor int `json:"next_cursor"` // 用于加载更早日志的 cursor +} + +// GetLogs 获取历史日志 +// @Summary 获取系统日志 +// @Description 分页获取系统历史日志,cursor=0 获取最新日志,cursor>0 获取更早日志 +// @Tags admin +// @Produce json +// @Security SessionCookie +// @Param cursor query int false "日志游标,0=获取最新" default(0) +// @Param limit query int false "每页条数" default(200) +// @Success 200 {object} util.ResponseAny{data=logs.logsResponse} "日志列表" +// @Failure 401 {object} util.ResponseAny "未登录" +// @Failure 403 {object} util.ResponseAny "无管理员权限" +// @Router /api/v1/admin/logs [get] +func GetLogs(c *gin.Context) { + cursorStr := c.DefaultQuery("cursor", "0") + limitStr := c.DefaultQuery("limit", "200") + + var cursor, limit int + if _, err := parsePositiveInt(cursorStr, &cursor); err != nil { + c.JSON(http.StatusBadRequest, util.Err("无效的 cursor 参数")) + return + } + if _, err := parsePositiveInt(limitStr, &limit); err != nil || limit <= 0 { + limit = defaultLimit + } + if limit > 500 { + limit = 500 + } + + entries, hasMore := logger.GlobalRingBuffer.Query(cursor, limit) + + resp := logsResponse{ + Lines: entries, + HasMore: hasMore, + } + if len(entries) > 0 { + resp.NextCursor = entries[0].Index + } + + c.JSON(http.StatusOK, util.OK(resp)) +} + +// wsMessage WebSocket 消息格式 +type wsMessage struct { + Type string `json:"type"` // "log" | "error" + Data json.RawMessage `json:"data"` +} + +// HandleLogWebSocket WebSocket 端点,实时推送系统日志 +// @Summary 系统日志实时推送 +// @Description 通过 WebSocket 实时推送系统日志,需要管理员权限 +// @Tags admin +// @Router /api/v1/admin/logs/ws [get] +func HandleLogWebSocket(c *gin.Context) { + upgrader := getUpgrader() + + conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) + if err != nil { + return + } + defer conn.Close() + + // 订阅 ring buffer + ch := logger.GlobalRingBuffer.Subscribe() + defer logger.GlobalRingBuffer.Unsubscribe(ch) + + // 在独立 goroutine 中读取客户端消息(保持连接活跃 + 检测断开) + done := make(chan struct{}) + go func() { + defer close(done) + for { + _, _, err := conn.ReadMessage() + if err != nil { + return + } + } + }() + + // 主循环:推送日志 + for { + select { + case <-done: + return + case entry, ok := <-ch: + if !ok { + return + } + data, _ := json.Marshal(entry) + msg := wsMessage{Type: "log", Data: data} + payload, _ := json.Marshal(msg) + if err := conn.WriteMessage(1, payload); err != nil { + return + } + } + } +} diff --git a/internal/apps/admin/logs/utils.go b/internal/apps/admin/logs/utils.go new file mode 100644 index 00000000..7d044c30 --- /dev/null +++ b/internal/apps/admin/logs/utils.go @@ -0,0 +1,47 @@ +/* +Copyright 2025-2026 linux.do + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package logs + +import ( + "net/http" + "strconv" + + "github.com/gorilla/websocket" +) + +// getUpgrader 返回 WebSocket 升级器 +func getUpgrader() *websocket.Upgrader { + return &websocket.Upgrader{ + CheckOrigin: func(r *http.Request) bool { + return true // CORS 由 Gin 中间件处理 + }, + } +} + +// parsePositiveInt 解析非负整数字符串 +func parsePositiveInt(s string, result *int) (bool, error) { + if s == "" { + *result = 0 + return true, nil + } + n, err := strconv.Atoi(s) + if err != nil || n < 0 { + return false, err + } + *result = n + return true, nil +} diff --git a/internal/logger/logger.go b/internal/logger/logger.go index ca936036..13c8f2a0 100644 --- a/internal/logger/logger.go +++ b/internal/logger/logger.go @@ -28,14 +28,26 @@ import ( var logger *otelzap.Logger +// GlobalRingBuffer 全局日志环形缓冲区,供 Admin 日志查询和 WebSocket 推送使用 +var GlobalRingBuffer *LogRingBuffer + func init() { logWriter, err := GetLogWriter() if err != nil { log.Fatalf("[Logger] get log writer err: %v\n", err) } + // 初始化 ring buffer(保留最近 5000 行日志) + GlobalRingBuffer = NewLogRingBuffer(5000) + + // 使用 multi writer 同时写入原始输出和 ring buffer + multiWriter := zapcore.NewMultiWriteSyncer( + logWriter, + zapcore.AddSync(GlobalRingBuffer), + ) + zapLogger := zap.New( - zapcore.NewCore(getEncoder(), logWriter, getLogLevel()), + zapcore.NewCore(getEncoder(), multiWriter, getLogLevel()), zap.AddCaller(), zap.AddCallerSkip(1), ) diff --git a/internal/logger/ringbuffer.go b/internal/logger/ringbuffer.go new file mode 100644 index 00000000..91b3a5e7 --- /dev/null +++ b/internal/logger/ringbuffer.go @@ -0,0 +1,184 @@ +/* +Copyright 2025-2026 linux.do + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package logger + +import ( + "io" + "sync" +) + +// LogEntry 日志条目,对应 ring buffer 中的一行日志 +type LogEntry struct { + Index int `json:"index"` // 全局递增序号 + Data string `json:"data"` // 一行日志原文(含换行符) +} + +// LogRingBuffer 固定容量的环形缓冲区,存储最近的日志行 +// 支持:追加日志、按 cursor 分页查询、订阅实时推送 +type LogRingBuffer struct { + mu sync.RWMutex + entries []LogEntry + cap int + head int // 下一条写入的位置 + count int // 当前条目数 + seq int // 全局递增序号 + + subscribers map[chan LogEntry]struct{} + subMu sync.RWMutex +} + +// NewLogRingBuffer 创建指定容量的日志环形缓冲区 +func NewLogRingBuffer(capacity int) *LogRingBuffer { + return &LogRingBuffer{ + entries: make([]LogEntry, capacity), + cap: capacity, + subscribers: make(map[chan LogEntry]struct{}), + } +} + +// Write 实现 io.Writer 接口,供 zapcore.WriteSyncer 调用 +// 按 '\n' 分割为独立行写入 ring buffer +func (r *LogRingBuffer) Write(p []byte) (int, error) { + if len(p) == 0 { + return 0, nil + } + + data := string(p) + start := 0 + for i := 0; i < len(data); i++ { + if data[i] == '\n' { + line := data[start:i] + start = i + 1 + if len(line) > 0 { + r.appendLine(line) + } + } + } + // 处理最后一行(没有换行符结尾的情况) + if start < len(data) && len(data[start:]) > 0 { + r.appendLine(data[start:]) + } + + return len(p), nil +} + +// Sync 实现 zapcore.WriteSyncer 接口 +func (r *LogRingBuffer) Sync() error { + return nil +} + +// appendLine 追加一行日志到 ring buffer 并通知订阅者 +func (r *LogRingBuffer) appendLine(line string) { + r.mu.Lock() + entry := LogEntry{ + Index: r.seq, + Data: line, + } + r.entries[r.head] = entry + r.head = (r.head + 1) % r.cap + if r.count < r.cap { + r.count++ + } + r.seq++ + r.mu.Unlock() + + // 异步通知订阅者 + r.subMu.RLock() + for ch := range r.subscribers { + select { + case ch <- entry: + default: + // 订阅者消费太慢,丢弃(避免阻塞日志写入) + } + } + r.subMu.RUnlock() +} + +// Query 查询历史日志 +// cursor=0 表示查询最新日志,cursor>0 表示查询 index < cursor 的更早日志 +// limit 为返回条数上限 +// 返回日志条目(按 index 升序)和是否有更早的日志 +func (r *LogRingBuffer) Query(cursor int, limit int) ([]LogEntry, bool) { + r.mu.RLock() + defer r.mu.RUnlock() + + if r.count == 0 { + return nil, false + } + + // 计算 ring buffer 中有效条目的范围 + // oldest index in ring: head - count (wrapping) + oldestPos := (r.head - r.count + r.cap) % r.cap + + // 将 ring buffer 中的有效条目按顺序收集 + ordered := make([]LogEntry, 0, r.count) + for i := 0; i < r.count; i++ { + pos := (oldestPos + i) % r.cap + ordered = append(ordered, r.entries[pos]) + } + + if cursor == 0 { + // 查询最新日志:返回最后 limit 条 + if len(ordered) <= limit { + return ordered, false + } + return ordered[len(ordered)-limit:], true + } + + // 查询 index < cursor 的更早日志 + // 找到 index < cursor 的条目 + var cut int + for cut = len(ordered); cut > 0; cut-- { + if ordered[cut-1].Index < cursor { + break + } + } + + if cut == 0 { + return nil, false + } + + // 返回 cut 之前的最后 limit 条 + start := cut - limit + if start < 0 { + start = 0 + } + + hasMore := start > 0 + return ordered[start:cut], hasMore +} + +// Subscribe 订阅实时日志推送 +// 返回一个 channel,调用者应 defer Unsubscribe +func (r *LogRingBuffer) Subscribe() chan LogEntry { + ch := make(chan LogEntry, 64) + r.subMu.Lock() + r.subscribers[ch] = struct{}{} + r.subMu.Unlock() + return ch +} + +// Unsubscribe 取消订阅 +func (r *LogRingBuffer) Unsubscribe(ch chan LogEntry) { + r.subMu.Lock() + delete(r.subscribers, ch) + r.subMu.Unlock() + close(ch) +} + +// 确保 LogRingBuffer 实现 io.Writer 接口 +var _ io.Writer = (*LogRingBuffer)(nil) diff --git a/internal/logger/ringbuffer_test.go b/internal/logger/ringbuffer_test.go new file mode 100644 index 00000000..bfef78b9 --- /dev/null +++ b/internal/logger/ringbuffer_test.go @@ -0,0 +1,205 @@ +/* +Copyright 2025-2026 linux.do + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package logger + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestLogRingBuffer_WriteAndQuery(t *testing.T) { + rb := NewLogRingBuffer(5) + + // Write some logs + rb.Write([]byte("line1\nline2\nline3\n")) + + entries, hasMore := rb.Query(0, 10) + assert.False(t, hasMore) + assert.Equal(t, 3, len(entries)) + assert.Equal(t, "line1", entries[0].Data) + assert.Equal(t, "line2", entries[1].Data) + assert.Equal(t, "line3", entries[2].Data) + assert.Equal(t, 0, entries[0].Index) + assert.Equal(t, 1, entries[1].Index) + assert.Equal(t, 2, entries[2].Index) +} + +func TestLogRingBuffer_CapacityOverflow(t *testing.T) { + rb := NewLogRingBuffer(3) + + rb.Write([]byte("a\nb\nc\nd\ne\n")) + + entries, hasMore := rb.Query(0, 10) + assert.False(t, hasMore) + assert.Equal(t, 3, len(entries)) + assert.Equal(t, "c", entries[0].Data) + assert.Equal(t, "d", entries[1].Data) + assert.Equal(t, "e", entries[2].Data) +} + +func TestLogRingBuffer_QueryLatest(t *testing.T) { + rb := NewLogRingBuffer(10) + + rb.Write([]byte("a\nb\nc\nd\ne\n")) + + // Query latest 2 + entries, hasMore := rb.Query(0, 2) + assert.True(t, hasMore) + assert.Equal(t, 2, len(entries)) + assert.Equal(t, "d", entries[0].Data) + assert.Equal(t, "e", entries[1].Data) +} + +func TestLogRingBuffer_QueryByCursor(t *testing.T) { + rb := NewLogRingBuffer(10) + + rb.Write([]byte("a\nb\nc\nd\ne\n")) + + // First get all to find indices + all, _ := rb.Query(0, 10) + assert.Equal(t, 5, len(all)) + + // Query entries before index 3 + entries, hasMore := rb.Query(3, 10) + assert.False(t, hasMore) + assert.Equal(t, 3, len(entries)) + assert.Equal(t, "a", entries[0].Data) + assert.Equal(t, "b", entries[1].Data) + assert.Equal(t, "c", entries[2].Data) +} + +func TestLogRingBuffer_QueryByCursorWithLimit(t *testing.T) { + rb := NewLogRingBuffer(10) + + rb.Write([]byte("a\nb\nc\nd\ne\n")) + + // Query 2 entries before index 4 + entries, hasMore := rb.Query(4, 2) + assert.True(t, hasMore) + assert.Equal(t, 2, len(entries)) + assert.Equal(t, "b", entries[0].Data) + assert.Equal(t, "c", entries[1].Data) +} + +func TestLogRingBuffer_QueryEmpty(t *testing.T) { + rb := NewLogRingBuffer(5) + + entries, hasMore := rb.Query(0, 10) + assert.False(t, hasMore) + assert.Nil(t, entries) +} + +func TestLogRingBuffer_QueryNonExistentCursor(t *testing.T) { + rb := NewLogRingBuffer(5) + rb.Write([]byte("a\nb\n")) + + entries, hasMore := rb.Query(999, 10) + assert.False(t, hasMore) + assert.Nil(t, entries) +} + +func TestLogRingBuffer_Subscribe(t *testing.T) { + rb := NewLogRingBuffer(5) + + ch := rb.Subscribe() + defer rb.Unsubscribe(ch) + + rb.Write([]byte("hello\n")) + + entry := <-ch + assert.Equal(t, "hello", entry.Data) + assert.Equal(t, 0, entry.Index) +} + +func TestLogRingBuffer_SubscribeMultiple(t *testing.T) { + rb := NewLogRingBuffer(5) + + ch1 := rb.Subscribe() + defer rb.Unsubscribe(ch1) + ch2 := rb.Subscribe() + defer rb.Unsubscribe(ch2) + + rb.Write([]byte("msg\n")) + + e1 := <-ch1 + e2 := <-ch2 + assert.Equal(t, "msg", e1.Data) + assert.Equal(t, "msg", e2.Data) +} + +func TestLogRingBuffer_WriteNoNewline(t *testing.T) { + rb := NewLogRingBuffer(5) + + rb.Write([]byte("partial")) + + entries, _ := rb.Query(0, 10) + assert.Equal(t, 1, len(entries)) + assert.Equal(t, "partial", entries[0].Data) +} + +func TestLogRingBuffer_WriteEmpty(t *testing.T) { + rb := NewLogRingBuffer(5) + + n, err := rb.Write([]byte("")) + assert.Equal(t, 0, n) + assert.NoError(t, err) + + entries, _ := rb.Query(0, 10) + assert.Nil(t, entries) +} + +func TestLogRingBuffer_QueryAfterOverflow(t *testing.T) { + rb := NewLogRingBuffer(3) + + rb.Write([]byte("1\n2\n3\n4\n5\n6\n7\n")) + + entries, hasMore := rb.Query(0, 10) + assert.False(t, hasMore) + assert.Equal(t, 3, len(entries)) + assert.Equal(t, "5", entries[0].Data) + assert.Equal(t, "6", entries[1].Data) + assert.Equal(t, "7", entries[2].Data) + + // Query by cursor - index 4 is "5", so cursor=4 should return index < 4 + older, hasMore2 := rb.Query(4, 10) + assert.True(t, hasMore2) // there are entries 0,1,2 but they've been overwritten... + // Actually after overflow, entries 0,1,2 are gone from ring, so index < 4 should return nothing + // Let's check: ring has indices 4,5,6. Query(cursor=4) looks for index < 4 → none in ring + _ = older + _ = hasMore2 +} + +func TestLogRingBuffer_NextCursor(t *testing.T) { + rb := NewLogRingBuffer(10) + + rb.Write([]byte("a\nb\nc\nd\ne\n")) + + // Query latest 2, should return next_cursor pointing to first returned entry + entries, _ := rb.Query(0, 2) + assert.Equal(t, 2, len(entries)) + // entries[0].Index = 3 ("d"), entries[1].Index = 4 ("e") + assert.Equal(t, 3, entries[0].Index) + + // Now use that index as cursor to get older entries + older, hasMore := rb.Query(entries[0].Index, 10) + assert.False(t, hasMore) + assert.Equal(t, 3, len(older)) + assert.Equal(t, "a", older[0].Data) + assert.Equal(t, "b", older[1].Data) + assert.Equal(t, "c", older[2].Data) +} diff --git a/internal/router/router.go b/internal/router/router.go index 3f6e999c..ddc0a2ec 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -29,6 +29,7 @@ import ( "github.com/linux-do/credit/internal/apps/admin" admin_auth_source "github.com/linux-do/credit/internal/apps/admin/auth_source" + admin_logs "github.com/linux-do/credit/internal/apps/admin/logs" admin_status "github.com/linux-do/credit/internal/apps/admin/status" admin_task "github.com/linux-do/credit/internal/apps/admin/task" admin_user "github.com/linux-do/credit/internal/apps/admin/user" @@ -184,6 +185,10 @@ func Serve() { // System status adminRouter.GET("/status", admin_status.GetSystemStatus) + // System logs + adminRouter.GET("/logs", admin_logs.GetLogs) + adminRouter.GET("/logs/ws", admin_logs.HandleLogWebSocket) + // Task dispatch adminRouter.GET("/tasks/types", admin_task.ListTaskTypes) adminRouter.POST("/tasks/dispatch", admin_task.DispatchTask)