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 (
+ <>
+
+ 待处理和失败的分享链接独立展示;成功后才会进入竞品作者监控。
+ {dateTime(job.created_at)}分享链接解析任务
+
+
+
- 自有 {accounts.length} 个 · 监测 {competitors.length} 个 · 当前显示{" "} - {visibleRows.length} 个 +
+ 共 {rows.length} 个{isMonitoring ? "监控" : "自有"}账号