This commit is contained in:
ryan
2026-06-08 16:28:18 +08:00
parent 5c7c995a95
commit c6cc8c3305
14 changed files with 1459 additions and 2 deletions
+1 -1
View File
@@ -6,7 +6,7 @@
.pnpm-store/
# logs
logs
/logs/
*.log
# config
+511
View File
@@ -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/<module>/`)内部按照以下约定组织文件:
| 文件名 | 职责 |
|--------|------|
| `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.<Section>.<Field>` 访问(全局单例)。
- 不允许在业务代码中使用 `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 手写颜色、背景、阴影等样式
<Button
type="button"
size="sm"
className="bg-indigo-600 hover:bg-indigo-700 text-white shadow-md shadow-indigo-600/10 transition-colors"
>
<Plus className="mr-1.5 size-3.5" />
新增认证源
</Button>
```
**正确示例**:
```tsx
// ✅ 使用 variant 属性,让组件遵循系统主题
<Button
type="button"
size="sm"
variant="secondary"
>
<Plus className="mr-1.5 size-3.5" />
新增认证源
</Button>
```
> **原则**:组件的视觉表现由 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/<module>/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/<module>/ 中创建:
- 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/<module>/
├── 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) |
---
+8
View File
@@ -0,0 +1,8 @@
"use client"
import {SystemLogs} from "@/components/common/admin/system-logs"
/* 系统日志页面 */
export default function LogsPage() {
return <SystemLogs />
}
@@ -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<Error | null>(null)
const [logs, setLogs] = useState<LogEntry[]>([])
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<HTMLDivElement>(null)
const wsRef = useRef<WebSocket | null>(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<typeof setTimeout>
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 <LoadingStateWithBorder />
if (error) return <ErrorInline error={error} onRetry={() => loadHistory(0)} />
return (
<div className="flex flex-col h-full">
{/* Header */}
<div className="flex items-center justify-between pb-4 border-b border-border/50">
<div className="flex items-center gap-2">
<Terminal className="size-5 text-muted-foreground" />
<h1 className="text-lg font-semibold">系统日志</h1>
</div>
<div className="flex items-center gap-2">
<Badge variant={connected ? "secondary" : "destructive"}>
{connected ? "已连接" : "未连接"}
</Badge>
{connected && (
<Button variant="outline" size="sm" onClick={togglePause}>
{paused
? <><Play className="size-3.5 mr-1.5" />恢复</>
: <><Pause className="size-3.5 mr-1.5" />暂停</>
}
</Button>
)}
{!connected && (
<Button variant="outline" size="sm" onClick={reconnect}>
重连
</Button>
)}
</div>
</div>
{/* Log viewer — fixed height, scrollable */}
<div
ref={containerRef}
onScroll={handleScroll}
onWheel={handleWheel}
onTouchStart={handleTouchStart}
className="mt-3 h-[calc(100vh-220px)] overflow-y-auto overflow-x-hidden rounded-md border border-border/50 bg-[#0d1117] font-mono text-[13px] leading-5 relative"
>
{/* Load older logs */}
{hasMore && (
<div className="sticky top-0 z-10 flex justify-center py-1.5 bg-[#0d1117]/90 backdrop-blur-sm">
<Button
variant="ghost"
size="sm"
onClick={handleLoadMore}
disabled={loadingMore}
className="text-muted-foreground hover:text-foreground h-7 text-xs"
>
{loadingMore
? <><Loader2 className="size-3 mr-1.5 animate-spin" />加载中...</>
: <><ChevronUp className="size-3 mr-1.5" />加载更早日志</>
}
</Button>
</div>
)}
{/* Log lines */}
<div className="px-3 py-2">
{logs.length === 0 ? (
<div className="text-center text-gray-500 py-8">暂无日志</div>
) : (
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 (
<div
key={entry.index}
className={`${color} whitespace-pre-wrap break-all hover:bg-white/5`}
>
{entry.data}
</div>
)
})
)}
</div>
</div>
{/* Floating "back to latest" button */}
{!autoScroll && (
<div className="absolute bottom-24 left-1/2 -translate-x-1/2 z-20">
<Button
variant="outline"
size="sm"
onClick={scrollToBottom}
className="shadow-lg bg-background/80 backdrop-blur-sm"
>
<ArrowDown className="size-3.5 mr-1.5" />
回到最新
</Button>
</div>
)}
</div>
)
}
+2
View File
@@ -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 },
],
@@ -347,4 +347,19 @@ export class AdminService extends BaseService {
static async getSystemStatus(): Promise<SystemStatus> {
return this.get<SystemStatus>('/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 });
}
}
+1
View File
@@ -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
+2
View File
@@ -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=
+131
View File
@@ -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
}
}
}
}
+47
View File
@@ -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
}
+13 -1
View File
@@ -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),
)
+184
View File
@@ -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)
+205
View File
@@ -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)
}
+5
View File
@@ -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)