From 27f7a27b3afec77c935a42b8b4e0f1862b8e8518 Mon Sep 17 00:00:00 2001 From: iceteaSA <171169159+iceteaSA@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:27:16 +0200 Subject: [PATCH] feat(protocol): share established route death decision Pin origin as recorded but not enforced until every compatible daemon has landed and deployed the origin bit. Keep SDK retries and provider classifications unchanged. CONSUMER-IMPACT: additive constants + predicate; SDK behavior unchanged; origin recorded not enforced --- Cargo.lock | 4 +- clients/subc-client/CHANGELOG.md | 4 ++ clients/subc-client/package.json | 2 +- clients/subc-client/src/client.ts | 27 ++++++++++- clients/subc-client/src/index.ts | 3 ++ clients/subc-client/src/provider.ts | 3 +- .../tests/golden-conformance.test.ts | 21 +++++++++ crates/subc-client-rs/Cargo.toml | 2 +- crates/subc-client-rs/src/consumer.rs | 24 ++++++---- .../subc-client-rs/tests/contract_tables.rs | 46 +++++++++++++++++++ crates/subc-protocol/Cargo.toml | 2 +- crates/subc-protocol/src/lib.rs | 23 ++++++++++ .../tests/golden/decision_tables.json | 20 ++++++++ crates/subc-protocol/tests/golden_json.rs | 29 ++++++++++++ 14 files changed, 193 insertions(+), 17 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 12eee1ba..3fffc22f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1807,7 +1807,7 @@ dependencies = [ [[package]] name = "subc-client-rs" -version = "0.18.4" +version = "0.18.5" dependencies = [ "async-trait", "serde", @@ -1918,7 +1918,7 @@ dependencies = [ [[package]] name = "subc-protocol" -version = "0.25.1" +version = "0.25.2" dependencies = [ "serde", "serde_json", diff --git a/clients/subc-client/CHANGELOG.md b/clients/subc-client/CHANGELOG.md index cbefa296..375da0dd 100644 --- a/clients/subc-client/CHANGELOG.md +++ b/clients/subc-client/CHANGELOG.md @@ -1,5 +1,9 @@ # Changelog +## 0.16.1 — 2026-09-24 + +- Export `UNKNOWN_CHANNEL`, `STALE_ROUTE_EPOCH`, and `isEstablishedRouteDead(flags, code)` for evict/reopen/resend-once decisions. The managed consumer and provider use the shared predicate without changing their current retry or `not_sent` behavior. The golden table records daemon-origin flags but does not enforce them until all compatible daemons emit the bit. + ## 0.16.0 — 2026-09-24 - Make `SubcProvider.closed` public: a `Promise` that resolves once, when the provider will serve no more. That happens after a channel-0 GOODBYE from the daemon (supervisor restart or stop, daemon shutdown), after `close()`, or after a connection loss the provider will not recover from. It never rejects. A supervised module should `await provider.closed` and then exit (`process.exit(0)`), so a restart no longer waits for the daemon to SIGTERM it at the drain deadline. The SDK itself never exits the process. diff --git a/clients/subc-client/package.json b/clients/subc-client/package.json index af09a7af..50605f9d 100644 --- a/clients/subc-client/package.json +++ b/clients/subc-client/package.json @@ -1,6 +1,6 @@ { "name": "@cortexkit/subc-client", - "version": "0.16.0", + "version": "0.16.1", "description": "TypeScript client for the subc daemon. Wire-compatible (byte-for-byte) with the Rust subc-transport handshake and subc-protocol envelope.", "type": "module", "exports": { diff --git a/clients/subc-client/src/client.ts b/clients/subc-client/src/client.ts index b53775e8..bdd57b61 100644 --- a/clients/subc-client/src/client.ts +++ b/clients/subc-client/src/client.ts @@ -320,6 +320,10 @@ export class SubcError extends Error { } } +// The wire flags belong to the frame, not the public error API. Keep them with +// that error across managed-call wrapping so Phase 2 can enforce origin here. +const errorFrameFlags = new WeakMap(); + function requireBinaryBody(body: unknown): Uint8Array { if (body instanceof Uint8Array) return body; const type = body === null ? "null" : Array.isArray(body) ? "array" : typeof body; @@ -692,7 +696,10 @@ export class SubcClient { // request was in flight. The daemon's contract for the code is // NOT-FORWARDED — dropped before delivery — so the retry is safe by // construction, and the remedy is identical: evict, re-open, resend once. - const deadBindCode = err.code === "unknown_channel" || err.code === "stale_route_epoch"; + const deadBindCode = isEstablishedRouteDead( + err.cause instanceof SubcError ? (errorFrameFlags.get(err.cause) ?? 0) : 0, + err.code, + ); if (deadBindCode && !retriedUnknownChannel && !this.closeStarted) { retriedUnknownChannel = true; this.evictRouteHandle(routeHandle); @@ -1776,7 +1783,9 @@ export class SubcClient { message?: string; detail?: unknown; }; - return new SubcError(parsed.message ?? "subc error", parsed.code, parsed.detail); + const error = new SubcError(parsed.message ?? "subc error", parsed.code, parsed.detail); + errorFrameFlags.set(error, frame.header.flags); + return error; } catch { return new SubcError(Buffer.from(frame.body).toString("utf8") || "subc error"); } @@ -1940,6 +1949,20 @@ export function isRetryableRouteOpenCode(code: string | undefined): boolean { ); } +export const UNKNOWN_CHANNEL = "unknown_channel"; +export const STALE_ROUTE_EPOCH = "stale_route_epoch"; + +/** Whether to evict this route, reopen it, and resend once. A consumer deriving + * custody, provider, or suspect verdicts from route death owns its own code list. + * Kept byte-identical to subc-protocol error_codes::is_established_route_dead. + * The origin bit is recorded but not enforced until every daemon consumers can + * meet sets it on these frames (landed and deployed); requiring it sooner would + * stop a new SDK from evicting genuine stale routes on older daemons. + */ +export function isEstablishedRouteDead(_flags: number, code: string | undefined): boolean { + return code === UNKNOWN_CHANNEL || code === STALE_ROUTE_EPOCH; +} + export async function connectionFileExists(path: string): Promise { try { await fs.access(path); diff --git a/clients/subc-client/src/index.ts b/clients/subc-client/src/index.ts index 0d56cd34..021be6d2 100644 --- a/clients/subc-client/src/index.ts +++ b/clients/subc-client/src/index.ts @@ -2,6 +2,9 @@ export { SubcClient, SubcError, SubcCallError, + UNKNOWN_CHANNEL, + STALE_ROUTE_EPOCH, + isEstablishedRouteDead, DEFAULT_RECONNECT_BACKOFF, SUBC_MODULE_ID_ENV, SUBC_LAUNCH_NONCE_ENV, diff --git a/clients/subc-client/src/provider.ts b/clients/subc-client/src/provider.ts index 89b6ac60..514aa060 100644 --- a/clients/subc-client/src/provider.ts +++ b/clients/subc-client/src/provider.ts @@ -3,6 +3,7 @@ import { Buffer } from "node:buffer"; import { AuthError, authenticateClient } from "./auth.js"; import { DEFAULT_RECONNECT_BACKOFF, + isEstablishedRouteDead, SUBC_LAUNCH_NONCE_ENV, SUBC_MODULE_ID_ENV, type BindIdentity, @@ -1325,7 +1326,7 @@ function providerErrorFromFrame(frame: Frame): SubcProviderError { return new SubcProviderError( body.message ?? "subc error", body.code, - body.code === "stale_route_epoch" || body.code === "unknown_channel" ? "not_sent" : "terminal", + isEstablishedRouteDead(frame.header.flags, body.code) ? "not_sent" : "terminal", body.detail, ); } catch { diff --git a/clients/subc-client/tests/golden-conformance.test.ts b/clients/subc-client/tests/golden-conformance.test.ts index 1c8a4393..77e1728c 100644 --- a/clients/subc-client/tests/golden-conformance.test.ts +++ b/clients/subc-client/tests/golden-conformance.test.ts @@ -5,6 +5,7 @@ import { join } from "node:path"; import { classifyRouteCloseReason, DEFAULT_REQUEST_TIMEOUT_MS, + isEstablishedRouteDead, isRetryableRouteOpenCode, LIVENESS_PROBE_WINDOW_MS, ROUTE_OPEN_RETRY_DEADLINE_MS, @@ -20,6 +21,7 @@ import { HEADER_LEN, MAX_FRAME_BODY_LEN, PROTOCOL_VERSION, + DAEMON_ORIGIN_FLAG, SUBSCRIPTION_FLAG, } from "../src/envelope"; import { machineIdFromHelloAck, type ModuleHelloAckBody } from "../src/provider"; @@ -432,6 +434,25 @@ describe("Rust golden fixtures", () => { } }); + test("execute established_route_dead decision table over every row", () => { + const tables = loadGolden<{ + established_route_dead: { + daemon_origin: string; + phase_2_order: string; + rows: { code: string; daemon_origin: boolean; route_dead: boolean }[]; + }; + }>("decision_tables"); + expect(tables.established_route_dead.daemon_origin).toBe("recorded_not_enforced"); + expect(tables.established_route_dead.phase_2_order).toBe( + "all_compatible_daemons_landed_and_deployed_before_origin_enforced", + ); + expect(tables.established_route_dead.rows.length).toBeGreaterThanOrEqual(14); + for (const row of tables.established_route_dead.rows) { + const flags = row.daemon_origin ? DAEMON_ORIGIN_FLAG : 0; + expect(isEstablishedRouteDead(flags, row.code)).toBe(row.route_dead); + } + }); + test("execute route_close_disposition decision table over every row", () => { interface DecisionTables { route_open_retryable: Record; diff --git a/crates/subc-client-rs/Cargo.toml b/crates/subc-client-rs/Cargo.toml index 1cbeb7b9..b9304514 100644 --- a/crates/subc-client-rs/Cargo.toml +++ b/crates/subc-client-rs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "subc-client-rs" -version = "0.18.4" +version = "0.18.5" edition = "2021" publish = true description = "Shared serve + consume client for Rust subc modules." diff --git a/crates/subc-client-rs/src/consumer.rs b/crates/subc-client-rs/src/consumer.rs index 0a20622d..59505eae 100644 --- a/crates/subc-client-rs/src/consumer.rs +++ b/crates/subc-client-rs/src/consumer.rs @@ -897,7 +897,7 @@ impl SubcConsumer { subc_ops, }); } - Ok(TerminalFrame::Error { body }) => return Err(CallError::Module(body)), + Ok(TerminalFrame::Error { body, .. }) => return Err(CallError::Module(body)), Ok(TerminalFrame::StreamEnd) => { return Err(CallError::not_sent("catalog.list returned StreamEnd")); } @@ -977,7 +977,7 @@ impl SubcConsumer { match result? { TerminalFrame::Response { body, .. } => Ok(body), TerminalFrame::StreamEnd => Ok(Vec::new()), - TerminalFrame::Error { body } => Err(CallError::Module(body)), + TerminalFrame::Error { body, .. } => Err(CallError::Module(body)), } } @@ -1222,8 +1222,8 @@ impl SubcConsumer { // contract is NOT-FORWARDED (dropped before delivery), so the retry // is safe by construction; the remedy is identical. // Parity with the TS client's retry-once in call(). - Ok(TerminalFrame::Error { body, .. }) - if (body.code == "unknown_channel" || body.code == "stale_route_epoch") + Ok(TerminalFrame::Error { body, flags }) + if error_codes::is_established_route_dead(flags, &body.code) && !retried_unknown_channel && Instant::now() < call_deadline => { @@ -3800,7 +3800,7 @@ impl PendingEntry { PendingCompletion::Subscription { closed, .. } => { let result = match terminal { PendingTerminal::Response { .. } | PendingTerminal::StreamEnd => Ok(()), - PendingTerminal::Error { body } => Err(CallError::Module(body)), + PendingTerminal::Error { body, .. } => Err(CallError::Module(body)), }; let _ = closed.send(result); } @@ -3901,7 +3901,7 @@ impl PendingResult { enum PendingTerminal { Response { generation: u64, body: Vec }, - Error { body: ErrorBody }, + Error { body: ErrorBody, flags: Flags }, StreamEnd, } @@ -3909,7 +3909,7 @@ impl PendingTerminal { fn into_terminal_frame(self) -> TerminalFrame { match self { Self::Response { generation, body } => TerminalFrame::Response { generation, body }, - Self::Error { body } => TerminalFrame::Error { body }, + Self::Error { body, flags } => TerminalFrame::Error { body, flags }, Self::StreamEnd => TerminalFrame::StreamEnd, } } @@ -3918,7 +3918,7 @@ impl PendingTerminal { #[derive(Debug)] enum TerminalFrame { Response { generation: u64, body: Vec }, - Error { body: ErrorBody }, + Error { body: ErrorBody, flags: Flags }, StreamEnd, } @@ -4175,7 +4175,13 @@ async fn dispatch_frame(shared: &Arc, generation: u64, frame: Frame) -> message: err.to_string(), detail: None, }); - shared.settle_pending(key, PendingTerminal::Error { body }); + shared.settle_pending( + key, + PendingTerminal::Error { + body, + flags: frame.header.flags, + }, + ); } FrameType::StreamEnd => shared.settle_pending(key, PendingTerminal::StreamEnd), FrameType::StreamData => shared.route_stream_data(key, frame.body), diff --git a/crates/subc-client-rs/tests/contract_tables.rs b/crates/subc-client-rs/tests/contract_tables.rs index 071f4776..ed5fdffb 100644 --- a/crates/subc-client-rs/tests/contract_tables.rs +++ b/crates/subc-client-rs/tests/contract_tables.rs @@ -18,6 +18,21 @@ struct BudgetRow { struct DecisionTables { route_open_retryable: std::collections::BTreeMap, route_close_disposition: std::collections::BTreeMap, + established_route_dead: RouteDeathTable, +} + +#[derive(Debug, Deserialize)] +struct RouteDeathTable { + daemon_origin: String, + phase_2_order: String, + rows: Vec, +} + +#[derive(Debug, Deserialize)] +struct RouteDeathRow { + code: String, + daemon_origin: bool, + route_dead: bool, } fn golden_path(name: &str) -> PathBuf { @@ -149,6 +164,37 @@ fn execute_route_open_retryable_decision_table() { } } +#[test] +fn execute_established_route_dead_decision_table() { + let content = + std::fs::read_to_string(golden_path("decision_tables")).expect("read decision tables"); + let tables: DecisionTables = serde_json::from_str(&content).expect("parse decision tables"); + assert_eq!( + tables.established_route_dead.daemon_origin, + "recorded_not_enforced" + ); + assert_eq!( + tables.established_route_dead.phase_2_order, + "all_compatible_daemons_landed_and_deployed_before_origin_enforced" + ); + assert!(tables.established_route_dead.rows.len() >= 14); + for row in &tables.established_route_dead.rows { + let flags = if row.daemon_origin { + subc_protocol::Flags::new(false, subc_protocol::Priority::Passive, false) + .with_daemon_origin() + } else { + subc_protocol::Flags::new(false, subc_protocol::Priority::Passive, false) + }; + assert_eq!( + subc_protocol::error_codes::is_established_route_dead(flags, &row.code), + row.route_dead, + "code={} daemon_origin={}", + row.code, + row.daemon_origin + ); + } +} + #[test] fn execute_route_close_disposition_decision_table() { let path = golden_path("decision_tables"); diff --git a/crates/subc-protocol/Cargo.toml b/crates/subc-protocol/Cargo.toml index 180f7999..60aff62f 100644 --- a/crates/subc-protocol/Cargo.toml +++ b/crates/subc-protocol/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "subc-protocol" -version = "0.25.1" +version = "0.25.2" edition = "2021" publish = true description = "Shared wire contract for subc <-> modules: the 21-byte envelope, the Frame (header + opaque body), channel-0 control bodies, route.bind/RouteTarget session shapes, and the capability manifest. Single source of truth, depended on by subc-core and AFT." diff --git a/crates/subc-protocol/src/lib.rs b/crates/subc-protocol/src/lib.rs index 60bba369..fdb505e3 100644 --- a/crates/subc-protocol/src/lib.rs +++ b/crates/subc-protocol/src/lib.rs @@ -47,6 +47,10 @@ pub mod tool_call; /// Error frames remain extensible strings, but these daemon-owned route-open /// outcomes need identical spelling across the daemon and SDK retry policies. pub mod error_codes { + use crate::Flags; + + pub const UNKNOWN_CHANNEL: &str = "unknown_channel"; + pub const STALE_ROUTE_EPOCH: &str = "stale_route_epoch"; pub const UNKNOWN_MODULE: &str = "unknown_module"; pub const MODULE_REMOVED: &str = "module_removed"; /// The target module's endpoint is draining for a reload, restart or disable. @@ -104,6 +108,25 @@ pub mod error_codes { MODULE_RELOADING | MODULE_WARMING | TARGET_UNAVAILABLE | MODULE_TIMEOUT ) } + + /// Whether the caller should evict this established route, reopen it, and + /// resend the request once. Consumers mapping route death to their own + /// custody, provider, or suspect verdict keep their own named code lists; + /// sharing this predicate for those meanings can change a verdict on a + /// protocol bump (prefrontal#59). + /// + /// Takes envelope flags so the daemon-origin check can be enabled here in + /// one place. Phase 2 requires first that every daemon a consumer can meet + /// sets DAEMON_ORIGIN on these error frames, both landed and deployed; + /// only then may this predicate require the bit. Requiring it sooner makes + /// a new SDK against an older daemon stop evicting on a genuine + /// `stale_route_epoch` — the failure this predicate is meant to prevent. + /// + /// The daemon emits these codes in `RouterError::to_error_frame` at + /// `crates/subc-daemon/src/router.rs` (UnknownChannel and StaleRouteEpoch). + pub fn is_established_route_dead(_flags: Flags, code: &str) -> bool { + matches!(code, UNKNOWN_CHANNEL | STALE_ROUTE_EPOCH) + } } pub use frame::{Frame, FrameBuildError}; diff --git a/crates/subc-protocol/tests/golden/decision_tables.json b/crates/subc-protocol/tests/golden/decision_tables.json index 4cb066fb..76bdbb28 100644 --- a/crates/subc-protocol/tests/golden/decision_tables.json +++ b/crates/subc-protocol/tests/golden/decision_tables.json @@ -16,6 +16,26 @@ "target_unavailable": "retryable", "unknown_module": "terminal" }, + "established_route_dead": { + "daemon_origin": "recorded_not_enforced", + "phase_2_order": "all_compatible_daemons_landed_and_deployed_before_origin_enforced", + "rows": [ + { "code": "unknown_channel", "daemon_origin": false, "route_dead": true }, + { "code": "unknown_channel", "daemon_origin": true, "route_dead": true }, + { "code": "stale_route_epoch", "daemon_origin": false, "route_dead": true }, + { "code": "stale_route_epoch", "daemon_origin": true, "route_dead": true }, + { "code": "module_reloading", "daemon_origin": false, "route_dead": false }, + { "code": "module_reloading", "daemon_origin": true, "route_dead": false }, + { "code": "unknown_module", "daemon_origin": false, "route_dead": false }, + { "code": "unknown_module", "daemon_origin": true, "route_dead": false }, + { "code": "target_unavailable", "daemon_origin": false, "route_dead": false }, + { "code": "target_unavailable", "daemon_origin": true, "route_dead": false }, + { "code": "module_warming", "daemon_origin": false, "route_dead": false }, + { "code": "module_warming", "daemon_origin": true, "route_dead": false }, + { "code": "handler_failed", "daemon_origin": false, "route_dead": false }, + { "code": "handler_failed", "daemon_origin": true, "route_dead": false } + ] + }, "route_close_disposition": { "capability_denied": "must_not_reopen", "crash": "must_not_reopen", diff --git a/crates/subc-protocol/tests/golden_json.rs b/crates/subc-protocol/tests/golden_json.rs index 2905941d..174b724c 100644 --- a/crates/subc-protocol/tests/golden_json.rs +++ b/crates/subc-protocol/tests/golden_json.rs @@ -553,6 +553,35 @@ fn route_open_retry_predicate_matches_the_decision_table() { )); } +#[test] +fn established_route_dead_predicate_matches_the_decision_table() { + let table: Value = + serde_json::from_str(&fs::read_to_string(golden_path("decision_tables")).unwrap()).unwrap(); + let route_death = &table["established_route_dead"]; + assert_eq!(route_death["daemon_origin"], "recorded_not_enforced"); + assert_eq!( + route_death["phase_2_order"], + "all_compatible_daemons_landed_and_deployed_before_origin_enforced" + ); + let rows = route_death["rows"].as_array().expect("route death rows"); + assert!(rows.len() >= 14); + for row in rows { + let origin = row["daemon_origin"].as_bool().expect("origin bit"); + let flags = if origin { + subc_protocol::Flags::new(false, subc_protocol::Priority::Passive, false) + .with_daemon_origin() + } else { + subc_protocol::Flags::new(false, subc_protocol::Priority::Passive, false) + }; + let code = row["code"].as_str().expect("code"); + assert_eq!( + error_codes::is_established_route_dead(flags, code), + row["route_dead"].as_bool().expect("route dead verdict"), + "code={code} daemon_origin={origin}" + ); + } +} + fn assert_golden(name: &str, value: &T) where T: Serialize + DeserializeOwned + PartialEq + Debug,