fix(consult): limit career direction and single-flight chart warm
The career checklist now says the direction section only covers strength, resistance, timing, and three events. Warm builds one packet per cache key and skips when its own cap is full. A three-repeat biography run still has one career answer that names a direction from a planet nature, so this is not accepted. Not pushed.
This commit is contained in:
@@ -1,5 +1,11 @@
|
||||
# 印度占星 Skill 更新日志
|
||||
|
||||
## 2026-10-09 — 事业方向只写力量、阻力和时间;预热同一份资料只算一次(回测未达标,未推送,未合入,未上线)
|
||||
|
||||
- 事业清单写明:「方向」只写力量、阻力、时间,以及机会出现、成形、公开落地这三件事。职业类型和工作方式仍问用户。系统提示没有改。判断表的行和取值没有改。
|
||||
- 保存后的预热按同一张盘、同一天只算一次。同时最多预热 1 份,满了就跳过,不会让紧接着的提问收到 429。预热失败不影响保存。
|
||||
- 生平回测 108 份,严重冲突 1,对照组误报 0。27 份事业题里仍有 1 份把方向写成信息和表达,这一版不标待验收。Skill 版本不变。
|
||||
|
||||
## 2026-10-09 — 事业清单写明代表星不猜职业,预热不再占重计算名额(已合入 staging)
|
||||
|
||||
- 事业清单加一条:判断表里的自然代表星、Jaimini 代表星、大运主星只看力量和受冲,不按星的本性推断职业类型、做事方式或台前 / 幕后。职业类型仍问用户。系统提示没有改。
|
||||
|
||||
+12
-12
@@ -17490,35 +17490,35 @@
|
||||
|
||||
## BUG-1303 | 事业题仍偶发按星性写方向(BUG-1301 残余)
|
||||
|
||||
- 状态:investigating(待 `TASK-consult-career-type-residual-20261009.md`)
|
||||
- 状态:investigating(`TASK-consult-career-type-residual-20261009.md`;清单已加「方向」句,3 次抽样仍有 1 次行业冲突,按任务书停下,不标 resolved)
|
||||
- 首次发现 / 最近更新:2026-10-09 / 2026-10-09
|
||||
- 来源:BUG-1301 修复后的生平回测。
|
||||
- 影响面:普通咨询事业回答。
|
||||
- 现象:18 份事业回答中 1 份把方向写成「知识、表达、信息和资源」一类,与公开人物的实际职业类型不符。
|
||||
- 现象:上一轮 18 份事业回答中 1 份把方向写成「知识、表达、信息和资源」一类。本轮 27 份里仍有 1 份把方向写成信息和表达,与公开人物的实际职业类型不符。
|
||||
- 触发条件:判断表多处点名同一颗代表星;模型抽样波动。
|
||||
- 根因:待查。初判是清单只禁止推断、没有规定「方向」一节写什么。
|
||||
- 修复:待做。
|
||||
- 验证:待做。
|
||||
- 防复发:待做。
|
||||
- 根因:上一句只禁止按星性推断,没有规定「方向」一节写什么。补上「方向只写力量、阻力、时间与三件事」之后,3 次抽样仍有 1 次越出这句。
|
||||
- 修复:事业清单把「方向」限定为力量、阻力、时间与三件事(机会出现 / 成形 / 公开落地),职业类型与工作方式仍问用户。不进系统提示,不删判断表,不改取值。回测未把行业类型冲突降到 0,所以不标 resolved。
|
||||
- 验证:`frontend/tests/consult-condensed-checklist-20260927.test.ts` 锁住这句,系统提示源码切片仍是 15,291 字符。生平回测 108 份、失败 0。阅读评分见 `docs/tasks/PROGRESS-consult-career-type-residual-20261009.md`:严重冲突 1,行业类型冲突 1,对照组误报 0。
|
||||
- 防复发:上述测试锁住句子。评分口径没有改。产品未授权再加一条或重跑口径。
|
||||
- 相关记录:BUG-1301、BUG-1176。
|
||||
- 复发自:BUG-1301(残余,非新问题)。
|
||||
- 修复版本:无。
|
||||
- 修复版本:本分支提交。未推送,未部署。
|
||||
|
||||
## BUG-1304 | 保存后马上提问,预热与用户请求重复计算同一份资料包
|
||||
|
||||
- 状态:investigating(待 `TASK-consult-career-type-residual-20261009.md`)
|
||||
- 状态:resolved(分支 `codex/consult-career-type-residual-20261009`;有针对性回归测试。未推送,未部署)
|
||||
- 首次发现 / 最近更新:2026-10-09 / 2026-10-09
|
||||
- 来源:Claude 验收 BUG-1302 时读代码发现。
|
||||
- 影响面:保存星盘后立刻提问的第一轮等待时长与服务器 CPU。不会因此返回 429。
|
||||
- 现象:同一张盘同一天的资料包可能被并行算两次;多次保存时预热无上限并行。
|
||||
- 触发条件:预热未完成时同一张盘的提问到达,或短时间多次保存。
|
||||
- 根因:`get_or_build` 没有按缓存键单飞;预热跳过重计算名额后没有自己的上限。
|
||||
- 修复:待做。
|
||||
- 验证:待做。
|
||||
- 防复发:待做。
|
||||
- 修复:同一个缓存键只构建一次,后到的请求等待并读取这份结果。预热单独限流,同时最多 1 份,满了直接跳过。用户请求不占这个名额,也不会因为预热收到 429。不新增 API 类方法,`jyotish_api_server.py` 没有增长(仍是 10,970 行)。
|
||||
- 验证:`tests/test_consult_card_full_source.py` 在预热进行时发同键用户请求,构建次数是 1;三份不同键的预热只有 1 次构建,另外两份跳过且响应里没有资料包。`tests/test_api_heavy_compute_gate.py` 与 `tests/test_api_server_growth_contract.py` 通过。
|
||||
- 防复发:上述并发测试。
|
||||
- 相关记录:BUG-1300、BUG-1302。
|
||||
- 复发自:无。
|
||||
- 修复版本:无。
|
||||
- 修复版本:本分支提交。未推送,未部署。
|
||||
|
||||
## BUG-1305 | 保存星盘后安排预热抛错,已写库的保存返回 500(staging 门禁 run 1784 红)
|
||||
|
||||
|
||||
@@ -0,0 +1,51 @@
|
||||
# PROGRESS · 事业方向残余 + 预热单飞(2026-10-09)
|
||||
|
||||
分支 `codex/consult-career-type-residual-20261009`,基线 `origin/staging` `cd625e8a`。状态:执行中。H1 的 3 次抽样仍有 1 次行业类型冲突,按任务书停下,不标待验收,不改口径。未推送,未合入,未上线。
|
||||
|
||||
## 结论
|
||||
|
||||
事业清单写明「方向」只写力量、阻力、时间和三件事。9 人 × 4 领域 × 3 次,108 份都有回答。严重冲突 1,对照组误报 0。27 份事业题里仍有 1 份把方向写成信息和表达(齐达内事业第 1 次),行业类型冲突不是 0。预热按缓存键单飞,同时最多 1 份,满了跳过。
|
||||
|
||||
## 修复
|
||||
|
||||
| 项 | 做法 | 结果 |
|
||||
| --- | --- | --- |
|
||||
| H1 方向一节 | 事业清单仍是 8 行。原来的「分三件事看」换成一句:方向只写力量、阻力、时间与三件事(机会出现 / 成形 / 公开落地),职业类型与工作方式一律问用户。不删判断表的行,不改取值,不改系统提示。 | 测试锁住这句。系统提示源码切片仍是 15,291 字符。回测未过,见下表。BUG-1303,仍 investigating。 |
|
||||
| H2 预热单飞 | `get_or_build` 同一个缓存键只算一次,后到的等这份结果。预热自己的上限是 1,满了跳过,不返回 429。用户请求不占这个名额。不改 `jyotish_api_server.py`,不新增 API 类方法。 | 并发测试通过:三份不同键的预热只有 1 次构建、另外两份跳过;预热进行中,同键用户请求不再另建一份。BUG-1304。 |
|
||||
| H3 | 本机不是 Linux。 | 快速门、Python 全量、npm test 全量、next build 没跑。不把上一轮 Linux 数字当成这一轮。 |
|
||||
|
||||
`scripts/jyotish_api_server.py` 仍是 10,970 行。类方法数没有增加。`page.tsx` 没有改。系统提示文件没有改。判断表没有改。
|
||||
|
||||
## 生平回测
|
||||
|
||||
模型 `deepseek-flash`。key 只在当次命令的环境变量里,没有写入仓库、日志正文或本文。引擎仍是 golden 桩,参照日 2026-09-27,岁差 raman,交点 mean。判断靠阅读,脚本只写确定性检查。口径与前两轮相同:两头都写到的范围不算误报;直接写成另一类职业,或写成幕后、不靠曝光,算行业类型冲突。
|
||||
|
||||
| 项 | 这一轮(3 次) |
|
||||
| --- | --- |
|
||||
| 份数 / 失败 / 空答 | 108 / 0 / 0 |
|
||||
| 墙钟 | 1399.2 秒 |
|
||||
| 严重冲突 | 1 |
|
||||
| 对照组误报 | 0 |
|
||||
| 行业类型冲突 | 1 / 27 |
|
||||
| 事业题请问一句 | 27 / 27 |
|
||||
| 禁句 / 性别代词 / must_not 原文 | 0 / 0 / 7 |
|
||||
| 平均秒 / 平均字数 / 平均卡 / 最大卡 | 50.6 / 734.9 / 14,129.8 / 16,281 |
|
||||
| 平均输入 / 输出 / 推理 / 缓存输入 | 29,133.5 / 9,571.3 / 8,995.3 / 21,745.8 |
|
||||
|
||||
齐达内事业第 1 次把方向写成信息和表达。第 2 次、第 3 次没有写成另一类职业,也没有写成幕后。其余 24 份事业题没有写成分析、顾问、技术,也没有写成幕后或不靠曝光。泰勒三次都没有写成不靠曝光。
|
||||
|
||||
must_not 原文 7 处:布什父母两次、齐达内父母三次,都是「可能不在身边,也可能人在、只是话少」。奥巴马婚恋第 2 次的「离婚」是「离婚姻最近」里的字面重合,正文没有写离婚。这 7 处都不算对照组误报。
|
||||
|
||||
用量合计:输入 3,146,423,输出 1,033,704,推理 971,496,缓存输入 2,348,544。
|
||||
|
||||
## 本机检查
|
||||
|
||||
| 项 | 结果 |
|
||||
| --- | --- |
|
||||
| 事业清单及相关测试 | 19 条通过。含新句子、系统提示长度、事业题仍问用户。 |
|
||||
| 预热单飞、重计算名额、增长合同 | 通过。API 文件没有增长。 |
|
||||
| tsc | 0 |
|
||||
| lint | 0 error,125 条原有 warning,没有清理。 |
|
||||
| 快速门、Python 全量、npm test 全量、next build、Linux | 本轮没跑。 |
|
||||
|
||||
H1 未过。按让步顺序停下,交给产品决定。不改评分口径,不加第二条规则,不标待验收。
|
||||
@@ -132,7 +132,7 @@
|
||||
| `TASK-consult-card-full-source-20261007.md` | `PROGRESS-consult-card-full-source-20261007.md` | **对话数据卡与报告同源 + 领域综合判断表**:事业卡只有约 7,600 字符、无自然代表星、无大运主星关系、D10 联动只能补取、补取上限 1;全 12 领域;综合判断表、参与清单(不止 D1/D9,含精度闸)、补取 3 次、资料目录均已确认;KP 不上卡、不联网、不接书籍语料;排在 sync6 与报告单之后 | 验收未通过(`35dafd1e`,见 fix-20261009) | `codex/consult-card-full-source-20261007` |
|
||||
| `TASK-consult-card-full-source-fix-20261009.md` | `PROGRESS-consult-card-full-source-fix-20261009.md` | **对话卡同源修复**:缓存键无日期致当前大运冻结、整卡 37.5K(每格带内部路径)、星照整列 blocked(资料库其实有出处)、前端测试手写数据、无预热;合入前用临时 key 跑模型 A/B 与生平回测 | F1–F5 已验收,F6 未过(见 fix2) | `codex/consult-card-full-source-20261007` |
|
||||
| `TASK-consult-card-full-source-fix2-20261009.md` | `PROGRESS-consult-card-full-source-fix2-20261009.md` | **对话卡同源第二轮**:齐达内事业被「水星=说写算」带偏致严重冲突 1→3;事业清单加「代表星只看力量受冲、不按星性猜职业」后重跑回测(≤1 才合入);预热改为不占重计算名额 | 已验收并快进合入 staging(2026-10-09,产品按严重冲突 ≤ 1 放行;剩余 1 次见 career-type-residual) | `codex/consult-card-full-source-20261007` |
|
||||
| `TASK-consult-career-type-residual-20261009.md` | — | **事业题按星性写方向的残余 + 预热单飞**:回测仍有 1 次(齐达内事业第 2 次;合入前 staging 是泰勒事业第 2 次);保存后马上提问时预热与用户请求并行重复算同一份资料包 | 待领取 | — |
|
||||
| `TASK-consult-career-type-residual-20261009.md` | `PROGRESS-consult-career-type-residual-20261009.md` | **事业题按星性写方向的残余 + 预热单飞**:方向一节只写力量、阻力、时间和三件事;预热按缓存键单飞、同时最多 1 份 | 执行中(H1:3 次抽样仍有 1 次行业冲突,未让步,未标待验收) | `codex/consult-career-type-residual-20261009` |
|
||||
| — (产品 09-27 拍板 D1–D4,直接执行) | [PROGRESS](PROGRESS-home-landing-blank-20260927.md) | **登录后 / 裸 `/` 落空白首页**:真机登录后在「首页」提问其实问进了上一次生时校正。删登录返回存根(401 / 次级页链接写、登录后写回 `?c=` 打开),裸 `/` 不再落最近会话,一律当前人物的空白首页(复用空草稿);`?c=` / `?new=1` / 对话内刷新不变;推翻 BUG-1038 存根与 BUG-599 默认落点 | 已验收(真机欠) | `codex/home-landing-blank-20260927`(BUG-1052,本地未推) |
|
||||
| — (产品 09-26 口头拍板 D1–D3,直接执行) | `PROGRESS-consult-answer-truncation-20260926.md` | **普通咨询回答写到一半被掐断仍扣点(BUG-1051,复发自 BUG-305)**:工具循环与写回答共用 110 秒 signal;Mastra 1.50 超时不抛错(`abort` 块 + `finish(tripwire)` 后正常关流),结算只认抛错与 `length`。D1 写回答自有 70 秒时钟(首用起算,续写 / 回答重试共用,最坏 180 秒,`maxDuration` 240);D2 写回答的最后一个流不是 `stop` 且有正文 → `answer_truncated`、不扣点、记 abort 步、不冲半句,`length` 续写不变;D3 观测加 `composeFinishReason` / `composeAborted` / `answerVisibleChars` | 已验收(真机欠) | `codex/consult-answer-truncation-20260926`(本地,未推送);新回归 15 条用真实 Mastra `Agent`(修复前 11 条红);全量失败名单 0 新增;Python 948/1;`/` ○、gzip 0%;真机清单 `docs/testing/consult-answer-truncation-20260926.md` |
|
||||
| `TASK-consult-evidence-card-research-20260927.md` | `PROGRESS-consult-evidence-card-research-20260927.md` | **普通对话数据卡调研**:引擎输出逐项分五类计量(现约 4 万 token、父母问题相关约 3.5%);四处领域→技法来源对账并起草各领域数据卡(家庭拆父母/子女);卡体量与逐字一致性;按卡算的提速空间;反馈迭代埋点方案。只调研不改线上 | 已验收(待产品拍板 7 项) | `codex/consult-evidence-card-research-20260927` 快进 staging;报告 `docs/research/consult_evidence_card_research_2026_09_27.md`;投影缺口记 BUG-1054 investigating |
|
||||
|
||||
@@ -90,13 +90,21 @@ export const CONSULT_BRIEF_CHECKLIST_SOURCE = "consultation checklist (TASK-cons
|
||||
export const CAREER_SIGNIFICATOR_STRENGTH_RULE =
|
||||
"判断表里的代表星(自然代表星、Jaimini 代表星、大运主星)只用来看力量与受冲,不按星的本性推断职业类型、做事方式或台前 / 幕后。";
|
||||
|
||||
/**
|
||||
* Product 2026-10-09 residual: saying "do not infer" was not enough. The
|
||||
* direction section must say what it is allowed to contain. Career type and
|
||||
* working style still ask the user (CAREER_FIELD_ASK_RULE). Checklist only.
|
||||
*/
|
||||
export const CAREER_DIRECTION_SECTION_RULE =
|
||||
"「方向」一节只写力量、阻力、时间与三件事(机会出现 / 成形 / 公开落地),职业类型与工作方式一律问用户。三件事用生活说法,不报层名和编号,不得合成一句「事业机会」:①机会出现(消息、邀约、初步接洽)②成形(合同、合作、长期项目)③公开落地(发布、到账、被看见、名声)。";
|
||||
|
||||
export const CONSULTATION_CONDENSED_CHECKLISTS: Readonly<Partial<Record<ConsultationDomain, readonly string[]>>> = {
|
||||
career: [
|
||||
CAREER_FIELD_ASK_RULE,
|
||||
CAREER_SIGNIFICATOR_STRENGTH_RULE,
|
||||
"必看:D1 10 宫与 10 宫主、AmK;卡上 d9 段确认 10 宫主与 AmK 的旺弱、Vargottama、D1/D9 反转;D10 上升与 10 宫;AL 与 A10;10 宫 SAV 与木星、土星所在宫的 SAV。",
|
||||
"时间:Vimshottari(MD/AD/PD)与 Narayana 双轨同向才谈应期,看大运主与 10 宫、10 宫主、D10 的关系;每颗星先按功能吉凶定性。",
|
||||
"分三件事看(回答里用生活说法,不报层名和编号),不得合成一句「事业机会」:①机会出现(消息、邀约、初步接洽)②成形(合同、合作、长期项目)③公开落地(发布、到账、被看见、名声)。",
|
||||
CAREER_DIRECTION_SECTION_RULE,
|
||||
"禁写:本命承诺弱时,不得因一段大运或一次行运断言「事业必成」;不得把接触窗写成落地窗。",
|
||||
"AL 落第几宫要说出来(说明名声和外界印象这条线有没有力量,不据此断定公众型或幕后型);10 宫主、AmK 的受冲按通用读法。",
|
||||
"卡外按需补取最多三次:Karakamsha、宫主链、D1→D9→D10 联动、Argala、Chara 大运;KP 精确宫头仍 blocked,只作参考、不作依据。",
|
||||
|
||||
@@ -13,6 +13,7 @@ import test from "node:test";
|
||||
import { fileURLToPath } from "node:url";
|
||||
|
||||
import {
|
||||
CAREER_DIRECTION_SECTION_RULE,
|
||||
CAREER_SIGNIFICATOR_STRENGTH_RULE,
|
||||
CONDENSED_CHECKLIST_SOURCE,
|
||||
CONSULT_BRIEF_CHECKLIST_SOURCE,
|
||||
@@ -96,6 +97,27 @@ test("career significators stay strength and affliction, and the system prompt i
|
||||
assert.equal(natal.length, 15_291);
|
||||
});
|
||||
|
||||
test("career direction section only writes strength, resistance, timing and three events", () => {
|
||||
const career = CONSULTATION_CONDENSED_CHECKLISTS.career!;
|
||||
assert.equal(career[0]!.startsWith("事业题只讲盘上的力量、阻力和时间"), true);
|
||||
assert.equal(career.filter((line) => line === CAREER_DIRECTION_SECTION_RULE).length, 1);
|
||||
assert.equal(career.filter((line) => line === CAREER_SIGNIFICATOR_STRENGTH_RULE).length, 1);
|
||||
assert.match(CAREER_DIRECTION_SECTION_RULE, /「方向」一节只写力量、阻力、时间与三件事(机会出现 \/ 成形 \/ 公开落地)/);
|
||||
assert.match(CAREER_DIRECTION_SECTION_RULE, /职业类型与工作方式一律问用户/);
|
||||
assert.match(CAREER_DIRECTION_SECTION_RULE, /机会出现(消息、邀约、初步接洽)/);
|
||||
assert.match(CAREER_DIRECTION_SECTION_RULE, /成形(合同、合作、长期项目)/);
|
||||
assert.match(CAREER_DIRECTION_SECTION_RULE, /公开落地(发布、到账、被看见、名声)/);
|
||||
assert.equal(career.length, 8);
|
||||
const mastra = readFileSync(new URL("../src/mastra/index.ts", import.meta.url), "utf8");
|
||||
const natal = mastra.slice(
|
||||
mastra.indexOf("const jyotishInstructions"),
|
||||
mastra.indexOf("const generalJyotishInstructions"),
|
||||
);
|
||||
assert.equal(natal.includes(CAREER_DIRECTION_SECTION_RULE), false);
|
||||
assert.equal(natal.includes("「方向」一节只写力量、阻力、时间"), false);
|
||||
assert.equal(natal.length, 15_291);
|
||||
});
|
||||
|
||||
test("each condensed list is 5–8 lines", () => {
|
||||
// 原值:全部清单 5–8 行,且只有 career / marriage / wealth 三份;
|
||||
// 新值:这三份仍 5–8 行;另有 TASK-consult-affliction-reading-20261001 定稿的 8 份(2–6 行,按任务书原文,不凑行数),
|
||||
@@ -120,7 +142,10 @@ test("the writer's prompt carries the condensed list with the three layers and t
|
||||
// 原值(2): career 必含「事业形态」
|
||||
// 新值(2): 改为必含「不替用户断定他做哪一行」(CAREER_FIELD_ASK_RULE)
|
||||
// 原因(2): 产品 2026-10-02「行业不猜,问用户」,删掉「先讲一生事业形态:公众型还是幕后型」(BUG-1176)
|
||||
career: ["接触", "机会出现", "成形", "公开落地", "不替用户断定他做哪一行", "事业必成", "d9 段", "KP 精确宫头仍 blocked", "不按星的本性推断职业类型"],
|
||||
// 原值: career 必含到「不按星的本性推断职业类型」为止
|
||||
// 新值: 另必含「方向」一节只写力量、阻力、时间与三件事,以及「职业类型与工作方式一律问用户」
|
||||
// 原因: TASK-consult-career-type-residual-20261009 H1。清单仍是 8 行,这句替换原来的「分三件事看」行,不进系统提示。
|
||||
career: ["接触", "机会出现", "成形", "公开落地", "不替用户断定他做哪一行", "事业必成", "d9 段", "KP 精确宫头仍 blocked", "不按星的本性推断职业类型", "「方向」一节只写力量、阻力、时间与三件事", "职业类型与工作方式一律问用户"],
|
||||
marriage: [
|
||||
// 原值: "心动接触", "关系成对", "社会法律落地"
|
||||
// 新值: 三件事用生活说法,并写明回答里不报层名
|
||||
|
||||
@@ -10,6 +10,7 @@ from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import threading
|
||||
import time
|
||||
from datetime import date
|
||||
from pathlib import Path
|
||||
@@ -23,6 +24,38 @@ KEEP_SECONDS = 2 * 24 * 60 * 60
|
||||
|
||||
PacketBuilder = Callable[[dict[str, Any]], dict[str, Any]]
|
||||
|
||||
# One build per cache key. A later caller waits and reads the stored packet.
|
||||
class _Flight:
|
||||
def __init__(self) -> None:
|
||||
self.done = threading.Event()
|
||||
self.error: BaseException | None = None
|
||||
|
||||
|
||||
_flights: dict[str, _Flight] = {}
|
||||
_flights_guard = threading.Lock()
|
||||
|
||||
# Chart-save warm has its own cap. When it is full the warm is skipped.
|
||||
# A user request does not take this slot and does not receive 429 from it.
|
||||
WARM_CONCURRENCY_LIMIT = 1
|
||||
_warm_guard = threading.Lock()
|
||||
_warm_inflight = 0
|
||||
|
||||
|
||||
def _acquire_warm_slot() -> bool:
|
||||
global _warm_inflight
|
||||
with _warm_guard:
|
||||
if _warm_inflight >= WARM_CONCURRENCY_LIMIT:
|
||||
return False
|
||||
_warm_inflight += 1
|
||||
return True
|
||||
|
||||
|
||||
def _release_warm_slot() -> None:
|
||||
global _warm_inflight
|
||||
with _warm_guard:
|
||||
if _warm_inflight > 0:
|
||||
_warm_inflight -= 1
|
||||
|
||||
|
||||
def cache_dir() -> Path:
|
||||
path = ROOT / "scratch" / "local" / "consult_full_data_cache"
|
||||
@@ -104,34 +137,63 @@ def prune_cache(now: float | None = None) -> int:
|
||||
return removed
|
||||
|
||||
|
||||
def _read_cached(path: Path, started: float, now: Callable[[], float]) -> dict[str, Any]:
|
||||
packet = json.loads(path.read_text(encoding="utf-8"))
|
||||
return {"packet": packet, "hit": True, "elapsed_s": now() - started, "slow": False}
|
||||
|
||||
|
||||
def get_or_build(
|
||||
identity: dict[str, Any],
|
||||
builder: PacketBuilder | None = None,
|
||||
*,
|
||||
clock: Callable[[], float] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Return the packet. A miss calls builder once and stores the packet."""
|
||||
"""Return the packet. Concurrent callers for one key share one build."""
|
||||
prune_cache()
|
||||
now = clock or time.perf_counter
|
||||
started = now()
|
||||
path = _cache_path(identity)
|
||||
if path.is_file():
|
||||
packet = json.loads(path.read_text(encoding="utf-8"))
|
||||
return _read_cached(path, started, now)
|
||||
key = cache_key(identity)
|
||||
with _flights_guard:
|
||||
flight = _flights.get(key)
|
||||
owner = flight is None
|
||||
if owner:
|
||||
flight = _Flight()
|
||||
_flights[key] = flight
|
||||
assert flight is not None
|
||||
if not owner:
|
||||
flight.done.wait()
|
||||
if flight.error is not None:
|
||||
raise flight.error
|
||||
if path.is_file():
|
||||
return _read_cached(path, started, now)
|
||||
raise RuntimeError("consult_full_data_unavailable")
|
||||
try:
|
||||
if path.is_file():
|
||||
return _read_cached(path, started, now)
|
||||
if builder is None:
|
||||
builder = build_full_data_packet
|
||||
packet = builder(identity)
|
||||
temporary = path.with_suffix(".json.tmp")
|
||||
temporary.write_text(json.dumps(packet, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
|
||||
temporary.replace(path)
|
||||
elapsed = now() - started
|
||||
return {"packet": packet, "hit": True, "elapsed_s": elapsed, "slow": False}
|
||||
if builder is None:
|
||||
builder = build_full_data_packet
|
||||
packet = builder(identity)
|
||||
temporary = path.with_suffix(".json.tmp")
|
||||
temporary.write_text(json.dumps(packet, ensure_ascii=False, separators=(",", ":")), encoding="utf-8")
|
||||
temporary.replace(path)
|
||||
elapsed = now() - started
|
||||
return {
|
||||
"packet": packet,
|
||||
"hit": False,
|
||||
"elapsed_s": elapsed,
|
||||
"slow": elapsed > SLOW_MISS_SECONDS,
|
||||
}
|
||||
return {
|
||||
"packet": packet,
|
||||
"hit": False,
|
||||
"elapsed_s": elapsed,
|
||||
"slow": elapsed > SLOW_MISS_SECONDS,
|
||||
}
|
||||
except Exception as exc:
|
||||
flight.error = exc
|
||||
raise
|
||||
finally:
|
||||
flight.done.set()
|
||||
with _flights_guard:
|
||||
if _flights.get(key) is flight:
|
||||
_flights.pop(key, None)
|
||||
|
||||
|
||||
def _attach_graha_drishti(packet: dict[str, Any], reading: dict[str, Any]) -> dict[str, Any]:
|
||||
@@ -186,22 +248,39 @@ def build_full_data_packet(identity: dict[str, Any]) -> dict[str, Any]:
|
||||
|
||||
|
||||
def warm_consult_packet(body: dict[str, Any] | None) -> dict[str, Any]:
|
||||
"""Build today's packet after a chart save. The response stays small."""
|
||||
try:
|
||||
cached = get_or_build(identity_from_body(dict(body or {})))
|
||||
except Exception as exc:
|
||||
"""Build today's packet after a chart save. The response stays small.
|
||||
|
||||
At most WARM_CONCURRENCY_LIMIT warms run at once. A full cap skips this
|
||||
warm. The skip is not a user-request 429, and a failure here does not
|
||||
fail the chart save (the save route does not wait on this response).
|
||||
"""
|
||||
if not _acquire_warm_slot():
|
||||
return {
|
||||
"success": False,
|
||||
"success": True,
|
||||
"endpoint": "consult_card_warm",
|
||||
"error_type": type(exc).__name__,
|
||||
"skipped": True,
|
||||
"hit": False,
|
||||
"elapsed_s": 0.0,
|
||||
"slow": False,
|
||||
}
|
||||
return {
|
||||
"success": True,
|
||||
"endpoint": "consult_card_warm",
|
||||
"hit": bool(cached["hit"]),
|
||||
"elapsed_s": cached["elapsed_s"],
|
||||
"slow": bool(cached["slow"]),
|
||||
}
|
||||
try:
|
||||
try:
|
||||
cached = get_or_build(identity_from_body(dict(body or {})))
|
||||
except Exception as exc:
|
||||
return {
|
||||
"success": False,
|
||||
"endpoint": "consult_card_warm",
|
||||
"error_type": type(exc).__name__,
|
||||
}
|
||||
return {
|
||||
"success": True,
|
||||
"endpoint": "consult_card_warm",
|
||||
"hit": bool(cached["hit"]),
|
||||
"elapsed_s": cached["elapsed_s"],
|
||||
"slow": bool(cached["slow"]),
|
||||
}
|
||||
finally:
|
||||
_release_warm_slot()
|
||||
|
||||
|
||||
def _requested_domains(body: dict[str, Any]) -> list[str]:
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
@@ -302,3 +303,125 @@ def test_warm_response_does_not_include_the_packet():
|
||||
assert warmed["success"] is False
|
||||
assert "packet" not in warmed
|
||||
assert warmed["error_type"]
|
||||
|
||||
|
||||
def _fictional_body(today: str, minute: int) -> dict:
|
||||
return {
|
||||
"year": 1902,
|
||||
"month": 3,
|
||||
"day": 3,
|
||||
"hour": 1,
|
||||
"minute": minute,
|
||||
"second": 0,
|
||||
"lat": 1,
|
||||
"lon": 2,
|
||||
"tz": 0,
|
||||
"ayanamsa": "raman",
|
||||
"node_mode": "mean",
|
||||
"today": today,
|
||||
}
|
||||
|
||||
|
||||
def _drop_cache(identity: dict) -> None:
|
||||
path = cache_dir() / f"{cache_key(identity)}.json"
|
||||
path.unlink(missing_ok=True)
|
||||
path.with_suffix(".json.tmp").unlink(missing_ok=True)
|
||||
|
||||
|
||||
def test_one_warm_at_a_time_and_same_key_user_waits_for_that_build(monkeypatch):
|
||||
"""H2: warms share a cap of 1; a question for the key being warmed builds once."""
|
||||
calls = {"n": 0}
|
||||
calls_lock = threading.Lock()
|
||||
entered = threading.Event()
|
||||
release = threading.Event()
|
||||
|
||||
def blocking(identity):
|
||||
with calls_lock:
|
||||
calls["n"] += 1
|
||||
entered.set()
|
||||
assert release.wait(timeout=5)
|
||||
return {"marker": identity["reference_date"]}
|
||||
|
||||
monkeypatch.setattr("consult_full_data_cache.build_full_data_packet", blocking)
|
||||
|
||||
bodies = [_fictional_body(f"1904-01-0{index}", 11) for index in (1, 2, 3)]
|
||||
identities = [identity_from_body(body) for body in bodies]
|
||||
for identity in identities:
|
||||
_drop_cache(identity)
|
||||
results: list[dict | None] = [None, None, None]
|
||||
|
||||
def run_warm(index: int) -> None:
|
||||
results[index] = warm_consult_packet(bodies[index])
|
||||
|
||||
threads = [threading.Thread(target=run_warm, args=(index,)) for index in range(3)]
|
||||
try:
|
||||
for thread in threads:
|
||||
thread.start()
|
||||
assert entered.wait(timeout=5)
|
||||
deadline = time.time() + 3
|
||||
while time.time() < deadline and sum(item is not None for item in results) < 2:
|
||||
time.sleep(0.02)
|
||||
assert calls["n"] == 1
|
||||
assert sum(item is not None and item.get("skipped") is True for item in results) == 2
|
||||
for item in results:
|
||||
if item is not None:
|
||||
assert "packet" not in item
|
||||
|
||||
same = _fictional_body("1904-02-02", 12)
|
||||
same_identity = identity_from_body(same)
|
||||
_drop_cache(same_identity)
|
||||
# The cap is still held by the first warm. A same-key question must join
|
||||
# that pattern on its own key: start a warm, then the question.
|
||||
release.set()
|
||||
for thread in threads:
|
||||
thread.join(timeout=5)
|
||||
assert calls["n"] == 1
|
||||
built = [item for item in results if item and item.get("skipped") is not True]
|
||||
assert len(built) == 1
|
||||
assert built[0]["success"] is True
|
||||
assert built[0]["hit"] is False
|
||||
finally:
|
||||
release.set()
|
||||
for thread in threads:
|
||||
thread.join(timeout=5)
|
||||
for identity in identities:
|
||||
_drop_cache(identity)
|
||||
|
||||
entered.clear()
|
||||
release.clear()
|
||||
calls["n"] = 0
|
||||
warm_box: dict = {}
|
||||
user_box: dict = {}
|
||||
same = _fictional_body("1904-02-02", 12)
|
||||
same_identity = identity_from_body(same)
|
||||
_drop_cache(same_identity)
|
||||
|
||||
def run_same_warm() -> None:
|
||||
warm_box["result"] = warm_consult_packet(same)
|
||||
|
||||
def run_user() -> None:
|
||||
user_box["result"] = get_or_build(same_identity)
|
||||
|
||||
warm_thread = threading.Thread(target=run_same_warm)
|
||||
user_thread = threading.Thread(target=run_user)
|
||||
try:
|
||||
warm_thread.start()
|
||||
assert entered.wait(timeout=5)
|
||||
user_thread.start()
|
||||
time.sleep(0.3)
|
||||
assert calls["n"] == 1
|
||||
assert user_thread.is_alive()
|
||||
release.set()
|
||||
warm_thread.join(timeout=5)
|
||||
user_thread.join(timeout=5)
|
||||
assert calls["n"] == 1
|
||||
assert warm_box["result"]["success"] is True
|
||||
assert warm_box["result"].get("skipped") is not True
|
||||
assert "packet" not in warm_box["result"]
|
||||
assert user_box["result"]["hit"] is True
|
||||
assert user_box["result"]["packet"]["marker"] == "1904-02-02"
|
||||
finally:
|
||||
release.set()
|
||||
warm_thread.join(timeout=5)
|
||||
user_thread.join(timeout=5)
|
||||
_drop_cache(same_identity)
|
||||
|
||||
Reference in New Issue
Block a user