use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use serde_json::{json, Value};
use supercode_harness::mcp;
use supercode_harness::server::{RpcEngine, RuntimeSubmitError};
use supercode_harness::tools::{ToolContext, ToolRegistry};
use supercode_harness::{
Agent, ChatMessage, ChatRequest, Config, FrontendElicitationAction, FrontendResponse,
FrontendRuntimeError, HarnessId, HarnessSessionService, HttpFrontendRuntime, Provider,
SdkErrorCode, SdkOperation, SdkRequest, SdkRuntime, SdkService, Usage,
};
fn repo_root() -> PathBuf {
Path::new(env!("CARGO_MANIFEST_DIR"))
.join("../..")
.canonicalize()
.unwrap()
}
fn fixture_locator() -> Value {
json!({
"harness": HarnessId::PI,
"session_id": "1e6f2a3b-0000-4000-8000-000000000001",
"storage": {
"kind": "file",
"path": repo_root().join("crates/harness/tests/fixtures/pi_session.jsonl"),
},
})
}
async fn json_sdk(
service: &mut HarnessSessionService,
operation: SdkOperation,
params: Value,
) -> Value {
let method = operation.method().unwrap();
service
.handle_async(json!({"jsonrpc":"2.0", "id":1, "method":method, "params":params}))
.await
}
async fn mcp_sdk(registry: &ToolRegistry, operation: &str, params: Value) -> Value {
let response = mcp::handle_request(
registry,
&ToolContext::new(repo_root()),
&json!({
"jsonrpc":"2.0",
"id":1,
"method":"tools/call",
"params": {
"name":"supercode_sdk",
"arguments":{"operation":operation, "params":params},
},
}),
)
.await
.unwrap();
if response["result"]["isError"] == true {
return response;
}
serde_json::from_str(response["result"]["content"][0]["text"].as_str().unwrap()).unwrap()
}
#[tokio::test]
async fn persisted_session_surfaces_share_identity_export_and_named_errors() {
let locator = fixture_locator();
let mut direct = HarnessSessionService::new();
let mut json_api = HarnessSessionService::new();
let mut registry = ToolRegistry::new();
mcp::register_sdk_tool(&mut registry);
let direct_load = direct
.execute(SdkRequest {
operation: SdkOperation::Load,
params: json!({"locator":locator.clone()}),
})
.await
.unwrap();
let api_load = json_sdk(
&mut json_api,
SdkOperation::Load,
json!({"locator":locator.clone()}),
)
.await["result"]
.clone();
let mcp_load = mcp_sdk(®istry, "load", json!({"locator":locator.clone()})).await;
let load_rows = [
("rust-sdk", direct_load),
("service-json-rpc", api_load),
("mcp-tool", mcp_load),
];
for (surface, value) in &load_rows {
assert_eq!(
value["session"]["session_id"], "1e6f2a3b-0000-4000-8000-000000000001",
"stable identity drifted on {surface}"
);
}
for (_, value) in &load_rows[1..] {
assert_eq!(value, &load_rows[0].1);
}
let mut exports = Vec::new();
for (surface, mut service) in [
("rust-sdk", HarnessSessionService::new()),
("service-json-rpc", HarnessSessionService::new()),
] {
let value = if surface == "rust-sdk" {
service
.execute(SdkRequest {
operation: SdkOperation::Export,
params: json!({"locator":locator.clone(), "target_harness":"pi"}),
})
.await
.unwrap()
} else {
json_sdk(
&mut service,
SdkOperation::Export,
json!({"locator":locator.clone(), "target_harness":"pi"}),
)
.await["result"]
.clone()
};
exports.push((
surface,
value["artifact"]["content"].as_str().unwrap().to_string(),
));
}
exports.push((
"mcp-tool",
mcp_sdk(
®istry,
"export",
json!({"locator":locator, "target_harness":"pi"}),
)
.await["artifact"]["content"]
.as_str()
.unwrap()
.to_string(),
));
for (surface, content) in &exports[1..] {
assert_eq!(content, &exports[0].1, "export drifted on {surface}");
}
let direct_error = HarnessSessionService::new()
.execute(SdkRequest {
operation: SdkOperation::Steer,
params: json!({}),
})
.await
.unwrap_err();
assert_eq!(direct_error.code(), SdkErrorCode::UnsupportedAction);
let json_error = json_sdk(
&mut HarnessSessionService::new(),
SdkOperation::Steer,
json!({}),
)
.await;
assert_eq!(json_error["error"]["name"], "unsupported_action");
let mcp_error = mcp_sdk(®istry, "steer", json!({})).await;
assert_eq!(mcp_error["result"]["isError"], true);
assert_eq!(
mcp_error["result"]["structuredContent"]["error"]["name"],
"unsupported_action"
);
assert_eq!(
mcp_error["result"]["structuredContent"]["error"]["operation"],
"steer"
);
}
struct BlockingProvider {
entered: tokio::sync::Notify,
release: tokio::sync::Notify,
}
impl BlockingProvider {
fn new() -> Arc<Self> {
Arc::new(Self {
entered: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
})
}
}
struct SharedProvider(Arc<BlockingProvider>);
#[async_trait]
impl Provider for SharedProvider {
async fn complete(
&self,
_request: &ChatRequest,
on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
self.0.entered.notify_one();
self.0.release.notified().await;
on_delta("same reply");
Ok((ChatMessage::assistant("same reply"), Usage::default()))
}
}
struct LiveObservation {
session_id: String,
events: Vec<(u64, String, Value)>,
history: Vec<ChatMessage>,
persisted: String,
}
async fn exercise_live(
runtime: Arc<dyn SdkRuntime>,
provider: Arc<BlockingProvider>,
persisted: Arc<Mutex<String>>,
) -> LiveObservation {
let mut attachment = runtime.attach(50).await.unwrap();
let session_id = attachment.descriptor.session_id.clone();
let active = {
let runtime = runtime.clone();
tokio::spawn(async move { runtime.submit("same prompt".into()).await })
};
provider.entered.notified().await;
assert!(matches!(
runtime.submit("busy prompt".into()).await,
Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy))
));
assert!(matches!(
runtime
.respond(FrontendResponse::Other {
request_id: 99,
action: FrontendElicitationAction::Cancel,
content: None,
})
.await,
Err(FrontendRuntimeError::UnsupportedAction("respond"))
));
provider.release.notify_waiters();
assert_eq!(active.await.unwrap().unwrap(), "same reply");
let mut events = Vec::new();
loop {
let event =
tokio::time::timeout(std::time::Duration::from_secs(2), attachment.next_event())
.await
.unwrap()
.unwrap();
let terminal = event.kind == "turn_succeeded";
events.push((event.sequence, event.kind, event.payload));
if terminal {
break;
}
}
let history = runtime.attach(50).await.unwrap().history;
let persisted = persisted.lock().unwrap().clone();
LiveObservation {
session_id,
events,
history,
persisted,
}
}
async fn exercise_attached_mcp(
runtime: Arc<RpcEngine>,
provider: Arc<BlockingProvider>,
persisted: Arc<Mutex<String>>,
) -> LiveObservation {
let mut registry = ToolRegistry::new();
registry.register(mcp::SdkMcpTool::attached(runtime.clone()).await.unwrap());
let registry = Arc::new(registry);
let active = {
let registry = registry.clone();
tokio::spawn(
async move { mcp_sdk(®istry, "input", json!({"prompt":"same prompt"})).await },
)
};
provider.entered.notified().await;
let busy = mcp_sdk(®istry, "input", json!({"prompt":"busy prompt"})).await;
assert_eq!(busy["result"]["structuredContent"]["error"]["name"], "busy");
assert_eq!(
busy["result"]["structuredContent"]["error"]["operation"],
"input"
);
let unsupported = mcp_sdk(
®istry,
"respond",
json!({"response":{
"kind":"other", "request_id":99, "action":"cancel", "content":null
}}),
)
.await;
assert_eq!(
unsupported["result"]["structuredContent"]["error"]["name"],
"unsupported_action"
);
assert_eq!(
unsupported["result"]["structuredContent"]["error"]["operation"],
"respond"
);
provider.release.notify_waiters();
let completed = active.await.unwrap();
assert_eq!(completed["session_id"], "stable-sdk-session");
assert_eq!(completed["reply"], "same reply");
let mut events = Vec::new();
loop {
let value = mcp_sdk(®istry, "events", json!({})).await;
assert_eq!(value["session_id"], "stable-sdk-session");
let event = &value["event"];
let terminal = event["kind"] == "turn_succeeded";
events.push((
event["sequence"].as_u64().unwrap(),
event["kind"].as_str().unwrap().to_string(),
event["payload"].clone(),
));
if terminal {
break;
}
}
LiveObservation {
session_id: "stable-sdk-session".into(),
events,
history: runtime.attach(50).await.unwrap().history,
persisted: persisted.lock().unwrap().clone(),
}
}
fn scripted_engine(
provider: Arc<BlockingProvider>,
persisted: Arc<Mutex<String>>,
cwd: &Path,
) -> Arc<RpcEngine> {
let agent = Agent::with_provider(
Config::builder()
.cwd(cwd)
.system_prompt("sdk conformance")
.build(),
Box::new(SharedProvider(provider)),
);
RpcEngine::new_named(
agent,
"stable-sdk-session",
Some(Box::new(move |agent: &supercode_harness::SdkAgent| {
*persisted.lock().unwrap() = serde_json::to_string(agent.history()).unwrap();
})),
)
}
#[tokio::test]
async fn local_and_http_live_surfaces_share_events_busy_and_persistence() {
let fixture: Value = serde_json::from_str(include_str!(
"../../../sdk/frontend/test/fixtures/conformance.json"
))
.unwrap();
assert_eq!(fixture["session_id"], "stable-sdk-session");
assert_eq!(fixture["system"], "sdk conformance");
assert_eq!(fixture["prompt"], "same prompt");
assert_eq!(fixture["competing_prompt"], "busy prompt");
assert_eq!(fixture["reply"], "same reply");
let cwd =
std::env::temp_dir().join(format!("supercode-sdk-conformance-{}", std::process::id()));
std::fs::create_dir_all(&cwd).unwrap();
let local_provider = BlockingProvider::new();
let local_persisted = Arc::new(Mutex::new(String::new()));
let local_engine = scripted_engine(local_provider.clone(), local_persisted.clone(), &cwd);
let local = exercise_live(local_engine.clone(), local_provider, local_persisted).await;
let http_provider = BlockingProvider::new();
let http_persisted = Arc::new(Mutex::new(String::new()));
let http_engine = scripted_engine(http_provider.clone(), http_persisted.clone(), &cwd);
let token: Arc<str> = "sdk-conformance-token".into();
let address =
supercode_harness::server::run_http(http_engine.clone(), "127.0.0.1:0", token.clone())
.await
.unwrap();
let http_runtime = HttpFrontendRuntime::connect(format!("http://{address}"), token.to_string())
.await
.unwrap();
let http = exercise_live(http_runtime, http_provider, http_persisted).await;
let mcp_provider = BlockingProvider::new();
let mcp_persisted = Arc::new(Mutex::new(String::new()));
let mcp_engine = scripted_engine(mcp_provider.clone(), mcp_persisted.clone(), &cwd);
let mcp = exercise_attached_mcp(mcp_engine.clone(), mcp_provider, mcp_persisted).await;
for (surface, observed) in [("local", &local), ("http", &http), ("mcp", &mcp)] {
assert_eq!(observed.session_id, "stable-sdk-session", "{surface}");
assert!(!observed.persisted.is_empty(), "{surface} did not persist");
assert_eq!(
observed
.events
.iter()
.map(|(_, kind, _)| kind.as_str())
.collect::<Vec<_>>(),
fixture["event_kinds"]
.as_array()
.unwrap()
.iter()
.map(|value| value.as_str().unwrap())
.collect::<Vec<_>>(),
"{surface}"
);
}
assert_eq!(local.events, http.events);
assert_eq!(
serde_json::to_value(&local.history).unwrap(),
serde_json::to_value(&http.history).unwrap(),
"local internal history may retain durable metadata that ChatMessage's public wire projection deliberately omits"
);
assert_eq!(local.persisted, http.persisted);
assert_eq!(local.events, mcp.events);
assert_eq!(
serde_json::to_value(&local.history).unwrap(),
serde_json::to_value(&mcp.history).unwrap(),
"local and MCP public history projections must remain identical"
);
assert_eq!(local.persisted, mcp.persisted);
local_engine.shutdown().await;
http_engine.shutdown().await;
mcp_engine.shutdown().await;
std::fs::remove_dir_all(cwd).ok();
}
#[tokio::test]
async fn simultaneous_http_inputs_have_one_controller_and_one_named_lease_error() {
let cwd = std::env::temp_dir().join(format!("supercode-sdk-input-race-{}", std::process::id()));
std::fs::create_dir_all(&cwd).unwrap();
let provider = BlockingProvider::new();
let persisted = Arc::new(Mutex::new(String::new()));
let engine = scripted_engine(provider.clone(), persisted, &cwd);
let token: Arc<str> = "sdk-input-race-token".into();
let address = supercode_harness::server::run_http(engine.clone(), "127.0.0.1:0", token.clone())
.await
.unwrap();
let first = HttpFrontendRuntime::connect(format!("http://{address}"), token.to_string())
.await
.unwrap();
let second = HttpFrontendRuntime::connect(format!("http://{address}"), token.to_string())
.await
.unwrap();
let mut attachment = first.attach(50).await.unwrap();
let barrier = Arc::new(tokio::sync::Barrier::new(3));
let first_input = {
let barrier = barrier.clone();
tokio::spawn(async move {
barrier.wait().await;
SdkRuntime::send_input(first, "simultaneous first".into()).await
})
};
let second_input = {
let barrier = barrier.clone();
tokio::spawn(async move {
barrier.wait().await;
SdkRuntime::send_input(second, "simultaneous second".into()).await
})
};
barrier.wait().await;
let results = [first_input.await.unwrap(), second_input.await.unwrap()];
assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
assert_eq!(
results
.iter()
.filter(|result| matches!(result, Err(FrontendRuntimeError::ControllerRequired { .. })))
.count(),
1
);
provider.entered.notified().await;
provider.release.notify_waiters();
loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "turn_succeeded" {
break;
}
}
engine.shutdown().await;
std::fs::remove_dir_all(cwd).ok();
}