From 137ed6f408eb4bf957027b88751aeb4d4d141fab Mon Sep 17 00:00:00 2001 From: jesse-ux Date: Mon, 5 Oct 2026 02:33:36 +0800 Subject: [PATCH] fix(consult): stop the consult-gate VedAstro wait and record step timing Consult turns reuse a same-day official snapshot inside the rectification gate, or record deferred_in_consultation instead of starting another 4-second snapshot. Classification, step 0, step 1, and the time to the first answer character go into the existing usage metadata and a new admin usage column. BUG-1231, BUG-1232 --- CHANGELOG.md | 7 + docs/BUG_HISTORY.md | 30 +++ ...RESS-consult-latency-quickwins-20261005.md | 115 +++++++++ docs/tasks/README.md | 2 +- frontend/DESIGN.md | 5 + frontend/src/app/api/admin/usage/route.ts | 68 ++++- frontend/src/app/api/consult/route.ts | 16 +- .../admin/billing-operations-resources.tsx | 9 + frontend/src/lib/agent-observability.ts | 24 ++ frontend/src/lib/consultation-step-timing.ts | 237 ++++++++++++++++++ frontend/src/lib/stream-agent-response.ts | 25 ++ frontend/src/mastra/consultation-tools.ts | 13 + .../consult-latency-timing-20261005.test.ts | 199 +++++++++++++++ scripts/jyotish_api_server.py | 18 +- scripts/vedastro_consultation_snapshot.py | 230 +++++++++++++++++ tests/test_vedastro_consultation_snapshot.py | 199 +++++++++++++++ 16 files changed, 1188 insertions(+), 9 deletions(-) create mode 100644 docs/tasks/PROGRESS-consult-latency-quickwins-20261005.md create mode 100644 frontend/src/lib/consultation-step-timing.ts create mode 100644 frontend/tests/consult-latency-timing-20261005.test.ts create mode 100644 scripts/vedastro_consultation_snapshot.py create mode 100644 tests/test_vedastro_consultation_snapshot.py diff --git a/CHANGELOG.md b/CHANGELOG.md index ea0a2c84..1ba07a0f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,12 @@ # 印度占星 Skill 更新日志 +## 2026-10-05 — 普通对话不再白等校正闸的外部快照,用量页能看分段时间(未上线) + +- 问父母、问年运这类普通对话,每个领域以前会在校正闸里再等一次外部星历快照,大约多 4 秒,失败了同一天还会再等。现在咨询路径复用当天已经拿到的官方结果;没有的话不再干等,记成「咨询里先不做」,同一天不重复等。生时校正自己的核对不变。顶层那一次外部核对仍会做(BUG-1231)。 +- 用量页多一列「分段」:分类耗时,第 0 步和第 1 步各自的耗时与推理、输出、输入、缓存输入 token,以及从第 1 步开始到写出第一个正文字的时间。供应商没返回的数显示为「—」。不改数据库表。对话提示词和数据卡没改(BUG-1232)。 +- 本机没有装 VedAstro SDK,任务书里「大约 4.6 秒降到 1 秒以内」没有在这台机器上重测。当前测试环境部署的仍是改前版本,真实分段要等这次上线后在用量页看。 +- Skill 版本不变;不改数据库结构;不改生时校正打分。 + ## 2026-10-04 — 门禁日志改成摘要,失败时还能看到是哪一项(未上线) - 测试环境的质量门禁以前把七项检查的全部明细打在同一步里,日志有四万行以上,失败页面打不开(BUG-1230)。 diff --git a/docs/BUG_HISTORY.md b/docs/BUG_HISTORY.md index eac3d91b..afc1b700 100644 --- a/docs/BUG_HISTORY.md +++ b/docs/BUG_HISTORY.md @@ -16627,3 +16627,33 @@ - 相关记录:BUG-995(前端测试要同时看 cancelled 和退出码)。 - 复发自:无 - 修复版本:待发布 + +## BUG-1231 | 咨询链校正闸同步白等 VedAstro 官方快照,失败还不缓存 + +- 状态:resolved(分支 `codex/consult-latency-quickwins-20261005`) +- 首次发现 / 最近更新:2026-10-04 / 2026-10-05 +- 来源:10-04 只读审计。任务书 `docs/tasks/TASK-consult-latency-quickwins-20261005.md`。写任务书时,在装了 vedastro SDK 的沙箱网络下按产品路径测得每个领域约 4.6 秒,第二轮仍是 4.6 秒。这组秒数不是本轮在这台机器上重测的。 +- 影响面:普通对话引擎里校正闸附带的 `rectification.vedastro_gateway`。顶层前台 VedAstro(BUG-301)不改。直接调用 `/api/rectification_gate`、请求里没有咨询标记时,仍走原来的官方快照。 +- 现象:同一份咨询响应里,顶层 `vedastro_gateway` 可以已经是 `official_verified`(BUG-727 的当日缓存),校正闸里的 `official_closure_reason` 仍是 `official_raw_response_missing`。校正闸每次再起一个最多 4 秒的子进程。失败结果不进正缓存,同一天每一轮都再等。屏蔽外网时引擎每个领域约 0.75 秒,这 4 秒是额外等待。生产机器是不是也在等,本轮没能在 staging 上确认(进度记录 T0)。 +- 根因:咨询工作流把「外部证据可以后做」这个标记只传给顶层前台任务。校正闸仍调用原来的官方快照。出生时间进校正闸时被规范成浮点,和正缓存的键对不上,所以当天已经有的官方结果也复用不到。 +- 修复:只在校正闸、且请求带了这个标记时,改走咨询专用路径。先用原始请求里的出生数据、岁差、交点、UTC 日期查正缓存,命中已验证结果就复用;否则查当天的负缓存;都没有就记 `official_closure_reason = deferred_in_consultation`,不起子进程,也不标成已验证。负缓存单独存放,只认当天的 UTC 日期。没有这个标记的校正请求仍调用原来的 gateway。 +- 验证:`tests/test_vedastro_consultation_snapshot.py`。咨询路径不调用快照子进程,也不调用原来的 gateway。负缓存同日命中、次日失效,邮箱不进缓存。当天已验证快照优先于负缓存,整数小时和规范后的浮点小时靠原始请求对齐。没有标记的校正面,返回包与现场结果逐项相同。增长合同通过,没有新增类方法。本机没有 vedastro SDK,4.6 秒降到 1 秒以内没有重测。v5 打分文件没有改动,本轮没有重跑 77 例。 +- 防复发:上述测试。开关只放在校正闸里,不放进顶层 gateway,避免把 BUG-301 的前台核对也跳掉。 +- 相关记录:BUG-727(当日正缓存)、BUG-301(顶层前台照旧,本单不推翻)。 +- 复发自:无 +- 修复版本:待发布 + +## BUG-1232 | 普通对话看不到第 0/1 步耗时和推理 token,分类耗时进不了用量 + +- 状态:resolved(分支 `codex/consult-latency-quickwins-20261005`) +- 首次发现 / 最近更新:2026-10-04 / 2026-10-05 +- 来源:10-04 只读审计。任务书同上。审计写答题步大约 1 万推理 token、大约 45 秒,但线上记录里没有逐步耗时和推理 token,后面的「限制推理强度」「精简说明」对不上数。这组数字不是本轮实测。 +- 影响面:`[agent-observability]` 日志,以及用量账本里已有的 JSON 字段。不改表,不加迁移。不改提示词和数据卡。 +- 现象:已有首字节、工具耗时、答案首输出、整轮耗时、总 token。没有第 0 步和第 1 步各自的耗时,也没有推理、输出、输入、缓存输入 token。分类耗时没有进用量记录。没有「第 1 步开始到第一个正文字」的时间。 +- 根因:流式过程没有按步骤记下供应商返回的 usage。用量记录和用量页都没有这些字段。 +- 修复:流式过程记下第 0 步和第 1 步的耗时和供应商 usage。供应商没返回的数写空,不估算。`answer.reasoning_ms` 从第 1 步开始计到第一个非空白正文字;第 0 步的旁白不计。分类耗时同时写入可观测事件和用量记录。步骤数字写入用量记录的 `modelSteps` 和 `answer.reasoning_ms`。用量页多一列「分段」。各领域工具耗时保持原字段。正文、提示词、出生资料、用户标识不写入这些字段。 +- 验证:`frontend/tests/consult-latency-timing-20261005.test.ts` 4 项通过。模拟步骤的日志和用量字段齐全,不含植入的正文。流式一轮把计时留在运行状态上,公开回执和响应正文里没有这些字段,也没有那段正文。`tsc --noEmit` 0 错。`npm run lint` 0 error,既有 warning 未动。用量页本轮没有登录打开。 +- 防复发:上述测试锁住对话路由、用量列表查询和页面列名。可观测结构是严格对象,多出来的正文字段进不了事件。 +- 相关记录:10-04 咨询耗时审计(任务书事故实证第 2 条)。 +- 复发自:无 +- 修复版本:待发布 diff --git a/docs/tasks/PROGRESS-consult-latency-quickwins-20261005.md b/docs/tasks/PROGRESS-consult-latency-quickwins-20261005.md new file mode 100644 index 00000000..efa74abf --- /dev/null +++ b/docs/tasks/PROGRESS-consult-latency-quickwins-20261005.md @@ -0,0 +1,115 @@ +# PROGRESS:普通对话耗时两项快修 — 2026-10-05 + +基线:任务书写成时已部署的是 `1252de3e`。开工时 `origin/staging` 为 `853772bc`(相对已部署提交只有文档)。执行中又快进到 `4c64733f`(再加一份纳迪研究任务书,产品代码与 `853772bc` 相同)。分支 `codex/consult-latency-quickwins-20261005`,工作树 `.worktrees/consult-latency-quickwins-20261005`。测试基线工作树停在 `853772bc`,没有改它。 + +未推送,未部署。用量页没有登录打开。修复版本:待发布。 + +## 做了什么 + +1. 咨询路径上,校正闸不再为官方快照干等。请求里已有「外部证据可以后做」这个标记时:先复用当天已经验证的官方快照;没有就查当天的失败记录;都没有就记 `deferred_in_consultation`,不起 4 秒子进程,也不标成已验证。直接做生时校正、请求里没有这个标记时,仍走原来的官方快照。顶层那一次前台核对仍做(BUG-301 不推翻)。BUG-1231。 +2. 普通对话补上分段计时。分类耗时、第 0 步、第 1 步的耗时和供应商 token,以及第 1 步开始到第一个正文字的时间,写入运行日志和用量账本里已有的 JSON 字段。用量页多一列「分段」。不改表。BUG-1232。 +3. 记录:BUG-1231、BUG-1232、CHANGELOG、本文件、状态板改为「已实现待验收」。 + +推理强度、精简说明、第 0 步关闭思考,按产品决定不在本单。 + +## T0 环境缺口 + +本轮没有登录态,也没有测试账号。公开接口测到: + +| 检查 | 结果 | +| --- | --- | +| `GET https://staging.jyotisha.chat/api/health` | 200。`deployment.gitCommit` = `1252de3e6363f3097668efa1b1d5d47186c9c2d0` | +| `POST https://staging.jyotisha.chat/api/consultation_workflow` | 404。公网边缘没有这条 Python 引擎路由 | +| 本机 `vedastro` 模块 | 未安装 | + +任务书里的 4.6 秒是沙箱网络下测的。当前测试环境跑的仍是改前版本。要确认生产机器是不是也在等这 4 秒,请在这次上线之前看现在的 `/admin/usage` 或容器日志:工具耗时是否接近「领域数 × 4 秒以上」。这次上线之后,同一处不应再出现这段等待;新的「分段」列用来看分类、第 0 步、第 1 步。 + +## T1 耗时 + +本机没有 VedAstro SDK,产品路径的冷热耗时没有重测。下表「改前」用任务书里的沙箱数字,「改后」是本轮测试能证明的行为。 + +| 场景 | 改前 | 改后 | +| --- | --- | --- | +| 装了 SDK、按产品路径,每个领域 | 约 4.6 秒,第二轮仍约 4.6 秒 | 本机测不了秒数。测试证明:不起快照子进程,状态是 `deferred_in_consultation`,不是已验证 | +| 同一天再问,上次是超时或失败 | 再等约 4.6 秒(失败不进缓存) | 负缓存同日命中,文件不重写,子进程不调用 | +| 下一个 UTC 日 | 任务书未单列 | 当天的负缓存失效,重新记 `deferred_in_consultation` | +| 当天已经有已验证的官方快照 | 顶层能用缓存,校正闸仍去等 | 校正闸复用同一份缓存。整数小时和进闸后变成的浮点小时,用原始请求对齐,避免对不上键 | +| 生时校正自己的请求(没有咨询标记) | 走原来的官方快照 | 仍走原来的 gateway,返回字段与现场结果逐项相同 | +| 顶层前台 VedAstro | 照旧(BUG-301) | 照旧。标记只作用在校正闸,不作用在顶层 gateway | + +3 位公开名人 × 父母 / 年运、冷热各两轮:环境缺口,没有秒数。 + +## T2 产品在哪里看 + +用量页 `/admin/usage` 的「分段」列,旧记录没有这些数时显示「—」。数字来自用量账本已有的 JSON 字段,没有新表、没有迁移。同一组数也在服务端日志 `[agent-observability]` 里。本轮没有登录,页面没有点开;列的文字由单元测试锁住。 + +| 你看到的 | 含义 | 没有数时 | +| --- | --- | --- | +| 分类 180 ms | 分类这一步的耗时 | 分类 — | +| 第 0 步 … · 推理 · 出 · 入 · 缓存入 | 决定要不要排盘的那一步。耗时是毫秒;后面四个是供应商返回的 token | 该项为 — | +| 第 1 步 … | 写答案的那一步,字段同上 | 同上 | +| 写到正文 450 ms | 第 1 步开始,到第一个正文字出现。第 0 步的旁白不计 | 写到正文 — | +| 各领域工具耗时 | 原有字段,本单没有改含义 | 原样 | + +日志和用量记录里的对应名字:`classification.durationMs`、`modelSteps`(只保留第 0 步和第 1 步)、`answer.reasoning_ms`、工具调用上原有的 `durationMs`。供应商没返回的 token 是空值,不估算。这些字段里没有提示词、答案正文、出生资料、用户标识。 + +## 开工预检 + +读了 `docs/research/pre_work_error_ledger.md`。本单不涉及碎片目录或镜像仓,没有再读两份碎片清扫。`scripts/pre_work_check.py --remote-timeout 8 --command-timeout 45` 的结果: + +| 项 | 结果 | +| --- | --- | +| Python 运行时 | 通过 | +| 外部引擎适配器 | 通过 | +| 远端可见 | 通过,远端是 `https://git.copse.top/root/Jyotisha.git` | +| 碎片扫描 | 90 秒超时 | +| 聚焦测试 | 45 秒超时 | +| 检查当时的分支 | 比 `origin/staging` 落后 1 个提交;随后已快进到 `4c64733f` | + +碎片扫描和聚焦测试超时是这次预检命令自己的时限,不是远端同步失败。没有新的镜像路径要记。 + +磁盘:开工时 G: 约 557 GB 空闲,D: 约 135 GB,C: 约 24 GB。 + +## 已完成的检查 + +| 项 | 结果 | +| --- | --- | +| 新增 Python 测试 + 增长合同 | 11 项通过 | +| 两个新 Python 文件的 ruff | 通过。`jyotish_api_server.py` 里原有的 ruff 问题没有顺手改 | +| 前端计时测试 | 4 项通过,0 失败,0 cancelled | +| `tsc --noEmit` | 0 错 | +| `npm run lint` | 0 error,126 条既有 warning,没有为消 warning 改业务代码 | +| v5 打分文件 | 工作区没有改动。没有重跑 77 例 | +| 生时校正打分 | 没改 | + +全量 `pytest tests`、英文对照、隐私测试、全量 `npm test`、快速门:见下面「套件对照」。对照完成前,不把本单说成套件已通过。 + +## 套件对照 + +前端全量 `npm test`(本机,不设 CI)。基线工作树 `853772bc`,本分支在快进后的产品代码上。两边都是退出码 1,cancelled 0,skipped 0。 + +| | tests | pass | fail | cancelled | +| --- | --- | --- | --- | --- | +| 开工基线 | 4844 | 4700 | 144 | 0 | +| 本分支 | 4848 | 4703 | 145 | 0 | + +多出来的 4 个测试是本单的计时测试,全量里 4 个都通过。144 条旧失败名字与基线逐条相同。多出来的那 1 条是 `database roles have no cluster privileges`:全量并行时 Docker 里的 Postgres 正在关闭,`psql` 报 `the database system is shutting down`。这条测试不读本单改过的文件。单独重跑该文件:2 项通过,0 失败。不是这次改动引进的。 + +英文对照 `tests/test_report_english_dictionary.py` 与隐私 `tests/test_repo_privacy_markers.py`:一起跑,退出码 0。 + +`tsc --noEmit`:0 错。`npm run lint`:0 error,126 条既有 warning。 + +全量 `pytest tests`(`--tb=line`,退出码以失败名单为准;日志里的短摘要计数行没有落盘)。用量页没有用浏览器点开。 + +| | 失败条数 | 与对方相比多出来的名字 | +| --- | --- | --- | +| 开工基线 `853772bc` | 130 | — | +| 本分支 | 131 | `tests/test_consultation_native_layers.py::test_the_new_layer_leaves_every_existing_output_unchanged` | + +其余 130 条名字与基线相同。多出来的这一条是既有的盘缓存抖动:缓存命中读回的 VedAstro 请求清单键序和现场重算不一样。这条测试会先跑掉新层再比整份输出,并且会剥掉 `vedastro_gateway`。它用的请求没有咨询标记,不走本单的校正闸开关。全量里失败一次。单独再跑:第一次失败,紧接着再跑一次通过。快速门里又失败一次,失败后用同一比较再跑,差异条数是 0。不记成新 bug。 + +v5 打分文件没有改动,没有重跑 77 例。 + +快速门 `scripts/run_quality_gate.py --profile quick`:退出码 1。1056 通过,1 跳过,1 失败。唯一失败仍是上面那条盘缓存测试。失败后立刻用同一比较再跑一遍,两边剥掉易变字段后差异条数是 0。这条失败跟着缓存时序走,不跟着本单的校正闸开关走。门禁跑完后,5 份待处理的 oracle 模板只多了回车换行,已还原,不进提交。 + +`next build`:这个工作树的 `frontend/node_modules` 是指到主检出的联接,Turbopack 拒绝「联接指向项目根以外」。改用门禁同一条 webpack 构建(`next build --webpack`):编译 3.6 分钟通过,TypeScript 2.5 分钟通过。收集页面数据时停在 `/api/consult`,原因是 Windows 不允许为 Skill 运行别名建符号链接(`EPERM`,`SKILL.md` → 临时目录)。路由表没有打出来,所以没能在这台机器上确认 `/` 仍是 Static。首页源文件没改。没有为了构建去改符号链接,也没有在工作树里重装依赖。 diff --git a/docs/tasks/README.md b/docs/tasks/README.md index 8642578e..41f37aa5 100644 --- a/docs/tasks/README.md +++ b/docs/tasks/README.md @@ -421,5 +421,5 @@ | `TASK-astrologer-rulings-batch5-20261004.md` | `PROGRESS-astrologer-rulings-batch5-20261004.md` | 占星师第六轮:Rath 双主星按 p.43(a)–(e)、罗计尊贵按 Rath Table 9(仅 Narayana 内)、子运方向 p.51 例外与书内分歧标注、Wadiyar 两对标软件特例;年主选不出时不再用 Muntha 主星顶替(对齐上游 03bea6ed);Rath 替换校正 v5 重试算(研究)(BUG 从 1223 起) | 已实现待验收(T2–T5 完成;T1 按红线停下 blocked;T6 研究分支已跑) | 分支 `codex/astrologer-rulings-batch5-20261004`(未推送;BUG-1223~1227;校正分数不变、未升版本);研究分支 `codex/narayana-rath-rectification-trial-20261004`(不合入) | | `TASK-astrologer-rulings-batch6-20261004.md` | `PROGRESS-astrologer-rulings-batch6-20261004.md` | 第七轮裁定(共享仓书面回复,产品采用;问 1 选 A):Rath 版双主星按 p.43 (a)–(e)(BUG-1227 解除 blocked,推翻第五批红线 2 的 Table 17 年数底线);第 5 级宫主度数只倒算计都(p.71 脚注 42);罗计旺陷 ±1 年;同宫两主比经度;BUG 从 1228 起 | 待验收 | 分支 `codex/astrologer-rulings-batch6-20261004`(BUG-1227~1229;校正分数文件未改、未升版本) | | `TASK-gate-log-volume-20261004.md` | `PROGRESS-gate-log-volume-20261004.md` | 门禁 validate 单步日志 4.4 万行 / 2 MB 网页打不开(run 3170):快速门只打摘要(失败给末 200 行 + 日志文件)、前端测试门禁上 dot + 失败汇总(本机仍 TAP)、拆分 validate 为 7 个 step(产品授权改 workflow,只限拆分与重定向);检查一项不少 | **已实现待验收**(BUG-1230;未推送、未部署;Gitea 各 step 页面是否打得开留待推送后由产品确认) | 分支 `codex/gate-log-volume-20261004` | -| `TASK-consult-latency-quickwins-20261005.md` | `PROGRESS-consult-latency-quickwins-20261005.md` | 普通对话耗时两项快修:咨询链校正闸不再同步白等 VedAstro 官方快照子进程(每域约 4 s,复用顶层缓存 + 负缓存,不推翻 BUG-301);补分段计时(第 0/1 步耗时、推理 token、分类耗时进日志与 usage)。推理强度/精简说明待模型 key 另单(BUG 从 1231 起) | 待领取 | 分支 `codex/consult-latency-quickwins-20261005` | +| `TASK-consult-latency-quickwins-20261005.md` | `PROGRESS-consult-latency-quickwins-20261005.md` | 普通对话耗时两项快修:咨询链校正闸不再同步白等 VedAstro 官方快照子进程(每域约 4 s,复用顶层缓存 + 负缓存,不推翻 BUG-301);补分段计时(第 0/1 步耗时、推理 token、分类耗时进日志与 usage)。推理强度/精简说明待模型 key 另单(BUG-1231、BUG-1232) | 已实现待验收 | 分支 `codex/consult-latency-quickwins-20261005` | | `TASK-rectification-nadi-seconds-research-20261005.md` | `PROGRESS-rectification-nadi-seconds-research-20261005.md` | **「纳迪秒级校准」可证伪检验(离线)**:竞品宣传「问前事到天 → 秒级」。本仓主链只用三层小运,Sookshma / Prana 与 D150 从未进评价集;v5 真值 52/77 是整 5 分钟(秒级无真值可对)。N0 五层小运 + D150 底座(前三层与主链对账 0 差)、N1 拟合率 vs 安慰剂日期(核心)、N2 留一件预测、N3 六题后区间内再细分能否提头名、N4 岁差 / 坐标 / 时间扰动的噪声地板、N5 D150 结构层(原文比对 blocked)、N6 结论 + 对外口径草稿。规则先登记再跑;不改生产代码;不得重调 BUG-1091 已关的权重 | 待领取 | 分支 `codex/rectification-nadi-seconds-research-20261005`(BUG-1240) | diff --git a/frontend/DESIGN.md b/frontend/DESIGN.md index 1990cb6e..327368e3 100644 --- a/frontend/DESIGN.md +++ b/frontend/DESIGN.md @@ -1,5 +1,10 @@ # Jyotisha Web Design System +## 用量页分段计时(2026-10-05,在途) + +`/admin/usage` 在原表格上多一列「分段」。格子里四行:分类耗时、第 0 步、第 1 步、第 1 步开始到第一个正文字符。每步带推理、输出、输入、缓存输入 token。供应商没返回的数显示「—」。不新开页面,不改表,数字来自 `usage_ledger.metadata` 里已有的 JSON。没有新的等待动画或输入框。 + + This file adapts the full visual analysis in `CLAUDE_DESIGN.md` to the shipped Jyotisha application. `CLAUDE_DESIGN.md` remains the upstream reference; this file is the implementation contract. ## 校正按盘型交付(2026-09-30,在途,BUG-1115~1117) diff --git a/frontend/src/app/api/admin/usage/route.ts b/frontend/src/app/api/admin/usage/route.ts index 8cf81c93..ad7d5d91 100644 --- a/frontend/src/app/api/admin/usage/route.ts +++ b/frontend/src/app/api/admin/usage/route.ts @@ -2,6 +2,68 @@ import { NextResponse } from "next/server"; import { requirePermission } from "@/lib/admin/auth"; import { pageOffset, queryAdminRows } from "@/lib/admin/database"; import { adminErrorResponse, invalidQueryResponse, parseListQuery } from "@/lib/admin/http"; -export const runtime="nodejs"; -type Row={id:string;user_id:string;email:string|null;request_id:string;feature_key:string;source:string;requested_model_id:string|null;actual_model_id:string|null;model_config_version:number|null;input_tokens:number;output_tokens:number;cost_microusd:string;duration_ms:number|null;created_at:Date;cache_read_tokens:string|null;total_count:string}; -export async function GET(request:Request){try{await requirePermission("billing.orders.read");const p=parseListQuery(request);if(!p.success)return invalidQueryResponse(p.error.flatten());const q=p.data.q?`%${p.data.q}%`:null;const rows=await queryAdminRows(`select l.id,l.user_id,u.email,l.request_id,l.feature_key,l.source,l.requested_model_id,l.actual_model_id,l.model_config_version,l.input_tokens,l.output_tokens,l.cost_microusd::text,l.duration_ms,l.created_at,l.metadata->'cache'->>'readTokens' as cache_read_tokens,count(*) over()::text total_count from public.usage_ledger l left join identity.users u on u.id=l.user_id where ($1::text is null or u.email ilike $1 or l.request_id ilike $1 or l.actual_model_id ilike $1) and ($2::text is null or l.source=$2 or l.feature_key=$2) order by l.created_at desc limit $3 offset $4`,[q,p.data.status??null,p.data.pageSize,pageOffset(p.data.page,p.data.pageSize)]);return NextResponse.json({data:rows.map(r=>({id:r.id,userId:r.user_id,email:r.email,requestId:r.request_id,featureKey:r.feature_key,source:r.source,requestedModelId:r.requested_model_id,actualModelId:r.actual_model_id,modelConfigVersion:r.model_config_version,inputTokens:r.input_tokens,outputTokens:r.output_tokens,costMicrousd:Number(r.cost_microusd),durationMs:r.duration_ms,cacheReadTokens:r.cache_read_tokens==null?null:Number(r.cache_read_tokens),createdAt:r.created_at.toISOString()})),total:Number(rows[0]?.total_count??0)});}catch(e){return adminErrorResponse(e)}} +import { usageTimingFromLedger } from "@/lib/consultation-step-timing"; + +export const runtime = "nodejs"; + +type Row = { + id: string; + user_id: string; + email: string | null; + request_id: string; + feature_key: string; + source: string; + requested_model_id: string | null; + actual_model_id: string | null; + model_config_version: number | null; + input_tokens: number; + output_tokens: number; + cost_microusd: string; + duration_ms: number | null; + created_at: Date; + cache_read_tokens: string | null; + model_steps: unknown; + classification_duration_ms: string | null; + answer_reasoning_ms: string | null; + total_count: string; +}; + +const LIST_SQL = `select l.id,l.user_id,u.email,l.request_id,l.feature_key,l.source,l.requested_model_id,l.actual_model_id,l.model_config_version,l.input_tokens,l.output_tokens,l.cost_microusd::text,l.duration_ms,l.created_at,l.metadata->'cache'->>'readTokens' as cache_read_tokens,l.metadata->'modelSteps' as model_steps,l.metadata->'classification'->>'durationMs' as classification_duration_ms,l.metadata->'answer'->>'reasoning_ms' as answer_reasoning_ms,count(*) over()::text total_count from public.usage_ledger l left join identity.users u on u.id=l.user_id where ($1::text is null or u.email ilike $1 or l.request_id ilike $1 or l.actual_model_id ilike $1) and ($2::text is null or l.source=$2 or l.feature_key=$2) order by l.created_at desc limit $3 offset $4`; + +export async function GET(request: Request) { + try { + await requirePermission("billing.orders.read"); + const parsed = parseListQuery(request); + if (!parsed.success) return invalidQueryResponse(parsed.error.flatten()); + const q = parsed.data.q ? `%${parsed.data.q}%` : null; + const rows = await queryAdminRows(LIST_SQL, [ + q, + parsed.data.status ?? null, + parsed.data.pageSize, + pageOffset(parsed.data.page, parsed.data.pageSize), + ]); + return NextResponse.json({ + data: rows.map((row) => ({ + id: row.id, + userId: row.user_id, + email: row.email, + requestId: row.request_id, + featureKey: row.feature_key, + source: row.source, + requestedModelId: row.requested_model_id, + actualModelId: row.actual_model_id, + modelConfigVersion: row.model_config_version, + inputTokens: row.input_tokens, + outputTokens: row.output_tokens, + costMicrousd: Number(row.cost_microusd), + durationMs: row.duration_ms, + cacheReadTokens: row.cache_read_tokens == null ? null : Number(row.cache_read_tokens), + timing: usageTimingFromLedger(row.model_steps, row.classification_duration_ms, row.answer_reasoning_ms), + createdAt: row.created_at.toISOString(), + })), + total: Number(rows[0]?.total_count ?? 0), + }); + } catch (error) { + return adminErrorResponse(error); + } +} diff --git a/frontend/src/app/api/consult/route.ts b/frontend/src/app/api/consult/route.ts index f1e2092a..e073a9e2 100644 --- a/frontend/src/app/api/consult/route.ts +++ b/frontend/src/app/api/consult/route.ts @@ -42,6 +42,7 @@ import { createAdminSupabaseClient } from "@/lib/supabase/admin"; import { createServerSupabaseClient } from "@/lib/supabase/server"; import { streamTextResponse } from "@/lib/stream-text-response"; import { classifyConsultationTurn, type SmalltalkUsage } from "@/lib/consultation-smalltalk"; +import { consultationUsageMetadata } from "@/lib/consultation-step-timing"; import { streamSmalltalkResponse } from "@/lib/stream-smalltalk-response"; import { streamFirstResponse } from "@/lib/stream-first-response"; import { consultationPublicActivityEvent, streamAgentResponse } from "@/lib/stream-agent-response"; @@ -662,6 +663,8 @@ export async function POST(request: Request) { const usageStartedAt = Date.now(); let classificationUsage: SmalltalkUsage = {}; let classificationOutcome: string | undefined; + let classificationDurationMs: number | null = null; + let latencyState: ReturnType | null = null; const turn = parsed.data.entrypoint === undefined ? await classifyConsultationTurn({ model: selectedModel, @@ -673,6 +676,7 @@ export async function POST(request: Request) { if (!observation.late) { classificationUsage = observation.usage ?? {}; classificationOutcome = observation.outcome; + classificationDurationMs = observation.durationMs; } const inputTokens = observation.usage?.inputTokens ?? 0; const outputTokens = observation.usage?.outputTokens ?? 0; @@ -682,6 +686,7 @@ export async function POST(request: Request) { modelVersion: String(selectedModel.configVersion), policyVersion: "consultation-smalltalk-v1", toolCalls: [], + classification: { durationMs: observation.durationMs }, contractPhases: [{ phase: observation.late ? "classification.late_usage" : `classification.${observation.outcome}`, durationMs: observation.durationMs, @@ -789,12 +794,17 @@ export async function POST(request: Request) { durationMs: Date.now() - usageStartedAt, metadata: { ...(cache ? { cache: { ...cache, hit: cache.readTokens > 0 } } : {}), - ...(classificationOutcome ? { classification: { - outcome: classificationOutcome, + ...(classificationOutcome || classificationDurationMs !== null ? { classification: { + ...(classificationOutcome ? { outcome: classificationOutcome } : {}), usageKnown: classificationUsage.inputTokens !== undefined || classificationUsage.outputTokens !== undefined, inputTokens: classificationUsage.inputTokens ?? null, outputTokens: classificationUsage.outputTokens ?? null, + durationMs: classificationDurationMs, } } : {}), + ...(latencyState?.modelStepTimings === undefined ? {} : consultationUsageMetadata({ + modelSteps: latencyState.modelStepTimings, + answerReasoningMs: latencyState.answerReasoningMs, + })), }, }; } @@ -939,6 +949,7 @@ export async function POST(request: Request) { // can produce, otherwise a run that exhausts its steps also truncates the // evidence of having done so. const state = createConsultationRuntimeState({ plannedSteps: AGENT_MAX_STEPS }); + latencyState = state; const hooks = createConsultationRuntimeHooks(state); const usages: Promise[] = []; const agentStartedAt = Date.now(); @@ -1012,6 +1023,7 @@ export async function POST(request: Request) { ], retryCount: Math.max(0, usages.length - 1), ...consultationModelStepTelemetry(state), + ...(classificationDurationMs === null ? {} : { classification: { durationMs: classificationDurationMs } }), ...(errorCode === undefined ? {} : { errorCode }), inputTokens, outputTokens, diff --git a/frontend/src/components/admin/billing-operations-resources.tsx b/frontend/src/components/admin/billing-operations-resources.tsx index 85b53959..6dd7b959 100644 --- a/frontend/src/components/admin/billing-operations-resources.tsx +++ b/frontend/src/components/admin/billing-operations-resources.tsx @@ -18,6 +18,7 @@ import { import { useState } from "react"; import { adminRequestJson, type AdminIdentity } from "@/lib/admin/providers"; +import { formatUsageTiming, type UsageTimingView } from "@/lib/consultation-step-timing"; import { formatAdminDate, ResourceTable } from "./resource-table"; const { Text } = Typography; @@ -416,6 +417,7 @@ type Usage = { costMicrousd: number; durationMs: number | null; cacheReadTokens: number | null; + timing: UsageTimingView | null; createdAt: string; }; @@ -468,6 +470,13 @@ export function UsageResource() { dataIndex: "durationMs", render: (value) => (value == null ? "—" : `${value} ms`), }, + { + title: "分段", + dataIndex: "timing", + render: (value: UsageTimingView | null) => ( + {formatUsageTiming(value)} + ), + }, { title: "时间", dataIndex: "createdAt", render: formatAdminDate }, ]; return ( diff --git a/frontend/src/lib/agent-observability.ts b/frontend/src/lib/agent-observability.ts index 1174d3ec..20af2308 100644 --- a/frontend/src/lib/agent-observability.ts +++ b/frontend/src/lib/agent-observability.ts @@ -99,6 +99,27 @@ export function settlementTelemetryOutcome( }; } +const nullableDurationSchema = durationMsSchema.nullable(); +const nullableTokenSchema = tokenCountSchema.nullable(); + +/** Step 0 decides the chart tool. Step 1 writes the answer. Tokens stay null when the provider omits them. */ +export const consultationModelStepTimingSchema = z.object({ + index: z.union([z.literal(0), z.literal(1)]), + durationMs: nullableDurationSchema, + reasoningTokens: nullableTokenSchema, + outputTokens: nullableTokenSchema, + inputTokens: nullableTokenSchema, + cachedInputTokens: nullableTokenSchema, +}).strict().readonly(); + +export const consultationAnswerTimingSchema = z.object({ + reasoning_ms: nullableDurationSchema, +}).strict().readonly(); + +export const consultationClassificationTimingSchema = z.object({ + durationMs: nullableDurationSchema, +}).strict().readonly(); + export const agentObservabilityToolCallSchema = z.object({ name: machineCodeSchema, durationMs: durationMsSchema, @@ -153,6 +174,9 @@ export const agentObservabilityEventSchema = z.object({ toolCalls: z.array(agentObservabilityToolCallSchema).max(64).optional(), contractPhases: z.array(agentObservabilityContractPhaseSchema).max(64).optional(), + classification: consultationClassificationTimingSchema.optional(), + modelSteps: z.array(consultationModelStepTimingSchema).max(2).optional(), + answer: consultationAnswerTimingSchema.optional(), retryCount: z.number().int().min(0).max(100).optional(), errorCode: machineCodeSchema.optional(), // Why the model stopped, and how many model steps the run consumed across diff --git a/frontend/src/lib/consultation-step-timing.ts b/frontend/src/lib/consultation-step-timing.ts new file mode 100644 index 00000000..eeb90203 --- /dev/null +++ b/frontend/src/lib/consultation-step-timing.ts @@ -0,0 +1,237 @@ +/** + * Per-step clocks for one ordinary consultation. + * + * Numbers and statuses only. Provider usage is copied when the provider sent + * the field; a missing field stays null and is never estimated. Answer text, + * prompts, and birth data are not fields here. + */ + +export type ConsultationModelStepTiming = { + index: 0 | 1; + durationMs: number | null; + reasoningTokens: number | null; + outputTokens: number | null; + inputTokens: number | null; + cachedInputTokens: number | null; +}; + +export type ConsultationLatencyRecord = { + steps: Array; + openIndex: number | null; + openStartedAt: number | null; + nextIndex: number; + step1StartedAt: number | null; + answerReasoningMs: number | null; +}; + +export type UsageTimingView = { + classificationDurationMs: number | null; + modelSteps: ConsultationModelStepTiming[]; + answerReasoningMs: number | null; +}; + +type UsageRecord = Record; + +export function createConsultationLatencyRecord(): ConsultationLatencyRecord { + return { + steps: [], + openIndex: null, + openStartedAt: null, + nextIndex: 0, + step1StartedAt: null, + answerReasoningMs: null, + }; +} + +function tokenOrNull(value: unknown): number | null { + if (typeof value !== "number" || !Number.isFinite(value) || value < 0) return null; + const truncated = Math.trunc(value); + return truncated > 1_000_000_000 ? null : truncated; +} + +function asRecord(value: unknown): UsageRecord | null { + return value && typeof value === "object" && !Array.isArray(value) + ? value as UsageRecord + : null; +} + +/** Read the provider's cached-input count. Absent means null, including when only a sibling field exists. */ +export function cachedInputTokensFromUsage(usage: UsageRecord): number | null { + if ("cachedInputTokens" in usage) return tokenOrNull(usage.cachedInputTokens); + const details = asRecord(usage.inputTokenDetails); + if (details && "cacheReadTokens" in details) return tokenOrNull(details.cacheReadTokens); + const promptDetails = asRecord(usage.prompt_tokens_details); + if (promptDetails && "cached_tokens" in promptDetails) return tokenOrNull(promptDetails.cached_tokens); + return null; +} + +function usageFromChunk(chunk: { payload?: unknown }): UsageRecord { + const payload = asRecord(chunk.payload); + const output = asRecord(payload?.output); + return asRecord(output?.usage) ?? asRecord(payload?.usage) ?? {}; +} + +function emptyStep(index: 0 | 1): ConsultationModelStepTiming { + return { + index, + durationMs: null, + reasoningTokens: null, + outputTokens: null, + inputTokens: null, + cachedInputTokens: null, + }; +} + +function closeOpen(record: ConsultationLatencyRecord, now: number, usage: UsageRecord): void { + const index = record.openIndex; + const startedAt = record.openStartedAt; + record.openIndex = null; + record.openStartedAt = null; + if (index === null) return; + record.nextIndex = index + 1; + if (index !== 0 && index !== 1) return; + const durationMs = startedAt === null ? null : Math.max(0, Math.trunc(now - startedAt)); + record.steps[index] = { + index, + durationMs: durationMs !== null && durationMs > 7 * 24 * 60 * 60 * 1000 ? null : durationMs, + reasoningTokens: tokenOrNull(usage.reasoningTokens), + outputTokens: tokenOrNull(usage.outputTokens), + inputTokens: tokenOrNull(usage.inputTokens), + cachedInputTokens: cachedInputTokensFromUsage(usage), + }; +} + +/** Fold one model-loop chunk into step 0 / step 1. Text payloads are ignored. */ +export function observeConsultationModelChunk( + record: ConsultationLatencyRecord, + chunk: { type?: string; payload?: unknown }, + now: number, +): void { + if (chunk.type === "step-start") { + if (record.openIndex !== null) closeOpen(record, now, {}); + const index = record.nextIndex; + record.openIndex = index; + record.openStartedAt = now; + if (index === 1) record.step1StartedAt = now; + return; + } + if (chunk.type === "step-finish") { + if (record.openIndex === null) { + record.openIndex = record.nextIndex; + record.openStartedAt = null; + } + closeOpen(record, now, usageFromChunk(chunk)); + } +} + +/** + * First answer-body character of step 1. Step 0 narration does not count, + * and a later continuation does not move the clock backward or forward. + */ +export function noteAnswerBodyCharacter(record: ConsultationLatencyRecord, now: number): void { + if (record.answerReasoningMs !== null) return; + if (record.openIndex !== 1 || record.step1StartedAt === null) return; + record.answerReasoningMs = Math.max(0, Math.trunc(now - record.step1StartedAt)); +} + +export function consultationModelSteps(record: ConsultationLatencyRecord): ConsultationModelStepTiming[] { + return [emptyStep(0), emptyStep(1)].map((blank, index) => record.steps[index] ?? blank); +} + +export function consultationUsageMetadata(input: { + modelSteps: readonly ConsultationModelStepTiming[] | null | undefined; + answerReasoningMs: number | null | undefined; +}): { modelSteps: ConsultationModelStepTiming[]; answer: { reasoning_ms: number | null } } { + const view = consultationUsageTiming({ + classificationDurationMs: null, + modelSteps: input.modelSteps, + answerReasoningMs: input.answerReasoningMs, + }); + return { + modelSteps: view.modelSteps, + answer: { reasoning_ms: view.answerReasoningMs }, + }; +} + +export function consultationUsageTiming(input: { + classificationDurationMs: number | null; + modelSteps: readonly ConsultationModelStepTiming[] | null | undefined; + answerReasoningMs: number | null | undefined; +}): UsageTimingView { + const steps = input.modelSteps ?? []; + return { + classificationDurationMs: input.classificationDurationMs, + modelSteps: [0, 1].map((index) => { + const found = steps.find((step) => step.index === index); + return found ?? emptyStep(index as 0 | 1); + }), + answerReasoningMs: input.answerReasoningMs ?? null, + }; +} + +function integerOrNull(value: unknown): number | null { + if (typeof value === "number" && Number.isFinite(value)) return tokenOrNull(value); + if (typeof value === "string" && value.trim() !== "") { + const parsed = Number(value); + return Number.isFinite(parsed) ? tokenOrNull(parsed) : null; + } + return null; +} + +/** Keep only the numeric timing fields stored on usage_ledger.metadata. */ +function parsedJson(value: unknown): unknown { + if (typeof value !== "string") return value; + try { + return JSON.parse(value) as unknown; + } catch { + return null; + } +} + +export function usageTimingFromLedger( + modelSteps: unknown, + classificationDurationMs: unknown, + answerReasoningMs: unknown, +): UsageTimingView { + const decoded = parsedJson(modelSteps); + const steps = Array.isArray(decoded) ? decoded : []; + const parsed = steps.map((item) => { + const record = asRecord(item); + if (!record || (record.index !== 0 && record.index !== 1)) return null; + return { + index: record.index as 0 | 1, + durationMs: integerOrNull(record.durationMs), + reasoningTokens: integerOrNull(record.reasoningTokens), + outputTokens: integerOrNull(record.outputTokens), + inputTokens: integerOrNull(record.inputTokens), + cachedInputTokens: integerOrNull(record.cachedInputTokens), + }; + }).filter((step): step is ConsultationModelStepTiming => step !== null); + return consultationUsageTiming({ + classificationDurationMs: integerOrNull(classificationDurationMs), + modelSteps: parsed, + answerReasoningMs: integerOrNull(answerReasoningMs), + }); +} + +function ms(value: number | null): string { + return value === null ? "—" : `${value} ms`; +} + +function tok(value: number | null): string { + return value === null ? "—" : value.toLocaleString("zh-CN"); +} + +/** One admin-table cell. Numbers only, so an old row with no timing still reads as dashes. */ +export function formatUsageTiming(view: UsageTimingView | null | undefined): string { + if (!view) return "—"; + const steps = view.modelSteps.length > 0 ? view.modelSteps : [emptyStep(0), emptyStep(1)]; + const lines = [ + `分类 ${ms(view.classificationDurationMs)}`, + ...steps.map((step) => ( + `第 ${step.index} 步 ${ms(step.durationMs)} · 推理 ${tok(step.reasoningTokens)} · 出 ${tok(step.outputTokens)} · 入 ${tok(step.inputTokens)} · 缓存入 ${tok(step.cachedInputTokens)}` + )), + `写到正文 ${ms(view.answerReasoningMs)}`, + ]; + return lines.join("\n"); +} diff --git a/frontend/src/lib/stream-agent-response.ts b/frontend/src/lib/stream-agent-response.ts index 60102364..0713cc3d 100644 --- a/frontend/src/lib/stream-agent-response.ts +++ b/frontend/src/lib/stream-agent-response.ts @@ -13,6 +13,13 @@ import { type ConsultationAgentPublicEvent, } from "./consultation-agent-events.ts"; import { toAgentModelFinishReason, type AgentModelFinishReason } from "./agent-observability.ts"; +import { + createConsultationLatencyRecord, + noteAnswerBodyCharacter, + observeConsultationModelChunk, + consultationModelSteps, + type ConsultationLatencyRecord, +} from "./consultation-step-timing.ts"; import { createVisibleTextTransformer } from "./stream-text-response.ts"; import { consultationWriteLabel } from "./consultation-activity-labels.ts"; import { logTruncatedReasoning } from "./consultation-budget.ts"; @@ -482,6 +489,12 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { // The calculation result the model saw, kept so a length continuation // writes with the same evidence (BUG-1053). Never sent to the client. let calculationEvidence: unknown; + const latencyRecord: ConsultationLatencyRecord = createConsultationLatencyRecord(); + const publishLatency = () => { + if (!options.state) return; + options.state.modelStepTimings = consultationModelSteps(latencyRecord); + options.state.answerReasoningMs = latencyRecord.answerReasoningMs; + }; // The one-shot evidence lookup's result, if the model used it; carried into // a length continuation with the card (TASK-consult-evidence-card-20260927). let lookupEvidence: unknown; @@ -751,6 +764,16 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { const stepCountBeforeAttempt = options.state.modelStepCount; try { for await (const chunk of readChunks(stream)) { + if (chunk.type === "step-start" || chunk.type === "step-finish") { + observeConsultationModelChunk(latencyRecord, chunk, Date.now()); + } + if ( + chunk.type === "text-delta" + && typeof chunk.payload?.text === "string" + && /\S/.test(chunk.payload.text) + ) { + noteAnswerBodyCharacter(latencyRecord, Date.now()); + } if (isSkillBindingTripwire(chunk)) { recordSkillBindingAbort(options); throw new Error(SKILL_BINDING_FAILED); @@ -863,6 +886,8 @@ export function streamAgentResponse(options: StreamAgentResponseOptions) { lastAttempt = outcome; if (outcome.contributed || outcome.answerPhase) answerTail = outcome; throw error; + } finally { + publishLatency(); } } diff --git a/frontend/src/mastra/consultation-tools.ts b/frontend/src/mastra/consultation-tools.ts index 6d0112d5..fd526756 100644 --- a/frontend/src/mastra/consultation-tools.ts +++ b/frontend/src/mastra/consultation-tools.ts @@ -288,6 +288,17 @@ export type ConsultationRuntimeState = { // strict and would reject them, so they never enter it. modelStepCount: number; modelFinishReason?: AgentModelFinishReason; + /** Step 0 and step 1 clocks. Observability only; absent until the stream records them. */ + modelStepTimings?: ReadonlyArray<{ + index: 0 | 1; + durationMs: number | null; + reasoningTokens: number | null; + outputTokens: number | null; + inputTokens: number | null; + cachedInputTokens: number | null; + }>; + /** Milliseconds from step 1 start to the first answer-body character. Null when that character never arrived. */ + answerReasoningMs?: number | null; // How the stream that wrote the answer ended (compose, continuation or answer // retry; the main loop when a path has no compose; the last stream when no // answer was written). "missing" means the stream closed without a finish @@ -351,6 +362,8 @@ export function consultationModelStepTelemetry(state: ConsultationRuntimeState) ...(state.composeFinishReason === undefined ? {} : { composeFinishReason: state.composeFinishReason }), ...(state.composeAborted === undefined ? {} : { composeAborted: state.composeAborted }), ...(state.answerVisibleChars === undefined ? {} : { answerVisibleChars: state.answerVisibleChars }), + ...(state.modelStepTimings === undefined ? {} : { modelSteps: state.modelStepTimings }), + ...(state.answerReasoningMs === undefined ? {} : { answer: { reasoning_ms: state.answerReasoningMs } }), }; } diff --git a/frontend/tests/consult-latency-timing-20261005.test.ts b/frontend/tests/consult-latency-timing-20261005.test.ts new file mode 100644 index 00000000..25804192 --- /dev/null +++ b/frontend/tests/consult-latency-timing-20261005.test.ts @@ -0,0 +1,199 @@ +import assert from "node:assert/strict"; +import { readFileSync } from "node:fs"; +import test from "node:test"; + +import { agentObservabilityEventSchema } from "../src/lib/agent-observability.ts"; +import { streamAgentResponse } from "../src/lib/stream-agent-response.ts"; +import { + consultationModelStepTelemetry, + consultationStepBudgetReceipt, + createConsultationRuntimeState, +} from "../src/mastra/consultation-tools.ts"; +import { + consultationModelSteps, + consultationUsageMetadata, + createConsultationLatencyRecord, + formatUsageTiming, + noteAnswerBodyCharacter, + observeConsultationModelChunk, + usageTimingFromLedger, +} from "../src/lib/consultation-step-timing.ts"; + +const SECRET = "这段正文不该出现在计时里"; + +test("simulated step results fill the timing fields and leave the answer text out", () => { + const record = createConsultationLatencyRecord(); + observeConsultationModelChunk(record, { type: "step-start" }, 1_000); + observeConsultationModelChunk(record, { + type: "step-finish", + payload: { + output: { + text: SECRET, + usage: { inputTokens: 11, outputTokens: 4, reasoningTokens: 9, cachedInputTokens: 2 }, + }, + }, + }, 1_600); + observeConsultationModelChunk(record, { type: "step-start" }, 5_000); + noteAnswerBodyCharacter(record, 5_450); + observeConsultationModelChunk(record, { + type: "step-finish", + payload: { + output: { + text: SECRET, + usage: { inputTokens: 20, outputTokens: 30, reasoningTokens: 100 }, + }, + }, + }, 9_000); + + const modelSteps = consultationModelSteps(record); + const usageMetadata = { + classification: { + outcome: "consult", + usageKnown: true, + inputTokens: 3, + outputTokens: 1, + durationMs: 180, + }, + ...consultationUsageMetadata({ + modelSteps, + answerReasoningMs: record.answerReasoningMs, + }), + }; + const event = agentObservabilityEventSchema.parse({ + requestId: "req-timing", + classification: { durationMs: 180 }, + modelSteps, + answer: { reasoning_ms: record.answerReasoningMs }, + toolCalls: [{ name: "run-jyotish-consultation", durationMs: 700, status: "completed" }], + }); + + assert.equal(event.classification?.durationMs, 180); + assert.deepEqual(event.modelSteps?.[0], { + index: 0, + durationMs: 600, + reasoningTokens: 9, + outputTokens: 4, + inputTokens: 11, + cachedInputTokens: 2, + }); + assert.deepEqual(event.modelSteps?.[1], { + index: 1, + durationMs: 4_000, + reasoningTokens: 100, + outputTokens: 30, + inputTokens: 20, + cachedInputTokens: null, + }); + assert.equal(event.answer?.reasoning_ms, 450); + assert.equal(event.toolCalls?.[0]?.durationMs, 700); + assert.equal(usageMetadata.classification.durationMs, 180); + assert.equal(usageMetadata.answer.reasoning_ms, 450); + assert.equal(JSON.stringify({ event, usageMetadata }).includes(SECRET), false); + + const ledger = usageTimingFromLedger( + JSON.stringify(usageMetadata.modelSteps), + "180", + "450", + ); + assert.equal(ledger.modelSteps[1]?.reasoningTokens, 100); + const shown = formatUsageTiming(ledger); + assert.match(shown, /分类 180 ms/); + assert.match(shown, /第 1 步 4000 ms/); + assert.match(shown, /写到正文 450 ms/); + assert.equal(shown.includes(SECRET), false); +}); + +test("a missing provider usage field stays null and step 0 text is not the answer clock", () => { + const record = createConsultationLatencyRecord(); + observeConsultationModelChunk(record, { type: "step-finish" }, 2_000); + noteAnswerBodyCharacter(record, 2_100); + observeConsultationModelChunk(record, { + type: "step-finish", + payload: { output: { usage: { outputTokens: 6 } } }, + }, 3_000); + const steps = consultationModelSteps(record); + assert.equal(steps[0]?.durationMs, null); + assert.equal(steps[0]?.inputTokens, null); + assert.equal(steps[1]?.durationMs, null); + assert.equal(steps[1]?.outputTokens, 6); + assert.equal(steps[1]?.reasoningTokens, null); + assert.equal(steps[1]?.cachedInputTokens, null); + assert.equal(record.answerReasoningMs, null); +}); + +test("the answer stream stores step timings on the run and not in the public receipt", async () => { + const state = createConsultationRuntimeState(); + state.consultationToolCallCount = 1; + state.consultationToolSuccessCount = 1; + state.consultationToolCompleted = true; + state.workflowReceipt = { route: "career", status: "ready", preciseTiming: "blocked", missingLayers: [] }; + const sentence = "事业方向的判断如下。"; + async function* chunks() { + yield { type: "step-start" }; + yield { + type: "step-finish", + payload: { + stepResult: { reason: "tool-calls" }, + output: { usage: { inputTokens: 8, outputTokens: 1, reasoningTokens: 4, cachedInputTokens: 0 }, text: sentence }, + }, + }; + yield { type: "step-start" }; + yield { type: "text-delta", payload: { text: sentence } }; + yield { + type: "step-finish", + payload: { + stepResult: { reason: "stop" }, + output: { usage: { reasoningTokens: 12 }, text: sentence }, + }, + }; + yield { type: "finish", payload: { stepResult: { reason: "stop" }, output: { usage: {}, steps: [{}, {}] } } }; + } + const response = streamAgentResponse({ + runId: "run", + requestId: "req", + state, + stream: chunks(), + requireTool: true, + toolStatus: () => "ready", + receipt: () => ({ + runId: "run", + runtime: "mastra-agentic" as const, + skill: { + name: "jyotish-vedic-astrology" as const, + loaded: state.jyotishSkillBound, + referenceReads: state.skillReferenceReadCount, + methodologySections: state.methodologySectionCount, + }, + steps: state.steps, + stepBudget: consultationStepBudgetReceipt(state), + workflow: state.workflowReceipt!, + }), + }); + const body = await response.text(); + const telemetry = consultationModelStepTelemetry(state); + const event = agentObservabilityEventSchema.parse({ requestId: "req-stream", ...telemetry }); + assert.equal(event.modelSteps?.[0]?.inputTokens, 8); + assert.equal(event.modelSteps?.[0]?.cachedInputTokens, 0); + assert.equal(event.modelSteps?.[0]?.reasoningTokens, 4); + assert.equal(event.modelSteps?.[1]?.reasoningTokens, 12); + assert.equal(event.modelSteps?.[1]?.inputTokens, null); + assert.equal(typeof event.answer?.reasoning_ms, "number"); + assert.equal(JSON.stringify(event).includes(sentence), false); + assert.equal(body.includes("modelSteps"), false); + assert.equal(body.includes("reasoning_ms"), false); +}); + +test("the consult route and the usage list keep the same timing fields", () => { + const route = readFileSync(new URL("../src/app/api/consult/route.ts", import.meta.url), "utf8"); + const list = readFileSync(new URL("../src/app/api/admin/usage/route.ts", import.meta.url), "utf8"); + const page = readFileSync(new URL("../src/components/admin/billing-operations-resources.tsx", import.meta.url), "utf8"); + assert.match(route, /classification: \{ durationMs: observation\.durationMs \}/); + assert.match(route, /durationMs: classificationDurationMs/); + assert.match(route, /consultationUsageMetadata/); + assert.match(list, /metadata->'cache'->>'readTokens'/); + assert.match(list, /metadata->'modelSteps'/); + assert.match(list, /metadata->'classification'->>'durationMs'/); + assert.match(list, /metadata->'answer'->>'reasoning_ms'/); + assert.match(page, /分段/); + assert.doesNotMatch(list, /alter table|create table/i); +}); diff --git a/scripts/jyotish_api_server.py b/scripts/jyotish_api_server.py index d7aace7d..67005a8a 100644 --- a/scripts/jyotish_api_server.py +++ b/scripts/jyotish_api_server.py @@ -2291,13 +2291,21 @@ def execute_consultation_workflow( elif step == 'run_rectification_gate': chart_planets = chart.get('planets') if isinstance(chart, dict) else {} chart_ascendant = chart.get('ascendant') if isinstance(chart, dict) else {} - rectification = handler._compute_rectification_gate({ + gate_body = { **birth_payload, 'planets': chart_planets if isinstance(chart_planets, dict) else {}, 'ascendant': chart_ascendant if isinstance(chart_ascendant, dict) else {}, 'declared_accuracy': body.get('declared_accuracy', body.get('accuracy', 'minute')), 'time_source': body.get('time_source', 'family_clear'), - }) + } + # Consult already set this flag for the foreground gateway. The gate + # must see it too, or it starts a second official snapshot that cannot + # finish inside the request. Rectification HTTP calls do not set it. + if defer_optional_external_evidence: + from scripts.vedastro_consultation_snapshot import snapshot_identity_from_request + gate_body['defer_optional_external_evidence'] = True + gate_body['vedastro_snapshot_identity'] = snapshot_identity_from_request(body) + rectification = handler._compute_rectification_gate(gate_body) executed_steps.append('run_rectification_gate') elif step == 'compute_chart': if not computed_chart: @@ -8471,7 +8479,11 @@ class JyotishAPIHandler(BaseHTTPRequestHandler, VedastroEvidenceMixin, SynastryM next_action = '保留原始出生记录来源;重要预测仍建议用 Dasha/Transit/案例验证交叉确认。' try: - vedastro_gateway = self._compute_vedastro_gateway_run(body) + if body.get('defer_optional_external_evidence'): + from scripts.vedastro_consultation_snapshot import consult_rectification_vedastro_gateway + vedastro_gateway = consult_rectification_vedastro_gateway(body) + else: + vedastro_gateway = self._compute_vedastro_gateway_run(body) except Exception as exc: # Keep rectification available when the external observation is down. vedastro_gateway = { 'scope': 'vedastro_gateway_run', diff --git a/scripts/vedastro_consultation_snapshot.py b/scripts/vedastro_consultation_snapshot.py new file mode 100644 index 00000000..8fd2cf4a --- /dev/null +++ b/scripts/vedastro_consultation_snapshot.py @@ -0,0 +1,230 @@ +"""Consult-only VedAstro snapshot for the rectification gate. + +The foreground gateway (BUG-301) still runs. This module is the other call, +inside the consultation rectification gate, which used to start the same +4-second official snapshot runner and then drop the result. + +When the request already asked to defer optional external evidence: + +- reuse a BUG-727 snapshot for the same birth, ayanamsa, node, and UTC date; +- otherwise return a negative-cache hit from earlier today; +- otherwise do not start the runner. Record + ``official_closure_reason=deferred_in_consultation`` until the end of + that UTC day. + +The rectification HTTP route does not set the flag, so it still calls the +real gateway. Negative entries live in their own directory and are never +served through the 7-day stale window. +""" +from __future__ import annotations + +import contextlib +import json +import os +import re +import tempfile +from datetime import datetime +from pathlib import Path +from typing import Any + +try: + from scripts.vedastro_snapshot_cache import ( + annotate_stale_gateway, + lookup_snapshot, + official_snapshot_reference_date, + snapshot_cache_key, + ) +except ModuleNotFoundError: # pragma: no cover - script execution path + from vedastro_snapshot_cache import ( + annotate_stale_gateway, + lookup_snapshot, + official_snapshot_reference_date, + snapshot_cache_key, + ) + +NEGATIVE_CACHE_SCHEMA = "vedastro_snapshot_negative_cache.v1" +DEFERRED_IN_CONSULTATION = "deferred_in_consultation" +_SHA256_NAME = re.compile(r"^[0-9a-f]{64}\.json$") +_IDENTITY_FIELDS = ( + "year", + "month", + "day", + "hour", + "minute", + "second", + "lat", + "lon", + "tz", + "ayanamsa", + "ayanamsa_name", + "ayanamsa_policy", + "node_mode", + "nodeMode", + "reference_date", + "today", + "transit_date", + "current_date", + "entrypoint", + "consult_entrypoint", +) +_IDENTITY_KEYS = { + "name", + "email", + "user_id", + "userid", + "user_email", + "session_id", + "sessionid", + "full_name", + "display_name", +} + + +def snapshot_identity_from_request(body: dict[str, Any] | None) -> dict[str, Any]: + """Fields the BUG-727 key reads, copied without normalizing numbers. + + The consultation birth payload turns hour ``5`` into ``5.0``. Those are + different cache keys, so the gate must look up the request's own values. + """ + payload = body if isinstance(body, dict) else {} + return {key: payload[key] for key in _IDENTITY_FIELDS if key in payload} + + +def negative_cache_dir() -> Path: + raw = str(os.environ.get("JYOTISH_VEDASTRO_SNAPSHOT_NEGATIVE_CACHE_DIR") or "").strip() + path = Path(raw) if raw else Path(__file__).resolve().parents[1] / "scratch" / "local" / "vedastro_snapshot_negative_cache" + path.mkdir(parents=True, exist_ok=True) + return path + + +def _negative_path(cache_key: str) -> Path: + if not re.fullmatch(r"[0-9a-f]{64}", cache_key): + raise ValueError("vedastro negative cache key must be sha256 hex") + return negative_cache_dir() / f"{cache_key}.json" + + +def _strip_identity(value: Any) -> Any: + if isinstance(value, dict): + return { + key: _strip_identity(item) + for key, item in value.items() + if str(key).strip().lower() not in _IDENTITY_KEYS + } + if isinstance(value, list): + return [_strip_identity(item) for item in value] + return value + + +def _cache_body(body: dict[str, Any]) -> dict[str, Any]: + identity = body.get("vedastro_snapshot_identity") + if isinstance(identity, dict): + return identity + return body + + +def _deferred_packet() -> dict[str, Any]: + return { + "scope": "vedastro_gateway_run", + "status": "official_blocked", + "official_closure_state": "official_blocked", + "official_closure_reason": DEFERRED_IN_CONSULTATION, + } + + +def lookup_negative_snapshot( + body: dict[str, Any] | None, + *, + today: str | None = None, +) -> dict[str, Any] | None: + """Same-day negative hit only. Tomorrow is a different key, so it misses.""" + payload = body if isinstance(body, dict) else {} + served = (today or official_snapshot_reference_date(payload))[:10] + path = _negative_path(snapshot_cache_key(payload, reference_date=served)) + if not _SHA256_NAME.match(path.name) or not path.is_file(): + return None + try: + record = json.loads(path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + return None + if not isinstance(record, dict) or record.get("schema") != NEGATIVE_CACHE_SCHEMA: + return None + if str(record.get("reference_date") or "")[:10] != served: + return None + gateway = record.get("gateway") + return dict(gateway) if isinstance(gateway, dict) else None + + +def store_negative_snapshot( + body: dict[str, Any] | None, + gateway: dict[str, Any], + *, + today: str | None = None, +) -> dict[str, Any] | None: + if not isinstance(gateway, dict): + return None + state = gateway.get("official_closure_state") or gateway.get("status") + if state == "official_verified": + return None + payload = body if isinstance(body, dict) else {} + served = (today or official_snapshot_reference_date(payload))[:10] + try: + datetime.strptime(served, "%Y-%m-%d") + except ValueError: + return None + cache_key = snapshot_cache_key(payload, reference_date=served) + record = { + "schema": NEGATIVE_CACHE_SCHEMA, + "cache_key": cache_key, + "reference_date": served, + "stored_at": datetime.utcnow().strftime("%Y-%m-%dT%H:%M:%SZ"), + "gateway": _strip_identity(gateway), + } + path = _negative_path(cache_key) + path.parent.mkdir(parents=True, exist_ok=True) + encoded = json.dumps(record, ensure_ascii=False, sort_keys=True) + fd, tmp_name = tempfile.mkstemp(prefix=f"{cache_key}.", suffix=".tmp", dir=str(path.parent)) + try: + with os.fdopen(fd, "w", encoding="utf-8") as handle: + handle.write(encoded) + os.replace(tmp_name, path) + except Exception: + with contextlib.suppress(OSError): + os.unlink(tmp_name) + raise + return record + + +def _verified_gateway(hit: dict[str, Any]) -> dict[str, Any] | None: + record = hit.get("record") if isinstance(hit, dict) else None + if not isinstance(record, dict): + return None + gateway = record.get("gateway") + if not isinstance(gateway, dict): + return None + state = gateway.get("official_closure_state") or gateway.get("status") + if state != "official_verified": + return None + if hit.get("freshness") == "stale": + return annotate_stale_gateway( + gateway, + reference_date=str(record.get("reference_date") or ""), + served_on_utc_date=str(hit.get("served_on_utc_date") or ""), + ) + return dict(gateway) + + +def consult_rectification_vedastro_gateway(body: dict[str, Any] | None) -> dict[str, Any]: + """Official snapshot for the consultation rectification gate. Never starts the runner.""" + payload = body if isinstance(body, dict) else {} + cache_body = _cache_body(payload) + today = official_snapshot_reference_date(cache_body) + reused = _verified_gateway(lookup_snapshot(cache_body, today=today) or {}) + if reused is not None: + return reused + cached = lookup_negative_snapshot(cache_body, today=today) + if cached is not None: + return cached + packet = _deferred_packet() + with contextlib.suppress(OSError): + store_negative_snapshot(cache_body, packet, today=today) + return packet diff --git a/tests/test_vedastro_consultation_snapshot.py b/tests/test_vedastro_consultation_snapshot.py new file mode 100644 index 00000000..76c6a7f2 --- /dev/null +++ b/tests/test_vedastro_consultation_snapshot.py @@ -0,0 +1,199 @@ +"""Consult rectification gate reuses or defers the official VedAstro snapshot.""" +from __future__ import annotations + +import json + +import pytest + +from scripts.jyotish_api_server import JyotishAPIHandler +from scripts.vedastro_consultation_snapshot import ( + consult_rectification_vedastro_gateway, + negative_cache_dir, + snapshot_identity_from_request, + store_negative_snapshot, +) +from scripts.vedastro_snapshot_cache import snapshot_cache_dir, store_snapshot + + +def _handler() -> JyotishAPIHandler: + return object.__new__(JyotishAPIHandler) + + +def _planets() -> dict: + return { + "Sun": {"lon": 80.0}, + "Moon": {"lon": 123.0}, + "Mars": {"lon": 210.0}, + "Mercury": {"lon": 75.0}, + "Jupiter": {"lon": 15.0}, + "Venus": {"lon": 102.0}, + "Saturn": {"lon": 330.0}, + "Rahu": {"lon": 5.0}, + "Ketu": {"lon": 185.0}, + } + + +def _body(**extra) -> dict: + payload = { + "year": 1990, + "month": 1, + "day": 1, + "hour": 5, + "minute": 0, + "second": 0, + "lat": 19.0, + "lon": 72.8, + "tz": 5.5, + "ayanamsa": "raman", + "node_mode": "mean", + "reference_date": "2026-10-05", + "today": "2026-10-05", + "current_date": "2026-10-05", + "planets": _planets(), + "ascendant": {"lon": 15.0}, + "declared_accuracy": "minute", + "time_source": "family_clear", + } + payload.update(extra) + return payload + + +def _verified(request_id: str = "same-day") -> dict: + return { + "scope": "vedastro_gateway_run", + "status": "official_verified", + "official_closure_state": "official_verified", + "official_closure_reason": "official_raw_response_present", + "official_raw_response": {"request_id": request_id, "natal": {"sun": "Leo"}}, + } + + +@pytest.fixture +def caches(tmp_path, monkeypatch): + monkeypatch.setenv("JYOTISH_VEDASTRO_SNAPSHOT_CACHE_DIR", str(tmp_path / "positive")) + monkeypatch.setenv("JYOTISH_VEDASTRO_SNAPSHOT_NEGATIVE_CACHE_DIR", str(tmp_path / "negative")) + return tmp_path + + +def _forbid_runner(monkeypatch) -> list[str]: + calls: list[str] = [] + + def runner(*_args, **_kwargs): + calls.append("runner") + raise AssertionError("snapshot runner must not start on the consult gate") + + def report(*_args, **_kwargs): + calls.append("report") + raise AssertionError("build_report must not start on the consult gate") + + monkeypatch.setattr( + "scripts.vedastro_service_adapter._try_official_capability_runner_snapshot_bundle", + runner, + ) + monkeypatch.setattr("scripts.vedastro_user_entrypoint.build_report", report) + return calls + + +def test_consult_gate_does_not_call_the_snapshot_runner(caches, monkeypatch) -> None: + calls = _forbid_runner(monkeypatch) + handler = _handler() + + def live_gateway(_body): + calls.append("gateway") + raise AssertionError("consult gate must not call the live gateway") + + handler._compute_vedastro_gateway_run = live_gateway # type: ignore[method-assign] + result = handler._compute_rectification_gate({ + **_body(), + "defer_optional_external_evidence": True, + }) + gateway = result["vedastro_gateway"] + assert calls == [] + assert gateway["official_closure_state"] == "official_blocked" + assert gateway["official_closure_reason"] == "deferred_in_consultation" + assert gateway["status"] != "official_verified" + assert "official_raw_response" not in gateway + assert result["endpoint"] == "rectification_gate" + + +def test_negative_cache_hits_the_same_utc_day_and_misses_the_next(caches, monkeypatch) -> None: + calls = _forbid_runner(monkeypatch) + stored = store_negative_snapshot(_body(), { + "scope": "vedastro_gateway_run", + "status": "official_blocked", + "official_closure_state": "official_blocked", + "official_closure_reason": "foreground_optional_evidence_timeout", + "email": "hidden@example.com", + }, today="2026-10-05") + assert stored is not None + assert "email" not in json.dumps(stored["gateway"]) + + same_day = consult_rectification_vedastro_gateway({ + **_body(), + "defer_optional_external_evidence": True, + }) + assert same_day["official_closure_reason"] == "foreground_optional_evidence_timeout" + before = sorted(path.read_bytes() for path in negative_cache_dir().glob("*.json")) + again = consult_rectification_vedastro_gateway({ + **_body(), + "defer_optional_external_evidence": True, + }) + after = sorted(path.read_bytes() for path in negative_cache_dir().glob("*.json")) + assert again["official_closure_reason"] == "foreground_optional_evidence_timeout" + assert before == after + + next_day = consult_rectification_vedastro_gateway({ + **_body(reference_date="2026-10-06", today="2026-10-06", current_date="2026-10-06"), + "defer_optional_external_evidence": True, + }) + assert next_day["official_closure_reason"] == "deferred_in_consultation" + assert calls == [] + assert negative_cache_dir().resolve() != snapshot_cache_dir().resolve() + + +def test_same_day_verified_snapshot_is_reused_ahead_of_the_negative_cache(caches, monkeypatch) -> None: + calls = _forbid_runner(monkeypatch) + request = _body() + store_snapshot(request, _verified("cached-official")) + store_negative_snapshot(request, { + "status": "official_blocked", + "official_closure_state": "official_blocked", + "official_closure_reason": "foreground_optional_evidence_timeout", + }, today="2026-10-05") + normalized = { + **request, + "hour": 5.0, + "minute": 0.0, + "defer_optional_external_evidence": True, + "vedastro_snapshot_identity": snapshot_identity_from_request(request), + } + reused = consult_rectification_vedastro_gateway(normalized) + assert reused["official_closure_state"] == "official_verified" + assert reused["official_raw_response"]["request_id"] == "cached-official" + missed = consult_rectification_vedastro_gateway({ + **request, + "hour": 5.0, + "defer_optional_external_evidence": True, + }) + assert missed["official_closure_reason"] == "deferred_in_consultation" + assert calls == [] + + +def test_rectification_surface_without_the_flag_still_uses_the_live_gateway(caches) -> None: + handler = _handler() + seen: list[dict] = [] + + def live_gateway(body): + seen.append(dict(body)) + return _verified("rectification-surface") + + handler._compute_vedastro_gateway_run = live_gateway # type: ignore[method-assign] + result = handler._compute_rectification_gate(_body()) + assert len(seen) == 1 + assert "defer_optional_external_evidence" not in seen[0] + gateway = result["vedastro_gateway"] + assert gateway["official_closure_state"] == "official_verified" + assert gateway["official_closure_reason"] == "official_raw_response_present" + assert gateway["official_raw_response"]["request_id"] == "rectification-surface" + assert result["endpoint"] == "rectification_gate" + assert result["success"] is True