use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use base64::Engine as _;
use serde_json::{json, Value};
use supercode::{
ChatMessage, FrontendApprovalDecision, FrontendAttachment, FrontendResponse,
FrontendRuntimeDescriptor, Role, SdkError, SdkEvent, SdkRuntime,
};
use url::Url;
pub const PROTOCOL_NAMESPACE: &str = "opencode_http/v1_2_15";
pub const OPENCODE_CLI_VERSION: &str = "1.2.15";
pub const HISTORY_LIMIT: usize = 4096;
const TRACED_ROUTES: &[&str] = &[
"GET /agent",
"GET /command",
"GET /config",
"GET /config/providers",
"GET /event",
"GET /experimental/resource",
"GET /formatter",
"GET /lsp",
"GET /mcp",
"GET /path",
"GET /provider",
"GET /provider/auth",
"GET /session/status",
"GET /session/{session_id}",
"GET /session/{session_id}/diff",
"GET /session/{session_id}/message?limit=100",
"GET /session/{session_id}/todo",
"GET /session?start={cursor}",
"GET /vcs",
"POST /permission/{permission_id}/reply",
"POST /session",
"POST /session/{session_id}/abort",
"POST /session/{session_id}/message",
];
#[derive(Debug, Clone, PartialEq)]
pub struct OpenCodeRequest {
pub method: String,
pub target: String,
pub body: Value,
}
impl OpenCodeRequest {
pub fn new(method: impl Into<String>, target: impl Into<String>) -> Self {
Self {
method: method.into(),
target: target.into(),
body: Value::Null,
}
}
pub fn with_body(mut self, body: Value) -> Self {
self.body = body;
self
}
}
pub enum ResponseBody {
Json(Value),
EventStream(Box<FrontendAttachment>),
}
pub struct OpenCodeResponse {
pub status: u16,
pub body: ResponseBody,
}
impl OpenCodeResponse {
fn json(status: u16, body: Value) -> Self {
Self {
status,
body: ResponseBody::Json(body),
}
}
fn events(attachment: FrontendAttachment) -> Self {
Self {
status: 200,
body: ResponseBody::EventStream(Box::new(attachment)),
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum AdapterError {
#[error("route not supported by this adapter version: {0}")]
UnsupportedRoute(String),
#[error("invalid request for `{route}`: {message}")]
InvalidRequest { route: String, message: String },
#[error("OpenCode session `{0}` is not attached to this runtime")]
UnknownSession(String),
#[error(transparent)]
Sdk(#[from] SdkError),
}
impl AdapterError {
pub fn response(&self) -> OpenCodeResponse {
let (status, name) = match self {
Self::UnsupportedRoute(_) => (404, "unsupported_route"),
Self::InvalidRequest { .. } | Self::UnknownSession(_) => (400, "invalid_request"),
Self::Sdk(error) if error.code() == supercode::SdkErrorCode::ControllerRequired => {
(409, "controller_required")
}
Self::Sdk(error) if error.code() == supercode::SdkErrorCode::LeaseExpired => {
(409, "lease_expired")
}
Self::Sdk(error) if error.code() == supercode::SdkErrorCode::Busy => (409, "busy"),
Self::Sdk(error) if error.code() == supercode::SdkErrorCode::UnsupportedAction => {
(409, "action_unavailable")
}
Self::Sdk(error) if error.code() == supercode::SdkErrorCode::Unauthorized => {
(403, "unauthorized")
}
Self::Sdk(error) if error.code() == supercode::SdkErrorCode::Unauthenticated => {
(401, "unauthenticated")
}
Self::Sdk(_) => (500, "sdk_error"),
};
OpenCodeResponse::json(status, json!({"name": name, "message": self.to_string()}))
}
}
pub struct OpenCodeAdapter {
runtime: Arc<dyn SdkRuntime>,
runtime_id: String,
workspace: PathBuf,
session_id: String,
history_ids: Arc<Mutex<HistoryIdentityState>>,
}
impl OpenCodeAdapter {
pub fn new(
runtime: Arc<dyn SdkRuntime>,
runtime_id: impl Into<String>,
workspace: impl Into<PathBuf>,
) -> Arc<Self> {
let runtime_id = runtime_id.into();
let session_id = stable_id("session", "ses", &runtime_id, 24);
Arc::new(Self {
runtime,
runtime_id,
workspace: workspace.into(),
session_id,
history_ids: Arc::new(Mutex::new(HistoryIdentityState::default())),
})
}
pub fn session_id(&self) -> &str {
&self.session_id
}
pub fn traced_routes() -> &'static [&'static str] {
TRACED_ROUTES
}
pub fn event_projection(&self, attachment: &FrontendAttachment) -> OpenCodeEventProjection {
let ordinals = self
.history_ids
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.project(&attachment.history, attachment.history_cursor);
let history_tail = attachment
.history
.iter()
.zip(ordinals)
.rev()
.find(|(message, _)| message.role == Role::Assistant)
.map(|(message, ordinal)| (ordinal, message_text(message)));
OpenCodeEventProjection::new(
self.session_id.clone(),
self.workspace.clone(),
&attachment.descriptor,
self.history_ids.clone(),
history_tail,
)
}
pub fn initial_events(&self, attachment: &FrontendAttachment) -> Vec<Value> {
vec![
json!({"type": "session.updated", "properties": {"info": self.session_info(attachment)}}),
json!({
"type": "session.status",
"properties": {
"sessionID": self.session_id,
"status": {"type": match attachment.descriptor.turn_state {
supercode::FrontendTurnState::Idle => "idle",
supercode::FrontendTurnState::Busy => "busy",
}}
}
}),
]
}
pub async fn handle(&self, request: OpenCodeRequest) -> OpenCodeResponse {
match self.handle_result(request).await {
Ok(response) => response,
Err(error) => error.response(),
}
}
async fn handle_result(
&self,
request: OpenCodeRequest,
) -> Result<OpenCodeResponse, AdapterError> {
let method = request.method.to_ascii_uppercase();
let target = request.target.as_str();
let path = target.split('?').next().unwrap_or(target);
match (method.as_str(), target) {
("GET", "/agent") => Ok(OpenCodeResponse::json(200, self.agents())),
("GET", "/command") => {
let descriptor = self.runtime.describe().await?;
Ok(OpenCodeResponse::json(200, commands(&descriptor)))
}
("GET", "/config") => {
let descriptor = self.runtime.describe().await?;
Ok(OpenCodeResponse::json(200, self.config(&descriptor)))
}
("GET", "/config/providers") => {
let descriptor = self.runtime.describe().await?;
Ok(OpenCodeResponse::json(
200,
provider_projection(&descriptor, false),
))
}
("GET", "/provider") => {
let descriptor = self.runtime.describe().await?;
Ok(OpenCodeResponse::json(
200,
provider_projection(&descriptor, true),
))
}
("GET", "/provider/auth") => Ok(OpenCodeResponse::json(200, json!({}))),
("GET", "/experimental/resource" | "/mcp" | "/vcs") => {
Ok(OpenCodeResponse::json(200, json!({})))
}
("GET", "/formatter" | "/lsp") => Ok(OpenCodeResponse::json(200, json!([]))),
("GET", "/path") => Ok(OpenCodeResponse::json(200, self.paths())),
("GET", "/session/status") => {
let descriptor = self.runtime.describe().await?;
let state = match descriptor.turn_state {
supercode::FrontendTurnState::Idle => json!({"type": "idle"}),
supercode::FrontendTurnState::Busy => json!({"type": "busy"}),
};
Ok(OpenCodeResponse::json(
200,
json!({self.session_id.clone(): state}),
))
}
("GET", "/event") => {
let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
Ok(OpenCodeResponse::events(attachment))
}
("GET", _) if session_start_cursor(target).is_some() => {
let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
Ok(OpenCodeResponse::json(
200,
json!([self.session_info(&attachment)]),
))
}
("POST", "/session") if request.body.is_null() => {
let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
Ok(OpenCodeResponse::json(200, self.session_info(&attachment)))
}
("POST", _) if target == path && self.abort_target(path) => self.interrupt().await,
("POST", _) if target == path && self.message_target(path) => {
self.submit_message(&request.body).await
}
("POST", _) if target == path && permission_request_id(path).is_some() => {
self.respond_permission(path, &request.body).await
}
("GET", _) if self.session_tail(path).is_some() => {
self.handle_session_get(path, target).await
}
_ => Err(AdapterError::UnsupportedRoute(format!(
"{method} {}",
request.target
))),
}
}
async fn handle_session_get(
&self,
path: &str,
target: &str,
) -> Result<OpenCodeResponse, AdapterError> {
let tail = self
.session_tail(path)
.expect("caller checked session prefix");
let traced_target = match tail {
"/message" => format!("{path}?limit=100"),
_ => path.to_string(),
};
if target != traced_target {
return Err(AdapterError::UnsupportedRoute(format!("GET {target}")));
}
let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
match tail {
"" => Ok(OpenCodeResponse::json(200, self.session_info(&attachment))),
"/message" => Ok(OpenCodeResponse::json(
200,
self.history_messages(&attachment),
)),
"/todo" | "/diff" => Ok(OpenCodeResponse::json(200, json!([]))),
_ => Err(AdapterError::UnsupportedRoute(format!("GET {path}"))),
}
}
fn session_tail<'a>(&self, path: &'a str) -> Option<&'a str> {
let prefix = "/session/";
let rest = path.strip_prefix(prefix)?;
let (id, tail) = rest
.split_once('/')
.map_or((rest, ""), |(id, tail)| (id, tail));
if id != self.session_id {
return None;
}
if tail.is_empty() {
Some("")
} else {
Some(
path.strip_prefix(&format!("{prefix}{id}"))
.expect("prefix matches"),
)
}
}
fn message_target(&self, path: &str) -> bool {
self.session_tail(path) == Some("/message")
}
fn abort_target(&self, path: &str) -> bool {
self.session_tail(path) == Some("/abort")
}
async fn interrupt(&self) -> Result<OpenCodeResponse, AdapterError> {
let descriptor = self.runtime.describe().await?;
if !descriptor.actions.interrupt {
return Err(AdapterError::Sdk(SdkError::UnsupportedAction("interrupt")));
}
let interrupted = self.runtime.interrupt().await?;
Ok(OpenCodeResponse::json(200, json!(interrupted)))
}
async fn submit_message(&self, body: &Value) -> Result<OpenCodeResponse, AdapterError> {
let descriptor = self.runtime.describe().await?;
validate_requested_model(body, &descriptor)?;
let submitted = prompt_from_message(body, &self.workspace)?;
let mut before = self.runtime.attach(HISTORY_LIMIT).await?;
self.history_ids
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.project(&before.history, before.history_cursor);
let attachment = match descriptor.turn_state {
supercode::FrontendTurnState::Idle if descriptor.actions.submit => {
self.history_ids
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.register_client_user(&submitted);
if submitted.image_urls.is_empty() {
self.runtime.submit(submitted.prompt).await?;
} else {
self.runtime
.submit_with_images(submitted.prompt, submitted.image_urls)
.await?;
}
self.runtime.attach(HISTORY_LIMIT).await?
}
supercode::FrontendTurnState::Busy if descriptor.actions.steer => {
let baseline_cursor = before.history_cursor;
self.runtime.steer(submitted.prompt).await?;
wait_for_completed_terminal(self.runtime.as_ref(), &mut before, baseline_cursor)
.await?
}
supercode::FrontendTurnState::Idle => {
return Err(AdapterError::Sdk(SdkError::UnsupportedAction("submit")))
}
supercode::FrontendTurnState::Busy => {
return Err(AdapterError::Sdk(SdkError::UnsupportedAction("steer")))
}
};
let messages = self.history_messages(&attachment);
let added = appended_message_count(&before.history, &attachment.history);
let assistant = messages
.as_array()
.and_then(|messages| {
messages
.get(messages.len().saturating_sub(added)..)
.unwrap_or_default()
.iter()
.rev()
.find(|message| message["info"]["role"] == "assistant")
})
.cloned()
.ok_or_else(|| AdapterError::InvalidRequest {
route: "POST /session/{session_id}/message".into(),
message: "runtime completed without an assistant history message".into(),
})?;
Ok(OpenCodeResponse::json(200, assistant))
}
async fn respond_permission(
&self,
path: &str,
body: &Value,
) -> Result<OpenCodeResponse, AdapterError> {
let request_id = permission_request_id(path).expect("route guard parsed request id");
let descriptor = self.runtime.describe().await?;
if !descriptor.actions.respond {
return Err(AdapterError::Sdk(SdkError::UnsupportedAction("respond")));
}
let reply = exact_string_field(body, &["reply"], "reply")?;
let decision = match reply {
"once" => FrontendApprovalDecision::Allow,
"always" => FrontendApprovalDecision::AllowForSession,
"reject" => FrontendApprovalDecision::Deny,
_ => {
return Err(AdapterError::InvalidRequest {
route: "POST /permission/{permission_id}/reply".into(),
message: "reply must be `once`, `always`, or `reject`".into(),
})
}
};
self.runtime
.respond(FrontendResponse::Approval {
request_id,
decision,
})
.await?;
Ok(OpenCodeResponse::json(200, json!(true)))
}
fn session_info(&self, attachment: &FrontendAttachment) -> Value {
let title = attachment
.history
.iter()
.rev()
.find_map(|message| message.content.as_deref())
.map(|text| truncate(text, 80))
.filter(|text| !text.is_empty())
.unwrap_or_else(|| format!("Supercode runtime {}", self.runtime_id));
json!({
"id": self.session_id,
"slug": format!("supercode-{}", &self.session_id[4..12]),
"projectID": "global",
"directory": path_text(&self.workspace),
"title": title,
"version": OPENCODE_CLI_VERSION,
"summary": {"additions": 0, "deletions": 0, "files": 0},
"time": {"created": 0, "updated": attachment.history_cursor},
})
}
fn agents(&self) -> Value {
json!([{
"name": "build",
"description": "Attached Supercode SDK runtime",
"options": {},
"permission": [],
"mode": "primary",
"native": false,
}])
}
fn config(&self, descriptor: &FrontendRuntimeDescriptor) -> Value {
let (provider, model) = provider_model(descriptor);
json!({
"$schema": "https://opencode.ai/config.json",
"share": "disabled",
"autoupdate": false,
"enabled_providers": [provider],
"model": format!("{provider}/{model}"),
"small_model": format!("{provider}/{model}"),
"provider": {},
"permission": {"edit": "ask"},
"agent": {},
"mode": {},
"plugin": [],
"command": {},
"username": "supercode",
})
}
fn paths(&self) -> Value {
let workspace = path_text(&self.workspace);
json!({
"home": workspace,
"state": workspace,
"config": workspace,
"worktree": workspace,
"directory": workspace,
})
}
fn history_messages(&self, attachment: &FrontendAttachment) -> Value {
let mut state = self
.history_ids
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let ordinals = state.project(&attachment.history, attachment.history_cursor);
let identities = ordinals
.iter()
.map(|ordinal| state.identity(&self.session_id, *ordinal))
.collect::<Vec<_>>();
drop(state);
history_messages(
&attachment.history,
&ordinals,
&identities,
&self.session_id,
&self.workspace,
&attachment.descriptor,
)
}
}
async fn wait_for_completed_terminal(
runtime: &dyn SdkRuntime,
attachment: &mut FrontendAttachment,
baseline_cursor: u64,
) -> Result<FrontendAttachment, AdapterError> {
loop {
let event = attachment
.next_event()
.await
.map_err(|error| AdapterError::Sdk(SdkError::Transport(error.to_string())))?;
if matches!(
event.kind.as_str(),
"turn_succeeded" | "turn_failed" | "turn_interrupted"
) {
let completed = runtime.attach(HISTORY_LIMIT).await?;
if completed.history_cursor > baseline_cursor
&& event.sequence > completed.history_cursor
{
return Ok(completed);
}
}
}
}
fn appended_message_count(previous: &[ChatMessage], current: &[ChatMessage]) -> usize {
let overlap = (0..=previous.len().min(current.len()))
.rev()
.find(|count| previous[previous.len() - *count..] == current[..*count])
.unwrap_or(0);
current.len().saturating_sub(overlap)
}
fn session_start_cursor(target: &str) -> Option<u64> {
let cursor = target.strip_prefix("/session?start=")?;
if cursor.is_empty() || !cursor.bytes().all(|byte| byte.is_ascii_digit()) {
return None;
}
cursor.parse().ok()
}
fn permission_request_id(path: &str) -> Option<u64> {
let encoded = path
.strip_prefix("/permission/per_")?
.strip_suffix("/reply")?;
(encoded.len() == 16)
.then(|| u64::from_str_radix(encoded, 16).ok())
.flatten()
}
fn permission_id(request_id: u64) -> String {
format!("per_{request_id:016x}")
}
struct SubmittedMessage {
prompt: String,
image_urls: Vec<String>,
message_id: String,
part_id: String,
}
fn prompt_from_message(body: &Value, workspace: &Path) -> Result<SubmittedMessage, AdapterError> {
let object = exact_object_fields(body, &["messageID", "agent", "model", "parts"])?;
for required in ["messageID", "agent", "model", "parts"] {
if !object.contains_key(required) {
return Err(invalid_message(format!("missing `{required}`")));
}
}
if object["agent"] != "build" {
return Err(invalid_message("agent must be `build`"));
}
let message_id = object["messageID"]
.as_str()
.filter(|id| !id.is_empty())
.ok_or_else(|| invalid_message("messageID must be a nonempty string"))?;
if !message_id.starts_with("msg_") {
return Err(invalid_message(
"messageID must use the traced `msg_` prefix",
));
}
let model = exact_object_fields(&object["model"], &["providerID", "modelID"])?;
for field in ["providerID", "modelID"] {
if model
.get(field)
.and_then(Value::as_str)
.is_none_or(str::is_empty)
{
return Err(invalid_message(format!(
"model.{field} must be a nonempty string"
)));
}
}
let parts = object["parts"]
.as_array()
.ok_or_else(|| invalid_message("parts must be an array"))?;
if parts.is_empty() {
return Err(invalid_message("message requires at least one part"));
}
let mut text = Vec::new();
let mut part_id = None;
let mut image_urls = Vec::new();
for part in parts {
match part.get("type").and_then(Value::as_str) {
Some("text") => {
let part = exact_object_fields(part, &["id", "type", "text"])?;
let value = part
.get("text")
.and_then(Value::as_str)
.ok_or_else(|| invalid_message("text part requires string text"))?;
let id = traced_part_id(part)?;
part_id.get_or_insert_with(|| id.to_string());
text.push(value.to_string());
}
Some("file") => {
let part = exact_object_fields(
part,
&["id", "type", "mime", "url", "filename", "source"],
)?;
if let Some(id) = part.get("id") {
let id = id
.as_str()
.filter(|id| id.starts_with("prt_") && id.len() > 4)
.ok_or_else(|| {
invalid_message("file part id must use the traced `prt_` prefix")
})?;
part_id.get_or_insert_with(|| id.to_string());
}
let resolved = resolve_file_part(part, workspace)?;
if let Some(block) = resolved.text_block {
text.push(block);
}
if let Some(image_url) = resolved.image_url {
image_urls.push(image_url);
}
}
Some(kind) => {
return Err(invalid_message(format!(
"unsupported OpenCode message part type `{kind}`"
)))
}
None => return Err(invalid_message("message part requires string `type`")),
}
}
let prompt = text.join("\n");
if prompt.is_empty() && image_urls.is_empty() {
return Err(invalid_message("message contains no usable input"));
}
Ok(SubmittedMessage {
prompt,
image_urls,
message_id: message_id.to_string(),
part_id: part_id.unwrap_or_else(|| stable_id("part", "prt", message_id, 24)),
})
}
fn traced_part_id(part: &serde_json::Map<String, Value>) -> Result<&str, AdapterError> {
part.get("id")
.and_then(Value::as_str)
.filter(|id| id.starts_with("prt_") && id.len() > 4)
.ok_or_else(|| invalid_message("text part requires a traced `prt_` string id"))
}
struct ResolvedFilePart {
text_block: Option<String>,
image_url: Option<String>,
}
fn resolve_file_part(
part: &serde_json::Map<String, Value>,
workspace: &Path,
) -> Result<ResolvedFilePart, AdapterError> {
const MAX_ATTACHMENT_BYTES: u64 = 10 * 1024 * 1024;
let mime = part
.get("mime")
.and_then(Value::as_str)
.filter(|mime| !mime.is_empty())
.ok_or_else(|| invalid_message("file part requires nonempty string `mime`"))?;
let raw_url = part
.get("url")
.and_then(Value::as_str)
.filter(|url| !url.is_empty())
.ok_or_else(|| invalid_message("file part requires nonempty string `url`"))?;
let filename = part
.get("filename")
.and_then(Value::as_str)
.filter(|name| !name.is_empty())
.unwrap_or("attachment");
if raw_url.starts_with("data:") {
let bytes = decode_data_uri(raw_url, mime)?;
if mime.starts_with("image/") {
return Ok(ResolvedFilePart {
text_block: None,
image_url: Some(raw_url.to_string()),
});
}
let value = String::from_utf8(bytes)
.map_err(|_| invalid_message("non-image data URI attachment must be UTF-8 text"))?;
return Ok(ResolvedFilePart {
text_block: Some(format!("[file: {filename}]\n{value}")),
image_url: None,
});
}
let parsed = Url::parse(raw_url)
.map_err(|_| invalid_message("file part url must be a data:, file:, or https: URL"))?;
if matches!(parsed.scheme(), "http" | "https") {
if !mime.starts_with("image/") {
return Err(invalid_message(
"remote non-image attachments are not fetched by the runtime",
));
}
return Ok(ResolvedFilePart {
text_block: None,
image_url: Some(raw_url.to_string()),
});
}
if parsed.scheme() != "file" {
return Err(invalid_message("file part URL scheme is not supported"));
}
let path = parsed
.to_file_path()
.map_err(|_| invalid_message("file part URL is not a valid runtime file path"))?;
let path = resolve_runtime_path(workspace, &path)?;
let metadata = std::fs::metadata(&path)
.map_err(|error| invalid_message(format!("runtime attachment is unreadable: {error}")))?;
if !metadata.is_file() || metadata.len() > MAX_ATTACHMENT_BYTES {
return Err(invalid_message(
"runtime attachment must be a file no larger than 10 MiB",
));
}
let bytes = std::fs::read(&path)
.map_err(|error| invalid_message(format!("runtime attachment is unreadable: {error}")))?;
if mime.starts_with("image/") {
let encoded = base64::engine::general_purpose::STANDARD.encode(bytes);
Ok(ResolvedFilePart {
text_block: None,
image_url: Some(format!("data:{mime};base64,{encoded}")),
})
} else {
let value = String::from_utf8(bytes)
.map_err(|_| invalid_message("non-image runtime attachment must be UTF-8 text"))?;
Ok(ResolvedFilePart {
text_block: Some(format!("[file: {filename}]\n{value}")),
image_url: None,
})
}
}
fn decode_data_uri(raw: &str, mime: &str) -> Result<Vec<u8>, AdapterError> {
let expected = format!("data:{mime};base64,");
let encoded = raw
.strip_prefix(&expected)
.ok_or_else(|| invalid_message("data URI mime must match file part mime and use base64"))?;
base64::engine::general_purpose::STANDARD
.decode(encoded)
.map_err(|_| invalid_message("file part contains invalid base64 data"))
}
fn resolve_runtime_path(workspace: &Path, requested: &Path) -> Result<PathBuf, AdapterError> {
let logical = logical_runtime_path(&path_text(workspace), &path_text(requested))?;
let workspace = std::fs::canonicalize(workspace).map_err(|error| {
invalid_message(format!("SDK runtime workspace is unavailable: {error}"))
})?;
let logical_path = PathBuf::from(logical);
let candidate = if logical_path.is_absolute() {
logical_path
} else {
workspace.join(logical_path)
};
let candidate = std::fs::canonicalize(candidate)
.map_err(|error| invalid_message(format!("runtime attachment is unavailable: {error}")))?;
if !candidate.starts_with(&workspace) {
return Err(invalid_message(
"attachment path escapes the SDK runtime workspace",
));
}
Ok(candidate)
}
fn logical_runtime_path(workspace: &str, requested: &str) -> Result<String, AdapterError> {
let workspace = normalized_logical_path(workspace)?;
let requested = requested.replace('\\', "/");
let absolute = requested.starts_with('/')
|| requested
.as_bytes()
.get(1)
.is_some_and(|separator| *separator == b':');
let candidate = if absolute {
requested
} else {
format!("{workspace}/{requested}")
};
let candidate = normalized_logical_path(&candidate)?;
let prefix = format!("{workspace}/");
if candidate != workspace && !candidate.starts_with(&prefix) {
return Err(invalid_message(
"attachment path escapes the SDK runtime workspace",
));
}
Ok(candidate)
}
fn normalized_logical_path(raw: &str) -> Result<String, AdapterError> {
let raw = raw.replace('\\', "/");
let (prefix, tail) = if raw.starts_with('/') {
("/".to_string(), raw.trim_start_matches('/'))
} else if raw
.as_bytes()
.get(1)
.is_some_and(|separator| *separator == b':')
{
(
raw[..2].to_ascii_uppercase(),
raw[2..].trim_start_matches('/'),
)
} else {
(String::new(), raw.as_str())
};
let mut segments = Vec::new();
for segment in tail.split('/') {
match segment {
"" | "." => {}
".." => {
if segments.pop().is_none() {
return Err(invalid_message("attachment path escapes its path root"));
}
}
value => segments.push(value),
}
}
let joined = segments.join("/");
Ok(match prefix.as_str() {
"/" => format!("/{joined}"),
"" => joined,
drive => format!("{drive}/{joined}"),
})
}
fn validate_requested_model(
body: &Value,
descriptor: &FrontendRuntimeDescriptor,
) -> Result<(), AdapterError> {
let model = body
.get("model")
.and_then(Value::as_object)
.ok_or_else(|| invalid_message("model must be an object"))?;
let (provider_id, model_id) = provider_model(descriptor);
if model.get("providerID").and_then(Value::as_str) != Some(provider_id.as_str())
|| model.get("modelID").and_then(Value::as_str) != Some(model_id.as_str())
{
return Err(invalid_message(
"model must match the attached SDK runtime descriptor",
));
}
Ok(())
}
fn exact_object_fields<'a>(
value: &'a Value,
allowed: &[&str],
) -> Result<&'a serde_json::Map<String, Value>, AdapterError> {
let object = value
.as_object()
.ok_or_else(|| invalid_message("body must be a JSON object"))?;
if let Some(unexpected) = object.keys().find(|key| !allowed.contains(&key.as_str())) {
return Err(invalid_message(format!("unexpected field `{unexpected}`")));
}
Ok(object)
}
fn exact_string_field<'a>(
value: &'a Value,
allowed: &[&str],
field: &str,
) -> Result<&'a str, AdapterError> {
let object = exact_object_fields(value, allowed)?;
object
.get(field)
.and_then(Value::as_str)
.ok_or_else(|| invalid_message(format!("`{field}` must be a string")))
}
fn invalid_message(message: impl Into<String>) -> AdapterError {
AdapterError::InvalidRequest {
route: "POST /session/{session_id}/message".into(),
message: message.into(),
}
}
#[derive(Default)]
struct HistoryIdentityState {
previous: Vec<ChatMessage>,
ordinals: Vec<u64>,
last_history_cursor: Option<u64>,
snapshots: BTreeMap<u64, Vec<u64>>,
pending_live: Vec<PendingHistoryIdentity>,
projected_live: Vec<PendingHistoryIdentity>,
live_sources: BTreeMap<(u64, bool), u64>,
client_users: Vec<SubmittedMessage>,
explicit_ids: BTreeMap<u64, (String, String)>,
next_ordinal: u64,
}
struct PendingHistoryIdentity {
ordinal: u64,
role: Role,
content: String,
projected_through: Option<u64>,
}
struct LiveMessage {
ordinal: u64,
message_id: String,
part_id: String,
text: String,
created: u64,
}
pub struct OpenCodeEventProjection {
session_id: String,
workspace: PathBuf,
provider: String,
model: String,
active: Option<LiveMessage>,
history_ids: Arc<Mutex<HistoryIdentityState>>,
history_tail: Option<LiveMessage>,
}
impl OpenCodeEventProjection {
fn new(
session_id: String,
workspace: PathBuf,
descriptor: &FrontendRuntimeDescriptor,
history_ids: Arc<Mutex<HistoryIdentityState>>,
history_tail: Option<(u64, String)>,
) -> Self {
let (provider, model) = provider_model(descriptor);
let history_tail = history_tail.map(|(ordinal, text)| {
let (message_id, part_id) = message_identity(&session_id, ordinal);
LiveMessage {
ordinal,
message_id,
part_id,
text,
created: ordinal,
}
});
Self {
session_id,
workspace,
provider,
model,
active: None,
history_ids,
history_tail,
}
}
pub fn project(&mut self, event: &SdkEvent) -> Vec<Value> {
match event.kind.as_str() {
"user_message" => self.user_message(event),
"turn_started" => {
let mut events = vec![self.status("busy")];
events.extend(self.ensure_active_events(event.sequence));
events
}
"text_delta" => self.text_delta(event),
"tool_call_started" | "tool_call_completed" => self.tool_event(event),
"request" => self.request(event),
"request_resolved" => self.request_resolved(event),
"turn_succeeded" => self.turn_finished(event, "stop"),
"turn_interrupted" => self.turn_finished(event, "abort"),
"turn_failed" => self.turn_finished(event, "error"),
_ => vec![json!({
"type": "supercode.event",
"properties": {"sequence": event.sequence, "kind": event.kind, "payload": event.payload}
})],
}
}
fn ensure_active_events(&mut self, sequence: u64) -> Vec<Value> {
if self.active.is_some() {
return Vec::new();
}
let active = self.new_live_message(sequence);
self.history_tail = None;
let events = vec![
json!({"type": "message.updated", "properties": {"info": self.assistant_info(&active, None, None)}}),
json!({"type": "message.part.updated", "properties": {"part": {
"id": active.part_id, "sessionID": self.session_id, "messageID": active.message_id,
"type": "text", "text": "", "time": {"start": active.created}
}}}),
];
self.active = Some(active);
events
}
fn user_message(&self, event: &SdkEvent) -> Vec<Value> {
let text = event
.payload
.get("text")
.and_then(Value::as_str)
.unwrap_or_default();
let (_ordinal, message_id, part_id) = self.reserve_live(Role::User, text, event.sequence);
vec![
json!({"type": "message.updated", "properties": {"info": {
"id": message_id, "sessionID": self.session_id, "role": "user",
"time": {"created": event.sequence}, "agent": "build",
"model": {"providerID": self.provider, "modelID": self.model}
}}}),
json!({"type": "message.part.updated", "properties": {"part": {
"id": part_id, "sessionID": self.session_id, "messageID": message_id,
"type": "text", "text": text
}}}),
]
}
fn text_delta(&mut self, event: &SdkEvent) -> Vec<Value> {
let delta = event
.payload
.get("text")
.and_then(Value::as_str)
.unwrap_or_default();
let mut events = self.ensure_active_events(event.sequence);
let session_id = self.session_id.clone();
let active = self.active.as_mut().expect("active message was created");
active.text.push_str(delta);
events.push(json!({"type": "message.part.delta", "properties": {
"sessionID": session_id, "messageID": active.message_id,
"partID": active.part_id, "field": "text", "delta": delta
}}));
events
}
fn tool_event(&mut self, event: &SdkEvent) -> Vec<Value> {
let mut events = self.ensure_active_events(event.sequence);
let session_id = self.session_id.clone();
let active = self.active.as_ref().expect("active message was created");
let call_id = event
.payload
.get("id")
.and_then(Value::as_str)
.unwrap_or("unknown");
let tool = event
.payload
.get("name")
.and_then(Value::as_str)
.unwrap_or("tool");
let completed = event.kind == "tool_call_completed";
let part_id = stable_id("live-tool-part", "prt", call_id, 24);
events.push(json!({"type": "message.part.updated", "properties": {"part": {
"id": part_id, "sessionID": session_id, "messageID": active.message_id,
"type": "tool", "callID": call_id, "tool": tool,
"state": if completed {
json!({"status": "completed", "input": event.payload.get("arguments").cloned().unwrap_or(Value::Null),
"output": event.payload.get("output").cloned().unwrap_or(Value::Null),
"time": {"start": event.sequence, "end": event.sequence}})
} else {
json!({"status": "running", "input": event.payload.get("arguments").cloned().unwrap_or(Value::Null),
"time": {"start": event.sequence}})
}
}}}));
events
}
fn request(&self, event: &SdkEvent) -> Vec<Value> {
let request = &event.payload["request"];
if request["kind"] != "approval" {
return vec![self.opaque(event)];
}
let Some(request_id) = request["id"].as_u64() else {
return vec![self.opaque(event)];
};
let payload = &request["payload"];
let subject = payload.get("subject").cloned().unwrap_or(Value::Null);
vec![json!({"type": "permission.asked", "properties": {
"id": permission_id(request_id), "sessionID": self.session_id,
"permission": payload.get("tool").cloned().unwrap_or_else(|| json!("tool")),
"patterns": if subject.is_null() { json!([]) } else { json!([subject]) },
"metadata": {"supercode": {"request": request, "sequence": event.sequence}},
"always": ["*"]
}})]
}
fn request_resolved(&self, event: &SdkEvent) -> Vec<Value> {
let Some(request_id) = event.payload.get("request_id").and_then(Value::as_u64) else {
return vec![self.opaque(event)];
};
let decision = event
.payload
.pointer("/response/decision")
.and_then(Value::as_str);
let reply = match decision {
Some("allow") => "once",
Some("allow_for_session") => "always",
_ => "reject",
};
vec![json!({"type": "permission.replied", "properties": {
"sessionID": self.session_id, "requestID": permission_id(request_id), "reply": reply
}})]
}
fn turn_finished(&mut self, event: &SdkEvent, finish: &str) -> Vec<Value> {
let reply = event
.payload
.get("reply")
.or_else(|| event.payload.get("message"))
.and_then(Value::as_str)
.unwrap_or_default();
let history_match = self.active.is_none()
&& finish == "stop"
&& self
.history_tail
.as_ref()
.is_some_and(|history| history.text == reply);
let mut events = if history_match {
Vec::new()
} else {
self.ensure_active_events(event.sequence)
};
let mut active = if history_match {
self.history_tail.take().expect("history match was checked")
} else {
self.active.take().expect("active message was created")
};
if active.text.is_empty() {
active.text = reply.to_string();
}
if !history_match {
self.update_live(active.ordinal, Role::Assistant, &active.text);
}
events.extend([
json!({"type": "message.part.updated", "properties": {"part": {
"id": active.part_id, "sessionID": self.session_id, "messageID": active.message_id,
"type": "text", "text": active.text,
"time": {"start": active.created, "end": event.sequence}
}}}),
json!({"type": "message.updated", "properties": {"info": self.assistant_info(&active, Some(event.sequence), Some(finish))}}),
self.status("idle"),
json!({"type": "session.idle", "properties": {"sessionID": self.session_id}}),
]);
if finish == "error" {
let completed = events.len() - 3;
events[completed]["properties"]["info"]["error"] = json!({
"name": "APIError", "data": {"message": event.payload.get("message").cloned().unwrap_or_default()}
});
}
events
}
fn new_live_message(&self, sequence: u64) -> LiveMessage {
let (ordinal, message_id, part_id) = self.reserve_live(Role::Assistant, "", sequence);
LiveMessage {
ordinal,
part_id,
message_id,
text: String::new(),
created: sequence,
}
}
fn assistant_info(
&self,
active: &LiveMessage,
completed: Option<u64>,
finish: Option<&str>,
) -> Value {
let mut info = json!({
"id": active.message_id, "sessionID": self.session_id, "role": "assistant",
"time": {"created": active.created},
"modelID": self.model, "providerID": self.provider, "mode": "build", "agent": "build",
"path": {"cwd": path_text(&self.workspace), "root": path_text(&self.workspace)},
"cost": 0, "tokens": {"input": 0, "output": 0, "reasoning": 0, "cache": {"read": 0, "write": 0}},
});
if let Some(completed) = completed {
info["time"]["completed"] = json!(completed);
}
if let Some(finish) = finish {
info["finish"] = json!(finish);
}
info
}
fn status(&self, status: &str) -> Value {
json!({"type": "session.status", "properties": {
"sessionID": self.session_id, "status": {"type": status}
}})
}
fn opaque(&self, event: &SdkEvent) -> Value {
json!({"type": "supercode.event", "properties": {
"sequence": event.sequence, "kind": event.kind, "payload": event.payload
}})
}
fn reserve_live(
&self,
role: Role,
content: &str,
source_sequence: u64,
) -> (u64, String, String) {
let mut state = self
.history_ids
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let ordinal = state.reserve_live(role, content, source_sequence);
let (message_id, part_id) = state.identity(&self.session_id, ordinal);
(ordinal, message_id, part_id)
}
fn update_live(&self, ordinal: u64, role: Role, content: &str) {
self.history_ids
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.update_live(ordinal, role, content);
}
}
impl HistoryIdentityState {
fn project(&mut self, current: &[ChatMessage], history_cursor: u64) -> Vec<u64> {
if let Some(projected) = self
.snapshots
.get(&history_cursor)
.filter(|projected| projected.len() == current.len())
{
return projected.clone();
}
if self
.last_history_cursor
.is_some_and(|cursor| history_cursor < cursor)
{
let overlap = (0..=self.previous.len().min(current.len()))
.rev()
.find(|count| current[current.len() - *count..] == self.previous[..*count])
.unwrap_or(0);
let missing = current.len().saturating_sub(overlap);
let mut projected = (0..missing).map(|_| self.allocate()).collect::<Vec<_>>();
projected.extend_from_slice(&self.ordinals[..overlap]);
self.remember_snapshot(history_cursor, &projected);
return projected;
}
let overlap = (0..=self.previous.len().min(current.len()))
.rev()
.find(|count| self.previous[self.previous.len() - *count..] == current[..*count])
.unwrap_or(0);
let mut projected = self.ordinals[self.ordinals.len() - overlap..].to_vec();
for message in ¤t[overlap..] {
let content = message_text(message);
let pending = self
.pending_live
.iter()
.position(|pending| pending.role == message.role && pending.content == content)
.map(|index| self.pending_live.remove(index).ordinal);
let allocated = pending.is_none();
let ordinal = pending.unwrap_or_else(|| self.allocate());
self.claim_client_user(ordinal, &message.role, &content);
if allocated && matches!(message.role, Role::User | Role::Assistant) {
self.projected_live.push(PendingHistoryIdentity {
ordinal,
role: message.role,
content,
projected_through: Some(history_cursor),
});
}
projected.push(ordinal);
}
self.projected_live
.retain(|pending| projected.binary_search(&pending.ordinal).is_ok());
self.previous = current.to_vec();
self.ordinals = projected.clone();
self.last_history_cursor = Some(history_cursor);
self.remember_snapshot(history_cursor, &projected);
projected
}
fn remember_snapshot(&mut self, history_cursor: u64, projected: &[u64]) {
const SNAPSHOT_LIMIT: usize = 64;
self.snapshots.insert(history_cursor, projected.to_vec());
while self.snapshots.len() > SNAPSHOT_LIMIT {
self.snapshots.pop_first();
}
}
fn register_client_user(&mut self, submitted: &SubmittedMessage) {
self.client_users.push(SubmittedMessage {
prompt: submitted.prompt.clone(),
image_urls: submitted.image_urls.clone(),
message_id: submitted.message_id.clone(),
part_id: submitted.part_id.clone(),
});
}
fn reserve_live(&mut self, role: Role, content: &str, source_sequence: u64) -> u64 {
let source = (source_sequence, role == Role::User);
if let Some(ordinal) = self.live_sources.get(&source) {
return *ordinal;
}
let matches = |pending: &PendingHistoryIdentity| {
pending.role == role
&& pending
.projected_through
.is_some_and(|cursor| source_sequence <= cursor)
};
let projected = if role == Role::Assistant && content.is_empty() {
self.projected_live.iter().rposition(matches)
} else {
self.projected_live
.iter()
.position(|pending| matches(pending) && pending.content == content)
};
let ordinal = projected
.map(|index| self.projected_live.remove(index).ordinal)
.unwrap_or_else(|| self.allocate());
self.live_sources.insert(source, ordinal);
self.claim_client_user(ordinal, &role, content);
if projected.is_none() {
self.pending_live.push(PendingHistoryIdentity {
ordinal,
role,
content: content.to_string(),
projected_through: None,
});
}
ordinal
}
fn update_live(&mut self, ordinal: u64, role: Role, content: &str) {
if let Some(pending) = self
.pending_live
.iter_mut()
.find(|pending| pending.ordinal == ordinal)
{
pending.role = role;
pending.content = content.to_string();
}
}
fn allocate(&mut self) -> u64 {
let ordinal = self.next_ordinal;
self.next_ordinal = self.next_ordinal.saturating_add(1);
ordinal
}
fn identity(&self, session_id: &str, ordinal: u64) -> (String, String) {
self.explicit_ids
.get(&ordinal)
.cloned()
.unwrap_or_else(|| message_identity(session_id, ordinal))
}
fn claim_client_user(&mut self, ordinal: u64, role: &Role, content: &str) {
if *role != Role::User || self.explicit_ids.contains_key(&ordinal) {
return;
}
if let Some(index) = self
.client_users
.iter()
.position(|submitted| submitted.prompt == content)
{
let submitted = self.client_users.remove(index);
self.explicit_ids
.insert(ordinal, (submitted.message_id, submitted.part_id));
}
}
}
fn provider_projection(descriptor: &FrontendRuntimeDescriptor, connected: bool) -> Value {
let (provider, model) = provider_model(descriptor);
let row = json!({
"id": provider,
"name": "Supercode runtime",
"env": [],
"options": {},
"source": "custom",
"models": {
model.clone(): {
"id": model,
"api": {"id": model, "npm": "@ai-sdk/openai-compatible"},
"status": "active",
"name": descriptor.model,
"providerID": provider,
"capabilities": {
"temperature": false,
"reasoning": true,
"attachment": false,
"toolcall": true,
"input": {"text": true, "audio": false, "image": false, "video": false, "pdf": false},
"output": {"text": true, "audio": false, "image": false, "video": false, "pdf": false},
"interleaved": false
},
"cost": {"input": 0, "output": 0, "cache": {"read": 0, "write": 0}},
"options": {},
"limit": {"context": 0, "output": 0},
"headers": {},
"family": "",
"release_date": "",
"variants": {}
}
}
});
if connected {
json!({"all": [row], "default": {provider.clone(): model}, "connected": [provider]})
} else {
json!({"providers": [row], "default": {provider: model}})
}
}
fn provider_model(descriptor: &FrontendRuntimeDescriptor) -> (String, String) {
descriptor
.model
.split_once('/')
.map(|(provider, model)| (sanitize_id(provider), sanitize_id(model)))
.unwrap_or_else(|| ("supercode".into(), sanitize_id(&descriptor.model)))
}
fn commands(descriptor: &FrontendRuntimeDescriptor) -> Value {
Value::Array(
descriptor
.commands
.iter()
.map(|command| {
json!({
"name": command.name,
"description": command.description,
"source": "command",
"template": "$ARGUMENTS",
"hints": ["$ARGUMENTS"],
})
})
.collect(),
)
}
fn history_messages(
history: &[ChatMessage],
ordinals: &[u64],
identities: &[(String, String)],
session_id: &str,
workspace: &Path,
descriptor: &FrontendRuntimeDescriptor,
) -> Value {
debug_assert_eq!(history.len(), ordinals.len());
debug_assert_eq!(history.len(), identities.len());
Value::Array(
history
.iter()
.zip(ordinals)
.zip(identities)
.map(|((message, ordinal), identity)| {
history_message(
message, session_id, workspace, *ordinal, identity, descriptor,
)
})
.collect(),
)
}
fn history_message(
message: &ChatMessage,
session_id: &str,
workspace: &Path,
ordinal: u64,
identity: &(String, String),
descriptor: &FrontendRuntimeDescriptor,
) -> Value {
let (message_id, part_id) = identity;
let (provider, model) = provider_model(descriptor);
let content = message_text(message);
match message.role {
Role::User => json!({
"info": {
"role": "user", "time": {"created": ordinal}, "summary": {"diffs": []},
"agent": "build", "model": {"providerID": provider, "modelID": model},
"id": message_id, "sessionID": session_id
},
"parts": [{"type": "text", "text": content, "id": part_id, "sessionID": session_id, "messageID": message_id}]
}),
Role::Assistant | Role::System | Role::Tool => {
let content = match message.role {
Role::System => format!("[system]\n{content}"),
Role::Tool => format!("[tool result]\n{content}"),
_ => content,
};
json!({
"info": {
"role": "assistant", "time": {"created": ordinal, "completed": ordinal},
"modelID": model, "providerID": provider, "mode": "build", "agent": "build",
"path": {"cwd": path_text(workspace), "root": path_text(workspace)},
"cost": 0, "tokens": {"input": 0, "output": 0, "reasoning": 0, "cache": {"read": 0, "write": 0}},
"finish": "stop", "id": message_id, "sessionID": session_id
},
"parts": [{"type": "text", "text": content, "time": {"start": ordinal, "end": ordinal}, "id": part_id, "sessionID": session_id, "messageID": message_id}]
})
}
}
}
fn message_identity(session_id: &str, ordinal: u64) -> (String, String) {
let message_id = stable_id("message", "msg", &format!("{session_id}:{ordinal}"), 24);
let part_id = stable_id("part", "prt", &format!("{message_id}:0"), 24);
(message_id, part_id)
}
fn message_text(message: &ChatMessage) -> String {
if let Some(content) = &message.content {
return content.clone();
}
if let Some(parts) = &message.content_parts {
return parts
.iter()
.map(Value::to_string)
.collect::<Vec<_>>()
.join("\n");
}
if let Some(tool_calls) = &message.tool_calls {
return serde_json::to_string(tool_calls).unwrap_or_else(|_| "[tool calls]".into());
}
String::new()
}
fn stable_id(domain: &str, prefix: &str, source: &str, digits: usize) -> String {
let mut input = Vec::with_capacity(domain.len() + source.len() + 1);
input.extend_from_slice(domain.as_bytes());
input.push(0);
input.extend_from_slice(source.as_bytes());
let digest = blake3::hash(&input).to_hex().to_string();
format!("{prefix}_{}", &digest[..digits])
}
fn sanitize_id(value: &str) -> String {
let sanitized = value
.chars()
.map(|character| {
if character.is_ascii_alphanumeric() || matches!(character, '-' | '_' | '.') {
character
} else {
'-'
}
})
.collect::<String>();
if sanitized.is_empty() {
"runtime".into()
} else {
sanitized
}
}
fn path_text(path: &Path) -> String {
path.to_string_lossy().replace('\\', "/")
}
fn truncate(text: &str, limit: usize) -> String {
text.chars().take(limit).collect::<String>()
}
#[cfg(test)]
mod tests {
use std::collections::{BTreeMap, VecDeque};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use async_trait::async_trait;
use supercode::{
FrontendActions, FrontendAttachSnapshot, FrontendConnectionState,
FrontendDisplayCapabilities, FrontendTurnState,
};
use tokio::sync::broadcast;
use super::*;
struct FixtureRuntime {
history: Mutex<Vec<ChatMessage>>,
descriptor: FrontendRuntimeDescriptor,
events: broadcast::Sender<supercode::FrontendEvent>,
submissions: Mutex<Vec<String>>,
image_submissions: Mutex<Vec<Vec<String>>>,
steers: Mutex<Vec<String>>,
responses: Mutex<Vec<FrontendResponse>>,
interrupts: AtomicUsize,
steer_accepting: AtomicBool,
}
impl FixtureRuntime {
fn new() -> Arc<Self> {
Self::with_turn_state(FrontendTurnState::Idle)
}
fn with_turn_state(turn_state: FrontendTurnState) -> Arc<Self> {
let (events, _) = broadcast::channel(16);
Arc::new(Self {
history: Mutex::new(vec![
ChatMessage::system("preserve system context"),
ChatMessage::user("hello from Claude"),
ChatMessage::assistant("continued through GLM"),
]),
descriptor: FrontendRuntimeDescriptor {
schema_version: 2,
session_id: "runtime-1".into(),
source_harness: Some("claude-code".into()),
emulation_profile: Some("claude-code".into()),
active_modules: Vec::new(),
commands: Vec::new(),
operations: Vec::new(),
actions: FrontendActions {
submit: true,
interrupt: true,
steer: true,
respond: true,
detach: true,
close: false,
},
display: FrontendDisplayCapabilities {
event_kinds: vec!["assistant_delta".into()],
opaque_fallback: true,
},
model: "openrouter/glm-5.2".into(),
turn_state,
connection_state: FrontendConnectionState::Connected,
extensions: BTreeMap::new(),
},
events,
submissions: Mutex::new(Vec::new()),
image_submissions: Mutex::new(Vec::new()),
steers: Mutex::new(Vec::new()),
responses: Mutex::new(Vec::new()),
interrupts: AtomicUsize::new(0),
steer_accepting: AtomicBool::new(true),
})
}
}
#[async_trait]
impl SdkRuntime for FixtureRuntime {
async fn describe(&self) -> Result<FrontendRuntimeDescriptor, SdkError> {
Ok(self.descriptor.clone())
}
async fn attach(&self, history_limit: usize) -> Result<FrontendAttachment, SdkError> {
let history = self
.history
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
let start = history.len().saturating_sub(history_limit);
let replay = if self.descriptor.turn_state == FrontendTurnState::Busy {
VecDeque::from([supercode::FrontendEvent {
sequence: 4,
kind: "turn_succeeded".into(),
payload: json!({
"type": "turn_succeeded",
"reply": "continued through GLM"
}),
}])
} else {
VecDeque::new()
};
Ok(FrontendAttachment::from_snapshot(
FrontendAttachSnapshot {
descriptor: self.descriptor.clone(),
history: history[start..].to_vec(),
history_cursor: history.len() as u64,
replay,
},
self.events.subscribe(),
))
}
async fn send_input(self: Arc<Self>, _prompt: String) -> Result<(), SdkError> {
Ok(())
}
async fn submit(&self, prompt: String) -> Result<String, SdkError> {
self.submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(prompt.clone());
let reply = format!("continued:{prompt}");
self.history
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.extend([
ChatMessage::user(prompt),
ChatMessage::assistant(reply.clone()),
]);
Ok(reply)
}
async fn submit_with_images(
&self,
prompt: String,
image_urls: Vec<String>,
) -> Result<String, SdkError> {
self.submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(prompt.clone());
self.image_submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(image_urls.clone());
let reply = format!("continued:{prompt}");
self.history
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.extend([
ChatMessage::user_with_images(prompt, &image_urls),
ChatMessage::assistant(reply.clone()),
]);
Ok(reply)
}
async fn interrupt(&self) -> Result<bool, SdkError> {
self.interrupts.fetch_add(1, Ordering::SeqCst);
Ok(true)
}
async fn steer(&self, prompt: String) -> Result<(), SdkError> {
if !self.steer_accepting.load(Ordering::SeqCst) {
return Err(SdkError::UnsupportedAction("steer"));
}
self.steers
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(prompt.clone());
let mut history = self
.history
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
history.push(ChatMessage::assistant(format!("steered:{prompt}")));
let sequence = history.len() as u64 + 1;
drop(history);
let _ = self.events.send(supercode::FrontendEvent {
sequence,
kind: "turn_succeeded".into(),
payload: json!({"type": "turn_succeeded", "reply": format!("steered:{prompt}")}),
});
Ok(())
}
async fn respond(&self, response: supercode::FrontendResponse) -> Result<(), SdkError> {
self.responses
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(response);
Ok(())
}
}
fn json_body(response: OpenCodeResponse) -> Value {
assert_eq!(response.status, 200);
match response.body {
ResponseBody::Json(body) => body,
ResponseBody::EventStream(_) => panic!("expected JSON response"),
}
}
#[test]
fn namespace_and_traced_route_set_are_exactly_pinned() {
assert_eq!(PROTOCOL_NAMESPACE, "opencode_http/v1_2_15");
assert_eq!(OPENCODE_CLI_VERSION, "1.2.15");
assert_eq!(TRACED_ROUTES.len(), 23);
assert!(TRACED_ROUTES.contains(&"GET /event"));
assert!(TRACED_ROUTES.contains(&"POST /session/{session_id}/abort"));
assert!(TRACED_ROUTES.contains(&"POST /session/{session_id}/message"));
}
#[test]
fn corpus_allowlist_and_provenance_pin_drive_the_namespace_exactly() {
let corpus =
Path::new(env!("CARGO_MANIFEST_DIR")).join("../../scripts/client-protocol-corpus");
let allowlist: Value = serde_json::from_slice(
&std::fs::read(corpus.join("allowlist.json")).expect("read corpus allowlist"),
)
.expect("parse corpus allowlist");
let mut observed = allowlist["clients"][PROTOCOL_NAMESPACE]["client_to_server"]
.as_object()
.expect("OpenCode route allowlist")
.keys()
.map(String::as_str)
.collect::<Vec<_>>();
let mut implemented = TRACED_ROUTES.to_vec();
observed.sort_unstable();
implemented.sort_unstable();
assert_eq!(implemented, observed);
let pins: Value = serde_json::from_slice(
&std::fs::read(corpus.join("pins.json")).expect("read corpus pins"),
)
.expect("parse corpus pins");
let pin = pins["clients"]
.as_array()
.expect("client pins")
.iter()
.find(|pin| pin["id"] == "opencode")
.expect("OpenCode pin");
assert_eq!(pin["version"], OPENCODE_CLI_VERSION);
assert_eq!(pin["namespace"], PROTOCOL_NAMESPACE);
assert_eq!(pin["commit"], "799b2623cbb1c0f19e045d87c2c8593e83678bc0");
assert_eq!(pin["license"]["spdx"], "MIT");
assert_eq!(
pin["contract"]["sha256"],
"cfb4d87bc11924794a1fe9f3daafb8b569cc412bccda85b5014240bc3afe7eff"
);
}
#[tokio::test]
async fn sanitized_stock_exchange_replays_every_request_and_reconnect_boundary() {
let runtime = FixtureRuntime::new();
let adapter = OpenCodeAdapter::new(runtime, "runtime-1", "/workspace");
let fixture = Path::new(env!("CARGO_MANIFEST_DIR"))
.join("../../scripts/client-protocol-corpus/fixtures/opencode_v1_2_15.jsonl");
let text = std::fs::read_to_string(fixture).expect("read sanitized OpenCode fixture");
let mut requests = 0_usize;
let mut event_attachments = Vec::new();
let mut last_attachment = 0_u64;
for line in text.lines() {
let row: Value = serde_json::from_str(line).expect("parse fixture row");
if row["direction"] != "client_to_server" {
continue;
}
requests += 1;
let attachment = row["attachment"].as_u64().expect("attachment id");
assert!(attachment >= last_attachment, "reconnect order regressed");
last_attachment = attachment;
let message = &row["message"];
let method = message["method"].as_str().expect("method");
let original = message["path"].as_str().expect("path");
let target = if let Some(suffix) = original.strip_prefix("/session/") {
let suffix = suffix
.find('/')
.map(|index| &suffix[index..])
.unwrap_or_default();
format!("/session/{}{suffix}", adapter.session_id())
} else if original.starts_with("/permission/") {
"/permission/per_000000000000002a/reply".into()
} else {
original.to_string()
};
let body = message["body"].as_str().unwrap_or_default();
let mut body = if body.is_empty() {
Value::Null
} else {
serde_json::from_str(body).expect("fixture JSON body")
};
if method == "POST" && target.ends_with("/message") {
body["model"] = json!({"providerID": "openrouter", "modelID": "glm-5.2"});
}
let response = adapter
.handle(OpenCodeRequest::new(method, &target).with_body(body))
.await;
assert_eq!(response.status, 200, "fixture request {method} {target}");
if original == "/event" {
assert!(matches!(response.body, ResponseBody::EventStream(_)));
event_attachments.push(attachment);
}
}
assert_eq!(requests, 43);
assert_eq!(event_attachments, vec![1, 2]);
}
#[test]
fn sanitized_stock_server_exchange_pins_live_creation_and_reconnect_identity() {
let fixture = Path::new(env!("CARGO_MANIFEST_DIR"))
.join("../../scripts/client-protocol-corpus/fixtures/opencode_v1_2_15.jsonl");
let text = std::fs::read_to_string(fixture).expect("read sanitized OpenCode fixture");
let mut live_events = Vec::new();
let mut reconnect_history = Vec::new();
for line in text.lines() {
let row: Value = serde_json::from_str(line).expect("parse fixture row");
if row["direction"] != "server_to_client" || row["event"] != "body_chunk" {
continue;
}
let Some(message) = row["message"].as_str() else {
continue;
};
if row["attachment"] == 1 {
for frame in message.split("\n\n") {
if let Some(data) = frame.strip_prefix("data: ") {
live_events.push(serde_json::from_str::<Value>(data).expect("SSE JSON"));
}
}
} else if row["attachment"] == 2 {
if let Ok(history) = serde_json::from_str::<Vec<Value>>(message) {
if history.len() > reconnect_history.len()
&& history.iter().all(|item| item.get("info").is_some())
{
reconnect_history = history;
}
}
}
}
let mut messages = std::collections::BTreeSet::new();
let mut parts = std::collections::BTreeSet::new();
let mut live_message_ids = std::collections::BTreeSet::new();
let mut deltas = 0_usize;
for event in &live_events {
match event["type"].as_str() {
Some("message.updated") => {
let info = &event["properties"]["info"];
if let Some(id) = info["id"].as_str() {
messages.insert(id.to_string());
if matches!(info["role"].as_str(), Some("user" | "assistant")) {
live_message_ids.insert(id.to_string());
}
}
}
Some("message.part.updated") => {
let part = &event["properties"]["part"];
let message_id = part["messageID"].as_str().expect("part message id");
assert!(
messages.contains(message_id),
"stock corpus updated part before creating message {message_id}"
);
parts.insert(part["id"].as_str().expect("part id").to_string());
}
Some("message.part.delta") => {
deltas += 1;
let properties = &event["properties"];
assert!(messages.contains(properties["messageID"].as_str().unwrap()));
assert!(parts.contains(properties["partID"].as_str().unwrap()));
}
_ => {}
}
}
assert!(deltas > 0, "stock corpus must contain live text deltas");
assert!(
!reconnect_history.is_empty(),
"reconnect history was not captured"
);
let reconnect_ids = reconnect_history
.iter()
.filter_map(|message| message["info"]["id"].as_str())
.collect::<std::collections::BTreeSet<_>>();
assert!(
live_message_ids
.iter()
.all(|id| reconnect_ids.contains(id.as_str())),
"stock reconnect must preserve every live message identity"
);
}
#[test]
fn deterministic_ids_are_domain_separated_and_collision_free_for_large_sample() {
let mut ids = std::collections::BTreeSet::new();
for index in 0..50_000 {
let source = format!("runtime-{index}");
assert!(ids.insert(stable_id("session", "ses", &source, 24)));
assert!(ids.insert(stable_id("message", "msg", &source, 24)));
assert!(ids.insert(stable_id("part", "prt", &source, 24)));
}
assert_eq!(ids.len(), 150_000);
assert_eq!(
stable_id("session", "id", "same-source", 24),
stable_id("session", "id", "same-source", 24)
);
assert_ne!(
stable_id("session", "id", "same-source", 24),
stable_id("message", "id", "same-source", 24)
);
}
#[test]
fn message_identity_survives_a_bounded_history_window_slide() {
let mut state = HistoryIdentityState::default();
let first = vec![
ChatMessage::user("a"),
ChatMessage::assistant("b"),
ChatMessage::user("c"),
ChatMessage::assistant("d"),
];
assert_eq!(state.project(&first, 4), vec![0, 1, 2, 3]);
let slid = vec![
ChatMessage::user("c"),
ChatMessage::assistant("d"),
ChatMessage::user("e"),
ChatMessage::assistant("f"),
];
let projected = state.project(&slid, 5);
assert_eq!(projected, vec![2, 3, 4, 5]);
assert_eq!(state.project(&slid, 5), projected);
let retained_before = stable_id("message", "msg", "ses_test:2", 24);
let retained_after = stable_id("message", "msg", &format!("ses_test:{}", projected[0]), 24);
assert_eq!(retained_before, retained_after);
}
#[tokio::test]
async fn read_routes_project_one_runtime_without_provider_credentials() {
let adapter = OpenCodeAdapter::new(FixtureRuntime::new(), "runtime-1", "/workspace");
let session_id = adapter.session_id().to_string();
let list = json_body(
adapter
.handle(OpenCodeRequest::new("GET", "/session?start=0"))
.await,
);
assert_eq!(list[0]["id"], session_id);
assert_eq!(list[0]["directory"], "/workspace");
let history = json_body(
adapter
.handle(OpenCodeRequest::new(
"GET",
format!("/session/{session_id}/message?limit=100"),
))
.await,
);
assert_eq!(history.as_array().unwrap().len(), 3);
assert_eq!(
history[0]["parts"][0]["text"],
"[system]\npreserve system context"
);
assert_eq!(history[1]["info"]["role"], "user");
assert_eq!(history[2]["parts"][0]["text"], "continued through GLM");
let provider = json_body(
adapter
.handle(OpenCodeRequest::new("GET", "/provider"))
.await,
);
assert_eq!(provider["connected"], json!(["openrouter"]));
let serialized = provider.to_string();
assert!(!serialized.contains("apiKey"));
assert!(!serialized.contains("OPENROUTER_API_KEY"));
}
#[tokio::test]
async fn unknown_routes_and_wrong_session_ids_fail_closed() {
let adapter = OpenCodeAdapter::new(FixtureRuntime::new(), "runtime-1", "/workspace");
let unknown = adapter
.handle(OpenCodeRequest::new("DELETE", "/session/anything"))
.await;
assert_eq!(unknown.status, 404);
let wrong = adapter
.handle(OpenCodeRequest::new("GET", "/session/ses_wrong"))
.await;
assert_eq!(wrong.status, 404);
let untraced_query = adapter
.handle(OpenCodeRequest::new(
"GET",
format!("/session/{}/message?limit=999", adapter.session_id()),
))
.await;
assert_eq!(untraced_query.status, 404);
for target in [
"/agent?x=1",
"/config?x=1",
"/event?x=1",
"/session?unexpected=start=1",
"/session?start=1&extra=1",
"/session?start=not-a-number",
] {
assert_eq!(
adapter
.handle(OpenCodeRequest::new("GET", target))
.await
.status,
404,
"accepted untraced target {target}"
);
}
assert_eq!(
adapter
.handle(
OpenCodeRequest::new("POST", "/session").with_body(json!({"unexpected": true}))
)
.await
.status,
404
);
}
#[tokio::test]
async fn event_route_returns_atomic_sdk_attachment_not_a_second_runtime() {
let runtime = FixtureRuntime::new();
let adapter = OpenCodeAdapter::new(runtime, "runtime-1", "/workspace");
let response = adapter.handle(OpenCodeRequest::new("GET", "/event")).await;
assert_eq!(response.status, 200);
match response.body {
ResponseBody::EventStream(attachment) => {
assert_eq!(attachment.descriptor.session_id, "runtime-1");
assert_eq!(attachment.history.len(), 3);
}
ResponseBody::Json(_) => panic!("expected event stream"),
}
}
#[tokio::test]
async fn traced_abort_is_descriptor_gated_and_invokes_sdk_interrupt() {
let runtime = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
let response = adapter
.handle(OpenCodeRequest::new(
"POST",
format!("/session/{}/abort", adapter.session_id()),
))
.await;
assert_eq!(json_body(response), json!(true));
assert_eq!(runtime.interrupts.load(Ordering::SeqCst), 1);
let mut unavailable = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
Arc::get_mut(&mut unavailable)
.expect("unshared fixture")
.descriptor
.actions
.interrupt = false;
let adapter = OpenCodeAdapter::new(unavailable.clone(), "runtime-2", "/workspace");
let denied = adapter
.handle(OpenCodeRequest::new(
"POST",
format!("/session/{}/abort", adapter.session_id()),
))
.await;
assert_eq!(denied.status, 409);
assert_eq!(unavailable.interrupts.load(Ordering::SeqCst), 0);
}
#[test]
fn runtime_attachment_paths_are_cross_platform_and_workspace_bounded() {
assert_eq!(
logical_runtime_path("/srv/project", "assets/pixel.png").unwrap(),
"/srv/project/assets/pixel.png"
);
assert_eq!(
logical_runtime_path(
"C:\\runtime\\project",
"C:\\runtime\\project\\assets\\pixel.png"
)
.unwrap(),
"C:/runtime/project/assets/pixel.png"
);
assert!(logical_runtime_path("/srv/project", "/Users/client/private.png").is_err());
assert!(
logical_runtime_path("C:\\runtime\\project", "C:\\Users\\client\\private.png").is_err()
);
assert!(logical_runtime_path("/srv/project", "../../private.png").is_err());
}
#[test]
fn every_attachment_uri_branch_is_runtime_resolved_and_fail_closed() {
let workspace = std::env::temp_dir().join(format!(
"supercode-opencode-attachment-branches-{}",
std::process::id()
));
let outside = std::env::temp_dir().join(format!(
"supercode-opencode-attachment-outside-{}",
std::process::id()
));
std::fs::remove_dir_all(&workspace).ok();
std::fs::remove_file(&outside).ok();
std::fs::create_dir_all(&workspace).unwrap();
std::fs::write(workspace.join("note.txt"), "runtime note").unwrap();
std::fs::write(&outside, "outside").unwrap();
let resolve = |part: Value| {
resolve_file_part(part.as_object().expect("file-part object"), &workspace)
};
for (label, part, expected_text, expected_image) in [
(
"data image",
json!({"type":"file","mime":"image/png","filename":"pixel.png","url":"data:image/png;base64,cG5n"}),
None,
Some("data:image/png;base64,cG5n"),
),
(
"data text",
json!({"type":"file","mime":"text/plain","filename":"note.txt","url":"data:text/plain;base64,aGVsbG8="}),
Some("[file: note.txt]\nhello"),
None,
),
(
"http image",
json!({"type":"file","mime":"image/png","filename":"remote.png","url":"http://example.test/remote.png"}),
None,
Some("http://example.test/remote.png"),
),
(
"https image",
json!({"type":"file","mime":"image/webp","filename":"remote.webp","url":"https://example.test/remote.webp"}),
None,
Some("https://example.test/remote.webp"),
),
(
"runtime text",
json!({"type":"file","mime":"text/plain","filename":"note.txt","url":Url::from_file_path(workspace.join("note.txt")).unwrap().to_string()}),
Some("[file: note.txt]\nruntime note"),
None,
),
] {
let resolved = resolve(part).unwrap_or_else(|error| panic!("{label}: {error}"));
assert_eq!(resolved.text_block.as_deref(), expected_text, "{label}");
assert_eq!(resolved.image_url.as_deref(), expected_image, "{label}");
}
for (label, part, expected) in [
(
"mismatched data mime",
json!({"type":"file","mime":"image/png","url":"data:image/jpeg;base64,cG5n"}),
"data URI mime must match",
),
(
"invalid image base64",
json!({"type":"file","mime":"image/png","url":"data:image/png;base64,%%%"}),
"invalid base64",
),
(
"non-UTF8 data text",
json!({"type":"file","mime":"text/plain","url":"data:text/plain;base64,/w=="}),
"must be UTF-8 text",
),
(
"remote text",
json!({"type":"file","mime":"text/plain","url":"https://example.test/private.txt"}),
"remote non-image attachments are not fetched",
),
(
"unsupported scheme",
json!({"type":"file","mime":"image/png","url":"ftp://example.test/pixel.png"}),
"URL scheme is not supported",
),
(
"invalid URL",
json!({"type":"file","mime":"image/png","url":"not a URL"}),
"must be a data:, file:, or https: URL",
),
(
"outside runtime path",
json!({"type":"file","mime":"text/plain","url":Url::from_file_path(&outside).unwrap().to_string()}),
"escapes the SDK runtime workspace",
),
] {
let error = resolve(part)
.err()
.unwrap_or_else(|| panic!("{label} was accepted"));
assert!(error.to_string().contains(expected), "{label}: {error}");
}
#[cfg(unix)]
{
std::os::unix::fs::symlink(&outside, workspace.join("escaped-link.txt")).unwrap();
let error = resolve(json!({
"type":"file", "mime":"text/plain",
"url":Url::from_file_path(workspace.join("escaped-link.txt")).unwrap().to_string()
}))
.err()
.expect("symlink escape must be rejected");
assert!(error
.to_string()
.contains("escapes the SDK runtime workspace"));
}
std::fs::remove_dir_all(workspace).unwrap();
std::fs::remove_file(outside).unwrap();
}
#[tokio::test]
async fn file_parts_are_resolved_on_the_sdk_runtime_without_client_path_expansion() {
let workspace = std::env::temp_dir().join(format!(
"supercode-opencode-runtime-attachment-{}",
std::process::id()
));
let outside = std::env::temp_dir().join(format!(
"supercode-opencode-client-attachment-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&workspace);
std::fs::create_dir_all(&workspace).unwrap();
std::fs::write(workspace.join("pixel.png"), b"png").unwrap();
std::fs::write(workspace.join("note.txt"), "runtime note").unwrap();
std::fs::write(&outside, b"client secret").unwrap();
let runtime = FixtureRuntime::new();
let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", &workspace);
let session_id = adapter.session_id().to_string();
let file_url = Url::from_file_path(workspace.join("pixel.png"))
.unwrap()
.to_string();
let accepted = adapter
.handle(
OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
json!({
"messageID": "msg_attachment", "agent": "build",
"model": {"providerID": "openrouter", "modelID": "glm-5.2"},
"parts": [
{"id": "prt_text", "type": "text", "text": "inspect"},
{"id": "prt_more", "type": "text", "text": "carefully"},
{"id": "prt_note", "type": "file", "mime": "text/plain", "filename": "note.txt", "url": Url::from_file_path(workspace.join("note.txt")).unwrap().to_string()},
{"id": "prt_image", "type": "file", "mime": "image/png", "filename": "pixel.png", "url": file_url}
]
}),
),
)
.await;
assert_eq!(accepted.status, 200);
assert_eq!(
*runtime
.submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec!["inspect\ncarefully\n[file: note.txt]\nruntime note".to_string()]
);
assert_eq!(
*runtime
.image_submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec![vec!["data:image/png;base64,cG5n".to_string()]]
);
let client_url = Url::from_file_path(&outside).unwrap().to_string();
let denied = adapter
.handle(
OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
json!({
"messageID": "msg_client_path", "agent": "build",
"model": {"providerID": "openrouter", "modelID": "glm-5.2"},
"parts": [{"id": "prt_client", "type": "file", "mime": "text/plain", "url": client_url}]
}),
),
)
.await;
assert_eq!(denied.status, 400);
assert_eq!(
runtime
.image_submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.len(),
1
);
std::fs::remove_dir_all(workspace).unwrap();
std::fs::remove_file(outside).unwrap();
}
#[tokio::test]
async fn terminal_replay_reuses_the_assistant_already_present_in_history() {
let runtime = FixtureRuntime::new();
let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
let attachment = runtime.attach(HISTORY_LIMIT).await.unwrap();
let history = adapter.history_messages(&attachment);
let existing_id = history.as_array().unwrap().last().unwrap()["info"]["id"].clone();
let sequence = attachment.history_cursor + 1;
let mut projection = adapter.event_projection(&attachment);
let replay = projection.project(&SdkEvent {
sequence,
kind: "turn_succeeded".into(),
payload: json!({"type": "turn_succeeded", "reply": "continued through GLM"}),
});
let completed = replay
.iter()
.find(|event| event["type"] == "message.updated")
.expect("terminal replay completes the history message");
assert_eq!(completed["properties"]["info"]["id"], existing_id);
assert_eq!(
replay
.iter()
.filter(|event| event["type"] == "message.updated")
.count(),
1
);
}
#[tokio::test]
async fn simultaneous_event_attachments_share_live_message_identity() {
let runtime = FixtureRuntime::new();
let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
let first = runtime.attach(HISTORY_LIMIT).await.unwrap();
let second = runtime.attach(HISTORY_LIMIT).await.unwrap();
let mut first = adapter.event_projection(&first);
let mut second = adapter.event_projection(&second);
let user = SdkEvent {
sequence: 40,
kind: "user_message".into(),
payload: json!({"type": "user_message", "text": "same event"}),
};
let started = SdkEvent {
sequence: 41,
kind: "turn_started".into(),
payload: json!({"type": "turn_started"}),
};
let first_user = first.project(&user);
let second_user = second.project(&user);
assert_eq!(
first_user[0]["properties"]["info"]["id"],
second_user[0]["properties"]["info"]["id"]
);
let first_assistant = first.project(&started);
let second_assistant = second.project(&started);
assert_eq!(
first_assistant[1]["properties"]["info"]["id"],
second_assistant[1]["properties"]["info"]["id"]
);
assert_eq!(
first_assistant[2]["properties"]["part"]["id"],
second_assistant[2]["properties"]["part"]["id"]
);
}
#[tokio::test]
async fn traced_mutations_submit_and_answer_only_exact_typed_requests() {
let runtime = FixtureRuntime::new();
let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
let session_id = adapter.session_id();
let submitted = adapter
.handle(
OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
json!({
"messageID": "msg_stock", "agent": "build",
"model": {"providerID": "openrouter", "modelID": "glm-5.2"},
"parts": [{"id": "prt_stock", "type": "text", "text": "continue through GLM"}]
}),
),
)
.await;
assert_eq!(submitted.status, 200);
assert_eq!(
*runtime
.submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec!["continue through GLM"]
);
let history = json_body(
adapter
.handle(OpenCodeRequest::new(
"GET",
format!("/session/{session_id}/message?limit=100"),
))
.await,
);
let submitted_user = history
.as_array()
.unwrap()
.iter()
.find(|message| message["parts"][0]["text"] == "continue through GLM")
.expect("submitted user message in history");
assert_eq!(submitted_user["info"]["id"], "msg_stock");
assert_eq!(submitted_user["parts"][0]["id"], "prt_stock");
let answered = adapter
.handle(
OpenCodeRequest::new("POST", "/permission/per_000000000000002a/reply")
.with_body(json!({"reply": "once"})),
)
.await;
assert_eq!(answered.status, 200);
assert_eq!(
*runtime
.responses
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec![FrontendResponse::Approval {
request_id: 42,
decision: FrontendApprovalDecision::Allow
}]
);
let extra = adapter
.handle(
OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
json!({
"messageID": "msg", "agent": "build", "model": {}, "parts": [],
"untraced": true
}),
),
)
.await;
assert_eq!(extra.status, 400);
assert_eq!(
runtime
.submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.len(),
1
);
for target in [
format!("/session/{session_id}/message?untraced=1"),
"/permission/per_000000000000002a/reply?untraced=1".into(),
] {
assert_eq!(
adapter
.handle(
OpenCodeRequest::new("POST", target).with_body(json!({"reply": "once"})),
)
.await
.status,
404
);
}
}
#[tokio::test]
async fn busy_steer_ignores_the_previous_turn_terminal_replay() {
let runtime = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
let response = adapter
.handle(
OpenCodeRequest::new("POST", format!("/session/{}/message", adapter.session_id()))
.with_body(json!({
"messageID": "msg_stock", "agent": "build",
"model": {"providerID": "openrouter", "modelID": "glm-5.2"},
"parts": [{"id": "prt_stock", "type": "text", "text": "change direction"}]
})),
)
.await;
assert_eq!(response.status, 200);
let response = match response.body {
ResponseBody::Json(body) => body,
ResponseBody::EventStream(_) => panic!("expected steer response"),
};
assert_eq!(response["parts"][0]["text"], "steered:change direction");
assert_ne!(response["parts"][0]["text"], "continued through GLM");
assert!(runtime
.submissions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_empty());
assert_eq!(
*runtime
.steers
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
vec!["change direction"]
);
}
#[tokio::test]
async fn busy_descriptor_race_propagates_atomic_steer_rejection() {
let runtime = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
runtime.steer_accepting.store(false, Ordering::SeqCst);
let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
let response = adapter
.handle(
OpenCodeRequest::new("POST", format!("/session/{}/message", adapter.session_id()))
.with_body(json!({
"messageID": "msg_stock", "agent": "build",
"model": {"providerID": "openrouter", "modelID": "glm-5.2"},
"parts": [{"id": "prt_stock", "type": "text", "text": "too late"}]
})),
)
.await;
assert_eq!(response.status, 409);
assert!(runtime
.steers
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_empty());
}
#[test]
fn delayed_sse_reuses_identities_already_projected_by_history() {
let descriptor = FixtureRuntime::new().descriptor.clone();
let identities = Arc::new(Mutex::new(HistoryIdentityState::default()));
let submitted = SubmittedMessage {
prompt: "continue".into(),
image_urls: Vec::new(),
message_id: "msg_history_first".into(),
part_id: "prt_history_first".into(),
};
let baseline = vec![
ChatMessage::user("continue"),
ChatMessage::assistant("prior identical prompt"),
];
let history = vec![
baseline[0].clone(),
baseline[1].clone(),
ChatMessage::user("continue"),
ChatMessage::assistant("tool request"),
ChatMessage::tool_result("call-1", "read", "tool result"),
ChatMessage::assistant("history won the race"),
];
let (history_user, history_assistant) = {
let mut state = identities
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.project(&baseline, 4);
state.register_client_user(&submitted);
let ordinals = state.project(&history, 10);
(
state.identity("ses_test", ordinals[2]),
state.identity("ses_test", ordinals[5]),
)
};
let mut future = OpenCodeEventProjection::new(
"ses_test".into(),
PathBuf::from("/workspace"),
&descriptor,
identities.clone(),
None,
);
let future_user = future.project(&SdkEvent {
sequence: 11,
kind: "user_message".into(),
payload: json!({"type": "user_message", "text": "continue"}),
});
let future_assistant = future.project(&SdkEvent {
sequence: 12,
kind: "turn_started".into(),
payload: json!({"type": "turn_started"}),
});
assert_ne!(future_user[0]["properties"]["info"]["id"], history_user.0);
assert_ne!(
future_assistant[1]["properties"]["info"]["id"],
history_assistant.0
);
let mut delayed = OpenCodeEventProjection::new(
"ses_test".into(),
PathBuf::from("/workspace"),
&descriptor,
identities,
None,
);
let user = delayed.project(&SdkEvent {
sequence: 6,
kind: "user_message".into(),
payload: json!({"type": "user_message", "text": "continue"}),
});
let assistant = delayed.project(&SdkEvent {
sequence: 7,
kind: "turn_started".into(),
payload: json!({"type": "turn_started"}),
});
assert_eq!(user[0]["properties"]["info"]["id"], history_user.0);
assert_eq!(user[1]["properties"]["part"]["id"], history_user.1);
assert_eq!(
assistant[1]["properties"]["info"]["id"],
history_assistant.0
);
assert_eq!(
assistant[2]["properties"]["part"]["id"],
history_assistant.1
);
}
#[test]
fn delayed_sse_reuses_history_ids_for_another_frontend_or_scheduler_turn() {
let descriptor = FixtureRuntime::new().descriptor.clone();
let identities = Arc::new(Mutex::new(HistoryIdentityState::default()));
let baseline = vec![
ChatMessage::user("existing"),
ChatMessage::assistant("existing reply"),
];
let current = vec![
baseline[0].clone(),
baseline[1].clone(),
ChatMessage::user("scheduled wakeup"),
ChatMessage::assistant("scheduler reply"),
];
let (history_user, history_assistant) = {
let mut state = identities
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.project(&baseline, 4);
let ordinals = state.project(¤t, 8);
(
state.identity("ses_test", ordinals[2]),
state.identity("ses_test", ordinals[3]),
)
};
let mut delayed = OpenCodeEventProjection::new(
"ses_test".into(),
PathBuf::from("/workspace"),
&descriptor,
identities,
None,
);
let user = delayed.project(&SdkEvent {
sequence: 5,
kind: "user_message".into(),
payload: json!({"type": "user_message", "text": "scheduled wakeup"}),
});
let assistant = delayed.project(&SdkEvent {
sequence: 6,
kind: "turn_started".into(),
payload: json!({"type": "turn_started"}),
});
assert_eq!(user[0]["properties"]["info"]["id"], history_user.0);
assert_eq!(user[1]["properties"]["part"]["id"], history_user.1);
assert_eq!(
assistant[1]["properties"]["info"]["id"],
history_assistant.0
);
assert_eq!(
assistant[2]["properties"]["part"]["id"],
history_assistant.1
);
}
#[test]
fn older_attachment_projection_cannot_rewind_newer_identity_state() {
let mut state = HistoryIdentityState::default();
let old = vec![
ChatMessage::user("old"),
ChatMessage::assistant("old reply"),
];
let newer = vec![
old[0].clone(),
old[1].clone(),
ChatMessage::user("new"),
ChatMessage::assistant("new reply"),
];
let newer_ids = state.project(&newer, 8);
assert_eq!(state.project(&old, 4), newer_ids[..2]);
let newest = vec![
newer[0].clone(),
newer[1].clone(),
newer[2].clone(),
newer[3].clone(),
ChatMessage::user("newest"),
ChatMessage::assistant("newest reply"),
];
let newest_ids = state.project(&newest, 12);
assert_eq!(&newest_ids[..4], &newer_ids);
assert_eq!(state.project(&newer, 8), newer_ids);
}
#[test]
fn canonical_events_project_to_stock_delta_and_reversible_permission_ids() {
let descriptor = FixtureRuntime::new().descriptor.clone();
let identities = Arc::new(Mutex::new(HistoryIdentityState::default()));
identities
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.register_client_user(&SubmittedMessage {
prompt: "continue".into(),
image_urls: Vec::new(),
message_id: "msg_stock_client".into(),
part_id: "prt_stock_client".into(),
});
let mut projection = OpenCodeEventProjection::new(
"ses_test".into(),
PathBuf::from("C:\\runtime\\worktree"),
&descriptor,
identities.clone(),
None,
);
let user = projection.project(&SdkEvent {
sequence: 6,
kind: "user_message".into(),
payload: json!({"type": "user_message", "text": "continue"}),
});
assert_eq!(user[0]["properties"]["info"]["id"], "msg_stock_client");
assert_eq!(user[1]["properties"]["part"]["id"], "prt_stock_client");
let started = projection.project(&SdkEvent {
sequence: 7,
kind: "turn_started".into(),
payload: json!({"type": "turn_started"}),
});
assert_eq!(started[0]["type"], "session.status");
assert_eq!(started[1]["type"], "message.updated");
assert_eq!(started[2]["type"], "message.part.updated");
let delta = projection.project(&SdkEvent {
sequence: 8,
kind: "text_delta".into(),
payload: json!({"type": "text_delta", "text": "hello"}),
});
assert_eq!(delta[0]["type"], "message.part.delta");
assert_eq!(delta[0]["properties"]["delta"], "hello");
assert_eq!(
delta[0]["properties"]["messageID"],
started[1]["properties"]["info"]["id"]
);
assert_eq!(
delta[0]["properties"]["partID"],
started[2]["properties"]["part"]["id"]
);
projection.project(&SdkEvent {
sequence: 9,
kind: "turn_succeeded".into(),
payload: json!({"type": "turn_succeeded", "reply": "hello"}),
});
let history = vec![
ChatMessage::user("continue"),
ChatMessage::assistant("hello"),
];
let mut identity_state = identities
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let ordinals = identity_state.project(&history, 9);
let projected_identities = ordinals
.iter()
.map(|ordinal| identity_state.identity("ses_test", *ordinal))
.collect::<Vec<_>>();
drop(identity_state);
let reconnected = history_messages(
&history,
&ordinals,
&projected_identities,
"ses_test",
Path::new("C:\\runtime\\worktree"),
&descriptor,
);
assert_eq!(
reconnected[0]["info"]["id"],
user[0]["properties"]["info"]["id"]
);
assert_eq!(
reconnected[1]["info"]["id"],
started[1]["properties"]["info"]["id"]
);
assert_eq!(
reconnected[1]["parts"][0]["id"],
started[2]["properties"]["part"]["id"]
);
assert_eq!(reconnected[1]["info"]["providerID"], "openrouter");
assert_eq!(reconnected[1]["info"]["modelID"], "glm-5.2");
let permission = projection.project(&SdkEvent {
sequence: 10,
kind: "request".into(),
payload: json!({"type": "request", "request": {
"id": 42, "kind": "approval",
"payload": {"tool": "edit", "subject": "C:\\runtime\\worktree\\probe.txt"}
}}),
});
assert_eq!(permission[0]["properties"]["id"], "per_000000000000002a");
assert_eq!(
permission_request_id("/permission/per_000000000000002a/reply"),
Some(42)
);
assert_eq!(
path_text(Path::new("C:\\runtime\\worktree")),
"C:/runtime/worktree"
);
}
#[test]
fn client_projection_types_never_enter_harness_source() {
let harness = std::fs::read_to_string(
Path::new(env!("CARGO_MANIFEST_DIR")).join("../harness/src/frontend.rs"),
)
.unwrap();
for forbidden in ["OpenCodeAdapter", "opencode_http", "OPENCODE_CLI_VERSION"] {
assert!(!harness.contains(forbidden), "harness contains {forbidden}");
}
}
}