fix(runner): stop after terminal result delivery failure
This commit is contained in:
parent
156c913779
commit
32b1124242
|
|
@ -280,6 +280,14 @@ pub fn run_durable_runner<E: CommandExecutor>(
|
|||
}
|
||||
lifecycle_after_reply = lifecycle_after_reply.merge(lifecycle);
|
||||
if let Err(error) = transport.send_json(&command_result_envelope(&state, &result)) {
|
||||
if lifecycle.durable_state().is_some() {
|
||||
return stop_after_terminal_result_delivery_failure(
|
||||
&mut state,
|
||||
&store,
|
||||
&mut executor,
|
||||
error,
|
||||
);
|
||||
}
|
||||
state.record_diagnostic(error.to_string());
|
||||
disconnected = true;
|
||||
break;
|
||||
|
|
@ -397,6 +405,14 @@ pub fn run_durable_runner<E: CommandExecutor>(
|
|||
if let Err(error) =
|
||||
transport.send_json(&command_result_envelope(&state, &result))
|
||||
{
|
||||
if lifecycle.durable_state().is_some() {
|
||||
return stop_after_terminal_result_delivery_failure(
|
||||
&mut state,
|
||||
&store,
|
||||
&mut executor,
|
||||
error,
|
||||
);
|
||||
}
|
||||
disconnected_since.get_or_insert_with(Instant::now);
|
||||
state.record_diagnostic(error.to_string());
|
||||
state.reconnect_count = state.reconnect_count.saturating_add(1);
|
||||
|
|
@ -499,6 +515,23 @@ fn persist_lifecycle_before_command_delivery(
|
|||
store.save(state)
|
||||
}
|
||||
|
||||
fn stop_after_terminal_result_delivery_failure<E: CommandExecutor>(
|
||||
state: &mut DurableState,
|
||||
store: &DurableStateStore,
|
||||
executor: &mut E,
|
||||
error: DurableRunnerError,
|
||||
) -> Result<(), DurableRunnerError> {
|
||||
// The terminal transition was committed before the attempted delivery.
|
||||
// Its result may or may not have reached the controller, but reconnecting
|
||||
// this process would overwrite that durable state as ready and admit work
|
||||
// after shutdown/suspend. Leave the result journaled for reconciliation by
|
||||
// a future authorized process instead.
|
||||
state.record_diagnostic(error.to_string());
|
||||
store.save(state)?;
|
||||
let _ = executor.shutdown();
|
||||
Err(error)
|
||||
}
|
||||
|
||||
fn poll_executor_events<E: CommandExecutor>(
|
||||
state: &mut DurableState,
|
||||
store: &DurableStateStore,
|
||||
|
|
@ -663,6 +696,10 @@ mod tests {
|
|||
|
||||
struct ShutdownFailingExecutor;
|
||||
|
||||
struct ShutdownCountingExecutor {
|
||||
shutdown_calls: usize,
|
||||
}
|
||||
|
||||
struct RetainingEventExecutor {
|
||||
events: VecDeque<PolledEvent>,
|
||||
fail_acknowledgement: bool,
|
||||
|
|
@ -696,6 +733,17 @@ mod tests {
|
|||
}
|
||||
}
|
||||
|
||||
impl CommandExecutor for ShutdownCountingExecutor {
|
||||
fn execute(&mut self, _command: &Command) -> Result<CommandExecution, DurableRunnerError> {
|
||||
Ok(CommandExecution::result(json!({"status": "completed"})))
|
||||
}
|
||||
|
||||
fn shutdown(&mut self) -> Result<(), DurableRunnerError> {
|
||||
self.shutdown_calls += 1;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl CommandExecutor for RetainingEventExecutor {
|
||||
fn execute(&mut self, _command: &Command) -> Result<CommandExecution, DurableRunnerError> {
|
||||
Ok(CommandExecution::result(json!({"status": "completed"})))
|
||||
|
|
@ -780,6 +828,40 @@ mod tests {
|
|||
fs::remove_dir_all(directory).unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn terminal_result_delivery_failure_stops_without_reopening_lifecycle() {
|
||||
let directory = std::env::temp_dir().join(format!(
|
||||
"paperclip-runner-terminal-result-delivery-{}",
|
||||
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();
|
||||
state.lifecycle = "stopped".to_owned();
|
||||
store.save(&state).unwrap();
|
||||
let mut executor = ShutdownCountingExecutor { shutdown_calls: 0 };
|
||||
|
||||
let error = stop_after_terminal_result_delivery_failure(
|
||||
&mut state,
|
||||
&store,
|
||||
&mut executor,
|
||||
DurableRunnerError::invalid("simulated result delivery failure"),
|
||||
)
|
||||
.expect_err("terminal result delivery failure remains observable");
|
||||
let (recovered, existed) = store.load_or_create(&config).unwrap();
|
||||
|
||||
assert!(error.to_string().contains("result delivery failure"));
|
||||
assert!(existed);
|
||||
assert_eq!(recovered.lifecycle, "stopped");
|
||||
assert_eq!(executor.shutdown_calls, 1);
|
||||
assert!(recovered
|
||||
.diagnostics
|
||||
.iter()
|
||||
.any(|diagnostic| diagnostic.contains("result delivery failure")));
|
||||
fs::remove_dir_all(directory).unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn event_batch_keeps_accepted_prefix_and_unacknowledged_suffix() {
|
||||
let directory = std::env::temp_dir().join(format!(
|
||||
|
|
|
|||
Loading…
Reference in New Issue