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_harness::{
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_harness::SdkErrorCode::ControllerRequired =>
{
(409, "controller_required")
}
Self::Sdk(error) if error.code() == supercode_harness::SdkErrorCode::LeaseExpired => {
(409, "lease_expired")
}
Self::Sdk(error) if error.code() == supercode_harness::SdkErrorCode::Busy => {
(409, "busy")
}
Self::Sdk(error)
if error.code() == supercode_harness::SdkErrorCode::UnsupportedAction =>
{
(409, "action_unavailable")
}
Self::Sdk(error) if error.code() == supercode_harness::SdkErrorCode::Unauthorized => {
(403, "unauthorized")
}
Self::Sdk(error)
if error.code() == supercode_harness::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_harness::FrontendTurnState::Idle => "idle",
supercode_harness::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_harness::FrontendTurnState::Idle => json!({"type": "idle"}),
supercode_harness::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_harness::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_harness::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_harness::FrontendTurnState::Idle => {
return Err(AdapterError::Sdk(SdkError::UnsupportedAction("submit")))
}
supercode_harness::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!("Volter Harness 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 Volter Harness 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": "Volter Harness 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>()
}