use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use async_trait::async_trait;
use serde_json::json;
use supercode_harness::server::{FrontendRequestBridge, RpcEngine};
use supercode_harness::subagents::BackgroundPromptsPolicy;
use supercode_harness::{
Agent, ApprovalPolicy, ChatMessage, ChatRequest, ClaudeRuntimeManifest, Config,
FrontendApprovalDecision, FrontendElicitationAction, FrontendRequest, FrontendResponse,
FrontendRuntime, FrontendRuntimeMetadata, FrontendTurnState, FunctionCall, Provider, Role,
RuntimeSubmitError, Session, ToolCall, Usage,
};
fn submit_req(id: i64, prompt: &str) -> supercode_harness::server::RpcRequest {
let v = json!({"id": id, "method": "submit", "params": {"prompt": prompt}});
serde_json::from_value(v).unwrap()
}
fn req(id: i64, method: &str) -> supercode_harness::server::RpcRequest {
let v = json!({"id": id, "method": method});
serde_json::from_value(v).unwrap()
}
fn temp_dir(tag: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!(
"supercode-server-engine-{tag}-{}",
std::process::id()
));
std::fs::create_dir_all(&dir).unwrap();
dir
}
struct SaysProvider(String);
#[async_trait]
impl Provider for SaysProvider {
async fn complete(
&self,
_req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
Ok((ChatMessage::assistant(self.0.clone()), Usage::default()))
}
}
#[tokio::test]
async fn submit_returns_the_reply_and_streams_a_turn_completed_event() {
let dir = temp_dir("submit-ok");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("hi there".to_string())));
let engine = RpcEngine::new(agent, None);
let mut events = engine.subscribe();
let resp = engine.handle_request(submit_req(1, "hello")).await;
assert_eq!(resp["id"], 1);
assert_eq!(resp["result"]["reply"], "hi there");
let mut saw_turn_completed = false;
while let Ok(v) = events.try_recv() {
if v["type"] == "turn_completed" {
saw_turn_completed = true;
}
}
assert!(
saw_turn_completed,
"expected a turn_completed event notification"
);
}
#[tokio::test]
async fn history_replays_canonical_messages_for_late_frontends() {
let dir = temp_dir("history");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("late reply".to_string())));
let engine = RpcEngine::new(agent, None);
engine.handle_request(submit_req(1, "late prompt")).await;
let request = serde_json::from_value(json!({
"id": 2, "method": "history", "params": {"limit": 2}
}))
.unwrap();
let response = engine.handle_request(request).await;
let messages = response["result"]["messages"].as_array().unwrap();
assert_eq!(messages.len(), 2);
assert_eq!(messages[0]["role"], "user");
assert_eq!(messages[0]["content"], "late prompt");
assert_eq!(messages[1]["role"], "assistant");
assert_eq!(messages[1]["content"], "late reply");
}
#[tokio::test]
async fn frontend_descriptor_reports_runtime_identity_without_agent_access() {
let dir = temp_dir("frontend-descriptor");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("ok".into())));
let engine = RpcEngine::new_named_with_frontend_metadata(
agent,
"frontend-session",
FrontendRuntimeMetadata {
source_harness: Some("claude-code".into()),
emulation_profile: Some("cc-parity".into()),
},
None,
);
let descriptor = FrontendRuntime::describe(engine.as_ref()).await.unwrap();
assert_eq!(descriptor.session_id, "frontend-session");
assert_eq!(descriptor.source_harness.as_deref(), Some("claude-code"));
assert_eq!(descriptor.emulation_profile.as_deref(), Some("cc-parity"));
assert_eq!(descriptor.schema_version, 2);
assert_eq!(descriptor.operations.len(), 1);
assert_eq!(descriptor.operations[0].id, "prompt:code-review");
assert_eq!(descriptor.turn_state, FrontendTurnState::Idle);
assert!(descriptor.actions.submit);
assert!(descriptor.actions.interrupt);
assert!(descriptor.actions.steer);
assert!(
descriptor.actions.close,
"the SDK owner advertises typed close before authorization narrows it"
);
assert!(
!descriptor.actions.respond,
"unsupported actions must be honest"
);
assert!(descriptor.display.opaque_fallback);
let mut wire = serde_json::to_value(&descriptor).unwrap();
assert_eq!(wire["schema_version"], 2);
assert!(wire.get("operations").is_some());
assert_eq!(
serde_json::to_string(&descriptor.commands[0]).unwrap(),
r#"{"name":"code-review","description":null}"#,
"the schema-v1 command projection must retain its exact byte shape"
);
wire.as_object_mut().unwrap().remove("operations");
wire["schema_version"] = serde_json::json!(1);
let v1: supercode_harness::FrontendRuntimeDescriptor = serde_json::from_value(wire).unwrap();
assert_eq!(v1.schema_version, 1);
assert!(v1.operations.is_empty());
assert!(v1.commands[0].argument_hint.is_none());
}
struct GatedSteeringProvider {
calls: AtomicUsize,
started: std::sync::Arc<tokio::sync::Notify>,
release: std::sync::Arc<tokio::sync::Notify>,
}
#[async_trait]
impl Provider for GatedSteeringProvider {
async fn complete(
&self,
req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
self.started.notify_one();
self.release.notified().await;
Ok((
ChatMessage {
role: Role::Assistant,
content: None,
content_parts: None,
tool_calls: Some(vec![ToolCall {
id: "steer-tool".into(),
kind: "function".into(),
function: FunctionCall {
name: "list_dir".into(),
arguments: "{}".into(),
},
}]),
tool_call_id: None,
name: None,
metadata: Default::default(),
},
Usage::default(),
))
} else {
assert!(req.messages.iter().any(|message| {
message.role == Role::User && message.content.as_deref() == Some("change direction")
}));
Ok((ChatMessage::assistant("steered"), Usage::default()))
}
}
}
#[tokio::test]
async fn frontend_steer_reaches_an_active_agent_without_waiting_for_its_lock() {
let started = std::sync::Arc::new(tokio::sync::Notify::new());
let release = std::sync::Arc::new(tokio::sync::Notify::new());
let agent = Agent::with_provider(
Config::builder().cwd(temp_dir("frontend-steer")).build(),
Box::new(GatedSteeringProvider {
calls: AtomicUsize::new(0),
started: started.clone(),
release: release.clone(),
}),
);
let engine = RpcEngine::new(agent, None);
let submit_engine = engine.clone();
let submit = tokio::spawn(async move { submit_engine.submit("start").await });
tokio::time::timeout(Duration::from_secs(2), started.notified())
.await
.expect("provider should hold the active agent lock");
tokio::time::timeout(
Duration::from_millis(100),
FrontendRuntime::steer(engine.as_ref(), "change direction".into()),
)
.await
.expect("steering must not wait for the active agent lock")
.unwrap();
release.notify_one();
assert_eq!(submit.await.unwrap().unwrap(), "steered");
}
struct GatedStreamingProvider {
started: std::sync::Arc<tokio::sync::Notify>,
release: std::sync::Arc<tokio::sync::Notify>,
}
#[async_trait]
impl Provider for GatedStreamingProvider {
async fn complete(
&self,
_req: &ChatRequest,
on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
on_delta("partial");
self.started.notify_one();
self.release.notified().await;
Ok((ChatMessage::assistant("final"), Usage::default()))
}
}
#[tokio::test]
async fn frontend_attach_bridges_replay_to_live_without_duplicates() {
let dir = temp_dir("frontend-boundary");
let started = std::sync::Arc::new(tokio::sync::Notify::new());
let release = std::sync::Arc::new(tokio::sync::Notify::new());
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(
config,
Box::new(GatedStreamingProvider {
started: started.clone(),
release: release.clone(),
}),
);
let engine = RpcEngine::new_named(agent, "frontend-boundary", None);
let submit_engine = engine.clone();
let submit = tokio::spawn(async move { submit_engine.submit("hello").await });
tokio::time::timeout(Duration::from_secs(2), started.notified())
.await
.expect("provider should reach the deterministic mid-turn gate");
let mut first = engine.frontend_attach(20).unwrap();
let mut second = engine.frontend_attach(20).unwrap();
assert_eq!(first.descriptor.turn_state, FrontendTurnState::Busy);
assert_eq!(second.descriptor.turn_state, FrontendTurnState::Busy);
for attachment in [&mut first, &mut second] {
let user = attachment.next_event().await.unwrap();
let started = attachment.next_event().await.unwrap();
let delta = attachment.next_event().await.unwrap();
assert_eq!(user.kind, "user_message");
assert_eq!(user.payload["text"], "hello");
assert_eq!(started.kind, "turn_started");
assert_eq!(started.payload["schema_version"], 1);
assert_eq!(delta.kind, "text_delta");
assert_eq!(delta.payload["text"], "partial");
assert!(delta.sequence > user.sequence);
}
release.notify_one();
assert_eq!(submit.await.unwrap().unwrap(), "final");
for attachment in [&mut first, &mut second] {
let completed = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "turn_completed" {
break event;
}
}
})
.await
.expect("both consumers should observe the same completion");
assert_eq!(completed.payload["type"], "turn_completed");
}
let late = engine.frontend_attach(2).unwrap();
assert_eq!(late.history.len(), 2);
assert_eq!(late.history[0].role, Role::User);
assert_eq!(late.history[0].content.as_deref(), Some("hello"));
assert_eq!(late.history[1].role, Role::Assistant);
assert_eq!(late.history[1].content.as_deref(), Some("final"));
drop(first);
release.notify_one();
assert_eq!(engine.submit("still live").await.unwrap(), "final");
let second_completion = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = second.next_event().await.unwrap();
if event.kind == "turn_completed" {
break event;
}
}
})
.await
.expect("the remaining consumer should survive another turn");
assert_eq!(second_completion.payload["type"], "turn_completed");
assert!(!engine.status().busy);
}
#[tokio::test]
async fn frontend_lifecycle_wraps_the_complete_sdk_turn_with_additive_versioning() {
let agent = Agent::with_provider(
Config::builder()
.cwd(temp_dir("frontend-lifecycle"))
.build(),
Box::new(SaysProvider("done".into())),
);
let engine = RpcEngine::new_named(agent, "frontend-lifecycle", None);
let mut frontend = engine.frontend_attach(20).unwrap();
assert_eq!(engine.submit("go").await.unwrap(), "done");
let mut kinds = Vec::new();
loop {
let event = tokio::time::timeout(Duration::from_secs(2), frontend.next_event())
.await
.unwrap()
.unwrap();
kinds.push(event.kind.clone());
if event.kind == "turn_started" || event.kind == "turn_succeeded" {
assert_eq!(event.payload["schema_version"], 1);
}
if event.kind == "turn_succeeded" {
break;
}
}
assert_eq!(
kinds,
[
"user_message",
"turn_started",
"usage",
"turn_completed",
"turn_succeeded"
]
);
}
#[tokio::test]
async fn two_frontends_share_one_busy_turn_and_one_persistence_boundary() {
let started = std::sync::Arc::new(tokio::sync::Notify::new());
let release = std::sync::Arc::new(tokio::sync::Notify::new());
let persisted = std::sync::Arc::new(AtomicUsize::new(0));
let persisted_hook = persisted.clone();
let agent = Agent::with_provider(
Config::builder()
.cwd(temp_dir("frontend-multi-client"))
.build(),
Box::new(GatedStreamingProvider {
started: started.clone(),
release: release.clone(),
}),
);
let engine = RpcEngine::new_named(
agent,
"one-session",
Some(Box::new(move |_| {
persisted_hook.fetch_add(1, Ordering::SeqCst);
})),
);
let mut first = engine.frontend_attach(20).unwrap();
let mut second = engine.frontend_attach(20).unwrap();
let submitting = engine.clone();
let turn =
tokio::spawn(
async move { FrontendRuntime::submit(submitting.as_ref(), "one".into()).await },
);
tokio::time::timeout(Duration::from_secs(2), started.notified())
.await
.expect("first frontend owns the active turn");
assert!(matches!(
FrontendRuntime::submit(engine.as_ref(), "two".into()).await,
Err(supercode_harness::FrontendRuntimeError::Submit(
supercode_harness::server::RuntimeSubmitError::Busy
))
));
release.notify_one();
assert_eq!(turn.await.unwrap().unwrap(), "final");
for attachment in [&mut first, &mut second] {
tokio::time::timeout(Duration::from_secs(2), async {
loop {
if attachment.next_event().await.unwrap().kind == "turn_succeeded" {
break;
}
}
})
.await
.expect("both frontends observe the single result");
}
assert_eq!(persisted.load(Ordering::SeqCst), 1);
assert_eq!(engine.session_id(), "one-session");
}
#[tokio::test]
async fn unknown_method_returns_a_named_error_not_a_silent_no_op() {
let dir = temp_dir("unknown-method");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("x".to_string())));
let engine = RpcEngine::new(agent, None);
let resp = engine.handle_request(req(1, "frobnicate")).await;
assert_eq!(resp["id"], 1);
assert!(resp.get("result").is_none());
assert!(resp["error"]["message"]
.as_str()
.unwrap()
.contains("frobnicate"));
}
#[tokio::test]
async fn submit_requires_a_prompt_param() {
let dir = temp_dir("missing-prompt");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("x".to_string())));
let engine = RpcEngine::new(agent, None);
let bad = serde_json::from_value(json!({"id": 1, "method": "submit", "params": {}})).unwrap();
let resp = engine.handle_request(bad).await;
assert!(resp["error"]["message"]
.as_str()
.unwrap()
.contains("prompt"));
}
struct SlowProvider(Duration, String);
#[async_trait]
impl Provider for SlowProvider {
async fn complete(
&self,
_req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
tokio::time::sleep(self.0).await;
Ok((ChatMessage::assistant(self.1.clone()), Usage::default()))
}
}
struct TwoStageProvider(AtomicUsize);
#[async_trait]
impl Provider for TwoStageProvider {
async fn complete(
&self,
_req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
let call = self.0.fetch_add(1, Ordering::SeqCst);
if call == 0 {
tokio::time::sleep(Duration::from_millis(100)).await;
}
Ok((
ChatMessage::assistant(if call == 0 { "first" } else { "second" }),
Usage::default(),
))
}
}
struct AbortThenSaysProvider {
calls: AtomicUsize,
started: std::sync::Arc<tokio::sync::Notify>,
}
#[async_trait]
impl Provider for AbortThenSaysProvider {
async fn complete(
&self,
_req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
self.started.notify_one();
std::future::pending().await
} else {
Ok((ChatMessage::assistant("clean"), Usage::default()))
}
}
}
#[tokio::test]
async fn sdk_steer_is_current_turn_bound_and_idle_rejection_never_leaks_forward() {
let dir = temp_dir("steer-current-turn");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(TwoStageProvider(AtomicUsize::new(0))));
let engine = RpcEngine::new(agent, None);
let idle = FrontendRuntime::steer(engine.as_ref(), "must not leak".into()).await;
assert!(idle.is_err(), "idle steering must be rejected atomically");
let active_engine = engine.clone();
let active = tokio::spawn(async move {
FrontendRuntime::submit(active_engine.as_ref(), "start".into()).await
});
tokio::time::sleep(Duration::from_millis(25)).await;
FrontendRuntime::steer(engine.as_ref(), "change direction".into())
.await
.expect("active turn accepts steering");
assert_eq!(active.await.unwrap().unwrap(), "second");
let history = FrontendRuntime::attach(engine.as_ref(), 16)
.await
.unwrap()
.history;
let text = history
.iter()
.filter(|message| message.role != Role::System)
.filter_map(|message| message.content.as_deref())
.collect::<Vec<_>>();
assert_eq!(text, ["start", "first", "change direction", "second"]);
let after = FrontendRuntime::submit(engine.as_ref(), "next".into())
.await
.unwrap();
assert_eq!(after, "second");
let history = FrontendRuntime::attach(engine.as_ref(), 16)
.await
.unwrap()
.history;
assert!(history
.iter()
.filter_map(|message| message.content.as_deref())
.all(|message| message != "must not leak"));
}
#[tokio::test]
async fn pre_loop_submit_failure_closes_sdk_steering_before_runtime_returns_idle() {
let dir = temp_dir("steer-pre-loop-failure");
let config = Config::builder().cwd(dir).build();
let mut agent = Agent::with_provider(config, Box::new(SaysProvider("clean".into())));
agent.set_context_limit(50_000);
let engine = RpcEngine::new(agent, None);
let failed = FrontendRuntime::submit(engine.as_ref(), "x".repeat(200_000)).await;
let failure = failed.expect_err("oversized prompt must fail before run_loop");
assert!(
failure.to_string().contains("context limit"),
"failure must come from the pre-loop context guard: {failure}"
);
let stale = FrontendRuntime::steer(engine.as_ref(), "must not be acknowledged".into()).await;
assert!(
stale.is_err(),
"an idle runtime after pre-loop failure must reject steering"
);
assert_eq!(
FrontendRuntime::submit(engine.as_ref(), "next".into())
.await
.unwrap(),
"clean"
);
let history = FrontendRuntime::attach(engine.as_ref(), 16)
.await
.unwrap()
.history;
assert!(history
.iter()
.filter_map(|message| message.content.as_deref())
.all(|message| message != "must not be acknowledged"));
}
#[tokio::test]
async fn dropping_a_claimed_submit_restores_the_complete_idle_runtime_state() {
let started = std::sync::Arc::new(tokio::sync::Notify::new());
let agent = Agent::with_provider(
Config::builder()
.cwd(temp_dir("steer-dropped-claim"))
.build(),
Box::new(AbortThenSaysProvider {
calls: AtomicUsize::new(0),
started: started.clone(),
}),
);
let engine = RpcEngine::new(agent, None);
let submit_engine = engine.clone();
let submit = tokio::spawn(async move {
FrontendRuntime::submit(submit_engine.as_ref(), "abort this turn".into()).await
});
tokio::time::timeout(Duration::from_secs(2), started.notified())
.await
.expect("first provider request must be in flight before abort");
assert!(engine.status().busy);
let mut live = engine
.frontend_attach(16)
.expect("frontend attaches to the claimed turn");
assert_eq!(live.descriptor.turn_state, FrontendTurnState::Busy);
assert_eq!(live.next_event().await.unwrap().kind, "user_message");
assert_eq!(live.next_event().await.unwrap().kind, "turn_started");
submit.abort();
assert!(submit.await.unwrap_err().is_cancelled());
tokio::task::yield_now().await;
assert!(!engine.status().busy, "dropped claim must restore idle");
let terminal = tokio::time::timeout(Duration::from_secs(2), live.next_event())
.await
.expect("attached frontend must receive a terminal lifecycle event")
.expect("frontend event stream remains connected");
assert_eq!(terminal.kind, "turn_interrupted");
let mut after_abort = engine
.frontend_attach(16)
.expect("late frontend attaches after the dropped claim");
assert_eq!(after_abort.descriptor.turn_state, FrontendTurnState::Idle);
assert_eq!(
[
after_abort.next_event().await.unwrap().kind,
after_abort.next_event().await.unwrap().kind,
after_abort.next_event().await.unwrap().kind,
],
["user_message", "turn_started", "turn_interrupted"],
"replay must close the abandoned turn before advertising idle"
);
assert!(!engine.interrupt().await, "cancel handle must be cleared");
assert!(
FrontendRuntime::steer(engine.as_ref(), "must not leak".into())
.await
.is_err()
);
assert_eq!(
FrontendRuntime::submit(engine.as_ref(), "next".into())
.await
.unwrap(),
"clean"
);
let history = FrontendRuntime::attach(engine.as_ref(), 16)
.await
.unwrap()
.history;
assert!(history
.iter()
.filter_map(|message| message.content.as_deref())
.all(|message| message != "must not leak"));
}
#[tokio::test]
async fn shutdown_is_a_bounded_quiescence_barrier_before_finalization() {
let engine = RpcEngine::new(
Agent::with_provider(
Config::builder()
.cwd(temp_dir("shutdown-finalization-barrier"))
.build(),
Box::new(SlowProvider(Duration::from_secs(30), "too late".into())),
),
None,
);
let submit_engine = engine.clone();
let submit = tokio::spawn(async move { submit_engine.submit("in flight").await });
tokio::time::timeout(Duration::from_secs(2), async {
while !engine.status().busy {
tokio::task::yield_now().await;
}
})
.await
.expect("turn must become active");
tokio::time::timeout(Duration::from_secs(2), engine.shutdown())
.await
.expect("shutdown must interrupt and quiesce the active turn");
assert_eq!(submit.await.unwrap(), Err(RuntimeSubmitError::Interrupted));
assert!(engine.is_shutting_down());
assert!(!engine.status().busy);
assert_eq!(
engine.submit("must not restart").await,
Err(RuntimeSubmitError::Interrupted),
"no user or scheduler claim may cross the shutdown barrier"
);
tokio::time::timeout(
Duration::from_secs(2),
engine.finalize_with(|agent| agent.history().len()),
)
.await
.expect("finalization must acquire the quiescent agent immediately");
}
#[tokio::test]
async fn status_is_answerable_while_a_turn_is_in_flight_and_reports_busy() {
let dir = temp_dir("status-busy");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(
config,
Box::new(SlowProvider(Duration::from_millis(300), "done".into())),
);
let engine = RpcEngine::new(agent, None);
let idle = engine.handle_request(req(0, "status")).await;
assert_eq!(idle["result"]["busy"], false);
let engine2 = engine.clone();
let handle = tokio::spawn(async move { engine2.handle_request(submit_req(1, "go")).await });
tokio::time::sleep(Duration::from_millis(50)).await;
let busy = engine.handle_request(req(2, "status")).await;
assert_eq!(busy["result"]["busy"], true);
assert_eq!(busy["result"]["model"], idle["result"]["model"]);
let submitted = handle.await.unwrap();
assert_eq!(submitted["result"]["reply"], "done");
}
#[tokio::test]
async fn history_is_answerable_while_a_turn_is_in_flight() {
let dir = temp_dir("history-busy");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(
config,
Box::new(SlowProvider(Duration::from_millis(300), "done".into())),
);
let engine = RpcEngine::new(agent, None);
engine.handle_request(submit_req(1, "first")).await;
let engine2 = engine.clone();
let active =
tokio::spawn(async move { engine2.handle_request(submit_req(2, "still running")).await });
tokio::time::sleep(Duration::from_millis(50)).await;
let request = serde_json::from_value(json!({
"id": 3, "method": "history", "params": {"limit": 10}
}))
.unwrap();
let response = tokio::time::timeout(Duration::from_millis(100), engine.handle_request(request))
.await
.expect("history must not wait for the active turn's agent lock");
let messages = response["result"]["messages"].as_array().unwrap();
assert_eq!(messages[messages.len() - 2]["content"], "first");
assert_eq!(messages[messages.len() - 1]["content"], "done");
assert!(
messages
.iter()
.all(|message| message["content"] != "still running"),
"an in-flight partial turn must not leak into the boundary snapshot"
);
let completed = active.await.unwrap();
assert_eq!(completed["result"]["reply"], "done");
}
#[tokio::test]
async fn a_second_submit_while_busy_is_refused_immediately_never_queued() {
let dir = temp_dir("busy-refused");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(
config,
Box::new(SlowProvider(Duration::from_millis(500), "done".into())),
);
let engine = RpcEngine::new(agent, None);
let engine2 = engine.clone();
let first = tokio::spawn(async move { engine2.handle_request(submit_req(1, "go")).await });
tokio::time::sleep(Duration::from_millis(50)).await;
let started = std::time::Instant::now();
let second = engine.handle_request(submit_req(2, "again")).await;
assert!(
started.elapsed() < Duration::from_millis(400),
"the busy-refusal must return immediately, not wait for the first turn"
);
assert!(
second.get("error").is_some(),
"expected a busy error, got {second}"
);
let first_resp = first.await.unwrap();
assert_eq!(first_resp["result"]["reply"], "done");
}
#[tokio::test]
async fn interrupt_cancels_an_in_flight_turn_well_before_it_would_finish() {
let dir = temp_dir("interrupt");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(
config,
Box::new(SlowProvider(Duration::from_secs(10), "too slow".into())),
);
let engine = RpcEngine::new(agent, None);
let engine2 = engine.clone();
let handle = tokio::spawn(async move { engine2.handle_request(submit_req(1, "go")).await });
tokio::time::sleep(Duration::from_millis(50)).await;
let started = std::time::Instant::now();
let interrupt_resp = engine.handle_request(req(2, "interrupt")).await;
assert_eq!(interrupt_resp["result"]["interrupted"], true);
let submit_resp = handle.await.unwrap();
assert!(
started.elapsed() < Duration::from_secs(2),
"interrupt must cancel the turn immediately, not wait out the 10s provider sleep"
);
assert!(
submit_resp["error"]["message"]
.as_str()
.unwrap()
.contains("interrupted"),
"expected an interrupted error, got {submit_resp}"
);
}
#[tokio::test]
async fn interrupt_with_nothing_in_flight_is_a_harmless_false_not_an_error() {
let dir = temp_dir("interrupt-idle");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("x".to_string())));
let engine = RpcEngine::new(agent, None);
let resp = engine.handle_request(req(1, "interrupt")).await;
assert_eq!(resp["result"]["interrupted"], false);
assert!(resp.get("error").is_none());
}
#[tokio::test]
async fn shutdown_sets_the_flag_and_wakes_every_waiter() {
let dir = temp_dir("shutdown");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("x".to_string())));
let engine = RpcEngine::new(agent, None);
assert!(!engine.is_shutting_down());
let engine2 = engine.clone();
let waiter = tokio::spawn(async move {
engine2.wait_for_shutdown().await;
});
let resp = engine.handle_request(req(1, "shutdown")).await;
assert_eq!(resp["result"]["shutting_down"], true);
assert!(engine.is_shutting_down());
tokio::time::timeout(Duration::from_secs(2), waiter)
.await
.expect("wait_for_shutdown must resolve promptly after shutdown")
.unwrap();
tokio::time::timeout(Duration::from_millis(500), engine.wait_for_shutdown())
.await
.expect("a late wait_for_shutdown call must return immediately");
}
#[tokio::test]
async fn sdk_runtime_owns_active_claude_scheduling_without_a_frontend_input_loop() {
let dir = temp_dir("sdk-scheduler");
let config = Config::builder().cwd(dir).build();
let mut agent = Agent::with_provider(
config,
Box::new(SaysProvider("scheduled by SDK".to_string())),
);
let source = [
json!({
"type": "user", "sessionId": "runtime-sdk", "cwd": "/tmp",
"timestamp": "2020-01-01T00:00:00Z",
"message": {"role": "user", "content": "initial"}
}),
json!({
"type": "assistant", "timestamp": "2020-01-01T00:00:01Z",
"message": {"role": "assistant", "content": [{
"type": "tool_use", "id": "wake-sdk", "name": "ScheduleWakeup",
"input": {"delaySeconds": 1, "prompt": "SDK_SCHEDULED_PROMPT"}
}]}
}),
json!({
"type": "user", "timestamp": "2020-01-01T00:00:02Z",
"message": {"role": "user", "content": [{
"type": "tool_result", "tool_use_id": "wake-sdk",
"content": "Next wakeup scheduled for 00:00:02 (in 1s)."
}]}
}),
json!({
"type": "queue-operation", "operation": "enqueue",
"timestamp": "2020-01-01T00:00:03Z",
"content": "SDK_QUEUED_PROMPT"
}),
]
.into_iter()
.map(|value| value.to_string())
.collect::<Vec<_>>()
.join("\n");
let session = Session::from_claude_code_str(&source).unwrap();
let mut manifest = ClaudeRuntimeManifest::from_session(&session).unwrap();
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
manifest.activate_scheduler(now).unwrap();
agent.set_claude_runtime_manifest(manifest);
let persisted = std::sync::Arc::new(std::sync::Mutex::new(None));
let persisted_hook = persisted.clone();
let engine = RpcEngine::new_named(
agent,
"runtime-sdk",
Some(Box::new(move |agent| {
*persisted_hook.lock().unwrap() = agent.claude_runtime_manifest().cloned();
})),
);
let mut events = engine.subscribe();
let mut frontend = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
for kind in [
"scheduled_prompt_started",
"scheduled_prompt_deferred",
"scheduled_prompt_completed",
"scheduler_error",
] {
assert!(
frontend
.descriptor
.display
.event_kinds
.iter()
.any(|item| item == kind),
"descriptor must advertise scheduler event {kind}"
);
}
assert!(engine.start_claude_scheduler());
assert!(!engine.start_claude_scheduler());
let completed_kinds = tokio::time::timeout(Duration::from_secs(2), async {
let mut completed = Vec::new();
loop {
let event = events.recv().await.unwrap();
if event["type"] == "scheduled_prompt_completed" {
completed.push(event["kind"].as_str().unwrap().to_string());
if completed.len() == 2 {
break completed;
}
}
}
})
.await
.expect("SDK scheduler should deliver an overdue wakeup without any attached frontend");
assert_eq!(completed_kinds, vec!["queue", "wakeup"]);
let frontend_kinds = tokio::time::timeout(Duration::from_secs(2), async {
let mut kinds = Vec::new();
let mut scheduled_completions = 0;
loop {
let event = frontend.next_event().await.unwrap();
kinds.push(event.kind.clone());
if event.kind == "scheduled_prompt_completed" {
scheduled_completions += 1;
if scheduled_completions == 2 {
break kinds;
}
}
}
})
.await
.expect("attached frontend should observe the complete scheduled lifecycle");
let started = frontend_kinds
.iter()
.position(|kind| kind == "scheduled_prompt_started")
.unwrap();
let turn_started = frontend_kinds
.iter()
.position(|kind| kind == "turn_started")
.unwrap();
let completed = frontend_kinds
.iter()
.position(|kind| kind == "scheduled_prompt_completed")
.unwrap();
assert!(
started < turn_started && turn_started < completed,
"{frontend_kinds:?}"
);
let final_manifest = persisted.lock().unwrap().clone().unwrap();
assert!(final_manifest.queue.pending.is_empty());
assert!(final_manifest.pending_wakeups.is_empty());
assert!(final_manifest.scheduler.deliveries.is_empty());
let history = engine.handle_request(req(99, "history")).await;
let delivered_prompts = history["result"]["messages"]
.as_array()
.unwrap()
.iter()
.filter(|message| message["role"] == "user")
.map(|message| message["content"].as_str().unwrap())
.collect::<Vec<_>>();
assert_eq!(
delivered_prompts,
vec!["SDK_QUEUED_PROMPT", "SDK_SCHEDULED_PROMPT"]
);
engine.shutdown().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn shutdown_joins_scheduler_before_the_final_persistence_boundary() {
let dir = temp_dir("scheduler-shutdown-join");
let config = Config::builder().cwd(dir).build();
let mut agent = Agent::with_provider(
config,
Box::new(SaysProvider("must not run after the seal".into())),
);
let source = [
json!({
"type": "user", "sessionId": "scheduler-shutdown", "cwd": "/tmp",
"timestamp": "2020-01-01T00:00:00Z",
"message": {"role": "user", "content": "initial"}
}),
json!({
"type": "assistant", "timestamp": "2020-01-01T00:00:01Z",
"message": {"role": "assistant", "content": [{
"type": "tool_use", "id": "wake-shutdown", "name": "ScheduleWakeup",
"input": {"delaySeconds": 1, "prompt": "MUST_NOT_CROSS_SHUTDOWN"}
}]}
}),
json!({
"type": "user", "timestamp": "2020-01-01T00:00:02Z",
"message": {"role": "user", "content": [{
"type": "tool_result", "tool_use_id": "wake-shutdown",
"content": "Next wakeup scheduled for 00:00:02 (in 1s)."
}]}
}),
]
.into_iter()
.map(|value| value.to_string())
.collect::<Vec<_>>()
.join("\n");
let session = Session::from_claude_code_str(&source).unwrap();
let mut manifest = ClaudeRuntimeManifest::from_session(&session).unwrap();
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs() as i64;
manifest.activate_scheduler(now).unwrap();
agent.set_claude_runtime_manifest(manifest);
let hook_calls = std::sync::Arc::new(AtomicUsize::new(0));
let hook_calls_for_callback = hook_calls.clone();
let (entered_tx, entered_rx) = std::sync::mpsc::sync_channel(1);
let (release_tx, release_rx) = std::sync::mpsc::sync_channel(1);
let release_rx = std::sync::Arc::new(std::sync::Mutex::new(release_rx));
let release_rx_for_callback = release_rx.clone();
let engine = RpcEngine::new_named(
agent,
"scheduler-shutdown",
Some(Box::new(move |_| {
if hook_calls_for_callback.fetch_add(1, Ordering::SeqCst) == 0 {
entered_tx.send(()).unwrap();
release_rx_for_callback.lock().unwrap().recv().unwrap();
}
})),
);
assert!(engine.start_claude_scheduler());
tokio::task::spawn_blocking(move || entered_rx.recv().unwrap())
.await
.unwrap();
let shutdown_engine = engine.clone();
let shutdown = tokio::spawn(async move { shutdown_engine.shutdown().await });
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(
!shutdown.is_finished(),
"shutdown returned while the scheduler still held a mutation hook"
);
release_tx.send(()).unwrap();
tokio::time::timeout(Duration::from_secs(2), shutdown)
.await
.expect("shutdown must join the released scheduler task")
.unwrap();
let calls_at_shutdown = hook_calls.load(Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
hook_calls.load(Ordering::SeqCst),
calls_at_shutdown,
"no scheduler persistence hook may run after shutdown returns"
);
assert!(!engine.start_claude_scheduler());
tokio::time::timeout(
Duration::from_secs(1),
engine.finalize_with(|agent| agent.history().len()),
)
.await
.expect("finalization must acquire the scheduler-quiescent agent");
}
#[tokio::test]
async fn run_stdio_returns_promptly_on_reader_eof_and_flushes_pending_output() {
let dir = temp_dir("stdio-eof");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("hi via stdio".to_string())));
let engine = RpcEngine::new(agent, None);
let input = format!(
"{}\n{}\n",
serde_json::json!({"id": 1, "method": "submit", "params": {"prompt": "hello"}}),
serde_json::json!({"id": 2, "method": "status"}),
);
let reader = tokio::io::BufReader::new(std::io::Cursor::new(input.into_bytes()));
let (server_writer, mut test_reader) = tokio::io::duplex(64 * 1024);
let outcome = tokio::time::timeout(
Duration::from_secs(5),
supercode_harness::server::run_stdio(engine, reader, server_writer),
)
.await
.expect(
"run_stdio must return within 5s once `reader` hits EOF — it hung forever before the \
teardown fix (event_task looping on events.recv() while `engine` stayed alive)",
);
outcome.expect("run_stdio should not itself error on a clean EOF");
let mut out = Vec::new();
tokio::io::AsyncReadExt::read_to_end(&mut test_reader, &mut out)
.await
.unwrap();
let out_text = String::from_utf8(out).unwrap();
assert!(
out_text.contains("\"reply\":\"hi via stdio\""),
"expected the buffered submit reply to have been flushed before return, got: {out_text}"
);
assert!(
out_text.contains("\"busy\":false"),
"expected the buffered status reply to have been flushed before return, got: {out_text}"
);
}
#[tokio::test]
async fn run_stdio_returns_promptly_after_an_explicit_shutdown_rpc() {
let dir = temp_dir("stdio-shutdown");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("unused".to_string())));
let engine = RpcEngine::new(agent, None);
let input = format!("{}\n", serde_json::json!({"id": 1, "method": "shutdown"}));
let reader = tokio::io::BufReader::new(std::io::Cursor::new(input.into_bytes()));
let (server_writer, mut test_reader) = tokio::io::duplex(64 * 1024);
let outcome = tokio::time::timeout(
Duration::from_secs(5),
supercode_harness::server::run_stdio(engine.clone(), reader, server_writer),
)
.await
.expect(
"run_stdio must return within 5s after an explicit `shutdown` RPC — it hung forever \
before the teardown fix (event_task never observed the shutdown signal)",
);
outcome.expect("run_stdio should not itself error after a clean shutdown");
assert!(
engine.is_shutting_down(),
"the shutdown RPC must still have taken effect on the engine"
);
let mut out = Vec::new();
tokio::io::AsyncReadExt::read_to_end(&mut test_reader, &mut out)
.await
.unwrap();
let out_text = String::from_utf8(out).unwrap();
assert!(
out_text.contains("\"shutting_down\":true"),
"expected the shutdown reply to have been flushed before return, got: {out_text}"
);
}
#[tokio::test]
async fn on_turn_complete_hook_fires_only_after_a_successful_submit() {
let dir = temp_dir("hook");
let config = Config::builder().cwd(dir).build();
let agent = Agent::with_provider(config, Box::new(SaysProvider("ok".to_string())));
let calls = std::sync::Arc::new(AtomicUsize::new(0));
let calls2 = calls.clone();
let engine = RpcEngine::new(
agent,
Some(Box::new(move |_agent: &supercode_harness::SdkAgent| {
calls2.fetch_add(1, Ordering::SeqCst);
})),
);
let resp = engine.handle_request(submit_req(1, "go")).await;
assert_eq!(resp["result"]["reply"], "ok");
assert_eq!(calls.load(Ordering::SeqCst), 1);
}
struct AsksForBash {
calls: AtomicUsize,
command: String,
expect_approved: bool,
}
#[async_trait]
impl Provider for AsksForBash {
async fn complete(
&self,
req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
let n = self.calls.fetch_add(1, Ordering::SeqCst);
if n == 0 {
let call = ChatMessage {
role: Role::Assistant,
content: None,
content_parts: None,
tool_calls: Some(vec![ToolCall {
id: "call_1".into(),
kind: "function".into(),
function: FunctionCall {
name: "bash".into(),
arguments: json!({"command": self.command}).to_string(),
},
}]),
tool_call_id: None,
name: None,
metadata: Default::default(),
};
Ok((call, Usage::default()))
} else {
let last = req.messages.last().unwrap();
assert_eq!(last.role, Role::Tool);
let content = last.content.clone().unwrap_or_default();
if self.expect_approved {
assert!(
!content.contains("was not approved for execution"),
"approved tool unexpectedly received a denial: {content}"
);
Ok((ChatMessage::assistant("approved"), Usage::default()))
} else {
assert!(
content.contains("was not approved for execution"),
"expected a denial message fed back to the model, got: {content}"
);
Ok((
ChatMessage::assistant("acknowledged, tool was denied".to_string()),
Usage::default(),
))
}
}
}
}
#[tokio::test]
async fn rpc_driven_tool_call_stays_denied_under_on_request_approval_with_no_handler() {
let dir = temp_dir("no-bypass");
let marker = dir.join("should-never-exist.marker");
let mut config = Config::builder()
.cwd(dir.clone())
.approval(ApprovalPolicy::OnRequest)
.build();
config.permissions_enabled = true;
let provider = AsksForBash {
calls: AtomicUsize::new(0),
command: format!("touch {}", marker.display()),
expect_approved: false,
};
let agent = Agent::with_provider(config, Box::new(provider));
let engine = RpcEngine::new(agent, None);
let resp = tokio::time::timeout(
Duration::from_secs(5),
engine.handle_request(submit_req(1, "please run that command")),
)
.await
.expect("submit must not hang waiting on an approval the RPC surface can never answer");
assert_eq!(resp["result"]["reply"], "acknowledged, tool was denied");
assert!(
!marker.exists(),
"the denied bash tool call must never have actually run"
);
std::fs::remove_dir_all(&dir).ok();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn local_frontend_can_answer_a_typed_approval_request_exactly_once() {
let dir = temp_dir("frontend-approval");
let marker = dir.join("approved.marker");
let mut config = Config::builder()
.cwd(dir.clone())
.approval(ApprovalPolicy::OnRequest)
.build();
config.permissions_enabled = true;
let agent = Agent::with_provider(
config,
Box::new(AsksForBash {
calls: AtomicUsize::new(0),
command: format!("touch {}", marker.display()),
expect_approved: true,
}),
);
let engine = RpcEngine::new_named_with_frontend_requests(
agent,
"frontend-approval",
FrontendRuntimeMetadata::default(),
None,
);
let mut attachment = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
let mut observer = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
assert!(attachment.descriptor.actions.respond);
let submit_engine = engine.clone();
let submit = tokio::spawn(async move { submit_engine.submit("run it").await });
let request: FrontendRequest = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "request" {
break serde_json::from_value(event.payload["request"].clone()).unwrap();
}
}
})
.await
.expect("approval request should reach the attached frontend");
assert_eq!(request.payload["tool"], "bash");
let response = FrontendResponse::Approval {
request_id: request.id,
decision: FrontendApprovalDecision::Allow,
};
FrontendRuntime::respond(engine.as_ref(), response.clone())
.await
.unwrap();
for stream in [&mut attachment, &mut observer] {
let resolved = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = stream.next_event().await.unwrap();
if event.kind == "request_resolved" {
break event;
}
}
})
.await
.expect("every attached frontend should observe the chosen resolution");
assert_eq!(resolved.payload["request_id"], request.id);
assert_eq!(
resolved.payload["response"],
serde_json::to_value(&response).unwrap()
);
}
assert_eq!(submit.await.unwrap().unwrap(), "approved");
assert!(marker.exists());
let mut late = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
let (late_request, late_resolution) = tokio::time::timeout(Duration::from_secs(2), async {
let mut replayed_request = None;
loop {
let event = late.next_event().await.unwrap();
if event.kind == "request" {
replayed_request = Some(event);
} else if event.kind == "request_resolved" {
break (replayed_request, event);
}
}
})
.await
.expect("completed-turn attachment should retain request decision history");
assert!(late_request.is_some());
assert_eq!(late_resolution.payload["request_id"], request.id);
assert_eq!(late_resolution.payload["response"]["decision"], "allow");
assert!(matches!(
FrontendRuntime::respond(engine.as_ref(), response).await,
Err(supercode_harness::FrontendRuntimeError::UnknownRequest(id)) if id == request.id
));
std::fs::remove_dir_all(&dir).ok();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn local_frontend_answers_legacy_approval_when_permissions_engine_is_disabled() {
let dir = temp_dir("frontend-legacy-approval");
let marker = dir.join("approved.marker");
let config = Config::builder()
.cwd(dir.clone())
.approval(ApprovalPolicy::OnRequest)
.build();
assert!(
!config.permissions_enabled,
"this regression must exercise the default legacy approval gate"
);
let agent = Agent::with_provider(
config,
Box::new(AsksForBash {
calls: AtomicUsize::new(0),
command: format!("touch {}", marker.display()),
expect_approved: true,
}),
);
let engine = RpcEngine::new_named_with_frontend_requests(
agent,
"frontend-legacy-approval",
FrontendRuntimeMetadata::default(),
None,
);
let mut attachment = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
let submit_engine = engine.clone();
let submit = tokio::spawn(async move { submit_engine.submit("run it").await });
let request: FrontendRequest = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "request" {
break serde_json::from_value(event.payload["request"].clone()).unwrap();
}
}
})
.await
.expect("legacy approval request should reach the attached frontend");
assert_eq!(request.payload["tool"], "bash");
assert_eq!(
request.payload["subject"],
format!("touch {}", marker.display())
);
FrontendRuntime::respond(
engine.as_ref(),
FrontendResponse::Approval {
request_id: request.id,
decision: FrontendApprovalDecision::Allow,
},
)
.await
.unwrap();
assert_eq!(submit.await.unwrap().unwrap(), "approved");
assert!(marker.exists(), "approved legacy tool call should execute");
std::fs::remove_dir_all(&dir).ok();
}
struct FrontendChildApprovalProvider {
calls: std::sync::Arc<AtomicUsize>,
}
#[async_trait]
impl Provider for FrontendChildApprovalProvider {
async fn complete(
&self,
req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
let call = |id: &str, name: &str, arguments: serde_json::Value| ChatMessage {
role: Role::Assistant,
content: None,
content_parts: None,
tool_calls: Some(vec![ToolCall {
id: id.into(),
kind: "function".into(),
function: FunctionCall {
name: name.into(),
arguments: arguments.to_string(),
},
}]),
tool_call_id: None,
name: None,
metadata: Default::default(),
};
match self.calls.fetch_add(1, Ordering::SeqCst) {
0 => Ok((
call(
"spawn",
"spawn_subagent",
json!({"task":"run bash", "background":true}),
),
Usage::default(),
)),
1 => Ok((ChatMessage::assistant("spawn kicked off"), Usage::default())),
2 => Ok((
call("child-bash", "bash", json!({"command":"echo child-ok"})),
Usage::default(),
)),
3 => {
let output = req
.messages
.last()
.and_then(|message| message.content.as_deref())
.unwrap_or_default();
assert!(
output.contains("child-ok") && !output.contains("not approved"),
"the actual child call should run after frontend approval: {output}"
);
Ok((ChatMessage::assistant("child complete"), Usage::default()))
}
n => panic!("unexpected provider call {n}"),
}
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn local_frontend_routes_a_real_background_child_approval() {
let dir = temp_dir("frontend-child-approval");
let mut config = Config::builder()
.cwd(dir.clone())
.subagents_enabled(true)
.subagents_background(true)
.subagents_background_prompts(BackgroundPromptsPolicy::Parent)
.build();
config.permissions_enabled = true;
config.permissions_ask_patterns = vec!["bash".into()];
let calls = std::sync::Arc::new(AtomicUsize::new(0));
let agent = Agent::with_provider(
config,
Box::new(FrontendChildApprovalProvider {
calls: calls.clone(),
}),
);
let engine = RpcEngine::new_named_with_frontend_requests(
agent,
"frontend-child-approval",
FrontendRuntimeMetadata::default(),
None,
);
let mut attachment = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
let submit_engine = engine.clone();
let submit = tokio::spawn(async move { submit_engine.submit("start child").await });
let request: FrontendRequest = tokio::time::timeout(Duration::from_secs(5), async {
loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "request" {
let request: FrontendRequest =
serde_json::from_value(event.payload["request"].clone()).unwrap();
if request.payload.get("child_agent_id").is_some() {
break request;
}
}
}
})
.await
.expect("a real background child approval should reach the SDK frontend");
assert_eq!(request.payload["tool"], "bash");
assert!(request.payload["child_agent_id"].is_string());
FrontendRuntime::respond(
engine.as_ref(),
FrontendResponse::Approval {
request_id: request.id,
decision: FrontendApprovalDecision::Allow,
},
)
.await
.unwrap();
assert_eq!(submit.await.unwrap().unwrap(), "spawn kicked off");
tokio::time::timeout(Duration::from_secs(5), async {
while calls.load(Ordering::SeqCst) < 4 {
tokio::task::yield_now().await;
}
})
.await
.expect("approved child should finish its scripted turn");
std::fs::remove_dir_all(&dir).ok();
}
#[tokio::test]
async fn local_frontend_can_answer_a_typed_mcp_elicitation() {
let bridge = FrontendRequestBridge::new();
let elicitation = bridge.elicitation_handler();
let agent = Agent::with_provider(
Config::builder()
.cwd(temp_dir("frontend-elicitation"))
.build(),
Box::new(SaysProvider("unused".into())),
);
let engine = RpcEngine::new_named_with_frontend_bridge(
agent,
"frontend-elicitation",
FrontendRuntimeMetadata::default(),
bridge,
None,
);
let mut attachment = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
let pending = tokio::spawn(async move {
elicitation
.handle(&supercode_harness::mcp::ElicitationRequest {
message: "Which environment?".into(),
requested_schema: json!({"type": "object"}),
})
.await
});
let request: FrontendRequest = loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "request" {
break serde_json::from_value(event.payload["request"].clone()).unwrap();
}
};
assert_eq!(request.payload["message"], "Which environment?");
FrontendRuntime::respond(
engine.as_ref(),
FrontendResponse::Elicitation {
request_id: request.id,
action: FrontendElicitationAction::Accept,
content: Some(json!({"environment": "staging"})),
},
)
.await
.unwrap();
let mut late = FrontendRuntime::attach(engine.as_ref(), 20).await.unwrap();
let mut saw_request = false;
let resolution = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = late.next_event().await.unwrap();
if event.kind == "request" {
saw_request = true;
}
if event.kind == "request_resolved" {
break event;
}
}
})
.await
.expect("a late attachment should replay the request and its resolution");
assert!(saw_request);
assert_eq!(resolution.payload["request_id"], request.id);
assert_eq!(resolution.payload["response"]["kind"], "elicitation");
assert_eq!(resolution.payload["response"]["action"], "accept");
let response = pending.await.unwrap();
assert_eq!(
response.action,
supercode_harness::mcp::ElicitationAction::Accept
);
assert_eq!(response.content, Some(json!({"environment": "staging"})));
}