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
This commit is contained in:
jesse-ux
2026-10-05 02:33:36 +08:00
parent 4c64733f91
commit 137ed6f408
16 changed files with 1188 additions and 9 deletions
+7
View File
@@ -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)。
+30
View File
@@ -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 条)。
- 复发自:无
- 修复版本:待发布
@@ -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。首页源文件没改。没有为了构建去改符号链接,也没有在工作树里重装依赖。
+1 -1
View File
@@ -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) |
+5
View File
@@ -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)
+65 -3
View File
@@ -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<Row>(`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<Row>(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);
}
}
+14 -2
View File
@@ -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<typeof createConsultationRuntimeState> | 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<Usage>[] = [];
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,
@@ -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) => (
<Text style={{ whiteSpace: "pre-line" }}>{formatUsageTiming(value)}</Text>
),
},
{ title: "时间", dataIndex: "createdAt", render: formatAdminDate },
];
return (
+24
View File
@@ -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
@@ -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<ConsultationModelStepTiming | undefined>;
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<string, unknown>;
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");
}
+25
View File
@@ -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();
}
}
+13
View File
@@ -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 } }),
};
}
@@ -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);
});
+15 -3
View File
@@ -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',
+230
View File
@@ -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
@@ -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