From e0dbd2ba8bdc22daa87f635576c6fc51d0127ad6 Mon Sep 17 00:00:00 2001 From: "macstudio-4[bot]" Date: Thu, 20 Aug 2026 04:59:17 +0800 Subject: [PATCH 1/2] fix: a stored codex failure is not adoptable by a later dispatch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two intake codex runs timed out against an unavailable upstream. Their failures were written to the codex-adoption records, and every later dispatch for the same prompt hash adopted the stored failure instead of running codex. `workflow_select` then dead-lettered continuously, reporting a 3600s timeout that took five seconds, and no retry could ever clear it. The engine already treats every codex failure as environmental: `boundary_resource::class_for_adapter_failure` has no deterministic-failure class and defaults to `provider-unavailable`. A stored failure therefore describes the environment at the time of that run, not a property of the prompt, and must not answer a dispatch that did not join the run producing it. - `run_adoptable_codex_request` reuses a completed record only when the run succeeded. `wait_for_adoption_result` is unchanged: a waiter attached to a run in flight is entitled to that run's failure, and rejecting it there would make waiters hang. - The outcome predicate is the exit code, not `error_kind`. `read_completed_adoption_result` builds records without `error_kind` through `CodexResult::success` even when the process exited non-zero. - `start_adoption_worker` discards the previous run's `result.json`, `effect.json`, and `effect-receipts.log` before publishing its intent. Those three are the recovery paths; without this a waiter on the fresh run promotes the previous failure straight back and the fix does nothing. `spawn_codex_sync_preserves_typed_jsonl_refusal_through_adoption` asserted that two dispatches spawn codex once. Its fixture emits `turn.failed` / "server overloaded" — the environmental shape this change is about — so that incidental assertion encoded the reported defect. The test's actual subject, typed metadata surviving the adoption machinery, is unchanged. Refs ChronoAIProject/fkst-packages#3919 Co-Authored-By: Claude Opus 5 --- SPEC.md | 2 +- crates/fkst-framework/src/sdk_codex.rs | 70 ++++++++++++++++++++- crates/fkst-framework/tests/sdk_codex.rs | 77 +++++++++++++++++++++++- 3 files changed, 144 insertions(+), 5 deletions(-) diff --git a/SPEC.md b/SPEC.md index 3ba62f3..b5f0438 100644 --- a/SPEC.md +++ b/SPEC.md @@ -113,7 +113,7 @@ The engine guarantees no graph reload, no atomic engine/package version pair, an - Department 默认以可靠方式消费队列;`M.spec.ephemeral = {"queue"}` 可将本 Department 对指定 consumed queue 的订阅降级为非可靠。`M.spec.retry = false` 只表示失败不重试;`M.spec.retry = { ... }` 可覆盖 `max_attempts`、`base`、`cap` 的任意子集,缺失字段从全局默认补齐。`M.spec.stall_window` 是可靠投递 lease 与续租窗口,不是 framework child 无输出 kill deadline。可靠订阅启动必须有 `FKST_DURABLE_ROOT`,缺失 fail-closed。可靠 source event 必须带 `SourceRef{kind,reference}`;cron 由 raiser 名派生,file_watch 由绝对路径派生,Department `RAISED` 进入可靠 queue 时继承上游 source_ref,缺失则 publish fail-closed 且上游 delivery 不 ack。可靠 consumer 由 Fanout wake + 定时 tick 调用 redb store `lease`,spawn framework 后仅在 exit 0 且 RAISED 批次中所有未被终态 tombstone 抑制的 publish 成功时 `ack`;终态 tombstone 抑制项继续该批次且不阻止 `ack`,非零退出、codex timeout、spawn error 或其余 RAISED publish 失败调用 `retry`,到 max attempts 写 redb dead 表并 best-effort publish `dead_letter`。当前 delivery 来自 `dead_letter` 时抑制再次发送 `dead_letter`。该机制不是新 source kind,不提供 exactly-once;语义是 at-least-once-until-ack,`Fanout::send` 在可靠路径只作进程内唤醒。 - An enqueue against an exact delivery identity retained by a terminal tombstone returns `TerminallySuppressed`. The router's `publish` compatibility entry point fails with `delivery terminally suppressed`; `publish_outcome` instead returns `PublishOutcome::TerminallySuppressed`, which Department RAISED batch publishing skips before continuing the batch. For the suppressed delivery, the router emits `event=delivery_suppressed reason=terminal_tombstone` and does not emit `reason=durable_enqueue` or wake that subscriber. A fresh delivery identity remains admissible. - If a pending durable delivery has no current subscriber, supervise persists `subscriber_absent_since_ms`; subscriber return before `FKST_SUBSCRIBER_ABSENT_DELIVERY_BUDGET` clears that clock and resumes delivery, while continuous absence past the budget moves the row to the replayable DLQ with `error_excerpt = "subscriber-absent"`. -- `spawn_codex_sync` and `spawn_codex` accept an optional `sandbox` value from the closed set `read-only`, `workspace-write`, and `danger-full-access`; any other value fails closed with `codex-sandbox-invalid`. An omitted `sandbox` preserves the byte-identical `codex exec --dangerously-bypass-approvals-and-sandbox` default required by implement and fix work. An explicit value emits `codex exec --sandbox ` without the bypass flag. The `read-only` mode denies model-generated filesystem writes and network access; prompt instructions are not a substitute for this structural boundary. Both calls also accept `timeout`, defaulting to 3600 seconds as the overall codex subprocess wall-clock cap; stdout/stderr activity is captured but does not extend it. A `spawn_codex` handle can only be joined by `await_all`; use `await_all({handle})` for one handle. First-result fanout and sleep timers are not part of the fixed Lua SDK surface. +- `spawn_codex_sync` and `spawn_codex` accept an optional `sandbox` value from the closed set `read-only`, `workspace-write`, and `danger-full-access`; any other value fails closed with `codex-sandbox-invalid`. An omitted `sandbox` preserves the byte-identical `codex exec --dangerously-bypass-approvals-and-sandbox` default required by implement and fix work. An explicit value emits `codex exec --sandbox ` without the bypass flag. The `read-only` mode denies model-generated filesystem writes and network access; prompt instructions are not a substitute for this structural boundary. Both calls also accept `timeout`, defaulting to 3600 seconds as the overall codex subprocess wall-clock cap; stdout/stderr activity is captured but does not extend it. A `spawn_codex` handle can only be joined by `await_all`; use `await_all({handle})` for one handle. First-result fanout and sleep timers are not part of the fixed Lua SDK surface. A completed codex-adoption record is reusable by a later dispatch only when that run succeeded; a stored failure is an environmental boundary condition rather than a property of the prompt, so it is never replayed to a dispatch that did not join the run that produced it, and starting a run discards the previous run's outcome artifacts. - `fkst.codex_runs()` is a read-only bounded observability surface over engine codex run records. It returns running and recent codex runs with `role`, `started_at`, `started_at_ms`, `timeout_seconds`, `lease_expires_at`, `lease_expires_at_ms`, `status` (`running` / `done` / `failed`), bounded `output_tail`, and optional `exit_code`; it does not expose runtime paths or unbounded stdout/stderr. The lease deadline is the codex run's own wall-clock timeout contract, derived as `started_at_ms + timeout_seconds * 1000`. `recent` is the last `50` completed codex runs by record time for bounded observability, not a time-bounded just-finished handoff window. - `fkst.observe([opts]) -> table` is the in-process adapter for generic durable delivery facts. Snapshot mode is the same projection emitted by `fkst-framework observe --json`; its `opts` may contain only `limit` (1..10000), `include = {"queues","errors","events","entities"}`, `since = `, and `page = { section = "dead_letters", after? = }`. Lineage mode is `lineage = { queue = , dept = , source_ref = { kind = , ref = } }` and cannot be combined with snapshot options. The lineage table and nested source reference reject unknown keys and empty required strings; supported kinds are `file`, `file_watch`, `cron`, `git`, and `external`. It returns only `live_delivery` and `terminal_dead_letter` when present, has no `limit` or truncation fields, and selects the earliest live delivery plus latest terminal tombstone in one read transaction through the ordered `(queue, dept, source_ref, record-kind, time, delivery-id)` index. Unsupported or business-shaped keys fail closed. In snapshot mode, `page` and `since` cannot be combined. Without `page`, `since` is applied to `delivery_id` before `limit` truncation, independently for live deliveries and DLQ entries, and request/result serialization remains byte-identical. With `page`, DLQ selection uses durable keyset order `(dead_at_ms, delivery_id)` and a `limit + 1` read; the snapshot adds `page = { section = "dead_letters", next? = }`, where `next` exists only when the lookahead proves another retained row exists. The cursor is versioned and bound to `"dead_letters"`; malformed encoding, shape, version, section, or delivery identity fails closed with `observe dead-letter cursor invalid`. Each call is one read transaction over the current mutable DLQ. Continuation traverses retained rows under that mutable-history model and does not provide cross-request snapshot isolation. The current snapshot contains queue depth state, queue `subscriber_status` (`"current"`, `"absent"`, or `"unknown"`), live delivery entries, DLQ entries, source metadata, limits, and truncation flags; it never emits full payload bodies and does not encode idle, board, audit, skip, workflow, or package-specific semantics. Live owner-process observe may report `"current"` or `"absent"` from the loaded graph; offline database observe reports `"unknown"` because durable rows are not current graph authority. Packages derive those meanings outside the engine. Production reads `FKST_DURABLE_ROOT` and uses the same live-socket/offline projection core as the CLI. - 边界资源必须静态枚举并经 adapter mediation 访问。当前 engine registry 锚点是 `fkst-framework boundary-resources` 与 `crates/fkst-framework/src/boundary_resource.rs`,覆盖 `codex.process`、`shell.process`、`git.process`、`runtime.filesystem` 与 `wall-clock`。可分类的 adapter failure 使用 `error_class`,值域为 `quota-exhausted`、`auth-degraded`、`provider-unavailable`、`provider-throttle`;package 不应从 stderr 文本反推边界状态。 diff --git a/crates/fkst-framework/src/sdk_codex.rs b/crates/fkst-framework/src/sdk_codex.rs index 4ff921d..3585771 100644 --- a/crates/fkst-framework/src/sdk_codex.rs +++ b/crates/fkst-framework/src/sdk_codex.rs @@ -937,7 +937,7 @@ fn run_adoptable_codex_request( host_root: &Path, config: &ConfigContext, ) -> Result { - if let Some(result) = read_completed_adoption_result(&paths)? { + if let Some(result) = read_reusable_adoption_result(&paths)? { return Ok(result); } std::fs::create_dir_all(&paths.dir).map_err(mlua::Error::external)?; @@ -950,11 +950,11 @@ fn run_adoptable_codex_request( .open(lock_path) .map_err(mlua::Error::external)?; flock(lock_file.as_raw_fd(), FlockArg::LockExclusive).map_err(mlua::Error::external)?; - if let Some(result) = read_completed_adoption_result(&paths)? { + if let Some(result) = read_reusable_adoption_result(&paths)? { drop(lock_file); return Ok(result); } - if let Some(result) = recover_completed_adoption_result_locked(&paths)? { + if let Some(result) = recover_reusable_adoption_result_locked(&paths)? { drop(lock_file); return Ok(result); } @@ -1173,6 +1173,47 @@ fn read_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result Result> { + let Some(result) = read_completed_adoption_result(paths)? else { + return Ok(None); + }; + if !adoption_result_succeeded(&result) { + return Ok(None); + } + Ok(Some(result)) +} + +/// A stored record without `error_kind` is still a failure when the process +/// exited non-zero: `read_completed_adoption_result` builds those through +/// `CodexResult::success`, so the exit code is the only reliable outcome signal. +fn adoption_result_succeeded(result: &CodexResult) -> bool { + result.error_kind.is_none() && result.exit_code == 0 +} + +fn recover_reusable_adoption_result_locked( + paths: &CodexAdoptionPaths, +) -> Result> { + let Some(result) = recover_completed_adoption_result_locked(paths)? else { + return Ok(None); + }; + if !adoption_result_succeeded(&result) { + return Ok(None); + } + Ok(Some(result)) +} + fn recover_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result> { let lock_file = open_adoption_run_lock(&paths.dir).map_err(mlua::Error::external)?; flock(lock_file.as_raw_fd(), FlockArg::LockExclusive).map_err(mlua::Error::external)?; @@ -1556,6 +1597,7 @@ fn start_adoption_worker( ) -> Result<()> { ensure_pool_with_context(host_root, config)?; std::fs::create_dir_all(&paths.dir).map_err(mlua::Error::external)?; + discard_previous_run_outcome(paths).map_err(mlua::Error::external)?; write_atomic(&paths.prompt, request.prompt.as_bytes()).map_err(mlua::Error::external)?; let intent = adoption_record_from_request(request, paths, CODEX_ADOPTION_STATUS_INTENT, None, None); @@ -1581,6 +1623,28 @@ fn start_adoption_worker( Ok(()) } +/// Starting a run invalidates the previous run's outcome artifacts. +/// +/// `status.json` is overwritten with the new intent, but the three recovery +/// paths — `result.json`, the effect record, and the effect log — are otherwise +/// never cleared, so a waiter on the fresh run would promote the previous run's +/// stored result straight back through `recover_completed_adoption_result`. +fn discard_previous_run_outcome(paths: &CodexAdoptionPaths) -> anyhow::Result<()> { + for path in [&paths.result, &paths.effect, &paths.effect_log] { + match std::fs::remove_file(path) { + Ok(()) => {} + Err(err) if err.kind() == std::io::ErrorKind::NotFound => {} + Err(err) => { + return Err(anyhow::anyhow!( + "codex adoption could not discard previous run outcome {}: {err}", + path.display() + )) + } + } + } + Ok(()) +} + fn adoption_record_from_request( request: &CodexRequest, paths: &CodexAdoptionPaths, diff --git a/crates/fkst-framework/tests/sdk_codex.rs b/crates/fkst-framework/tests/sdk_codex.rs index 03e4957..e48c504 100644 --- a/crates/fkst-framework/tests/sdk_codex.rs +++ b/crates/fkst-framework/tests/sdk_codex.rs @@ -1000,6 +1000,77 @@ printf '{"type":"item.completed","item":{"type":"agent_message","text":"result-% assert!(!worktree.join(".fkst-codex").exists()); } +#[cfg(unix)] +#[test] +fn spawn_codex_sync_does_not_adopt_a_failed_result_for_a_later_dispatch() { + // A stored codex failure describes the environment at the time of that run, + // not a property of the prompt: every codex failure is classified as a + // boundary condition. Adopting one for a later dispatch turns a single + // transient upstream fault into a permanent answer that no retry can clear. + let tmp = tempfile::tempdir().unwrap(); + let bin_dir = tmp.path().join("bin"); + let worktree = tmp.path().join("wt"); + let capture_dir = tmp.path().join("capture"); + std::fs::create_dir_all(&worktree).unwrap(); + std::fs::create_dir_all(&capture_dir).unwrap(); + install_codex_script( + &bin_dir, + r#"#!/bin/sh +cat >/dev/null +count_file="$CAPTURE_DIR/spawns" +count=0 +if [ -f "$count_file" ]; then + count=$(cat "$count_file") +fi +count=$((count + 1)) +printf '%s' "$count" > "$count_file" +if [ "$count" = "1" ]; then + echo "upstream refused the request: 503 service unavailable" >&2 + exit 1 +fi +printf '{"type":"item.completed","item":{"type":"agent_message","text":"result-%s"}}\n' "$count" +"#, + ); + + let mut sandbox = ProcessSandbox::new(); + sandbox.enter_cwd(tmp.path()).runtime_root(".fkst/runtime"); + sandbox.prepend_path(&bin_dir); + sandbox.set_env("CAPTURE_DIR", capture_dir.to_string_lossy().into_owned()); + sandbox.set_env(CODEX_WORKER_BIN_ENV, framework_bin()); + sandbox.runtime_log_dir(tmp.path().join("runtime")); + let (_lock, _guard) = sandbox.enter(); + + let lua = Lua::new(); + register(&lua).unwrap(); + let spawn: mlua::Function = lua.globals().get("spawn_codex_sync").unwrap(); + + let first_opts = lua_opts(&lua, "same-work"); + first_opts + .set("worktree", worktree.to_string_lossy().into_owned()) + .unwrap(); + first_opts.set("dedup_key", "same-dedup").unwrap(); + let first: Table = spawn.call(first_opts).unwrap(); + assert_ne!(first.get::("exit_code").unwrap(), 0); + + let second_opts = lua_opts(&lua, "same-work"); + second_opts + .set("worktree", worktree.to_string_lossy().into_owned()) + .unwrap(); + second_opts.set("dedup_key", "same-dedup").unwrap(); + let second: Table = spawn.call(second_opts).unwrap(); + assert_eq!( + second.get::("exit_code").unwrap(), + 0, + "the second dispatch must run codex again instead of replaying the stored failure" + ); + assert_eq!(second.get::("stdout").unwrap(), "result-2"); + assert_eq!( + std::fs::read_to_string(capture_dir.join("spawns")).unwrap(), + "2", + "a stored failure must not suppress the respawn" + ); +} + #[cfg(unix)] #[test] fn spawn_codex_sync_uses_runtime_adoption_dir_for_read_only_worktree() { @@ -3204,6 +3275,10 @@ esac first.get::("codex_error_info").unwrap(), "server_overloaded" ); + // The typed refusal this fixture emits is `turn.failed` / "server overloaded", + // an environmental condition. A later dispatch runs codex again rather than + // replaying it; the typed metadata this test is about still survives the + // adoption machinery on every run. let second: Table = spawn.call(opts).unwrap(); assert_eq!(second.get::("exit_code").unwrap(), 1); assert_eq!( @@ -3212,7 +3287,7 @@ esac ); assert_eq!( std::fs::read_to_string(capture_dir.join("spawns")).unwrap(), - "1" + "2" ); let adoption = adoption_work_dir(tmp.path()); From 5e7498c1ee3612b23cb1819bba72ed5be7f55451 Mon Sep 17 00:00:00 2001 From: "macstudio-4[bot]" Date: Thu, 20 Aug 2026 16:53:51 +0800 Subject: [PATCH 2/2] fix: a codex result says whether it was produced or adopted The devloop refused ChronoAIProject/fkst-packages#3919 as `wrong-layer` and named the gap in the first commit on this branch: it replaces completed failures, but `CodexResult`/`into_lua_table` still expose neither `run_id` nor produced/adopted provenance, so no contract-complete substrate revision exists for packages to pin. #3919 asks for both halves. This is the second: an adopted five-second read and a real hour-long run reported the same thing, which is why identifying the defect took a measurement of `ELAPSED_MS` rather than a reading of the message. - `CodexResult` carries `provenance` (`produced` / `adopted`) and the `run_id` of the run that produced the record; `into_lua_table` exposes both. - Provenance is decided by comparing the stored `run_id` against the requesting run's, not by which code path returned the result. In the adoption protocol the producing call also reads its own outcome back through `read_completed_adoption_result`, so a path-based rule marks a producer `adopted`. The added test caught exactly that. - The requesting run id is threaded through the read and recover paths. Refs ChronoAIProject/fkst-packages#3919 Co-Authored-By: Claude Opus 5 --- SPEC.md | 2 +- crates/fkst-framework/src/sdk_codex.rs | 77 ++++++++++++++++++++---- crates/fkst-framework/tests/sdk_codex.rs | 70 +++++++++++++++++++++ 3 files changed, 136 insertions(+), 13 deletions(-) diff --git a/SPEC.md b/SPEC.md index b5f0438..05f3f2c 100644 --- a/SPEC.md +++ b/SPEC.md @@ -113,7 +113,7 @@ The engine guarantees no graph reload, no atomic engine/package version pair, an - Department 默认以可靠方式消费队列;`M.spec.ephemeral = {"queue"}` 可将本 Department 对指定 consumed queue 的订阅降级为非可靠。`M.spec.retry = false` 只表示失败不重试;`M.spec.retry = { ... }` 可覆盖 `max_attempts`、`base`、`cap` 的任意子集,缺失字段从全局默认补齐。`M.spec.stall_window` 是可靠投递 lease 与续租窗口,不是 framework child 无输出 kill deadline。可靠订阅启动必须有 `FKST_DURABLE_ROOT`,缺失 fail-closed。可靠 source event 必须带 `SourceRef{kind,reference}`;cron 由 raiser 名派生,file_watch 由绝对路径派生,Department `RAISED` 进入可靠 queue 时继承上游 source_ref,缺失则 publish fail-closed 且上游 delivery 不 ack。可靠 consumer 由 Fanout wake + 定时 tick 调用 redb store `lease`,spawn framework 后仅在 exit 0 且 RAISED 批次中所有未被终态 tombstone 抑制的 publish 成功时 `ack`;终态 tombstone 抑制项继续该批次且不阻止 `ack`,非零退出、codex timeout、spawn error 或其余 RAISED publish 失败调用 `retry`,到 max attempts 写 redb dead 表并 best-effort publish `dead_letter`。当前 delivery 来自 `dead_letter` 时抑制再次发送 `dead_letter`。该机制不是新 source kind,不提供 exactly-once;语义是 at-least-once-until-ack,`Fanout::send` 在可靠路径只作进程内唤醒。 - An enqueue against an exact delivery identity retained by a terminal tombstone returns `TerminallySuppressed`. The router's `publish` compatibility entry point fails with `delivery terminally suppressed`; `publish_outcome` instead returns `PublishOutcome::TerminallySuppressed`, which Department RAISED batch publishing skips before continuing the batch. For the suppressed delivery, the router emits `event=delivery_suppressed reason=terminal_tombstone` and does not emit `reason=durable_enqueue` or wake that subscriber. A fresh delivery identity remains admissible. - If a pending durable delivery has no current subscriber, supervise persists `subscriber_absent_since_ms`; subscriber return before `FKST_SUBSCRIBER_ABSENT_DELIVERY_BUDGET` clears that clock and resumes delivery, while continuous absence past the budget moves the row to the replayable DLQ with `error_excerpt = "subscriber-absent"`. -- `spawn_codex_sync` and `spawn_codex` accept an optional `sandbox` value from the closed set `read-only`, `workspace-write`, and `danger-full-access`; any other value fails closed with `codex-sandbox-invalid`. An omitted `sandbox` preserves the byte-identical `codex exec --dangerously-bypass-approvals-and-sandbox` default required by implement and fix work. An explicit value emits `codex exec --sandbox ` without the bypass flag. The `read-only` mode denies model-generated filesystem writes and network access; prompt instructions are not a substitute for this structural boundary. Both calls also accept `timeout`, defaulting to 3600 seconds as the overall codex subprocess wall-clock cap; stdout/stderr activity is captured but does not extend it. A `spawn_codex` handle can only be joined by `await_all`; use `await_all({handle})` for one handle. First-result fanout and sleep timers are not part of the fixed Lua SDK surface. A completed codex-adoption record is reusable by a later dispatch only when that run succeeded; a stored failure is an environmental boundary condition rather than a property of the prompt, so it is never replayed to a dispatch that did not join the run that produced it, and starting a run discards the previous run's outcome artifacts. +- `spawn_codex_sync` and `spawn_codex` accept an optional `sandbox` value from the closed set `read-only`, `workspace-write`, and `danger-full-access`; any other value fails closed with `codex-sandbox-invalid`. An omitted `sandbox` preserves the byte-identical `codex exec --dangerously-bypass-approvals-and-sandbox` default required by implement and fix work. An explicit value emits `codex exec --sandbox ` without the bypass flag. The `read-only` mode denies model-generated filesystem writes and network access; prompt instructions are not a substitute for this structural boundary. Both calls also accept `timeout`, defaulting to 3600 seconds as the overall codex subprocess wall-clock cap; stdout/stderr activity is captured but does not extend it. A `spawn_codex` handle can only be joined by `await_all`; use `await_all({handle})` for one handle. First-result fanout and sleep timers are not part of the fixed Lua SDK surface. A completed codex-adoption record is reusable by a later dispatch only when that run succeeded; a stored failure is an environmental boundary condition rather than a property of the prompt, so it is never replayed to a dispatch that did not join the run that produced it, and starting a run discards the previous run's outcome artifacts. Every codex result carries `provenance` (`produced` when this call's own run produced it, `adopted` when it read another run's completed record) and, when known, the `run_id` of the run that produced it, so an adopted read is distinguishable from a run that just happened. - `fkst.codex_runs()` is a read-only bounded observability surface over engine codex run records. It returns running and recent codex runs with `role`, `started_at`, `started_at_ms`, `timeout_seconds`, `lease_expires_at`, `lease_expires_at_ms`, `status` (`running` / `done` / `failed`), bounded `output_tail`, and optional `exit_code`; it does not expose runtime paths or unbounded stdout/stderr. The lease deadline is the codex run's own wall-clock timeout contract, derived as `started_at_ms + timeout_seconds * 1000`. `recent` is the last `50` completed codex runs by record time for bounded observability, not a time-bounded just-finished handoff window. - `fkst.observe([opts]) -> table` is the in-process adapter for generic durable delivery facts. Snapshot mode is the same projection emitted by `fkst-framework observe --json`; its `opts` may contain only `limit` (1..10000), `include = {"queues","errors","events","entities"}`, `since = `, and `page = { section = "dead_letters", after? = }`. Lineage mode is `lineage = { queue = , dept = , source_ref = { kind = , ref = } }` and cannot be combined with snapshot options. The lineage table and nested source reference reject unknown keys and empty required strings; supported kinds are `file`, `file_watch`, `cron`, `git`, and `external`. It returns only `live_delivery` and `terminal_dead_letter` when present, has no `limit` or truncation fields, and selects the earliest live delivery plus latest terminal tombstone in one read transaction through the ordered `(queue, dept, source_ref, record-kind, time, delivery-id)` index. Unsupported or business-shaped keys fail closed. In snapshot mode, `page` and `since` cannot be combined. Without `page`, `since` is applied to `delivery_id` before `limit` truncation, independently for live deliveries and DLQ entries, and request/result serialization remains byte-identical. With `page`, DLQ selection uses durable keyset order `(dead_at_ms, delivery_id)` and a `limit + 1` read; the snapshot adds `page = { section = "dead_letters", next? = }`, where `next` exists only when the lookahead proves another retained row exists. The cursor is versioned and bound to `"dead_letters"`; malformed encoding, shape, version, section, or delivery identity fails closed with `observe dead-letter cursor invalid`. Each call is one read transaction over the current mutable DLQ. Continuation traverses retained rows under that mutable-history model and does not provide cross-request snapshot isolation. The current snapshot contains queue depth state, queue `subscriber_status` (`"current"`, `"absent"`, or `"unknown"`), live delivery entries, DLQ entries, source metadata, limits, and truncation flags; it never emits full payload bodies and does not encode idle, board, audit, skip, workflow, or package-specific semantics. Live owner-process observe may report `"current"` or `"absent"` from the loaded graph; offline database observe reports `"unknown"` because durable rows are not current graph authority. Packages derive those meanings outside the engine. Production reads `FKST_DURABLE_ROOT` and uses the same live-socket/offline projection core as the CLI. - 边界资源必须静态枚举并经 adapter mediation 访问。当前 engine registry 锚点是 `fkst-framework boundary-resources` 与 `crates/fkst-framework/src/boundary_resource.rs`,覆盖 `codex.process`、`shell.process`、`git.process`、`runtime.filesystem` 与 `wall-clock`。可分类的 adapter failure 使用 `error_class`,值域为 `quota-exhausted`、`auth-degraded`、`provider-unavailable`、`provider-throttle`;package 不应从 stderr 文本反推边界状态。 diff --git a/crates/fkst-framework/src/sdk_codex.rs b/crates/fkst-framework/src/sdk_codex.rs index 3585771..15c093f 100644 --- a/crates/fkst-framework/src/sdk_codex.rs +++ b/crates/fkst-framework/src/sdk_codex.rs @@ -233,6 +233,28 @@ pub(crate) struct CodexResult { error: Option, codex_error_info: Option, diagnostics: Vec, + /// Identity of the codex run that produced this result. A caller that adopted + /// another run's record needs it to tell whose outcome it is holding. + run_id: Option, + /// Whether this call ran codex (`produced`) or read another run's completed + /// record (`adopted`). Without it, an adopted five-second read and a real + /// hour-long run report the same thing. + provenance: CodexProvenance, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum CodexProvenance { + Produced, + Adopted, +} + +impl CodexProvenance { + const fn label(self) -> &'static str { + match self { + Self::Produced => "produced", + Self::Adopted => "adopted", + } + } } // spawn_codex returns a pipeline-local opaque handle consumed by await_all. @@ -319,6 +341,8 @@ impl CodexResult { error: None, codex_error_info: None, diagnostics: Vec::new(), + run_id: None, + provenance: CodexProvenance::Produced, } } @@ -376,6 +400,8 @@ impl CodexResult { error: Some(message), codex_error_info: None, diagnostics: Vec::new(), + run_id: None, + provenance: CodexProvenance::Produced, } } @@ -396,6 +422,8 @@ impl CodexResult { error: Some(message), codex_error_info: None, diagnostics: Vec::new(), + run_id: None, + provenance: CodexProvenance::Produced, } } @@ -415,6 +443,10 @@ impl CodexResult { } else if self.exit_code != 0 { t.set("codex_error_info", "UNKNOWN")?; } + t.set("provenance", self.provenance.label())?; + if let Some(run_id) = self.run_id { + t.set("run_id", run_id)?; + } Ok(t) } } @@ -937,7 +969,7 @@ fn run_adoptable_codex_request( host_root: &Path, config: &ConfigContext, ) -> Result { - if let Some(result) = read_reusable_adoption_result(&paths)? { + if let Some(result) = read_reusable_adoption_result(&paths, &request.run_id)? { return Ok(result); } std::fs::create_dir_all(&paths.dir).map_err(mlua::Error::external)?; @@ -950,11 +982,11 @@ fn run_adoptable_codex_request( .open(lock_path) .map_err(mlua::Error::external)?; flock(lock_file.as_raw_fd(), FlockArg::LockExclusive).map_err(mlua::Error::external)?; - if let Some(result) = read_reusable_adoption_result(&paths)? { + if let Some(result) = read_reusable_adoption_result(&paths, &request.run_id)? { drop(lock_file); return Ok(result); } - if let Some(result) = recover_reusable_adoption_result_locked(&paths)? { + if let Some(result) = recover_reusable_adoption_result_locked(&paths, &request.run_id)? { drop(lock_file); return Ok(result); } @@ -1133,7 +1165,10 @@ fn adoption_request_hash(request: &CodexRequest) -> String { hex } -fn read_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result> { +fn read_completed_adoption_result( + paths: &CodexAdoptionPaths, + requesting_run_id: &str, +) -> Result> { let Some(record) = read_adoption_record(&paths.status)? else { return Ok(None); }; @@ -1151,6 +1186,12 @@ fn read_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result Result Result Result> { - let Some(result) = read_completed_adoption_result(paths)? else { +fn read_reusable_adoption_result( + paths: &CodexAdoptionPaths, + requesting_run_id: &str, +) -> Result> { + let Some(result) = read_completed_adoption_result(paths, requesting_run_id)? else { return Ok(None); }; if !adoption_result_succeeded(&result) { @@ -1204,8 +1252,9 @@ fn adoption_result_succeeded(result: &CodexResult) -> bool { fn recover_reusable_adoption_result_locked( paths: &CodexAdoptionPaths, + requesting_run_id: &str, ) -> Result> { - let Some(result) = recover_completed_adoption_result_locked(paths)? else { + let Some(result) = recover_completed_adoption_result_locked(paths, requesting_run_id)? else { return Ok(None); }; if !adoption_result_succeeded(&result) { @@ -1214,23 +1263,27 @@ fn recover_reusable_adoption_result_locked( Ok(Some(result)) } -fn recover_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result> { +fn recover_completed_adoption_result( + paths: &CodexAdoptionPaths, + requesting_run_id: &str, +) -> Result> { let lock_file = open_adoption_run_lock(&paths.dir).map_err(mlua::Error::external)?; flock(lock_file.as_raw_fd(), FlockArg::LockExclusive).map_err(mlua::Error::external)?; - let result = recover_completed_adoption_result_locked(paths); + let result = recover_completed_adoption_result_locked(paths, requesting_run_id); drop(lock_file); result } fn recover_completed_adoption_result_locked( paths: &CodexAdoptionPaths, + requesting_run_id: &str, ) -> Result> { let recovered = recover_completed_adoption_result_locked_from_disk(paths).map_err(mlua::Error::external)?; if !recovered { return Ok(None); }; - read_completed_adoption_result(paths) + read_completed_adoption_result(paths, requesting_run_id) } fn recover_completed_adoption_result_locked_from_disk( @@ -1778,10 +1831,10 @@ fn wait_for_adoption_result( .unwrap_or_else(|| Duration::from_secs(DEFAULT_CODEX_TIMEOUT_SECONDS as u64)) + CODEX_ADOPTION_TIMEOUT_GRACE; loop { - if let Some(result) = read_completed_adoption_result(paths)? { + if let Some(result) = read_completed_adoption_result(paths, &request.run_id)? { return Ok(result); } - if let Some(result) = recover_completed_adoption_result(paths)? { + if let Some(result) = recover_completed_adoption_result(paths, &request.run_id)? { return Ok(result); } if Instant::now() >= deadline { diff --git a/crates/fkst-framework/tests/sdk_codex.rs b/crates/fkst-framework/tests/sdk_codex.rs index e48c504..330d693 100644 --- a/crates/fkst-framework/tests/sdk_codex.rs +++ b/crates/fkst-framework/tests/sdk_codex.rs @@ -1000,6 +1000,76 @@ printf '{"type":"item.completed","item":{"type":"agent_message","text":"result-% assert!(!worktree.join(".fkst-codex").exists()); } +#[cfg(unix)] +#[test] +fn spawn_codex_sync_reports_whether_a_result_was_produced_or_adopted() { + // An adopted five-second read and a real hour-long run otherwise report the + // same thing. #3919 took a measurement of ELAPSED_MS to identify precisely + // because the result carried no provenance. + let tmp = tempfile::tempdir().unwrap(); + let bin_dir = tmp.path().join("bin"); + let worktree = tmp.path().join("wt"); + let capture_dir = tmp.path().join("capture"); + std::fs::create_dir_all(&worktree).unwrap(); + std::fs::create_dir_all(&capture_dir).unwrap(); + install_codex_script( + &bin_dir, + r#"#!/bin/sh +cat >/dev/null +count_file="$CAPTURE_DIR/spawns" +count=0 +if [ -f "$count_file" ]; then + count=$(cat "$count_file") +fi +count=$((count + 1)) +printf '%s' "$count" > "$count_file" +printf '{"type":"item.completed","item":{"type":"agent_message","text":"result-%s"}}\n' "$count" +"#, + ); + + let mut sandbox = ProcessSandbox::new(); + sandbox.enter_cwd(tmp.path()).runtime_root(".fkst/runtime"); + sandbox.prepend_path(&bin_dir); + sandbox.set_env("CAPTURE_DIR", capture_dir.to_string_lossy().into_owned()); + sandbox.set_env(CODEX_WORKER_BIN_ENV, framework_bin()); + sandbox.runtime_log_dir(tmp.path().join("runtime")); + let (_lock, _guard) = sandbox.enter(); + + let lua = Lua::new(); + register(&lua).unwrap(); + let spawn: mlua::Function = lua.globals().get("spawn_codex_sync").unwrap(); + + let first_opts = lua_opts(&lua, "same-work"); + first_opts + .set("worktree", worktree.to_string_lossy().into_owned()) + .unwrap(); + first_opts.set("dedup_key", "provenance-dedup").unwrap(); + let first: Table = spawn.call(first_opts).unwrap(); + assert_eq!(first.get::("exit_code").unwrap(), 0); + assert_eq!(first.get::("provenance").unwrap(), "produced"); + + let second_opts = lua_opts(&lua, "same-work"); + second_opts + .set("worktree", worktree.to_string_lossy().into_owned()) + .unwrap(); + second_opts.set("dedup_key", "provenance-dedup").unwrap(); + let second: Table = spawn.call(second_opts).unwrap(); + assert_eq!(second.get::("stdout").unwrap(), "result-1"); + assert_eq!( + second.get::("provenance").unwrap(), + "adopted", + "a dispatch that read another run's record must say so" + ); + let run_id = second + .get::("run_id") + .expect("an adopted result must name the run that produced it"); + assert!(!run_id.is_empty()); + assert_eq!( + std::fs::read_to_string(capture_dir.join("spawns")).unwrap(), + "1" + ); +} + #[cfg(unix)] #[test] fn spawn_codex_sync_does_not_adopt_a_failed_result_for_a_later_dispatch() {