Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
137ed6f408 | ||
|
|
4c64733f91 | ||
|
|
853772bcea | ||
|
|
1252de3e63 |
@@ -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)。
|
||||
|
||||
@@ -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。首页源文件没改。没有为了构建去改符号链接,也没有在工作树里重装依赖。
|
||||
@@ -68,3 +68,11 @@
|
||||
| 推送 / 部署 / Gitea 页面 | 没做 |
|
||||
|
||||
快速门跑的时候改写了 `references/oracle/artifacts/pending_packets/` 下 5 个模板文件。那不是本单的改动,收尾时还原,不提交。
|
||||
|
||||
## Claude 验收与补修(2026-10-04)
|
||||
|
||||
- 实测:快速门 18 行(改前约 14,000 行);`CI=true npm test` 346 行,24 条已知失败的名字和报错都列出,退出码 1;本机 Node 22 `npm test` 输出 TAP,失败名单与开工基线逐条相同;快速门子步骤失败(假命令 exit 3)时打出 ✗、末 200 行、`gate-logs/01-fake-fail.log` 和退出码;`test_run_quality_gate_output.py` 10 条通过;workflow 合同测试唯二的失败(live staging sync 两条)本来就在基线 24 条里。
|
||||
- workflow 拆成 7 步,原 7 条命令一条不少、参数不变,没有吞退出码。
|
||||
- **补修(违反红线 4)**:`run-tests.mjs` 把 `tests/*.test.ts` 原样交给 `node --test`。Node 22 会自己展开,Node 20 不会,报「Could not find tests/*.test.ts」,一条测试都不跑;本机默认的 Node 就是 v20.19.2。改为脚本自己读 `tests/` 目录,`*.test.ts` 排序后接 `*.test.tsx` 排序,与 shell 展开的顺序相同。
|
||||
- 补修后,Node 20 下 4,932 项、69 条失败,与旧写法 `tsx --test tests/*.test.ts tests/*.test.tsx` 在 Node 20 下的失败名单逐条相同(都是 Node 20 自身的旧问题)。
|
||||
- Node 22 下 4,932 项、24 条失败,名单与基线相同;`CI=true` 时 346 行。
|
||||
|
||||
@@ -421,3 +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、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) |
|
||||
|
||||
@@ -0,0 +1,101 @@
|
||||
# TASK:普通对话耗时——外部服务不再白等 + 补分段计时 — 2026-10-05
|
||||
|
||||
## 基线
|
||||
|
||||
- 开工时 `git fetch origin --prune` 之后的最新 `origin/staging`。写本单时是 `1252de3e`,已部署。
|
||||
- 分支 `codex/consult-latency-quickwins-20261005`,工作树 `.worktrees/consult-latency-quickwins-20261005`。
|
||||
- 依据:10-04 只读审计(Claude 子代理)。测量脚本与原始数据在 Claude scratchpad `bench/`,不入库;结论摘要见下面「事故实证」。
|
||||
|
||||
## 事故实证
|
||||
|
||||
一轮普通对话(问父母、问年运)全程串行,分为:分类(thinking 关,≤3 s)→ 第 0 步(thinking 开,决定调用排盘工具)→ 引擎 `/api/consultation_workflow`(每个领域一次,最多 2 个领域)→ 第 1 步写答案(thinking 开)。
|
||||
|
||||
1. **校正闸每个领域白等 4 秒。**
|
||||
- 咨询链里的 rectification gate(`jyotish_api_server._compute_rectification_gate` → `vedastro_gateway.run_gateway_packet`)会走到 `vedastro_service_adapter` 中调用 `_try_official_capability_runner_snapshot_bundle` 的那一段(按符号定位)。它起一个子进程,撞上 `DEFAULT_TIMEOUT_SECONDS = 4`。
|
||||
- 本机装上 vedastro SDK、按产品路径(`defer_optional_external_evidence: true`)实测,每次 4.6 s,第二轮仍是 4.6 s。原因:失败结果不写缓存。
|
||||
- 同一份响应里,顶层 `vedastro_gateway` 是 `official_verified`(BUG-727 的缓存在起作用),但 `rectification.vedastro_gateway.official_closure_reason = official_raw_response_missing`。
|
||||
- 屏蔽外网时,引擎每个领域 0.75 s;这 4 s 是额外的等待。
|
||||
- **这些数字在沙箱网络下测得,生产 VPS 是否同样撞超时,须在 staging 先确认(T0)。**
|
||||
2. **看不到分段时间。**
|
||||
- 现有埋点(`[agent-observability]`、`/admin/usage`)有 `first_byte`、工具耗时、`answer.first_output`、`run.total`、总 token。
|
||||
- 没有:推理 token 数;第 0 步、第 1 步各自的耗时;分类耗时进不了 usage。
|
||||
- 推理慢的主因(答题步约 1 万推理 token、约 45 秒)因此无法在线上逐轮核对,后续「限制推理强度」「精简说明」两单也没有办法验证。
|
||||
|
||||
## 决策记录(产品负责人 2026-10-05)
|
||||
|
||||
1. 做「不白等外部服务」和「补分段计时」两项,也就是审计建议的 ③ 和 ⑤。
|
||||
2. 「限制推理强度」「精简说明」「第 0 步关 thinking」**不在本单**,等产品提供临时模型 key、跑名人回测之后另开单。
|
||||
3. **不推翻 BUG-301**:顶层前台 VedAstro 照旧调用。本单只处理咨询链里校正闸那一处注定拿不到结果的等待。
|
||||
|
||||
## 硬红线
|
||||
|
||||
1. 遵守 `jyotish_api_server.py` 的增长冻结(AGENTS §6):类方法数不增长,`JyotishAPIHandler.__new__` 伪造点不增长。改动放进 `scripts/vedastro_service_adapter.py` / `scripts/vedastro_gateway.py`,或新模块。
|
||||
2. 不改生时校正的打分:v5 77 例与改动前逐项相同;冻结身份文件(`references/rectification_sealed_holdout.v1.json` → `production_scoring_files`)原则上不动。确实要动,就按 ERR-110 重新冻结,并同步 `frontend/tests/rectification-confirmation-gate.test.ts` 的两处路径(三栏)和 `tests/test_sealed_holdout_contract_freshness.py`。
|
||||
3. 生时校正面自己调用 VedAstro 的行为(若与咨询共用同一函数)不得改变;只对咨询路径生效,或者用开关区分。改动前后,生时校正的官方证据字段要逐项相同。
|
||||
4. 埋点只记**数字和状态**(毫秒、token 数、步骤号、是否超时),不记提示词、答案正文、出生资料、用户 ID(AGENTS §8、§5.6)。
|
||||
5. 不改普通对话提示词和数据卡内容。
|
||||
6. 验收要跑 Python 全量 `pytest tests` 并与开工基线对比失败名单(BUG-1220 的教训),不能只跑快速门。
|
||||
7. 不用 `git stash`;不切换别人的工作树;开工前先 `df -h /`。
|
||||
|
||||
## 任务分解
|
||||
|
||||
### T0 staging 先确认(只读,不改代码)
|
||||
- 用 staging 的公开接口或测试账号,对一位公开名人发一次 `/api/consultation_workflow`(或经普通对话发一次),记录:
|
||||
- `rectification.vedastro_gateway.official_closure_reason`;
|
||||
- 该请求的耗时。
|
||||
- 拿不到登录态时,在进度记录里写明「T0 环境缺口」,并给产品一条可照做的检查方法:在 `/admin/usage` 或容器日志里看工具耗时是否接近「领域数 × 4 s 以上」。
|
||||
- T0 不阻断 T1。
|
||||
|
||||
### T1 校正闸不再同步白等官方快照子进程
|
||||
- 咨询路径下(以请求里已有的 `defer_optional_external_evidence` 或同等标记区分),校正闸遇到官方快照 runner 时:
|
||||
- **优先**:复用顶层 `vedastro_gateway` 已经拿到的官方结果(同一份「出生数据 + 岁差 + 交点 + UTC 日期」缓存键,BUG-727)。
|
||||
- 拿不到时,不再前台起 4 秒子进程,直接标 `official_closure_reason = deferred_in_consultation`(或同义的明确状态),不伪装成 `verified`。
|
||||
- 失败或超时的结果按同一个缓存键写入**负缓存**,有效期到当日 UTC 结束,避免同日每轮重复等待。
|
||||
- **验收**:
|
||||
- 本机装 vedastro SDK 后,按产品路径实测:每个领域的耗时从约 4.6 s 降到不超过 1.0 s(测 3 位名人 × 父母 / 年运,冷、热各两轮),数字写进进度记录;
|
||||
- 咨询响应里的 `rectification.vedastro_gateway` 状态如实;
|
||||
- 生时校正面的 VedAstro 证据字段改动前后逐项相同;
|
||||
- 新增测试:咨询路径不调用 snapshot runner 子进程(mock 断言);负缓存同日命中;跨日失效。
|
||||
|
||||
### T2 补分段计时
|
||||
- `[agent-observability]`(`frontend/src/lib/agent-observability.ts`、`stream-agent-response.ts`、普通对话 route)新增字段:
|
||||
- 分类:`classification.durationMs`(已有的话,确认它进了 usage);
|
||||
- 第 0 步、第 1 步各自的 `durationMs`、`reasoningTokens`、`outputTokens`、`inputTokens`、`cachedInputTokens`(取供应商返回的 usage;拿不到的写 null,不估算);
|
||||
- 工具调用:各领域的 `durationMs`(已有的话保持不变);
|
||||
- `answer.reasoning_ms`:第 1 步开始到第一个正文字符的时间。
|
||||
- 同样的数字写进 `/admin/usage` 记录的 metadata(**不动表结构**,只用已有的 JSON 字段;若没有可用的 JSON 字段,写进进度记录,留给另单,本单不加迁移)。
|
||||
- `/admin/usage` 页面能看到这些数字最好;需要改 UI 时只加一列或一个展开区,并同步 `frontend/DESIGN.md`。
|
||||
- **验收**:
|
||||
- 新增单元测试:给定一组模拟的步骤结果,日志 / usage metadata 的字段齐全,并且不含正文;
|
||||
- 进度记录写明产品在哪里看、每个字段的含义(一张表)。
|
||||
|
||||
### T3 记录
|
||||
- `docs/BUG_HISTORY.md`:T1、T2 各一条(T1 关联 BUG-727、BUG-301;T2 关联 10-04 审计)。
|
||||
- `CHANGELOG.md`;`docs/tasks/PROGRESS-consult-latency-quickwins-20261005.md`(T0 结果或环境缺口、T1 前后耗时表、T2 字段表);`docs/tasks/README.md` 改为「已实现待验收」。
|
||||
|
||||
## 让步顺序
|
||||
|
||||
T2 → T1 → T0 → T3。T2 风险最低且是后续两单的前置;T1 若发现生时校正与咨询共用路径、难以隔离,先交 T2,T1 写清阻塞点。
|
||||
|
||||
## 开工前置命令
|
||||
|
||||
```bash
|
||||
cd /workspace/Jyotisha && git status -sb | head -1
|
||||
git fetch origin --prune
|
||||
git worktree add -b codex/consult-latency-quickwins-20261005 .worktrees/consult-latency-quickwins-20261005 origin/staging
|
||||
df -h /
|
||||
grep -oE "BUG-1[0-9]{3}" docs/BUG_HISTORY.md | sort -u | tail -1
|
||||
python3 scripts/pre_work_check.py --remote-timeout 8 --command-timeout 45
|
||||
```
|
||||
|
||||
外部引擎改动需先读 `docs/research/pre_work_error_ledger.md`(AGENTS §9),并检索 `docs/BUG_HISTORY.md` 里的 VedAstro / BUG-301 / BUG-727 记录。
|
||||
|
||||
## 验收口径
|
||||
|
||||
- Python:快速门 + **全量 `pytest tests` 与基线对比** + 英文对照 + 隐私测试;v5 逐项相同。
|
||||
- 前端:`tsc` 0 错;lint 0 error;`npm test` 失败名单与开工基线逐条相同。
|
||||
- 不 push,由 Claude 验收后推 staging;部署后由产品在 `/admin/usage` 看一轮真实分段。
|
||||
|
||||
## BUG 编号起点
|
||||
|
||||
BUG-1231(10-05 实测最大号 BUG-1230;开工时再核对)。
|
||||
@@ -0,0 +1,167 @@
|
||||
# TASK · 「纳迪秒级校准」可证伪检验(研究单,2026-10-05)
|
||||
|
||||
## 基线
|
||||
|
||||
- `origin/staging` @ `853772bc`(写作时 head;开工时以最新 `origin/staging` 为准)。
|
||||
- 分支 `codex/rectification-nadi-seconds-research-20261005`,工作树 `.worktrees/rectification-nadi-seconds-research-20261005`。
|
||||
- **离线研究,不改生产代码**:不动 `scripts/research/sealed_holdout_rerun.py::PRODUCTION_FILES` 里的冻结文件,不动 `scripts/active_rectification_event_engine.py`、`scripts/dasha_calculator_enhanced.py`、`scripts/divisional_charts_extended.py`、任何常数,也不动 `frontend/`。研究代码只放在 `scripts/research/` 和 `tests/`,引擎函数只能 import,不能改。
|
||||
- 数据:`references/real_case_calibration/minute_rectification_holdout_v5.json`(77 例 Rodden AA 公开名人,964 件事;BUG-1090)。数据集声明的口径是 `ayanamsa = raman`、`node_mode = mean`,本单默认沿用;换口径只在 N4 里做。
|
||||
|
||||
## 先读
|
||||
|
||||
- `docs/research/rectification_scoring_research_2026_09_29.md`:打分层三条 no_benefit,42 个特征里 36–39 个是噪声。
|
||||
- `docs/research/rectification_varga_resolution_2026_09_30.md`:盘型口径;±10 六题后 D9 / D10 头段 = 真值 88% / 86%。
|
||||
- `docs/research/holdout_v5_build_2026_09_29.md`:v5 基线;六题后头名 0.64 / 0.49 / 0.26,真值在区间内 76/77、76/77、75/77。
|
||||
- `docs/research/rectification_minute_resolution_closure_2026_09_14.md` §1–§2
|
||||
|
||||
## 事故实证
|
||||
|
||||
### 起因
|
||||
|
||||
竞品 2026-10 宣传「纳迪辅助校准,符合条件时精度最高可达秒级」,验证方式是**问前事、具体到天**。产品问:我们为什么做不到。
|
||||
|
||||
Claude 2026-10-05 的判断是:秒级的盘我们能算,做不到的是**验证**。但这个判断有一个没测过的缺口,所以写本单把它补上。
|
||||
|
||||
### 物理量(Claude 用 swisseph 实算;样本为上海、1990-06-15,Lahiri)
|
||||
|
||||
| 量 | 数值 | 含义 |
|
||||
| --- | --- | --- |
|
||||
| 上升点移动速度 | 0.21–0.37°/分钟 | 这个纬度和日期,每 4 秒约走 1′ |
|
||||
| D150 一段持续多久 | 32–57 秒 | 纳迪段本身就是「不到一分钟」的量级 |
|
||||
| D60 一段持续多久 | 81–143 秒 | |
|
||||
| 出生时间差 1 秒,Vimshottari 时间轴平移 | 0.025 天(太阳大运)到 0.084 天(金星大运) | 想把一件事对准到天,出生时间就要准到约 12–40 秒 |
|
||||
| 岁差 Lahiri 与 KP 相差 | 0.097°(5.8′) | ≈ 上升点 27 秒 |
|
||||
| 出生地东西方向差 10 公里(北纬 31°) | — | ≈ 恒星时 25 秒 |
|
||||
|
||||
### 本仓现状
|
||||
|
||||
1. **主链只用到第 3 层小运**:`scripts/active_rectification_event_engine.py::_active_vimshottari` 只返回 MD / AD / PD 三层主星,`_score_event` 也只对这三层计分。第 4 层 Sookshma、第 5 层 Prana 的算法在 `scripts/dasha_calculator_enhanced.py::calculate_five_level_dasha` 里已经有,但**从来没有进过评价集**。
|
||||
2. **D150 只有等分实现**:`scripts/divisional_charts_extended.py` 中的 `calc_custom_varga(degree, 150)` 只是把每宫等分成 150 段。古典纳迪段(Chandra Kala Nadi 体系)按动 / 固 / 变宫区分顺排还是逆排,而且每段都有名称和描述文本。这两样本仓都没有。
|
||||
3. **v5 的日精度事件**:一共 244 件。43/77 例有 ≥3 件,19/77 例有 ≥5 件,4/77 例一件也没有(年精度 577 件,月精度 143 件)。
|
||||
4. **v5 的「真值出生时间」大多是取整过的**:77 例里有 **52 例的分钟数是 5 的整数倍**(:00 有 14 例,:30 有 10 例)。随机分布下只该有约 20%(约 15 例)。也就是说,**标准答案本身多数只精确到 5 分钟左右,秒级没有可对照的真值**。分钟数不是 5 的倍数的只有 25 例。
|
||||
|
||||
## 根因(为什么至今不能对外说秒级)
|
||||
|
||||
1. **没有检验过**:Sookshma / Prana 小运和 D150 这两层,是「秒级」说法唯一可能的来源,但从来没在已知出生时间的人身上测过。
|
||||
2. **「问前事对上了」这种验证分辨不出真假**:小运分 5 层,每层 9 颗主星,再乘上十几张分盘,几乎任何一秒都能找到一条解释过去某天的路径。用事件定出时间,再用同一批事件证明这个时间,是循环论证。必须有安慰剂对照和留出预测。
|
||||
3. **输入本身的误差已经比秒大**:岁差流派、出生地坐标、「出生时刻」的定义、出生证明取整,每一项都在几十秒到几分钟之间。
|
||||
|
||||
## 决策记录(产品 2026-10-05)
|
||||
|
||||
- **批准本研究单**:对竞品式的「纳迪 / 小运对日子 → 秒级」做可证伪检验。**只做研究,不做实现**。只有通过本单的「过门标准」,才另立实现单。
|
||||
- **本单不算重开 BUG-1091**:BUG-1091(打分层,closed_by_design)关的是「给已有 42 个特征重新调权重」。本单测的是**从来没进过评价集的新层**(Sookshma / Prana 小运、D150),目的是证伪,不是调参。执行方**不得**借本单去调已有特征的权重,也不得改 `DOMAIN_CONFIG` 的宫位映射。
|
||||
- **结果出来之前,产品对外不说「秒级」**。如果本单不过门,产品要的是一份对外口径草稿(见 N6),而不是功能。
|
||||
- 纳迪段的**原文描述**(Chandra Kala Nadi 各段的命运文字):仓库里没有合法来源。本单**不转录、不爬取、不让模型凭记忆生成**这类文本。按描述文字比对经历这一路记为 `blocked`,D150 只做结构层检验。
|
||||
|
||||
## 硬红线
|
||||
|
||||
1. **规则先登记再跑**:N1–N3 用到的「事件被解释」规则族、阈值和随机种子,必须先写进 `docs/research/nadi_seconds_preregistration_2026_10_05.json` 并单独 commit,然后才能第一次在 v5 上跑。PROGRESS 里写明这个 commit 的 SHA。看过结果后再改的规则,只能放进「事后」列,不能用来判过门。
|
||||
2. **真值不可见**:排序器与拟合器不得读取 `true_minute`,也不得读取出生资料里的 `time` 字段作为输入。只有评测函数可以读。新增测试要断言这一点,沿用 `truth_hidden_from_ranker` 的检查方式。
|
||||
3. **必须有安慰剂对照**:每个「对上了」的指标,都要同时报告用安慰剂日期跑出的同一指标,置换检验以「整例」为单位抽样。只报真实日期、不报安慰剂的数字,一律不算数。
|
||||
4. **不得写「秒级准确率」**:所有准确率都以「与记录时间相差 ≤1 分钟 / ≤2 分钟」为口径。取整组(52 例)和非取整组(25 例)分开列,不得合并后只报有利的一组。
|
||||
5. 生产代码零改动;`PYTHONHASHSEED=0`,所有随机数带固定种子;所有 JSON 两次复跑逐字节一致。
|
||||
6. 快速门结果与开工基线逐条一致;隐私守卫全绿;只用 v5 已有的公开名人数据,不新增人物,不访问 astro.com。
|
||||
7. 负结果完整写出。任何一格不过门,都把数字写进 PROGRESS,不得通过调阈值凑出通过。
|
||||
|
||||
## 任务分解
|
||||
|
||||
### N0 · 研究底座与真值审计
|
||||
|
||||
- 新建 `scripts/research/nadi_seconds_lib.py`:
|
||||
- 对任意候选时刻,返回 v5 每件事件当天的**五层** Vimshottari 主星(调用 `calculate_five_level_dasha` 或同一条主链的递归切分)。
|
||||
- 返回 D150 段号,两种口径:等分 `calc_custom_varga`,以及古典纳迪段排序(动宫顺排、固宫逆排、变宫从中段起排)。古典排序的出处写在代码注释里;出处不确定的写 `variant_unverified`,不得冒充定论。
|
||||
- 网格:±10 分钟窗口按 5 秒一步(241 个候选),另在记录时间 ±2 分钟内按 1 秒一步。
|
||||
- **对账**:前三层主星必须与 `_active_vimshottari` 逐例、逐事件 0 差,对账结果写进 JSON。
|
||||
- **真值审计表**:77 例的分钟数分布、取整组与非取整组名单(只列 case_id)、每例日精度事件数。
|
||||
- 验收:`tests/test_nadi_seconds_research.py` 覆盖以下几点:
|
||||
- 前三层与主链对账 0 差;
|
||||
- 五层边界首尾相接;
|
||||
- D150 段号在段边界两侧正确翻转;
|
||||
- 跨午夜正确;
|
||||
- 不读 `true_minute`;
|
||||
- 两次复跑逐字节一致。
|
||||
|
||||
### N1 · 竞品式拟合率与安慰剂(核心)
|
||||
|
||||
- **规则族**:在预登记里固定,至少包括以下两族:
|
||||
- (a) 事件当天,第 k 层小运主星 ∈ 该领域目标宫主(`DOMAIN_CONFIG` + `_house_lords`,只 import,不改),k 分别取 4 和 5;
|
||||
- (b) 第 1–5 层里有 ≥m 层命中。
|
||||
- **对每例、每个候选秒**:计算这一秒能「解释」多少件日精度事件。
|
||||
- **报三组数字**:
|
||||
- 真值分钟内各秒的拟合率;
|
||||
- 窗口内随机秒的拟合率;
|
||||
- **安慰剂日期**(每件事的日期按固定种子平移 ±30–180 天;另做一组把事件在案例之间互换)下的拟合率。
|
||||
- 过门:真值秒的拟合率高于安慰剂,整例置换检验 p < 0.05。样本只用 ≥3 件日精度事件的 43 例。
|
||||
- 验收:表格写进研究文档。另外报一个关键数:「窗口内能解释全部日精度事件的秒占多大比例」。这个比例接近 100%,就说明「问前事对上了」分辨不出真假。
|
||||
|
||||
### N2 · 留出预测(竞品验证方式的诚实版)
|
||||
|
||||
- 用 ≥4 件日精度事件的案例做留一件:用其余事件拟合出最佳秒(并列的全部保留),再看被留出的那件事在这些秒上是否被解释。
|
||||
- 对照两组:同窗口内随机秒;被留出事件换成安慰剂日期。
|
||||
- 过门:留出事件的解释率高于随机秒,整例 bootstrap 95% 置信区间不含 0。
|
||||
|
||||
### N3 · 能不能找回记录时间
|
||||
|
||||
- **N3a 单独排序**:只按纳迪层拟合排序,看头名与记录时间相差 ≤1 分钟 / ≤2 分钟的比例。对照均匀随机(±10 窗口下 ≤1 分钟约 3/21)和引擎先验头名(v5 ±10 为 0.18)。
|
||||
- **N3b 线上六题后交付区间里再细分**(这是产品真正关心的问题):沿用六题回放(`scripts/research/futile_collect_stop_replay.py` 的口径),在线上交付区间内用纳迪层排序,看头名 ≤1 分钟的命中率能否往上提。0.64 是 v5 文档里「头名」口径的数,不一定等于「≤1 分钟」口径;**先在 ≤1 分钟口径下复现线上基线,再做比较**。
|
||||
- 过门(两条都要满足):
|
||||
- N3b 头名 ≤1 分钟命中的提升,整例 bootstrap 95% 置信区间不含 0;
|
||||
- 真值仍在区间内的比例不低于 76/77。
|
||||
- 取整组和非取整组分开列;±30 / ±60 的数字只报告,不参与判定。
|
||||
|
||||
### N4 · 输入噪声地板(不论 N1–N3 结果如何都要做)
|
||||
|
||||
- 在**记录时间**上,统计每件日精度事件当天的第 4 / 第 5 层主星,以及本命 D150 段号,在下列扰动下有多少比例发生变化:
|
||||
- 岁差 Raman ↔ Lahiri ↔ KP ↔ True Chitra;
|
||||
- 出生时间 ±15 秒、±30 秒、±60 秒;
|
||||
- 出生地坐标东西方向 ±10 公里。
|
||||
- 交点 mean ↔ true 只影响罗睺 / 计都的位置,不影响 Vimshottari,单列说明即可。
|
||||
- 判定(只决定对外口径,不决定是否过门):只换一个岁差流派,第 5 层主星或 D150 段号就有 >20% 的事件 / 例子跟着变,结论就是「秒级结果跨软件不可复现」。这种情况下,即使 N1–N3 全部通过,对外文案也不得出现「秒级」。
|
||||
|
||||
### N5 · D150 结构层(低优先)
|
||||
|
||||
- 不使用文本,只测结构:本命 D150 段主星是否落在事件领域目标宫主里,以及 D150 上升星座的宫主星在事件当天的小运里是否激活。规则先登记,方法同 N1(要有安慰剂)。
|
||||
- 按描述文字比对经历这一路写 `blocked`,原因写「无合法文本来源」。
|
||||
|
||||
### N6 · 结论、对外口径与记录
|
||||
|
||||
- `docs/research/rectification_nadi_seconds_2026_10_05.md`:
|
||||
- 一句话结论表:N1–N5 各一行,写明是否过门、关键数字、安慰剂对照;
|
||||
- N4 噪声地板表;
|
||||
- 结论:通过的话,写实现单要点(改哪些模块、不改哪些模块);不通过,就 `closed_by_design`。
|
||||
- **对外口径草稿**(不论结论如何都要写,对照 `frontend/docs/VOICE.md`):三到五句,向用户说明我们的校正能做到什么精度、靠什么数据,以及为什么不说「秒级」。**不得点名竞品。**
|
||||
- 同步更新:
|
||||
- `docs/research/ACTIVE_FRONTS.md` 索引;
|
||||
- `docs/tasks/PROGRESS-rectification-nadi-seconds-research-20261005.md`;
|
||||
- `docs/BUG_HISTORY.md` 对应编号的状态;
|
||||
- `docs/tasks/README.md` 状态板这一行。
|
||||
|
||||
## 让步顺序
|
||||
|
||||
N0 > N1 > N4 > N2 > N3 > N5 > N6 里的图表。
|
||||
|
||||
- N0、N1、N4 缺一项都不能合入。
|
||||
- N2 / N3 / N5 没做的,写 `not_started`,不得写成结论。
|
||||
- 时间不够时,先砍 N5,再砍 N3 的 ±30 / ±60 列。
|
||||
|
||||
## 开工前置命令
|
||||
|
||||
```bash
|
||||
git fetch origin --prune
|
||||
git worktree add -b codex/rectification-nadi-seconds-research-20261005 .worktrees/rectification-nadi-seconds-research-20261005 origin/staging
|
||||
cd .worktrees/rectification-nadi-seconds-research-20261005
|
||||
export PYTHONHASHSEED=0
|
||||
python3 -m pytest -q tests/test_scoring_research.py tests/test_holdout_v5_schema.py tests/test_varga_resolution_research.py tests/test_rectification_validation_integrity_gate.py tests/test_repo_privacy_markers.py # 基线
|
||||
python3 scripts/research/futile_collect_stop_replay.py --help
|
||||
python3 -c "import swisseph; print(swisseph.version)"
|
||||
```
|
||||
|
||||
长任务每完成一个 N 就 commit 一次检查点,因为 09-29 有子代理被限额打断后成果没落盘的教训。研究分支只 commit,不推 staging;推送由验收方完成。
|
||||
|
||||
## BUG 编号
|
||||
|
||||
开工时核对 `docs/BUG_HISTORY.md` 当前最大号(写作时是 `BUG-1230`)。同日 `TASK-consult-latency-quickwins-20261005` 已从 BUG-1231 起编号,为避免撞号,本单预留 **BUG-1240**:「秒级 / 纳迪校准从未经过检验:Sookshma / Prana 小运与 D150 不在评价集内,真值时间多为取整值」。如果开工时 1240 已被占用,顺延到下一个空号,并在 PROGRESS 里写明。
|
||||
|
||||
执行方以 `investigating` 登记。研究结束后:通过的改为 `resolved`(另立实现单);不通过的改为 `closed_by_design`,并附上对外口径。
|
||||
|
||||
关联:BUG-560、BUG-1090、BUG-1091、BUG-1105。
|
||||
@@ -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)
|
||||
|
||||
@@ -10,14 +10,22 @@
|
||||
* The process exit code is the test runner's exit code (BUG-995).
|
||||
*/
|
||||
import { spawn } from "node:child_process";
|
||||
import { mkdirSync, readFileSync } from "node:fs";
|
||||
import { mkdirSync, readFileSync, readdirSync } from "node:fs";
|
||||
import path from "node:path";
|
||||
import { fileURLToPath } from "node:url";
|
||||
|
||||
const frontendDir = path.resolve(path.dirname(fileURLToPath(import.meta.url)), "..");
|
||||
const repoRoot = path.resolve(frontendDir, "..");
|
||||
const tsxCli = path.join(frontendDir, "node_modules", "tsx", "dist", "cli.mjs");
|
||||
const testGlobs = ["tests/*.test.ts", "tests/*.test.tsx"];
|
||||
// Expanded here instead of handing `tests/*.test.ts` to `node --test`: Node 20
|
||||
// does not expand test globs ("Could not find tests/*.test.ts"), while the old
|
||||
// `tsx --test tests/*.test.ts tests/*.test.tsx` script relied on the shell. Same
|
||||
// order as the shell: every *.test.ts sorted, then every *.test.tsx sorted.
|
||||
function testFiles() {
|
||||
const names = readdirSync(path.join(frontendDir, "tests"));
|
||||
const pick = (suffix) => names.filter((name) => name.endsWith(suffix)).sort().map((name) => `tests/${name}`);
|
||||
return [...pick(".test.ts"), ...pick(".test.tsx")];
|
||||
}
|
||||
const failureBlockLimit = 60;
|
||||
|
||||
function envEnabled(name) {
|
||||
@@ -193,7 +201,7 @@ function invokedDirectly() {
|
||||
async function runTests() {
|
||||
const extra = process.argv.slice(2);
|
||||
if (!gateMode()) {
|
||||
const code = await runRunner(["--test", ...testGlobs, ...extra]);
|
||||
const code = await runRunner(["--test", ...testFiles(), ...extra]);
|
||||
process.exit(code);
|
||||
}
|
||||
|
||||
@@ -208,7 +216,7 @@ async function runTests() {
|
||||
"--test-reporter=tap",
|
||||
"--test-reporter-destination",
|
||||
tapPath,
|
||||
...testGlobs,
|
||||
...testFiles(),
|
||||
...extra,
|
||||
],
|
||||
{ suppressFailedTestsTrailer: true },
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 (
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
@@ -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',
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user