use std::io::{BufRead, BufReader, Write};
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use serde_json::{json, Value};
use supercode::server::{run_http, RpcEngine};
use supercode::{
register_live_runtime, Agent, ChatMessage, ChatRequest, Config, FrontendRuntimeMetadata,
LiveRuntimeSource, Provider, Session, SessionFormat, Usage,
};
static SDK_PROCESS_ENV_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
const SURFACE_SENTINELS: [&str; 5] = [
"CLI_EXPORT_SENTINEL",
"MCP_EXPORT_SENTINEL",
"TYPESCRIPT_EXPORT_SENTINEL",
"ACP_ENVELOPE_SENTINEL",
"TERMINAL_STYLE_SENTINEL",
];
fn bin() -> PathBuf {
PathBuf::from(env!("CARGO_BIN_EXE_supercode"))
}
fn repo_root() -> PathBuf {
Path::new(env!("CARGO_MANIFEST_DIR"))
.join("../..")
.canonicalize()
.unwrap()
}
fn fixture() -> PathBuf {
repo_root().join("crates/harness/tests/fixtures/pi_session.jsonl")
}
fn assert_pi_disk_reload_clean(content: &str) {
let nonce = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let path = std::env::temp_dir().join(format!("supercode-surface-export-{nonce}.jsonl"));
std::fs::write(&path, content).unwrap();
let reloaded = Session::load(&path).unwrap();
let normalized = reloaded.to_jsonl(SessionFormat::Pi).unwrap();
for sentinel in SURFACE_SENTINELS {
assert!(!content.contains(sentinel), "export leaked {sentinel}");
assert!(
!normalized.contains(sentinel),
"disk reload leaked {sentinel}"
);
}
std::fs::remove_file(path).ok();
}
fn locator() -> Value {
json!({
"harness": "pi",
"session_id": "1e6f2a3b-0000-4000-8000-000000000001",
"storage": {"kind": "file", "path": fixture()},
})
}
fn request(
child: &mut std::process::Child,
reader: &mut BufReader<std::process::ChildStdout>,
value: Value,
) -> Value {
let request_id = value.get("id").cloned().unwrap_or(Value::Null);
writeln!(child.stdin.as_mut().unwrap(), "{value}").unwrap();
child.stdin.as_mut().unwrap().flush().unwrap();
loop {
let mut line = String::new();
reader.read_line(&mut line).unwrap();
assert!(
!line.is_empty(),
"service closed before response {request_id}"
);
let response: Value = serde_json::from_str(&line).unwrap();
if response.get("id") == Some(&request_id) {
return response;
}
}
}
struct BlockingProvider {
entered: tokio::sync::Notify,
release: tokio::sync::Notify,
}
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::Result<(ChatMessage, Usage)> {
self.0.entered.notify_one();
self.0.release.notified().await;
on_delta("same reply");
Ok((ChatMessage::assistant("same reply"), Usage::default()))
}
}
fn spawn_harness(
workspace: &Path,
home: &Path,
) -> (std::process::Child, BufReader<std::process::ChildStdout>) {
let mut child = Command::new(bin())
.args([
"--cwd",
workspace.to_str().unwrap(),
"harness",
"serve",
"--poll-ms",
"5",
])
.env("SUPERCODE_HOME", home)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let reader = BufReader::new(child.stdout.take().unwrap());
(child, reader)
}
fn attach_live(
child: &mut std::process::Child,
reader: &mut BufReader<std::process::ChildStdout>,
workspace: &Path,
endpoint: &str,
id: u64,
) -> (String, Value) {
let opened = request(
child,
reader,
json!({
"jsonrpc":"2.0", "id":id, "method":"harness.v1.runtimes.attach_existing",
"params":{
"harness":"pi", "runtime_id":"source-session", "cwd":workspace,
"base_url":endpoint,
}
}),
);
assert!(opened.get("error").is_none(), "{opened}");
(
opened["result"]["connection"].as_str().unwrap().to_string(),
opened,
)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cli_process_preserves_the_complete_live_sdk_script_contract() {
let _environment = SDK_PROCESS_ENV_LOCK.lock().await;
let nonce = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let root = std::env::temp_dir().join(format!("supercode-cli-sdk-{nonce}"));
let workspace = root.join("workspace");
let home = root.join("home");
std::fs::create_dir_all(&workspace).unwrap();
std::fs::create_dir_all(&home).unwrap();
std::env::set_var("SUPERCODE_HOME", &home);
let provider = Arc::new(BlockingProvider {
entered: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let persisted = Arc::new(Mutex::new(String::new()));
let hook = persisted.clone();
let runtime = RpcEngine::new_named_with_frontend_metadata(
Agent::with_provider(
Config::builder()
.cwd(&workspace)
.system_prompt("sdk conformance")
.build(),
Box::new(SharedProvider(provider.clone())),
),
"stable-sdk-session",
FrontendRuntimeMetadata {
source_harness: Some("pi".into()),
emulation_profile: None,
},
Some(Box::new(move |agent: &supercode::SdkAgent| {
*hook.lock().unwrap() = serde_json::to_string(agent.history()).unwrap();
})),
);
let token: Arc<str> = "cli-live-token".into();
let address = run_http(runtime.clone(), "127.0.0.1:0", token.clone())
.await
.unwrap();
let registration = register_live_runtime(
runtime.session_id(),
LiveRuntimeSource {
harness: "pi".into(),
session_id: "source-session".into(),
workspace: workspace.clone(),
},
format!("http://{address}"),
token.to_string(),
)
.unwrap();
let (mut first, mut first_out) = spawn_harness(&workspace, &home);
let (first_connection, opened) = attach_live(
&mut first,
&mut first_out,
&workspace,
registration.endpoint().as_str(),
1,
);
assert_eq!(
opened["result"]["handle"]["runtime_id"],
"stable-sdk-session"
);
let first_turn = tokio::task::spawn_blocking(move || {
let sent = request(
&mut first,
&mut first_out,
json!({
"jsonrpc":"2.0", "id":2, "method":"harness.v1.runtimes.send_input",
"params":{"connection":first_connection, "text":"same prompt"}
}),
);
(first, first_out, sent)
});
provider.entered.notified().await;
let (mut second, mut second_out) = spawn_harness(&workspace, &home);
let (second_connection, _) = attach_live(
&mut second,
&mut second_out,
&workspace,
registration.endpoint().as_str(),
3,
);
let denied = request(
&mut second,
&mut second_out,
json!({
"jsonrpc":"2.0", "id":4, "method":"harness.v1.runtimes.send_input",
"params":{"connection":second_connection, "text":"busy prompt"}
}),
);
assert_eq!(denied["error"]["name"], "controller_required", "{denied}");
assert_eq!(denied["error"]["operation"], Value::Null);
let unsupported = request(
&mut second,
&mut second_out,
json!({
"jsonrpc":"2.0", "id":5, "method":"harness.v1.runtimes.respond",
"params":{"connection":second_connection, "request_id":99, "response":{
"kind":"other", "request_id":99, "action":"cancel", "content":null
}}
}),
);
assert_eq!(
unsupported["error"]["name"], "unsupported_action",
"{unsupported}"
);
assert_eq!(unsupported["error"]["operation"], "respond");
provider.release.notify_waiters();
let (mut first, mut first_out, sent) = first_turn.await.unwrap();
assert!(sent.get("error").is_none(), "{sent}");
let mut events = Vec::new();
while events.last().map(|(_, kind): &(u64, String)| kind.as_str()) != Some("turn_succeeded") {
let mut line = String::new();
first_out.read_line(&mut line).unwrap();
assert!(
!line.is_empty(),
"CLI closed before emitting runtime events"
);
let value: Value = serde_json::from_str(&line).unwrap();
if value["method"] == "harness.v1.runtimes.event" {
assert_eq!(value["params"]["session_id"], "stable-sdk-session");
let event = &value["params"]["event"];
let initial_system_history =
event["kind"] == "message" && event["payload"]["role"] == "system";
if !initial_system_history {
events.push((
value["params"]["sequence"].as_u64().unwrap(),
event["kind"].as_str().unwrap().to_string(),
));
}
}
}
assert_eq!(
events,
vec![
(1, "user_message".into()),
(2, "turn_started".into()),
(3, "text_delta".into()),
(4, "usage".into()),
(5, "turn_completed".into()),
(6, "turn_succeeded".into()),
]
);
let expected_persisted = serde_json::to_string(&vec![
ChatMessage::system("sdk conformance"),
ChatMessage::user("same prompt"),
ChatMessage::assistant("same reply"),
])
.unwrap();
assert_eq!(*persisted.lock().unwrap(), expected_persisted);
drop(first.stdin.take());
drop(second.stdin.take());
assert!(first.wait().unwrap().success());
assert!(second.wait().unwrap().success());
runtime.shutdown().await;
drop(registration);
std::env::remove_var("SUPERCODE_HOME");
std::fs::remove_dir_all(root).ok();
}
#[tokio::test]
async fn cli_and_opt_in_mcp_project_the_real_sdk_service() {
let _environment = SDK_PROCESS_ENV_LOCK.lock().await;
let root = repo_root();
let mut cli = Command::new(bin())
.args([
"--cwd",
root.to_str().unwrap(),
"harness",
"serve",
"--poll-ms",
"10",
])
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let mut cli_out = BufReader::new(cli.stdout.take().unwrap());
let loaded = request(
&mut cli,
&mut cli_out,
json!({
"jsonrpc":"2.0", "id":"CLI_RPC_SENTINEL", "method":"harness.v1.sessions.load",
"params":{"locator":locator(), "_surface":"CLI_PARAM_SENTINEL"},
}),
);
assert_eq!(
loaded["result"]["session"]["session_id"],
locator()["session_id"]
);
assert!(!loaded["result"].to_string().contains("CLI_PARAM_SENTINEL"));
let exported = request(
&mut cli,
&mut cli_out,
json!({
"jsonrpc":"2.0", "id":3, "method":"harness.v1.sessions.export",
"params":{
"locator":locator(), "target_harness":"pi",
"_surface":"CLI_EXPORT_SENTINEL",
},
}),
);
assert!(!exported["result"]["artifact"]["content"]
.as_str()
.unwrap()
.contains("CLI_EXPORT_SENTINEL"));
assert_pi_disk_reload_clean(exported["result"]["artifact"]["content"].as_str().unwrap());
let unsupported = request(
&mut cli,
&mut cli_out,
json!({
"jsonrpc":"2.0", "id":2, "method":"harness.v1.runtimes.steer", "params":{},
}),
);
assert_eq!(unsupported["error"]["name"], "unsupported_action");
drop(cli.stdin.take());
assert!(cli.wait().unwrap().success());
let mut mcp = Command::new(bin())
.args([
"--cwd",
root.to_str().unwrap(),
"mcp",
"serve",
"--sdk-readonly",
])
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let mut mcp_out = BufReader::new(mcp.stdout.take().unwrap());
let loaded = request(
&mut mcp,
&mut mcp_out,
json!({
"jsonrpc":"2.0", "id":"MCP_CALL_SENTINEL", "method":"tools/call",
"params":{"name":"supercode_sdk", "arguments":{"operation":"load", "params":{
"locator":locator(), "_surface":"MCP_PARAM_SENTINEL"
}}},
}),
);
let payload: Value =
serde_json::from_str(loaded["result"]["content"][0]["text"].as_str().unwrap()).unwrap();
assert_eq!(payload["session"]["session_id"], locator()["session_id"]);
assert!(!payload.to_string().contains("MCP_PARAM_SENTINEL"));
let exported = request(
&mut mcp,
&mut mcp_out,
json!({
"jsonrpc":"2.0", "id":3, "method":"tools/call",
"params":{"name":"supercode_sdk", "arguments":{"operation":"export", "params":{
"locator":locator(), "target_harness":"pi", "_surface":"MCP_EXPORT_SENTINEL"
}}},
}),
);
let payload: Value =
serde_json::from_str(exported["result"]["content"][0]["text"].as_str().unwrap()).unwrap();
assert!(!payload["artifact"]["content"]
.as_str()
.unwrap()
.contains("MCP_EXPORT_SENTINEL"));
assert_pi_disk_reload_clean(payload["artifact"]["content"].as_str().unwrap());
let unsupported = request(
&mut mcp,
&mut mcp_out,
json!({
"jsonrpc":"2.0", "id":2, "method":"tools/call",
"params":{"name":"supercode_sdk", "arguments":{"operation":"steer", "params":{}}},
}),
);
assert_eq!(
unsupported["result"]["structuredContent"]["error"]["name"],
"unsupported_action"
);
drop(mcp.stdin.take());
assert!(mcp.wait().unwrap().success());
}
#[tokio::test]
async fn typescript_adapter_drives_the_real_cli_sdk_service() {
let _environment = SDK_PROCESS_ENV_LOCK.lock().await;
let export_path = std::env::temp_dir().join(format!(
"supercode-typescript-export-{}.jsonl",
std::process::id()
));
let status = Command::new("node")
.args(["--test", "sdk/typescript/test/live-cli.test.mjs"])
.current_dir(repo_root())
.env("SUPERCODE_LIVE_BIN", bin())
.env("SUPERCODE_LIVE_FIXTURE", fixture())
.env("SUPERCODE_LIVE_WORKSPACE", repo_root())
.env("SUPERCODE_TS_EXPORT", &export_path)
.status()
.expect("Node.js is required for the TypeScript SDK conformance proof");
assert!(status.success());
assert_pi_disk_reload_clean(&std::fs::read_to_string(&export_path).unwrap());
std::fs::remove_file(export_path).ok();
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn typescript_process_preserves_the_complete_live_sdk_script_contract() {
let _environment = SDK_PROCESS_ENV_LOCK.lock().await;
let nonce = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let root = std::env::temp_dir().join(format!("supercode-ts-sdk-{nonce}"));
let workspace = root.join("workspace");
let home = root.join("home");
let gate = root.join("release-gate");
let typescript_export = root.join("typescript-export.jsonl");
std::fs::create_dir_all(&workspace).unwrap();
std::fs::create_dir_all(&home).unwrap();
std::env::set_var("SUPERCODE_HOME", &home);
let provider = Arc::new(BlockingProvider {
entered: tokio::sync::Notify::new(),
release: tokio::sync::Notify::new(),
});
let persisted = Arc::new(Mutex::new(String::new()));
let hook = persisted.clone();
let runtime = RpcEngine::new_named_with_frontend_metadata(
Agent::with_provider(
Config::builder()
.cwd(&workspace)
.system_prompt("sdk conformance")
.build(),
Box::new(SharedProvider(provider.clone())),
),
"stable-sdk-session",
FrontendRuntimeMetadata {
source_harness: Some("pi".into()),
emulation_profile: None,
},
Some(Box::new(move |agent: &supercode::SdkAgent| {
*hook.lock().unwrap() = serde_json::to_string(agent.history()).unwrap();
})),
);
let token: Arc<str> = "typescript-live-token".into();
let address = run_http(runtime.clone(), "127.0.0.1:0", token.clone())
.await
.unwrap();
let registration = register_live_runtime(
runtime.session_id(),
LiveRuntimeSource {
harness: "pi".into(),
session_id: "source-session".into(),
workspace: workspace.clone(),
},
format!("http://{address}"),
token.to_string(),
)
.unwrap();
let node_bin = bin();
let node_fixture = fixture();
let node_root = repo_root();
let node_workspace = workspace.clone();
let node_home = home.clone();
let node_gate = gate.clone();
let node_export = typescript_export.clone();
let endpoint = registration.endpoint().as_str().to_string();
let node = tokio::task::spawn_blocking(move || {
Command::new("node")
.args(["--test", "sdk/typescript/test/live-cli.test.mjs"])
.current_dir(&node_root)
.env("SUPERCODE_LIVE_BIN", node_bin)
.env("SUPERCODE_LIVE_FIXTURE", node_fixture)
.env("SUPERCODE_LIVE_WORKSPACE", node_workspace)
.env("SUPERCODE_LIVE_ENDPOINT", endpoint)
.env("SUPERCODE_LIVE_GATE", node_gate)
.env("SUPERCODE_TS_EXPORT", node_export)
.env("SUPERCODE_HOME", node_home)
.status()
.expect("Node.js is required for the TypeScript SDK conformance proof")
});
provider.entered.notified().await;
tokio::time::timeout(std::time::Duration::from_secs(10), async {
while !gate.exists() {
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
})
.await
.expect("TypeScript should observe busy and unsupported errors");
provider.release.notify_waiters();
assert!(node.await.unwrap().success());
assert_pi_disk_reload_clean(&std::fs::read_to_string(&typescript_export).unwrap());
let expected_persisted = serde_json::to_string(&vec![
ChatMessage::system("sdk conformance"),
ChatMessage::user("same prompt"),
ChatMessage::assistant("same reply"),
])
.unwrap();
assert_eq!(*persisted.lock().unwrap(), expected_persisted);
runtime.shutdown().await;
drop(registration);
std::env::remove_var("SUPERCODE_HOME");
std::fs::remove_dir_all(root).ok();
}