Expose negotiated ACP turn controls through durable and direct sessions

Propagate live Pi steering and native queued follow-up capabilities after session open and lazy warm prompt initialization. Preserve conservative static declarations, reject forged capabilities and malformed control modes, and admit only profile-bound candidate credentials at the Rust sidecar boundary.

Co-Authored-By: Paperclip <noreply@paperclip.ing>
This commit is contained in:
DottaandPaperclip committed 2026-09-28 12:05:11 -05:00
1 parent 51db509899
commit a5ff475ecb
22 files changed
+417 -40

No files matched your search

@@ -14,7 +14,7 @@ use sha2::{Digest, Sha256};
use crate::acpx_provider_session::{
AcpxPermissionMode, AcpxProviderRuntimePolicy, AcpxProviderSession, AcpxProviderSessionConfig,
AcpxProviderSessionIdentity,
AcpxProviderSessionIdentity, AcpxTurnControlCapabilities,
};
use crate::acpx_sidecar_transport::AcpxSidecarTransportConfig;
#[cfg(test)]
@@ -422,6 +422,8 @@ struct AcpxDurableState {
#[serde(default)]
goal_projection: Value,
#[serde(default)]
turn_controls: AcpxTurnControlCapabilities,
#[serde(default)]
goal_revision: u64,
#[serde(default)]
goal_source_revision: Option<u64>,
@@ -451,6 +453,7 @@ impl AcpxDurableState {
provider_exit_unconfirmed: false,
semantic_result: None,
goal_projection: Value::Null,
turn_controls: AcpxTurnControlCapabilities::default(),
goal_revision: 0,
goal_source_revision: None,
pending_events: VecDeque::new(),
@@ -823,15 +826,17 @@ impl AcpxCommandExecutor {
let session = self.start_session(true)?;
let identity = session.identity().clone();
let process_id = session.process_id();
let turn_controls = session.turn_control_capabilities();
let state = self
.state
.as_mut()
.expect("ACPX state remains available during recovery");
state.lifecycle = "session_open".to_owned();
state.turn_controls = turn_controls;
state.push(NormalizedProviderEvent {
event_type: "session.resumed".to_owned(),
priority: EventPriority::P0,
payload: session_event_payload(&state.descriptor, &identity, process_id),
payload: session_event_payload(&state.descriptor, &identity, process_id, turn_controls),
})?;
self.session = Some(session);
self.save_state()
@@ -1018,6 +1023,7 @@ impl AcpxCommandExecutor {
.expect("ACPX session exists after successful start");
let identity = session.identity().clone();
let process_id = session.process_id();
let turn_controls = session.turn_control_capabilities();
let resumed = self
.state
.as_ref()
@@ -1031,7 +1037,9 @@ impl AcpxCommandExecutor {
state.identity = Some(identity.clone());
state.active_turn_id = None;
state.lifecycle = "session_open".to_owned();
let payload = session_event_payload(&state.descriptor, &identity, process_id);
state.turn_controls = turn_controls;
let payload =
session_event_payload(&state.descriptor, &identity, process_id, turn_controls);
let goal = self.goal_control("session.goal.get", &json!({}))?;
self.save_state()?;
Ok(CommandExecution {
@@ -1129,6 +1137,7 @@ impl AcpxCommandExecutor {
projection["requestId"] = request_id.clone();
}
state.goal_projection = projection.clone();
let turn_controls = state.turn_controls;
self.save_state()?;
let event_type = match command {
"session.goal.clear" => "session.goal.cleared",
@@ -1141,7 +1150,7 @@ impl AcpxCommandExecutor {
(
"session.capabilities.updated".to_owned(),
EventPriority::P0,
json!({"sessionGoals":projection["sessionGoals"]}),
json!({"sessionGoals":projection["sessionGoals"], "turnControls":turn_controls}),
),
(event_type.to_owned(), EventPriority::P0, projection),
],
@@ -1219,6 +1228,12 @@ impl AcpxCommandExecutor {
.as_mut()
.expect("ACPX state exists after turn start");
state.lifecycle = "turn_active".to_owned();
state.turn_controls = self
.session
.as_ref()
.expect("live ACPX session")
.turn_control_capabilities();
let capabilities = json!({"sessionGoals": state.goal_projection["sessionGoals"], "turnControls":state.turn_controls});
self.save_state()?;
let mut events = Vec::with_capacity(if provider_process_will_be_replaced {
2
@@ -1241,6 +1256,11 @@ impl AcpxCommandExecutor {
),
));
}
events.push((
"session.capabilities.updated".to_owned(),
EventPriority::P0,
capabilities,
));
events.push((
"turn.started".to_owned(),
EventPriority::P0,
@@ -1270,10 +1290,15 @@ impl AcpxCommandExecutor {
.get("turnId")
.and_then(Value::as_str)
.ok_or_else(|| DurableRunnerError::invalid("turn.steer payload.turnId is required"))?;
let mode = payload
.get("mode")
.and_then(Value::as_str)
.unwrap_or("steer");
let mode = match payload.get("mode") {
None => "steer",
Some(Value::String(mode)) => mode.as_str(),
_ => {
return Err(DurableRunnerError::invalid(
"turn.steer payload.mode must be a string",
))
}
};
if !matches!(mode, "steer" | "follow_up")
|| text.trim().is_empty()
|| text.len() > 65_536
@@ -1941,11 +1966,14 @@ fn session_event_payload(
descriptor: &AcpxProviderDescriptor,
identity: &AcpxProviderSessionIdentity,
process_id: u32,
turn_controls: AcpxTurnControlCapabilities,
) -> Value {
let mut public_descriptor = descriptor.public_descriptor(Some(identity));
public_descriptor["turnControls"] = json!(turn_controls);
json!({
"provider": "acpx",
"driver": "acpx_runtime",
"providerDescriptor": descriptor.public_descriptor(Some(identity)),
"providerDescriptor": public_descriptor,
"runtimeIdentity": {
"executionKind": "local_process",
"processId": process_id,
@@ -2201,6 +2229,24 @@ mod tests {
.is_empty());
}
#[test]
fn turn_control_rejects_nonstring_mode_before_provider_access() {
let directory = temporary_directory("control-mode");
let config = test_config(&directory, None);
let mut executor = AcpxCommandExecutor::with_runner_config(&directory, &config);
for mode in [Value::Null, json!(1), json!(false), json!({})] {
let error = executor
.steer_turn(
"command-1",
&json!({"text":"change direction", "turnId":"turn-1", "mode":mode}),
)
.err()
.unwrap();
assert!(error.to_string().contains("mode must be a string"));
}
fs::remove_dir_all(directory).unwrap();
}
#[test]
fn retained_events_exposes_terminal_suffix_without_restoring_provider() {
let directory = temporary_directory("retained-terminal-suffix");
@@ -224,6 +224,30 @@ impl AcpxProviderSessionIdentity {
}
}
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct AcpxTurnControlCapabilities {
pub steering: bool,
pub queued_follow_up: bool,
}
fn verified_turn_controls(
value: Option<&Value>,
agent: &str,
) -> Result<AcpxTurnControlCapabilities, LocalRunnerError> {
let Some(value) = value else {
return Ok(AcpxTurnControlCapabilities::default());
};
let controls: AcpxTurnControlCapabilities = serde_json::from_value(value.clone())
.map_err(|_| LocalRunnerError::invalid("ACPX negotiated turn controls are malformed"))?;
if agent != "pi" && (controls.steering || controls.queued_follow_up) {
return Err(LocalRunnerError::invalid(
"ACPX profile cannot advertise these turn controls",
));
}
Ok(controls)
}
pub struct AcpxProviderSession {
transport: AcpxSidecarTransport,
config: AcpxProviderSessionConfig,
@@ -231,6 +255,7 @@ pub struct AcpxProviderSession {
tool_bridge: ProviderToolBridge,
reserved_tool_bridge: ProviderToolBridge,
identity: AcpxProviderSessionIdentity,
turn_controls: AcpxTurnControlCapabilities,
catalog_revision: u64,
working_directory: PathBuf,
closed: bool,
@@ -250,7 +275,7 @@ impl AcpxProviderSession {
let mut transport =
AcpxSidecarTransport::start_for_agent(&config.transport, &config.agent)?;
let bootstrap = bootstrap(&mut transport, config);
let (identity, state) = match bootstrap {
let (identity, state, turn_controls) = match bootstrap {
Ok(value) => value,
Err(error) => {
let cleanup = transport.shutdown();
@@ -264,6 +289,7 @@ impl AcpxProviderSession {
tool_bridge,
reserved_tool_bridge,
identity,
turn_controls,
catalog_revision: config.catalog_revision,
working_directory: config.working_directory.clone(),
closed: false,
@@ -279,6 +305,10 @@ impl AcpxProviderSession {
&self.identity
}
pub fn turn_control_capabilities(&self) -> AcpxTurnControlCapabilities {
self.turn_controls
}
pub fn state(&self) -> &AcpxProviderState {
&self.state
}
@@ -393,6 +423,11 @@ impl AcpxProviderSession {
"ACPX sidecar did not confirm the requested turn",
)));
}
self.turn_controls =
match verified_turn_controls(response.get("turnControls"), &self.config.agent) {
Ok(controls) => controls,
Err(error) => return Err(self.fail_closed(error)),
};
if let Err(error) = self.state.begin_turn(turn_id) {
return Err(self.fail_closed(error));
}
@@ -420,6 +455,16 @@ impl AcpxProviderSession {
"ACPX turn control violates its bounded contract",
));
}
let supported = if mode == "steer" {
self.turn_controls.steering
} else {
self.turn_controls.queued_follow_up
};
if !supported {
return Err(LocalRunnerError::invalid(
"ACPX turn control was not negotiated",
));
}
self.state.reserve_turn_control(turn_id, control_id)?;
// The sidecar checks the negotiated live capability. An error may follow
// delivery, so the reservation must never be released for automatic retry.
@@ -892,12 +937,13 @@ impl AcpxProviderSession {
&restart_config.transport,
&restart_config.agent,
)?;
let (replacement_identity, _) = match bootstrap(&mut replacement, &restart_config) {
Ok(value) => value,
Err(error) => {
return Err(self.reject_replacement(replacement, error));
}
};
let (replacement_identity, _, turn_controls) =
match bootstrap(&mut replacement, &restart_config) {
Ok(value) => value,
Err(error) => {
return Err(self.reject_replacement(replacement, error));
}
};
if replacement_identity != self.identity {
return Err(self.reject_replacement(
replacement,
@@ -907,6 +953,7 @@ impl AcpxProviderSession {
));
}
self.transport = replacement;
self.turn_controls = turn_controls;
self.transport_terminated = false;
Ok(())
}
@@ -1143,7 +1190,14 @@ impl Drop for AcpxProviderSession {
fn bootstrap(
transport: &mut AcpxSidecarTransport,
config: &AcpxProviderSessionConfig,
) -> Result<(AcpxProviderSessionIdentity, AcpxProviderState), LocalRunnerError> {
) -> Result<
(
AcpxProviderSessionIdentity,
AcpxProviderState,
AcpxTurnControlCapabilities,
),
LocalRunnerError,
> {
let sidecar_tools = sidecar_run_tool_operations(&config.tool_set);
let initialized = transport.request(
GeneratedAcpxSidecarCommand::Initialize,
@@ -1169,6 +1223,7 @@ fn bootstrap(
}),
)?;
let identity = verify_open_response(&opened, transport.process_id(), config)?;
let turn_controls = verified_turn_controls(opened.get("turnControls"), &config.agent)?;
let attached = transport.request(
GeneratedAcpxSidecarCommand::RunAttach,
@@ -1185,7 +1240,11 @@ fn bootstrap(
"ACPX sidecar did not confirm the requested run attachment",
));
}
Ok((identity, AcpxProviderState::new(&config.run_id)?))
Ok((
identity,
AcpxProviderState::new(&config.run_id)?,
turn_controls,
))
}
fn verify_initialize_response(value: &Value, process_id: u32) -> Result<(), LocalRunnerError> {
@@ -1377,6 +1436,35 @@ fn with_cleanup_error(
#[cfg(test)]
mod tests {
#[test]
fn turn_controls_require_exact_live_pi_capability_fields() {
use super::*;
assert_eq!(
verified_turn_controls(None, "pi").unwrap(),
AcpxTurnControlCapabilities::default()
);
assert!(
verified_turn_controls(Some(&json!({"steering":true,"queuedFollowUp":true})), "pi")
.unwrap()
.steering
);
for value in [
Value::Null,
json!({}),
json!({"steering":1,"queuedFollowUp":false}),
json!({"steering":true,"queuedFollowUp":true,"extra":true}),
] {
assert!(verified_turn_controls(Some(&value), "pi").is_err());
}
for agent in ["codex", "claude", "cursor", "copilot"] {
assert!(verified_turn_controls(
Some(&json!({"steering":true,"queuedFollowUp":false})),
agent
)
.is_err());
}
}
use super::*;
#[test]
@@ -98,9 +98,12 @@ impl AcpxSidecarTransport {
let credential_keys: &[&str] = match agent {
"claude" => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
"codex" => &["OPENAI_API_KEY", "CODEX_API_KEY"],
"pi" => &["OPENROUTER_API_KEY"],
"cursor" => &["CURSOR_API_KEY", "CURSOR_AUTH_TOKEN"],
"copilot" => &["COPILOT_GITHUB_TOKEN"],
_ => {
return Err(LocalRunnerError::invalid(
"ACPX sidecar credentials require a qualified claude or codex agent",
"ACPX sidecar credentials require a known agent profile",
))
}
};
@@ -148,6 +148,7 @@ fn run() -> Result<(), Box<dyn std::error::Error>> {
"url": std::env::var("PAPERCLIP_NATIVE_MCP_URL").ok(),
"hasToken": std::env::var("PAPERCLIP_NATIVE_MCP_TOKEN").is_ok(),
"hasUnrelatedSecret": std::env::var("UNRELATED_EVAL_SECRET").is_ok(),
"credentialKeys": (["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN", "OPENAI_API_KEY", "CODEX_API_KEY", "OPENROUTER_API_KEY", "CURSOR_API_KEY", "CURSOR_AUTH_TOKEN", "COPILOT_GITHUB_TOKEN", "GITHUB_TOKEN", "GH_TOKEN"].into_iter().filter(|key| std::env::var(key).is_ok()).collect::<Vec<_>>()),
}
}),
)?;
@@ -206,6 +207,7 @@ fn run() -> Result<(), Box<dyn std::error::Error>> {
| "bootstrap-wrong-run"
| "controls"
| "controls-wrong-ack"
| "controls-lazy"
| "turns"
| "turns-wrong-turn"
| "turns-wrong-cancel"
@@ -730,6 +732,7 @@ fn bootstrap_success(
"providerLifetimeFenceCandidates": [60001, 60002, 60003],
},
"status": {},
"turnControls": {"steering": matches!(mode, "controls" | "controls-wrong-ack"), "queuedFollowUp":matches!(mode, "controls" | "controls-wrong-ack")},
})
}
"run.attach" => json!({
@@ -737,6 +740,7 @@ fn bootstrap_success(
"catalogRevision": params.get("catalogRevision"),
}),
"turn.start" => json!({
"turnControls": {"steering": matches!(mode, "controls" | "controls-wrong-ack" | "controls-lazy"), "queuedFollowUp":matches!(mode, "controls" | "controls-wrong-ack" | "controls-lazy")},
"turnId": if mode == "turns-wrong-turn" { "wrong-turn" } else { params.get("turnId").and_then(Value::as_str).unwrap_or("missing") },
}),
"turn.steer" => json!({
@@ -36,8 +36,18 @@ fn config(mode: &str) -> AcpxProviderSessionConfig {
request_timeout: Duration::from_secs(1),
shutdown_grace: Duration::from_millis(100),
},
agent: "codex".to_owned(),
model: "gpt-5.6-sol".to_owned(),
agent: if mode.starts_with("controls") {
"pi"
} else {
"codex"
}
.to_owned(),
model: if mode.starts_with("controls") {
"openrouter/deepseek/deepseek-v4-flash-0731"
} else {
"gpt-5.6-sol"
}
.to_owned(),
run_id: "run-1".to_owned(),
catalog_revision: 1,
runtime_directory: std::env::temp_dir(),
@@ -45,7 +55,15 @@ fn config(mode: &str) -> AcpxProviderSessionConfig {
working_directory: std::env::temp_dir(),
permission_mode: AcpxPermissionMode::ApproveReads,
permission_mode_pinned: true,
provider_policy: None,
provider_policy: if mode.starts_with("controls") {
Some(
paperclip_runner_core::acpx_provider_session::AcpxProviderRuntimePolicy {
read_only: false,
},
)
} else {
None
},
system_instructions: "Complete the supplied task.".to_owned(),
runtime_context: serde_json::Value::Null,
tool_set: tool_set(),
@@ -744,3 +762,19 @@ fn turn_controls_fail_closed_on_mismatched_acknowledgement() {
.steer_turn("turn-1", "control-2", "follow_up", "Do not replay")
.is_err());
}
#[test]
fn lazy_warm_handshake_updates_live_turn_control_discovery() {
let mut session = AcpxProviderSession::start(&config("controls-lazy")).unwrap();
assert!(!session.turn_control_capabilities().steering);
assert!(!session.turn_control_capabilities().queued_follow_up);
session
.start_turn("turn-lazy", "Work", &std::env::temp_dir())
.unwrap();
assert!(session.turn_control_capabilities().steering);
assert!(session.turn_control_capabilities().queued_follow_up);
session
.steer_turn("turn-lazy", "control-1", "follow_up", "Then validate")
.unwrap();
session.shutdown("verified live handshake").unwrap();
}
@@ -233,12 +233,28 @@ fn assigned_gateway_binding_reaches_qualified_sidecar_without_unrelated_secrets(
"fixture-token-never-returned-in-test-output",
)
.env("UNRELATED_EVAL_SECRET", "must-not-cross-boundary")
.envs(
[
"ANTHROPIC_API_KEY",
"CLAUDE_CODE_OAUTH_TOKEN",
"OPENAI_API_KEY",
"CODEX_API_KEY",
"OPENROUTER_API_KEY",
"CURSOR_API_KEY",
"CURSOR_AUTH_TOKEN",
"COPILOT_GITHUB_TOKEN",
"GITHUB_TOKEN",
"GH_TOKEN",
]
.into_iter()
.map(|key| (key, "fixture-credential")),
)
.status()
.unwrap();
assert!(status.success(), "isolated gateway environment test failed");
return;
}
for agent in ["claude", "codex"] {
for agent in ["claude", "codex", "pi", "cursor", "copilot"] {
let mut sidecar = AcpxSidecarTransport::start_for_agent(
&AcpxSidecarTransportConfig {
command: PathBuf::from(env!("CARGO_BIN_EXE_fake-acpx-sidecar")),
@@ -260,6 +276,15 @@ fn assigned_gateway_binding_reaches_qualified_sidecar_without_unrelated_secrets(
"assigned gateway credential was dropped"
);
assert_eq!(response["hasUnrelatedSecret"], false);
let expected = match agent {
"claude" => vec!["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
"codex" => vec!["OPENAI_API_KEY", "CODEX_API_KEY"],
"pi" => vec!["OPENROUTER_API_KEY"],
"cursor" => vec!["CURSOR_API_KEY", "CURSOR_AUTH_TOKEN"],
"copilot" => vec!["COPILOT_GITHUB_TOKEN"],
_ => unreachable!(),
};
assert_eq!(response["credentialKeys"], json!(expected));
sidecar.shutdown().unwrap();
}
}
@@ -97,6 +97,24 @@ const driver: HarnessDriver = {
};
describe("HarnessDriverBackend", () => {
it("refreshes negotiated controls after a warm handshake and separates follow-ups", async () => {
let controls = { steering: false, queuedFollowUp: false };
const steer = vi.fn(async () => {});
const session = Object.assign(new FakeHarnessSession(), { steer, turnControlCapabilities: () => ({ ...controls }) });
const backend = new HarnessDriverBackend({ ...driver, openSession: async () => session });
const opened = await backend.openSession({ identity: {
runId: "run-1", sessionId: "session-1", companyId: "company-1", issueId: "issue-1", agentId: "agent-1",
} });
expect(await opened.capabilities()).toMatchObject({ steering: false, queuedFollowUp: false });
controls = { steering: false, queuedFollowUp: true };
expect(await opened.capabilities()).toMatchObject(controls);
expect(() => opened.steer!({ turnId: "turn-1", message: { role: "user", text: "steer" }, correlationId: "c-1" })).toThrow();
await opened.steer!({ turnId: "turn-1", message: { role: "user", text: "follow" }, correlationId: "c-2", mode: "follow_up" });
expect(steer).toHaveBeenCalledWith(expect.objectContaining({ mode: "follow_up" }));
controls = { steering: true, queuedFollowUp: false };
expect(await opened.capabilities()).toMatchObject(controls);
expect(() => opened.steer!({ turnId: "turn-1", message: { role: "user", text: "follow" }, correlationId: "c-3", mode: "follow_up" })).toThrow();
});
it.each([false, true])("honors driver steering support even when the transport exposes a steer method (%s)", async supported => {
const steer = vi.fn(async () => { throw new Error("provider does not support steering"); });
const session = Object.assign(new FakeHarnessSession(), { steer });
@@ -450,10 +450,12 @@ class HarnessNativeSession implements NativeSession {
}
async capabilities() {
const negotiated = this.#session.turnControlCapabilities?.();
return {
resume: true,
typedEvents: true,
steering: this.#steeringSupported && this.#session.steer !== undefined,
steering: (negotiated?.steering ?? this.#steeringSupported) && this.#session.steer !== undefined,
queuedFollowUp: negotiated?.queuedFollowUp === true && this.#session.steer !== undefined,
interruption: this.#session.interrupt !== undefined,
structuredResult: true,
read: this.#session.read !== undefined,
@@ -678,13 +680,15 @@ class HarnessNativeSession implements NativeSession {
}
steer(input: {
mode?: "steer" | "follow_up";
turnId: string;
message: { role: "user"; text: string };
correlationId?: string;
}) {
this.#assertProtocolIntegrity();
if (!this.#steeringSupported || this.#session.steer === undefined)
throw new Error("steering is unavailable");
const negotiated = this.#session.turnControlCapabilities?.();
const supported = input.mode === "follow_up" ? negotiated?.queuedFollowUp === true : (negotiated?.steering ?? this.#steeringSupported);
if (!supported || this.#session.steer === undefined) throw new Error("steering or queued follow-up is unavailable");
return this.#withProtocolIntegrity(() => this.#session.steer!(input));
}
@@ -4,6 +4,7 @@ import { fileURLToPath } from "node:url";
import { afterEach, describe, expect, it, vi } from "vitest";
import { normalizeAcpxPermission } from "../drivers/acpx/acp-permission-adapter.js";
import { ACPX_SIDECAR_PROTOCOL_VERSION } from "../drivers/acpx/sidecar-protocol.js";
import {
awaitSidecarCleanupWithin,
@@ -30,6 +31,29 @@ afterEach(async () => {
});
describe("qualified ACPX runtime sidecar", () => {
it.each(["pi", "copilot", "cursor", "codex", "claude"])("offers verified session permission grants only for %s", async agent => {
const source = readFileSync(fileURLToPath(new URL("./acpx-runtime-sidecar.ts", import.meta.url)), "utf8");
const start = source.indexOf(" const { signal } = context;", source.indexOf("async function waitForPermission"));
const end = source.indexOf("\nasync function waitForInput", start);
expect(start).toBeGreaterThan(0);
const permissions = new Map<string, unknown>();
const emitted: Array<{ choices: Array<{ key: string }> }> = [];
const wait = new Function("permissions", "openParams", "normalizeAcpxPermission", "emit",
`let turnId = "turn-1", requestSequence = 0; const MAX_PENDING_INPUTS = 512;
const stableRequestId = () => "request-1"; const requireAcpxResponseDelivery = c => c.responseDelivery;
return async function(activeTurnId, request, context) { ${source.slice(start, end)}`)(
permissions, { agent }, normalizeAcpxPermission, (_event: string, payload: { choices: Array<{ key: string }> }) => emitted.push(payload),
);
const abort = new AbortController();
const pending = wait("turn-1", { sessionId: "session", inferredKind: "edit", raw: {
sessionId: "session", toolCall: { toolCallId: "call", title: "Edit file" },
options: ["allow_once", "allow_always", "reject_once"].map(kind => ({ kind, optionId: kind, name: kind })),
} }, { signal: abort.signal, responseDelivery: Promise.resolve() });
expect(emitted[0]!.choices.some(choice => choice.key === "accept_for_session")).toBe(agent === "pi" || agent === "copilot");
abort.abort();
await expect(pending).resolves.toEqual({ outcome: "cancel" });
expect(permissions.size).toBe(0);
});
it.each(["paperclip_finish", "paperclip_block"])(
"bounds pending %s calls before reserved handling and resumes admission",
async (operationId) => {
@@ -329,6 +329,7 @@ async function dispatch(
),
sidecarPid: process.pid,
status: opened.status,
turnControls: openedHost.steeringCapability() ?? { steering: false, queuedFollowUp: false },
};
}
if (request.command === "run.attach") {
@@ -382,7 +383,10 @@ async function dispatch(
throw error;
}
void pumpTurn(currentTurnId, runtimeTurn, activeHost, usageBefore, extensions.drain);
return { turnId: currentTurnId };
// Warm sessions may defer initialize until their first prompt. Publish only
// the capabilities of that live initialized connection, never old disk state.
await runtimeTurn.promptStarted;
return { turnId: currentTurnId, turnControls: activeHost.steeringCapability() ?? { steering: false, queuedFollowUp: false } };
}
if (request.command === "turn.steer") {
const control = parseAcpxTurnControl(request.params);
@@ -738,7 +742,7 @@ async function waitForPermission(
if (turnId !== activeTurnId || signal.aborted || permissions.size >= MAX_PENDING_INPUTS) {
return { outcome: "cancel" };
}
const normalized = normalizeAcpxPermission(request, openParams?.agent === "pi" ? { allowAlwaysScope: "session" } : {});
const normalized = normalizeAcpxPermission(request, ["pi", "copilot"].includes(openParams?.agent ?? "") ? { allowAlwaysScope: "session" } : {});
const responseDelivery = requireAcpxResponseDelivery(context);
const requestId = stableRequestId(activeTurnId, ++requestSequence, normalized.toolCallId);
return await new Promise((settle) => {
@@ -2,7 +2,7 @@ import type {
PrpEvent,
PrpStructuredRunResult,
} from "../protocol/replay-contract.js";
import type { NativeSessionCapabilities, NativeUserMessage } from "./types.js";
import type { NativeSessionCapabilities, NativeTurnControlCapabilities, NativeUserMessage } from "./types.js";
import {
PAPERCLIP_RUNTIME_REQUEST_SCHEMA_V2,
parsePaperclipQuestionResponse,
@@ -495,6 +495,7 @@ export interface HarnessSessionRecoveryResult {
}
export interface HarnessSession {
turnControlCapabilities?(): NativeTurnControlCapabilities | null;
ids(): {
driverSessionId: string;
providerSessionId?: string | null;
@@ -154,6 +154,7 @@ export interface NativeSession {
effectiveCollaborationMode?: "default" | "plan";
}>;
steer?(input: {
mode?: "steer" | "follow_up";
turnId: string;
message: NativeUserMessage;
correlationId?: string;
@@ -28,6 +28,7 @@ export interface NativeSessionCapabilities {
typedEvents: boolean;
typedEventFamilies?: TypedEventFamilyCapability[];
steering: boolean;
queuedFollowUp?: boolean;
interruption: boolean;
structuredResult: boolean;
read?: boolean;
@@ -47,3 +48,9 @@ export interface NativeUserMessage {
text: string;
}
import type { TypedEventFamilyCapability } from "../provider-events.js";
/** Live provider handshake; queued follow-up is distinct from active steering. */
export interface NativeTurnControlCapabilities {
steering: boolean;
queuedFollowUp: boolean;
}
@@ -23,10 +23,11 @@ import type { AcpxRecoveryWorkspaceLease } from "./runtime-sandbox.js";
describe("Codex ACPX harness driver", () => {
it("delivers negotiated steering and follow-up distinctly, once, for the active turn", async () => {
const fixture = driverFixture({ providerPolicy: { readOnly: true } });
const fixture = driverFixture({ agent: "pi", model: "openrouter/deepseek/deepseek-v4-flash-0731", providerPolicy: { readOnly: true } });
fixture.host.steeringCapability.mockReturnValue({ steering: true, queuedFollowUp: true });
const session = await fixture.driver.openSession({ runId: "run-controls", normalizedSessionId: "session-1", workingDirectory: "/workspace" });
expect(fixture.hostOptions?.providerPolicy).toEqual({ readOnly: true });
expect(session.turnControlCapabilities?.()).toEqual({ steering: true, queuedFollowUp: true });
const { turnId } = await session.startTurn({ message: { text: "Work" } });
const turnInput = fixture.host.startTurn.mock.calls[0]![0];
const context = { requestId: 0, signal: new AbortController().signal };
@@ -48,7 +49,7 @@ describe("Codex ACPX harness driver", () => {
});
it("rejects unnegotiated controls and retains ambiguous delivery attempts", async () => {
const fixture = driverFixture();
const fixture = driverFixture({ agent: "pi", model: "openrouter/deepseek/deepseek-v4-flash-0731", providerPolicy: { readOnly: true } });
const session = await fixture.driver.openSession({ runId: "run-controls", normalizedSessionId: "session-1", workingDirectory: "/workspace" });
const { turnId } = await session.startTurn({ message: { text: "Work" } });
await expect(session.steer!({ turnId, message: { text: "No handshake" } })).rejects.toThrow(/did not negotiate/);
@@ -62,6 +62,7 @@ import {
} from "./acp-question-adapter.js";
import {
acpxDriverDescriptor,
acpxCapabilities,
validateAcpxDriverConfig,
} from "./driver-profile.js";
import {
@@ -315,7 +316,7 @@ export class CodexAcpxDriver implements HarnessDriver {
resume: true,
runtimeRequestResolution: true,
runtimeRequestHandoff: true,
unsupported: ["steering", "goals", "threadLineage"],
unsupported: descriptor.capabilities.unsupported,
},
};
}
@@ -842,6 +843,11 @@ class CodexAcpxSession implements HarnessSession {
return this.#events;
}
turnControlCapabilities() {
const capability = acpxCapabilities(this.#agent, this.#host.steeringCapability?.());
return { steering: capability.steering, queuedFollowUp: capability.queuedFollowUp === true };
}
async startTurn(input: {
message: NativeUserMessage;
}): Promise<{ turnId: string }> {
@@ -923,7 +929,7 @@ class CodexAcpxSession implements HarnessSession {
this.#assertOpen();
if (input.turnId !== this.#activeTurnId || this.#pendingTerminal) throw new HarnessStaleTurnError(input.turnId);
const mode = input.mode ?? "steer";
const capability = this.#host.steeringCapability?.();
const capability = this.turnControlCapabilities();
const operation = mode === "steer" ? this.#host.steerActiveTurn : this.#host.queueFollowUp;
if (!(mode === "steer" ? capability?.steering : capability?.queuedFollowUp) || !operation) {
throw new HarnessCapabilityUnavailableError(mode, "the ACP provider did not negotiate this turn control");
@@ -7,6 +7,14 @@ import {
} from "./driver-profile.js";
describe("ACPX driver profile", () => {
it("enables only negotiated Pi controls and keeps the static profile conservative", () => {
expect(acpxCapabilities("pi")).toMatchObject({ steering: false, queuedFollowUp: false });
expect(acpxCapabilities("pi", { steering: true, queuedFollowUp: false })).toMatchObject({ steering: true, queuedFollowUp: false });
expect(acpxCapabilities("pi", { steering: false, queuedFollowUp: true })).toMatchObject({ steering: false, queuedFollowUp: true });
for (const agent of ["cursor", "copilot", "codex", "claude"] as const) {
expect(acpxCapabilities(agent, { steering: true, queuedFollowUp: true })).toMatchObject({ steering: false, queuedFollowUp: false });
}
});
it.each([
["codex", "available"],
["claude", "available"],
@@ -3,7 +3,7 @@ import type {
HarnessDriverDescriptor,
} from "../../contracts/harness-driver.js";
import type { NativeAcpxPermissionMode } from "../../contracts/native-execution.js";
import type { NativeSessionCapabilities } from "../../contracts/types.js";
import type { NativeSessionCapabilities, NativeTurnControlCapabilities } from "../../contracts/types.js";
import { providerFamilyCapabilities } from "../../provider-events.js";
import {
ACPX_DRIVER_KIND,
@@ -31,7 +31,9 @@ export interface ValidatedAcpxDriverConfig extends Record<string, unknown> {
export function acpxCapabilities(
agent: QualifiedAcpxAgent,
negotiated?: NativeTurnControlCapabilities | null,
): NativeSessionCapabilities {
const controls = agent === "pi" ? negotiated : null;
const profile = ACPX_CAPABILITY_PROFILES[agent];
return {
resume: profile.recovery === "session-load",
@@ -44,7 +46,8 @@ export function acpxCapabilities(
provider_notice: "available",
artifact: "policy_disabled",
}),
steering: false,
steering: controls?.steering === true,
queuedFollowUp: controls?.queuedFollowUp === true,
interruption: true,
structuredResult: true,
read: true,
@@ -55,7 +58,7 @@ export function acpxCapabilities(
runtimeRequestHandoff: true,
goals: false,
threadLineage: false,
unsupported: ["steering", "goals", "threadLineage"],
unsupported: [...(controls?.steering ? [] : ["steering"]), "goals", "threadLineage"],
};
}
@@ -1,3 +1,4 @@
import type { NativeTurnControlCapabilities } from "../../contracts/types.js";
import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
import type { HarnessRuntimeRequestResolution } from "../../contracts/harness-driver.js";
import { githubCredentialEnvironment } from "../../github-credential-environment.js";
@@ -30,6 +31,8 @@ export type CodexServerRequestHandler = (
) => Promise<Record<string, unknown>>;
export interface CodexAppServerTransport {
/** Runner-owned live capability projection; absent on native Codex transports. */
turnControlCapabilities?(): NativeTurnControlCapabilities | null;
request(
method: string,
params: Record<string, unknown>,
@@ -1486,6 +1486,38 @@ describe("Codex app-server Codex driver", () => {
).rejects.toThrow("cannot be a filesystem root");
});
it("uses negotiated ACP controls through the Codex transport facade", async () => {
let controls = { steering: false, queuedFollowUp: false };
const transport = Object.assign(new FakeCodexTransport(), { turnControlCapabilities: () => ({ ...controls }) });
const session = await makeDriver([transport], {
driverIdentity: { kind: "acpx_runtime", displayName: "Pi ACP", version: "test" },
capabilities: { steering: false },
}).openSession({ runId: "run-pi", normalizedSessionId: "session-pi", workingDirectory: TEST_WORKING_DIRECTORY });
try {
expect(session.turnControlCapabilities?.()).toEqual(controls);
const { turnId } = await session.startTurn({ message: { role: "user", text: "Work" } });
await expect(session.steer?.({ turnId, message: { role: "user", text: "Change" } })).rejects.toThrow("not negotiated");
controls = { steering: true, queuedFollowUp: true };
expect(session.turnControlCapabilities?.()).toEqual(controls);
const input = { turnId, correlationId: "queued-1", mode: "follow_up" as const, message: { role: "user" as const, text: "Then validate" } };
await session.steer?.(input);
expect(transport.calls.find(call => call.method === "turn/steer")?.params).toMatchObject({ expectedTurnId: turnId, mode: "follow_up", correlationId: "queued-1" });
await expect(session.steer?.(input)).rejects.toThrow("already acknowledged");
expect(transport.calls.filter(call => call.method === "turn/steer")).toHaveLength(1);
} finally { await session.close({ reason: "verified" }); }
});
it("keeps native follow-up disabled for the Codex app-server driver", async () => {
const transport = Object.assign(new FakeCodexTransport(), { turnControlCapabilities: () => ({ steering: true, queuedFollowUp: true }) });
const session = await makeDriver([transport]).openSession({ runId: "run-codex", normalizedSessionId: "session-codex", workingDirectory: TEST_WORKING_DIRECTORY });
try {
expect(session.turnControlCapabilities?.()).toEqual({ steering: true, queuedFollowUp: false });
const { turnId } = await session.startTurn({ message: { role: "user", text: "Work" } });
await expect(session.steer?.({ turnId, mode: "follow_up", message: { role: "user", text: "Queue" } })).rejects.toThrow("not negotiated");
expect(transport.calls.filter(call => call.method === "turn/steer")).toHaveLength(0);
} finally { await session.close({ reason: "verified" }); }
});
it("steers and interrupts an active turn without replacing the session", async () => {
const transport = new FakeCodexTransport();
const session = await makeDriver([transport]).openSession({
@@ -72,6 +72,13 @@ export class CodexHarnessSession
}
}
turnControlCapabilities() {
if (this.driverKind === "acpx_runtime") {
return this.transport.turnControlCapabilities?.() ?? { steering: false, queuedFollowUp: false };
}
return { steering: this.capabilities.steering, queuedFollowUp: false };
}
ids(): ReturnType<HarnessSession["ids"]> {
return {
driverSessionId: this.opened.threadId,
@@ -306,14 +313,17 @@ export class CodexHarnessSession
correlationId?: string;
}): Promise<void> {
this.assertProtocolIntegrity();
if (input.mode === "follow_up") throw this.unsupported("steering", "queued follow-up is not exposed by this driver");
this.requireCapability("steering");
const controls = this.turnControlCapabilities();
if (!(input.mode === "follow_up" ? controls.queuedFollowUp : controls.steering)) {
throw this.unsupported("steering", "requested turn control was not negotiated");
}
this.requireActiveTurn(input.turnId, "steering");
if (input.correlationId) {
const acknowledgedTurnId = this.acknowledgedSteeringCorrelations.get(
input.correlationId,
);
if (acknowledgedTurnId) {
if (this.driverKind === "acpx_runtime") throw new Error("ACP turn control correlation was already acknowledged");
if (acknowledgedTurnId !== input.turnId)
throw new HarnessOperationAlreadyTerminalError("steering");
return;
@@ -325,6 +335,7 @@ export class CodexHarnessSession
input: [userInput(input.message)],
expectedTurnId: input.turnId,
correlationId: input.correlationId,
...(input.mode === undefined ? {} : { mode: input.mode }),
});
if (this.activeTurnId !== input.turnId) {
throw new HarnessOperationAlreadyTerminalError("steering");
@@ -339,7 +350,8 @@ export class CodexHarnessSession
"item.completed",
{
kind: "steering_acknowledgement",
text: "Steering acknowledged for the active turn.",
text: input.mode === "follow_up" ? "Follow-up queued by the active provider." : "Steering acknowledged for the active turn.",
...(this.driverKind === "acpx_runtime" ? { mode: input.mode ?? "steer" } : {}),
status: "acknowledged",
},
{
@@ -56,6 +56,7 @@ import { releaseMaterializedNativeRuntimeSkills } from "../drivers/runtime-conte
import { RUNNERD_CANONICAL_ITEM } from "../drivers/codex/codex-driver-values.js";
import {
parseAcpxTurnControlCapabilities,
authorizedToolSetForProvider,
createCapabilityRunnerdCodexTransport,
createCapabilityRunnerdProviderEnvironment,
@@ -7205,3 +7206,15 @@ it("resolves explicit skills to the remote provider home and rejects unassigned
}
expect(() => resolveRunnerdCodexSkillInputs([skill], null, "/runner/codex-home")).toThrow("assigned runtime skill");
});
it("admits only exact live Pi turn controls", () => {
expect(parseAcpxTurnControlCapabilities(undefined, "pi")).toEqual({ steering: false, queuedFollowUp: false });
expect(parseAcpxTurnControlCapabilities({ steering: true, queuedFollowUp: true }, "pi")).toEqual({ steering: true, queuedFollowUp: true });
for (const value of [null, [], {}, { steering: 1, queuedFollowUp: false }, { steering: true, queuedFollowUp: "true" }, { steering: true, queuedFollowUp: true, arbitrary: true }]) {
expect(() => parseAcpxTurnControlCapabilities(value, "pi")).toThrow("malformed");
}
for (const agent of ["cursor", "copilot", "codex", "claude", undefined]) {
expect(() => parseAcpxTurnControlCapabilities({ steering: true, queuedFollowUp: false }, agent)).toThrow("cannot advertise");
}
});
@@ -38,7 +38,7 @@ import type {
DurableRecoveryCommittedEvent,
DurableRecoveryIdentity,
} from "../contracts/durable-recovery.js";
import type { NativeRunIdentity } from "../contracts/types.js";
import type { NativeRunIdentity, NativeTurnControlCapabilities } from "../contracts/types.js";
import type { PrpEvent } from "../protocol/replay-contract.js";
import { NativeSessionCloseUnrecoverableError } from "../contracts/native-session-backend.js";
import type {
@@ -3275,6 +3275,27 @@ function unwrapToolResponse(response: Record<string, unknown>): {
};
}
/** Only a live, admitted Pi ACP session can enable native turn controls. */
export function parseAcpxTurnControlCapabilities(
value: unknown,
agent: unknown,
): NativeTurnControlCapabilities {
const unsupported = { steering: false, queuedFollowUp: false };
if (value === undefined) return unsupported;
if (!value || typeof value !== "object" || Array.isArray(value)) {
throw new Error("ACPX turn control capabilities are malformed");
}
const controls = value as Record<string, unknown>;
if (Object.keys(controls).some(key => key !== "steering" && key !== "queuedFollowUp")
|| typeof controls.steering !== "boolean" || typeof controls.queuedFollowUp !== "boolean") {
throw new Error("ACPX turn control capabilities are malformed");
}
if (agent !== "pi" && (controls.steering || controls.queuedFollowUp)) {
throw new Error("ACPX profile cannot advertise these turn controls");
}
return { steering: controls.steering, queuedFollowUp: controls.queuedFollowUp };
}
class DurablePrpCodexTransport implements CodexAppServerTransport {
readonly #root: string;
readonly #ownsRoot: boolean;
@@ -3307,6 +3328,13 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
#threadId = "";
#sessionId: string | null = null;
#providerIdentity: Record<string, unknown> | null = null;
#turnControls: NativeTurnControlCapabilities = { steering: false, queuedFollowUp: false };
turnControlCapabilities(): NativeTurnControlCapabilities | null {
if (this.options.provider !== "acpx") return null;
if (this.#closed || this.#failure) return { steering: false, queuedFollowUp: false };
return { ...this.#turnControls };
}
#providerIdentityEventType:
"harness.ready" | "session.started" | "session.resumed" | null = null;
#checkpointProviderIdentityExpectation: {
@@ -5891,6 +5919,12 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
this.#applyProviderIdentityEvent(event);
continue;
}
if (event.eventType === "session.capabilities.updated" && this.options.provider === "acpx") {
const capabilities = record(record(event.envelope.payload).payload);
if (capabilities.turnControls !== undefined) {
this.#turnControls = parseAcpxTurnControlCapabilities(capabilities.turnControls, this.#evidence.acpxAgent);
}
}
if (event.eventType === "harness.diagnostic") {
const diagnostic = record(record(event.envelope.payload).payload);
if (
@@ -6171,6 +6205,12 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
} = resolveRunnerdSessionIdentity(started);
const providerIdentity = record(started.providerIdentity);
this.#confirmCheckpointProviderIdentity(started, event.eventType);
if (this.options.provider === "acpx") {
if (descriptor.turnControls !== undefined && descriptor.agent !== (this.options.acpxAgent ?? "codex")) {
throw new Error("ACPX capability identity differs from the admitted profile");
}
this.#turnControls = parseAcpxTurnControlCapabilities(descriptor.turnControls, descriptor.agent);
}
if (pid !== null) {
this.#evidence.providerPid = pid;
this.#evidence.providerProcessStartedAt = readLocalProcessStartedAt(pid);