From 4f736cba21773a36512595378717759affc7beb6 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 21 Sep 2026 16:44:09 +0800 Subject: [PATCH] feat: split account management --- AGENTS.md | 99 +++--- docs/e2e-test-plan.md | 17 +- docs/plan01.md | 2 +- .../controlplane/api/app_migrated_test.go | 3 +- internal/controlplane/api/creator.go | 100 ++++-- .../api/creator_route_validation_test.go | 3 +- internal/creator/competitor_share_jobs.go | 233 +++++++++++++ .../creator/competitor_share_jobs_test.go | 29 ++ internal/creator/competitor_upsert.go | 27 ++ internal/creator/content.go | 12 +- internal/creator/integration_test.go | 65 ++++ .../migrations/037_competitor_share_jobs.sql | 18 + internal/creator/models.go | 21 ++ internal/creator/store.go | 4 + web/src/app/Layout.jsx | 10 +- web/src/app/Layout.test.jsx | 25 +- .../features/accounts/AccountImportPage.jsx | 295 ++++++++++++++++ web/src/features/accounts/AccountsPage.jsx | 121 ++++--- .../features/accounts/AccountsPage.test.jsx | 52 ++- .../features/creator/CreatorPages.test.jsx | 186 ++++------ .../competitors/CreatorCompetitorsPage.jsx | 216 +----------- .../creator/settings/CreatorSettingsPage.jsx | 326 ++++++++++++------ web/src/routes/_app/accounts/import.tsx | 4 + web/src/routes/_app/accounts/monitoring.tsx | 4 + web/src/shared/api/dataProvider.js | 4 +- 25 files changed, 1246 insertions(+), 630 deletions(-) create mode 100644 internal/creator/competitor_share_jobs.go create mode 100644 internal/creator/competitor_share_jobs_test.go create mode 100644 internal/creator/competitor_upsert.go create mode 100644 internal/creator/migrations/037_competitor_share_jobs.sql create mode 100644 web/src/features/accounts/AccountImportPage.jsx create mode 100644 web/src/routes/_app/accounts/import.tsx create mode 100644 web/src/routes/_app/accounts/monitoring.tsx diff --git a/AGENTS.md b/AGENTS.md index bd1422b..c782d00 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,9 +1,7 @@ -# 仓库规范 - -## 宪法(通用工程约束) +## 宪法 - 任何涉及文件的调研或修改,如果当前是 git 仓库,需要先同步远程提交到本地,避免调研过时问题。 -- 基于 TDD 进行功能的开发与业务变更,单元测试覆盖率要保证 65% 以上。 +- 基于 TDD 进行功能的开发与业务变更,单元测试覆盖率要保证 65% 以上 - 任何时候我提出任何需求均需要理解并**结构化复述后与我进行确认,避免理解偏差**。 - 不要在代码里藏兜底逻辑来吞掉错误、隐藏问题。出了问题就应该让它爆出来,否则你永远找不到真实问题。 - 当一个问题出现时,不要用各种 small fix、针对性补丁来掩盖它。**必须定位真实根因,彻底修复**。在 bug 上糊纸只会让系统积累你不知道的危险暗病。 @@ -11,18 +9,47 @@ - 始终注意在关键路径上给自己留足排查日志,确保每一个**关键节点都是可追溯**的。 - 当项目关键技术栈或产品方向发生变更时,同步更新 agents.md。文档必须随代码一起演进,不能让它变成过时的谎言。 - 大规模重构或实验性改动前,必须先切新分支。 -- 不以维护向后兼容性为目标。**对于已经废弃的代码路径,应直接移除**,不再通过兼容层、回退机制或迁移方案予以保留。(注:开发阶段数据兼容豁免见下,本条针对代码路径。) + +- 不以维护向后兼容性为目标。**对于已经废弃的代码路径,应直接移除**,不再通过兼容层、回退机制或迁移方案予以保留。 - 在充分满足当前需求的前提下,采用**尽可能简单的实现方案**。避免引入缺乏实际需求依据的抽象、配置项和间接层。 - **采用渐进式、分层的方式构建系统**。首先完成能够端到端运行的最小版本,再基于稳定可用的产品逐步增加功能。不要以尚未成熟的复杂性取代已经可用的产品。 - **保持组件的模块化**,并明确划分不同职责与关注点。 - 当成熟且维护良好的库能够降低整体复杂度或提高可靠性时,应优先采用。除非有明确理由,不要重复实现通用功能。 -- 在自行实现功能或新增依赖之前,应优先评估项目现有依赖的能力。应先查阅相关文档与类型定义,不应未经确认就认定某个库不具备所需能力。 +- 在自行实现功能或新增依赖之前,应优先评估项目现有依赖的能力。应先查阅相关文档和类型定义,不应未经确认就认定某个库不具备所需能力。 - 架构决策**应着眼于长期演进**。不要采用仅能解决当前问题、且预期需要在后续替换的权宜方案。 - 在设计解决方案之前,**先研究成熟产品如何解决同类问题**。优先采用经过验证的模式和约定,避免从零开始另行设计一套方案。 +## 禁止清单(不主动考虑、不主动提议、不实现,遇到只记入 TODO 技术债列表) + +1. 法律合规:商业库授权、开源协议合规、GDPR/个保、隐私政策(法务负责)。 +2. 依赖安全:NPM 及第三方包漏洞、安全补丁、依赖升级策略。 +3. 访问安全:服务只需支持局域网访问(host 绑定 0.0.0.0 即可),不考虑公网暴露、HTTPS、认证/权限体系(登录、RBAC)、限流、防爬、数据加密、审计日志。 + +## 红线清单(快速阶段也不能省,现在便宜、以后极贵) + +1. 数据模型/表结构:认真设计,建表慎重——改表成本远高于写代码。 +2. 目录结构与模块边界:保持简单清晰,不堆一坨代码。 +3. 基础错误日志:出错时至少能看到发生了什么。 +4. Git:小步提交,保持历史清晰。 +5. 基础输入校验:仅防止程序崩溃,不做安全加固。 +6. 环境差异配置(端口、地址等)与代码分离(.env 或配置项)。 + +## 技术栈 + +前端框架:Refine+shadcn+Tailwind CSS; + +后台管理:[TanstackStarter ](https://github.com/kiranism/tanstack-start-dashboard) 框架 + +图标 RemixIcon; + +后端 Go 库 fiber/v3; logus; viper; cobra;samber/lo; samber/mo + +用户交互需要使用 skill:impeccable 优化交互操作 + ## 沟通方式 - 向使用者回报时,使用清楚直白的语言说明做了什么、结果如何。最终回复禁用术语、技术实现细节与工程腔。写法是:对一个聪明但没在看代码的人解释。 + - 实际执行过程(思考、规划、写程序、除错、解决问题)保持完整的技术严谨度,这条规范只适用于对使用者的沟通方式。 ## 回复风格 @@ -35,9 +62,11 @@ - 不使用口语化表达,说重点,简单明了 - 需要时搭配条列式与表格加强输出可读性 -## 多选项决策 +## 决策规则 -- 当方案有多个选项时,列出每个选项的优缺点,并明确指出推荐选项与原因。 +- 当方案有多个选项时,列出每个选项的优缺点,并明确指出推荐选项与原因,先问我。 +- 有多种实现方式时,选最简单能跑通的。 +- 遇到"禁止清单"中的问题:不展开、不实现,追加到 TODO 技术债列表即可。 ## Sub-Agent 使用时机 @@ -47,63 +76,15 @@ - 各子任务职责明确分离,合并执行会造成 context 混杂 - 大量结构相同的重复性任务(可用 `spawn_agents_on_csv` batch 执行) - 各子任务需要不同的 model 配置或 sandbox 权限,例如: + - 探索型任务使用轻量 model + `read-only` sandbox - 审查型任务使用高推理 model + `read-only` sandbox - 修改型任务使用执行导向 model + `workspace-write` sandbox -## 开发阶段原则 +## 验证标准 -- 开发阶段仅关注业务功能:不实现访问限制、认证、网络隔离等安全策略,安全由用户自行把控。 -- 开发阶段不做数据兼容:数据库 schema 可随时破坏性重建,不编写迁移兼容、存量回填或双写代码;网关与环境数据均可删除重来。 -- 保持变更小而独立可评审,并附带覆盖该变更的最小相关检查。 -- 使用下方已批准的技术栈;在栈内优先复用现有代码、标准库和平台原生能力,而非新增依赖或抽象。 -- 在集成边界保持幂等性和向后兼容;文档化重试与失败行为。 -- 永不提交密钥、生产凭据或个人账号数据。 -- 引用上游项目时,记录其来源与许可证;除非许可证明确允许复用,否则须独立实现。 - -## 当前产品方向 - -- 当前业务范围与验收以 [docs/plan01.md](docs/plan01.md) 为准:竞品分析、账号与大小号响应、环境与代理、评论线索和私信;先完成抖音完整流程,再完成小红书。 -- 与旧探索规划或现有功能冲突时,以上述需求及使用者最新确认为准。自动响应按策略和 UID 冷却执行;人工发送逐次确认,两者不得混淆。自有账号互动与私信采用事件监听,不用轮询或 Mock 冒充实际能力。 -- 上述内容是目标范围,不代表所有平台能力都已通过真实账号验收。平台能力缺口、新的业务歧义必须向使用者确认;禁止自行删减需求、隐藏失败或以未要求的通用框架扩大实现范围。控制面继续使用 Go,浏览器 gateway 使用 Python。 -- 浏览器环境由各机器 gateway 管理宿主机浏览器进程与 Xvfb,不再通过 Docker 创建浏览器环境;复用现有多机控制契约。采集结束(含失败、取消、超时)回收任务创建的临时资源,保留账号 Profile、正式结果和合法长期监听。评审、实施与验收见 [变更评审](docs/native-browser-change-review.md)、[实施计划](docs/native-browser-implementation-plan.md)、[验证文档](docs/native-browser-verification.md)。当前 worktree 实施单节点 native browser;多节点、A/B、跨机故障恢复和性能对比另行验收。 - -## 已批准的技术栈 - -### Go 控制面 - -- 使用 Go 1.26、Fiber v3(HTTP 路由与服务生命周期)、Viper(配置)、Logrus(应用日志)、Cobra(可执行入口)。依赖版本由 `go.mod` 和 `go.sum` 精确锁定。 -- 保留成熟的标准库集成,如反向代理,不重复造轮子;去 Docker 改造完成后移除浏览器链路的 Docker HTTP 客户端。`net/http` handler 跨越 Fiber 边界时,使用 Fiber 官方适配器。 -- 创建局部 `viper.New()` 实例,只绑定支持的输入,显式应用默认值,并在产生网络、文件系统或 Docker 副作用前完成全部配置校验。未经评审的需求批准,不使用 Viper 全局单例、远程 provider 或热加载。 -- 通过 Logrus 输出结构化 JSON 日志,保持 `service` 等稳定字段。可复用代码只返回错误,并在服务边界记录一次。 -- 每个服务只保留一个最小化的 Cobra 根命令。仅当存在真实的运维工作流需求时,才添加子命令、持久化 flag、代码生成器或补全。 -- 未经评审的需求批准,不添加 ORM、Redis、任务框架或另一套 HTTP/配置/日志/CLI 技术栈。 - -### Python 浏览器 gateway - -- 浏览器 gateway 使用 Python 3.12+;优先使用标准库 HTTP、进程管理、socket/ssl/asyncio 与显式输入校验,复用已有 CDP/WebSocket 能力。当前 native gateway 由宿主机非 root systemd user service 管理 Xvfb、浏览器、Profile、代理和运行代次;仅在实际需求明确且批准后新增锁定的浏览器依赖。 -- 每台 gateway 仅管理本机浏览器/Xvfb 生命周期、运行代次、Profile、临时资源、代理及页面动作,复用现有多机路由;账号身份必须在每次写操作前核对。浏览器链路不依赖 Docker socket,不创建浏览器容器、卷、网络或镜像,也不保留 Docker 回退。 -- 任务清理必须核对节点、运行代次和资源归属,不能误删账号登录资料、正式素材或他人会话;控制面与 gateway 各自处理本机创建的临时文件,清理失败必须可见。 -- Python 依赖必须写入锁定文件;不允许自动登录、任意 CDP、Cookie/验证码/密码回显或把不确定写结果转换为成功。 - -### React 前端 - -- 使用 React 19、TanStack Start、TanStack Router、TanStack Query、Vite 7、shadcn/ui 与 Tailwind CSS。业务数据访问使用 `web/src/shared/api/dataProvider.js`,不再引入 React Router 或 Refine;路由树由 `web/scripts/generate-route-tree.mjs` 生成并在 `web/package-lock.json` 中锁定精确版本。 -- 视觉体系由 CreatorHub 自有掌控:品牌 token 保留在项目主题与 Tailwind 配置中,导航/外壳保留在项目自有的 Layout 组件中。shadcn/ui 组件按需生成并纳入仓库自有代码,不引入 Ant Design、Arco、MUI 等其他组件体系。 -- 图标统一使用 RemixIcon,不为差异化自创奇怪图标。 -- 将启动、停止、回收等领域动作保留为显式动作;不得为迎合通用 CRUD 惯例而伪装成资源更新。 +开始任务前先定义完成标准。交付前依此验证,发现问题就修好再测,不把未完成的工作交回给使用者。只有确认完成,或遇到真正需要使用者介入的障碍时,才回报。 ## UI 设计 默认收敛、克制、常规;尺寸与间距根据界面类型、信息密度、平台习惯、使用频率和视觉层级判断,不写死统一规格,也不主动放大。辅助入口、设置、开关、工具按钮不应抢视觉中心。常见功能必须使用大众通用、用户一眼可识别的图标隐喻,优先成熟图标库、系统图标或行业通用符号,不为差异化自创奇怪图标;自定义图标也必须保持常见轮廓、比例和语义。除非明确要求强调,否则优先用位置、分组、轻微颜色、hover、tooltip、分隔线和状态反馈表达层级,避免夸张尺寸、重色块、大圆角、厚边框、强阴影、装饰性渐变和营销页式布局。实现后必须与同屏元素对比检查,若显得突兀、过大、过重或破坏信息密度,应主动收敛。 - -## 验证与交付 - -- 功能完成后的功能验收不使用浏览器或其他自动化操作;只启动可联调的测试环境,并提供清晰的手工验证步骤,由使用者完成实际功能验证。 -- 测试阶段默认不构建 Docker 镜像或容器;优先直接裸启动本地 Go control-plane 与 Vite 前端,保证本地服务可运行、可联调。 -- 测试环境必须支持局域网手工验证:前端与本地 control-plane 监听 `0.0.0.0`,不得只绑定 `127.0.0.1`,并提供局域网访问地址。 -- 开始任务前先定义完成标准。交付前依此验证,发现问题就修好再测,不把未完成的工作交回给使用者。只有确认完成,或遇到真正需要使用者介入的障碍时,才回报。 -- 每个非平凡行为变更附带最小的回归测试,且该测试在无此变更时会失败。在信任与集成边界覆盖成功、校验、失败和兼容路径;单元测试覆盖率保证 65% 以上。 -- 控制面变更必须通过 `go test ./...`、`go vet ./...`,并构建 `./cmd/control-plane`;涉及并发、生命周期或共享状态的变更须运行 `go test -race ./...`。Python gateway 必须通过其非交互式单元测试与覆盖率检查;只有明确涉及 gateway/Docker/Compose 变更且获得使用者同意时,才执行 Compose 构建/健康检查。Docker 或 Compose 变更还须通过 `docker compose config --quiet`。 -- 前端变更必须从 lockfile 安装、通过仓库的非交互式测试命令,并通过 `npm --prefix web run build`。主题、Layout、导航、资源动作或 data provider 的变更需要聚焦的交互覆盖,包括适用的错误与禁用状态。 -- 除非 issue 明确批准契约变更,保持既有 API 行为不变;浏览器生命周期已按批准的 native runtime 契约替换旧 Docker 生命周期。在 PR 中文档化任何状态码、载荷、配置、迁移、安全或重试方面的影响。 diff --git a/docs/e2e-test-plan.md b/docs/e2e-test-plan.md index cdab99c..2984bc7 100644 --- a/docs/e2e-test-plan.md +++ b/docs/e2e-test-plan.md @@ -834,14 +834,14 @@ gateway 使用宿主机已核验的默认浏览器运行时;运行时安装、 ### 5.12 竞品、作品、评论与指标采集 -#### COL-01 竞品登记与本地链接预览 [UI/API] +#### COL-01 分享链接入队与竞品作者登记 [UI/API] -- 前置:获准C主页、真实稳定标识;附录F。 +- 前置:获准作品分享链接;附录F。 - 操作: - 1. 抖音主页填入后点解析链接预览,对照人工确认标识;填写昵称加入监测 - 2. 试错误host、无路径、非法标识和非法URL;重复登记同平台标识 - 3. 读取列表/详情,确认同步状态与启用标志 -- 预期:预览仅本地候选,不冒充身份认证;重复/非法明确失败;成功清表单但保留平台,未同步不显示伪作品。 + 1. 粘贴抖音或小红书分享内容,确认只创建分享链接解析任务,不同步等待作者信息 + 2. 查询任务列表,确认状态从待解析/解析中变化;后台成功后竞品作者才出现在竞品列表并开始监控 + 3. 构造无效链接或网关不可用场景,确认自动最多尝试3次,失败原因可查询,并可点击“重新入队” +- 预期:分享链接任务与竞品作者列表分离;失败不创建半成品作者,不隐藏真实原因;重复成功解析通过平台账号标识幂等更新。 #### COL-02 暂停恢复与同步前置 [UI/API] @@ -1387,7 +1387,10 @@ E-strategy: | 方法 / 路径 | body或查询样例 | 入口 / 用例 | | --- | --- | --- | | GET `/api/creator/competitors` | `?platform=douyin` | UI、COL-01 | -| POST `/api/creator/competitors` | F-competitor | UI、COL-01 | +| GET `/api/creator/competitor-share-jobs` | `?status=failed&platform=douyin` | UI、COL-01 | +| POST `/api/creator/competitor-share-jobs` | `{"platform":"douyin","share_url":"https://v.douyin.com/...","tags":["重点监测"]}` | UI、COL-01;立即入队,后台最多尝试3次 | +| GET `/api/creator/competitor-share-jobs/{J}` | 无 | UI、COL-01;查询状态与失败原因 | +| POST `/api/creator/competitor-share-jobs/{J}/retry` | 无 | UI、COL-01;失败任务重新入队 | | GET `/api/creator/competitors/{C}` | 无 | **API-only详情**、COL-01/02 | | POST `/api/creator/competitors/{C}/pause` | 无 | UI、COL-02 | | POST `/api/creator/competitors/{C}/resume` | 无 | UI、COL-02 | diff --git a/docs/plan01.md b/docs/plan01.md index 4511fe8..fe7fafd 100644 --- a/docs/plan01.md +++ b/docs/plan01.md @@ -299,7 +299,7 @@ | 页面与推荐入口 | 列表、详情及关键字段 | 编辑、按钮与跳转 | 对应验收 | | --- | --- | --- | --- | -| 竞品 `/competitors`,详情 `/competitors/:id` | 平台、昵称/稳定标识、主页、监测启停、最近/下次采集、进度和失败;详情用“作品/采集记录”页签 | 导入链接→解析预览→确认保存/取消;启用/暂停、手动更新;账号进入作品列表 | AC-C1、AC-C2、AC-C3、AC-U1 | +| 竞品 `/competitors`,分享任务 `/creator/competitors` | 竞品作者展示平台、昵称/稳定标识、主页、监测启停、最近/下次采集、进度和失败;分享任务展示链接、平台、处理状态、尝试次数和失败原因 | 分享链接直接入队→后台最多尝试3次解析→成功后自动保存作者并监控;失败可查询原因并重新入队;竞品作者启用/暂停、手动更新 | AC-C1、AC-C2、AC-C3、AC-U1 | | 作品(竞品详情作品页签,跨账号筛选复用同一列表) | C1 字段、C3 筛选、多页加载、已采页/条数及总量未知提示;详情展示原文、指标时间、下一指标计划/停止原因 | 筛选/清空、分页、打开原平台、选取素材;详情抽屉返回保留筛选与页码,进度不以首屏当完成 | AC-C2、AC-C4、AC-C5、AC-C6、AC-U1 | | 素材(作品详情“素材/文稿”页签) | 选中来源、分步状态、视频/音频预览、转写、无音轨/无语音结果、产物引用与失败原因 | 第一次确认选取后才准备;只重试失败步骤;准备完成后填要求并第二次确认仿写/放弃;标题/口播编辑、保存,保存失败保留输入,无发布按钮 | AC-C7、AC-C8、AC-C9、AC-C10、AC-U2 | | 账号 `/accounts`、`/accounts/:id` | A1 资料、业务/登录/环境状态分列;详情“资料/登录/大小号与策略/监听记录”页签 | 新增/编辑抽屉、保存/取消;密码只显示已配置;打开绑定的同一浏览器人工登录、重核身份,冲突显示期望/实际标识并停止;跳转环境与任务 | AC-A1、AC-A2、AC-A3、AC-U3 | diff --git a/internal/controlplane/api/app_migrated_test.go b/internal/controlplane/api/app_migrated_test.go index e394123..923174d 100644 --- a/internal/controlplane/api/app_migrated_test.go +++ b/internal/controlplane/api/app_migrated_test.go @@ -135,8 +135,7 @@ func TestCreatorRouteValidationCoverage(t *testing.T) { {http.MethodPost, "/api/creator/strategies/missing/disable"}, {http.MethodDelete, "/api/creator/strategies/missing"}, {http.MethodPut, "/api/creator/strategies/missing"}, - {http.MethodPost, "/api/creator/competitors/preview"}, - {http.MethodPost, "/api/creator/competitors"}, + {http.MethodPost, "/api/creator/competitor-share-jobs"}, {http.MethodPost, "/api/creator/competitors/missing/pause"}, {http.MethodPost, "/api/creator/competitors/missing/resume"}, {http.MethodPost, "/api/creator/competitors/missing/sync"}, diff --git a/internal/controlplane/api/creator.go b/internal/controlplane/api/creator.go index 1dc92db..03acb0a 100644 --- a/internal/controlplane/api/creator.go +++ b/internal/controlplane/api/creator.go @@ -273,36 +273,45 @@ func registerCreatorWithServices(app *fiber.App, store *creator.Store, phaseASto } return c.JSON(items) }) - app.Post("/api/creator/competitors/preview", func(c fiber.Ctx) error { - var input competitorShareRequest - if err := decodeCreator(c, &input); err != nil { - return creatorError(c, err) - } - preview, err := previewCompetitorShare(c.Context(), store, phaseAStore, hubStore, input.AccountID, input.Platform, input.ShareURL) + app.Get("/api/creator/competitor-share-jobs", func(c fiber.Ctx) error { + items, err := store.ListCompetitorShareJobs(c.Context(), c.Query("platform"), c.Query("status")) if err != nil { return creatorError(c, err) } - return c.JSON(preview) + return c.JSON(items) }) - app.Post("/api/creator/competitors", func(c fiber.Ctx) error { - var request competitorShareRequest + app.Get("/api/creator/competitor-share-jobs/:id", func(c fiber.Ctx) error { + item, err := store.GetCompetitorShareJob(c.Context(), c.Params("id")) + if err != nil { + return creatorError(c, err) + } + return c.JSON(item) + }) + app.Post("/api/creator/competitor-share-jobs", func(c fiber.Ctx) error { + var request creator.CompetitorShareJobInput if err := decodeCreator(c, &request); err != nil { return creatorError(c, err) } - preview, err := previewCompetitorShare(c.Context(), store, phaseAStore, hubStore, request.AccountID, request.Platform, request.ShareURL) + platform, err := competitorSharePlatform(request.ShareURL) if err != nil { return creatorError(c, err) } - input := preview.input() - input.Tags = request.Tags - if err := validateXiaohongshuCompetitor(input); err != nil { - return creatorError(c, err) + if request.Platform != "" && request.Platform != platform { + return creatorError(c, creator.ErrInvalid) } - item, err := store.CreateCompetitor(c.Context(), input) + request.Platform = platform + item, err := store.CreateCompetitorShareJob(c.Context(), request) if err != nil { return creatorError(c, err) } - return c.Status(fiber.StatusCreated).JSON(item) + return c.Status(fiber.StatusAccepted).JSON(item) + }) + app.Post("/api/creator/competitor-share-jobs/:id/retry", func(c fiber.Ctx) error { + item, err := store.RetryCompetitorShareJob(c.Context(), c.Params("id")) + if err != nil { + return creatorError(c, err) + } + return c.Status(fiber.StatusAccepted).JSON(item) }) app.Get("/api/creator/competitors/:id", func(c fiber.Ctx) error { item, err := store.GetCompetitor(c.Context(), c.Params("id")) @@ -1524,13 +1533,6 @@ func (browser creatorGatewayBrowser) Media(ctx context.Context, target, destinat return writeCreatorMedia(destination, data) } -type competitorShareRequest struct { - AccountID string `json:"account_id"` - Platform string `json:"platform"` - ShareURL string `json:"share_url"` - Tags []string `json:"tags"` -} - type competitorSharePreview struct { AccountID string `json:"account_id"` Platform string `json:"platform"` @@ -1952,6 +1954,49 @@ func syncCreatorCompetitorWithClaim(ctx context.Context, store *creator.Store, p return report, nil } +func processCompetitorShareJob(ctx context.Context, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store, jobID string) error { + for { + job, leaseToken, claimed, err := store.ClaimCompetitorShareJob(ctx, jobID, time.Now().UTC()) + if err != nil { + return err + } + if !claimed { + return nil + } + + preview, processErr := previewCompetitorShare(ctx, store, phaseAStore, hubStore, "", job.Platform, job.ShareURL) + competitorID := "" + if processErr == nil { + input := preview.input() + input.Tags = job.Tags + if processErr = validateXiaohongshuCompetitor(input); processErr == nil { + var competitor creator.Competitor + competitor, processErr = store.UpsertCompetitor(ctx, input) + competitorID = competitor.ID + } + } + if processErr == nil { + if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, creator.CompetitorShareJobSucceeded, competitorID, ""); err != nil { + return err + } + logrus.WithFields(logrus.Fields{"job_id": job.ID, "competitor_id": competitorID, "attempts": job.Attempts}).Info("creator competitor share job completed") + return nil + } + + status := creator.CompetitorShareJobQueued + if job.Attempts >= creator.MaxCompetitorShareJobAttempts { + status = creator.CompetitorShareJobFailed + } + if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, status, "", processErr.Error()); err != nil { + return errors.Join(processErr, err) + } + logrus.WithFields(logrus.Fields{"job_id": job.ID, "attempts": job.Attempts, "status": status}).WithError(processErr).Warn("creator competitor share job attempt failed") + if status == creator.CompetitorShareJobFailed || ctx.Err() != nil { + return nil + } + } +} + func runCreatorScheduleOnce(ctx context.Context, store *creator.Store, phaseAStore *accountdomain.Store, hubStore *hub.Store) error { if store == nil { return creator.ErrUnavailable @@ -1961,6 +2006,15 @@ func runCreatorScheduleOnce(ctx context.Context, store *creator.Store, phaseASto return err } now := time.Now().UTC() + shareJobs, err := store.ListDueCompetitorShareJobs(ctx, now) + if err != nil { + return err + } + for _, job := range shareJobs { + if err := processCompetitorShareJob(ctx, store, phaseAStore, hubStore, job.ID); err != nil { + logrus.WithError(err).WithField("job_id", job.ID).Warn("creator competitor share job failed") + } + } competitors, err := store.ListDueCompetitors(ctx, now) if err != nil { return err diff --git a/internal/controlplane/api/creator_route_validation_test.go b/internal/controlplane/api/creator_route_validation_test.go index b492638..b5e8add 100644 --- a/internal/controlplane/api/creator_route_validation_test.go +++ b/internal/controlplane/api/creator_route_validation_test.go @@ -52,8 +52,7 @@ func TestCreatorWriteRoutesRejectMalformedInputBeforeStoreAccess(t *testing.T) { {method: http.MethodPost, path: "/api/creator/relations"}, {method: http.MethodPost, path: "/api/creator/accounts/account-1/strategies"}, {method: http.MethodPut, path: "/api/creator/strategies/strategy-1"}, - {method: http.MethodPost, path: "/api/creator/competitors/preview"}, - {method: http.MethodPost, path: "/api/creator/competitors"}, + {method: http.MethodPost, path: "/api/creator/competitor-share-jobs"}, {method: http.MethodPut, path: "/api/creator/competitors/competitor-1"}, {method: http.MethodPost, path: "/api/creator/competitors/competitor-1/sync"}, {method: http.MethodPost, path: "/api/creator/xiaohongshu/search"}, diff --git a/internal/creator/competitor_share_jobs.go b/internal/creator/competitor_share_jobs.go new file mode 100644 index 0000000..91288f8 --- /dev/null +++ b/internal/creator/competitor_share_jobs.go @@ -0,0 +1,233 @@ +package creator + +import ( + "context" + "database/sql" + "errors" + "net/url" + "strconv" + "strings" + "time" + "unicode/utf8" + + "github.com/jackc/pgx/v5/pgtype" +) + +const ( + CompetitorShareJobQueued = "queued" + CompetitorShareJobProcessing = "processing" + CompetitorShareJobSucceeded = "succeeded" + CompetitorShareJobFailed = "failed" + + MaxCompetitorShareJobAttempts = 3 +) + +func validateCompetitorShareJobInput(input CompetitorShareJobInput) error { + if !ValidatePlatform(input.Platform) || input.ShareURL == "" || utf8.RuneCountInString(input.ShareURL) > 2000 { + return ErrInvalid + } + parsed, err := url.Parse(input.ShareURL) + if err != nil || parsed.Scheme != "https" || parsed.Hostname() == "" || parsed.User != nil || parsed.Port() != "" || parsed.Fragment != "" { + return ErrInvalid + } + return validateCreatorTags(input.Tags) +} + +func normalizeCompetitorShareJobInput(input CompetitorShareJobInput) (CompetitorShareJobInput, error) { + input.Platform = strings.TrimSpace(input.Platform) + input.ShareURL = strings.TrimSpace(input.ShareURL) + if input.Tags == nil { + input.Tags = []string{} + } + if err := validateCompetitorShareJobInput(input); err != nil { + return CompetitorShareJobInput{}, err + } + return input, nil +} + +func (s *Store) CreateCompetitorShareJob(ctx context.Context, input CompetitorShareJobInput) (CompetitorShareJob, error) { + input, err := normalizeCompetitorShareJobInput(input) + if err != nil { + return CompetitorShareJob{}, err + } + id := newID("competitor-share-job") + if _, err := s.db.ExecContext(ctx, ` + INSERT INTO creator_competitor_share_job (id, platform, share_url, tags) + VALUES ($1, $2, $3, $4)`, id, input.Platform, input.ShareURL, input.Tags); err != nil { + return CompetitorShareJob{}, databaseError(err) + } + return s.GetCompetitorShareJob(ctx, id) +} + +func scanCompetitorShareJob(scanner interface{ Scan(...any) error }) (CompetitorShareJob, error) { + var result CompetitorShareJob + var tags pgtype.FlatArray[string] + var competitorID sql.NullString + var leaseUntil, lastAttemptAt, completedAt sql.NullTime + if err := scanner.Scan(&result.ID, &result.Platform, &result.ShareURL, pgtype.NewMap().SQLScanner(&tags), + &result.Status, &result.Attempts, &competitorID, &result.FailureReason, &leaseUntil, &lastAttemptAt, + &completedAt, &result.CreatedAt, &result.UpdatedAt); err != nil { + return CompetitorShareJob{}, err + } + result.Tags = []string(tags) + if competitorID.Valid { + result.CompetitorID = competitorID.String + } + result.LastAttemptAt = nullableTime(lastAttemptAt) + result.CompletedAt = nullableTime(completedAt) + return result, nil +} + +const competitorShareJobSelect = `SELECT id, platform, share_url, tags, status, attempts, competitor_id, + failure_reason, lease_until, last_attempt_at, completed_at, created_at, updated_at + FROM creator_competitor_share_job` + +func (s *Store) GetCompetitorShareJob(ctx context.Context, id string) (CompetitorShareJob, error) { + result, err := scanCompetitorShareJob(s.db.QueryRowContext(ctx, competitorShareJobSelect+` WHERE id = $1`, id)) + return result, rowError(err) +} + +func validateCompetitorShareJobStatus(status string) bool { + switch status { + case CompetitorShareJobQueued, CompetitorShareJobProcessing, CompetitorShareJobSucceeded, CompetitorShareJobFailed: + return true + default: + return false + } +} + +func (s *Store) ListCompetitorShareJobs(ctx context.Context, platform, status string) ([]CompetitorShareJob, error) { + if platform != "" && !ValidatePlatform(platform) || status != "" && !validateCompetitorShareJobStatus(status) { + return nil, ErrInvalid + } + query := competitorShareJobSelect + conditions := make([]string, 0, 2) + args := make([]any, 0, 2) + if platform != "" { + args = append(args, platform) + conditions = append(conditions, "platform = $"+strconv.Itoa(len(args))) + } + if status != "" { + args = append(args, status) + conditions = append(conditions, "status = $"+strconv.Itoa(len(args))) + } + if len(conditions) > 0 { + query += ` WHERE ` + strings.Join(conditions, ` AND `) + } + query += ` ORDER BY created_at DESC, id` + rows, err := s.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, databaseError(err) + } + defer rows.Close() + result := make([]CompetitorShareJob, 0) + for rows.Next() { + item, err := scanCompetitorShareJob(rows) + if err != nil { + return nil, err + } + result = append(result, item) + } + return result, rows.Err() +} + +func (s *Store) ListDueCompetitorShareJobs(ctx context.Context, now time.Time) ([]CompetitorShareJob, error) { + if now.IsZero() { + return nil, ErrInvalid + } + rows, err := s.db.QueryContext(ctx, competitorShareJobSelect+` WHERE + (status = 'queued' AND attempts < $1) OR + (status = 'processing' AND lease_until IS NOT NULL AND lease_until <= $2) + ORDER BY created_at, id`, MaxCompetitorShareJobAttempts, now.UTC()) + if err != nil { + return nil, databaseError(err) + } + defer rows.Close() + result := make([]CompetitorShareJob, 0) + for rows.Next() { + item, err := scanCompetitorShareJob(rows) + if err != nil { + return nil, err + } + result = append(result, item) + } + return result, rows.Err() +} + +func (s *Store) ClaimCompetitorShareJob(ctx context.Context, id string, now time.Time) (CompetitorShareJob, string, bool, error) { + if strings.TrimSpace(id) == "" || now.IsZero() { + return CompetitorShareJob{}, "", false, ErrInvalid + } + token := newID("competitor-share-lease") + claimed, err := scanCompetitorShareJob(s.db.QueryRowContext(ctx, ` + UPDATE creator_competitor_share_job + SET status = 'processing', + attempts = CASE WHEN status = 'queued' THEN attempts + 1 ELSE attempts END, + lease_token = $3, + lease_until = $2 + interval '10 minutes', + last_attempt_at = $2, + updated_at = $2 + WHERE id = $1 AND ( + (status = 'queued' AND attempts < $4) OR + (status = 'processing' AND lease_until IS NOT NULL AND lease_until <= $2) + ) + RETURNING id, platform, share_url, tags, status, attempts, competitor_id, + failure_reason, lease_until, last_attempt_at, completed_at, created_at, updated_at`, id, now.UTC(), token, MaxCompetitorShareJobAttempts)) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + return CompetitorShareJob{}, "", false, nil + } + return CompetitorShareJob{}, "", false, databaseError(err) + } + return claimed, token, true, nil +} + +func (s *Store) MarkCompetitorShareJob(ctx context.Context, id, leaseToken, status, competitorID, failureReason string) error { + if strings.TrimSpace(id) == "" || strings.TrimSpace(leaseToken) == "" || + (status != CompetitorShareJobQueued && status != CompetitorShareJobSucceeded && status != CompetitorShareJobFailed) { + return ErrInvalid + } + if status == CompetitorShareJobSucceeded && strings.TrimSpace(competitorID) == "" || status != CompetitorShareJobSucceeded && competitorID != "" { + return ErrInvalid + } + result, err := s.db.ExecContext(ctx, ` + UPDATE creator_competitor_share_job + SET status = $3, competitor_id = $4, failure_reason = $5, + lease_token = '', lease_until = NULL, + completed_at = CASE WHEN $3 = 'queued' THEN NULL ELSE now() END, + updated_at = now() + WHERE id = $1 AND lease_token = $2 AND status = 'processing'`, id, leaseToken, status, nullableString(competitorID), failureReason) + if err != nil { + return databaseError(err) + } + if affected, err := result.RowsAffected(); err != nil { + return databaseError(err) + } else if affected != 1 { + return ErrConflict + } + return nil +} + +func (s *Store) RetryCompetitorShareJob(ctx context.Context, id string) (CompetitorShareJob, error) { + id = strings.TrimSpace(id) + if id == "" { + return CompetitorShareJob{}, ErrInvalid + } + result, err := s.db.ExecContext(ctx, ` + UPDATE creator_competitor_share_job + SET status = 'queued', attempts = 0, competitor_id = NULL, + lease_token = '', lease_until = NULL, completed_at = NULL, updated_at = now() + WHERE id = $1 AND status = 'failed'`, id) + if err != nil { + return CompetitorShareJob{}, databaseError(err) + } + if affected, err := result.RowsAffected(); err != nil { + return CompetitorShareJob{}, databaseError(err) + } else if affected != 1 { + if _, getErr := s.GetCompetitorShareJob(ctx, id); getErr != nil { + return CompetitorShareJob{}, getErr + } + return CompetitorShareJob{}, ErrConflict + } + return s.GetCompetitorShareJob(ctx, id) +} diff --git a/internal/creator/competitor_share_jobs_test.go b/internal/creator/competitor_share_jobs_test.go new file mode 100644 index 0000000..4c606b1 --- /dev/null +++ b/internal/creator/competitor_share_jobs_test.go @@ -0,0 +1,29 @@ +package creator + +import ( + "errors" + "testing" +) + +func TestNormalizeCompetitorShareJobInput(t *testing.T) { + job, err := normalizeCompetitorShareJobInput(CompetitorShareJobInput{ + Platform: " douyin ", + ShareURL: " https://v.douyin.com/share-a/ ", + }) + if err != nil { + t.Fatal(err) + } + if job.Platform != PlatformDouyin || job.ShareURL != "https://v.douyin.com/share-a/" || job.Tags == nil { + t.Fatalf("normalized job = %+v", job) + } + + for _, input := range []CompetitorShareJobInput{ + {Platform: PlatformDouyin, ShareURL: "http://v.douyin.com/share-a/"}, + {Platform: "unsupported", ShareURL: "https://v.douyin.com/share-a/"}, + {Platform: PlatformDouyin, ShareURL: "https://v.douyin.com/share-a/#fragment"}, + } { + if _, err := normalizeCompetitorShareJobInput(input); !errors.Is(err, ErrInvalid) { + t.Fatalf("input %+v returned err=%v", input, err) + } + } +} diff --git a/internal/creator/competitor_upsert.go b/internal/creator/competitor_upsert.go new file mode 100644 index 0000000..6dd7905 --- /dev/null +++ b/internal/creator/competitor_upsert.go @@ -0,0 +1,27 @@ +package creator + +import "context" + +func (s *Store) UpsertCompetitor(ctx context.Context, input CompetitorInput) (Competitor, error) { + input, err := normalizeCompetitorInput(input) + if err != nil { + return Competitor{}, err + } + id := newID("competitor") + var returnedID string + if err := s.db.QueryRowContext(ctx, ` + INSERT INTO creator_competitor (id, platform, platform_account_key, unique_id, nickname, avatar_url, homepage_url, tags, next_sync_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, now()) + ON CONFLICT (platform, platform_account_key) DO UPDATE SET + unique_id = CASE WHEN EXCLUDED.unique_id = '' THEN creator_competitor.unique_id ELSE EXCLUDED.unique_id END, + nickname = CASE WHEN EXCLUDED.nickname = '' THEN creator_competitor.nickname ELSE EXCLUDED.nickname END, + avatar_url = CASE WHEN EXCLUDED.avatar_url = '' THEN creator_competitor.avatar_url ELSE EXCLUDED.avatar_url END, + homepage_url = EXCLUDED.homepage_url, + tags = CASE WHEN cardinality(EXCLUDED.tags) = 0 THEN creator_competitor.tags ELSE EXCLUDED.tags END, + updated_at = now() + RETURNING id`, id, input.Platform, input.PlatformAccountKey, input.UniqueID, input.Nickname, + input.AvatarURL, input.HomepageURL, input.Tags).Scan(&returnedID); err != nil { + return Competitor{}, databaseError(err) + } + return s.GetCompetitor(ctx, returnedID) +} diff --git a/internal/creator/content.go b/internal/creator/content.go index 82a0c31..8c633e5 100644 --- a/internal/creator/content.go +++ b/internal/creator/content.go @@ -36,7 +36,7 @@ func validateCreatorTags(tags []string) error { return nil } -func (s *Store) CreateCompetitor(ctx context.Context, input CompetitorInput) (Competitor, error) { +func normalizeCompetitorInput(input CompetitorInput) (CompetitorInput, error) { input.Platform = strings.TrimSpace(input.Platform) input.PlatformAccountKey = strings.TrimSpace(input.PlatformAccountKey) input.UniqueID = strings.TrimSpace(input.UniqueID) @@ -50,7 +50,15 @@ func (s *Store) CreateCompetitor(ctx context.Context, input CompetitorInput) (Co utf8.RuneCountInString(input.PlatformAccountKey) > 255 || utf8.RuneCountInString(input.UniqueID) > 255 || utf8.RuneCountInString(input.Nickname) > 255 || utf8.RuneCountInString(input.AvatarURL) > 1000 || validateHomepage(input.HomepageURL) != nil || validateCreatorTags(input.Tags) != nil { - return Competitor{}, ErrInvalid + return CompetitorInput{}, ErrInvalid + } + return input, nil +} + +func (s *Store) CreateCompetitor(ctx context.Context, input CompetitorInput) (Competitor, error) { + input, err := normalizeCompetitorInput(input) + if err != nil { + return Competitor{}, err } id := newID("competitor") if _, err := s.db.ExecContext(ctx, ` diff --git a/internal/creator/integration_test.go b/internal/creator/integration_test.go index 83a26b0..46489c6 100644 --- a/internal/creator/integration_test.go +++ b/internal/creator/integration_test.go @@ -131,6 +131,71 @@ func createIntegrationAccount(t *testing.T, ctx context.Context, phaseAStore *ac return id } +func TestCompetitorShareJobLifecycle(t *testing.T) { + store, _, ctx := openCreatorIntegrationStore(t) + stamp := fmt.Sprintf("%d", time.Now().UnixNano()) + job, err := store.CreateCompetitorShareJob(ctx, CompetitorShareJobInput{ + Platform: PlatformDouyin, + ShareURL: "https://v.douyin.com/share-" + stamp + "/", + Tags: []string{"重点监测"}, + }) + if err != nil || job.Status != CompetitorShareJobQueued || job.Attempts != 0 { + t.Fatalf("create share job: job=%+v err=%v", job, err) + } + due, err := store.ListDueCompetitorShareJobs(ctx, time.Now().UTC()) + if err != nil || len(due) != 1 || due[0].ID != job.ID { + t.Fatalf("list due share jobs: jobs=%+v err=%v", due, err) + } + claimed, leaseToken, ok, err := store.ClaimCompetitorShareJob(ctx, job.ID, time.Now().UTC()) + if err != nil || !ok || leaseToken == "" || claimed.Attempts != 1 || claimed.Status != CompetitorShareJobProcessing { + t.Fatalf("claim share job: job=%+v token=%q claimed=%v err=%v", claimed, leaseToken, ok, err) + } + if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, CompetitorShareJobQueued, "", "temporary failure"); err != nil { + t.Fatal(err) + } + claimed, leaseToken, ok, err = store.ClaimCompetitorShareJob(ctx, job.ID, time.Now().UTC()) + if err != nil || !ok || leaseToken == "" || claimed.Attempts != 2 { + t.Fatalf("claim retry share job: job=%+v token=%q claimed=%v err=%v", claimed, leaseToken, ok, err) + } + if err := store.MarkCompetitorShareJob(ctx, job.ID, leaseToken, CompetitorShareJobFailed, "", "final failure"); err != nil { + t.Fatal(err) + } + failed, err := store.GetCompetitorShareJob(ctx, job.ID) + if err != nil || failed.Status != CompetitorShareJobFailed || failed.Attempts != 2 || failed.FailureReason != "final failure" { + t.Fatalf("failed share job: job=%+v err=%v", failed, err) + } + requeued, err := store.RetryCompetitorShareJob(ctx, job.ID) + if err != nil || requeued.Status != CompetitorShareJobQueued || requeued.Attempts != 0 || requeued.FailureReason != "final failure" { + t.Fatalf("requeue share job: job=%+v err=%v", requeued, err) + } +} + +func TestCompetitorUpsertIsIdempotent(t *testing.T) { + store, _, ctx := openCreatorIntegrationStore(t) + stamp := fmt.Sprintf("%d", time.Now().UnixNano()) + input := CompetitorInput{ + Platform: PlatformDouyin, + PlatformAccountKey: "sec_uid_upsert_" + stamp, + UniqueID: "upsert_" + stamp, + Nickname: "Competitor", + HomepageURL: "https://www.douyin.com/user/sec_uid_upsert_" + stamp, + Tags: []string{"重点监测"}, + } + first, err := store.UpsertCompetitor(ctx, input) + if err != nil { + t.Fatal(err) + } + second, err := store.UpsertCompetitor(ctx, CompetitorInput{ + Platform: input.Platform, + PlatformAccountKey: input.PlatformAccountKey, + Nickname: "Updated Competitor", + HomepageURL: input.HomepageURL, + }) + if err != nil || second.ID != first.ID || second.Nickname != "Updated Competitor" || len(second.Tags) != 1 { + t.Fatalf("upsert competitor: first=%+v second=%+v err=%v", first, second, err) + } +} + func TestCreatorPostgresContentAndWorkflow(t *testing.T) { store, phaseAStore, ctx := openCreatorIntegrationStore(t) stamp := fmt.Sprintf("%d", time.Now().UnixNano()) diff --git a/internal/creator/migrations/037_competitor_share_jobs.sql b/internal/creator/migrations/037_competitor_share_jobs.sql new file mode 100644 index 0000000..056d7bd --- /dev/null +++ b/internal/creator/migrations/037_competitor_share_jobs.sql @@ -0,0 +1,18 @@ +CREATE TABLE IF NOT EXISTS creator_competitor_share_job ( + id text PRIMARY KEY, + platform text NOT NULL CHECK (platform IN ('douyin', 'xiaohongshu')), + share_url text NOT NULL, + tags text[] NOT NULL DEFAULT '{}', + status text NOT NULL DEFAULT 'queued' CHECK (status IN ('queued', 'processing', 'succeeded', 'failed')), + attempts integer NOT NULL DEFAULT 0 CHECK (attempts >= 0 AND attempts <= 3), + competitor_id text REFERENCES creator_competitor(id) ON DELETE SET NULL, + failure_reason text NOT NULL DEFAULT '', + lease_token text NOT NULL DEFAULT '', + lease_until timestamptz, + last_attempt_at timestamptz, + completed_at timestamptz, + created_at timestamptz NOT NULL DEFAULT now(), + updated_at timestamptz NOT NULL DEFAULT now() +); +CREATE INDEX IF NOT EXISTS creator_competitor_share_job_due_idx + ON creator_competitor_share_job (status, lease_until, created_at); diff --git a/internal/creator/models.go b/internal/creator/models.go index 09b5bbb..ef800d7 100644 --- a/internal/creator/models.go +++ b/internal/creator/models.go @@ -133,6 +133,27 @@ type CompetitorInput struct { Tags []string `json:"tags"` } +type CompetitorShareJob struct { + ID string `json:"id"` + Platform string `json:"platform"` + ShareURL string `json:"share_url"` + Tags []string `json:"tags"` + Status string `json:"status"` + Attempts int `json:"attempts"` + CompetitorID string `json:"competitor_id,omitempty"` + FailureReason string `json:"failure_reason,omitempty"` + LastAttemptAt *time.Time `json:"last_attempt_at,omitempty"` + CompletedAt *time.Time `json:"completed_at,omitempty"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +type CompetitorShareJobInput struct { + Platform string `json:"platform"` + ShareURL string `json:"share_url"` + Tags []string `json:"tags"` +} + type WorkSource struct { Platform string `json:"platform"` SourceType string `json:"source_type"` diff --git a/internal/creator/store.go b/internal/creator/store.go index 8b9689c..82e4077 100644 --- a/internal/creator/store.go +++ b/internal/creator/store.go @@ -80,6 +80,9 @@ var migration035 string //go:embed migrations/036_competitor_unique_id.sql var migration036 string +//go:embed migrations/037_competitor_share_jobs.sql +var migration037 string + type SecretReference struct { ID string Provider string @@ -177,6 +180,7 @@ func (s *Store) migrate(ctx context.Context) error { {version: 34, sql: migration034}, {version: 35, sql: migration035}, {version: 36, sql: migration036}, + {version: 37, sql: migration037}, } for _, migration := range migrations { var applied bool diff --git a/web/src/app/Layout.jsx b/web/src/app/Layout.jsx index d872db2..7d21ac7 100644 --- a/web/src/app/Layout.jsx +++ b/web/src/app/Layout.jsx @@ -23,7 +23,9 @@ const menu = [ { section: '运营', items: [ - { to: '/accounts', icon: 'ri-account-circle-line', label: '账号管理' }, + { to: '/accounts/import', icon: 'ri-upload-cloud-2-line', label: '账号导入' }, + { to: '/accounts', icon: 'ri-account-circle-line', label: '我的账号' }, + { to: '/accounts/monitoring', icon: 'ri-eye-line', label: '监控账号' }, { to: '/creator/competitors', icon: 'ri-line-chart-line', label: '竞品分析' }, { to: '/creator/workbench', icon: 'ri-chat-3-line', label: '运营工作台' }, { to: '/creator/settings', icon: 'ri-settings-3-line', label: '采集设置' }, @@ -113,11 +115,13 @@ function AppSidebar({ pathname }) { } const pageMetadataRules = [ - { test: (path) => path === '/accounts', title: '账号管理', subtitle: '统一管理自有账号与监测账号,按账号类型筛选。' }, + { test: (path) => path === '/accounts/import', title: '账号导入', subtitle: '管理分享链接解析任务,解析成功后自动加入监控账号。' }, + { test: (path) => path === '/accounts', title: '我的账号', subtitle: '创建和管理自己维护的账号。' }, + { test: (path) => path === '/accounts/monitoring', title: '监控账号', subtitle: '管理需要持续跟踪的竞品账号及采集状态。' }, { test: (path) => path === '/accounts/new', title: '创建社媒账号', subtitle: '创建账号后,再在编辑页配置登录身份与账号策略。' }, { test: (path) => path.endsWith('/edit') && path.startsWith('/accounts/'), title: '编辑社媒账号', subtitle: '维护账号资料、登录核验与自动响应策略。' }, { test: (path) => path.startsWith('/accounts/'), title: '账号详情', subtitle: '查看账号状态、登录身份与运行环境绑定。' }, - { test: (path) => path === '/creator/competitors', title: '竞品分析', subtitle: '粘贴作品分享内容,自动识别作者并加入监听队列。' }, + { test: (path) => path === '/creator/competitors', title: '竞品分析', subtitle: '查看竞品作品与指标,分享链接导入请前往账号导入。' }, { test: (path) => path === '/creator/workbench', title: '运营工作台', subtitle: '评论、线索、私信与写操作均保留来源和明确结果。' }, { test: (path) => path === '/creator/settings', title: '采集设置', subtitle: '统一配置采集窗口、指标采集和已批准的服务。' }, { test: (path) => path === '/competitors' || path.startsWith('/competitors/'), title: '竞品分析', subtitle: '粘贴作品分享内容,自动识别作者并加入监听队列。' }, diff --git a/web/src/app/Layout.test.jsx b/web/src/app/Layout.test.jsx index 0b2d70e..509e593 100644 --- a/web/src/app/Layout.test.jsx +++ b/web/src/app/Layout.test.jsx @@ -18,7 +18,7 @@ afterEach(() => { }); describe("CreatorHubLayout", () => { - it("exposes one merged account-management menu item", () => { + it("exposes three account-management menu items", () => { Object.defineProperty(window, "matchMedia", { configurable: true, value: () => ({ matches: true }), @@ -33,17 +33,18 @@ describe("CreatorHubLayout", () => { , ); - expect( - screen.getByRole("link", { name: "账号管理" }).getAttribute("href"), - ).toBe("/accounts"); - expect(screen.queryByRole("link", { name: "竞品账号" })).toBeNull(); - expect(screen.queryByRole("link", { name: "账号策略" })).toBeNull(); - expect(screen.queryByRole("link", { name: "浏览器版本" })).toBeNull(); - expect(screen.getByRole("heading", { name: "账号管理" })).toBeTruthy(); - expect( - screen.getByText("统一管理自有账号与监测账号,按账号类型筛选。"), - ).toBeTruthy(); - expect(screen.getByRole("link", { name: "账号管理" }).querySelector("i")) + expect(screen.getByRole("link", { name: "账号导入" }).getAttribute("href")).toBe( + "/accounts/import", + ); + expect(screen.getByRole("link", { name: "我的账号" }).getAttribute("href")).toBe( + "/accounts", + ); + expect(screen.getByRole("link", { name: "监控账号" }).getAttribute("href")).toBe( + "/accounts/monitoring", + ); + expect(screen.getByRole("heading", { name: "我的账号" })).toBeTruthy(); + expect(screen.getByText("创建和管理自己维护的账号。")).toBeTruthy(); + expect(screen.getByRole("link", { name: "我的账号" }).querySelector("i")) .not.toBeNull(); }); diff --git a/web/src/features/accounts/AccountImportPage.jsx b/web/src/features/accounts/AccountImportPage.jsx new file mode 100644 index 0000000..369f4bb --- /dev/null +++ b/web/src/features/accounts/AccountImportPage.jsx @@ -0,0 +1,295 @@ +import { useEffect, useState } from "react"; +import { useDataProvider } from "../../shared/hooks/dataHooks.js"; +import { Alert } from "../../components/ui/alert"; +import { Button } from "../../components/ui/button"; +import { Card, CardContent } from "../../components/ui/card"; +import { Dialog, DialogContent, DialogFooter, DialogHeader, DialogTitle } from "../../components/ui/dialog"; +import { Table, TableBody, TableCell, TableHead, TableHeader, TableRow } from "../../components/ui/table"; +import { Textarea } from "../../components/ui/textarea"; +import { Field, PageHeader, PageState, StatusPill, TagInput, conflictMessage, dateTime } from "../../shared/ui/ui.jsx"; +import { useUnsavedChanges } from "../../shared/hooks/hooks.js"; + +const platformOptions = [ + { value: "douyin", label: "抖音" }, + { value: "xiaohongshu", label: "小红书" }, +]; + +const shareJobStatus = { + queued: { label: "待解析", tone: "warning" }, + processing: { label: "解析中", tone: "warning" }, + succeeded: { label: "已加入监控", tone: "success" }, + failed: { label: "解析失败", tone: "danger" }, +}; + +const platformLabel = (value) => + platformOptions.find((option) => option.value === value)?.label || value || "—"; + +const shareURLPattern = /https?:\/\/[^\s<>"'`]+/giu; +const trailingShareURLPunctuation = /[,。;!?、)》)】\]}>,.?!:;]+$/u; + +export function extractShareURL(value) { + const candidates = String(value || "") + .match(shareURLPattern) + ?.map((candidate) => candidate.replace(trailingShareURLPunctuation, "")) + .filter(Boolean); + return candidates?.find((candidate) => platformForShareURL(candidate)) || candidates?.[0] || ""; +} + +export function platformForShareURL(value) { + try { + const hostname = new URL(value).hostname.toLowerCase(); + if (["www.douyin.com", "v.douyin.com"].includes(hostname)) { + return "douyin"; + } + if (["www.xiaohongshu.com", "xhslink.com", "www.xhslink.com"].includes(hostname)) { + return "xiaohongshu"; + } + } catch { + return ""; + } + return ""; +} + +export function AccountImportPage() { + const dataProvider = useDataProvider()("default"); + const [createOpen, setCreateOpen] = useState(false); + const [shareText, setShareText] = useState(""); + const [shareTags, setShareTags] = useState([]); + const [shareJobs, setShareJobs] = useState([]); + const [shareJobsPending, setShareJobsPending] = useState(true); + const [shareJobsError, setShareJobsError] = useState(null); + const [shareJobActionID, setShareJobActionID] = useState(null); + const [busy, setBusy] = useState(false); + const [notice, setNotice] = useState(null); + + const loadShareJobs = async () => { + setShareJobsPending(true); + setShareJobsError(null); + try { + const result = await dataProvider.getList({ + resource: "creator-competitor-share-jobs", + filters: [], + }); + setShareJobs(result.data ?? []); + } catch (loadError) { + setShareJobsError(loadError); + } finally { + setShareJobsPending(false); + } + }; + + useEffect(() => { + loadShareJobs(); + }, [dataProvider]); + + const extractedShareURL = extractShareURL(shareText); + const sharePlatform = platformForShareURL(extractedShareURL); + + const closeCreate = () => { + if (busy) return; + setCreateOpen(false); + setShareText(""); + setShareTags([]); + setNotice(null); + }; + + const create = async () => { + if (!extractedShareURL) { + setNotice({ variant: "default", text: "请粘贴包含分享链接的内容。" }); + return; + } + if (!sharePlatform) { + setNotice({ variant: "default", text: "链接不是支持的平台分享链接。" }); + return; + } + setBusy(true); + setNotice(null); + try { + const variables = { + platform: sharePlatform, + share_url: extractedShareURL, + }; + if (shareTags.length) variables.tags = shareTags; + await dataProvider.create({ + resource: "creator-competitor-share-jobs", + variables, + }); + await loadShareJobs(); + setCreateOpen(false); + setShareText(""); + setShareTags([]); + setNotice({ variant: "default", text: "分享链接已加入作者解析队列。" }); + } catch (createError) { + setNotice({ + variant: "destructive", + text: conflictMessage(createError, "分享链接入队失败"), + }); + } finally { + setBusy(false); + } + }; + + const retryShareJob = async (jobID) => { + setShareJobActionID(jobID); + setNotice(null); + try { + await dataProvider.creatorAction( + `/creator/competitor-share-jobs/${encodeURIComponent(jobID)}/retry`, + ); + await loadShareJobs(); + setNotice({ variant: "default", text: "任务已重新加入解析队列。" }); + } catch (retryError) { + setNotice({ + variant: "destructive", + text: conflictMessage(retryError, "任务重新入队失败"), + }); + } finally { + setShareJobActionID(null); + } + }; + + useUnsavedChanges(Boolean(shareText || shareTags.length)); + + return ( + <> + +