Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <mode>` 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 <mode>` 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 = <delivery_id>`, and `page = { section = "dead_letters", after? = <opaque cursor> }`. Lineage mode is `lineage = { queue = <queue>, dept = <department>, source_ref = { kind = <kind>, ref = <reference> } }` 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? = <opaque cursor> }`, 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 文本反推边界状态。
Expand Down
135 changes: 126 additions & 9 deletions crates/fkst-framework/src/sdk_codex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -233,6 +233,28 @@ pub(crate) struct CodexResult {
error: Option<String>,
codex_error_info: Option<String>,
diagnostics: Vec<CodexDiagnostic>,
/// 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<String>,
/// 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.
Expand Down Expand Up @@ -319,6 +341,8 @@ impl CodexResult {
error: None,
codex_error_info: None,
diagnostics: Vec::new(),
run_id: None,
provenance: CodexProvenance::Produced,
}
}

Expand Down Expand Up @@ -376,6 +400,8 @@ impl CodexResult {
error: Some(message),
codex_error_info: None,
diagnostics: Vec::new(),
run_id: None,
provenance: CodexProvenance::Produced,
}
}

Expand All @@ -396,6 +422,8 @@ impl CodexResult {
error: Some(message),
codex_error_info: None,
diagnostics: Vec::new(),
run_id: None,
provenance: CodexProvenance::Produced,
}
}

Expand All @@ -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)
}
}
Expand Down Expand Up @@ -937,7 +969,7 @@ fn run_adoptable_codex_request(
host_root: &Path,
config: &ConfigContext,
) -> Result<CodexResult> {
if let Some(result) = read_completed_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)?;
Expand All @@ -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_completed_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_completed_adoption_result_locked(&paths)? {
if let Some(result) = recover_reusable_adoption_result_locked(&paths, &request.run_id)? {
drop(lock_file);
return Ok(result);
}
Expand Down Expand Up @@ -1133,7 +1165,10 @@ fn adoption_request_hash(request: &CodexRequest) -> String {
hex
}

fn read_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result<Option<CodexResult>> {
fn read_completed_adoption_result(
paths: &CodexAdoptionPaths,
requesting_run_id: &str,
) -> Result<Option<CodexResult>> {
let Some(record) = read_adoption_record(&paths.status)? else {
return Ok(None);
};
Expand All @@ -1151,6 +1186,12 @@ fn read_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result<Option<C
let exit_code = record.exit_code.unwrap_or(-1);
let codex_error_info = record.codex_error_info.clone();
let diagnostics = record.diagnostics.clone();
let produced_by = record.run_id.clone();
let provenance = if produced_by == requesting_run_id {
CodexProvenance::Produced
} else {
CodexProvenance::Adopted
};
if let Some(kind) = record.error_kind.as_deref() {
let mut result = CodexResult::failure(
kind,
Expand All @@ -1165,31 +1206,84 @@ fn read_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result<Option<C
);
result.codex_error_info = codex_error_info;
result.diagnostics = diagnostics;
result.run_id = Some(produced_by.clone());
result.provenance = provenance;
return Ok(Some(result));
}
let mut result = CodexResult::success(stdout, stderr, exit_code, record.log_path);
result.codex_error_info = codex_error_info;
result.diagnostics = diagnostics;
result.run_id = Some(produced_by);
result.provenance = provenance;
Ok(Some(result))
}

/// A completed adoption record is reusable by a *later* dispatch only when the run
/// succeeded.
///
/// Every codex failure is classified as an environmental boundary condition —
/// `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 replaying it to a later dispatch freezes one transient fault into
/// a permanent answer that no retry can ever clear.
///
/// Waiters attached to a run in flight keep reading the completed record through
/// `read_completed_adoption_result`; they are entitled to their own run's failure.
fn read_reusable_adoption_result(
paths: &CodexAdoptionPaths,
requesting_run_id: &str,
) -> Result<Option<CodexResult>> {
let Some(result) = read_completed_adoption_result(paths, requesting_run_id)? else {
return Ok(None);
};
if !adoption_result_succeeded(&result) {
return Ok(None);
}
Ok(Some(result))
}

fn recover_completed_adoption_result(paths: &CodexAdoptionPaths) -> Result<Option<CodexResult>> {
/// 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,
requesting_run_id: &str,
) -> Result<Option<CodexResult>> {
let Some(result) = recover_completed_adoption_result_locked(paths, requesting_run_id)? else {
return Ok(None);
};
if !adoption_result_succeeded(&result) {
return Ok(None);
}
Ok(Some(result))
}

fn recover_completed_adoption_result(
paths: &CodexAdoptionPaths,
requesting_run_id: &str,
) -> Result<Option<CodexResult>> {
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<Option<CodexResult>> {
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(
Expand Down Expand Up @@ -1556,6 +1650,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);
Expand All @@ -1581,6 +1676,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,
Expand Down Expand Up @@ -1714,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 {
Expand Down
Loading
Loading