fix(runner): reconcile terminal result delivery

This commit is contained in:
Dotta 2026-09-02 23:38:26 -05:00
parent 76c82e642c
commit d582fbe1aa
2 changed files with 193 additions and 6 deletions

View File

@ -6,7 +6,7 @@ use serde_json::{json, Value};
use super::state::{
Command, CommandDisposition, DurableState, DurableStateStore, EventPriority,
StoredCommandResult,
PendingTerminalDelivery, StoredCommandResult,
};
use super::transport::{
current_unix_ms, validate_control_identity, AuthenticatedTransport, ConnectionMetadata,
@ -141,7 +141,9 @@ pub fn run_durable_runner<E: CommandExecutor>(
config.validate()?;
let store = DurableStateStore::new(&config.state_dir)?;
let (mut state, recovered) = store.load_or_create(&config)?;
if state.lifecycle == "revoked" || state.lifecycle == "stopped" {
if state.lifecycle == "revoked"
|| (state.lifecycle == "stopped" && state.pending_terminal_delivery.is_none())
{
return Ok(());
}
if recovered {
@ -265,6 +267,16 @@ pub fn run_durable_runner<E: CommandExecutor>(
if let Some(acked_source_seq) = welcome.acked_source_seq {
state.apply_ack(acked_source_seq)?;
}
if state.pending_terminal_delivery.is_some() {
return reconcile_pending_terminal_delivery(
&mut state,
&store,
&config,
&mut executor,
&mut transport,
&welcome.pending_commands,
);
}
state.lifecycle = "ready".to_owned();
state.recoverable_failure = None;
store.save(&state)?;
@ -276,7 +288,12 @@ pub fn run_durable_runner<E: CommandExecutor>(
let (result, lifecycle) =
process_command(&mut state, &store, &config, &mut executor, &command)?;
if let Some(durable_lifecycle) = lifecycle.durable_state() {
persist_lifecycle_before_command_delivery(&mut state, &store, durable_lifecycle)?;
persist_lifecycle_before_command_delivery(
&mut state,
&store,
durable_lifecycle,
&result,
)?;
}
lifecycle_after_reply = lifecycle_after_reply.merge(lifecycle);
if let Err(error) = transport.send_json(&command_result_envelope(&state, &result)) {
@ -293,6 +310,7 @@ pub fn run_durable_runner<E: CommandExecutor>(
break;
}
if lifecycle.durable_state().is_some() {
complete_terminal_delivery_after_send(&mut state, &store, &result)?;
// A terminal lifecycle command is the final command this
// process may accept. Flush its already-durable outbox below,
// then release the executor without observing later commands.
@ -400,6 +418,7 @@ pub fn run_durable_runner<E: CommandExecutor>(
&mut state,
&store,
durable_lifecycle,
&result,
)?;
}
if let Err(error) =
@ -419,6 +438,9 @@ pub fn run_durable_runner<E: CommandExecutor>(
store.save(&state)?;
break;
}
if lifecycle.durable_state().is_some() {
complete_terminal_delivery_after_send(&mut state, &store, &result)?;
}
if let Err(error) = send_outbox(&mut transport, &state, &mut sent_source_seq) {
state.record_diagnostic(
"outbox delivery failed; unacknowledged suffix remains durable",
@ -497,7 +519,8 @@ fn persist_lifecycle_before_shutdown<E: CommandExecutor>(
executor: &mut E,
lifecycle: &str,
) -> Result<(), DurableRunnerError> {
persist_lifecycle_before_command_delivery(state, store, lifecycle)?;
state.lifecycle = lifecycle.to_owned();
store.save(state)?;
executor.shutdown()
}
@ -505,6 +528,7 @@ fn persist_lifecycle_before_command_delivery(
state: &mut DurableState,
store: &DurableStateStore,
lifecycle: &str,
result: &StoredCommandResult,
) -> Result<(), DurableRunnerError> {
// A terminal command result is already durable before this boundary. Save
// its matching lifecycle before exposing that result to the controller,
@ -512,9 +536,89 @@ fn persist_lifecycle_before_command_delivery(
// Recovery can therefore never observe a ready runner after the controller
// has already observed its terminal command result.
state.lifecycle = lifecycle.to_owned();
state.pending_terminal_delivery = Some(PendingTerminalDelivery {
command_id: result.command_id.clone(),
controller_seq: result.controller_seq,
command_type: result.command_type.clone(),
lifecycle: lifecycle.to_owned(),
});
store.save(state)
}
fn complete_terminal_delivery_after_send(
state: &mut DurableState,
store: &DurableStateStore,
result: &StoredCommandResult,
) -> Result<(), DurableRunnerError> {
let pending = state.pending_terminal_delivery.as_ref().ok_or_else(|| {
DurableRunnerError::invalid("terminal result delivery has no durable recovery fence")
})?;
if pending.command_id != result.command_id
|| pending.controller_seq != result.controller_seq
|| pending.command_type != result.command_type
|| state.lifecycle != pending.lifecycle
{
return Err(DurableRunnerError::invalid(
"terminal result delivery does not match its durable recovery fence",
));
}
state.pending_terminal_delivery = None;
store.save(state)
}
fn reconcile_pending_terminal_delivery<E: CommandExecutor>(
state: &mut DurableState,
store: &DurableStateStore,
config: &DurableRunnerConfig,
executor: &mut E,
transport: &mut AuthenticatedTransport,
pending_commands: &[Command],
) -> Result<(), DurableRunnerError> {
let pending = state.pending_terminal_delivery.clone().ok_or_else(|| {
DurableRunnerError::invalid("terminal result reconciliation has no durable fence")
})?;
if let Some(command) = pending_commands
.iter()
.find(|command| command.command_id == pending.command_id)
{
if command.controller_seq != pending.controller_seq
|| command.command_type != pending.command_type
{
return Err(DurableRunnerError::invalid(
"controller changed the pending terminal command identity",
));
}
let (result, lifecycle) = process_command(state, store, config, executor, command)?;
if lifecycle.durable_state() != Some(pending.lifecycle.as_str()) {
return Err(DurableRunnerError::invalid(
"pending terminal command did not replay its durable lifecycle",
));
}
if let Err(error) = transport.send_json(&command_result_envelope(state, &result)) {
return stop_after_terminal_result_delivery_failure(state, store, executor, error);
}
complete_terminal_delivery_after_send(state, store, &result)?;
} else {
// An authenticated welcome is the controller's authoritative pending
// set. Absence means the prior write reached the controller even if
// the runner did not observe transport success before it exited.
state.record_diagnostic(
"controller confirmed the pending terminal result was already delivered",
);
state.pending_terminal_delivery = None;
store.save(state)?;
}
let mut sent_source_seq = state.acked_source_seq;
if let Err(error) = send_outbox(transport, state, &mut sent_source_seq) {
state.record_diagnostic("outbox delivery failed after terminal result reconciliation");
store.save(state)?;
let _ = executor.shutdown();
return Err(error);
}
executor.shutdown()
}
fn stop_after_terminal_result_delivery_failure<E: CommandExecutor>(
state: &mut DurableState,
store: &DurableStateStore,
@ -838,9 +942,17 @@ mod tests {
let config = config(directory.clone());
let store = DurableStateStore::new(&directory).unwrap();
let (mut state, _) = store.load_or_create(&config).unwrap();
state.lifecycle = "stopped".to_owned();
store.save(&state).unwrap();
let mut executor = ShutdownCountingExecutor { shutdown_calls: 0 };
let command = command("runner.shutdown");
let (result, lifecycle) =
process_command(&mut state, &store, &config, &mut executor, &command).unwrap();
persist_lifecycle_before_command_delivery(
&mut state,
&store,
lifecycle.durable_state().unwrap(),
&result,
)
.unwrap();
let error = stop_after_terminal_result_delivery_failure(
&mut state,
@ -854,6 +966,13 @@ mod tests {
assert!(error.to_string().contains("result delivery failure"));
assert!(existed);
assert_eq!(recovered.lifecycle, "stopped");
assert_eq!(
recovered
.pending_terminal_delivery
.as_ref()
.map(|pending| pending.command_id.as_str()),
Some("command_1")
);
assert_eq!(executor.shutdown_calls, 1);
assert!(recovered
.diagnostics
@ -862,6 +981,38 @@ mod tests {
fs::remove_dir_all(directory).unwrap();
}
#[test]
fn successful_terminal_result_delivery_clears_the_recovery_fence() {
let directory = std::env::temp_dir().join(format!(
"paperclip-runner-terminal-result-delivered-{}",
std::process::id()
));
let _ = fs::remove_dir_all(&directory);
let config = config(directory.clone());
let store = DurableStateStore::new(&directory).unwrap();
let (mut state, _) = store.load_or_create(&config).unwrap();
let mut executor = CountingExecutor { calls: 0 };
let command = command("runner.suspend");
let (result, lifecycle) =
process_command(&mut state, &store, &config, &mut executor, &command).unwrap();
persist_lifecycle_before_command_delivery(
&mut state,
&store,
lifecycle.durable_state().unwrap(),
&result,
)
.unwrap();
assert!(state.pending_terminal_delivery.is_some());
complete_terminal_delivery_after_send(&mut state, &store, &result).unwrap();
let (recovered, existed) = store.load_or_create(&config).unwrap();
assert!(existed);
assert_eq!(recovered.lifecycle, "suspended");
assert!(recovered.pending_terminal_delivery.is_none());
fs::remove_dir_all(directory).unwrap();
}
#[test]
fn event_batch_keeps_accepted_prefix_and_unacknowledged_suffix() {
let directory = std::env::temp_dir().join(format!(

View File

@ -151,6 +151,15 @@ pub struct StoredCommandResult {
pub result: Value,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub(crate) struct PendingTerminalDelivery {
pub(crate) command_id: String,
pub(crate) controller_seq: u64,
pub(crate) command_type: String,
pub(crate) lifecycle: String,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
struct ExecutorEventReceipt {
@ -189,6 +198,8 @@ pub struct DurableState {
#[serde(default)]
pub processed_command_fingerprints: BTreeMap<String, String>,
#[serde(default)]
pub(crate) pending_terminal_delivery: Option<PendingTerminalDelivery>,
#[serde(default)]
executor_event_receipts: BTreeMap<String, ExecutorEventReceipt>,
pub diagnostics: Vec<String>,
pub backpressure: bool,
@ -217,6 +228,7 @@ impl DurableState {
outbox: Vec::new(),
processed_commands: BTreeMap::new(),
processed_command_fingerprints: BTreeMap::new(),
pending_terminal_delivery: None,
executor_event_receipts: BTreeMap::new(),
diagnostics: Vec::new(),
backpressure: false,
@ -965,6 +977,29 @@ fn validate_binding(
&& receipt.source_seq <= state.highest_source_seq()
&& executor_receipt_sequences.insert(receipt.source_seq)
});
let pending_terminal_delivery_is_valid =
state
.pending_terminal_delivery
.as_ref()
.map_or(true, |pending| {
let expected_lifecycle = match pending.command_type.as_str() {
"runner.suspend" => "suspended",
"runner.shutdown" => "stopped",
_ => return false,
};
pending.lifecycle == expected_lifecycle
&& state.lifecycle == expected_lifecycle
&& pending.controller_seq == state.last_controller_command_seq
&& state
.processed_commands
.get(&pending.command_id)
.is_some_and(|result| {
result.command_id == pending.command_id
&& result.controller_seq == pending.controller_seq
&& result.command_type == pending.command_type
&& result.status != "pending"
})
});
command_sequences.sort_unstable();
let command_cursors_are_valid = match (command_sequences.first(), command_sequences.last()) {
(None, None) => state.compacted_through_controller_seq == state.last_controller_command_seq,
@ -991,6 +1026,7 @@ fn validate_binding(
|| !command_cursors_are_valid
|| !command_fingerprints_are_valid
|| !executor_event_receipts_are_valid
|| !pending_terminal_delivery_is_valid
{
return Err(DurableRunnerError::invalid(
"durable state cursors, bounds, or journals are inconsistent",