diff --git a/CHANGELOG.md b/CHANGELOG.md index 9c0f3f6..035e3cb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,24 @@ and this project adheres to [Semantic Versioning](https://semver.org/). ## [Unreleased] +### Changed +- 会话列表默认启用 5 秒 WS 发现,补齐最近消息并按活动时间倒序; + 缺失时间和未知未读数返回 null,列表完整性与最近消息上下文状态分别说明。 +- 历史消息保留 CID、消息 ID 和创建时间;支持 limit 限定最近条数及 timeout, + 按服务端消息 ID 去重,缺少 ID 的记录不按正文或时间推断为重复。 + +### Fixed +- 多图发布仅最终提交计一次商品预算,媒体上传独立计数;已上传收据、已确认类目和地址可复用,失败保留准备信息与提交状态。 +- 本机限流和熔断按账号持久化,跨进程事务加锁、原子保存,损坏状态拒绝写入;风控保存失败保留原始错误。 +- 搜索从真实页面响应读取语义字段,保留未知标签,区分昵称、地域与原价;无命中时排除推荐卡片,分页响应绑定查询与页码,部分失败在 CLI/MCP 明确报告。 +- 商品详情从 itemDO/sellerDO 提取,保留价格和状态 0;类目分数对应选中分类,发布地址使用 selectedPoi。 +- 单聊结合 userInfo 与 ownerInfo 及当前账号确认对端,避免买卖家角色对换时 + 选到自己;WS 会话的角色信息与消息参加者共同用于身份确认。 +- MCP 输出 schema 与统一的 ok/data 或错误对象一致,避免历史工具返回 list + 的声明与适配器实际对象冲突。 +- 消息读取在超时、连接失败或分页游标不推进时明确报告失败;列表保留已有 + 结果并标注未完成的摘要,取消时清理连接和心跳。 + ## [0.4.0] - 2026-09-07 ### Added diff --git a/README.md b/README.md index 10aab61..b4046f9 100644 --- a/README.md +++ b/README.md @@ -66,7 +66,7 @@ - 🔐 **17 个命令覆盖核心链路**:发布、下架、查询、图片上传、AI 类目识别、默认地址、IM 收发 + 会话列表、skills 安装 - 📡 **真·实时 IM**:WebSocket 长连 + 自动重连 + **三类事件分类输出** - `event=message`(收到消息)· `event=read`(已读回执)· `event=new_msg`(轻量通知) -- 🛡 **内置风控护栏**:令牌桶限流(1 写/分钟)+ RGV587 自动熔断 +- 🛡 **内置风控护栏**:账号与业务桶限流(经营 1 次/分钟、媒体 9 次/分钟)+ RGV587 自动熔断 - 🧠 **AI-first I/O**:`--format json/yaml/table/md/csv`,给 LLM 喂 JSON、给人看表格 - ⚡ **一次定义,三种入口**:CLI / MCP / Skill 共享同一 registry - ✅ **真实端到端验证**:每个命令都跑过真实账号 @@ -180,7 +180,7 @@ $ goofish list-commands --format table | `media upload` | 上传图片到闲鱼 CDN | ✅ | | `category recommend` | AI 识别商品类目 | ❌ | | `location default` | 获取默认发布地址 | ❌ | -| `message list-chats` | 拉取会话列表(左栏;`--watch-secs N` 叠加 WS 历史推送补漏) | ❌ | +| `message list-chats` | 会话列表(默认 HTTP + 短时 WS,补最近消息并按活动时间倒序) | ❌ | | `search items` | 搜索闲鱼商品(浏览器路径 Playwright + 系统 Chrome) | ❌ | | `item view` | 浏览器视角看商品详情(字段完整,抗风控;`item get` 的姊妹版) | ❌ | | `message history` | 拉取会话历史消息 | ❌ | @@ -189,6 +189,22 @@ $ goofish list-commands --format table +会话读取使用真实消息时间,而非会话创建时间: + +```bash +goofish message list-chats --format json +goofish message history --limit 20 --timeout 30 --format json +``` + +`list-chats` 默认 `--watch-secs 5`,在一个已就绪的连接中补齐真人会话的最近消息页; +`--watch-secs 0` 关闭 WS 会话发现,结果可能遗漏 HTTP 未返回的会话。 +`--timeout` 限制摘要读取阶段。历史输出保留 `cid`、`message_id` 和毫秒时间 `created_at`, +`--limit 0` 保留完整翻页行为,其他正数只读最近 N 条。 + +`metadata_status` 标明摘要是当前、为空、不完整或读取失败;无法确定的时间为 `null`。 +`has_more` 只代表 HTTP 分页,`enumeration_complete=false` 表明短时同步不能证明全部历史会话已覆盖。 +接口契约与异常边界见 [会话同步](docs/conversation-sync.md)。 +
goofish auth status — 登录态健康检查 @@ -245,25 +261,21 @@ $ goofish message send \ goofish item publish — 发布商品(含风控护栏) ```bash -$ goofish item publish \ - --title "男士毛呢大衣 驼色长款" \ - --desc "全新未拆封 原价 2999 现 999" \ - --images ./a.png,./b.png \ - --price 999 +goofish item publish "男士毛呢大衣 驼色长款" \ + "全新未拆封 原价 2999 现 999" ./a.png ./b.png 999 --format json ``` -流程: -1. `media upload` 每张图 → CDN URL + 尺寸 -2. `category recommend` 拿 AI 识别的 catId -3. `location default` 拿默认地址 -4. `mtop.idle.pc.idleitem.publish` 落库 +流程:上传图片 → 推荐类目 → 默认地址 → 提交。已有上传收据可通过 +`--images-json '[{"url":"https://example.alicdn.com/image.png","width":1024,"height":1024}]'` +复用,省略本地图片参数;已确认类目和地址分别通过 `--category-json`、`--location-json` +传入对应工具返回的完整 DTO。 -返回: -```json -{"ok": true, "itemId": "1046118265141", "status": "published"} -``` +返回 `item_id`、`status="accepted"` 和 `requires_readback=true`;使用 +`goofish item get ` 和自己的商品列表确认保存。失败也保留已上传图片和准备信息; +`submission_unknown` 必须先核对商品列表,不自动重复发布。 -**触发令牌桶限流**(1 写/分钟)。高频调用会被本地拒绝,避免被闲鱼风控。 +发布与下架共用账号的 `item.write` 预算,默认 1 次/分钟;媒体上传独立计数,默认 +9 次/分钟。一个多图商品只消耗一次商品预算。
--- @@ -297,7 +309,7 @@ Claude 会自动把全部命令看成 tool:`goofish_item_get` / `goofish_item_ | WebSocket 批量 push 全量解码 | 一帧多条消息全部还原,不丢单 | | WebSocket 自动重连 | 断线自退避重连,长跑无感知 | | 已读回执 / typing / 新消息通知分类 | `/s/sync` 元事件结构化为三类 JSONL | -| 全局限流 + 风控熔断 | 令牌桶 1 写/分钟 + RGV587 自动熔断 | +| 全局限流 + 风控熔断 | 账号业务预算 + RGV587 自动熔断 | | 单元测试 | 33 个,ruff 零告警 | | 包分发 | `pip install goofish-cli` / `uvx goofish-cli` | diff --git a/docs/architecture.md b/docs/architecture.md index e438beb..bd23093 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -37,7 +37,7 @@ │ core/sign.py pyexecjs → goofish_js_version_2.js │ │ core/session.py cookie 加载 + requests.Session │ │ core/mtop.py 统一 mtop 调用 + 错误分类 │ -│ core/limiter.py 令牌桶限流(1 写/分钟,可配) │ +│ core/limiter.py 账号业务桶限流(经营/媒体分开) │ │ core/guard.py 风控熔断(RGV587 → trip) │ │ core/output.py 统一渲染 json/yaml/table/md/csv │ │ core/errors.py 异常体系 + exit_code │ @@ -69,21 +69,26 @@ def get(item_id: str) -> dict: - **namespace + name** → CLI 路径 `goofish item get`,MCP tool 名 `item_get` - **columns** → 输出契约(table/csv/md 场景的列顺序) - **strategy** → 认证要求(PUBLIC / COOKIE / WS) -- **write=True** → 自动触发限流 + 风控熔断 +- **write=True** → 写操作元数据;实际写端点必须进入 `write_operation(session, bucket)` ## 风控护栏 -1. **令牌桶**(`core/limiter.py`):默认 1 写/分钟,持久化在 `~/.goofish-cli/limiter.json` -2. **熔断**(`core/guard.py`):`watch()` 上下文内命中 `RiskControlError` → 写 `circuit.json`,默认熔断 10 分钟 -3. **响应体识别**(`core/mtop.py`):自动扫 `RGV587_ERROR / punish / FAIL_SYS_USER_VALIDATE` 等关键字 +1. **滑动窗口预算**(`core/limiter.py`):按账号和业务桶计数。发布与下架共用 `item.write`,消息用 `message.write`,默认各 1 次/分钟;上传用独立 `media.write`,默认 9 次/分钟。经营预算通过 `GOOFISH_WRITE_RPM` 配置,媒体通过 `GOOFISH_MEDIA_WRITE_RPM` 配置。这是本机预算,不代表平台额度。 +2. **熔断**(`core/guard.py`):所有写端点先检查账号熔断;命中 `RiskControlError` 后持久化,默认 10 分钟。旧版共享熔断和预算须自然过期。保存失败保留原始风控错误并提示停止写入。 +3. **持久化**(`core/local_state.py`):稳定锁文件覆盖读、检查和原子替换事务;跨 CLI 进程和 MCP 线程共享。损坏状态拒绝写入,不静默清零。可通过 `GOOFISH_LIMITER_PATH`、`GOOFISH_GUARD_PATH` 选择隔离状态路径。 +4. **响应识别**:MTop 与上传端点根据真实响应确认成功或拒绝。失败尝试不退预算;结果未知时先回读,不自动重发。 -这些护栏在 **底层强制**,Agent 层无法绕过。 +发布可接受本地图片,也可通过 `images_json` 复用 `media_upload` 的完整收据;`category_json` 和 `location_json` 接收相应工具返回的已确认 DTO。准备或提交失败保留收据与准备信息,分别报告 `not_submitted`、`submission_rejected` 或 `submission_unknown`。返回商品 ID 代表 `accepted`,详情与商品列表回读负责确认实际保存。 + +搜索观察网页实际发出的搜索响应,按 `keyword/pageNumber` 与请求发起代次关联;分页前清除上一页状态,忽略迟到响应。只有平台控制字段确认的查询结果进入 `items`;无命中时网页的推荐卡片单独识别。未标明语义的标签保留在 `labels`,品牌、成色等未知属性返回 null,卖家昵称与地域分开。翻页失败保留部分结果,CLI 非零退出,MCP 返回 `ok=false` 和数据及错误。 + +商品详情使用 `itemDO/sellerDO`,不使用埋点推测价格和状态;价格 0 和状态码 0 保留,`defaultPrice` 布尔值只表达议价类型。分类模型分数只从选中分类或分类卡的复合身份获取,不把其他属性分数当置信度,也不把缺失分数当 0。 ## 实现要点 - **`t` 毫秒位**:`int(time.time() * 1000)` 取真实毫秒(避免 `int(time.time()) * 1000` 把末三位抹成 000 的精度陷阱) -- **默认地址**:`commonAddresses[0]` 兜底,参数化接口支持显式指定 `addressId` +- **默认地址**:优先 `selectedPoi`,缺失时使用第一个常用地址;发布使用返回的 `selected` DTO,显式地址通过 `location_json` 提供。 - **风控识别**:扫描响应体关键字(`RGV587_ERROR` / `punish` / `FAIL_SYS_USER_VALIDATE`),命中即熔断 -- **限流**:令牌桶持久化在 `~/.goofish-cli/limiter.json` +- **限流**:滑动窗口预算持久化在 `~/.goofish-cli/limiter.json` - **形态**:CLI + MCP(+ Skill 规划),同一份 registry 输出三种形态 - **命令组织**:每命令一个文件,`@command(...)` 装饰器自注册,参照 opencli diff --git a/docs/conversation-sync.md b/docs/conversation-sync.md new file mode 100644 index 0000000..a4d6dc7 --- /dev/null +++ b/docs/conversation-sync.md @@ -0,0 +1,38 @@ +# 会话同步契约 + +CLI 与 MCP 消费同一注册函数。MCP 适配器返回 `{ok, data}` 或 `{ok, error_type, message}`,输出 schema 声明该包裹对象,业务函数内部返回 list 时也遵守同一契约。`message list-chats` 以 HTTP 会话摘要作为基线,默认收集 5 秒 WS 会话发现,再用一个连接逐个读取真人和未分类会话的最近消息页。已分类系统频道保留其服务端摘要,不调用真人历史接口。 + +## 身份与活动 + +- `session_id` 和历史中的 `cid` 去除 `@goofish` 后缀,作为合并与去重标识。 +- 最近活动取历史消息的 `createAt`,保留为毫秒时间 `created_at` / `ts`。会话创建时间不能替代最近消息时间;未知时间为 `null`,排在已知时间之后。 +- HTTP 单聊同时检查 userInfo 与 ownerInfo,排除当前账号,只有一个候选时选为对端;不能假定 userInfo 永远是对方。 +- WS 会话角色须包含当前账号,且只有一个非当前账号成员,才能据此确定对端。同时核对 HTTP 的对端字段;角色与 HTTP 候选冲突时不赋值。缺少两路身份时再使用消息发送者/接收者,仍有多个候选时不赋值。 +- 商品的 `itemSellerId`、`squadName` 和消息正文不能当作可靠对端昵称。消息中的发送者标签只用于对应的发送者,不能把自己的昵称赋给对端。 +- 同一 CID 合并 HTTP 与 WS 信息,最近消息 ID 和预览来自新读取的消息页,按最近活动倒序输出。 + +## 有界读取与结果状态 + +`message history --limit N` 只读取最近 N 条,按时间正序返回;`limit=0` 完整翻页。输出保留消息 ID、CID 和时间,跨页以消息 ID 去重。游标不推进、响应结构无效或服务端拒绝均报错,不作为空历史成功返回。 + +摘要读取复用一个已鉴权和就绪的 WS 连接,每次请求按 mid 匹配响应,只有一个接收者。超时取消会关闭连接并等待心跳任务退出。历史整体 WS 阶段和列表摘要阶段各受 `--timeout` 约束;HTTP 请求另有基础客户端超时,WS 发现受 `watch_secs + 15` 秒约束。 + +列表保留成功的摘要,逐会话暴露读取失败: + +| 字段 | 含义 | +|---|---| +| `metadata_status` | `current` 最近消息已读取;`partial` 身份、时间或正文未齐;`empty` 当前历史为空;`stale` 摘要读取失败;`http` / `system` 系统摘要 | +| `metadata_missing` | 未能确定的字段,不以默认零值伪装 | +| `metadata_errors` | 摘要或发现阶段的失败及对应会话 | +| `metadata_complete` | 本次真人和未分类会话的最近消息与对端上下文是否完整;不包含未读计数或历史枚举完整性 | +| `unknown_activity_count` | 没有可确认活动时间的条目数 | +| `has_more_scope` | `http`,has_more 仅对应 HTTP 分页 | +| `enumeration_complete` | false;HTTP 和短时 WS 发现不保证覆盖全部历史会话 | + +`--watch-secs 0` 显式关闭 WS 发现,不保证 HTTP 未返回的会话能被列出。WS 新发现会话的 unread 仍可能为 null;即使最近消息与对端上下文完整,也不能据此推导未读数。模型调用方应依据字段状态决定是否继续,而非仅依据退出码判断上下文完整。 + +## 维护与验证 + +消息页读取与解析集中在 `core/message_history.py`;连接、握手、ACK 和会话发现由 `core/ws.py` 提供。列表层只承担会话合并、身份选择、摘要构建与排序,不实现另一套协议接收循环。 + +修改后按 CONTRIBUTING.md 运行现有检查,并使用同一真实会话验证 CLI 默认列表、最近/完整历史和 MCP stdio 调用。对照原始 createAt、消息 ID、正文及当前账号身份;同一会话不得重复,系统会话不能误当真人,缺失或超时须明确暴露。可在隔离副本禁用 WS、摘要补齐或排序,验证各步骤对真实结果的独立贡献。 diff --git a/pyproject.toml b/pyproject.toml index 06704c2..780bd63 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -35,6 +35,7 @@ dependencies = [ "browser-cookie3>=0.20", "playwright>=1.58.0", "cryptography>=41.0", + "filelock>=3.16,<4", ] [project.optional-dependencies] diff --git a/skills/goofish-overview/SKILL.md b/skills/goofish-overview/SKILL.md index 999a061..d07fc25 100644 --- a/skills/goofish-overview/SKILL.md +++ b/skills/goofish-overview/SKILL.md @@ -90,8 +90,7 @@ metadata: Agent 在 goofish 任务里请遵守: 1. **任何写操作前先读 `auth_status`**,避免 token 已过期却继续写 → 白忙一场。 -2. **`item_publish / media_upload / message_send` 都有速率限制**(令牌桶 - 1 写/分钟),不要短时间连发。RGV587 触发后需用户手动 `goofish auth reset-guard`。 +2. **`item_publish / media_upload / message_send` 都有速率限制**(账号业务预算),不要短时间连发。RGV587 触发后需用户手动 `goofish auth reset-guard`。 3. **所有对外发送(发商品、发消息)务必让用户先确认文案**,不要自动提交。 4. **多步任务中途报状态**:上一步完成了什么、下一步准备做什么。闲鱼任务 通常是 3-6 步,用户期望有进度感。 diff --git a/skills/goofish-overview/references/mcp-tools-index.md b/skills/goofish-overview/references/mcp-tools-index.md index de829f1..4f63a4b 100644 --- a/skills/goofish-overview/references/mcp-tools-index.md +++ b/skills/goofish-overview/references/mcp-tools-index.md @@ -20,17 +20,17 @@ | `item_get` | HTTP 视角拉详情(只读) | `item_id` | 无 | | `item_list` | 查看当前账号的在售商品 | `limit?` | 无 | | `item_view` | 浏览器视角拉详情(字段更全、抗风控) | `item_id` | 触发 Playwright 启动 | -| `item_publish` | 发布商品(自动类目+默认地址) | `title, desc, price, image_urls, cat_id?, addr?` | **写操作**,令牌桶 1 写/分钟 | +| `item_publish` | 发布商品(自动类目+默认地址) | `title, desc, price, images?, images_json?, category_json?, location_json?` | **写操作**,item.write 默认 1 次/分钟 | | `item_delete` | 下架/删除商品 | `item_id` | **写操作** + 风控护栏 | -**发布前强依赖**:`category_recommend` 拿 catId、`media_upload` 拿 image_urls、`location_default` 兜底 addr。 +**发布前强依赖**:`category_recommend` 拿 catId、`media_upload` 拿图片收据、`location_default` 兜底 location_json。 ## Message(4 个) | 工具 | 用途 | 典型入参 | 备注 | |---|---|---|---| -| `message_list_chats` | 会话列表(左栏) | `limit?, watch_secs?` | session.sync v3.0 是阉割版,默认会启 WebSocket 增量补齐 cid | -| `message_history` | 某 cid 的历史消息(翻页到底) | `cid` | 拉上下文做意图分类用 | +| `message_list_chats` | 按最近活动排序的会话摘要 | `fetch_num?, watch_secs?, timeout?` | 默认 5 秒 WS 发现并读取最近消息;检查 metadata_status,不能把 has_more 当全部会话完整性 | +| `message_history` | 某 cid 的历史消息 | `cid, limit_per_page?, limit?, timeout?` | limit=0 翻页到底;正数只读最近 N 条,保留 message_id/created_at | | `message_watch` | 常驻 IM 长连接,事件以 JSONL 输出 | `secs?` | 阻塞式;Skill 里慎用,短连接更合适 | | `message_send` | 发消息(text/image) | `cid, text` 或 `cid, image_url, w, h` | **写操作** + 外联词风控 | @@ -41,7 +41,7 @@ | `category_recommend` | AI 识别类目(输入标题+图片返回 catId/catName) | 发布前必跑,类目错放会降权 | | `media_upload` | 上传图片到闲鱼 CDN | 返回 `{url, width, height}`,给 `item_publish` 直接用 | | `search_items` | 搜闲鱼商品(浏览器路径,抗风控) | 诊断限流时用"卖家视角 vs 买家视角"对比 | -| `location_default` | 账号默认发布地址 | `item_publish` 不传 addr 时的兜底 | +| `location_default` | 账号默认发布地址 | `item_publish` 不传 location_json 时的兜底 | ## Skills(1 个,辅助类) @@ -66,12 +66,15 @@ auth_status → message_list_chats → message_history (按 cid) → [LLM 意图 auth_status → search_items (用自家核心词,买家视角) → item_view (拉自家详情) → item_get (历史元数据对比) ``` -## 写操作全局节流 +## 写操作预算与错误恢复 -`item_publish / item_delete / media_upload / message_send` 共用令牌桶: -- 容量 1、每分钟补 1 -- 短时连发会返回 `RATE_LIMITED` -- RGV587(服务端风控)触发后需用户从浏览器重新导 cookie(带 x5sec/mtop_partitioned_detect),`auth_reset_guard` 只解本地熔断 +- 同账号发布和下架使用 `item.write`,消息使用 `message.write`,默认各 1 次/分钟。 +- 上传使用独立 `media.write`,默认 9 次/分钟;上传收据通过 `images_json` 复用。 +- CLI 进程与 MCP 线程共享加锁、原子保存的状态;损坏文件拒绝写入。 +- 本机预算不代表平台额度;平台风控触发账号熔断,`auth_reset_guard` 只改变本机状态。 +- 发布失败保留图片、类目、地址及提交状态;`submission_unknown` 先回读,不能自动重发。 +- 搜索 `complete=false` 是未完成,CLI 非零退出,MCP `ok=false` 仍带部分 `data`。 +- 搜索的品牌、成色未知时是 null;用 labels 查看原始标签,不能按顺序猜属性。空查询命中不含网页推荐。 ## 错误类型映射 diff --git a/skills/goofish-publish-item/SKILL.md b/skills/goofish-publish-item/SKILL.md index fb4bbfb..a3bae6b 100644 --- a/skills/goofish-publish-item/SKILL.md +++ b/skills/goofish-publish-item/SKILL.md @@ -67,8 +67,8 @@ Claude Code / Cursor 调 `mcp__goofish__<逻辑名>`。frontmatter 的 `allowed- ### Step 3 · 类目识别 ``` -调 category_recommend(title=候选标题, images=[首图]) - → 拿到 catId / catName +调 category_recommend(title=候选标题, images_json="[]") + → 拿到完整类目 DTO;已有上传收据时用 images_json=JSON.stringify(收据数组) ``` **结果校验**: @@ -96,8 +96,8 @@ Claude Code / Cursor 调 `mcp__goofish__<逻辑名>`。frontmatter 的 `allowed- ### Step 6 · 图片处理 ``` -for image in images: - 调 media_upload(image) +for image in 本地图片路径: + 调 media_upload(path=image) 收集返回的 {url, width, height} 检查: @@ -130,12 +130,16 @@ for image in images: title=..., desc=..., price=..., - image_urls=[...], - cat_id=..., - addr=...(没有就省略,让服务端用 location_default) + images_json=JSON.stringify([{url, width, height}, ...]), + category_json=JSON.stringify(category_recommend 的完整 DTO), + location_json=JSON.stringify(location_default 的完整 DTO) ) ``` +发布也可用 `images=[本地路径, ...]` 自动上传,不能把 CDN URL 当本地路径。已提前上传时复用 `images_json`,避免重复上传。类目 DTO 必须包含 cat_id/cat_name/channel_cat_id/tb_cat_id,地址采用 selected 字段。缺少身份的类目不能猜补。模型分数不是校准概率,未知值为 null。 + +成功返回 accepted 与 item_id 后,用 item_get 回读确认。失败返回的 images/category/location 可复用;submission_unknown 必须先核对自己的商品列表,禁止盲目重试。发布和下架共用 item.write 默认 1 次/分钟,上传独立 media.write 默认 9 次/分钟。 + ## 反模式(不要做) ❌ **跳过 category_recommend**,自己猜 catId diff --git a/skills/goofish-publish-item/references/category-selection.md b/skills/goofish-publish-item/references/category-selection.md index b6179dd..8f19586 100644 --- a/skills/goofish-publish-item/references/category-selection.md +++ b/skills/goofish-publish-item/references/category-selection.md @@ -13,7 +13,7 @@ category_recommend( ) ``` -返回 `{catId, catName, confidence?}`。 +返回 `{cat_id, cat_name, channel_cat_id, tb_cat_id, confidence, confidence_source}`。confidence 为模型分数,不是校准概率;缺失为 null。 ## 结果处理 @@ -75,4 +75,4 @@ category_recommend( ## 本模块不调其它工具 -只调 `category_recommend`。不需要 `item_get` 也不需要反查——API 给的 catId 直接传给 `item_publish` 就够。 +只调 `category_recommend`。不需要 `item_get` 也不需要反查——将完整类目 DTO 作为 category_json 传给 item_publish;缺少 cat_id/channel_cat_id/tb_cat_id 时不得猜补。 diff --git a/skills/goofish-reply-buyer/SKILL.md b/skills/goofish-reply-buyer/SKILL.md index dc7e524..2700484 100644 --- a/skills/goofish-reply-buyer/SKILL.md +++ b/skills/goofish-reply-buyer/SKILL.md @@ -60,10 +60,13 @@ auth_status → valid=false 直接停,提示用户 auth login ### Step 2 · 拉未读会话 ``` -message_list_chats(limit=20) +message_list_chats(fetch_num=20) ``` -返回会话列表,每项含 cid / 对方昵称 / 最后一条消息预览 / 未读数。 +返回按活动时间倒序的会话列表,每项含 session_id / 对方昵称 / 最后一条消息预览 / 未读数。 +默认短时 WS 会话发现;先检查 metadata_status 和 metadata_errors。身份缺失或摘要过期的会话先补信息,不猜对方身份。 +HTTP 的 has_more 不代表整份会话列表完整,短时同步也不能保证发现全部历史会话。 +WS 条目的 unread=null 表示未知,不能按零忽略,也不能直接与数字比较。 **筛选策略**: - 优先未读数 > 0 的 @@ -74,7 +77,7 @@ message_list_chats(limit=20) ``` 对每个选中的 cid: - message_history(cid=..., limit=20) + message_history(cid=session_id, limit=20) ``` 20 条基本够分类意图了。太长的会话(50+ 条)拉 30 条。 diff --git a/src/goofish_cli/cli.py b/src/goofish_cli/cli.py index fa1dd2e..48e438d 100644 --- a/src/goofish_cli/cli.py +++ b/src/goofish_cli/cli.py @@ -9,6 +9,7 @@ from loguru import logger from goofish_cli.core import GoofishError, iter_commands +from goofish_cli.core.errors import PartialResultError from goofish_cli.core.output import Format, render from goofish_cli.core.registry import Command, discover @@ -51,6 +52,10 @@ def wrapper(**kwargs): fmt = Format(kwargs.pop("format")) try: result = cmd.func(**kwargs) + except PartialResultError as e: + render(e.data, fmt=fmt, columns=cmd.columns or None) + typer.secho(f"[{e.error.get('error_type', 'PartialResultError')}] {e}", fg=typer.colors.RED, err=True) + raise typer.Exit(code=e.exit_code) from e except GoofishError as e: typer.secho(f"[{type(e).__name__}] {e}", fg=typer.colors.RED, err=True) raise typer.Exit(code=e.exit_code) from e diff --git a/src/goofish_cli/commands/category/recommend.py b/src/goofish_cli/commands/category/recommend.py index 785a588..f4fcc42 100644 --- a/src/goofish_cli/commands/category/recommend.py +++ b/src/goofish_cli/commands/category/recommend.py @@ -4,12 +4,34 @@ """ import json +import math +import time from typing import Any +from goofish_cli.commands.media.upload import validate_image_receipt from goofish_cli.core import Session, Strategy, command +from goofish_cli.core.errors import GoofishError from goofish_cli.core.mtop import call +def _selected_score(data: dict, predict: dict) -> tuple[float | None, str | None]: + explicit = predict.get("confidence") + if isinstance(explicit, (float, int)) and not isinstance(explicit, bool) and math.isfinite(explicit): + return float(explicit), "categoryPredictResult.confidence" + values = [] + for card in data.get("cardList") or []: + body = card.get("cardData") or {} + if str(body.get("propertyId")) != "-10000": + continue + for candidate in body.get("valuesList") or []: + keys = ("channelCatId", "tbCatId", "catId") + if all(str(candidate.get(key)) == str(predict[key]) for key in keys if predict.get(key) is not None): + score = candidate.get("score") + if isinstance(score, (float, int)) and not isinstance(score, bool) and math.isfinite(score): + values.append(float(score)) + return (values[0], "classification_card.score") if len(values) == 1 else (None, None) + + @command( namespace="category", name="recommend", @@ -22,8 +44,17 @@ def recommend( images_json: str = "[]", ) -> dict[str, Any]: """images_json 是 JSON 字符串:[{"url":"...","width":1024,"height":1024}, ...]""" + if not title.strip(): + raise GoofishError("类目推荐标题不能为空") + try: + images = json.loads(images_json) if images_json else [] + except ValueError: + raise GoofishError("images-json 必须是有效 JSON 图片收据数组") from None + if not isinstance(images, list) or len(images) > 9: + raise GoofishError("images-json 必须是至多 9 张图片收据的数组") + images = [validate_image_receipt(value) for value in images] session = Session.load() - images = json.loads(images_json) if images_json else [] + unique_code = str(time.time_ns() // 1000) image_infos: list[dict[str, Any]] = [] for img in images: image_infos.append({ @@ -48,17 +79,23 @@ def recommend( "scene": "newPublishChoice", "description": title, "imageInfos": image_infos, - "uniqueCode": "1775905618164677", + "uniqueCode": unique_code, }, version="2.0", spm_cnt="a21ybx.publish.0.0", ) predict = (raw.get("data", {}) or {}).get("categoryPredictResult", {}) or {} + if not predict.get("catId") or not predict.get("channelCatId") or not predict.get("tbCatId"): + raise GoofishError("类目推荐响应缺少选中分类身份,不能继续自动发布", raw=raw) + score, source = _selected_score(raw.get("data") or {}, predict) return { "cat_id": str(predict.get("catId", "")), "cat_name": predict.get("catName", ""), "channel_cat_id": str(predict.get("channelCatId", "")), "tb_cat_id": str(predict.get("tbCatId", "")), - "confidence": predict.get("confidence", 0), + "confidence": score, + "confidence_source": source, + "confidence_kind": "model_score_not_calibrated_probability", + "unique_code": unique_code, "raw": raw, } diff --git a/src/goofish_cli/commands/item/delete.py b/src/goofish_cli/commands/item/delete.py index b6ad019..b7515cb 100644 --- a/src/goofish_cli/commands/item/delete.py +++ b/src/goofish_cli/commands/item/delete.py @@ -1,9 +1,8 @@ """item delete — 下架商品。接口 com.taobao.idle.item.delete v1.1(写操作)""" from goofish_cli.core import Session, Strategy, command -from goofish_cli.core.guard import watch -from goofish_cli.core.limiter import acquire from goofish_cli.core.mtop import call +from goofish_cli.core.write_operation import write_operation @command( @@ -16,7 +15,7 @@ ) def delete(item_id: str) -> dict[str, object]: session = Session.load() - with acquire("item.write"), watch(): + with write_operation(session, "item.write"): raw = call( session, api="com.taobao.idle.item.delete", diff --git a/src/goofish_cli/commands/item/get.py b/src/goofish_cli/commands/item/get.py index 3e91322..cadea3c 100644 --- a/src/goofish_cli/commands/item/get.py +++ b/src/goofish_cli/commands/item/get.py @@ -3,6 +3,7 @@ from typing import Any from goofish_cli.core import Session, Strategy, command +from goofish_cli.core.item_data import normalize_item from goofish_cli.core.mtop import call @@ -22,13 +23,4 @@ def get(item_id: str) -> dict[str, Any]: version="1.0", spm_cnt="a21ybx.item.0.0", ) - data = raw.get("data", {}) or {} - track = data.get("trackParams", {}) or {} - return { - "item_id": track.get("id", item_id), - "title": track.get("title", ""), - "price": track.get("soldPrice") or track.get("price", ""), - "seller_nick": track.get("seller_nick", ""), - "status": track.get("itemStatus", ""), - "raw": raw, - } + return normalize_item(raw, str(item_id)) diff --git a/src/goofish_cli/commands/item/publish.py b/src/goofish_cli/commands/item/publish.py index 245b4d7..0779506 100644 --- a/src/goofish_cli/commands/item/publish.py +++ b/src/goofish_cli/commands/item/publish.py @@ -5,15 +5,18 @@ """ import json +import time +from decimal import Decimal, InvalidOperation from typing import Any, Literal from goofish_cli.commands.category.recommend import recommend from goofish_cli.commands.location.default import default as get_default_location -from goofish_cli.commands.media.upload import upload +from goofish_cli.commands.media.upload import upload, validate_image_path, validate_image_receipt from goofish_cli.core import Session, Strategy, command +from goofish_cli.core.errors import GoofishError, PartialResultError from goofish_cli.core.guard import watch -from goofish_cli.core.limiter import acquire from goofish_cli.core.mtop import call +from goofish_cli.core.write_operation import write_operation @command( @@ -23,63 +26,116 @@ strategy=Strategy.COOKIE, columns=["item_id", "title", "price", "cat_name", "ok"], write=True, + arguments=["images", "price"], ) def publish( title: str, desc: str, - images: list[str], + images: list[str] | None = None, + *, price: float, original_price: float | None = None, delivery: Literal["包邮", "按距离计费", "一口价", "无需邮寄"] = "无需邮寄", post_price: float = 0, can_self_pickup: bool = True, + images_json: str = "[]", + category_json: str = "", + location_json: str = "", ) -> dict[str, Any]: + # 所有确定性的输入错误在任何上传或经营提交前拒绝。 + if not title.strip() or not desc.strip(): + raise GoofishError("标题和描述不能为空") + _money(price) + _money(post_price) + if original_price is not None: + _money(original_price) + try: + prepared = json.loads(images_json) + cat = json.loads(category_json) if category_json else None + loc = json.loads(location_json) if location_json else None + except ValueError: + raise GoofishError("图片、类目或地址 JSON 格式无效") from None + if not isinstance(prepared, list): + raise GoofishError("images-json 必须是上传图片收据数组") + infos = [validate_image_receipt(value) for value in prepared] + paths = [validate_image_path(path) for path in images or []] + if not 1 <= len(infos) + len(paths) <= 9: + raise GoofishError("发布必须有 1–9 张图片,可使用本地文件或上传收据") + if cat is not None: + _category(cat) + if loc is not None: + _address(loc) session = Session.load() - - # 1. 上传图片(每张独立限流) - image_infos: list[dict[str, Any]] = [] - for img_path in images: - with acquire("item.write"): - r = upload(img_path) - image_infos.append({"url": r["url"], "width": r["width"], "height": r["height"]}) - - # 2. AI 类目 - cat = recommend(title, json.dumps(image_infos)) - - # 3. 默认地址 - loc = get_default_location() - - # 4. 发布(最后一步,走熔断) - data = _build_publish_data( - title=title, - desc=desc, - image_infos=image_infos, - price=price, - original_price=original_price, - delivery=delivery, - post_price=post_price, - can_self_pickup=can_self_pickup, - cat_info=cat, - location=loc, - ) - - with acquire("item.write"), watch(): - raw = call( - session, - api="mtop.idle.pc.idleitem.publish", - data=data, - version="1.0", - spm_cnt="a21ybx.publish.0.0", - ) - - data_out = raw.get("data", {}) or {} - return { - "item_id": data_out.get("itemId", ""), - "title": title, - "price": price, - "cat_name": cat["cat_name"], - "ok": any("SUCCESS" in r for r in raw.get("ret", [])), - } + stage = "upload" + attempted = False + try: + with watch(account=session.unb): + for path in paths: + infos.append(upload(str(path))) + stage = "category" + cat = cat if cat is not None else recommend(title, json.dumps(infos)) + stage = "location" + loc = loc if loc is not None else get_default_location() + _category(cat) + _address(loc) + data = _build_publish_data(title=title, desc=desc, image_infos=infos, price=price, + original_price=original_price, delivery=delivery, + post_price=post_price, can_self_pickup=can_self_pickup, + cat_info=cat, location=loc) + stage = "submit" + with write_operation(session, "item.write"): + attempted = True + raw = call(session, api="mtop.idle.pc.idleitem.publish", data=data, + version="1.0", spm_cnt="a21ybx.publish.0.0") + result = raw.get("data") if isinstance(raw, dict) else None + item_id = result.get("itemId") if isinstance(result, dict) else None + if item_id is None or not str(item_id).strip(): + raise ValueError("missing_item_id") + except (GoofishError, OSError, ValueError, TypeError, AttributeError) as exc: + unknown = attempted and not isinstance(exc, GoofishError) + status = "submission_unknown" if unknown else "submission_rejected" if attempted else "not_submitted" + message = "发布结果未知;先核对自己的商品列表,勿自动重复发布" if unknown else str(exc) if isinstance(exc, GoofishError) else "发布准备响应结构或连接失败,未提交商品" + progress = {"images": infos, "category": {k: v for k, v in (cat or {}).items() if k != "raw"}, + "location": loc, "status": status, "failed_stage": stage, + "requires_readback": unknown} + error = {"error_type": type(exc).__name__, "phase": stage, "submission_status": status} + if hasattr(exc, "retry_after"): + error["retry_after"] = exc.retry_after + raise PartialResultError(message, progress, error, exit_code=getattr(exc, "exit_code", 1)) from exc + return {"item_id": str(item_id), "title": title, "price": price, + "cat_name": cat["cat_name"], "ok": True, + "status": "accepted", "requires_readback": True, "images": infos} + + +def _money(value: float) -> str: + try: + amount = Decimal(str(value)) + except InvalidOperation: + raise GoofishError("金额必须是有效数字") from None + try: + valid = amount.is_finite() and amount >= 0 and amount == amount.quantize(Decimal("0.01")) + except InvalidOperation: + valid = False + if not valid: + raise GoofishError("金额必须非负且最多两位小数") + return str(int(amount * 100)) + + +def _category(value: Any) -> None: + if not isinstance(value, dict) or any(not value.get(k) for k in ("cat_id", "cat_name", "channel_cat_id", "tb_cat_id")): + raise GoofishError("已确认类目必须包含 cat_id/cat_name/channel_cat_id/tb_cat_id") + + +def _address(value: Any) -> dict: + if not isinstance(value, dict): + raise GoofishError("发布地址必须是 location default 的 DTO") + selected = value.get("selected") + if selected is None: + addresses = value.get("all") or [] + selected = next((x for x in addresses if isinstance(x, dict) and str(x.get("divisionId", "")) == str(value.get("division_id", ""))), None) + if not isinstance(selected, dict) or not selected.get("divisionId"): + raise GoofishError("未确认发布地址,请提供 location-json") + return selected def _build_publish_data( @@ -102,11 +158,11 @@ def _build_publish_data( "url": img["url"], "heightSize": img["height"], "widthSize": img["width"], - "major": True, + "major": index == 0, "type": 0, "status": "done", } - for img in image_infos + for index, img in enumerate(image_infos) ] post_fee: dict[str, Any] = { @@ -122,7 +178,7 @@ def _build_publish_data( post_fee["templateId"] = "-100" elif delivery == "一口价": post_fee["supportFreight"] = True - post_fee["postPriceInCent"] = str(int(post_price * 100)) + post_fee["postPriceInCent"] = _money(post_price) post_fee["templateId"] = "0" elif delivery == "无需邮寄": post_fee["templateId"] = "0" @@ -130,23 +186,18 @@ def _build_publish_data( price_dto: dict[str, str] = {} default_price = price <= 0 if not default_price: - price_dto["priceInCent"] = str(int(price * 100)) + price_dto["priceInCent"] = _money(price) if original_price and original_price > 0: - price_dto["origPriceInCent"] = str(int(original_price * 100)) - - item_addr: dict[str, Any] = {} - if location.get("division_id"): - all_addrs = location.get("all", []) or [] - first = all_addrs[0] if all_addrs else {} - item_addr = { - "area": first.get("area", ""), - "city": first.get("city", ""), - "divisionId": first.get("divisionId", ""), - "gps": f"{first.get('longitude', '')},{first.get('latitude', '')}", - "poiId": first.get("poiId", ""), - "poiName": first.get("poi", ""), - "prov": first.get("prov", ""), - } + price_dto["origPriceInCent"] = _money(original_price) + + first = _address(location) + item_addr = { + "area": first.get("area", ""), "city": first.get("city", ""), + "divisionId": first.get("divisionId", ""), + "gps": f"{first.get('longitude', '')},{first.get('latitude', '')}", + "poiId": first.get("poiId", ""), "poiName": first.get("poi", ""), + "prov": first.get("prov", ""), + } return { "freebies": False, @@ -168,7 +219,7 @@ def _build_publish_data( "tbCatId": cat_info["tb_cat_id"], }, "onlyTakeSelf": can_self_pickup, - "uniqueCode": "1775897582791680", + "uniqueCode": cat_info.get("unique_code") or str(time.time_ns() // 1000), "sourceId": "pcMainPublish", "bizcode": "pcMainPublish", "publishScene": "pcMainPublish", diff --git a/src/goofish_cli/commands/item/view.py b/src/goofish_cli/commands/item/view.py index e5e892c..37ac6e7 100644 --- a/src/goofish_cli/commands/item/view.py +++ b/src/goofish_cli/commands/item/view.py @@ -1,6 +1,6 @@ """item view — 浏览器视角的商品详情。对标 OpenCLI `xianyu/item.js`。 -`item get` 已经在 CLI 里直签调 `mtop.taobao.idle.pc.detail` v1.0,字段浅抽(5 项)。 +`item get` 在 CLI 里直签调 `mtop.taobao.idle.pc.detail` v1.0。 `item view` 在真实 Chrome 的商品页上下文里调同一 API(走 `window.lib.mtop.request`), 好处: @@ -105,12 +105,15 @@ def _build_item_url(item_id: str) -> str: item_id: clean(item.itemId || itemId), title: clean(item.title || ''), description: clean(item.desc || ''), - price: clean('¥' + (item.soldPrice || item.defaultPrice || '')).replace(/^¥\s*$/, ''), + price: item.soldPrice == null || typeof item.soldPrice === 'boolean' || item.soldPrice === '' + ? null : '¥' + clean(item.soldPrice), + price_kind: item.defaultPrice === true ? 'negotiable' : 'fixed', original_price: clean(item.originalPrice || ''), want_count: String(item.wantCnt ?? ''), collect_count: String(item.collectCnt ?? ''), browse_count: String(item.browseCnt ?? ''), status: clean(item.itemStatusStr || ''), + status_code: item.itemStatus ?? null, condition: clean(findLabel('成色')), brand: clean(findLabel('品牌')), category: clean(findLabel('分类')), diff --git a/src/goofish_cli/commands/location/default.py b/src/goofish_cli/commands/location/default.py index 06de86b..33a127d 100644 --- a/src/goofish_cli/commands/location/default.py +++ b/src/goofish_cli/commands/location/default.py @@ -36,5 +36,6 @@ def default(longitude: float = 121.4737, latitude: float = 31.2304) -> dict[str, "area": selected.get("area", ""), "poi": selected.get("poi", ""), "division_id": str(selected.get("divisionId", "")), + "selected": selected, "all": addrs or [selected], } diff --git a/src/goofish_cli/commands/media/upload.py b/src/goofish_cli/commands/media/upload.py index 2220d0e..67ceb2c 100644 --- a/src/goofish_cli/commands/media/upload.py +++ b/src/goofish_cli/commands/media/upload.py @@ -1,53 +1,77 @@ -"""media upload — 上传图片到闲鱼 CDN。""" +"""图片上传入口拥有独立媒体预算,只有服务端确认的完整图片收据才返回。""" +from __future__ import annotations -import os +import mimetypes +from pathlib import Path +from typing import Any +from urllib.parse import urlsplit from goofish_cli.core import Session, Strategy, command +from goofish_cli.core.errors import GoofishError, RiskControlError from goofish_cli.core.session import USER_AGENT +from goofish_cli.core.write_operation import write_operation UPLOAD_URL = "https://stream-upload.goofish.com/api/upload.api" -@command( - namespace="media", - name="upload", - description="上传图片到闲鱼 CDN,返回图片 URL + 尺寸", - strategy=Strategy.COOKIE, - columns=["url", "width", "height", "size"], - write=True, -) -def upload(path: str) -> dict[str, object]: +def validate_image_path(path: str) -> Path: + source = Path(path).expanduser() + if not source.is_file() or source.stat().st_size == 0: + raise GoofishError("图片文件不存在或为空") + mime = mimetypes.guess_type(str(source))[0] + if not mime or not mime.startswith("image/"): + raise GoofishError("无法识别图片格式,请使用有正确扩展名的图片") + return source + + +def validate_image_receipt(value: Any) -> dict[str, Any]: + if not isinstance(value, dict): + raise GoofishError("已上传图片必须包含 url、width、height") + url = value.get("url") + parsed = urlsplit(url) if isinstance(url, str) else None + host = parsed.hostname if parsed else "" + if not parsed or parsed.scheme != "https" or parsed.username or not host or not ( + host == "alicdn.com" or host.endswith(".alicdn.com") + or host == "goofish.com" or host.endswith(".goofish.com") + ): + raise GoofishError("图片 URL 必须是闲鱼上传返回的 HTTPS CDN 地址") + dimensions = [] + for key in ("width", "height"): + v = value.get(key) + if isinstance(v, bool) or not isinstance(v, (int, str)) or not str(v).isdigit() or int(v) <= 0: + raise GoofishError("图片宽高必须是正整数") + dimensions.append(int(v)) + return {"url": url, "width": dimensions[0], "height": dimensions[1]} + + +@command(namespace="media", name="upload", description="上传图片,独立媒体预算,返回可复用图片收据", + strategy=Strategy.COOKIE, columns=["url", "width", "height", "size"], write=True) +def upload(path: str) -> dict[str, Any]: + source = validate_image_path(path) session = Session.load() - abs_path = os.path.expanduser(path) - if not os.path.exists(abs_path): - raise FileNotFoundError(f"图片不存在:{abs_path}") - - headers = { - "accept": "*/*", - "origin": "https://www.goofish.com", - "referer": "https://www.goofish.com/", - "user-agent": USER_AGENT, - } + headers = {"accept": "*/*", "origin": "https://www.goofish.com", "referer": "https://www.goofish.com/", "user-agent": USER_AGENT} params = {"floderId": "0", "appkey": "xy_chat", "_input_charset": "utf-8"} - - with open(abs_path, "rb") as f: - resp = session.http.post( - UPLOAD_URL, - headers=headers, - params=params, - files={"file": (os.path.basename(abs_path), f, "image/png")}, - timeout=60, - ) - raw = resp.json() - obj = raw.get("object") or {} - pix = str(obj.get("pix", "0x0")) try: - width, height = map(int, pix.split("x")) - except ValueError: - width = height = 0 - return { - "url": obj.get("url", ""), - "width": width, - "height": height, - "size": obj.get("size", 0), - } + with source.open("rb") as stream, write_operation(session, "media.write"): + resp = session.http.post(UPLOAD_URL, headers=headers, params=params, + files={"file": (source.name, stream, mimetypes.guess_type(str(source))[0])}, timeout=60) + resp.raise_for_status() + raw = resp.json() + if not isinstance(raw, dict): + raise GoofishError("上传响应结构无效;未确认成功") + if raw.get("success") is not True: + marker = str(raw.get("status", "")) + if any(k in marker for k in ("RGV587", "USER_VALIDATE", "ILLEGAL_ACCESS")): + raise RiskControlError("图片上传被平台风控拒绝") + raise GoofishError("图片上传未被服务端确认成功") + obj = raw.get("object") or {} + if not isinstance(obj, dict): + raise GoofishError("上传响应缺少图片对象") + try: + width, height = map(int, str(obj.get("pix", "")).split("x")) + except (ValueError, TypeError): + raise GoofishError("上传响应缺少有效图片尺寸") from None + receipt = validate_image_receipt({"url": obj.get("url"), "width": width, "height": height}) + except (OSError, ValueError): + raise GoofishError("图片上传连接或响应失败;未确认成功,不自动重传") from None + return {**receipt, "size": obj.get("size")} diff --git a/src/goofish_cli/commands/message/history.py b/src/goofish_cli/commands/message/history.py index 65aa0af..3a3b642 100644 --- a/src/goofish_cli/commands/message/history.py +++ b/src/goofish_cli/commands/message/history.py @@ -16,8 +16,18 @@ name="history", description="拉取指定 cid 会话的历史消息(翻页到底)", strategy=Strategy.COOKIE, - columns=["send_user_id", "send_user_name", "message"], + columns=["created_at", "message_id", "send_user_id", "send_user_name", "message"], ) -def history(cid: str, limit_per_page: int = 20) -> list[dict[str, Any]]: +def history( + cid: str, limit_per_page: int = 20, limit: int = 0, timeout: float = 30.0 +) -> list[dict[str, Any]]: session = Session.load() - return asyncio.run(list_user_messages(session, cid, limit_per_page=limit_per_page)) + return asyncio.run( + list_user_messages( + session, + cid, + limit_per_page=limit_per_page, + limit=limit, + timeout=timeout, + ) + ) diff --git a/src/goofish_cli/commands/message/list_chats.py b/src/goofish_cli/commands/message/list_chats.py index fa8aa66..3eb5e57 100644 --- a/src/goofish_cli/commands/message/list_chats.py +++ b/src/goofish_cli/commands/message/list_chats.py @@ -1,110 +1,225 @@ -"""message list-chats — 拉取会话列表(左栏)。 +"""合并 HTTP 与有界 WS 会话发现,读取最近消息后按真实活动时间排序。""" -h5 接口 `mtop.taobao.idlemessage.pc.session.sync` v3.0 是阉割版,只返回活跃 -Top N 会话;网页左栏看到的完整列表其实是靠 ACCS 长连累积的,单次 HTTP 拿不到。 -所以提供 `--watch-secs N` 可选开关:短时连 WS + `ackDiff(pts=0)` 拉历史推送, -从中抽取会话激活事件 + new_msg 通知里的 cid,补齐 baseline 漏掉的会话。 - -数据来源两路合并: - -1. `mtop.taobao.idlemessage.pc.session.sync` v3.0 —— baseline,字段齐全 - (peer_nick / peer_user_id / unread / last_msg / ts / session_type / item_id)。 -2. `--watch-secs N`(可选)—— watch,只有 cid 骨架 - (session_id / session_type / item_id / ts),`peer_nick` / `peer_user_id` / - `last_msg` / `unread` 都填空值。要正文请自己调 `message history `。 - -输出 record 带 `source` 字段区分 `baseline` 和 `watch`,shape 一致,方便调用方统一处理。 -""" +from __future__ import annotations import asyncio +import math from typing import Any +from websockets.exceptions import ConnectionClosed, InvalidHandshake + from goofish_cli.core import Session, Strategy, command +from goofish_cli.core.errors import GoofishError +from goofish_cli.core.message_history import normalize_id, recent_messages, timestamp from goofish_cli.core.mtop import call +from goofish_cli.core.ws import collect_session_cids -def _pick(d: dict[str, Any], *path: str, default: Any = "") -> Any: - cur: Any = d - for key in path: - if not isinstance(cur, dict): - return default - cur = cur.get(key) - if cur is None: - return default - return cur - - -def _parse_session(item: dict[str, Any]) -> dict[str, Any]: +def _parse_session(item: dict[str, Any], myid: str = "") -> dict[str, Any]: session = item.get("session") or {} user_info = session.get("userInfo") or {} - summary = _pick(item, "message", "summary", default={}) or {} + if myid and session.get("sessionType") == 1: + people = [user_info, session.get("ownerInfo") or {}] + peers = { + normalize_id(p.get("userId")): p + for p in people + if p.get("userId") and normalize_id(p["userId"]) != myid + } + user_info = next(iter(peers.values())) if len(peers) == 1 else {} + summary = (item.get("message") or {}).get("summary") or {} return { - "session_id": str(session.get("sessionId", "")), + "session_id": normalize_id(session.get("sessionId")), "peer_nick": user_info.get("nick", "") or user_info.get("fishNick", ""), - "peer_user_id": str(user_info.get("userId", "")), + "peer_user_id": normalize_id(user_info.get("userId")), "unread": summary.get("unread", 0), "last_msg": summary.get("summary", ""), - "ts": summary.get("ts", 0), + "ts": timestamp(summary.get("ts")), "session_type": session.get("sessionType", 0), "item_id": "", + "message_id": "", "source": "baseline", + "metadata_status": "http", } -def _watch_record(w: dict[str, Any]) -> dict[str, Any]: - """把 WS 收集到的裸 cid 包成跟 baseline 同形状的 record。""" - ts_raw = w.get("last_msg_ts") or 0 - try: - ts = int(ts_raw) - except (TypeError, ValueError): - ts = 0 +def _watch_record(item: dict[str, Any]) -> dict[str, Any]: return { - "session_id": str(w["cid"]), + "session_id": normalize_id(item["cid"]), "peer_nick": "", - "peer_user_id": str(w.get("peer_user_id", "")), - "unread": 0, + "peer_user_id": normalize_id(item.get("peer_user_id")), + "unread": None, "last_msg": "", - "ts": ts, - "session_type": int(w.get("session_type") or 0), - "item_id": str(w.get("item_id", "")), + "ts": timestamp(item.get("last_msg_ts")), + "session_type": int(item.get("session_type") or 0), + "item_id": str(item.get("item_id") or ""), + "message_id": str(item.get("last_msg_id") or ""), "source": "watch", + "metadata_status": "pending", } +def _preview(message: Any) -> str: + if not isinstance(message, dict): + return str(message or "") + text = message.get("text") or {} + if isinstance(text, dict) and text.get("text"): + return str(text["text"]) + if message.get("contentType") == 2: + return "[图片]" + return str(message.get("summary") or "[非文本消息]") + + +def _enrich( + row: dict[str, Any], messages: list[dict[str, Any]], participants: list[str], myid: str +) -> None: + if not messages: + row.update(metadata_status="empty", ts=None, last_msg="", message_id="") + if row["peer_user_id"] == myid: + row.update(peer_user_id="", peer_nick="") + missing = [key for key in ("peer_user_id", "peer_nick") if not row[key]] + if missing: + row.update(metadata_status="partial", metadata_missing=missing) + return + latest = max(messages, key=lambda m: m["created_at"] or 0) + row.update( + ts=latest["created_at"], + message_id=latest["message_id"], + last_msg=_preview(latest["message"]), + metadata_status="current", + ) + if row["session_type"] == 0 and latest.get("session_type"): + row["session_type"] = latest["session_type"] + if row["session_type"] not in (0, 1): + row["metadata_status"] = "system" + return + members = {normalize_id(value) for value in participants if value} + peers = members - {myid} if myid in members else set() + if row["peer_user_id"] and row["peer_user_id"] != myid: + peers.add(row["peer_user_id"]) + if not peers: + peers = { + m["send_user_id"] for m in messages if m["send_user_id"] and m["send_user_id"] != myid + } + peers.update(uid for m in messages for uid in m["receiver_user_ids"] if uid and uid != myid) + if len(peers) == 1: + peer = next(iter(peers)) + row["peer_user_id"] = peer + labels = [ + m + for m in messages + if m["send_user_id"] == peer + and m["send_user_name"] + and isinstance(m["message"], dict) + and m["message"].get("contentType") in (1, 2) + ] + if labels: + row["peer_nick"] = max(labels, key=lambda m: m["created_at"] or 0)["send_user_name"] + else: + row["peer_user_id"] = "" + row["peer_nick"] = "" + missing = [key for key in ("ts", "message_id", "peer_user_id", "peer_nick") if not row[key]] + if latest.get("parse_error"): + missing.append("message_content") + if missing: + row.update(metadata_status="partial", metadata_missing=missing) + + +async def _snapshot( + session: Session, baseline: list[dict[str, Any]], watch_secs: float, timeout: float +) -> tuple[list[dict[str, Any]], int, list[dict[str, str]]]: + records = {row["session_id"]: row for row in baseline if row["session_id"]} + participants: dict[str, list[str]] = {} + errors: list[dict[str, str]] = [] + discovered = 0 + if watch_secs: + try: + async with asyncio.timeout(watch_secs + 15): + pushed = await collect_session_cids(session, duration=watch_secs) + for item in pushed: + cid = normalize_id(item["cid"]) + if not cid: + continue + participants[cid] = item.get("participant_user_ids", []) + if cid not in records: + records[cid] = _watch_record(item) + discovered += 1 + else: + records[cid]["source"] = "baseline+watch" + if item.get("item_id"): + records[cid]["item_id"] = str(item["item_id"]) + except (TimeoutError, ConnectionClosed, OSError, InvalidHandshake) as exc: + errors.append({"stage": "discovery", "error": type(exc).__name__}) + + cids = [cid for cid, row in records.items() if row["session_type"] in (0, 1)] + for row in records.values(): + if row["source"] == "watch" and row["session_type"] not in (0, 1): + row["metadata_status"] = "system" + cids.sort(key=lambda cid: "watch" not in records[cid]["source"]) + histories, failures = await recent_messages(session, cids, timeout=timeout) + for cid in cids: + if cid in failures: + records[cid]["metadata_status"] = "stale" + errors.append({"stage": "recent_message", "session_id": cid, "error": failures[cid]}) + else: + _enrich(records[cid], histories[cid], participants.get(cid, []), session.unb) + rows = sorted(records.values(), key=lambda r: (r["ts"] or 0, r["session_id"]), reverse=True) + return rows, discovered, errors + + @command( namespace="message", name="list-chats", - description="拉取会话列表(左栏):session.sync 基线 + 可选 WS 增量补 cid", + description="会话列表:HTTP + 默认短时 WS 发现,最近消息补齐并按活动时间倒序", strategy=Strategy.COOKIE, columns=[ - "session_id", "peer_nick", "peer_user_id", - "unread", "last_msg", "ts", "source", + "session_id", + "peer_nick", + "peer_user_id", + "unread", + "last_msg", + "ts", + "source", + "metadata_status", ], ) -def list_chats(fetch_num: int = 50, watch_secs: float = 0.0) -> dict[str, Any]: +def list_chats( + fetch_num: int = 50, watch_secs: float = 5.0, timeout: float = 30.0 +) -> dict[str, Any]: + if ( + fetch_num < 1 + or watch_secs < 0 + or timeout <= 0 + or not math.isfinite(watch_secs) + or not math.isfinite(timeout) + ): + raise GoofishError("fetch-num 和 timeout 必须大于 0,watch-secs 不得为负") session = Session.load() - raw = call( - session, - api="mtop.taobao.idlemessage.pc.session.sync", - data={"fetchNum": int(fetch_num)}, - version="3.0", - spm_cnt="a21ybx.im.0.0", - ) + try: + raw = call( + session, + api="mtop.taobao.idlemessage.pc.session.sync", + data={"fetchNum": fetch_num}, + version="3.0", + spm_cnt="a21ybx.im.0.0", + ) + except OSError as exc: + raise GoofishError(f"HTTP 会话基线连接失败:{type(exc).__name__}") from None data = raw.get("data") or {} - baseline = [_parse_session(s) for s in data.get("sessions") or []] - known = {b["session_id"] for b in baseline} - - extras: list[dict[str, Any]] = [] - if watch_secs > 0: - from goofish_cli.core.ws import collect_session_cids - - pushed = asyncio.run(collect_session_cids(session, duration=float(watch_secs))) - extras = [_watch_record(w) for w in pushed if str(w["cid"]) not in known] - + baseline = [_parse_session(item, session.unb) for item in data.get("sessions") or []] + rows, discovered, errors = asyncio.run(_snapshot(session, baseline, watch_secs, timeout)) return { - "sessions": baseline + extras, + "sessions": rows, "has_more": bool(data.get("hasMore")), - "total": len(baseline) + len(extras), + "has_more_scope": "http", + "total": len(rows), "from_baseline": len(baseline), - "from_watch": len(extras), + "from_watch": discovered, + "ws_enabled": watch_secs > 0, + "enumeration_complete": False, + "enumeration_note": "HTTP 与短时 WS 是有界会话发现,不能保证覆盖所有历史会话", + "metadata_complete": not errors + and all(r["metadata_status"] not in ("pending", "partial", "stale") for r in rows), + "metadata_scope": "personal_and_unclassified", + "metadata_errors": errors, + "unknown_activity_count": sum(r["ts"] is None for r in rows), } diff --git a/src/goofish_cli/commands/message/send.py b/src/goofish_cli/commands/message/send.py index d75b67f..b7a6e8b 100644 --- a/src/goofish_cli/commands/message/send.py +++ b/src/goofish_cli/commands/message/send.py @@ -14,9 +14,8 @@ from goofish_cli.core import Session, Strategy, command from goofish_cli.core.errors import GoofishError -from goofish_cli.core.guard import watch -from goofish_cli.core.limiter import acquire from goofish_cli.core.token import get_access_token +from goofish_cli.core.write_operation import write_operation from goofish_cli.core.ws import ( connect, create_chat, @@ -49,7 +48,7 @@ def send( item_id: str = "", ) -> dict[str, Any]: session = Session.load() - with acquire("message.write"), watch(): + with write_operation(session, "message.write"): return asyncio.run(_send( session, cid=cid, diff --git a/src/goofish_cli/commands/search/search.py b/src/goofish_cli/commands/search/search.py index 0ba45b8..cd853fa 100644 --- a/src/goofish_cli/commands/search/search.py +++ b/src/goofish_cli/commands/search/search.py @@ -1,15 +1,14 @@ """search — 搜索闲鱼商品。对标 OpenCLI `xianyu/search.js`。 思路:打开 `https://www.goofish.com/search?q=xxx` 让页面自己渲染,autoScroll 触发 -懒加载,再在 page context 里跑 DOM 选择器提卡片。**不走 mtop 直签**: -- search 没对外 API,只有 HTML 卡片 + 动态加载 -- 浏览器真实渲染天然抗风控 +懒加载,再在 page context 里跑 DOM 选择器提卡片。观察页面实际发出的搜索响应及其控制字段,DOM 用于分页与异常观察。 +搜索不另行直签,也不把页面推荐当作查询命中。 分页(2026-08 实测):搜索页**不是无限滚动**——窄查询滚动到底卡片数不再增长; 翻页靠 DOM 里的 `search-pagination-container`(宽查询下 1..50 页),点击右箭头后 SPA 内部重渲染、URL 不变,`?page=N` URL 参数被服务端忽略。所以跨页抓取 = 点击 翻页箭头 + 按 item_id 去重累积,终止条件:达到 --pages 上限 / 右箭头 disabled / -连续一页无新增(防御)。 +连续一页无新增(防御);跨页失败同时返回部分数据及非成功状态。 字段参考 OpenCLI:`item_id / rank / title / price / original_price / condition / brand / location / badge / url / extra`。 @@ -17,12 +16,15 @@ from __future__ import annotations import asyncio +import json import re from typing import Any +from urllib.parse import parse_qs, urlsplit from goofish_cli.core import Strategy, command from goofish_cli.core.browser import auto_scroll, goofish_page -from goofish_cli.core.errors import AuthRequiredError, GoofishError +from goofish_cli.core.errors import AuthRequiredError, GoofishError, PartialResultError +from goofish_cli.core.search_data import normalize_search # limit 是**跨页总上限**(去重后条数)。站点每页 30 卡,50 页满配远超此值, # 200 是给 MCP 调用方的运行时护栏(每页 ~4s,200 条 ≈ 7 页 ≈ 40s)。 @@ -92,7 +94,7 @@ def _build_search_url(query: str) -> str: priceWrap: '[class*="price-wrap"]', priceNum: '[class*="number"]', priceDec: '[class*="decimal"]', - priceDesc: '[class*="price-desc"] [title], [class*="price-desc"] [style*="line-through"]', + priceDesc: '[class*="price-desc"] [style*="line-through"]', sellerWrap: '[class*="row4-wrap-seller"]', sellerText: '[class*="seller-text"]', badge: '[class*="credit-container"] [title], [class*="credit-container"] span', @@ -131,11 +133,14 @@ def _build_search_url(query: str) -> str: title, url: href, price: clean('¥' + priceNumber + priceDecimal).replace(/^¥\s*$/, ''), - original_price: clean(originalPriceNode?.getAttribute('title') || originalPriceNode?.textContent || ''), - condition: attrs[0] || '', - brand: attrs[1] || '', + original_price: /^([¥¥]\s*)?\d+(\.\d{1,2})?$/.test(clean(originalPriceNode?.textContent)) ? clean(originalPriceNode?.textContent) : null, + condition: null, + brand: null, + attributes: attrs, extra: attrs.slice(2).join(' | '), - location, + location: null, + seller_text: location, + source: 'dom_unclassified', badge: clean(badgeNode?.getAttribute('title') || badgeNode?.textContent || ''), }; }) @@ -178,6 +183,91 @@ def _build_search_url(query: str) -> str: """ +class _SearchPage: + """响应同时绑定真实 keyword/pageNumber 与请求发起代次,拒绝旧页迟到结果。""" + + def __init__(self, page: Any): + self.page = page + self.payload: dict | None = None + self.generation = 0 + self.query = "" + self.number = 1 + self.requests: dict[Any, int] = {} + self.tasks: set[asyncio.Task] = set() + self.ready = asyncio.Event() + page.on("request", self._request) + page.on("response", self._schedule) + + def begin(self, query: str, number: int) -> None: + self.generation += 1 + self.query, self.number = query, number + self.payload = None + self.requests.clear() + self.ready.clear() + + def _request(self, request: Any) -> None: + if urlsplit(request.url).path != "/h5/mtop.taobao.idlemtopsearch.pc.search/1.0/": + return + try: + fields = parse_qs(request.post_data or "") + if "data" not in fields: + fields = parse_qs(urlsplit(request.url).query) + data = json.loads(fields.get("data", ["{}"])[0]) + if data.get("keyword") == self.query and int(data.get("pageNumber", 0)) == self.number: + self.requests[request] = self.generation + except (ValueError, TypeError, AttributeError): + return + + def _schedule(self, response: Any) -> None: + generation = self.requests.pop(response.request, None) + if generation is not None: + task = asyncio.create_task(self._capture(response, generation)) + self.tasks.add(task) + task.add_done_callback(self.tasks.discard) + + async def _capture(self, response: Any, generation: int) -> None: + try: + raw = await response.json() + if not isinstance(raw, dict): + raise ValueError("invalid_search_response") + values = raw.get("ret") or [] + ret = " | ".join(values) if isinstance(values, list) else str(values) + if ret and "SUCCESS" not in ret: + auth = any(t in ret for t in ("SESSION_EXPIRED", "TOKEN_EXOIRED", "TOKEN_EMPTY")) + blocked = any(t in ret for t in ("RGV587", "USER_VALIDATE", "ILLEGAL_ACCESS")) + payload = {"items": [], "requiresAuth": auth, "blocked": blocked, + "response_error": "authentication_required" if auth else "blocked" if blocked else "search_response_rejected"} + elif not isinstance(raw.get("data"), dict) or "resultList" not in raw["data"]: + payload = {"items": [], "schema_unknown": True} + else: + payload = normalize_search(raw) + except Exception: + payload = {"items": [], "schema_unknown": True} + if generation == self.generation: + self.payload = payload + self.ready.set() + + async def extract(self) -> dict[str, Any]: + try: + await asyncio.wait_for(self.ready.wait(), timeout=8) + except TimeoutError: + return {"items": [], "schema_unknown": True, "response_error": "search_response_timeout"} + dom = await self.page.evaluate(_EXTRACT_JS, MAX_LIMIT) + if not isinstance(dom, dict) or self.payload is None: + return {"items": [], "schema_unknown": True} + return {**dom, **self.payload, + "recommendations_present": bool(self.payload.get("empty") and dom.get("items"))} + + def __getattr__(self, name: str): + return getattr(self.page, name) + + +async def _extract_page(page: Any) -> dict[str, Any]: + if isinstance(page, _SearchPage): + return await page.extract() + return await page.evaluate(_EXTRACT_JS, MAX_LIMIT) + + def _first_card_id(items: list[dict[str, Any]]) -> str: """当前页首卡 id——翻页等待的"内容已变化"基准。""" for it in items: @@ -188,8 +278,12 @@ def _first_card_id(items: list[dict[str, Any]]) -> str: def _raise_for_failed_page(payload: dict[str, Any]) -> None: - """首页级失败(0 卡片)按原语义抛错;翻页途中的失败由调用方优雅终止。""" + """首页级失败(0 卡片)按原语义抛错;翻页途中的失败由调用方保留部分数据并报告失败。""" items = payload.get("items") or [] + if payload.get("schema_unknown"): + raise GoofishError("未观察到可确认的搜索响应结构,不能把推荐卡片当搜索命中") + if payload.get("response_error") and not payload.get("requiresAuth") and not payload.get("blocked"): + raise GoofishError("搜索服务拒绝请求") # "登录后" 在页脚也会出现——只有在"没拿到卡片 && 命中关键词"时才判定 auth 失败 if not items and payload.get("requiresAuth"): raise AuthRequiredError("www.goofish.com 搜索结果页要求登录,cookies 可能失效") @@ -210,11 +304,12 @@ async def _walk_pages( fetched_pages: int, pages: int, limit: int, + failures: list[dict] | None = None, ) -> tuple[int, int | None, str]: """翻页状态机:点右箭头 → 等重渲染 → 去重累积。 返回 (fetched_pages, total_pages, stopped_reason)。任何翻页途中的 - Playwright 异常(SPA 重渲染销毁执行上下文等)都被吞掉并优雅终止—— + Playwright 异常(SPA 重渲染销毁执行上下文等)会保留已累积数据并记录失败—— 调用方拿到已累积的部分结果,stopped_reason 说明终止原因。 锚点(cur_first_id)始终取**未过滤**的下一页 payload 首卡: @@ -222,6 +317,11 @@ async def _walk_pages( 用 fresh 的首卡做锚点会与 DOM 实况脱节,导致 page-change 等待 假阳性、提前终止(review P1-2)。 """ + def failure(kind: str, phase: str, exception_type: str = "") -> None: + if failures is not None: + failures.append({"error_type": kind, "phase": phase, + "target_page": fetched_pages + 1, "exception_type": exception_type}) + anchor_id = _first_card_id(items) or "" total_pages: int | None = None stopped_reason = "pages_reached" @@ -229,10 +329,12 @@ async def _walk_pages( while fetched_pages < pages and len(items) < limit: try: pag = await page.evaluate(_PAGINATION_JS) - except Exception: # noqa: BLE001 — 执行上下文销毁等,保留已抓结果 + except Exception as exc: # noqa: BLE001 + failure("browser_error", "pagination", type(exc).__name__) return fetched_pages, total_pages, "error" # 结构突变(None / 非dict)同样按优雅终止处理,不能 AttributeError 穿透 if not isinstance(pag, dict): + failure("schema_changed", "pagination") return fetched_pages, total_pages, "error" if total_pages is None and pag.get("totalPages") is not None: @@ -242,16 +344,23 @@ async def _walk_pages( try: arrow = page.locator('[class*="pagination-arrow-container"]').nth(1) + if isinstance(page, _SearchPage): + page.begin(page.query, fetched_pages + 1) await arrow.click(timeout=5000) # 等待"首卡 id 不再等于上一页 DOM 首卡"。changed=False 是等待超时 # (重渲染大概率失败);先看提取结果再定性。 changed = await page.evaluate(_WAIT_PAGE_CHANGE_JS, anchor_id) await page.wait_for_timeout(PAGE_STABLE_MS) - nxt = await page.evaluate(_EXTRACT_JS, MAX_LIMIT) - except Exception: # noqa: BLE001 — 同上,部分结果优先 + nxt = await _extract_page(page) + except Exception as exc: # noqa: BLE001 + failure("browser_error", "page_transition", type(exc).__name__) return fetched_pages, total_pages, "error" if not isinstance(nxt, dict): + failure("schema_changed", "extract") return fetched_pages, total_pages, "error" + if nxt.get("requiresAuth") or nxt.get("blocked") or nxt.get("schema_unknown") or nxt.get("response_error"): + failure("authentication_required" if nxt.get("requiresAuth") else "blocked" if nxt.get("blocked") else nxt.get("response_error") or "schema_changed", "extract") + return fetched_pages, total_pages, "blocked" if nxt.get("blocked") else "auth_required" if nxt.get("requiresAuth") else "error" # 锚点更新为**未过滤** payload 的首卡(= 本页 DOM 实际首卡)。 # fresh 是去重后的列表:若本页首卡与上一页重复,会被过滤掉, @@ -267,6 +376,8 @@ async def _walk_pages( if not fresh: if not nxt_items and nxt.get("blocked"): return fetched_pages, total_pages, "blocked" + if not changed: + failure("page_stale", "page_transition") return fetched_pages, total_pages, "no_new" if changed else "stale" for it in fresh: @@ -288,15 +399,21 @@ async def _run(query: str, limit: int, pages: int) -> dict[str, Any]: seen: set[str] = set() fetched_pages = 0 total_pages: int | None = None + failures: list[dict] = [] + result_kind = "unknown" + reported_match_count = None + recommendations_present = False - async with goofish_page() as page: + async with goofish_page() as raw_page: + page = _SearchPage(raw_page) # ---- 第 1 页(保留瞬时登录墙重试:整页重新导航)---- payload: dict[str, Any] | None = None for attempt in range(1, AUTH_WALL_ATTEMPTS + 1): + page.begin(query, 1) await page.goto(url, wait_until="domcontentloaded") await page.wait_for_timeout(2000) await auto_scroll(page, times=2) - raw = await page.evaluate(_EXTRACT_JS, MAX_LIMIT) + raw = await _extract_page(page) if not isinstance(raw, dict): raise GoofishError("搜索页返回结构非预期") payload = raw @@ -308,6 +425,9 @@ async def _run(query: str, limit: int, pages: int) -> dict[str, Any]: assert payload is not None _raise_for_failed_page(payload) + result_kind = payload.get("result_kind", "unknown") + reported_match_count = payload.get("reported_match_count") + recommendations_present = payload.get("recommendations_present", False) # item_id 是输出的稳定主键:解析不出数字 id 的卡片直接跳过, # 否则同一张坏卡跨页会重复追加(seen 只记非空 id) @@ -323,7 +443,7 @@ async def _run(query: str, limit: int, pages: int) -> dict[str, Any]: # ---- 翻页:点右箭头 → 等重渲染 → 去重累积(状态机,见 _walk_pages)---- fetched_pages, total_pages, stopped_reason = await _walk_pages( - page, items, seen, fetched_pages, pages, limit + page, items, seen, fetched_pages, pages, limit, failures ) # 循环提前退出(到 limit / 末页)时补一次终态读取 @@ -335,7 +455,7 @@ async def _run(query: str, limit: int, pages: int) -> dict[str, Any]: total_pages = None items = items[:limit] - return { + result = { "items": [ {"rank": i + 1, "item_id": _item_id_from_url(it.get("url", "")), **it} for i, it in enumerate(items) @@ -344,8 +464,18 @@ async def _run(query: str, limit: int, pages: int) -> dict[str, Any]: "pages_fetched": fetched_pages, "pages_total": total_pages, "stopped_reason": stopped_reason, - "query": "", + "query": query, + "result_kind": result_kind, + "reported_match_count": reported_match_count, + "recommendations_present": recommendations_present, + "complete": not failures, + "errors": failures, } + if failures: + error = failures[0] + code = 77 if error["error_type"] == "authentication_required" else 76 if error["error_type"] == "blocked" else 1 + raise PartialResultError("搜索未完成,已保留部分结果", result, error, exit_code=code) + return result @command( diff --git a/src/goofish_cli/core/errors.py b/src/goofish_cli/core/errors.py index 5918e52..bdfbe50 100644 --- a/src/goofish_cli/core/errors.py +++ b/src/goofish_cli/core/errors.py @@ -49,3 +49,13 @@ class EmptyResultError(GoofishError): class BlockedError(GoofishError): """请求被拦截(验证码页 / 安全验证 / 异常访问)。对标 opencli 的 blocked 分支。""" exit_code = 79 + + +class PartialResultError(GoofishError): + """失败的操作仍有可消费的部分数据,两个入口同时保留数据和失败状态。""" + + def __init__(self, message: str, data: dict, error: dict, *, exit_code: int = 1): + super().__init__(message) + self.data = data + self.error = error + self.exit_code = exit_code diff --git a/src/goofish_cli/core/guard.py b/src/goofish_cli/core/guard.py index deb8359..527914a 100644 --- a/src/goofish_cli/core/guard.py +++ b/src/goofish_cli/core/guard.py @@ -1,18 +1,26 @@ -"""风控熔断。检测到 RiskControlError 后写入熔断时间戳,后续请求直接拒绝。""" +"""账号熔断状态原子持久化;旧共享熔断自然过期,手动重置只改变本机状态。""" from __future__ import annotations -import json +import hashlib +import math import os import time from contextlib import contextmanager from pathlib import Path -from goofish_cli.core.errors import RiskControlError +from filelock import FileLock + +from goofish_cli.core.errors import GoofishError, RiskControlError +from goofish_cli.core.local_state import atomic_write, locked_state STATE_PATH = Path.home() / ".goofish-cli" / "circuit.json" DEFAULT_BREAK_MINUTES = 10 +def _path() -> Path: + return Path(os.environ.get("GOOFISH_GUARD_PATH", str(STATE_PATH))).expanduser() + + def _break_seconds() -> int: try: return max(60, int(os.environ.get("GOOFISH_CIRCUIT_BREAK_MINUTES", DEFAULT_BREAK_MINUTES)) * 60) @@ -20,44 +28,53 @@ def _break_seconds() -> int: return DEFAULT_BREAK_MINUTES * 60 -def _load() -> float: - if not STATE_PATH.exists(): - return 0.0 - try: - return float(json.loads(STATE_PATH.read_text()).get("until", 0)) - except (json.JSONDecodeError, OSError, ValueError): - return 0.0 +def _key(account: str) -> str: + return hashlib.sha256(account.encode()).hexdigest() -def _save(until: float) -> None: - STATE_PATH.parent.mkdir(parents=True, exist_ok=True) - STATE_PATH.write_text(json.dumps({"until": until})) +def _validate(state: dict) -> None: + values = state.get("accounts", {}) + if not isinstance(values, dict) or any(not isinstance(key, str) for key in values): + raise GoofishError("熔断状态格式无效,未清空状态") + for value in [state.get("until", 0), *values.values()]: + if isinstance(value, bool) or not isinstance(value, (float, int)) or not math.isfinite(value): + raise GoofishError("熔断时间格式无效,未发起写请求") -def check() -> None: - until = _load() - if until and time.time() < until: - remain = int(until - time.time()) - raise RiskControlError( - f"风控熔断中,剩余 {remain}s。触发后自动冷却,可通过 `goofish auth reset-guard` 手动解除。" - ) +def check(*, account: str = "") -> None: + with locked_state(_path()) as state: + _validate(state) + until = max(state.get("until", 0), state.get("accounts", {}).get(_key(account), 0)) + if until and time.time() < until: + remain = max(1, int(until - time.time())) + raise RiskControlError(f"风控熔断中,剩余 {remain}s。等待冷却或由操作者检查后使用 auth reset-guard。") -def trip(reason: str = "") -> None: - _save(time.time() + _break_seconds()) +def trip(reason: str = "", *, account: str = "") -> None: + with locked_state(_path()) as state: + _validate(state) + until = time.time() + _break_seconds() + if account: + state.setdefault("accounts", {})[_key(account)] = until + else: + state["until"] = until def reset() -> None: - if STATE_PATH.exists(): - STATE_PATH.unlink() + path = _path() + path.parent.mkdir(parents=True, exist_ok=True) + with FileLock(str(path) + ".lock", timeout=5, mode=0o600): + atomic_write(path, {}) @contextmanager -def watch(): - """包住写操作:命中 RGV587 自动熔断。""" - check() +def watch(*, account: str = ""): + check(account=account) try: yield - except RiskControlError: - trip() + except RiskControlError as exc: + try: + trip(account=account) + except GoofishError: + exc.args = (str(exc) + ";本机熔断未能持久化,请停止写入并检查护栏目录",) raise diff --git a/src/goofish_cli/core/item_data.py b/src/goofish_cli/core/item_data.py new file mode 100644 index 0000000..d2e8cb1 --- /dev/null +++ b/src/goofish_cli/core/item_data.py @@ -0,0 +1,40 @@ +"""商品业务 DTO 适配;埋点字段和默认价格标志不作为商品数据。""" +from __future__ import annotations + +from typing import Any + +from goofish_cli.core.errors import GoofishError + + +def normalize_item(raw: dict[str, Any], requested_id: str) -> dict[str, Any]: + data = raw.get("data") or {} + item = data.get("itemDO") or {} + seller = data.get("sellerDO") or {} + if not isinstance(item, dict) or not item.get("itemId"): + raise GoofishError("商品响应缺少 itemDO/itemId,不能确认详情") + item_id = str(item["itemId"]) + if item_id != str(requested_id): + raise GoofishError("商品响应身份与请求不一致") + labels = item.get("itemLabelExtList") or [] + properties = {str(label.get("propertyText")): label.get("text") + for label in labels if isinstance(label, dict) and label.get("propertyText")} + amount = item.get("soldPrice") + if isinstance(amount, bool): + amount = None + return { + "item_id": item_id, + "title": item.get("title"), + "description": item.get("desc"), + "price": str(amount) if amount is not None and amount != "" else None, + "price_kind": "negotiable" if item.get("defaultPrice") is True else "fixed", + "original_price": item.get("originalPrice"), + "status": item.get("itemStatusStr"), + "status_code": item.get("itemStatus"), + "condition": properties.get("成色"), + "brand": properties.get("品牌"), + "category": properties.get("分类"), + "seller_nick": seller.get("nick") or seller.get("uniqueName"), + "seller_id": str(seller["sellerId"]) if seller.get("sellerId") is not None else None, + "location": seller.get("publishCity") or seller.get("city"), + "raw": raw, + } diff --git a/src/goofish_cli/core/limiter.py b/src/goofish_cli/core/limiter.py index 4771bfb..e07c1de 100644 --- a/src/goofish_cli/core/limiter.py +++ b/src/goofish_cli/core/limiter.py @@ -1,60 +1,56 @@ -"""令牌桶限流。单账号 + 单命名空间,默认 1 写/分钟(可配)。 - -写入 ~/.goofish-cli/limiter.json 做进程间共享(单机多进程场景)。 -""" +"""账号与业务桶限流,锁住读取/检查/写入事务,写尝试失败不退还预算。""" from __future__ import annotations -import json +import hashlib +import math import os import time from contextlib import contextmanager from pathlib import Path -from goofish_cli.core.errors import RateLimitedError +from goofish_cli.core.errors import GoofishError, RateLimitedError +from goofish_cli.core.local_state import locked_state STATE_PATH = Path.home() / ".goofish-cli" / "limiter.json" DEFAULT_WRITE_RPM = 1 +DEFAULT_MEDIA_RPM = 9 -def _rpm() -> int: +def _rpm(bucket: str = "") -> int: + variable = "GOOFISH_MEDIA_WRITE_RPM" if bucket == "media.write" else "GOOFISH_WRITE_RPM" + default = DEFAULT_MEDIA_RPM if bucket == "media.write" else DEFAULT_WRITE_RPM try: - return max(1, int(os.environ.get("GOOFISH_WRITE_RPM", DEFAULT_WRITE_RPM))) + return max(1, int(os.environ.get(variable, default))) except ValueError: - return DEFAULT_WRITE_RPM - - -def _load() -> dict[str, list[float]]: - if not STATE_PATH.exists(): - return {} - try: - return json.loads(STATE_PATH.read_text()) - except (json.JSONDecodeError, OSError): - return {} - - -def _save(state: dict[str, list[float]]) -> None: - STATE_PATH.parent.mkdir(parents=True, exist_ok=True) - STATE_PATH.write_text(json.dumps(state)) + return default -def check(bucket: str) -> None: - """消耗一个令牌。超限抛 RateLimitedError。""" - now = time.time() +def check(bucket: str, *, account: str = "") -> None: window = 60.0 - rpm = _rpm() - state = _load() - hits = [t for t in state.get(bucket, []) if now - t < window] - if len(hits) >= rpm: - wait = window - (now - hits[0]) - raise RateLimitedError( - f"限流:bucket={bucket} 每 {window:.0f}s 上限 {rpm},再等 {wait:.1f}s" - ) - hits.append(now) - state[bucket] = hits - _save(state) + rpm = _rpm(bucket) + path = Path(os.environ.get("GOOFISH_LIMITER_PATH", str(STATE_PATH))).expanduser() + key = f"account/{hashlib.sha256(account.encode()).hexdigest()}/{bucket}" if account else bucket + with locked_state(path) as state: + now = time.time() + for name, values in state.items(): + if not isinstance(name, str) or not isinstance(values, list) or any( + isinstance(value, bool) or not isinstance(value, (float, int)) + or not math.isfinite(value) or value < 0 for value in values + ): + raise GoofishError("限流状态格式无效;未清空预算,未发起写请求") + hits = [t for t in state.get(key, []) if now - t < window] + # 旧版没有账号字段,最近的共享预算仍须自然过期,不能借升级绕过。 + legacy = [t for t in state.get(bucket, []) if now - t < window] if account else [] + active = sorted(hits + legacy) + if len(active) >= rpm: + wait = max(0, window - (now - active[0])) + error = RateLimitedError(f"限流:bucket={bucket} 每 {window:.0f}s 上限 {rpm},再等 {wait:.1f}s") + error.retry_after = wait + raise error + state[key] = hits + [now] @contextmanager -def acquire(bucket: str): - check(bucket) +def acquire(bucket: str, *, account: str = ""): + check(bucket, account=account) yield diff --git a/src/goofish_cli/core/local_state.py b/src/goofish_cli/core/local_state.py new file mode 100644 index 0000000..2d0ecf7 --- /dev/null +++ b/src/goofish_cli/core/local_state.py @@ -0,0 +1,55 @@ +"""本机护栏状态:稳定锁文件覆盖事务,临时文件落盘后原子替换。""" +from __future__ import annotations + +import json +import os +import tempfile +from collections.abc import Iterator +from contextlib import contextmanager +from pathlib import Path +from typing import Any + +from filelock import FileLock, Timeout + +from goofish_cli.core.errors import GoofishError + + +def read_json(path: Path) -> dict[str, Any]: + if not path.exists(): + return {} + try: + value = json.loads(path.read_text()) + except (OSError, ValueError) as exc: + raise GoofishError("本机护栏状态不可读取;保留原文件,请检查配置或恢复备份") from exc + if not isinstance(value, dict): + raise GoofishError("本机护栏状态格式无效;未清空已有状态") + return value + + +def atomic_write(path: Path, value: dict[str, Any]) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + fd, temporary = tempfile.mkstemp(prefix=path.name + ".", dir=path.parent) + try: + with os.fdopen(fd, "w") as stream: + json.dump(value, stream, allow_nan=False) + stream.flush() + os.fsync(stream.fileno()) + os.chmod(temporary, 0o600) + os.replace(temporary, path) + finally: + if os.path.exists(temporary): + os.unlink(temporary) + + +@contextmanager +def locked_state(path: Path) -> Iterator[dict[str, Any]]: + path.parent.mkdir(parents=True, exist_ok=True) + try: + with FileLock(str(path) + ".lock", timeout=5, mode=0o600): + value = read_json(path) + yield value + atomic_write(path, value) + except Timeout as exc: + raise GoofishError("本机护栏状态锁超时,未发起写请求") from exc + except OSError: + raise GoofishError("本机护栏状态无法持久化,请检查目录权限或磁盘空间") from None diff --git a/src/goofish_cli/core/message_history.py b/src/goofish_cli/core/message_history.py new file mode 100644 index 0000000..8e70697 --- /dev/null +++ b/src/goofish_cli/core/message_history.py @@ -0,0 +1,177 @@ +"""消息页读取:复用已就绪的连接,保留身份、消息 ID 和活动时间。""" + +from __future__ import annotations + +import asyncio +import base64 +import json +import math +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager, suppress +from typing import Any + +from websockets.exceptions import ConnectionClosed, InvalidHandshake + +from goofish_cli.core.errors import GoofishError +from goofish_cli.core.session import Session +from goofish_cli.core.sign import generate_mid +from goofish_cli.core.token import get_access_token +from goofish_cli.core.ws import connect, heartbeat_loop, recv_ack, register, wait_ready + + +def normalize_id(value: Any) -> str: + """LWP 与 HTTP 使用同一不带域后缀的标识。""" + return str(value or "").removesuffix("@goofish").strip() + + +def timestamp(value: Any) -> int | None: + if isinstance(value, bool): + return None + try: + result = int(value) + except (ValueError, TypeError, OverflowError): + return None + return result if result > 0 else None + + +def parse_message(model: dict[str, Any], cid: str) -> dict[str, Any]: + if not isinstance(model, dict) or not isinstance(model.get("message"), dict): + raise GoofishError("消息记录缺少 message 对象") + message = model.get("message") or {} + extension = message.get("extension") or {} + content = message.get("content") or {} + if not isinstance(extension, dict) or not isinstance(content, dict): + raise GoofishError("消息扩展或正文格式无效") + payload = content + parse_error = "" + custom = content.get("custom") or {} + if not isinstance(custom, dict): + raise GoofishError("消息 custom 正文格式无效") + if custom.get("data"): + try: + payload = json.loads(base64.b64decode(custom["data"])) + except (ValueError, TypeError, UnicodeDecodeError): + payload = { + "contentType": content.get("contentType"), + "summary": custom.get("summary", ""), + } + parse_error = "消息正文无法解码" + sender = message.get("sender") or {} + receivers = message.get("receivers") or [] + if not isinstance(sender, dict) or not isinstance(receivers, list): + raise GoofishError("消息参加者格式无效") + result = { + "cid": normalize_id(message.get("cid")) or cid, + "message_id": str(message.get("messageId") or ""), + "created_at": timestamp(message.get("createAt")), + "session_type": timestamp(extension.get("sessionType")), + "send_user_id": normalize_id(extension.get("senderUserId") or sender.get("uid")), + "send_user_name": extension.get("reminderTitle", ""), + "receiver_user_ids": [ + normalize_id(r.get("uid") if isinstance(r, dict) else r) for r in receivers + ], + "message": payload, + } + if parse_error: + result["parse_error"] = parse_error + return result + + +@asynccontextmanager +async def history_connection(session: Session) -> AsyncIterator[Any]: + token = get_access_token(session) + async with connect(session) as ws: + mids = await register(ws, session, token) + heartbeat = asyncio.create_task(heartbeat_loop(ws)) + try: + if not await wait_ready(ws, mids=mids): + raise GoofishError("IM 消息读取连接未就绪") + yield ws + finally: + heartbeat.cancel() + with suppress(asyncio.CancelledError): + await heartbeat + + +async def read_page(ws: Any, cid: str, cursor: Any, page_size: int) -> dict[str, Any]: + mid = generate_mid() + await ws.send( + json.dumps( + { + "lwp": "/r/MessageManager/listUserMessages", + "headers": {"mid": mid}, + "body": [f"{cid}@goofish", False, cursor, page_size, False], + } + ) + ) + ack = await recv_ack(ws, mid, timeout=10.0) + if ack is None: + raise GoofishError("读取消息页超时") + if ack.get("code") != 200: + raise GoofishError(f"读取消息页被拒绝:code={ack.get('code')}") + body = ack.get("body") + if not isinstance(body, dict) or not isinstance(body.get("userMessageModels"), list): + raise GoofishError("消息页响应缺少 userMessageModels") + return body + + +async def read_history( + session: Session, cid: str, page_size: int = 20, limit: int = 0, timeout: float = 30.0 +) -> list[dict[str, Any]]: + cid = normalize_id(cid) + if not cid or "@" in cid: + raise GoofishError("cid 必须是有效会话标识") + if not 1 <= page_size <= 100 or limit < 0 or timeout <= 0 or not math.isfinite(timeout): + raise GoofishError("limit-per-page 必须为 1–100,limit 不得为负,timeout 必须大于 0") + messages: dict[str | tuple[str, int], dict[str, Any]] = {} + cursor: Any = 9007199254740991 + cursors: set[str] = set() + try: + async with asyncio.timeout(timeout), history_connection(session) as ws: + while True: + if str(cursor) in cursors: + raise GoofishError("消息分页游标没有推进") + cursors.add(str(cursor)) + count = min(page_size, limit - len(messages)) if limit else page_size + body = await read_page(ws, cid, cursor, count) + for model in body["userMessageModels"]: + item = parse_message(model, cid) + # 缺少服务端 ID 时保留每次出现,不能以相同正文/时间推断同一消息。 + key = item["message_id"] or ("missing_id", len(messages)) + messages[key] = item + if (limit and len(messages) >= limit) or body.get("hasMore") not in (1, "1", True): + break + cursor = body.get("nextCursor") + if cursor is None or not body["userMessageModels"]: + raise GoofishError("消息分页声称有下一页,但没有有效游标或消息") + except TimeoutError as exc: + raise GoofishError("读取会话历史超时,可增大 --timeout 或使用 --limit") from exc + except (ConnectionClosed, OSError, InvalidHandshake) as exc: + raise GoofishError(f"IM 消息连接失败或中断:{type(exc).__name__}") from None + ordered = sorted(messages.values(), key=lambda m: m["created_at"] or 0) + return ordered[-limit:] if limit else ordered + + +async def recent_messages( + session: Session, cids: list[str], timeout: float = 30.0 +) -> tuple[dict[str, list[dict[str, Any]]], dict[str, str]]: + """单连接逐会话取最近一页;失败保留已读摘要,并明确未完成的会话。""" + histories: dict[str, list[dict[str, Any]]] = {} + errors: dict[str, str] = {} + pending = list(dict.fromkeys(normalize_id(cid) for cid in cids)) + if not pending: + return histories, errors + try: + async with asyncio.timeout(timeout), history_connection(session) as ws: + for cid in pending: + try: + body = await read_page(ws, cid, 9007199254740991, 20) + histories[cid] = [parse_message(m, cid) for m in body["userMessageModels"]] + except GoofishError as exc: + errors[cid] = str(exc) + except (TimeoutError, ConnectionClosed, OSError, InvalidHandshake, GoofishError) as exc: + reason = str(exc) if isinstance(exc, GoofishError) else "会话摘要读取超时或连接断开" + errors.update( + {cid: reason for cid in pending if cid not in histories and cid not in errors} + ) + return histories, errors diff --git a/src/goofish_cli/core/search_data.py b/src/goofish_cli/core/search_data.py new file mode 100644 index 0000000..544cf69 --- /dev/null +++ b/src/goofish_cli/core/search_data.py @@ -0,0 +1,68 @@ +"""浏览器实际搜索响应的语义适配,保留标签,未知属性不按位置猜。""" +from __future__ import annotations + +import re +from typing import Any + + +def _price(value: Any) -> str | None: + if isinstance(value, bool) or value is None: + return None + text = str(value).strip() + return text if re.fullmatch(r"(?:¥|¥)?\s*\d+(?:\.\d{1,2})?", text) else None + + +def normalize_search(raw: dict[str, Any]) -> dict[str, Any]: + data = raw.get("data") or {} + info = data.get("resultInfo") or {} + control = info.get("searchResControlFields") if isinstance(info, dict) else None + if not isinstance(control, dict) or not isinstance(data.get("resultList"), list): + return {"items": [], "schema_unknown": True, "result_kind": "unknown"} + empty = control.get("hasItems") is False or control.get("srpFeedsItemsDataEmpty") is True or control.get("numFound") == 0 + if not empty and control.get("hasItems") is not True: + return {"items": [], "schema_unknown": True, "result_kind": "unknown"} + kind = "empty" if empty else "related_search" if control.get("similar") is True else "search" + items = [] + for row in data.get("resultList") or []: + main = (row.get("data") or {}).get("item", {}).get("main", {}) + ex = main.get("exContent") or {} + if not ex.get("itemId") or not ex.get("title"): + continue + tags = [] + attributes = {} + original = _price(ex.get("originalPrice")) + badge = None + for group, values in (ex.get("fishTags") or {}).items(): + for entry in values.get("tagList") or []: + td = entry.get("data") or {} + text = td.get("content") + if not isinstance(text, str) or td.get("type") == "img": + continue + tags.append({"group": group, "type": td.get("type"), "text": text}) + name = td.get("propertyText") or td.get("propertyName") + if name: + attributes[str(name)] = text + if td.get("type") in ("strikethroughText", "strikeThroughText", "lineThroughText") or td.get("strikethrough") is True: + original = _price(text) or original + if group == "r4" and badge is None: + badge = text + detail = ex.get("detailParams") or {} + amount = _price(detail.get("soldPrice")) + if amount is None: + parts = ex.get("price") or [] + amount = _price("".join(str(p.get("text") or "") for p in parts + if p.get("type") in ("sign", "integer", "decimal"))) + items.append({ + "item_id": str(ex["itemId"]), "title": ex["title"], + "url": f"https://www.goofish.com/item?id={ex['itemId']}", + "price": amount if amount is None or amount.startswith(("¥", "¥")) else "¥" + amount, + "original_price": original, + "condition": attributes.get("成色"), "brand": attributes.get("品牌"), + "attributes": attributes, "labels": tags, + "location": ex.get("area") or None, + "seller_nick": ex.get("userNickName") or detail.get("userNick") or None, + "badge": badge, "extra": " | ".join(t["text"] for t in tags), + "source": "browser_search_response", "result_kind": kind, + }) + return {"items": [] if empty else items, "recommendations": items if empty else [], "empty": empty, "result_kind": kind, + "reported_match_count": control.get("numFound"), "requiresAuth": False, "blocked": False} diff --git a/src/goofish_cli/core/write_operation.py b/src/goofish_cli/core/write_operation.py new file mode 100644 index 0000000..62529c1 --- /dev/null +++ b/src/goofish_cli/core/write_operation.py @@ -0,0 +1,15 @@ +"""实际写入口共用护栏:先检查熔断,再占一次账号与业务桶的写预算。""" +from __future__ import annotations + +from collections.abc import Iterator +from contextlib import contextmanager + +from goofish_cli.core.guard import watch +from goofish_cli.core.limiter import acquire +from goofish_cli.core.session import Session + + +@contextmanager +def write_operation(session: Session, bucket: str) -> Iterator[None]: + with watch(account=session.unb), acquire(bucket, account=session.unb): + yield diff --git a/src/goofish_cli/core/ws.py b/src/goofish_cli/core/ws.py index c51002d..fd2fb09 100644 --- a/src/goofish_cli/core/ws.py +++ b/src/goofish_cli/core/ws.py @@ -308,24 +308,22 @@ async def create_chat(ws: ClientConnection, *, myid: str, toid: str, item_id: st return mid -async def collect_session_cids( - session: Session, duration: float = 5.0 -) -> list[dict[str, Any]]: +async def collect_session_cids(session: Session, duration: float = 5.0) -> list[dict[str, Any]]: """连 WS + /reg + ackDiff(pts=0),收 duration 秒,返回所有 push 到的 session cid。 涵盖两路: - `/s/vulcan` 里 `operation.sessionInfo`(历史会话激活事件) - `extract_meta_event` 的 `new_msg`(最新未读通知) - 返回字段只有 `cid / session_type / item_id / last_msg_ts / last_msg_id`, - **没有** peer_user_id / 昵称 / 消息正文。原因:sessionInfo.extensions.extUserId/ - itemSellerId 是卖家 ID(在登录账号作为卖家时就是自己),不能无脑当 peer。 - 上游想拿 peer/正文要自己走 `message history `。 + 返回会话标识、类型、商品、消息时间通知和 owner/extUser 角色参加者。 + 角色字段不能直接当对端,也不包含昵称和消息正文;调用方结合当前账号与 + 最近消息页确认身份。itemSellerId 与 squadName 不用于对端身份推断。 """ token = get_access_token(session) acc: dict[str, dict[str, Any]] = {} async with connect(session) as ws: + reg_mid = generate_mid() reg = { "lwp": "/reg", "headers": { @@ -337,7 +335,7 @@ async def collect_session_cids( "wv": "im:3,au:3,sy:6", "sync": "0,0;0;0;", "did": session.device_id, - "mid": generate_mid(), + "mid": reg_mid, }, } await ws.send(json.dumps(reg)) @@ -373,6 +371,8 @@ async def collect_session_cids( msg = json.loads(raw) except json.JSONDecodeError: continue + if (msg.get("headers") or {}).get("mid") == reg_mid and msg.get("code") != 200: + raise GoofishError(f"IM 会话发现注册失败:code={msg.get('code')}") with suppress(Exception): await ws.send(json.dumps(build_ack(msg))) @@ -387,9 +387,18 @@ async def collect_session_cids( ext = sess_info.get("extensions") or {} entry = acc.setdefault(cid, {"cid": cid}) entry["session_type"] = ( - sess_info.get("sessionType") or decoded.get("chatType") or entry.get("session_type", 0) + sess_info.get("sessionType") + or decoded.get("chatType") + or entry.get("session_type", 0) ) entry["item_id"] = str(ext.get("itemId") or entry.get("item_id", "")) + entry["participant_user_ids"] = list( + dict.fromkeys( + str(ext[k]).removesuffix("@goofish") + for k in ("ownerUserId", "extUserId") + if ext.get(k) + ) + ) continue # b) new_msg:{"1":"cid@goofish","2":1,"3":msgId,"4":ts} @@ -401,80 +410,36 @@ async def collect_session_cids( entry["last_msg_ts"] = meta.get("ts", "") finally: hb.cancel() + with suppress(asyncio.CancelledError): + await hb # 统一字段 + 填默认值 out: list[dict[str, Any]] = [] for cid, e in acc.items(): - out.append({ - "cid": cid, - "session_type": int(e.get("session_type") or 0), - "item_id": e.get("item_id", "") or "", - "last_msg_id": e.get("last_msg_id", ""), - "last_msg_ts": e.get("last_msg_ts", ""), - }) + out.append( + { + "cid": cid, + "session_type": int(e.get("session_type") or 0), + "item_id": e.get("item_id", "") or "", + "last_msg_id": e.get("last_msg_id", ""), + "last_msg_ts": e.get("last_msg_ts", ""), + "participant_user_ids": e.get("participant_user_ids", []), + } + ) return out async def list_user_messages( - session: Session, cid: str, limit_per_page: int = 20 + session: Session, + cid: str, + limit_per_page: int = 20, + limit: int = 0, + timeout: float = 30.0, ) -> list[dict[str, Any]]: - """一次性拉指定会话的历史消息(翻页直到 hasMore=0)。""" - token = get_access_token(session) - messages: list[dict[str, Any]] = [] - send_mid = generate_mid() - req = { - "lwp": "/r/MessageManager/listUserMessages", - "headers": {"mid": send_mid}, - "body": [f"{cid}@goofish", False, 9007199254740991, limit_per_page, False], - } - - async with connect(session) as ws: - await register(ws, session, token) - hb = asyncio.create_task(heartbeat_loop(ws)) - try: - async for raw in ws: - try: - msg = json.loads(raw) - except json.JSONDecodeError: - continue - with suppress(Exception): - await ws.send(json.dumps(build_ack(msg))) - - lwp = msg.get("lwp") - if lwp == "/s/vulcan": - await ws.send(json.dumps(req)) - continue - - recv_mid = (msg.get("headers") or {}).get("mid", "") - if recv_mid != send_mid: - continue + """读取历史并保留消息标识与真实时间;limit=0 翻页到底。""" + from goofish_cli.core.message_history import read_history - body = msg.get("body") or {} - models = body.get("userMessageModels") or [] - for um in models: - try: - ext = um["message"]["extension"] - data_b64 = um["message"]["content"]["custom"]["data"] - payload = json.loads(base64.b64decode(data_b64).decode("utf-8")) - messages.insert(0, { - "send_user_id": ext.get("senderUserId", ""), - "send_user_name": ext.get("reminderTitle", ""), - "message": payload, - }) - except Exception as e: # noqa: BLE001 - logger.debug(f"parse history item failed: {e}") - - has_more = body.get("hasMore") == 1 - if has_more: - send_mid = generate_mid() - req["headers"]["mid"] = send_mid - req["body"][2] = body.get("nextCursor") - await ws.send(json.dumps(req)) - else: - break - finally: - hb.cancel() - return messages + return await read_history(session, cid, limit_per_page, limit, timeout) def _decode_one(raw: str) -> dict[str, Any] | None: diff --git a/src/goofish_cli/mcp_server.py b/src/goofish_cli/mcp_server.py index b96c2a4..71c2594 100644 --- a/src/goofish_cli/mcp_server.py +++ b/src/goofish_cli/mcp_server.py @@ -16,6 +16,7 @@ 注意:`uvx goofish-mcp` 单写会因 PyPI 无同名包而解析失败——不是别名能救的,这是 uvx 按 command 名查包的默认行为决定的。 """ + from __future__ import annotations import asyncio @@ -26,6 +27,7 @@ from mcp.server.fastmcp import FastMCP from goofish_cli.core import GoofishError, iter_commands +from goofish_cli.core.errors import PartialResultError from goofish_cli.core.registry import discover mcp = FastMCP("goofish") @@ -57,12 +59,15 @@ async def handler(**kwargs: Any) -> dict[str, Any]: try: result = await asyncio.to_thread(cmd.func, **kwargs) return {"ok": True, "data": result} + except PartialResultError as e: + return {"ok": False, "data": e.data, "error": e.error, "message": str(e)} except GoofishError as e: return {"ok": False, "error_type": type(e).__name__, "message": str(e)} handler.__name__ = tool_name handler.__doc__ = doc - handler.__signature__ = sig # type: ignore[attr-defined] + # 输入沿用业务签名,输出是适配器的统一包裹,不能声明成业务层的 list。 + handler.__signature__ = sig.replace(return_annotation=dict[str, Any]) # type: ignore[attr-defined] mcp.tool(name=tool_name, description=doc)(handler) diff --git a/tests/test_guard.py b/tests/test_guard.py index 71cddf0..d580669 100644 --- a/tests/test_guard.py +++ b/tests/test_guard.py @@ -26,9 +26,11 @@ def test_reset_clears_circuit(tmp_path, monkeypatch): from goofish_cli.core import guard guard.trip("test") - assert guard._load() > 0 + from goofish_cli.core.errors import RiskControlError + with pytest.raises(RiskControlError): + guard.check() guard.reset() - assert guard._load() == 0 + guard.check() # reset 后 watch 不再抛 with guard.watch(): pass diff --git a/tests/test_list_chats.py b/tests/test_list_chats.py index 4a981da..4ae3e64 100644 --- a/tests/test_list_chats.py +++ b/tests/test_list_chats.py @@ -1,4 +1,5 @@ """验 message list-chats 的 session 结构解析。""" + from __future__ import annotations from goofish_cli.commands.message.list_chats import _parse_session, _watch_record @@ -57,17 +58,8 @@ def test_parse_session_falls_back_to_fish_nick(): def test_parse_session_missing_fields_defaults(): out = _parse_session({}) - assert out == { - "session_id": "", - "peer_nick": "", - "peer_user_id": "", - "unread": 0, - "last_msg": "", - "ts": 0, - "session_type": 0, - "item_id": "", - "source": "baseline", - } + assert out["ts"] is None + assert out["peer_user_id"] == "" def test_parse_session_null_summary(): @@ -82,14 +74,16 @@ def test_parse_session_null_summary(): def test_watch_record_shape(): - out = _watch_record({ - "cid": "test-cid", - "session_type": 1, - "peer_user_id": "test-user", - "item_id": "900123", - "last_msg_id": "msg-x", - "last_msg_ts": "1776780018537", - }) + out = _watch_record( + { + "cid": "test-cid", + "session_type": 1, + "peer_user_id": "test-user", + "item_id": "900123", + "last_msg_id": "msg-x", + "last_msg_ts": "1776780018537", + } + ) assert out["session_id"] == "test-cid" assert out["peer_user_id"] == "test-user" assert out["item_id"] == "900123" @@ -100,7 +94,7 @@ def test_watch_record_shape(): assert out["last_msg"] == "" -def test_watch_record_invalid_ts_defaults_zero(): +def test_watch_record_missing_activity_time_is_unknown(): out = _watch_record({"cid": "123"}) - assert out["ts"] == 0 + assert out["ts"] is None assert out["source"] == "watch" diff --git a/tests/test_search.py b/tests/test_search.py index 7a48642..1f0b295 100644 --- a/tests/test_search.py +++ b/tests/test_search.py @@ -2,14 +2,11 @@ from __future__ import annotations import asyncio -from unittest.mock import patch import pytest -import goofish_cli.commands.search.search as search_mod from goofish_cli.commands.search.search import __test__ as t -from goofish_cli.commands.search.search import _item_id_from_url, search -from goofish_cli.core.errors import AuthRequiredError, GoofishError +from goofish_cli.commands.search.search import _item_id_from_url def test_normalize_limit_clamps_and_defaults(): @@ -336,80 +333,3 @@ def test_walk_pages_non_dict_pagination_state_is_error(): _walk(FakePage([bad]), [_item("100")], {"100"}, 1, pages=3, limit=200) ) assert (fetched, reason) == (1, "error"), f"bad={bad!r}" - - -# --------------------------------------------------------------------------- -# 公开 search() 入口回归(#28 曾发生 _run 返回 None 而测试全绿的事故, -# 顶层组装层必须钉住:items/rank/query/分页元数据/错误传播) -# --------------------------------------------------------------------------- -class FakePageCtx: - """让 `_run` 真正执行的 goofish_page 替身:返回同一个 FakePage。""" - - def __init__(self, page: FakePage) -> None: - self._page = page - - async def __aenter__(self) -> FakePage: - return self._page - - async def __aexit__(self, *exc) -> None: - return None - - -def test_search_returns_ranked_items_with_metadata(): - """成功路径:_run 组装 rank/item_id/分页元数据,search 填充 query。""" - # 第 1 页(extract)→ pagination → 等待 → 第 2 页 extract → 终态 pagination - page = FakePage( - [ - _page_payload(["100", "101"]), - _pagination(True, 50), - True, - _page_payload(["102"]), - _pagination(False, 50), - ] - ) - with patch.object(search_mod, "goofish_page", lambda **kw: FakePageCtx(page)): - result = search("X570", limit=50, pages=2) - assert result["query"] == "X570" - assert result["total"] == 3 - assert result["pages_fetched"] == 2 - assert result["pages_total"] == 50 - assert result["stopped_reason"] == "pages_reached" - assert [it["rank"] for it in result["items"]] == [1, 2, 3] - assert [it["item_id"] for it in result["items"]] == ["100", "101", "102"] - assert result["items"][0]["title"] == "item-100" - - -def test_search_single_page_default_contract(): - """默认参数(pages=1, limit=20)行为与分页前一致:只抓第 1 页。""" - page = FakePage( - [ - _page_payload([str(i) for i in range(100, 110)]), - _pagination(True, 50), # 终态读取(total_pages 补取) - ] - ) - with patch.object(search_mod, "goofish_page", lambda **kw: FakePageCtx(page)): - result = search("X570") - assert result["pages_fetched"] == 1 - assert result["total"] == 10 - assert result["query"] == "X570" - - -def test_search_propagates_auth_required_error(): - """登录墙错误从 _run 传播到 search() 上层(顶层契约的一部分)。""" - # AUTH_WALL_ATTEMPTS=2:两次导航各消费一份 extract,仍零卡片才抛 - page = FakePage([_page_payload([], requiresAuth=True), _page_payload([], requiresAuth=True)]) - with ( - patch.object(search_mod, "goofish_page", lambda **kw: FakePageCtx(page)), - pytest.raises(AuthRequiredError, match="www.goofish.com"), - ): - search("X570") - - -def test_search_propagates_structure_error(): - """结构突变(0 卡片且非 auth/empty/blocked)按原语义从顶层抛出。""" - page = FakePage([_page_payload([], bodyPreview="weird page")]) - with ( - patch.object(search_mod, "goofish_page", lambda **kw: FakePageCtx(page)), - pytest.raises(GoofishError, match="DOM 结构已变"), - ): - search("X570")