fix(runner): settle parked provider before suspend
This commit is contained in:
parent
f810246e8b
commit
2cf19a8de0
|
|
@ -942,6 +942,45 @@ impl AcpxCommandExecutor {
|
|||
})))
|
||||
}
|
||||
|
||||
fn stop_turn_for_suspension(
|
||||
&mut self,
|
||||
reason: &str,
|
||||
) -> Result<CommandExecution, DurableRunnerError> {
|
||||
let turn_id = self
|
||||
.state
|
||||
.as_ref()
|
||||
.and_then(|state| state.active_turn_id.clone());
|
||||
let Some(turn_id) = turn_id else {
|
||||
return Ok(CommandExecution::result(json!({
|
||||
"status": "already_settled",
|
||||
"reason": reason,
|
||||
})));
|
||||
};
|
||||
self.session
|
||||
.as_mut()
|
||||
.ok_or_else(|| DurableRunnerError::invalid("ACPX session is unavailable"))?
|
||||
.terminate_active_turn_for_suspension(&turn_id)
|
||||
.map_err(|error| {
|
||||
DurableRunnerError::invalid(format!(
|
||||
"failed to terminate ACPX turn at the suspension boundary: {error}"
|
||||
))
|
||||
})?;
|
||||
self.session = None;
|
||||
let state = self
|
||||
.state
|
||||
.as_mut()
|
||||
.expect("ACPX state remains available after provider termination");
|
||||
state.active_turn_id = None;
|
||||
state.lifecycle = "suspended".to_owned();
|
||||
self.save_state()?;
|
||||
Ok(CommandExecution::result(json!({
|
||||
"status": "stopped",
|
||||
"providerTurnId": turn_id,
|
||||
"reason": reason,
|
||||
"providerExitConfirmed": true,
|
||||
})))
|
||||
}
|
||||
|
||||
fn resolve_request(&mut self, payload: &Value) -> Result<CommandExecution, DurableRunnerError> {
|
||||
let request_id = payload
|
||||
.get("requestId")
|
||||
|
|
@ -1171,9 +1210,8 @@ impl CommandExecutor for AcpxCommandExecutor {
|
|||
"code": "provider_command_unavailable",
|
||||
"message": "ACPX does not support steering an active turn",
|
||||
}))),
|
||||
"turn.interrupt" | "turn.stop" | "run.cancel" => {
|
||||
self.interrupt_turn(&command.command_type)
|
||||
}
|
||||
"turn.interrupt" | "run.cancel" => self.interrupt_turn(&command.command_type),
|
||||
"turn.stop" => self.stop_turn_for_suspension(&command.command_type),
|
||||
"request.resolve" => self.resolve_request(&command.payload),
|
||||
"semantic_tool.result" => self.deliver_tool_result(&command.payload),
|
||||
"session.snapshot" => self.snapshot(),
|
||||
|
|
|
|||
|
|
@ -695,6 +695,30 @@ impl AcpxProviderSession {
|
|||
}
|
||||
}
|
||||
|
||||
/// Reaps an active provider generation at a controller-owned suspension
|
||||
/// boundary without waiting for the provider's graceful close protocol.
|
||||
///
|
||||
/// A governed Paperclip result can settle the run while the model is still
|
||||
/// waiting for its semantic-tool callback to unwind. In that state the
|
||||
/// ordinary sidecar close path may wait for the callback longer than the
|
||||
/// server process that owns this runner. Process-group termination is the
|
||||
/// fail-closed boundary: only after it proves the old generation is gone
|
||||
/// may runnerd persist the durable session as attachable by the next run.
|
||||
pub fn terminate_active_turn_for_suspension(
|
||||
&mut self,
|
||||
turn_id: &str,
|
||||
) -> Result<(), LocalRunnerError> {
|
||||
self.ensure_open()?;
|
||||
validate_stable_id(turn_id, DURABLE_STABLE_ID_CHARS, "ACPX turn id")?;
|
||||
if self.state.active_turn_id() != Some(turn_id) {
|
||||
return Err(LocalRunnerError::invalid(
|
||||
"ACPX suspension termination named a stale or inactive turn",
|
||||
));
|
||||
}
|
||||
self.closed = true;
|
||||
self.terminate_transport()
|
||||
}
|
||||
|
||||
fn terminate_transport(&mut self) -> Result<(), LocalRunnerError> {
|
||||
if self.transport_terminated {
|
||||
return Ok(());
|
||||
|
|
|
|||
|
|
@ -2120,6 +2120,80 @@ impl CodexCommandExecutor {
|
|||
})))
|
||||
}
|
||||
|
||||
fn stop_turn_for_suspension(
|
||||
&mut self,
|
||||
reason: &str,
|
||||
) -> Result<CommandExecution, DurableRunnerError> {
|
||||
self.restore_provider_if_needed()?;
|
||||
let provider_turn_id = self
|
||||
.state
|
||||
.as_ref()
|
||||
.and_then(|state| state.active_provider_turn_id.clone());
|
||||
let Some(provider_turn_id) = provider_turn_id else {
|
||||
return Ok(CommandExecution::result(json!({
|
||||
"status": "already_settled",
|
||||
"reason": reason,
|
||||
})));
|
||||
};
|
||||
|
||||
// The cooperative interrupt is useful to the provider, but its RPC
|
||||
// acknowledgement is not proof that an active turn stopped. A
|
||||
// controller issues turn.stop only while closing a run whose result is
|
||||
// already durable, so terminate the exact process generation before
|
||||
// publishing the provider state as attachable by a successor run.
|
||||
let interrupt_accepted = self.interrupt_turn(reason).is_ok();
|
||||
let provider_shutdown_failed = self
|
||||
.provider
|
||||
.as_mut()
|
||||
.is_some_and(|provider| provider.shutdown().is_err());
|
||||
if provider_shutdown_failed {
|
||||
return Err(DurableRunnerError::invalid(
|
||||
"failed to prove provider termination at the suspension boundary",
|
||||
));
|
||||
}
|
||||
self.provider = None;
|
||||
|
||||
let identity = self.event_identity()?;
|
||||
let state = self
|
||||
.state
|
||||
.as_mut()
|
||||
.expect("Codex state remains available after provider termination");
|
||||
state.settle_active_provider_turn_identity()?;
|
||||
let settled = state
|
||||
.tool_bridge
|
||||
.settle_turn("provider_turn_stopped_for_suspension")
|
||||
.map_err(|error| {
|
||||
DurableRunnerError::invalid(format!(
|
||||
"failed to settle semantic tools at the suspension boundary: {error}"
|
||||
))
|
||||
})?;
|
||||
for result in settled {
|
||||
state.push_terminal_event(semantic_result_event(&identity, &result))?;
|
||||
}
|
||||
state.active_provider_turn_id = None;
|
||||
state.ambiguous_turn_start_pending = false;
|
||||
state.completed_turn_authoritative = false;
|
||||
state.completed_turn_process_generation = None;
|
||||
state.completed_provider_turn_id = None;
|
||||
state.receipt_limit_diagnostic_emitted = false;
|
||||
state.receipt_limit_interrupt_pending = false;
|
||||
state.receipt_limit_interrupt_accepted = false;
|
||||
state.receipt_limit_interrupt_attempts = 0;
|
||||
state.receipt_limit_interrupt_deadline_unix_ms = None;
|
||||
state.active_provider_result_fingerprint = None;
|
||||
state.active_provider_result_disposition = None;
|
||||
state.last_agent_message = None;
|
||||
state.lifecycle = "provider_exited".to_owned();
|
||||
self.save_state()?;
|
||||
Ok(CommandExecution::result(json!({
|
||||
"status": "stopped",
|
||||
"providerTurnId": provider_turn_id,
|
||||
"reason": reason,
|
||||
"interruptAccepted": interrupt_accepted,
|
||||
"providerExitConfirmed": true,
|
||||
})))
|
||||
}
|
||||
|
||||
fn steer_turn(&mut self, payload: &Value) -> Result<CommandExecution, DurableRunnerError> {
|
||||
let text = payload
|
||||
.get("text")
|
||||
|
|
@ -2912,9 +2986,8 @@ impl CommandExecutor for CodexCommandExecutor {
|
|||
"session.open" => self.open_session(),
|
||||
"turn.start" => self.start_turn(&command.payload),
|
||||
"turn.steer" => self.steer_turn(&command.payload),
|
||||
"turn.interrupt" | "turn.stop" | "run.cancel" => {
|
||||
self.interrupt_turn(&command.command_type)
|
||||
}
|
||||
"turn.interrupt" | "run.cancel" => self.interrupt_turn(&command.command_type),
|
||||
"turn.stop" => self.stop_turn_for_suspension(&command.command_type),
|
||||
"request.resolve" => self.resolve_request(&command.payload),
|
||||
"semantic_tool.result" => self.deliver_semantic_result(&command.payload),
|
||||
"session.snapshot" => self.snapshot(),
|
||||
|
|
|
|||
|
|
@ -59,6 +59,19 @@ fn rejects_suspension_during_an_active_turn_without_closing_the_session() {
|
|||
session.shutdown("test complete").unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reaps_an_active_provider_generation_at_the_suspension_boundary() {
|
||||
let mut session = AcpxProviderSession::start(&config("suspend")).unwrap();
|
||||
session
|
||||
.start_turn("turn-1", "Please help", &std::env::temp_dir())
|
||||
.unwrap();
|
||||
|
||||
session
|
||||
.terminate_active_turn_for_suspension("turn-1")
|
||||
.unwrap();
|
||||
assert!(session.shutdown("already terminated").is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn fails_closed_when_the_suspension_acknowledgement_does_not_match() {
|
||||
for mode in ["suspend-wrong-ack", "suspend-wrong-identity"] {
|
||||
|
|
|
|||
|
|
@ -2079,3 +2079,44 @@ it("rejects the notification stream promptly when runnerd exits after accepting
|
|||
await rm(stateDirectory, { recursive: true, force: true });
|
||||
}
|
||||
}, 30_000);
|
||||
|
||||
it("persists an active provider as settled before bounded suspension", async () => {
|
||||
const stateDirectory = await mkdtemp(
|
||||
join(tmpdir(), "runnerd-active-suspension-"),
|
||||
);
|
||||
const bundle = createCapabilityRunnerdCodexTransport({
|
||||
runnerBinary: defaultCapabilityRunnerdBinary(),
|
||||
codexCommand: fakeCodex,
|
||||
codexArgs: fakeCodexArgs(stateDirectory, "--linger-after-turn-start"),
|
||||
stateDirectory,
|
||||
});
|
||||
bundle.transport.setServerRequestHandler(async () => ({
|
||||
success: true,
|
||||
contentItems: [],
|
||||
}));
|
||||
try {
|
||||
await bundle.transport.request("initialize", {});
|
||||
await bundle.transport.request("thread/start", {
|
||||
cwd: tmpdir(),
|
||||
dynamicTools: [],
|
||||
});
|
||||
await bundle.transport.request("turn/start", {
|
||||
input: [{ type: "text", text: "Wait for another instruction." }],
|
||||
});
|
||||
await bundle.transport.close();
|
||||
|
||||
const providerState = JSON.parse(
|
||||
await readFile(
|
||||
join(stateDirectory, "runner", "codex-provider-state.json"),
|
||||
"utf8",
|
||||
),
|
||||
) as Record<string, unknown>;
|
||||
expect(providerState).toMatchObject({
|
||||
lifecycle: "provider_exited",
|
||||
activeProviderTurnId: null,
|
||||
});
|
||||
} finally {
|
||||
await bundle.transport.close();
|
||||
await rm(stateDirectory, { recursive: true, force: true });
|
||||
}
|
||||
}, 30_000);
|
||||
|
|
|
|||
Loading…
Reference in New Issue