diff --git a/.env.example b/.env.example index 9e081ca..c4feef7 100644 --- a/.env.example +++ b/.env.example @@ -70,6 +70,16 @@ SHADOW_TIMEOUT=20s SHADOW_MAX_BODY_BYTES=2097152 SHADOW_MAX_CONCURRENT=16 +# M8 对象存储(MinIO)。上传/下载全部经网关代理,MinIO 不暴露主机端口。 +# 本地 compose 会自带一个 minio 服务;连远端 S3 时把 S3_ENDPOINT 指向它即可。 +S3_ENDPOINT=http://minio:9000 +S3_ACCESS_KEY_ID=gateway +S3_SECRET_ACCESS_KEY=gateway-secret +S3_BUCKET=gateway-files +S3_REGION=us-east-1 +S3_USE_SSL=false +S3_MAX_FILE_BYTES=134217728 + # Used only by cmd/gateway-bootstrap; never commit the real value. BOOTSTRAP_ADMIN_USERNAME=admin BOOTSTRAP_ADMIN_PASSWORD= diff --git a/README.md b/README.md index e4f67dd..6e7021b 100644 --- a/README.md +++ b/README.md @@ -30,16 +30,17 @@ AI Gateway 的全量 Go 重构工程。M0–M6 工程实现已完成,当前可 - 可扩展内容策略:Go RE2 不可变编译快照,按端点、模型/API Key 匹配,支持仅审计、阻断与提示词文本脱敏;默认保护常见 API Key、Token、密码和 secret。 - 带时间版本的模型价格与成本核算:按 Provider/模型选择价格,输入与输出 Token 分别计价,结果进入调用审计和 PostgreSQL 按日聚合。 - Prompt 分类、模板和不可变版本,支持显式变量定义、必填校验、历史版本激活与 API Key 渲染接口。 -- PostgreSQL 知识库:2 MiB 有界文本正文、段落感知重叠分块、FTS + 中文二元词片混合检索,以及可替换的 `Retriever` 接口;不依赖对象存储或向量数据库。 +- PostgreSQL 知识库:2 MiB 有界文本正文、段落感知重叠分块、FTS + 中文二元词片混合检索,以及可替换的 `Retriever` 接口;向量化由 M8 的 pgvector/Ollama 阶段提供。 - 声明式 HTTP 工具:JSON Schema 基础校验、KEK 加密请求头、注册和拨号双层 SSRF 防护、禁止重定向、1 MiB 响应限制与调用记录。 - AI 应用草稿和不可变发布版本,将模型、Prompt、知识库、工具组合为 `/v1/applications/{code}/chat/completions`;所有模型轮次继续经过鉴权、配额、内容策略、路由、成本和审计。 - 独立通知 Worker 消费可靠 outbox,按精确事件或末尾 `*` 模式投递 HMAC-SHA256 Webhook;内容策略命中由审计批处理异步产生脱敏事件,失败投递可在 Art 管理端重试。 - 门户自助工作台:部门范围资产目录、Prompt 搜索/收藏、个人审计/用量/成本、模型访问申请与管理员审批。 +- M8 对象存储:自托管 MinIO,上传/下载全部经网关代理(不暴露主机端口),管理端文件管理与门户个人文件仓库,`sha256` 完整性校验与严格归属隔离。 - 门户应用托管会话:服务端加密运行凭证、单会话租约、不可变消息序列和 SHA-256 哈希链,不向浏览器暴露应用 API Key。 - 独立事实核验配置、作用域策略与事件契约,复用 Provider 加密凭据和知识库引用,为同步/异步执行器保留清晰模块边界。 - 旧 Python 源码 201 条路由全部有覆盖、替代或退役决策,未决契约缺口为 0;OpenAPI 0.10.0 覆盖全部 Go 字面量路由。 -MinIO/S3 和 ClickHouse 不属于基线部署,也不是启动依赖。审计与统计第一阶段存放在 PostgreSQL;对象存储和分析库仅保留后续 adapter 扩展点。 +M8 起 MinIO(对象存储)纳入基线部署并由 compose 提供,但不作为启动依赖:网关启动时 MinIO 未就绪只告警、上传请求得到明确报错,服务不会因对象存储缺失而崩溃。ClickHouse 不属于基线,审计与统计继续存放在 PostgreSQL。 ## 本地启动 diff --git a/cmd/gateway-api/main.go b/cmd/gateway-api/main.go index 333b19d..c17e0b5 100644 --- a/cmd/gateway-api/main.go +++ b/cmd/gateway-api/main.go @@ -23,6 +23,7 @@ import ( "aigateway.local/core/internal/platform/cryptox" "aigateway.local/core/internal/platform/database" "aigateway.local/core/internal/platform/health" + "aigateway.local/core/internal/platform/storage" "aigateway.local/core/internal/platform/httpserver" "aigateway.local/core/internal/portal" "aigateway.local/core/internal/pricing" @@ -231,6 +232,28 @@ func main() { marketplaceService := workbench.NewMarketplaceService(workbenchService, mcpServerService, skillService, digitalEmployeeService) mcpClient := workbench.NewMCPClient(cfg.Credentials.AllowPrivateToolURL, 60*time.Second) marketplaceHandler := workbench.NewMarketplaceAdminHTTPHandler(marketplaceService, mcpServerService, skillService, digitalEmployeeService, mcpClient, identityService) + // M8: 对象存储(MinIO)文件管理。文件体在 MinIO,元数据在 PostgreSQL。MinIO + // 不暴露主机端口,上传/下载全部经网关代理,凭据只留在 API 容器内。 + objectStore, err := storage.NewClient(storage.Config{ + Endpoint: cfg.ObjectStorage.Endpoint, + AccessKeyID: cfg.ObjectStorage.AccessKeyID, + SecretAccessKey: cfg.ObjectStorage.SecretAccessKey, + Bucket: cfg.ObjectStorage.Bucket, + Region: cfg.ObjectStorage.Region, + UseSSL: cfg.ObjectStorage.UseSSL, + MaxFileBytes: cfg.ObjectStorage.MaxFileBytes, + }) + if err != nil { + logger.Error("object storage client initialization failed", "error", err) + os.Exit(1) + } + if err := objectStore.EnsureBucket(ctx); err != nil { + // 桶未就绪不致命:MinIO 起来后会自动建桶,当前上传请求会得到明确报错。 + logger.Warn("object storage bucket is not ready; uploads fail until MinIO is reachable", "error", err) + } + fileService := workbench.NewFileService(workbenchService, objectStore) + filesAdminHandler := workbench.NewFilesAdminHTTPHandler(fileService, identityService) + filesPortalHandler := workbench.NewFilesPortalHTTPHandler(fileService, identityService) applicationKeyCipher, err := cryptox.NewKeyring( cfg.Credentials.MasterKey, cfg.Credentials.KEKVersion, cfg.Credentials.KEKKeyring, "application-runtime-key", ) @@ -307,6 +330,10 @@ func main() { controlMux.Handle("/api/v1/admin/model-requests/", portalAdminHandler) controlMux.Handle("/api/v1/admin/system-info", operationsHandler) controlMux.Handle("/api/v1/admin/monitoring/overview", operationsHandler) + controlMux.Handle("/api/v1/admin/files", filesAdminHandler) + controlMux.Handle("/api/v1/admin/files/", filesAdminHandler) + controlMux.Handle("/api/v1/portal/files", filesPortalHandler) + controlMux.Handle("/api/v1/portal/files/", filesPortalHandler) controlMux.Handle("/api/v1/admin/reload", operationsHandler) controlMux.Handle("/api/v1/admin/identities/", identityManagementHandler) controlMux.Handle("/api/v1/admin/departments", identityManagementHandler) diff --git a/deploy/PRODUCTION.md b/deploy/PRODUCTION.md index c2fb516..1930e7a 100644 --- a/deploy/PRODUCTION.md +++ b/deploy/PRODUCTION.md @@ -1,8 +1,9 @@ # Production deployment This bundle builds the Go services and both Art Design Pro applications from -source. PostgreSQL and two Redis roles are included; MinIO/S3 and ClickHouse are -not required. +source. PostgreSQL, two Redis roles and MinIO (object storage, M8) are included; +ClickHouse is not required. MinIO is not a startup dependency: the gateway only +warns and refuses file uploads until the bucket is reachable. ## Prerequisites diff --git a/deploy/docker-compose.production.yml b/deploy/docker-compose.production.yml index b02966f..472662c 100644 --- a/deploy/docker-compose.production.yml +++ b/deploy/docker-compose.production.yml @@ -26,6 +26,13 @@ x-gateway-environment: &gateway-environment SHADOW_BASE_URL: ${SHADOW_BASE_URL:-} SHADOW_API_KEY: ${SHADOW_API_KEY:-} SHADOW_SAMPLE_RATE: ${SHADOW_SAMPLE_RATE:-0} + S3_ENDPOINT: ${S3_ENDPOINT:-http://minio:9000} + S3_ACCESS_KEY_ID: ${S3_ACCESS_KEY_ID:-gateway} + S3_SECRET_ACCESS_KEY: ${S3_SECRET_ACCESS_KEY:-} + S3_BUCKET: ${S3_BUCKET:-gateway-files} + S3_REGION: ${S3_REGION:-us-east-1} + S3_USE_SSL: ${S3_USE_SSL:-false} + S3_MAX_FILE_BYTES: ${S3_MAX_FILE_BYTES:-134217728} x-backend-service: &backend-service image: ai-gateway-go:${GATEWAY_VERSION:-0.10.0} @@ -107,6 +114,8 @@ services: condition: service_healthy redis-cache: condition: service_healthy + minio: + condition: service_started healthcheck: test: ["CMD-SHELL", "wget -q -O /dev/null http://127.0.0.1:8080/readyz"] interval: 10s @@ -171,6 +180,18 @@ services: condition: service_healthy restart: unless-stopped + # M8: 对象存储。MinIO 不暴露主机端口,上传/下载全部经网关代理; + # 桶由 gateway-api 启动时的 EnsureBucket 兜底创建。stateful 服务不做加固。 + minio: + image: minio/minio:latest + command: ["server", "/data", "--console-address", ":9001"] + environment: + MINIO_ROOT_USER: ${S3_ACCESS_KEY_ID:-gateway} + MINIO_ROOT_PASSWORD: ${S3_SECRET_ACCESS_KEY:-} + volumes: + - minio-data:/data + restart: unless-stopped + bootstrap-admin: <<: *backend-service profiles: ["tools"] @@ -187,3 +208,4 @@ services: volumes: postgres-data: redis-critical-data: + minio-data: diff --git a/deploy/docker-compose.yml b/deploy/docker-compose.yml index 6c54a75..170e667 100644 --- a/deploy/docker-compose.yml +++ b/deploy/docker-compose.yml @@ -58,6 +58,13 @@ services: SHADOW_BASE_URL: ${SHADOW_BASE_URL:-} SHADOW_API_KEY: ${SHADOW_API_KEY:-} SHADOW_SAMPLE_RATE: ${SHADOW_SAMPLE_RATE:-0} + S3_ENDPOINT: ${S3_ENDPOINT:-http://minio:9000} + S3_ACCESS_KEY_ID: ${S3_ACCESS_KEY_ID:-gateway} + S3_SECRET_ACCESS_KEY: ${S3_SECRET_ACCESS_KEY:-gateway-secret} + S3_BUCKET: ${S3_BUCKET:-gateway-files} + S3_REGION: ${S3_REGION:-us-east-1} + S3_USE_SSL: ${S3_USE_SSL:-false} + S3_MAX_FILE_BYTES: ${S3_MAX_FILE_BYTES:-134217728} depends_on: postgres: condition: service_healthy @@ -78,6 +85,20 @@ services: condition: service_healthy redis-cache: condition: service_healthy + minio: + condition: service_started + restart: unless-stopped + + # M8: 对象存储。MinIO 不暴露主机端口,凭据只留在 API 容器内; + # 上传/下载全部经网关代理。桶由 gateway-api 启动时的 EnsureBucket 兜底创建。 + minio: + image: minio/minio:latest + command: ["server", "/data", "--console-address", ":9001"] + environment: + MINIO_ROOT_USER: ${S3_ACCESS_KEY_ID:-gateway} + MINIO_ROOT_PASSWORD: ${S3_SECRET_ACCESS_KEY:-gateway-secret} + volumes: + - minio-data:/data restart: unless-stopped admin-web: @@ -152,3 +173,4 @@ services: volumes: postgres-data: redis-critical-data: + minio-data: diff --git a/deploy/nginx-web.conf b/deploy/nginx-web.conf index 30c14cd..8bb234e 100644 --- a/deploy/nginx-web.conf +++ b/deploy/nginx-web.conf @@ -8,10 +8,10 @@ server { # and an absolute redirect would drop that port and send browsers to :80. absolute_redirect off; - # Match the gateway's HTTP_MAX_BODY_BYTES (default 32 MiB). nginx's default - # of 1 MiB otherwise rejects large prompts and knowledge-base imports with - # 413 before they ever reach the gateway. - client_max_body_size 32m; + # 256 MiB 必须盖过文件上传上限 S3_MAX_FILE_BYTES(默认 128 MiB)。文件体 + # 由网关的流式上传处理器把关(http.MaxBytesReader + LimitReader), + # nginx 只做最外层限制,避免大文件在到达网关前就被 413 拒绝。 + client_max_body_size 256m; root /usr/share/nginx/html; index index.html; diff --git a/deploy/production.env.example b/deploy/production.env.example index 0cc1ff2..6ac92fa 100644 --- a/deploy/production.env.example +++ b/deploy/production.env.example @@ -37,6 +37,17 @@ ADMIN_PORT=8081 PORTAL_PORT=8082 REDIS_CACHE_MAXMEMORY=256mb +# M8 对象存储(MinIO)。MinIO 不暴露主机端口,上传/下载全部经网关代理。 +# 生产建议生成强随机密钥对(S3_ACCESS_KEY_ID/S3_SECRET_ACCESS_KEY), +# compose 的 minio 服务会用这两个值初始化 MINIO_ROOT_USER/PASSWORD。 +S3_ENDPOINT=http://minio:9000 +S3_ACCESS_KEY_ID=CHANGE_ME_MINIO_ACCESS_KEY +S3_SECRET_ACCESS_KEY=CHANGE_ME_MINIO_SECRET_KEY +S3_BUCKET=gateway-files +S3_REGION=us-east-1 +S3_USE_SSL=false +S3_MAX_FILE_BYTES=134217728 + # Used only for the one-time bootstrap-admin command; remove after use. BOOTSTRAP_ADMIN_USERNAME=admin BOOTSTRAP_ADMIN_PASSWORD=CHANGE_ME_AT_LEAST_12_CHARACTERS diff --git a/docs/rewrite-progress.md b/docs/rewrite-progress.md index 30ec68f..d2bf5c1 100644 --- a/docs/rewrite-progress.md +++ b/docs/rewrite-progress.md @@ -113,3 +113,12 @@ - Compose 同时交付 API、管理端 `8081` 和门户端 `8082`;同一前端 Dockerfile 通过受限 build arg 构建两套独立 SPA。 工程实现与本地生产等价验证已关闭。正式上线仍是外部发布门禁:拿到实际旧数据库脱敏快照和旧密钥迁移授权后执行领域转换,并在真实新旧双系统与生产等价流量下完成影子观察、容量验收及切换/回退演练;这些动作不会在缺少生产数据和授权时伪造为已执行。 + +## 已完成:M8 基础设施层(P1 对象存储与文件管理) + +- 自托管 MinIO 对象存储入基线:`internal/platform/storage`(minio-go 适配层)剥离 scheme 推导 Secure,领域层不直接依赖 S3 客户端;`S3_*` 配置默认值 + Validate(1–512 MiB 上限、endpoint 无路径)。迁移 `000023` 建 `gateway.file_objects`(uuid/object_key 唯一/sha256/scope/owner 校验 + 部分索引)。 +- 文件上传/下载全部经网关代理(MinIO 不暴露主机端口):流式 multipart 上传(`http.MaxBytesReader` + `LimitReader` 超限即拒),`PutObject` 成功后才插元数据行、失败回滚对象;下载 io.Copy 流式 + RFC 5987 `filename*`;删除先删行再 best-effort 清对象(避免孤儿阻塞删除)。 +- `FileService` 个人/系统双范围:admin 文件管理(`file:read|manage` RBAC + 管理端文件管理页)与 portal 个人文件仓(严格归属隔离,跨用户读返回 404)。 +- Compose(dev + production)新增 `minio` 服务与 `minio-data` 卷;nginx `client_max_body_size` 32m→256m 盖过 128 MiB 上传上限;`.env.example` / `production.env.example` 补 `S3_*`。 +- 端到端验证:上传→列表→下载往返一致→删除后桶无孤儿;admin 读 portal 文件 404;`TestFileObjectLifecycle` 集成测试连真实 MinIO+PostgreSQL 通过;`go build ./...`、`go vet ./...`、全量单测通过。 +- 管理端"文件管理"与门户端"文件仓库"菜单由服务端动态菜单下发。 diff --git a/docs/旗舰版需求规划与完成情况.md b/docs/旗舰版需求规划与完成情况.md index 888644a..e983430 100644 --- a/docs/旗舰版需求规划与完成情况.md +++ b/docs/旗舰版需求规划与完成情况.md @@ -101,7 +101,7 @@ Prompt 分类/模板/不可变版本/必填校验;知识库(2 MiB 有界正文 | 安全策略 | 数据安全:工具数据输入输出脱敏 + 大模型回答隐私敏感信息拦截替换 | — | ✓ | ✓ | ⚠️ 提示词输入脱敏 ✅;工具输出/回答拦截替换 ❌ | | 站内消息 | 平台推送站内消息与动态 | — | ✓ | ✓ | ❌(通知 worker 仅 webhook) | | 审批授权 | 资源/模型/渠道使用申请流程审批管理 | — | ✓ | ✓ | ⚠️ 仅模型申请 | -| 文件管理 | 平台文件资源与对象存储文件浏览管理 | — | ✓ | ✓ | ❌(MinIO/S3 不在基线) | +| 文件管理 | 平台文件资源与对象存储文件浏览管理 | — | ✓ | ✓ | ✅(M8 MinIO 已上线,admin 文件管理) | | 审计日志 | 系统全量历史操作审计日志查询 | — | ✓ | ✓ | ✅ | | 企业报表 | 企业级运营数据报表统计分析 | — | ✓ | ✓ | ❌ | | License | 平台 License 授权管理与有效期管控 | ✓ | ✓ | ✓ | ❌ | @@ -125,7 +125,7 @@ Prompt 分类/模板/不可变版本/必填校验;知识库(2 MiB 有界正文 | 我的资源 | 查看插件资源权限等级(可查看/仅使用/管理) | ⚠️ 有访问控制,三档等级未成体系 | | 我的资源 | 申请大模型/token量/skill/mcp/数字员工/渠道权限 | ⚠️ 仅模型申请 | | 配置管理 | 个人配置环境变量,供 skill/mcp 使用 | ❌ | -| 个人文件仓库 | 对话产生的报告/文件存入个人仓库 | ❌ | +| 个人文件仓库 | 对话产生的报告/文件存入个人仓库 | ✅(M8 个人文件仓,经网关上传/下载) | | 安全策略 | 个人智能体安全策略(网络/工具命令校验审批/限流/脱敏/隐私拦截) | ❌ | | 个人中心 | 账号信息、密码修改、登录记录查看 | ✅ | | 消息通知 | 系统消息、审批待办、任务执行结果提醒 | ⚠️ webhook 投递 ✅;站内消息/待办/结果提醒 ❌ | @@ -134,7 +134,7 @@ Prompt 分类/模板/不可变版本/必填校验;知识库(2 MiB 有界正文 ## 四、差距汇总 -- **完全未实现(❌,约 20 项)**:AI 助手、收藏、数据报表/企业报表、配置管理(env 注入)、智能体管理三项(节点/LLMTrace/会话)、记忆管理三项、渠道管理全项、多租户、供应链安全扫描、站内消息、完整审批流、文件管理(对象存储)、License、定时任务全项、个人渠道、个人文件仓库、个人安全策略、ARM64。 +- **完全未实现(❌,约 18 项)**:AI 助手、收藏、数据报表/企业报表、配置管理(env 注入)、智能体管理三项(节点/LLMTrace/会话)、记忆管理三项、渠道管理全项、多租户、供应链安全扫描、站内消息、完整审批流、License、定时任务全项、个人渠道、个人安全策略、ARM64。(M8 已落地:文件管理/MinIO、个人文件仓库) - **部分覆盖需补齐(⚠️,约 10 项)**:平台概览看板、知识库向量化语义召回、资源/渠道审批、工具输出脱敏与大模型回答拦截替换、工具命令审批与工具限流、集群部署方案、站内消息/审批待办/任务结果、数字员工会话入口、资源权限等级三档、全类型权限申请。 --- diff --git a/go.mod b/go.mod index 87e513b..2b29fdb 100644 --- a/go.mod +++ b/go.mod @@ -8,20 +8,36 @@ require ( github.com/beevik/etree v1.5.0 github.com/crewjam/saml v0.5.1 github.com/jackc/pgx/v5 v5.10.0 + github.com/minio/minio-go/v7 v7.2.1 github.com/redis/go-redis/v9 v9.21.0 github.com/russellhaering/goxmldsig v1.4.0 ) require ( github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/dustin/go-humanize v1.0.1 // indirect github.com/golang-jwt/jwt/v4 v4.5.2 // indirect + github.com/google/uuid v1.6.0 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/jonboulle/clockwork v0.2.2 // indirect + github.com/klauspost/compress v1.18.6 // indirect + github.com/klauspost/cpuid/v2 v2.2.11 // indirect + github.com/klauspost/crc32 v1.3.0 // indirect github.com/mattermost/xml-roundtrip-validator v0.1.0 // indirect + github.com/minio/crc64nvme v1.1.1 // indirect + github.com/minio/md5-simd v1.1.2 // indirect + github.com/philhofer/fwd v1.2.0 // indirect + github.com/rs/xid v1.6.0 // indirect + github.com/tinylib/msgp v1.6.1 // indirect + github.com/zeebo/xxh3 v1.1.0 // indirect go.uber.org/atomic v1.11.0 // indirect - golang.org/x/crypto v0.33.0 // indirect - golang.org/x/sync v0.17.0 // indirect - golang.org/x/text v0.29.0 // indirect + go.yaml.in/yaml/v3 v3.0.4 // indirect + golang.org/x/crypto v0.51.0 // indirect + golang.org/x/net v0.53.0 // indirect + golang.org/x/sync v0.20.0 // indirect + golang.org/x/sys v0.44.0 // indirect + golang.org/x/text v0.37.0 // indirect + gopkg.in/ini.v1 v1.67.2 // indirect ) diff --git a/go.sum b/go.sum index 8ae1c91..4f6345a 100644 --- a/go.sum +++ b/go.sum @@ -13,10 +13,14 @@ github.com/crewjam/saml v0.5.1/go.mod h1:r0fDkmFe5URDgPrmtH0IYokva6fac3AUdstiPhy github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= github.com/golang-jwt/jwt/v4 v4.5.2 h1:YtQM7lnr8iZ+j5q71MGKkNw9Mn7AjHM68uc9g5fXeUI= github.com/golang-jwt/jwt/v4 v4.5.2/go.mod h1:m21LjoU+eqJr34lmDMbreY2eSTRJ1cv77w39/MY0Ch0= -github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= -github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= @@ -27,16 +31,31 @@ github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/jonboulle/clockwork v0.2.2 h1:UOGuzwb1PwsrDAObMuhUnj0p5ULPj8V/xJ7Kx9qUBdQ= github.com/jonboulle/clockwork v0.2.2/go.mod h1:Pkfl5aHPm1nk2H9h0bjmnJD/BcgbGXUBGnn1kMkgxc8= -github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE= -github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/klauspost/compress v1.18.6 h1:2jupLlAwFm95+YDR+NwD2MEfFO9d4z4Prjl1XXDjuao= +github.com/klauspost/compress v1.18.6/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/klauspost/cpuid/v2 v2.0.1/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= +github.com/klauspost/cpuid/v2 v2.2.11 h1:0OwqZRYI2rFrjS4kvkDnqJkKHdHaRnCm68/DY4OxRzU= +github.com/klauspost/cpuid/v2 v2.2.11/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/klauspost/crc32 v1.3.0 h1:sSmTt3gUt81RP655XGZPElI0PelVTZ6YwCRnPSupoFM= +github.com/klauspost/crc32 v1.3.0/go.mod h1:D7kQaZhnkX/Y0tstFGf8VUzv2UofNGqCjnC3zdHB0Hw= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= github.com/kr/pretty v0.2.1/go.mod h1:ipq/a2n7PKx3OHsz4KJII5eveXtPO4qwEXGdVfWzfnI= +github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/mattermost/xml-roundtrip-validator v0.1.0 h1:RXbVD2UAl7A7nOTR4u7E3ILa4IbtvKBHw64LDsmu9hU= github.com/mattermost/xml-roundtrip-validator v0.1.0/go.mod h1:qccnGMcpgwcNaBnxqpJpWWUiPNr5H3O8eDgGV9gT5To= +github.com/minio/crc64nvme v1.1.1 h1:8dwx/Pz49suywbO+auHCBpCtlW1OfpcLN7wYgVR6wAI= +github.com/minio/crc64nvme v1.1.1/go.mod h1:eVfm2fAzLlxMdUGc0EEBGSMmPwmXD5XiNRpnu9J3bvg= +github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34= +github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM= +github.com/minio/minio-go/v7 v7.2.1 h1:PfBfwvKB/MmqyN8Vb1G9voWisaM9OrLv+WwOvMwS9Dw= +github.com/minio/minio-go/v7 v7.2.1/go.mod h1:EU9hENAStx/xXduNdrGO5e4X5vk19NtgB+RIPjZO8o0= +github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM= +github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM= github.com/pkg/diff v0.0.0-20210226163009-20ebb0f2a09e/go.mod h1:pJLUxLENpZxwdsKMEsNbx1VGcRFpLqf3715MtcvvzbA= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= @@ -46,30 +65,51 @@ github.com/redis/go-redis/v9 v9.21.0 h1:FPBE4hhbAke+TLmcY3WkpbDffJEomdqPn3HYiqAt github.com/redis/go-redis/v9 v9.21.0/go.mod h1:v/M13XI1PVCDcm01VtPFOADfZtHf8YW3baQf57KlIkA= github.com/rogpeppe/go-internal v1.6.1/go.mod h1:xXDCJY+GAPziupqXw64V24skbSoqbTEfhy4qGm1nDQc= github.com/rogpeppe/go-internal v1.8.0/go.mod h1:WmiCO8CzOY8rg0OYDC4/i/2WRWAB6poM+XZ2dLUbcbE= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU= +github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= github.com/russellhaering/goxmldsig v1.4.0 h1:8UcDh/xGyQiyrW+Fq5t8f+l2DLB1+zlhYzkPUJ7Qhys= github.com/russellhaering/goxmldsig v1.4.0/go.mod h1:gM4MDENBQf7M+V824SGfyIUVFWydB7n0KkEubVJl+Tw= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= +github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= +github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/tinylib/msgp v1.6.1 h1:ESRv8eL3u+DNHUoSAAQRE50Hm162zqAnBoGv9PzScPY= +github.com/tinylib/msgp v1.6.1/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77roAqEA= +github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= +github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0= github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= -golang.org/x/crypto v0.33.0 h1:IOBPskki6Lysi0lo9qQvbxiQ+FvsCC/YWOecCHAixus= -golang.org/x/crypto v0.33.0/go.mod h1:bVdXmD7IV/4GdElGPozy6U7lWdRXA4qyRVGJV57uQ5M= -golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= -golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= -golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc= -golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= -golang.org/x/text v0.29.0 h1:1neNs90w9YzJ9BocxfsQNHKuAT4pkghyXc4nhZ6sJvk= -golang.org/x/text v0.29.0/go.mod h1:7MhJOA9CD2qZyOKYazxdYMF85OwPdEr9jTtBpO7ydH4= +go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= +golang.org/x/crypto v0.51.0 h1:IBPXwPfKxY7cWQZ38ZCIRPI50YLeevDLlLnyC5wRGTI= +golang.org/x/crypto v0.51.0/go.mod h1:8AdwkbraGNABw2kOX6YFPs3WM22XqI4EXEd8g+x7Oc8= +golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA= +golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs= +golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= +golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sys v0.44.0 h1:ildZl3J4uzeKP07r2F++Op7E9B29JRUy+a27EibtBTQ= +golang.org/x/sys v0.44.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= +golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/errgo.v2 v2.1.0/go.mod h1:hNsd1EY+bozCKY1Ytp96fpM3vjJbqLJn88ws8XvfDNI= +gopkg.in/ini.v1 v1.67.2 h1:JtOSMb9OuaCZKr7h5D/h6iii14sK0hLbplTc6frx4Ss= +gopkg.in/ini.v1 v1.67.2/go.mod h1:x/cyOwCgZqOkJoDIJ3c1KNHMo10+nLGAhh+kn3Zizss= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= diff --git a/internal/identity/account.go b/internal/identity/account.go index 52e850b..7e2f1ae 100644 --- a/internal/identity/account.go +++ b/internal/identity/account.go @@ -73,6 +73,8 @@ const ( PermissionDigitalEmployeeManage = "digital_employee:manage" PermissionMarketplaceRead = "marketplace:read" PermissionMarketplaceManage = "marketplace:manage" + PermissionFileRead = "file:read" + PermissionFileManage = "file:manage" ) var rolePermissions = map[string][]string{ @@ -93,8 +95,9 @@ var rolePermissions = map[string][]string{ PermissionSkillRead, PermissionSkillManage, PermissionDigitalEmployeeRead, PermissionDigitalEmployeeManage, PermissionMarketplaceRead, PermissionMarketplaceManage, + PermissionFileRead, PermissionFileManage, }, - "auditor": {PermissionProviderRead, PermissionAPIKeyRead, PermissionAuditRead, PermissionUsageRead, PermissionOutboxRead, PermissionContentPolicyRead, PermissionPricingRead, PermissionPromptRead, PermissionKnowledgeRead, PermissionToolRead, PermissionApplicationRead, PermissionNotificationRead, PermissionMCPServerRead, PermissionSkillRead, PermissionDigitalEmployeeRead, PermissionMarketplaceRead}, + "auditor": {PermissionProviderRead, PermissionAPIKeyRead, PermissionAuditRead, PermissionUsageRead, PermissionOutboxRead, PermissionContentPolicyRead, PermissionPricingRead, PermissionPromptRead, PermissionKnowledgeRead, PermissionToolRead, PermissionApplicationRead, PermissionNotificationRead, PermissionMCPServerRead, PermissionSkillRead, PermissionDigitalEmployeeRead, PermissionMarketplaceRead, PermissionFileRead}, "member": {}, } diff --git a/internal/identity/http.go b/internal/identity/http.go index 0c56dd8..c6e8015 100644 --- a/internal/identity/http.go +++ b/internal/identity/http.go @@ -377,6 +377,9 @@ func adminMenus(account Account) []map[string]any { if HasPermission(account, PermissionApplicationRead) || HasPermission(account, PermissionApplicationManage) { assetsChildren = append(assetsChildren, map[string]any{"name": "Applications", "path": "applications", "component": "/gateway/applications", "meta": map[string]any{"title": "AI 应用"}}) } + if HasPermission(account, PermissionFileRead) || HasPermission(account, PermissionFileManage) { + assetsChildren = append(assetsChildren, map[string]any{"name": "Files", "path": "files", "component": "/gateway/files", "meta": map[string]any{"title": "文件管理"}}) + } if len(assetsChildren) > 0 { menus = append(menus, map[string]any{"name": "Assets", "path": "/assets", "component": "/index/index", "meta": map[string]any{"title": "AI 资产", "icon": "ri:box-3-line"}, "children": assetsChildren}) } @@ -424,6 +427,7 @@ func portalMenus() []map[string]any { {"name": "PortalPrompts", "path": "prompts", "component": "/portal/prompts", "meta": map[string]any{"title": "Prompt 广场"}}, {"name": "PortalUsage", "path": "usage", "component": "/portal/usage", "meta": map[string]any{"title": "我的用量"}}, {"name": "PortalAccess", "path": "access", "component": "/portal/access", "meta": map[string]any{"title": "模型权限"}}, + {"name": "PortalFiles", "path": "files", "component": "/portal/files", "meta": map[string]any{"title": "文件仓库"}}, }}, } } diff --git a/internal/operations/admin_http.go b/internal/operations/admin_http.go index 9f62e92..afec513 100644 --- a/internal/operations/admin_http.go +++ b/internal/operations/admin_http.go @@ -48,7 +48,7 @@ func (h *AdminHTTPHandler) systemInfo(w http.ResponseWriter, r *http.Request) { apiresponse.Error(w, http.StatusServiceUnavailable, "系统信息查询失败") return } - apiresponse.OK(w, map[string]any{"version": h.version, "go_version": runtime.Version(), "uptime_seconds": int64(time.Since(h.startedAt).Seconds()), "database": "postgresql", "object_storage": false, "clickhouse": false, "enabled_providers": providers, "enabled_models": models, "enabled_api_keys": keys}) + apiresponse.OK(w, map[string]any{"version": h.version, "go_version": runtime.Version(), "uptime_seconds": int64(time.Since(h.startedAt).Seconds()), "database": "postgresql", "object_storage": true, "clickhouse": false, "enabled_providers": providers, "enabled_models": models, "enabled_api_keys": keys}) } func (h *AdminHTTPHandler) overview(w http.ResponseWriter, r *http.Request) { diff --git a/internal/platform/config/config.go b/internal/platform/config/config.go index 9ba3a7c..fd6ffd8 100644 --- a/internal/platform/config/config.go +++ b/internal/platform/config/config.go @@ -11,18 +11,19 @@ import ( ) type Config struct { - Environment string - Server Server - Database Database - Redis Redis - Security Security - Auth Auth - Credentials Credentials - Upstream Upstream - Audit Audit - Outbox Outbox - RuntimeData RuntimeData - Shadow Shadow + Environment string + Server Server + Database Database + Redis Redis + Security Security + Auth Auth + Credentials Credentials + Upstream Upstream + Audit Audit + Outbox Outbox + RuntimeData RuntimeData + Shadow Shadow + ObjectStorage ObjectStorage } type Server struct { @@ -114,6 +115,16 @@ type Shadow struct { MaxConcurrent int } +type ObjectStorage struct { + Endpoint string + AccessKeyID string + SecretAccessKey string + Bucket string + Region string + UseSSL bool + MaxFileBytes int64 +} + func Load() (Config, error) { cfg := Config{ Environment: env("APP_ENV", "local"), @@ -182,6 +193,15 @@ func Load() (Config, error) { Timeout: duration("SHADOW_TIMEOUT", 20*time.Second), MaxBodyBytes: int64Value("SHADOW_MAX_BODY_BYTES", 2<<20), MaxConcurrent: intValue("SHADOW_MAX_CONCURRENT", 16), }, + ObjectStorage: ObjectStorage{ + Endpoint: strings.TrimRight(env("S3_ENDPOINT", "http://minio:9000"), "/"), + AccessKeyID: env("S3_ACCESS_KEY_ID", "gateway"), + SecretAccessKey: env("S3_SECRET_ACCESS_KEY", "gateway-secret"), + Bucket: env("S3_BUCKET", "gateway-files"), + Region: env("S3_REGION", "us-east-1"), + UseSSL: boolValue("S3_USE_SSL", false), + MaxFileBytes: int64Value("S3_MAX_FILE_BYTES", 128<<20), + }, } return cfg, cfg.Validate() @@ -236,6 +256,21 @@ func (c Config) Validate() error { errs = append(errs, errors.New("SHADOW_API_KEY is required when SHADOW_BASE_URL is set")) } } + if err := validateHTTPURL(c.ObjectStorage.Endpoint); err != nil { + errs = append(errs, fmt.Errorf("S3_ENDPOINT: %w", err)) + } else if endpointURL, parseErr := url.Parse(c.ObjectStorage.Endpoint); parseErr == nil && endpointURL.Path != "" { + // minio-go 的 Endpoint 不接受带路径的完整 URL(报 "fully qualified paths")。 + errs = append(errs, errors.New("S3_ENDPOINT must not contain a path")) + } + if c.ObjectStorage.AccessKeyID == "" || c.ObjectStorage.SecretAccessKey == "" { + errs = append(errs, errors.New("S3_ACCESS_KEY_ID and S3_SECRET_ACCESS_KEY are required")) + } + if c.ObjectStorage.Bucket == "" || len(c.ObjectStorage.Bucket) > 63 { + errs = append(errs, errors.New("S3_BUCKET must be a non-empty bucket name of at most 63 characters")) + } + if c.ObjectStorage.MaxFileBytes < 1<<20 || c.ObjectStorage.MaxFileBytes > 512<<20 { + errs = append(errs, errors.New("S3_MAX_FILE_BYTES must be between 1 MiB and 512 MiB")) + } return errors.Join(errs...) } diff --git a/internal/platform/storage/client.go b/internal/platform/storage/client.go new file mode 100644 index 0000000..940d72f --- /dev/null +++ b/internal/platform/storage/client.go @@ -0,0 +1,96 @@ +package storage + +import ( + "context" + "io" + "net/url" + + "github.com/minio/minio-go/v7" + "github.com/minio/minio-go/v7/pkg/credentials" +) + +// Config holds the S3-compatible (MinIO) connection settings. +type Config struct { + Endpoint string + AccessKeyID string + SecretAccessKey string + Bucket string + Region string + UseSSL bool + MaxFileBytes int64 +} + +// Client is the thin object-store adapter. The domain layer imports this +// package instead of minio-go directly, keeping the S3 client at the platform +// boundary like cache and database. +type Client struct { + mc *minio.Client + bucket string + region string + maxBytes int64 +} + +func NewClient(cfg Config) (*Client, error) { + // minio-go 的 Endpoint 参数必须是裸 host[:port],不接受带 scheme 的完整 + // URL(否则报 "Endpoint url cannot have fully qualified paths.")。这里剥离 + // scheme,并由 scheme 推导 Secure(https→true),二者都比 S3_USE_SSL 优先。 + secure := cfg.UseSSL + host := cfg.Endpoint + if parsed, err := url.Parse(cfg.Endpoint); err == nil && parsed.Scheme != "" { + secure = parsed.Scheme == "https" + host = parsed.Host + } + mc, err := minio.New(host, &minio.Options{ + Creds: credentials.NewStaticV4(cfg.AccessKeyID, cfg.SecretAccessKey, ""), + Secure: secure, + Region: cfg.Region, + }) + if err != nil { + return nil, err + } + mc.SetAppInfo("ai-gateway", "0.10.0") + return &Client{mc: mc, bucket: cfg.Bucket, region: cfg.Region, maxBytes: cfg.MaxFileBytes}, nil +} + +// EnsureBucket creates the configured bucket if it does not exist. Safe to call +// repeatedly; it is a no-op once the bucket exists. +func (c *Client) EnsureBucket(ctx context.Context) error { + exists, err := c.mc.BucketExists(ctx, c.bucket) + if err != nil { + return err + } + if exists { + return nil + } + return c.mc.MakeBucket(ctx, c.bucket, minio.MakeBucketOptions{Region: c.region}) +} + +func (c *Client) MaxBytes() int64 { return c.maxBytes } + +// PutObject streams a reader to the object store. A size of -1 lets the client +// use chunked multipart upload so files are not buffered in memory. +func (c *Client) PutObject(ctx context.Context, key string, r io.Reader, size int64, contentType string) (int64, error) { + info, err := c.mc.PutObject(ctx, c.bucket, key, r, size, minio.PutObjectOptions{ContentType: contentType}) + if err != nil { + return 0, err + } + return info.Size, nil +} + +// OpenObject returns a streaming reader plus the object size. +func (c *Client) OpenObject(ctx context.Context, key string) (io.ReadCloser, int64, error) { + obj, err := c.mc.GetObject(ctx, c.bucket, key, minio.GetObjectOptions{}) + if err != nil { + return nil, 0, err + } + stat, err := obj.Stat() + if err != nil { + _ = obj.Close() + return nil, 0, err + } + return obj, stat.Size, nil +} + +func (c *Client) DeleteObject(ctx context.Context, key string) error { + return c.mc.RemoveObject(ctx, c.bucket, key, minio.RemoveObjectOptions{}) +} diff --git a/internal/workbench/files.go b/internal/workbench/files.go new file mode 100644 index 0000000..e748d2d --- /dev/null +++ b/internal/workbench/files.go @@ -0,0 +1,194 @@ +package workbench + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "io" + "strings" + "time" + + "aigateway.local/core/internal/platform/storage" + "github.com/jackc/pgx/v5" +) + +// FileObject is the metadata row for one object stored in MinIO. The object +// body lives in the bucket; this table is the searchable index plus the access +// control (personal scope is bound to a portal user, system scope to admins). +type FileObject struct { + ID string `json:"id"` + ObjectKey string `json:"-"` + OriginalName string `json:"original_name"` + ContentType string `json:"content_type"` + SizeBytes int64 `json:"size_bytes"` + ContentSHA256 string `json:"content_sha256"` + Scope string `json:"scope"` + OwnerUserID *string `json:"owner_user_id,omitempty"` + CreatedBy string `json:"created_by"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +// FileService stores object metadata in PostgreSQL and the body in MinIO. Both +// writes are kept consistent: the object is uploaded first and only then is the +// metadata row inserted; any failure rolls the object back. +type FileService struct { + assets *Service + store *storage.Client +} + +func NewFileService(assets *Service, store *storage.Client) *FileService { + return &FileService{assets: assets, store: store} +} + +func (s *FileService) MaxBytes() int64 { return s.store.MaxBytes() } + +// countingReader tallies how many bytes flow through it so uploads can enforce +// the configured size ceiling without buffering the whole file in memory. +type countingReader struct { + r io.Reader + n *int64 +} + +func (c *countingReader) Read(p []byte) (int, error) { + n, err := c.r.Read(p) + *c.n += int64(n) + return n, err +} + +// Upload streams the body into the object store and records its metadata. The +// owner must be set for the personal scope and nil for the system scope. +func (s *FileService) Upload(ctx context.Context, scope string, ownerUserID *string, createdBy, originalName, contentType string, body io.Reader) (FileObject, error) { + originalName = strings.TrimSpace(originalName) + if originalName == "" || len(originalName) > 255 { + return FileObject{}, errors.New("文件名无效") + } + if scope != "personal" && scope != "system" { + return FileObject{}, errors.New("无效的文件范围") + } + if contentType == "" { + contentType = "application/octet-stream" + } + objectKey, err := newUUID() + if err != nil { + return FileObject{}, err + } + objectKey = scope + "/" + objectKey + var size int64 + hasher := sha256.New() + limited := io.LimitReader(body, s.store.MaxBytes()+1) + counted := &countingReader{r: limited, n: &size} + rollbackObject := func() { _ = s.store.DeleteObject(ctx, objectKey) } + if _, err = s.store.PutObject(ctx, objectKey, io.TeeReader(counted, hasher), -1, contentType); err != nil { + return FileObject{}, err + } + if size > s.store.MaxBytes() { + rollbackObject() + return FileObject{}, errors.New("文件超过大小上限") + } + id, err := newUUID() + if err != nil { + rollbackObject() + return FileObject{}, err + } + var owner any + if ownerUserID != nil && strings.TrimSpace(*ownerUserID) != "" { + owner = *ownerUserID + } + var obj FileObject + err = s.assets.pool.QueryRow(ctx, `INSERT INTO gateway.file_objects(id,object_key,original_name,content_type,size_bytes,content_sha256,scope,owner_user_id,created_by) VALUES($1,$2,$3,$4,$5,$6,$7,$8,$9) RETURNING id::text,object_key,original_name,content_type,size_bytes,content_sha256,scope,owner_user_id::text,created_by,created_at,updated_at`, id, objectKey, originalName, contentType, size, hex.EncodeToString(hasher.Sum(nil)), scope, owner, createdBy).Scan(&obj.ID, &obj.ObjectKey, &obj.OriginalName, &obj.ContentType, &obj.SizeBytes, &obj.ContentSHA256, &obj.Scope, &obj.OwnerUserID, &obj.CreatedBy, &obj.CreatedAt, &obj.UpdatedAt) + if err != nil { + rollbackObject() + return FileObject{}, err + } + return obj, nil +} + +const fileObjectSelect = `SELECT id::text,object_key,original_name,content_type,size_bytes,content_sha256,scope,owner_user_id::text,created_by,created_at,updated_at FROM gateway.file_objects` + +func scanFileObject(row pgx.Row) (FileObject, error) { + var obj FileObject + err := row.Scan(&obj.ID, &obj.ObjectKey, &obj.OriginalName, &obj.ContentType, &obj.SizeBytes, &obj.ContentSHA256, &obj.Scope, &obj.OwnerUserID, &obj.CreatedBy, &obj.CreatedAt, &obj.UpdatedAt) + return obj, mapNotFound(err) +} + +func (s *FileService) ListPersonal(ctx context.Context, ownerUserID string) ([]FileObject, error) { + rows, err := s.assets.pool.Query(ctx, fileObjectSelect+` WHERE scope='personal' AND owner_user_id=$1 ORDER BY created_at DESC`, ownerUserID) + if err != nil { + return nil, err + } + defer rows.Close() + items := []FileObject{} + for rows.Next() { + obj, err := scanFileObject(rows) + if err != nil { + return nil, err + } + items = append(items, obj) + } + return items, rows.Err() +} + +func (s *FileService) ListSystem(ctx context.Context) ([]FileObject, error) { + rows, err := s.assets.pool.Query(ctx, fileObjectSelect+` WHERE scope='system' ORDER BY created_at DESC`) + if err != nil { + return nil, err + } + defer rows.Close() + items := []FileObject{} + for rows.Next() { + obj, err := scanFileObject(rows) + if err != nil { + return nil, err + } + items = append(items, obj) + } + return items, rows.Err() +} + +func (s *FileService) GetPersonal(ctx context.Context, ownerUserID, id string) (FileObject, error) { + return scanFileObject(s.assets.pool.QueryRow(ctx, fileObjectSelect+` WHERE id=$1 AND scope='personal' AND owner_user_id=$2`, id, ownerUserID)) +} + +func (s *FileService) GetSystem(ctx context.Context, id string) (FileObject, error) { + return scanFileObject(s.assets.pool.QueryRow(ctx, fileObjectSelect+` WHERE id=$1 AND scope='system'`, id)) +} + +// Open streams the object body for download. +func (s *FileService) Open(ctx context.Context, obj FileObject) (io.ReadCloser, int64, error) { + return s.store.OpenObject(ctx, obj.ObjectKey) +} + +func (s *FileService) DeletePersonal(ctx context.Context, ownerUserID, id string) error { + return s.deleteFile(ctx, id, `scope='personal' AND owner_user_id=$2`, ownerUserID) +} + +func (s *FileService) DeleteSystem(ctx context.Context, id string) error { + return s.deleteFile(ctx, id, `scope='system'`, nil) +} + +// deleteFile removes the metadata row and, best-effort, the object body. The +// object removal is best-effort because an orphan in the bucket is recoverable +// and must never block the delete that the user asked for. +func (s *FileService) deleteFile(ctx context.Context, id, where string, owner any) error { + // 只传 SQL 中真实出现的占位符参数:system 范围没有 $2,多传 nil 会让 + // Postgres 报 "bind message supplies 2 parameters" 错误。 + args := []any{id} + if owner != nil { + args = append(args, owner) + } + var key string + if err := s.assets.pool.QueryRow(ctx, `SELECT object_key FROM gateway.file_objects WHERE id=$1 AND `+where, args...).Scan(&key); err != nil { + return mapNotFound(err) + } + tag, err := s.assets.pool.Exec(ctx, `DELETE FROM gateway.file_objects WHERE id=$1 AND `+where, args...) + if err != nil { + return err + } + if tag.RowsAffected() == 0 { + return ErrNotFound + } + _ = s.store.DeleteObject(ctx, key) + return nil +} diff --git a/internal/workbench/files_admin_http.go b/internal/workbench/files_admin_http.go new file mode 100644 index 0000000..db228b3 --- /dev/null +++ b/internal/workbench/files_admin_http.go @@ -0,0 +1,165 @@ +package workbench + +import ( + "errors" + "io" + "mime/multipart" + "net/http" + "net/url" + "strconv" + "strings" + + "aigateway.local/core/internal/identity" + "aigateway.local/core/internal/platform/apiresponse" +) + +// FilesAdminHTTPHandler exposes the system-scoped file store to the admin +// console. Objects are uploaded through the gateway (never exposing MinIO), so +// the S3 credentials stay inside the API container. +type FilesAdminHTTPHandler struct { + files *FileService + identity *identity.Service + mux *http.ServeMux +} + +func NewFilesAdminHTTPHandler(files *FileService, identityService *identity.Service) *FilesAdminHTTPHandler { + h := &FilesAdminHTTPHandler{files: files, identity: identityService, mux: http.NewServeMux()} + h.mux.HandleFunc("POST /api/v1/admin/files", h.upload) + h.mux.HandleFunc("GET /api/v1/admin/files", h.list) + h.mux.HandleFunc("GET /api/v1/admin/files/{id}/download", h.download) + h.mux.HandleFunc("DELETE /api/v1/admin/files/{id}", h.delete) + return h +} + +func (h *FilesAdminHTTPHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { h.mux.ServeHTTP(w, r) } + +func (h *FilesAdminHTTPHandler) require(w http.ResponseWriter, r *http.Request, permission string) (identity.Account, bool) { + account, err := h.identity.Authenticate(r.Context(), identity.KindAdmin, r.Header.Get("Authorization")) + if err != nil { + apiresponse.Error(w, http.StatusUnauthorized, "登录状态无效") + return identity.Account{}, false + } + if !identity.HasPermission(account, permission) { + apiresponse.Error(w, http.StatusForbidden, "缺少文件管理权限") + return identity.Account{}, false + } + return account, true +} + +// firstFilePart pulls the multipart "file" part out of a streaming request so +// uploads never need to buffer the whole body in memory. +func firstFilePart(r *http.Request) (*multipart.Part, error) { + reader, err := r.MultipartReader() + if err != nil { + return nil, err + } + for { + part, err := reader.NextPart() + if err == io.EOF { + return nil, errors.New("no file part") + } + if err != nil { + return nil, err + } + if part.FormName() == "file" { + return part, nil + } + } +} + +// upload accepts either multipart/form-data (browser) or a raw body with a +// ?filename= query parameter (scripting). The size ceiling is enforced by the +// service's LimitReader, not by buffering. +func (h *FilesAdminHTTPHandler) upload(w http.ResponseWriter, r *http.Request) { + a, ok := h.require(w, r, identity.PermissionFileManage) + if !ok { + return + } + var ( + body io.Reader + originalName = strings.TrimSpace(r.URL.Query().Get("filename")) + contentType = r.Header.Get("Content-Type") + ) + if strings.HasPrefix(contentType, "multipart/form-data") { + part, err := firstFilePart(r) + if err != nil { + apiresponse.Error(w, http.StatusBadRequest, "缺少上传文件") + return + } + defer part.Close() + body = part + originalName = part.FileName() + if ct := part.Header.Get("Content-Type"); ct != "" { + contentType = ct + } + } + if originalName == "" { + apiresponse.Error(w, http.StatusBadRequest, "缺少文件名") + return + } + obj, err := h.files.Upload(r.Context(), "system", nil, a.ID, originalName, contentType, body) + if err != nil { + fileError(w, err) + return + } + apiresponse.OK(w, obj) +} + +func (h *FilesAdminHTTPHandler) list(w http.ResponseWriter, r *http.Request) { + if _, ok := h.require(w, r, identity.PermissionFileRead); !ok { + return + } + items, err := h.files.ListSystem(r.Context()) + if err != nil { + fileError(w, err) + return + } + apiresponse.OK(w, items) +} + +func (h *FilesAdminHTTPHandler) download(w http.ResponseWriter, r *http.Request) { + if _, ok := h.require(w, r, identity.PermissionFileRead); !ok { + return + } + obj, err := h.files.GetSystem(r.Context(), r.PathValue("id")) + if err != nil { + fileError(w, err) + return + } + serveFileContent(w, r, h.files, obj) +} + +func (h *FilesAdminHTTPHandler) delete(w http.ResponseWriter, r *http.Request) { + if _, ok := h.require(w, r, identity.PermissionFileManage); !ok { + return + } + if err := h.files.DeleteSystem(r.Context(), r.PathValue("id")); err != nil { + fileError(w, err) + return + } + apiresponse.OK(w, map[string]bool{"deleted": true}) +} + +// serveFileContent streams the object body straight to the client with a +// Content-Disposition so the original filename is preserved on download. +func serveFileContent(w http.ResponseWriter, r *http.Request, files *FileService, obj FileObject) { + reader, size, err := files.Open(r.Context(), obj) + if err != nil { + apiresponse.Error(w, http.StatusServiceUnavailable, "文件内容暂不可用") + return + } + defer reader.Close() + w.Header().Set("Content-Disposition", "attachment; filename*=UTF-8''"+url.PathEscape(obj.OriginalName)) + w.Header().Set("Content-Type", obj.ContentType) + w.Header().Set("Content-Length", strconv.FormatInt(size, 10)) + _, _ = io.Copy(w, reader) +} + +func fileError(w http.ResponseWriter, err error) { + switch { + case errors.Is(err, ErrNotFound): + apiresponse.Error(w, http.StatusNotFound, "文件不存在") + default: + apiresponse.Error(w, http.StatusBadRequest, err.Error()) + } +} diff --git a/internal/workbench/files_integration_test.go b/internal/workbench/files_integration_test.go new file mode 100644 index 0000000..f088079 --- /dev/null +++ b/internal/workbench/files_integration_test.go @@ -0,0 +1,142 @@ +package workbench + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "io" + "os" + "testing" + + "aigateway.local/core/internal/platform/config" + "aigateway.local/core/internal/platform/database" + "aigateway.local/core/internal/platform/storage" +) + +// TestFileObjectLifecycle runs against real PostgreSQL + MinIO. Set +// WORKBENCH_TEST_DATABASE_URL and WORKBENCH_TEST_S3_ENDPOINT to enable it; with +// the deploy compose up these are reachable from the build container via the +// deploy_default network. +func TestFileObjectLifecycle(t *testing.T) { + databaseURL := os.Getenv("WORKBENCH_TEST_DATABASE_URL") + s3Endpoint := os.Getenv("WORKBENCH_TEST_S3_ENDPOINT") + if databaseURL == "" || s3Endpoint == "" { + t.Skip("WORKBENCH_TEST_DATABASE_URL and WORKBENCH_TEST_S3_ENDPOINT are not set") + } + ctx := context.Background() + pool, err := database.Open(ctx, config.Database{URL: databaseURL, MaxConns: 8, MinConns: 0}) + if err != nil { + t.Fatal(err) + } + defer pool.Close() + + store, err := storage.NewClient(storage.Config{ + Endpoint: s3Endpoint, + AccessKeyID: os.Getenv("WORKBENCH_TEST_S3_ACCESS_KEY"), + SecretAccessKey: os.Getenv("WORKBENCH_TEST_S3_SECRET_KEY"), + Bucket: os.Getenv("WORKBENCH_TEST_S3_BUCKET"), + Region: "us-east-1", + MaxFileBytes: 8 << 20, + }) + if err != nil { + t.Fatal(err) + } + if err := store.EnsureBucket(ctx); err != nil { + t.Fatal(err) + } + adminID := "22222222-2222-4222-8222-222222222222" + portalID := "33333333-3333-4333-8333-333333333333" + if _, err := pool.Exec(ctx, `INSERT INTO gateway.admin_accounts(id,username,password_hash,role) VALUES($1,'m8-files-admin','test','superadmin') ON CONFLICT(id) DO NOTHING`, adminID); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `INSERT INTO gateway.portal_users(id,account,password_hash,name,active) VALUES($1,'m8-files-user','test','m8',true) ON CONFLICT(id) DO NOTHING`, portalID); err != nil { + t.Fatal(err) + } + cleanup := func() { + _, _ = pool.Exec(ctx, `DELETE FROM gateway.file_objects`) + } + cleanup() + defer cleanup() + + files := NewFileService(NewService(pool), store) + + content := []byte("M8 对象存储文件管理集成测试 payload\nline two\n") + owner := portalID + sysObj, err := files.Upload(ctx, "system", nil, adminID, "报告.docx", "application/vnd.openxmlformats-officedocument.wordprocessingml.document", bytes.NewReader(content)) + if err != nil { + t.Fatal(err) + } + personalObj, err := files.Upload(ctx, "personal", &owner, portalID, "notes.txt", "text/plain", bytes.NewReader(content)) + if err != nil { + t.Fatal(err) + } + + // sha256 往返一致 + hasher := sha256.Sum256(content) + if personalObj.ContentSHA256 != hex.EncodeToString(hasher[:]) || personalObj.SizeBytes != int64(len(content)) { + t.Fatalf("integrity mismatch: sha=%s size=%d", personalObj.ContentSHA256, personalObj.SizeBytes) + } + + // 流式读回 + reader, size, err := files.Open(ctx, personalObj) + if err != nil { + t.Fatal(err) + } + got, err := io.ReadAll(reader) + _ = reader.Close() + if err != nil || int64(len(got)) != size || !bytes.Equal(got, content) { + t.Fatalf("roundtrip mismatch: len=%d size=%d err=%v", len(got), size, err) + } + + // 列表按范围隔离 + sysList, err := files.ListSystem(ctx) + if err != nil || len(sysList) != 1 || sysList[0].ID != sysObj.ID { + t.Fatalf("system list mismatch: %#v err=%v", sysList, err) + } + myList, err := files.ListPersonal(ctx, portalID) + if err != nil || len(myList) != 1 || myList[0].ID != personalObj.ID { + t.Fatalf("personal list mismatch: %#v err=%v", myList, err) + } + otherList, err := files.ListPersonal(ctx, "44444444-4444-4444-8444-444444444444") + if err != nil || len(otherList) != 0 { + t.Fatalf("other personal list should be empty: %#v err=%v", otherList, err) + } + + // 归属隔离:A 不能读 B 的文件 + if _, err := files.GetPersonal(ctx, adminID, personalObj.ID); !errors.Is(err, ErrNotFound) { + t.Fatalf("expected ErrNotFound for cross-owner read, got %v", err) + } + if _, err := files.GetSystem(ctx, personalObj.ID); !errors.Is(err, ErrNotFound) { + t.Fatalf("expected ErrNotFound for system read of personal file, got %v", err) + } + + // 超限上传被拒(8MiB 上限,用无 EOF 占位 reader 凑 9MiB) + if _, err := files.Upload(ctx, "system", nil, adminID, "big.bin", "application/octet-stream", io.LimitReader(zeroReader{}, 9<<20)); err == nil { + t.Fatal("expected oversized upload to fail") + } + + // 删除后元数据与对象都消失 + if err := files.DeleteSystem(ctx, sysObj.ID); err != nil { + t.Fatal(err) + } + if _, err := files.GetSystem(ctx, sysObj.ID); !errors.Is(err, ErrNotFound) { + t.Fatalf("expected ErrNotFound after delete, got %v", err) + } + if err := files.DeletePersonal(ctx, portalID, personalObj.ID); err != nil { + t.Fatal(err) + } + if _, err := files.GetPersonal(ctx, portalID, personalObj.ID); !errors.Is(err, ErrNotFound) { + t.Fatalf("expected ErrNotFound after personal delete, got %v", err) + } +} + +type zeroReader struct{} + +func (zeroReader) Read(p []byte) (int, error) { + for i := range p { + p[i] = 0 + } + return len(p), nil +} diff --git a/internal/workbench/portal_files_http.go b/internal/workbench/portal_files_http.go new file mode 100644 index 0000000..756a6a5 --- /dev/null +++ b/internal/workbench/portal_files_http.go @@ -0,0 +1,133 @@ +package workbench + +import ( + "errors" + "io" + "mime/multipart" + "net/http" + "strings" + + "aigateway.local/core/internal/identity" + "aigateway.local/core/internal/platform/apiresponse" +) + +// FilesPortalHTTPHandler serves each portal user's personal file store. Unlike +// the admin handler every object is scoped to the authenticated account, so a +// user can only ever see and open their own files. +type FilesPortalHTTPHandler struct { + files *FileService + identity *identity.Service + mux *http.ServeMux +} + +func NewFilesPortalHTTPHandler(files *FileService, identityService *identity.Service) *FilesPortalHTTPHandler { + h := &FilesPortalHTTPHandler{files: files, identity: identityService, mux: http.NewServeMux()} + h.mux.HandleFunc("POST /api/v1/portal/files", h.upload) + h.mux.HandleFunc("GET /api/v1/portal/files", h.list) + h.mux.HandleFunc("GET /api/v1/portal/files/{id}/download", h.download) + h.mux.HandleFunc("DELETE /api/v1/portal/files/{id}", h.delete) + return h +} + +func (h *FilesPortalHTTPHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { h.mux.ServeHTTP(w, r) } + +func (h *FilesPortalHTTPHandler) account(w http.ResponseWriter, r *http.Request) (identity.Account, bool) { + account, err := h.identity.Authenticate(r.Context(), identity.KindPortal, r.Header.Get("Authorization")) + if err != nil { + apiresponse.Error(w, http.StatusUnauthorized, "登录状态无效或已过期") + return identity.Account{}, false + } + return account, true +} + +func (h *FilesPortalHTTPHandler) upload(w http.ResponseWriter, r *http.Request) { + a, ok := h.account(w, r) + if !ok { + return + } + var ( + body io.Reader + originalName = strings.TrimSpace(r.URL.Query().Get("filename")) + contentType = r.Header.Get("Content-Type") + ) + if strings.HasPrefix(contentType, "multipart/form-data") { + part, err := h.firstFilePart(r) + if err != nil { + apiresponse.Error(w, http.StatusBadRequest, "缺少上传文件") + return + } + defer part.Close() + body = part + originalName = part.FileName() + if ct := part.Header.Get("Content-Type"); ct != "" { + contentType = ct + } + } + if originalName == "" { + apiresponse.Error(w, http.StatusBadRequest, "缺少文件名") + return + } + obj, err := h.files.Upload(r.Context(), "personal", &a.ID, a.ID, originalName, contentType, body) + if err != nil { + fileError(w, err) + return + } + apiresponse.OK(w, obj) +} + +func (h *FilesPortalHTTPHandler) firstFilePart(r *http.Request) (*multipart.Part, error) { + reader, err := r.MultipartReader() + if err != nil { + return nil, err + } + for { + part, err := reader.NextPart() + if err == io.EOF { + return nil, errors.New("no file part") + } + if err != nil { + return nil, err + } + if part.FormName() == "file" { + return part, nil + } + } +} + +func (h *FilesPortalHTTPHandler) list(w http.ResponseWriter, r *http.Request) { + a, ok := h.account(w, r) + if !ok { + return + } + items, err := h.files.ListPersonal(r.Context(), a.ID) + if err != nil { + fileError(w, err) + return + } + apiresponse.OK(w, items) +} + +func (h *FilesPortalHTTPHandler) download(w http.ResponseWriter, r *http.Request) { + a, ok := h.account(w, r) + if !ok { + return + } + obj, err := h.files.GetPersonal(r.Context(), a.ID, r.PathValue("id")) + if err != nil { + fileError(w, err) + return + } + serveFileContent(w, r, h.files, obj) +} + +func (h *FilesPortalHTTPHandler) delete(w http.ResponseWriter, r *http.Request) { + a, ok := h.account(w, r) + if !ok { + return + } + if err := h.files.DeletePersonal(r.Context(), a.ID, r.PathValue("id")); err != nil { + fileError(w, err) + return + } + apiresponse.OK(w, map[string]bool{"deleted": true}) +} diff --git a/migrations/000023_minio_file_objects.sql b/migrations/000023_minio_file_objects.sql new file mode 100644 index 0000000..77ff81d --- /dev/null +++ b/migrations/000023_minio_file_objects.sql @@ -0,0 +1,32 @@ +-- M8 基础设施:对象存储文件管理 +-- 元数据行落 PostgreSQL,文件体存 MinIO(S3 兼容)。对象键即桶内路径,body 永不经数据库。 + +-- file_objects:文件元数据索引 + 访问控制 +-- scope='personal' 绑定 portal 用户(个人文件仓);scope='system' 为平台文件(admin 文件管理)。 +CREATE TABLE IF NOT EXISTS gateway.file_objects ( + id uuid PRIMARY KEY, + object_key text NOT NULL UNIQUE, + original_name text NOT NULL CHECK (length(original_name) BETWEEN 1 AND 255), + content_type text NOT NULL DEFAULT 'application/octet-stream', + size_bytes bigint NOT NULL DEFAULT 0 CHECK (size_bytes BETWEEN 0 AND 536870912), + content_sha256 char(64) NOT NULL, + scope text NOT NULL CHECK (scope IN ('personal', 'system')), + owner_user_id uuid REFERENCES gateway.portal_users(id) ON DELETE CASCADE, + created_by uuid NOT NULL, + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + updated_at timestamptz NOT NULL DEFAULT clock_timestamp(), + CONSTRAINT file_objects_personal_owner CHECK ( + (scope = 'personal' AND owner_user_id IS NOT NULL) + OR (scope = 'system' AND owner_user_id IS NULL) + ) +); + +-- 部分索引:个人文件按用户倒序,系统文件按时间倒序,sha256 去重查找 +CREATE INDEX IF NOT EXISTS file_objects_personal_idx + ON gateway.file_objects (owner_user_id, created_at DESC) + WHERE scope = 'personal'; +CREATE INDEX IF NOT EXISTS file_objects_system_idx + ON gateway.file_objects (created_at DESC) + WHERE scope = 'system'; +CREATE INDEX IF NOT EXISTS file_objects_sha256_idx + ON gateway.file_objects (content_sha256); diff --git a/web/apps/admin/src/api/files.ts b/web/apps/admin/src/api/files.ts new file mode 100644 index 0000000..4c78d9c --- /dev/null +++ b/web/apps/admin/src/api/files.ts @@ -0,0 +1,41 @@ +import axios from 'axios' +import request from '@/utils/http' +import { useUserStore } from '@/store/modules/user' + +export interface FileObject { + id: string + original_name: string + content_type: string + size_bytes: number + content_sha256: string + scope: string + owner_user_id?: string + created_by: string + created_at: string + updated_at: string +} + +export const fetchFiles = () => request.get({ url: '/api/v1/admin/files' }) +export const uploadFile = (file: File) => { + const form = new FormData() + form.append('file', file) + return request.post({ url: '/api/v1/admin/files', data: form, timeout: 120000 }) +} +export const deleteFile = (id: string) => request.del({ url: `/api/v1/admin/files/${id}` }) + +// 文件体以 blob 流式下载,绕开统一响应拦截器(其期望 JSON envelope)。 +export async function downloadFile(id: string, filename: string) { + const { accessToken } = useUserStore() + const token = accessToken.startsWith('Bearer ') ? accessToken : `Bearer ${accessToken}` + const res = await axios.get(`/api/v1/admin/files/${id}/download`, { + baseURL: import.meta.env.VITE_API_URL, + headers: { Authorization: token }, + responseType: 'blob' + }) + const url = URL.createObjectURL(res.data) + const anchor = document.createElement('a') + anchor.href = url + anchor.download = filename + anchor.click() + URL.revokeObjectURL(url) +} diff --git a/web/apps/admin/src/views/gateway/files/index.vue b/web/apps/admin/src/views/gateway/files/index.vue new file mode 100644 index 0000000..6659080 --- /dev/null +++ b/web/apps/admin/src/views/gateway/files/index.vue @@ -0,0 +1,55 @@ + + diff --git a/web/apps/portal/src/api/portal.ts b/web/apps/portal/src/api/portal.ts index 2b29795..9cb823a 100644 --- a/web/apps/portal/src/api/portal.ts +++ b/web/apps/portal/src/api/portal.ts @@ -1,4 +1,6 @@ +import axios from 'axios' import request from '@/utils/http' +import { useUserStore } from '@/store/modules/user' export interface Application { id:string;code:string;name:string;description:string;published_version?:number } export interface Knowledge { id:string;name:string;description:string;document_count:number;chunk_count:number } @@ -21,6 +23,31 @@ export const fetchStats=(days:number)=>request.get({url:'/api/v1/portal/s export const fetchLogs=(days:number,limit=50)=>request.get<{items:AuditEvent[]}>({url:'/api/v1/portal/logs',params:{days,limit}}) export const changePassword=(params:{old_password:string;new_password:string})=>request.post({url:'/api/v1/portal/password',params}) +// --- 个人文件仓库(M8 对象存储)--- +export interface FileObject { id:string;original_name:string;content_type:string;size_bytes:number;content_sha256:string;scope:string;created_at:string;updated_at:string } +export const fetchMyFiles=()=>request.get({url:'/api/v1/portal/files'}) +export const uploadMyFile=(file:File)=>{ + const form=new FormData() + form.append('file',file) + return request.post({url:'/api/v1/portal/files',data:form,timeout:120000}) +} +export const deleteMyFile=(id:string)=>request.del({url:`/api/v1/portal/files/${id}`}) +export async function downloadMyFile(id:string,filename:string){ + const { accessToken } = useUserStore() + const token = accessToken.startsWith('Bearer ')?accessToken:`Bearer ${accessToken}` + const res = await axios.get(`/api/v1/portal/files/${id}/download`,{ + baseURL:import.meta.env.VITE_API_URL, + headers:{Authorization:token}, + responseType:'blob' + }) + const url = URL.createObjectURL(res.data) + const anchor = document.createElement('a') + anchor.href = url + anchor.download = filename + anchor.click() + URL.revokeObjectURL(url) +} + // --- 资源市场 --- export interface MarketItem { type:string;code:string;name:string;description:string;category_id?:string;category_name:string;tags:string[];department_ids:string[];updated_at:string } diff --git a/web/apps/portal/src/views/portal/files/index.vue b/web/apps/portal/src/views/portal/files/index.vue new file mode 100644 index 0000000..7451004 --- /dev/null +++ b/web/apps/portal/src/views/portal/files/index.vue @@ -0,0 +1,54 @@ + +