1use std::collections::BTreeMap;
4use std::path::{Path, PathBuf};
5use std::sync::{Arc, Mutex};
6
7use base64::Engine as _;
8use serde_json::{json, Value};
9use supercode::{
10 ChatMessage, FrontendApprovalDecision, FrontendAttachment, FrontendResponse,
11 FrontendRuntimeDescriptor, Role, SdkError, SdkEvent, SdkRuntime,
12};
13use url::Url;
14
15pub const PROTOCOL_NAMESPACE: &str = "opencode_http/v1_2_15";
17pub const OPENCODE_CLI_VERSION: &str = "1.2.15";
19pub const HISTORY_LIMIT: usize = 4096;
21
22const TRACED_ROUTES: &[&str] = &[
23 "GET /agent",
24 "GET /command",
25 "GET /config",
26 "GET /config/providers",
27 "GET /event",
28 "GET /experimental/resource",
29 "GET /formatter",
30 "GET /lsp",
31 "GET /mcp",
32 "GET /path",
33 "GET /provider",
34 "GET /provider/auth",
35 "GET /session/status",
36 "GET /session/{session_id}",
37 "GET /session/{session_id}/diff",
38 "GET /session/{session_id}/message?limit=100",
39 "GET /session/{session_id}/todo",
40 "GET /session?start={cursor}",
41 "GET /vcs",
42 "POST /permission/{permission_id}/reply",
43 "POST /session",
44 "POST /session/{session_id}/abort",
45 "POST /session/{session_id}/message",
46];
47
48#[derive(Debug, Clone, PartialEq)]
50pub struct OpenCodeRequest {
51 pub method: String,
52 pub target: String,
53 pub body: Value,
54}
55
56impl OpenCodeRequest {
57 pub fn new(method: impl Into<String>, target: impl Into<String>) -> Self {
58 Self {
59 method: method.into(),
60 target: target.into(),
61 body: Value::Null,
62 }
63 }
64
65 pub fn with_body(mut self, body: Value) -> Self {
66 self.body = body;
67 self
68 }
69}
70
71pub enum ResponseBody {
74 Json(Value),
75 EventStream(Box<FrontendAttachment>),
76}
77
78pub struct OpenCodeResponse {
80 pub status: u16,
81 pub body: ResponseBody,
82}
83
84impl OpenCodeResponse {
85 fn json(status: u16, body: Value) -> Self {
86 Self {
87 status,
88 body: ResponseBody::Json(body),
89 }
90 }
91
92 fn events(attachment: FrontendAttachment) -> Self {
93 Self {
94 status: 200,
95 body: ResponseBody::EventStream(Box::new(attachment)),
96 }
97 }
98}
99
100#[derive(Debug, thiserror::Error)]
103pub enum AdapterError {
104 #[error("route not supported by this adapter version: {0}")]
105 UnsupportedRoute(String),
106 #[error("invalid request for `{route}`: {message}")]
107 InvalidRequest { route: String, message: String },
108 #[error("OpenCode session `{0}` is not attached to this runtime")]
109 UnknownSession(String),
110 #[error(transparent)]
111 Sdk(#[from] SdkError),
112}
113
114impl AdapterError {
115 pub fn response(&self) -> OpenCodeResponse {
116 let (status, name) = match self {
117 Self::UnsupportedRoute(_) => (404, "unsupported_route"),
118 Self::InvalidRequest { .. } | Self::UnknownSession(_) => (400, "invalid_request"),
119 Self::Sdk(error) if error.code() == supercode::SdkErrorCode::ControllerRequired => {
120 (409, "controller_required")
121 }
122 Self::Sdk(error) if error.code() == supercode::SdkErrorCode::LeaseExpired => {
123 (409, "lease_expired")
124 }
125 Self::Sdk(error) if error.code() == supercode::SdkErrorCode::Busy => (409, "busy"),
126 Self::Sdk(error) if error.code() == supercode::SdkErrorCode::UnsupportedAction => {
127 (409, "action_unavailable")
128 }
129 Self::Sdk(error) if error.code() == supercode::SdkErrorCode::Unauthorized => {
130 (403, "unauthorized")
131 }
132 Self::Sdk(error) if error.code() == supercode::SdkErrorCode::Unauthenticated => {
133 (401, "unauthenticated")
134 }
135 Self::Sdk(_) => (500, "sdk_error"),
136 };
137 OpenCodeResponse::json(status, json!({"name": name, "message": self.to_string()}))
138 }
139}
140
141pub struct OpenCodeAdapter {
143 runtime: Arc<dyn SdkRuntime>,
144 runtime_id: String,
145 workspace: PathBuf,
146 session_id: String,
147 history_ids: Arc<Mutex<HistoryIdentityState>>,
148}
149
150impl OpenCodeAdapter {
151 pub fn new(
152 runtime: Arc<dyn SdkRuntime>,
153 runtime_id: impl Into<String>,
154 workspace: impl Into<PathBuf>,
155 ) -> Arc<Self> {
156 let runtime_id = runtime_id.into();
157 let session_id = stable_id("session", "ses", &runtime_id, 24);
158 Arc::new(Self {
159 runtime,
160 runtime_id,
161 workspace: workspace.into(),
162 session_id,
163 history_ids: Arc::new(Mutex::new(HistoryIdentityState::default())),
164 })
165 }
166
167 pub fn session_id(&self) -> &str {
168 &self.session_id
169 }
170
171 pub fn traced_routes() -> &'static [&'static str] {
172 TRACED_ROUTES
173 }
174
175 pub fn event_projection(&self, attachment: &FrontendAttachment) -> OpenCodeEventProjection {
176 let ordinals = self
177 .history_ids
178 .lock()
179 .unwrap_or_else(std::sync::PoisonError::into_inner)
180 .project(&attachment.history, attachment.history_cursor);
181 let history_tail = attachment
182 .history
183 .iter()
184 .zip(ordinals)
185 .rev()
186 .find(|(message, _)| message.role == Role::Assistant)
187 .map(|(message, ordinal)| (ordinal, message_text(message)));
188 OpenCodeEventProjection::new(
189 self.session_id.clone(),
190 self.workspace.clone(),
191 &attachment.descriptor,
192 self.history_ids.clone(),
193 history_tail,
194 )
195 }
196
197 pub fn initial_events(&self, attachment: &FrontendAttachment) -> Vec<Value> {
198 vec![
199 json!({"type": "session.updated", "properties": {"info": self.session_info(attachment)}}),
200 json!({
201 "type": "session.status",
202 "properties": {
203 "sessionID": self.session_id,
204 "status": {"type": match attachment.descriptor.turn_state {
205 supercode::FrontendTurnState::Idle => "idle",
206 supercode::FrontendTurnState::Busy => "busy",
207 }}
208 }
209 }),
210 ]
211 }
212
213 pub async fn handle(&self, request: OpenCodeRequest) -> OpenCodeResponse {
216 match self.handle_result(request).await {
217 Ok(response) => response,
218 Err(error) => error.response(),
219 }
220 }
221
222 async fn handle_result(
223 &self,
224 request: OpenCodeRequest,
225 ) -> Result<OpenCodeResponse, AdapterError> {
226 let method = request.method.to_ascii_uppercase();
227 let target = request.target.as_str();
228 let path = target.split('?').next().unwrap_or(target);
229 match (method.as_str(), target) {
230 ("GET", "/agent") => Ok(OpenCodeResponse::json(200, self.agents())),
231 ("GET", "/command") => {
232 let descriptor = self.runtime.describe().await?;
233 Ok(OpenCodeResponse::json(200, commands(&descriptor)))
234 }
235 ("GET", "/config") => {
236 let descriptor = self.runtime.describe().await?;
237 Ok(OpenCodeResponse::json(200, self.config(&descriptor)))
238 }
239 ("GET", "/config/providers") => {
240 let descriptor = self.runtime.describe().await?;
241 Ok(OpenCodeResponse::json(
242 200,
243 provider_projection(&descriptor, false),
244 ))
245 }
246 ("GET", "/provider") => {
247 let descriptor = self.runtime.describe().await?;
248 Ok(OpenCodeResponse::json(
249 200,
250 provider_projection(&descriptor, true),
251 ))
252 }
253 ("GET", "/provider/auth") => Ok(OpenCodeResponse::json(200, json!({}))),
254 ("GET", "/experimental/resource" | "/mcp" | "/vcs") => {
255 Ok(OpenCodeResponse::json(200, json!({})))
256 }
257 ("GET", "/formatter" | "/lsp") => Ok(OpenCodeResponse::json(200, json!([]))),
258 ("GET", "/path") => Ok(OpenCodeResponse::json(200, self.paths())),
259 ("GET", "/session/status") => {
260 let descriptor = self.runtime.describe().await?;
261 let state = match descriptor.turn_state {
262 supercode::FrontendTurnState::Idle => json!({"type": "idle"}),
263 supercode::FrontendTurnState::Busy => json!({"type": "busy"}),
264 };
265 Ok(OpenCodeResponse::json(
266 200,
267 json!({self.session_id.clone(): state}),
268 ))
269 }
270 ("GET", "/event") => {
271 let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
272 Ok(OpenCodeResponse::events(attachment))
273 }
274 ("GET", _) if session_start_cursor(target).is_some() => {
275 let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
276 Ok(OpenCodeResponse::json(
277 200,
278 json!([self.session_info(&attachment)]),
279 ))
280 }
281 ("POST", "/session") if request.body.is_null() => {
282 let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
283 Ok(OpenCodeResponse::json(200, self.session_info(&attachment)))
284 }
285 ("POST", _) if target == path && self.abort_target(path) => self.interrupt().await,
286 ("POST", _) if target == path && self.message_target(path) => {
287 self.submit_message(&request.body).await
288 }
289 ("POST", _) if target == path && permission_request_id(path).is_some() => {
290 self.respond_permission(path, &request.body).await
291 }
292 ("GET", _) if self.session_tail(path).is_some() => {
293 self.handle_session_get(path, target).await
294 }
295 _ => Err(AdapterError::UnsupportedRoute(format!(
296 "{method} {}",
297 request.target
298 ))),
299 }
300 }
301
302 async fn handle_session_get(
303 &self,
304 path: &str,
305 target: &str,
306 ) -> Result<OpenCodeResponse, AdapterError> {
307 let tail = self
308 .session_tail(path)
309 .expect("caller checked session prefix");
310 let traced_target = match tail {
311 "/message" => format!("{path}?limit=100"),
312 _ => path.to_string(),
313 };
314 if target != traced_target {
315 return Err(AdapterError::UnsupportedRoute(format!("GET {target}")));
316 }
317 let attachment = self.runtime.attach(HISTORY_LIMIT).await?;
318 match tail {
319 "" => Ok(OpenCodeResponse::json(200, self.session_info(&attachment))),
320 "/message" => Ok(OpenCodeResponse::json(
321 200,
322 self.history_messages(&attachment),
323 )),
324 "/todo" | "/diff" => Ok(OpenCodeResponse::json(200, json!([]))),
325 _ => Err(AdapterError::UnsupportedRoute(format!("GET {path}"))),
326 }
327 }
328
329 fn session_tail<'a>(&self, path: &'a str) -> Option<&'a str> {
330 let prefix = "/session/";
331 let rest = path.strip_prefix(prefix)?;
332 let (id, tail) = rest
333 .split_once('/')
334 .map_or((rest, ""), |(id, tail)| (id, tail));
335 if id != self.session_id {
336 return None;
337 }
338 if tail.is_empty() {
339 Some("")
340 } else {
341 Some(
342 path.strip_prefix(&format!("{prefix}{id}"))
343 .expect("prefix matches"),
344 )
345 }
346 }
347
348 fn message_target(&self, path: &str) -> bool {
349 self.session_tail(path) == Some("/message")
350 }
351
352 fn abort_target(&self, path: &str) -> bool {
353 self.session_tail(path) == Some("/abort")
354 }
355
356 async fn interrupt(&self) -> Result<OpenCodeResponse, AdapterError> {
357 let descriptor = self.runtime.describe().await?;
358 if !descriptor.actions.interrupt {
359 return Err(AdapterError::Sdk(SdkError::UnsupportedAction("interrupt")));
360 }
361 let interrupted = self.runtime.interrupt().await?;
362 Ok(OpenCodeResponse::json(200, json!(interrupted)))
363 }
364
365 async fn submit_message(&self, body: &Value) -> Result<OpenCodeResponse, AdapterError> {
366 let descriptor = self.runtime.describe().await?;
367 validate_requested_model(body, &descriptor)?;
368 let submitted = prompt_from_message(body, &self.workspace)?;
369 let mut before = self.runtime.attach(HISTORY_LIMIT).await?;
374 self.history_ids
375 .lock()
376 .unwrap_or_else(std::sync::PoisonError::into_inner)
377 .project(&before.history, before.history_cursor);
378 let attachment = match descriptor.turn_state {
379 supercode::FrontendTurnState::Idle if descriptor.actions.submit => {
380 self.history_ids
381 .lock()
382 .unwrap_or_else(std::sync::PoisonError::into_inner)
383 .register_client_user(&submitted);
384 if submitted.image_urls.is_empty() {
385 self.runtime.submit(submitted.prompt).await?;
386 } else {
387 self.runtime
388 .submit_with_images(submitted.prompt, submitted.image_urls)
389 .await?;
390 }
391 self.runtime.attach(HISTORY_LIMIT).await?
392 }
393 supercode::FrontendTurnState::Busy if descriptor.actions.steer => {
394 let baseline_cursor = before.history_cursor;
395 self.runtime.steer(submitted.prompt).await?;
396 wait_for_completed_terminal(self.runtime.as_ref(), &mut before, baseline_cursor)
397 .await?
398 }
399 supercode::FrontendTurnState::Idle => {
400 return Err(AdapterError::Sdk(SdkError::UnsupportedAction("submit")))
401 }
402 supercode::FrontendTurnState::Busy => {
403 return Err(AdapterError::Sdk(SdkError::UnsupportedAction("steer")))
404 }
405 };
406 let messages = self.history_messages(&attachment);
407 let added = appended_message_count(&before.history, &attachment.history);
408 let assistant = messages
409 .as_array()
410 .and_then(|messages| {
411 messages
412 .get(messages.len().saturating_sub(added)..)
413 .unwrap_or_default()
414 .iter()
415 .rev()
416 .find(|message| message["info"]["role"] == "assistant")
417 })
418 .cloned()
419 .ok_or_else(|| AdapterError::InvalidRequest {
420 route: "POST /session/{session_id}/message".into(),
421 message: "runtime completed without an assistant history message".into(),
422 })?;
423 Ok(OpenCodeResponse::json(200, assistant))
424 }
425
426 async fn respond_permission(
427 &self,
428 path: &str,
429 body: &Value,
430 ) -> Result<OpenCodeResponse, AdapterError> {
431 let request_id = permission_request_id(path).expect("route guard parsed request id");
432 let descriptor = self.runtime.describe().await?;
433 if !descriptor.actions.respond {
434 return Err(AdapterError::Sdk(SdkError::UnsupportedAction("respond")));
435 }
436 let reply = exact_string_field(body, &["reply"], "reply")?;
437 let decision = match reply {
438 "once" => FrontendApprovalDecision::Allow,
439 "always" => FrontendApprovalDecision::AllowForSession,
440 "reject" => FrontendApprovalDecision::Deny,
441 _ => {
442 return Err(AdapterError::InvalidRequest {
443 route: "POST /permission/{permission_id}/reply".into(),
444 message: "reply must be `once`, `always`, or `reject`".into(),
445 })
446 }
447 };
448 self.runtime
449 .respond(FrontendResponse::Approval {
450 request_id,
451 decision,
452 })
453 .await?;
454 Ok(OpenCodeResponse::json(200, json!(true)))
455 }
456
457 fn session_info(&self, attachment: &FrontendAttachment) -> Value {
458 let title = attachment
459 .history
460 .iter()
461 .rev()
462 .find_map(|message| message.content.as_deref())
463 .map(|text| truncate(text, 80))
464 .filter(|text| !text.is_empty())
465 .unwrap_or_else(|| format!("Supercode runtime {}", self.runtime_id));
466 json!({
467 "id": self.session_id,
468 "slug": format!("supercode-{}", &self.session_id[4..12]),
469 "projectID": "global",
470 "directory": path_text(&self.workspace),
471 "title": title,
472 "version": OPENCODE_CLI_VERSION,
473 "summary": {"additions": 0, "deletions": 0, "files": 0},
474 "time": {"created": 0, "updated": attachment.history_cursor},
475 })
476 }
477
478 fn agents(&self) -> Value {
479 json!([{
480 "name": "build",
481 "description": "Attached Supercode SDK runtime",
482 "options": {},
483 "permission": [],
484 "mode": "primary",
485 "native": false,
486 }])
487 }
488
489 fn config(&self, descriptor: &FrontendRuntimeDescriptor) -> Value {
490 let (provider, model) = provider_model(descriptor);
491 json!({
492 "$schema": "https://opencode.ai/config.json",
493 "share": "disabled",
494 "autoupdate": false,
495 "enabled_providers": [provider],
496 "model": format!("{provider}/{model}"),
497 "small_model": format!("{provider}/{model}"),
498 "provider": {},
499 "permission": {"edit": "ask"},
500 "agent": {},
501 "mode": {},
502 "plugin": [],
503 "command": {},
504 "username": "supercode",
505 })
506 }
507
508 fn paths(&self) -> Value {
509 let workspace = path_text(&self.workspace);
510 json!({
511 "home": workspace,
512 "state": workspace,
513 "config": workspace,
514 "worktree": workspace,
515 "directory": workspace,
516 })
517 }
518
519 fn history_messages(&self, attachment: &FrontendAttachment) -> Value {
520 let mut state = self
521 .history_ids
522 .lock()
523 .unwrap_or_else(std::sync::PoisonError::into_inner);
524 let ordinals = state.project(&attachment.history, attachment.history_cursor);
525 let identities = ordinals
526 .iter()
527 .map(|ordinal| state.identity(&self.session_id, *ordinal))
528 .collect::<Vec<_>>();
529 drop(state);
530 history_messages(
531 &attachment.history,
532 &ordinals,
533 &identities,
534 &self.session_id,
535 &self.workspace,
536 &attachment.descriptor,
537 )
538 }
539}
540
541async fn wait_for_completed_terminal(
542 runtime: &dyn SdkRuntime,
543 attachment: &mut FrontendAttachment,
544 baseline_cursor: u64,
545) -> Result<FrontendAttachment, AdapterError> {
546 loop {
547 let event = attachment
548 .next_event()
549 .await
550 .map_err(|error| AdapterError::Sdk(SdkError::Transport(error.to_string())))?;
551 if matches!(
552 event.kind.as_str(),
553 "turn_succeeded" | "turn_failed" | "turn_interrupted"
554 ) {
555 let completed = runtime.attach(HISTORY_LIMIT).await?;
560 if completed.history_cursor > baseline_cursor
561 && event.sequence > completed.history_cursor
562 {
563 return Ok(completed);
564 }
565 }
566 }
567}
568
569fn appended_message_count(previous: &[ChatMessage], current: &[ChatMessage]) -> usize {
570 let overlap = (0..=previous.len().min(current.len()))
571 .rev()
572 .find(|count| previous[previous.len() - *count..] == current[..*count])
573 .unwrap_or(0);
574 current.len().saturating_sub(overlap)
575}
576
577fn session_start_cursor(target: &str) -> Option<u64> {
578 let cursor = target.strip_prefix("/session?start=")?;
579 if cursor.is_empty() || !cursor.bytes().all(|byte| byte.is_ascii_digit()) {
580 return None;
581 }
582 cursor.parse().ok()
583}
584
585fn permission_request_id(path: &str) -> Option<u64> {
586 let encoded = path
587 .strip_prefix("/permission/per_")?
588 .strip_suffix("/reply")?;
589 (encoded.len() == 16)
590 .then(|| u64::from_str_radix(encoded, 16).ok())
591 .flatten()
592}
593
594fn permission_id(request_id: u64) -> String {
595 format!("per_{request_id:016x}")
596}
597
598struct SubmittedMessage {
599 prompt: String,
600 image_urls: Vec<String>,
601 message_id: String,
602 part_id: String,
603}
604
605fn prompt_from_message(body: &Value, workspace: &Path) -> Result<SubmittedMessage, AdapterError> {
606 let object = exact_object_fields(body, &["messageID", "agent", "model", "parts"])?;
607 for required in ["messageID", "agent", "model", "parts"] {
608 if !object.contains_key(required) {
609 return Err(invalid_message(format!("missing `{required}`")));
610 }
611 }
612 if object["agent"] != "build" {
613 return Err(invalid_message("agent must be `build`"));
614 }
615 let message_id = object["messageID"]
616 .as_str()
617 .filter(|id| !id.is_empty())
618 .ok_or_else(|| invalid_message("messageID must be a nonempty string"))?;
619 if !message_id.starts_with("msg_") {
620 return Err(invalid_message(
621 "messageID must use the traced `msg_` prefix",
622 ));
623 }
624 let model = exact_object_fields(&object["model"], &["providerID", "modelID"])?;
625 for field in ["providerID", "modelID"] {
626 if model
627 .get(field)
628 .and_then(Value::as_str)
629 .is_none_or(str::is_empty)
630 {
631 return Err(invalid_message(format!(
632 "model.{field} must be a nonempty string"
633 )));
634 }
635 }
636 let parts = object["parts"]
637 .as_array()
638 .ok_or_else(|| invalid_message("parts must be an array"))?;
639 if parts.is_empty() {
640 return Err(invalid_message("message requires at least one part"));
641 }
642 let mut text = Vec::new();
643 let mut part_id = None;
644 let mut image_urls = Vec::new();
645 for part in parts {
646 match part.get("type").and_then(Value::as_str) {
647 Some("text") => {
648 let part = exact_object_fields(part, &["id", "type", "text"])?;
649 let value = part
650 .get("text")
651 .and_then(Value::as_str)
652 .ok_or_else(|| invalid_message("text part requires string text"))?;
653 let id = traced_part_id(part)?;
654 part_id.get_or_insert_with(|| id.to_string());
655 text.push(value.to_string());
656 }
657 Some("file") => {
658 let part = exact_object_fields(
659 part,
660 &["id", "type", "mime", "url", "filename", "source"],
661 )?;
662 if let Some(id) = part.get("id") {
663 let id = id
664 .as_str()
665 .filter(|id| id.starts_with("prt_") && id.len() > 4)
666 .ok_or_else(|| {
667 invalid_message("file part id must use the traced `prt_` prefix")
668 })?;
669 part_id.get_or_insert_with(|| id.to_string());
670 }
671 let resolved = resolve_file_part(part, workspace)?;
672 if let Some(block) = resolved.text_block {
673 text.push(block);
674 }
675 if let Some(image_url) = resolved.image_url {
676 image_urls.push(image_url);
677 }
678 }
679 Some(kind) => {
680 return Err(invalid_message(format!(
681 "unsupported OpenCode message part type `{kind}`"
682 )))
683 }
684 None => return Err(invalid_message("message part requires string `type`")),
685 }
686 }
687 let prompt = text.join("\n");
688 if prompt.is_empty() && image_urls.is_empty() {
689 return Err(invalid_message("message contains no usable input"));
690 }
691 Ok(SubmittedMessage {
692 prompt,
693 image_urls,
694 message_id: message_id.to_string(),
695 part_id: part_id.unwrap_or_else(|| stable_id("part", "prt", message_id, 24)),
696 })
697}
698
699fn traced_part_id(part: &serde_json::Map<String, Value>) -> Result<&str, AdapterError> {
700 part.get("id")
701 .and_then(Value::as_str)
702 .filter(|id| id.starts_with("prt_") && id.len() > 4)
703 .ok_or_else(|| invalid_message("text part requires a traced `prt_` string id"))
704}
705
706struct ResolvedFilePart {
707 text_block: Option<String>,
708 image_url: Option<String>,
709}
710
711fn resolve_file_part(
712 part: &serde_json::Map<String, Value>,
713 workspace: &Path,
714) -> Result<ResolvedFilePart, AdapterError> {
715 const MAX_ATTACHMENT_BYTES: u64 = 10 * 1024 * 1024;
716 let mime = part
717 .get("mime")
718 .and_then(Value::as_str)
719 .filter(|mime| !mime.is_empty())
720 .ok_or_else(|| invalid_message("file part requires nonempty string `mime`"))?;
721 let raw_url = part
722 .get("url")
723 .and_then(Value::as_str)
724 .filter(|url| !url.is_empty())
725 .ok_or_else(|| invalid_message("file part requires nonempty string `url`"))?;
726 let filename = part
727 .get("filename")
728 .and_then(Value::as_str)
729 .filter(|name| !name.is_empty())
730 .unwrap_or("attachment");
731
732 if raw_url.starts_with("data:") {
733 let bytes = decode_data_uri(raw_url, mime)?;
734 if mime.starts_with("image/") {
735 return Ok(ResolvedFilePart {
736 text_block: None,
737 image_url: Some(raw_url.to_string()),
738 });
739 }
740 let value = String::from_utf8(bytes)
741 .map_err(|_| invalid_message("non-image data URI attachment must be UTF-8 text"))?;
742 return Ok(ResolvedFilePart {
743 text_block: Some(format!("[file: {filename}]\n{value}")),
744 image_url: None,
745 });
746 }
747
748 let parsed = Url::parse(raw_url)
749 .map_err(|_| invalid_message("file part url must be a data:, file:, or https: URL"))?;
750 if matches!(parsed.scheme(), "http" | "https") {
751 if !mime.starts_with("image/") {
752 return Err(invalid_message(
753 "remote non-image attachments are not fetched by the runtime",
754 ));
755 }
756 return Ok(ResolvedFilePart {
757 text_block: None,
758 image_url: Some(raw_url.to_string()),
759 });
760 }
761 if parsed.scheme() != "file" {
762 return Err(invalid_message("file part URL scheme is not supported"));
763 }
764 let path = parsed
765 .to_file_path()
766 .map_err(|_| invalid_message("file part URL is not a valid runtime file path"))?;
767 let path = resolve_runtime_path(workspace, &path)?;
768 let metadata = std::fs::metadata(&path)
769 .map_err(|error| invalid_message(format!("runtime attachment is unreadable: {error}")))?;
770 if !metadata.is_file() || metadata.len() > MAX_ATTACHMENT_BYTES {
771 return Err(invalid_message(
772 "runtime attachment must be a file no larger than 10 MiB",
773 ));
774 }
775 let bytes = std::fs::read(&path)
776 .map_err(|error| invalid_message(format!("runtime attachment is unreadable: {error}")))?;
777 if mime.starts_with("image/") {
778 let encoded = base64::engine::general_purpose::STANDARD.encode(bytes);
779 Ok(ResolvedFilePart {
780 text_block: None,
781 image_url: Some(format!("data:{mime};base64,{encoded}")),
782 })
783 } else {
784 let value = String::from_utf8(bytes)
785 .map_err(|_| invalid_message("non-image runtime attachment must be UTF-8 text"))?;
786 Ok(ResolvedFilePart {
787 text_block: Some(format!("[file: {filename}]\n{value}")),
788 image_url: None,
789 })
790 }
791}
792
793fn decode_data_uri(raw: &str, mime: &str) -> Result<Vec<u8>, AdapterError> {
794 let expected = format!("data:{mime};base64,");
795 let encoded = raw
796 .strip_prefix(&expected)
797 .ok_or_else(|| invalid_message("data URI mime must match file part mime and use base64"))?;
798 base64::engine::general_purpose::STANDARD
799 .decode(encoded)
800 .map_err(|_| invalid_message("file part contains invalid base64 data"))
801}
802
803fn resolve_runtime_path(workspace: &Path, requested: &Path) -> Result<PathBuf, AdapterError> {
804 let logical = logical_runtime_path(&path_text(workspace), &path_text(requested))?;
805 let workspace = std::fs::canonicalize(workspace).map_err(|error| {
806 invalid_message(format!("SDK runtime workspace is unavailable: {error}"))
807 })?;
808 let logical_path = PathBuf::from(logical);
809 let candidate = if logical_path.is_absolute() {
810 logical_path
811 } else {
812 workspace.join(logical_path)
813 };
814 let candidate = std::fs::canonicalize(candidate)
815 .map_err(|error| invalid_message(format!("runtime attachment is unavailable: {error}")))?;
816 if !candidate.starts_with(&workspace) {
817 return Err(invalid_message(
818 "attachment path escapes the SDK runtime workspace",
819 ));
820 }
821 Ok(candidate)
822}
823
824fn logical_runtime_path(workspace: &str, requested: &str) -> Result<String, AdapterError> {
825 let workspace = normalized_logical_path(workspace)?;
826 let requested = requested.replace('\\', "/");
827 let absolute = requested.starts_with('/')
828 || requested
829 .as_bytes()
830 .get(1)
831 .is_some_and(|separator| *separator == b':');
832 let candidate = if absolute {
833 requested
834 } else {
835 format!("{workspace}/{requested}")
836 };
837 let candidate = normalized_logical_path(&candidate)?;
838 let prefix = format!("{workspace}/");
839 if candidate != workspace && !candidate.starts_with(&prefix) {
840 return Err(invalid_message(
841 "attachment path escapes the SDK runtime workspace",
842 ));
843 }
844 Ok(candidate)
845}
846
847fn normalized_logical_path(raw: &str) -> Result<String, AdapterError> {
848 let raw = raw.replace('\\', "/");
849 let (prefix, tail) = if raw.starts_with('/') {
850 ("/".to_string(), raw.trim_start_matches('/'))
851 } else if raw
852 .as_bytes()
853 .get(1)
854 .is_some_and(|separator| *separator == b':')
855 {
856 (
857 raw[..2].to_ascii_uppercase(),
858 raw[2..].trim_start_matches('/'),
859 )
860 } else {
861 (String::new(), raw.as_str())
862 };
863 let mut segments = Vec::new();
864 for segment in tail.split('/') {
865 match segment {
866 "" | "." => {}
867 ".." => {
868 if segments.pop().is_none() {
869 return Err(invalid_message("attachment path escapes its path root"));
870 }
871 }
872 value => segments.push(value),
873 }
874 }
875 let joined = segments.join("/");
876 Ok(match prefix.as_str() {
877 "/" => format!("/{joined}"),
878 "" => joined,
879 drive => format!("{drive}/{joined}"),
880 })
881}
882
883fn validate_requested_model(
884 body: &Value,
885 descriptor: &FrontendRuntimeDescriptor,
886) -> Result<(), AdapterError> {
887 let model = body
888 .get("model")
889 .and_then(Value::as_object)
890 .ok_or_else(|| invalid_message("model must be an object"))?;
891 let (provider_id, model_id) = provider_model(descriptor);
892 if model.get("providerID").and_then(Value::as_str) != Some(provider_id.as_str())
893 || model.get("modelID").and_then(Value::as_str) != Some(model_id.as_str())
894 {
895 return Err(invalid_message(
896 "model must match the attached SDK runtime descriptor",
897 ));
898 }
899 Ok(())
900}
901
902fn exact_object_fields<'a>(
903 value: &'a Value,
904 allowed: &[&str],
905) -> Result<&'a serde_json::Map<String, Value>, AdapterError> {
906 let object = value
907 .as_object()
908 .ok_or_else(|| invalid_message("body must be a JSON object"))?;
909 if let Some(unexpected) = object.keys().find(|key| !allowed.contains(&key.as_str())) {
910 return Err(invalid_message(format!("unexpected field `{unexpected}`")));
911 }
912 Ok(object)
913}
914
915fn exact_string_field<'a>(
916 value: &'a Value,
917 allowed: &[&str],
918 field: &str,
919) -> Result<&'a str, AdapterError> {
920 let object = exact_object_fields(value, allowed)?;
921 object
922 .get(field)
923 .and_then(Value::as_str)
924 .ok_or_else(|| invalid_message(format!("`{field}` must be a string")))
925}
926
927fn invalid_message(message: impl Into<String>) -> AdapterError {
928 AdapterError::InvalidRequest {
929 route: "POST /session/{session_id}/message".into(),
930 message: message.into(),
931 }
932}
933
934#[derive(Default)]
935struct HistoryIdentityState {
936 previous: Vec<ChatMessage>,
937 ordinals: Vec<u64>,
938 last_history_cursor: Option<u64>,
939 snapshots: BTreeMap<u64, Vec<u64>>,
940 pending_live: Vec<PendingHistoryIdentity>,
941 projected_live: Vec<PendingHistoryIdentity>,
942 live_sources: BTreeMap<(u64, bool), u64>,
943 client_users: Vec<SubmittedMessage>,
944 explicit_ids: BTreeMap<u64, (String, String)>,
945 next_ordinal: u64,
946}
947
948struct PendingHistoryIdentity {
949 ordinal: u64,
950 role: Role,
951 content: String,
952 projected_through: Option<u64>,
953}
954
955struct LiveMessage {
956 ordinal: u64,
957 message_id: String,
958 part_id: String,
959 text: String,
960 created: u64,
961}
962
963pub struct OpenCodeEventProjection {
966 session_id: String,
967 workspace: PathBuf,
968 provider: String,
969 model: String,
970 active: Option<LiveMessage>,
971 history_ids: Arc<Mutex<HistoryIdentityState>>,
972 history_tail: Option<LiveMessage>,
973}
974
975impl OpenCodeEventProjection {
976 fn new(
977 session_id: String,
978 workspace: PathBuf,
979 descriptor: &FrontendRuntimeDescriptor,
980 history_ids: Arc<Mutex<HistoryIdentityState>>,
981 history_tail: Option<(u64, String)>,
982 ) -> Self {
983 let (provider, model) = provider_model(descriptor);
984 let history_tail = history_tail.map(|(ordinal, text)| {
985 let (message_id, part_id) = message_identity(&session_id, ordinal);
986 LiveMessage {
987 ordinal,
988 message_id,
989 part_id,
990 text,
991 created: ordinal,
992 }
993 });
994 Self {
995 session_id,
996 workspace,
997 provider,
998 model,
999 active: None,
1000 history_ids,
1001 history_tail,
1002 }
1003 }
1004
1005 pub fn project(&mut self, event: &SdkEvent) -> Vec<Value> {
1006 match event.kind.as_str() {
1007 "user_message" => self.user_message(event),
1008 "turn_started" => {
1009 let mut events = vec![self.status("busy")];
1010 events.extend(self.ensure_active_events(event.sequence));
1011 events
1012 }
1013 "text_delta" => self.text_delta(event),
1014 "tool_call_started" | "tool_call_completed" => self.tool_event(event),
1015 "request" => self.request(event),
1016 "request_resolved" => self.request_resolved(event),
1017 "turn_succeeded" => self.turn_finished(event, "stop"),
1018 "turn_interrupted" => self.turn_finished(event, "abort"),
1019 "turn_failed" => self.turn_finished(event, "error"),
1020 _ => vec![json!({
1021 "type": "supercode.event",
1022 "properties": {"sequence": event.sequence, "kind": event.kind, "payload": event.payload}
1023 })],
1024 }
1025 }
1026
1027 fn ensure_active_events(&mut self, sequence: u64) -> Vec<Value> {
1028 if self.active.is_some() {
1029 return Vec::new();
1030 }
1031 let active = self.new_live_message(sequence);
1032 self.history_tail = None;
1033 let events = vec![
1034 json!({"type": "message.updated", "properties": {"info": self.assistant_info(&active, None, None)}}),
1035 json!({"type": "message.part.updated", "properties": {"part": {
1036 "id": active.part_id, "sessionID": self.session_id, "messageID": active.message_id,
1037 "type": "text", "text": "", "time": {"start": active.created}
1038 }}}),
1039 ];
1040 self.active = Some(active);
1041 events
1042 }
1043
1044 fn user_message(&self, event: &SdkEvent) -> Vec<Value> {
1045 let text = event
1046 .payload
1047 .get("text")
1048 .and_then(Value::as_str)
1049 .unwrap_or_default();
1050 let (_ordinal, message_id, part_id) = self.reserve_live(Role::User, text, event.sequence);
1051 vec![
1052 json!({"type": "message.updated", "properties": {"info": {
1053 "id": message_id, "sessionID": self.session_id, "role": "user",
1054 "time": {"created": event.sequence}, "agent": "build",
1055 "model": {"providerID": self.provider, "modelID": self.model}
1056 }}}),
1057 json!({"type": "message.part.updated", "properties": {"part": {
1058 "id": part_id, "sessionID": self.session_id, "messageID": message_id,
1059 "type": "text", "text": text
1060 }}}),
1061 ]
1062 }
1063
1064 fn text_delta(&mut self, event: &SdkEvent) -> Vec<Value> {
1065 let delta = event
1066 .payload
1067 .get("text")
1068 .and_then(Value::as_str)
1069 .unwrap_or_default();
1070 let mut events = self.ensure_active_events(event.sequence);
1071 let session_id = self.session_id.clone();
1072 let active = self.active.as_mut().expect("active message was created");
1073 active.text.push_str(delta);
1074 events.push(json!({"type": "message.part.delta", "properties": {
1075 "sessionID": session_id, "messageID": active.message_id,
1076 "partID": active.part_id, "field": "text", "delta": delta
1077 }}));
1078 events
1079 }
1080
1081 fn tool_event(&mut self, event: &SdkEvent) -> Vec<Value> {
1082 let mut events = self.ensure_active_events(event.sequence);
1083 let session_id = self.session_id.clone();
1084 let active = self.active.as_ref().expect("active message was created");
1085 let call_id = event
1086 .payload
1087 .get("id")
1088 .and_then(Value::as_str)
1089 .unwrap_or("unknown");
1090 let tool = event
1091 .payload
1092 .get("name")
1093 .and_then(Value::as_str)
1094 .unwrap_or("tool");
1095 let completed = event.kind == "tool_call_completed";
1096 let part_id = stable_id("live-tool-part", "prt", call_id, 24);
1097 events.push(json!({"type": "message.part.updated", "properties": {"part": {
1098 "id": part_id, "sessionID": session_id, "messageID": active.message_id,
1099 "type": "tool", "callID": call_id, "tool": tool,
1100 "state": if completed {
1101 json!({"status": "completed", "input": event.payload.get("arguments").cloned().unwrap_or(Value::Null),
1102 "output": event.payload.get("output").cloned().unwrap_or(Value::Null),
1103 "time": {"start": event.sequence, "end": event.sequence}})
1104 } else {
1105 json!({"status": "running", "input": event.payload.get("arguments").cloned().unwrap_or(Value::Null),
1106 "time": {"start": event.sequence}})
1107 }
1108 }}}));
1109 events
1110 }
1111
1112 fn request(&self, event: &SdkEvent) -> Vec<Value> {
1113 let request = &event.payload["request"];
1114 if request["kind"] != "approval" {
1115 return vec![self.opaque(event)];
1116 }
1117 let Some(request_id) = request["id"].as_u64() else {
1118 return vec![self.opaque(event)];
1119 };
1120 let payload = &request["payload"];
1121 let subject = payload.get("subject").cloned().unwrap_or(Value::Null);
1122 vec![json!({"type": "permission.asked", "properties": {
1123 "id": permission_id(request_id), "sessionID": self.session_id,
1124 "permission": payload.get("tool").cloned().unwrap_or_else(|| json!("tool")),
1125 "patterns": if subject.is_null() { json!([]) } else { json!([subject]) },
1126 "metadata": {"supercode": {"request": request, "sequence": event.sequence}},
1127 "always": ["*"]
1128 }})]
1129 }
1130
1131 fn request_resolved(&self, event: &SdkEvent) -> Vec<Value> {
1132 let Some(request_id) = event.payload.get("request_id").and_then(Value::as_u64) else {
1133 return vec![self.opaque(event)];
1134 };
1135 let decision = event
1136 .payload
1137 .pointer("/response/decision")
1138 .and_then(Value::as_str);
1139 let reply = match decision {
1140 Some("allow") => "once",
1141 Some("allow_for_session") => "always",
1142 _ => "reject",
1143 };
1144 vec![json!({"type": "permission.replied", "properties": {
1145 "sessionID": self.session_id, "requestID": permission_id(request_id), "reply": reply
1146 }})]
1147 }
1148
1149 fn turn_finished(&mut self, event: &SdkEvent, finish: &str) -> Vec<Value> {
1150 let reply = event
1151 .payload
1152 .get("reply")
1153 .or_else(|| event.payload.get("message"))
1154 .and_then(Value::as_str)
1155 .unwrap_or_default();
1156 let history_match = self.active.is_none()
1157 && finish == "stop"
1158 && self
1159 .history_tail
1160 .as_ref()
1161 .is_some_and(|history| history.text == reply);
1162 let mut events = if history_match {
1163 Vec::new()
1164 } else {
1165 self.ensure_active_events(event.sequence)
1166 };
1167 let mut active = if history_match {
1168 self.history_tail.take().expect("history match was checked")
1169 } else {
1170 self.active.take().expect("active message was created")
1171 };
1172 if active.text.is_empty() {
1173 active.text = reply.to_string();
1174 }
1175 if !history_match {
1176 self.update_live(active.ordinal, Role::Assistant, &active.text);
1177 }
1178 events.extend([
1179 json!({"type": "message.part.updated", "properties": {"part": {
1180 "id": active.part_id, "sessionID": self.session_id, "messageID": active.message_id,
1181 "type": "text", "text": active.text,
1182 "time": {"start": active.created, "end": event.sequence}
1183 }}}),
1184 json!({"type": "message.updated", "properties": {"info": self.assistant_info(&active, Some(event.sequence), Some(finish))}}),
1185 self.status("idle"),
1186 json!({"type": "session.idle", "properties": {"sessionID": self.session_id}}),
1187 ]);
1188 if finish == "error" {
1189 let completed = events.len() - 3;
1190 events[completed]["properties"]["info"]["error"] = json!({
1191 "name": "APIError", "data": {"message": event.payload.get("message").cloned().unwrap_or_default()}
1192 });
1193 }
1194 events
1195 }
1196
1197 fn new_live_message(&self, sequence: u64) -> LiveMessage {
1198 let (ordinal, message_id, part_id) = self.reserve_live(Role::Assistant, "", sequence);
1199 LiveMessage {
1200 ordinal,
1201 part_id,
1202 message_id,
1203 text: String::new(),
1204 created: sequence,
1205 }
1206 }
1207
1208 fn assistant_info(
1209 &self,
1210 active: &LiveMessage,
1211 completed: Option<u64>,
1212 finish: Option<&str>,
1213 ) -> Value {
1214 let mut info = json!({
1215 "id": active.message_id, "sessionID": self.session_id, "role": "assistant",
1216 "time": {"created": active.created},
1217 "modelID": self.model, "providerID": self.provider, "mode": "build", "agent": "build",
1218 "path": {"cwd": path_text(&self.workspace), "root": path_text(&self.workspace)},
1219 "cost": 0, "tokens": {"input": 0, "output": 0, "reasoning": 0, "cache": {"read": 0, "write": 0}},
1220 });
1221 if let Some(completed) = completed {
1222 info["time"]["completed"] = json!(completed);
1223 }
1224 if let Some(finish) = finish {
1225 info["finish"] = json!(finish);
1226 }
1227 info
1228 }
1229
1230 fn status(&self, status: &str) -> Value {
1231 json!({"type": "session.status", "properties": {
1232 "sessionID": self.session_id, "status": {"type": status}
1233 }})
1234 }
1235
1236 fn opaque(&self, event: &SdkEvent) -> Value {
1237 json!({"type": "supercode.event", "properties": {
1238 "sequence": event.sequence, "kind": event.kind, "payload": event.payload
1239 }})
1240 }
1241
1242 fn reserve_live(
1243 &self,
1244 role: Role,
1245 content: &str,
1246 source_sequence: u64,
1247 ) -> (u64, String, String) {
1248 let mut state = self
1249 .history_ids
1250 .lock()
1251 .unwrap_or_else(std::sync::PoisonError::into_inner);
1252 let ordinal = state.reserve_live(role, content, source_sequence);
1253 let (message_id, part_id) = state.identity(&self.session_id, ordinal);
1254 (ordinal, message_id, part_id)
1255 }
1256
1257 fn update_live(&self, ordinal: u64, role: Role, content: &str) {
1258 self.history_ids
1259 .lock()
1260 .unwrap_or_else(std::sync::PoisonError::into_inner)
1261 .update_live(ordinal, role, content);
1262 }
1263}
1264
1265impl HistoryIdentityState {
1266 fn project(&mut self, current: &[ChatMessage], history_cursor: u64) -> Vec<u64> {
1270 if let Some(projected) = self
1271 .snapshots
1272 .get(&history_cursor)
1273 .filter(|projected| projected.len() == current.len())
1274 {
1275 return projected.clone();
1276 }
1277 if self
1278 .last_history_cursor
1279 .is_some_and(|cursor| history_cursor < cursor)
1280 {
1281 let overlap = (0..=self.previous.len().min(current.len()))
1285 .rev()
1286 .find(|count| current[current.len() - *count..] == self.previous[..*count])
1287 .unwrap_or(0);
1288 let missing = current.len().saturating_sub(overlap);
1289 let mut projected = (0..missing).map(|_| self.allocate()).collect::<Vec<_>>();
1290 projected.extend_from_slice(&self.ordinals[..overlap]);
1291 self.remember_snapshot(history_cursor, &projected);
1292 return projected;
1293 }
1294 let overlap = (0..=self.previous.len().min(current.len()))
1295 .rev()
1296 .find(|count| self.previous[self.previous.len() - *count..] == current[..*count])
1297 .unwrap_or(0);
1298 let mut projected = self.ordinals[self.ordinals.len() - overlap..].to_vec();
1299 for message in ¤t[overlap..] {
1300 let content = message_text(message);
1301 let pending = self
1302 .pending_live
1303 .iter()
1304 .position(|pending| pending.role == message.role && pending.content == content)
1305 .map(|index| self.pending_live.remove(index).ordinal);
1306 let allocated = pending.is_none();
1307 let ordinal = pending.unwrap_or_else(|| self.allocate());
1308 self.claim_client_user(ordinal, &message.role, &content);
1309 if allocated && matches!(message.role, Role::User | Role::Assistant) {
1310 self.projected_live.push(PendingHistoryIdentity {
1316 ordinal,
1317 role: message.role,
1318 content,
1319 projected_through: Some(history_cursor),
1320 });
1321 }
1322 projected.push(ordinal);
1323 }
1324 self.projected_live
1325 .retain(|pending| projected.binary_search(&pending.ordinal).is_ok());
1326 self.previous = current.to_vec();
1327 self.ordinals = projected.clone();
1328 self.last_history_cursor = Some(history_cursor);
1329 self.remember_snapshot(history_cursor, &projected);
1330 projected
1331 }
1332
1333 fn remember_snapshot(&mut self, history_cursor: u64, projected: &[u64]) {
1334 const SNAPSHOT_LIMIT: usize = 64;
1335 self.snapshots.insert(history_cursor, projected.to_vec());
1336 while self.snapshots.len() > SNAPSHOT_LIMIT {
1337 self.snapshots.pop_first();
1338 }
1339 }
1340
1341 fn register_client_user(&mut self, submitted: &SubmittedMessage) {
1342 self.client_users.push(SubmittedMessage {
1343 prompt: submitted.prompt.clone(),
1344 image_urls: submitted.image_urls.clone(),
1345 message_id: submitted.message_id.clone(),
1346 part_id: submitted.part_id.clone(),
1347 });
1348 }
1349
1350 fn reserve_live(&mut self, role: Role, content: &str, source_sequence: u64) -> u64 {
1351 let source = (source_sequence, role == Role::User);
1352 if let Some(ordinal) = self.live_sources.get(&source) {
1353 return *ordinal;
1354 }
1355 let matches = |pending: &PendingHistoryIdentity| {
1356 pending.role == role
1357 && pending
1358 .projected_through
1359 .is_some_and(|cursor| source_sequence <= cursor)
1360 };
1361 let projected = if role == Role::Assistant && content.is_empty() {
1362 self.projected_live.iter().rposition(matches)
1365 } else {
1366 self.projected_live
1367 .iter()
1368 .position(|pending| matches(pending) && pending.content == content)
1369 };
1370 let ordinal = projected
1371 .map(|index| self.projected_live.remove(index).ordinal)
1372 .unwrap_or_else(|| self.allocate());
1373 self.live_sources.insert(source, ordinal);
1374 self.claim_client_user(ordinal, &role, content);
1375 if projected.is_none() {
1376 self.pending_live.push(PendingHistoryIdentity {
1377 ordinal,
1378 role,
1379 content: content.to_string(),
1380 projected_through: None,
1381 });
1382 }
1383 ordinal
1384 }
1385
1386 fn update_live(&mut self, ordinal: u64, role: Role, content: &str) {
1387 if let Some(pending) = self
1388 .pending_live
1389 .iter_mut()
1390 .find(|pending| pending.ordinal == ordinal)
1391 {
1392 pending.role = role;
1393 pending.content = content.to_string();
1394 }
1395 }
1396
1397 fn allocate(&mut self) -> u64 {
1398 let ordinal = self.next_ordinal;
1399 self.next_ordinal = self.next_ordinal.saturating_add(1);
1400 ordinal
1401 }
1402
1403 fn identity(&self, session_id: &str, ordinal: u64) -> (String, String) {
1404 self.explicit_ids
1405 .get(&ordinal)
1406 .cloned()
1407 .unwrap_or_else(|| message_identity(session_id, ordinal))
1408 }
1409
1410 fn claim_client_user(&mut self, ordinal: u64, role: &Role, content: &str) {
1411 if *role != Role::User || self.explicit_ids.contains_key(&ordinal) {
1412 return;
1413 }
1414 if let Some(index) = self
1415 .client_users
1416 .iter()
1417 .position(|submitted| submitted.prompt == content)
1418 {
1419 let submitted = self.client_users.remove(index);
1420 self.explicit_ids
1421 .insert(ordinal, (submitted.message_id, submitted.part_id));
1422 }
1423 }
1424}
1425
1426fn provider_projection(descriptor: &FrontendRuntimeDescriptor, connected: bool) -> Value {
1427 let (provider, model) = provider_model(descriptor);
1428 let row = json!({
1429 "id": provider,
1430 "name": "Supercode runtime",
1431 "env": [],
1432 "options": {},
1433 "source": "custom",
1434 "models": {
1435 model.clone(): {
1436 "id": model,
1437 "api": {"id": model, "npm": "@ai-sdk/openai-compatible"},
1438 "status": "active",
1439 "name": descriptor.model,
1440 "providerID": provider,
1441 "capabilities": {
1442 "temperature": false,
1443 "reasoning": true,
1444 "attachment": false,
1445 "toolcall": true,
1446 "input": {"text": true, "audio": false, "image": false, "video": false, "pdf": false},
1447 "output": {"text": true, "audio": false, "image": false, "video": false, "pdf": false},
1448 "interleaved": false
1449 },
1450 "cost": {"input": 0, "output": 0, "cache": {"read": 0, "write": 0}},
1451 "options": {},
1452 "limit": {"context": 0, "output": 0},
1453 "headers": {},
1454 "family": "",
1455 "release_date": "",
1456 "variants": {}
1457 }
1458 }
1459 });
1460 if connected {
1461 json!({"all": [row], "default": {provider.clone(): model}, "connected": [provider]})
1462 } else {
1463 json!({"providers": [row], "default": {provider: model}})
1464 }
1465}
1466
1467fn provider_model(descriptor: &FrontendRuntimeDescriptor) -> (String, String) {
1468 descriptor
1469 .model
1470 .split_once('/')
1471 .map(|(provider, model)| (sanitize_id(provider), sanitize_id(model)))
1472 .unwrap_or_else(|| ("supercode".into(), sanitize_id(&descriptor.model)))
1473}
1474
1475fn commands(descriptor: &FrontendRuntimeDescriptor) -> Value {
1476 Value::Array(
1477 descriptor
1478 .commands
1479 .iter()
1480 .map(|command| {
1481 json!({
1482 "name": command.name,
1483 "description": command.description,
1484 "source": "command",
1485 "template": "$ARGUMENTS",
1486 "hints": ["$ARGUMENTS"],
1487 })
1488 })
1489 .collect(),
1490 )
1491}
1492
1493fn history_messages(
1494 history: &[ChatMessage],
1495 ordinals: &[u64],
1496 identities: &[(String, String)],
1497 session_id: &str,
1498 workspace: &Path,
1499 descriptor: &FrontendRuntimeDescriptor,
1500) -> Value {
1501 debug_assert_eq!(history.len(), ordinals.len());
1502 debug_assert_eq!(history.len(), identities.len());
1503 Value::Array(
1504 history
1505 .iter()
1506 .zip(ordinals)
1507 .zip(identities)
1508 .map(|((message, ordinal), identity)| {
1509 history_message(
1510 message, session_id, workspace, *ordinal, identity, descriptor,
1511 )
1512 })
1513 .collect(),
1514 )
1515}
1516
1517fn history_message(
1518 message: &ChatMessage,
1519 session_id: &str,
1520 workspace: &Path,
1521 ordinal: u64,
1522 identity: &(String, String),
1523 descriptor: &FrontendRuntimeDescriptor,
1524) -> Value {
1525 let (message_id, part_id) = identity;
1526 let (provider, model) = provider_model(descriptor);
1527 let content = message_text(message);
1528 match message.role {
1529 Role::User => json!({
1530 "info": {
1531 "role": "user", "time": {"created": ordinal}, "summary": {"diffs": []},
1532 "agent": "build", "model": {"providerID": provider, "modelID": model},
1533 "id": message_id, "sessionID": session_id
1534 },
1535 "parts": [{"type": "text", "text": content, "id": part_id, "sessionID": session_id, "messageID": message_id}]
1536 }),
1537 Role::Assistant | Role::System | Role::Tool => {
1538 let content = match message.role {
1539 Role::System => format!("[system]\n{content}"),
1540 Role::Tool => format!("[tool result]\n{content}"),
1541 _ => content,
1542 };
1543 json!({
1544 "info": {
1545 "role": "assistant", "time": {"created": ordinal, "completed": ordinal},
1546 "modelID": model, "providerID": provider, "mode": "build", "agent": "build",
1547 "path": {"cwd": path_text(workspace), "root": path_text(workspace)},
1548 "cost": 0, "tokens": {"input": 0, "output": 0, "reasoning": 0, "cache": {"read": 0, "write": 0}},
1549 "finish": "stop", "id": message_id, "sessionID": session_id
1550 },
1551 "parts": [{"type": "text", "text": content, "time": {"start": ordinal, "end": ordinal}, "id": part_id, "sessionID": session_id, "messageID": message_id}]
1552 })
1553 }
1554 }
1555}
1556
1557fn message_identity(session_id: &str, ordinal: u64) -> (String, String) {
1558 let message_id = stable_id("message", "msg", &format!("{session_id}:{ordinal}"), 24);
1559 let part_id = stable_id("part", "prt", &format!("{message_id}:0"), 24);
1560 (message_id, part_id)
1561}
1562
1563fn message_text(message: &ChatMessage) -> String {
1564 if let Some(content) = &message.content {
1565 return content.clone();
1566 }
1567 if let Some(parts) = &message.content_parts {
1568 return parts
1569 .iter()
1570 .map(Value::to_string)
1571 .collect::<Vec<_>>()
1572 .join("\n");
1573 }
1574 if let Some(tool_calls) = &message.tool_calls {
1575 return serde_json::to_string(tool_calls).unwrap_or_else(|_| "[tool calls]".into());
1576 }
1577 String::new()
1578}
1579
1580fn stable_id(domain: &str, prefix: &str, source: &str, digits: usize) -> String {
1581 let mut input = Vec::with_capacity(domain.len() + source.len() + 1);
1582 input.extend_from_slice(domain.as_bytes());
1583 input.push(0);
1584 input.extend_from_slice(source.as_bytes());
1585 let digest = blake3::hash(&input).to_hex().to_string();
1586 format!("{prefix}_{}", &digest[..digits])
1587}
1588
1589fn sanitize_id(value: &str) -> String {
1590 let sanitized = value
1591 .chars()
1592 .map(|character| {
1593 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_' | '.') {
1594 character
1595 } else {
1596 '-'
1597 }
1598 })
1599 .collect::<String>();
1600 if sanitized.is_empty() {
1601 "runtime".into()
1602 } else {
1603 sanitized
1604 }
1605}
1606
1607fn path_text(path: &Path) -> String {
1608 path.to_string_lossy().replace('\\', "/")
1609}
1610
1611fn truncate(text: &str, limit: usize) -> String {
1612 text.chars().take(limit).collect::<String>()
1613}
1614
1615#[cfg(test)]
1616mod tests {
1617 use std::collections::{BTreeMap, VecDeque};
1618 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
1619
1620 use async_trait::async_trait;
1621 use supercode::{
1622 FrontendActions, FrontendAttachSnapshot, FrontendConnectionState,
1623 FrontendDisplayCapabilities, FrontendTurnState,
1624 };
1625 use tokio::sync::broadcast;
1626
1627 use super::*;
1628
1629 struct FixtureRuntime {
1630 history: Mutex<Vec<ChatMessage>>,
1631 descriptor: FrontendRuntimeDescriptor,
1632 events: broadcast::Sender<supercode::FrontendEvent>,
1633 submissions: Mutex<Vec<String>>,
1634 image_submissions: Mutex<Vec<Vec<String>>>,
1635 steers: Mutex<Vec<String>>,
1636 responses: Mutex<Vec<FrontendResponse>>,
1637 interrupts: AtomicUsize,
1638 steer_accepting: AtomicBool,
1639 }
1640
1641 impl FixtureRuntime {
1642 fn new() -> Arc<Self> {
1643 Self::with_turn_state(FrontendTurnState::Idle)
1644 }
1645
1646 fn with_turn_state(turn_state: FrontendTurnState) -> Arc<Self> {
1647 let (events, _) = broadcast::channel(16);
1648 Arc::new(Self {
1649 history: Mutex::new(vec![
1650 ChatMessage::system("preserve system context"),
1651 ChatMessage::user("hello from Claude"),
1652 ChatMessage::assistant("continued through GLM"),
1653 ]),
1654 descriptor: FrontendRuntimeDescriptor {
1655 schema_version: 2,
1656 session_id: "runtime-1".into(),
1657 source_harness: Some("claude-code".into()),
1658 emulation_profile: Some("claude-code".into()),
1659 active_modules: Vec::new(),
1660 commands: Vec::new(),
1661 operations: Vec::new(),
1662 actions: FrontendActions {
1663 submit: true,
1664 interrupt: true,
1665 steer: true,
1666 respond: true,
1667 detach: true,
1668 close: false,
1669 },
1670 display: FrontendDisplayCapabilities {
1671 event_kinds: vec!["assistant_delta".into()],
1672 opaque_fallback: true,
1673 },
1674 model: "openrouter/glm-5.2".into(),
1675 turn_state,
1676 connection_state: FrontendConnectionState::Connected,
1677 extensions: BTreeMap::new(),
1678 },
1679 events,
1680 submissions: Mutex::new(Vec::new()),
1681 image_submissions: Mutex::new(Vec::new()),
1682 steers: Mutex::new(Vec::new()),
1683 responses: Mutex::new(Vec::new()),
1684 interrupts: AtomicUsize::new(0),
1685 steer_accepting: AtomicBool::new(true),
1686 })
1687 }
1688 }
1689
1690 #[async_trait]
1691 impl SdkRuntime for FixtureRuntime {
1692 async fn describe(&self) -> Result<FrontendRuntimeDescriptor, SdkError> {
1693 Ok(self.descriptor.clone())
1694 }
1695
1696 async fn attach(&self, history_limit: usize) -> Result<FrontendAttachment, SdkError> {
1697 let history = self
1698 .history
1699 .lock()
1700 .unwrap_or_else(std::sync::PoisonError::into_inner)
1701 .clone();
1702 let start = history.len().saturating_sub(history_limit);
1703 let replay = if self.descriptor.turn_state == FrontendTurnState::Busy {
1704 VecDeque::from([supercode::FrontendEvent {
1705 sequence: 4,
1706 kind: "turn_succeeded".into(),
1707 payload: json!({
1708 "type": "turn_succeeded",
1709 "reply": "continued through GLM"
1710 }),
1711 }])
1712 } else {
1713 VecDeque::new()
1714 };
1715 Ok(FrontendAttachment::from_snapshot(
1716 FrontendAttachSnapshot {
1717 descriptor: self.descriptor.clone(),
1718 history: history[start..].to_vec(),
1719 history_cursor: history.len() as u64,
1720 replay,
1721 },
1722 self.events.subscribe(),
1723 ))
1724 }
1725
1726 async fn send_input(self: Arc<Self>, _prompt: String) -> Result<(), SdkError> {
1727 Ok(())
1728 }
1729
1730 async fn submit(&self, prompt: String) -> Result<String, SdkError> {
1731 self.submissions
1732 .lock()
1733 .unwrap_or_else(std::sync::PoisonError::into_inner)
1734 .push(prompt.clone());
1735 let reply = format!("continued:{prompt}");
1736 self.history
1737 .lock()
1738 .unwrap_or_else(std::sync::PoisonError::into_inner)
1739 .extend([
1740 ChatMessage::user(prompt),
1741 ChatMessage::assistant(reply.clone()),
1742 ]);
1743 Ok(reply)
1744 }
1745
1746 async fn submit_with_images(
1747 &self,
1748 prompt: String,
1749 image_urls: Vec<String>,
1750 ) -> Result<String, SdkError> {
1751 self.submissions
1752 .lock()
1753 .unwrap_or_else(std::sync::PoisonError::into_inner)
1754 .push(prompt.clone());
1755 self.image_submissions
1756 .lock()
1757 .unwrap_or_else(std::sync::PoisonError::into_inner)
1758 .push(image_urls.clone());
1759 let reply = format!("continued:{prompt}");
1760 self.history
1761 .lock()
1762 .unwrap_or_else(std::sync::PoisonError::into_inner)
1763 .extend([
1764 ChatMessage::user_with_images(prompt, &image_urls),
1765 ChatMessage::assistant(reply.clone()),
1766 ]);
1767 Ok(reply)
1768 }
1769
1770 async fn interrupt(&self) -> Result<bool, SdkError> {
1771 self.interrupts.fetch_add(1, Ordering::SeqCst);
1772 Ok(true)
1773 }
1774
1775 async fn steer(&self, prompt: String) -> Result<(), SdkError> {
1776 if !self.steer_accepting.load(Ordering::SeqCst) {
1777 return Err(SdkError::UnsupportedAction("steer"));
1778 }
1779 self.steers
1780 .lock()
1781 .unwrap_or_else(std::sync::PoisonError::into_inner)
1782 .push(prompt.clone());
1783 let mut history = self
1784 .history
1785 .lock()
1786 .unwrap_or_else(std::sync::PoisonError::into_inner);
1787 history.push(ChatMessage::assistant(format!("steered:{prompt}")));
1788 let sequence = history.len() as u64 + 1;
1789 drop(history);
1790 let _ = self.events.send(supercode::FrontendEvent {
1791 sequence,
1792 kind: "turn_succeeded".into(),
1793 payload: json!({"type": "turn_succeeded", "reply": format!("steered:{prompt}")}),
1794 });
1795 Ok(())
1796 }
1797
1798 async fn respond(&self, response: supercode::FrontendResponse) -> Result<(), SdkError> {
1799 self.responses
1800 .lock()
1801 .unwrap_or_else(std::sync::PoisonError::into_inner)
1802 .push(response);
1803 Ok(())
1804 }
1805 }
1806
1807 fn json_body(response: OpenCodeResponse) -> Value {
1808 assert_eq!(response.status, 200);
1809 match response.body {
1810 ResponseBody::Json(body) => body,
1811 ResponseBody::EventStream(_) => panic!("expected JSON response"),
1812 }
1813 }
1814
1815 #[test]
1816 fn namespace_and_traced_route_set_are_exactly_pinned() {
1817 assert_eq!(PROTOCOL_NAMESPACE, "opencode_http/v1_2_15");
1818 assert_eq!(OPENCODE_CLI_VERSION, "1.2.15");
1819 assert_eq!(TRACED_ROUTES.len(), 23);
1820 assert!(TRACED_ROUTES.contains(&"GET /event"));
1821 assert!(TRACED_ROUTES.contains(&"POST /session/{session_id}/abort"));
1822 assert!(TRACED_ROUTES.contains(&"POST /session/{session_id}/message"));
1823 }
1824
1825 #[test]
1826 fn corpus_allowlist_and_provenance_pin_drive_the_namespace_exactly() {
1827 let corpus =
1828 Path::new(env!("CARGO_MANIFEST_DIR")).join("../../scripts/client-protocol-corpus");
1829 let allowlist: Value = serde_json::from_slice(
1830 &std::fs::read(corpus.join("allowlist.json")).expect("read corpus allowlist"),
1831 )
1832 .expect("parse corpus allowlist");
1833 let mut observed = allowlist["clients"][PROTOCOL_NAMESPACE]["client_to_server"]
1834 .as_object()
1835 .expect("OpenCode route allowlist")
1836 .keys()
1837 .map(String::as_str)
1838 .collect::<Vec<_>>();
1839 let mut implemented = TRACED_ROUTES.to_vec();
1840 observed.sort_unstable();
1841 implemented.sort_unstable();
1842 assert_eq!(implemented, observed);
1843
1844 let pins: Value = serde_json::from_slice(
1845 &std::fs::read(corpus.join("pins.json")).expect("read corpus pins"),
1846 )
1847 .expect("parse corpus pins");
1848 let pin = pins["clients"]
1849 .as_array()
1850 .expect("client pins")
1851 .iter()
1852 .find(|pin| pin["id"] == "opencode")
1853 .expect("OpenCode pin");
1854 assert_eq!(pin["version"], OPENCODE_CLI_VERSION);
1855 assert_eq!(pin["namespace"], PROTOCOL_NAMESPACE);
1856 assert_eq!(pin["commit"], "799b2623cbb1c0f19e045d87c2c8593e83678bc0");
1857 assert_eq!(pin["license"]["spdx"], "MIT");
1858 assert_eq!(
1859 pin["contract"]["sha256"],
1860 "cfb4d87bc11924794a1fe9f3daafb8b569cc412bccda85b5014240bc3afe7eff"
1861 );
1862 }
1863
1864 #[tokio::test]
1865 async fn sanitized_stock_exchange_replays_every_request_and_reconnect_boundary() {
1866 let runtime = FixtureRuntime::new();
1867 let adapter = OpenCodeAdapter::new(runtime, "runtime-1", "/workspace");
1868 let fixture = Path::new(env!("CARGO_MANIFEST_DIR"))
1869 .join("../../scripts/client-protocol-corpus/fixtures/opencode_v1_2_15.jsonl");
1870 let text = std::fs::read_to_string(fixture).expect("read sanitized OpenCode fixture");
1871 let mut requests = 0_usize;
1872 let mut event_attachments = Vec::new();
1873 let mut last_attachment = 0_u64;
1874 for line in text.lines() {
1875 let row: Value = serde_json::from_str(line).expect("parse fixture row");
1876 if row["direction"] != "client_to_server" {
1877 continue;
1878 }
1879 requests += 1;
1880 let attachment = row["attachment"].as_u64().expect("attachment id");
1881 assert!(attachment >= last_attachment, "reconnect order regressed");
1882 last_attachment = attachment;
1883 let message = &row["message"];
1884 let method = message["method"].as_str().expect("method");
1885 let original = message["path"].as_str().expect("path");
1886 let target = if let Some(suffix) = original.strip_prefix("/session/") {
1887 let suffix = suffix
1888 .find('/')
1889 .map(|index| &suffix[index..])
1890 .unwrap_or_default();
1891 format!("/session/{}{suffix}", adapter.session_id())
1892 } else if original.starts_with("/permission/") {
1893 "/permission/per_000000000000002a/reply".into()
1894 } else {
1895 original.to_string()
1896 };
1897 let body = message["body"].as_str().unwrap_or_default();
1898 let mut body = if body.is_empty() {
1899 Value::Null
1900 } else {
1901 serde_json::from_str(body).expect("fixture JSON body")
1902 };
1903 if method == "POST" && target.ends_with("/message") {
1904 body["model"] = json!({"providerID": "openrouter", "modelID": "glm-5.2"});
1905 }
1906 let response = adapter
1907 .handle(OpenCodeRequest::new(method, &target).with_body(body))
1908 .await;
1909 assert_eq!(response.status, 200, "fixture request {method} {target}");
1910 if original == "/event" {
1911 assert!(matches!(response.body, ResponseBody::EventStream(_)));
1912 event_attachments.push(attachment);
1913 }
1914 }
1915 assert_eq!(requests, 43);
1916 assert_eq!(event_attachments, vec![1, 2]);
1917 }
1918
1919 #[test]
1920 fn sanitized_stock_server_exchange_pins_live_creation_and_reconnect_identity() {
1921 let fixture = Path::new(env!("CARGO_MANIFEST_DIR"))
1922 .join("../../scripts/client-protocol-corpus/fixtures/opencode_v1_2_15.jsonl");
1923 let text = std::fs::read_to_string(fixture).expect("read sanitized OpenCode fixture");
1924 let mut live_events = Vec::new();
1925 let mut reconnect_history = Vec::new();
1926 for line in text.lines() {
1927 let row: Value = serde_json::from_str(line).expect("parse fixture row");
1928 if row["direction"] != "server_to_client" || row["event"] != "body_chunk" {
1929 continue;
1930 }
1931 let Some(message) = row["message"].as_str() else {
1932 continue;
1933 };
1934 if row["attachment"] == 1 {
1935 for frame in message.split("\n\n") {
1936 if let Some(data) = frame.strip_prefix("data: ") {
1937 live_events.push(serde_json::from_str::<Value>(data).expect("SSE JSON"));
1938 }
1939 }
1940 } else if row["attachment"] == 2 {
1941 if let Ok(history) = serde_json::from_str::<Vec<Value>>(message) {
1942 if history.len() > reconnect_history.len()
1943 && history.iter().all(|item| item.get("info").is_some())
1944 {
1945 reconnect_history = history;
1946 }
1947 }
1948 }
1949 }
1950
1951 let mut messages = std::collections::BTreeSet::new();
1952 let mut parts = std::collections::BTreeSet::new();
1953 let mut live_message_ids = std::collections::BTreeSet::new();
1954 let mut deltas = 0_usize;
1955 for event in &live_events {
1956 match event["type"].as_str() {
1957 Some("message.updated") => {
1958 let info = &event["properties"]["info"];
1959 if let Some(id) = info["id"].as_str() {
1960 messages.insert(id.to_string());
1961 if matches!(info["role"].as_str(), Some("user" | "assistant")) {
1962 live_message_ids.insert(id.to_string());
1963 }
1964 }
1965 }
1966 Some("message.part.updated") => {
1967 let part = &event["properties"]["part"];
1968 let message_id = part["messageID"].as_str().expect("part message id");
1969 assert!(
1970 messages.contains(message_id),
1971 "stock corpus updated part before creating message {message_id}"
1972 );
1973 parts.insert(part["id"].as_str().expect("part id").to_string());
1974 }
1975 Some("message.part.delta") => {
1976 deltas += 1;
1977 let properties = &event["properties"];
1978 assert!(messages.contains(properties["messageID"].as_str().unwrap()));
1979 assert!(parts.contains(properties["partID"].as_str().unwrap()));
1980 }
1981 _ => {}
1982 }
1983 }
1984 assert!(deltas > 0, "stock corpus must contain live text deltas");
1985 assert!(
1986 !reconnect_history.is_empty(),
1987 "reconnect history was not captured"
1988 );
1989 let reconnect_ids = reconnect_history
1990 .iter()
1991 .filter_map(|message| message["info"]["id"].as_str())
1992 .collect::<std::collections::BTreeSet<_>>();
1993 assert!(
1994 live_message_ids
1995 .iter()
1996 .all(|id| reconnect_ids.contains(id.as_str())),
1997 "stock reconnect must preserve every live message identity"
1998 );
1999 }
2000
2001 #[test]
2002 fn deterministic_ids_are_domain_separated_and_collision_free_for_large_sample() {
2003 let mut ids = std::collections::BTreeSet::new();
2004 for index in 0..50_000 {
2005 let source = format!("runtime-{index}");
2006 assert!(ids.insert(stable_id("session", "ses", &source, 24)));
2007 assert!(ids.insert(stable_id("message", "msg", &source, 24)));
2008 assert!(ids.insert(stable_id("part", "prt", &source, 24)));
2009 }
2010 assert_eq!(ids.len(), 150_000);
2011 assert_eq!(
2012 stable_id("session", "id", "same-source", 24),
2013 stable_id("session", "id", "same-source", 24)
2014 );
2015 assert_ne!(
2016 stable_id("session", "id", "same-source", 24),
2017 stable_id("message", "id", "same-source", 24)
2018 );
2019 }
2020
2021 #[test]
2022 fn message_identity_survives_a_bounded_history_window_slide() {
2023 let mut state = HistoryIdentityState::default();
2024 let first = vec![
2025 ChatMessage::user("a"),
2026 ChatMessage::assistant("b"),
2027 ChatMessage::user("c"),
2028 ChatMessage::assistant("d"),
2029 ];
2030 assert_eq!(state.project(&first, 4), vec![0, 1, 2, 3]);
2031
2032 let slid = vec![
2033 ChatMessage::user("c"),
2034 ChatMessage::assistant("d"),
2035 ChatMessage::user("e"),
2036 ChatMessage::assistant("f"),
2037 ];
2038 let projected = state.project(&slid, 5);
2039 assert_eq!(projected, vec![2, 3, 4, 5]);
2040 assert_eq!(state.project(&slid, 5), projected);
2041
2042 let retained_before = stable_id("message", "msg", "ses_test:2", 24);
2043 let retained_after = stable_id("message", "msg", &format!("ses_test:{}", projected[0]), 24);
2044 assert_eq!(retained_before, retained_after);
2045 }
2046
2047 #[tokio::test]
2048 async fn read_routes_project_one_runtime_without_provider_credentials() {
2049 let adapter = OpenCodeAdapter::new(FixtureRuntime::new(), "runtime-1", "/workspace");
2050 let session_id = adapter.session_id().to_string();
2051 let list = json_body(
2052 adapter
2053 .handle(OpenCodeRequest::new("GET", "/session?start=0"))
2054 .await,
2055 );
2056 assert_eq!(list[0]["id"], session_id);
2057 assert_eq!(list[0]["directory"], "/workspace");
2058
2059 let history = json_body(
2060 adapter
2061 .handle(OpenCodeRequest::new(
2062 "GET",
2063 format!("/session/{session_id}/message?limit=100"),
2064 ))
2065 .await,
2066 );
2067 assert_eq!(history.as_array().unwrap().len(), 3);
2068 assert_eq!(
2069 history[0]["parts"][0]["text"],
2070 "[system]\npreserve system context"
2071 );
2072 assert_eq!(history[1]["info"]["role"], "user");
2073 assert_eq!(history[2]["parts"][0]["text"], "continued through GLM");
2074
2075 let provider = json_body(
2076 adapter
2077 .handle(OpenCodeRequest::new("GET", "/provider"))
2078 .await,
2079 );
2080 assert_eq!(provider["connected"], json!(["openrouter"]));
2081 let serialized = provider.to_string();
2082 assert!(!serialized.contains("apiKey"));
2083 assert!(!serialized.contains("OPENROUTER_API_KEY"));
2084 }
2085
2086 #[tokio::test]
2087 async fn unknown_routes_and_wrong_session_ids_fail_closed() {
2088 let adapter = OpenCodeAdapter::new(FixtureRuntime::new(), "runtime-1", "/workspace");
2089 let unknown = adapter
2090 .handle(OpenCodeRequest::new("DELETE", "/session/anything"))
2091 .await;
2092 assert_eq!(unknown.status, 404);
2093 let wrong = adapter
2094 .handle(OpenCodeRequest::new("GET", "/session/ses_wrong"))
2095 .await;
2096 assert_eq!(wrong.status, 404);
2097 let untraced_query = adapter
2098 .handle(OpenCodeRequest::new(
2099 "GET",
2100 format!("/session/{}/message?limit=999", adapter.session_id()),
2101 ))
2102 .await;
2103 assert_eq!(untraced_query.status, 404);
2104 for target in [
2105 "/agent?x=1",
2106 "/config?x=1",
2107 "/event?x=1",
2108 "/session?unexpected=start=1",
2109 "/session?start=1&extra=1",
2110 "/session?start=not-a-number",
2111 ] {
2112 assert_eq!(
2113 adapter
2114 .handle(OpenCodeRequest::new("GET", target))
2115 .await
2116 .status,
2117 404,
2118 "accepted untraced target {target}"
2119 );
2120 }
2121 assert_eq!(
2122 adapter
2123 .handle(
2124 OpenCodeRequest::new("POST", "/session").with_body(json!({"unexpected": true}))
2125 )
2126 .await
2127 .status,
2128 404
2129 );
2130 }
2131
2132 #[tokio::test]
2133 async fn event_route_returns_atomic_sdk_attachment_not_a_second_runtime() {
2134 let runtime = FixtureRuntime::new();
2135 let adapter = OpenCodeAdapter::new(runtime, "runtime-1", "/workspace");
2136 let response = adapter.handle(OpenCodeRequest::new("GET", "/event")).await;
2137 assert_eq!(response.status, 200);
2138 match response.body {
2139 ResponseBody::EventStream(attachment) => {
2140 assert_eq!(attachment.descriptor.session_id, "runtime-1");
2141 assert_eq!(attachment.history.len(), 3);
2142 }
2143 ResponseBody::Json(_) => panic!("expected event stream"),
2144 }
2145 }
2146
2147 #[tokio::test]
2148 async fn traced_abort_is_descriptor_gated_and_invokes_sdk_interrupt() {
2149 let runtime = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
2150 let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
2151 let response = adapter
2152 .handle(OpenCodeRequest::new(
2153 "POST",
2154 format!("/session/{}/abort", adapter.session_id()),
2155 ))
2156 .await;
2157 assert_eq!(json_body(response), json!(true));
2158 assert_eq!(runtime.interrupts.load(Ordering::SeqCst), 1);
2159
2160 let mut unavailable = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
2161 Arc::get_mut(&mut unavailable)
2162 .expect("unshared fixture")
2163 .descriptor
2164 .actions
2165 .interrupt = false;
2166 let adapter = OpenCodeAdapter::new(unavailable.clone(), "runtime-2", "/workspace");
2167 let denied = adapter
2168 .handle(OpenCodeRequest::new(
2169 "POST",
2170 format!("/session/{}/abort", adapter.session_id()),
2171 ))
2172 .await;
2173 assert_eq!(denied.status, 409);
2174 assert_eq!(unavailable.interrupts.load(Ordering::SeqCst), 0);
2175 }
2176
2177 #[test]
2178 fn runtime_attachment_paths_are_cross_platform_and_workspace_bounded() {
2179 assert_eq!(
2180 logical_runtime_path("/srv/project", "assets/pixel.png").unwrap(),
2181 "/srv/project/assets/pixel.png"
2182 );
2183 assert_eq!(
2184 logical_runtime_path(
2185 "C:\\runtime\\project",
2186 "C:\\runtime\\project\\assets\\pixel.png"
2187 )
2188 .unwrap(),
2189 "C:/runtime/project/assets/pixel.png"
2190 );
2191 assert!(logical_runtime_path("/srv/project", "/Users/client/private.png").is_err());
2192 assert!(
2193 logical_runtime_path("C:\\runtime\\project", "C:\\Users\\client\\private.png").is_err()
2194 );
2195 assert!(logical_runtime_path("/srv/project", "../../private.png").is_err());
2196 }
2197
2198 #[test]
2199 fn every_attachment_uri_branch_is_runtime_resolved_and_fail_closed() {
2200 let workspace = std::env::temp_dir().join(format!(
2201 "supercode-opencode-attachment-branches-{}",
2202 std::process::id()
2203 ));
2204 let outside = std::env::temp_dir().join(format!(
2205 "supercode-opencode-attachment-outside-{}",
2206 std::process::id()
2207 ));
2208 std::fs::remove_dir_all(&workspace).ok();
2209 std::fs::remove_file(&outside).ok();
2210 std::fs::create_dir_all(&workspace).unwrap();
2211 std::fs::write(workspace.join("note.txt"), "runtime note").unwrap();
2212 std::fs::write(&outside, "outside").unwrap();
2213
2214 let resolve = |part: Value| {
2215 resolve_file_part(part.as_object().expect("file-part object"), &workspace)
2216 };
2217 for (label, part, expected_text, expected_image) in [
2218 (
2219 "data image",
2220 json!({"type":"file","mime":"image/png","filename":"pixel.png","url":"data:image/png;base64,cG5n"}),
2221 None,
2222 Some("data:image/png;base64,cG5n"),
2223 ),
2224 (
2225 "data text",
2226 json!({"type":"file","mime":"text/plain","filename":"note.txt","url":"data:text/plain;base64,aGVsbG8="}),
2227 Some("[file: note.txt]\nhello"),
2228 None,
2229 ),
2230 (
2231 "http image",
2232 json!({"type":"file","mime":"image/png","filename":"remote.png","url":"http://example.test/remote.png"}),
2233 None,
2234 Some("http://example.test/remote.png"),
2235 ),
2236 (
2237 "https image",
2238 json!({"type":"file","mime":"image/webp","filename":"remote.webp","url":"https://example.test/remote.webp"}),
2239 None,
2240 Some("https://example.test/remote.webp"),
2241 ),
2242 (
2243 "runtime text",
2244 json!({"type":"file","mime":"text/plain","filename":"note.txt","url":Url::from_file_path(workspace.join("note.txt")).unwrap().to_string()}),
2245 Some("[file: note.txt]\nruntime note"),
2246 None,
2247 ),
2248 ] {
2249 let resolved = resolve(part).unwrap_or_else(|error| panic!("{label}: {error}"));
2250 assert_eq!(resolved.text_block.as_deref(), expected_text, "{label}");
2251 assert_eq!(resolved.image_url.as_deref(), expected_image, "{label}");
2252 }
2253
2254 for (label, part, expected) in [
2255 (
2256 "mismatched data mime",
2257 json!({"type":"file","mime":"image/png","url":"data:image/jpeg;base64,cG5n"}),
2258 "data URI mime must match",
2259 ),
2260 (
2261 "invalid image base64",
2262 json!({"type":"file","mime":"image/png","url":"data:image/png;base64,%%%"}),
2263 "invalid base64",
2264 ),
2265 (
2266 "non-UTF8 data text",
2267 json!({"type":"file","mime":"text/plain","url":"data:text/plain;base64,/w=="}),
2268 "must be UTF-8 text",
2269 ),
2270 (
2271 "remote text",
2272 json!({"type":"file","mime":"text/plain","url":"https://example.test/private.txt"}),
2273 "remote non-image attachments are not fetched",
2274 ),
2275 (
2276 "unsupported scheme",
2277 json!({"type":"file","mime":"image/png","url":"ftp://example.test/pixel.png"}),
2278 "URL scheme is not supported",
2279 ),
2280 (
2281 "invalid URL",
2282 json!({"type":"file","mime":"image/png","url":"not a URL"}),
2283 "must be a data:, file:, or https: URL",
2284 ),
2285 (
2286 "outside runtime path",
2287 json!({"type":"file","mime":"text/plain","url":Url::from_file_path(&outside).unwrap().to_string()}),
2288 "escapes the SDK runtime workspace",
2289 ),
2290 ] {
2291 let error = resolve(part)
2292 .err()
2293 .unwrap_or_else(|| panic!("{label} was accepted"));
2294 assert!(error.to_string().contains(expected), "{label}: {error}");
2295 }
2296
2297 #[cfg(unix)]
2298 {
2299 std::os::unix::fs::symlink(&outside, workspace.join("escaped-link.txt")).unwrap();
2300 let error = resolve(json!({
2301 "type":"file", "mime":"text/plain",
2302 "url":Url::from_file_path(workspace.join("escaped-link.txt")).unwrap().to_string()
2303 }))
2304 .err()
2305 .expect("symlink escape must be rejected");
2306 assert!(error
2307 .to_string()
2308 .contains("escapes the SDK runtime workspace"));
2309 }
2310
2311 std::fs::remove_dir_all(workspace).unwrap();
2312 std::fs::remove_file(outside).unwrap();
2313 }
2314
2315 #[tokio::test]
2316 async fn file_parts_are_resolved_on_the_sdk_runtime_without_client_path_expansion() {
2317 let workspace = std::env::temp_dir().join(format!(
2318 "supercode-opencode-runtime-attachment-{}",
2319 std::process::id()
2320 ));
2321 let outside = std::env::temp_dir().join(format!(
2322 "supercode-opencode-client-attachment-{}",
2323 std::process::id()
2324 ));
2325 let _ = std::fs::remove_dir_all(&workspace);
2326 std::fs::create_dir_all(&workspace).unwrap();
2327 std::fs::write(workspace.join("pixel.png"), b"png").unwrap();
2328 std::fs::write(workspace.join("note.txt"), "runtime note").unwrap();
2329 std::fs::write(&outside, b"client secret").unwrap();
2330
2331 let runtime = FixtureRuntime::new();
2332 let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", &workspace);
2333 let session_id = adapter.session_id().to_string();
2334 let file_url = Url::from_file_path(workspace.join("pixel.png"))
2335 .unwrap()
2336 .to_string();
2337 let accepted = adapter
2338 .handle(
2339 OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
2340 json!({
2341 "messageID": "msg_attachment", "agent": "build",
2342 "model": {"providerID": "openrouter", "modelID": "glm-5.2"},
2343 "parts": [
2344 {"id": "prt_text", "type": "text", "text": "inspect"},
2345 {"id": "prt_more", "type": "text", "text": "carefully"},
2346 {"id": "prt_note", "type": "file", "mime": "text/plain", "filename": "note.txt", "url": Url::from_file_path(workspace.join("note.txt")).unwrap().to_string()},
2347 {"id": "prt_image", "type": "file", "mime": "image/png", "filename": "pixel.png", "url": file_url}
2348 ]
2349 }),
2350 ),
2351 )
2352 .await;
2353 assert_eq!(accepted.status, 200);
2354 assert_eq!(
2355 *runtime
2356 .submissions
2357 .lock()
2358 .unwrap_or_else(std::sync::PoisonError::into_inner),
2359 vec!["inspect\ncarefully\n[file: note.txt]\nruntime note".to_string()]
2360 );
2361 assert_eq!(
2362 *runtime
2363 .image_submissions
2364 .lock()
2365 .unwrap_or_else(std::sync::PoisonError::into_inner),
2366 vec![vec!["data:image/png;base64,cG5n".to_string()]]
2367 );
2368
2369 let client_url = Url::from_file_path(&outside).unwrap().to_string();
2370 let denied = adapter
2371 .handle(
2372 OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
2373 json!({
2374 "messageID": "msg_client_path", "agent": "build",
2375 "model": {"providerID": "openrouter", "modelID": "glm-5.2"},
2376 "parts": [{"id": "prt_client", "type": "file", "mime": "text/plain", "url": client_url}]
2377 }),
2378 ),
2379 )
2380 .await;
2381 assert_eq!(denied.status, 400);
2382 assert_eq!(
2383 runtime
2384 .image_submissions
2385 .lock()
2386 .unwrap_or_else(std::sync::PoisonError::into_inner)
2387 .len(),
2388 1
2389 );
2390
2391 std::fs::remove_dir_all(workspace).unwrap();
2392 std::fs::remove_file(outside).unwrap();
2393 }
2394
2395 #[tokio::test]
2396 async fn terminal_replay_reuses_the_assistant_already_present_in_history() {
2397 let runtime = FixtureRuntime::new();
2398 let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
2399 let attachment = runtime.attach(HISTORY_LIMIT).await.unwrap();
2400 let history = adapter.history_messages(&attachment);
2401 let existing_id = history.as_array().unwrap().last().unwrap()["info"]["id"].clone();
2402 let sequence = attachment.history_cursor + 1;
2403 let mut projection = adapter.event_projection(&attachment);
2404 let replay = projection.project(&SdkEvent {
2405 sequence,
2406 kind: "turn_succeeded".into(),
2407 payload: json!({"type": "turn_succeeded", "reply": "continued through GLM"}),
2408 });
2409 let completed = replay
2410 .iter()
2411 .find(|event| event["type"] == "message.updated")
2412 .expect("terminal replay completes the history message");
2413 assert_eq!(completed["properties"]["info"]["id"], existing_id);
2414 assert_eq!(
2415 replay
2416 .iter()
2417 .filter(|event| event["type"] == "message.updated")
2418 .count(),
2419 1
2420 );
2421 }
2422
2423 #[tokio::test]
2424 async fn simultaneous_event_attachments_share_live_message_identity() {
2425 let runtime = FixtureRuntime::new();
2426 let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
2427 let first = runtime.attach(HISTORY_LIMIT).await.unwrap();
2428 let second = runtime.attach(HISTORY_LIMIT).await.unwrap();
2429 let mut first = adapter.event_projection(&first);
2430 let mut second = adapter.event_projection(&second);
2431 let user = SdkEvent {
2432 sequence: 40,
2433 kind: "user_message".into(),
2434 payload: json!({"type": "user_message", "text": "same event"}),
2435 };
2436 let started = SdkEvent {
2437 sequence: 41,
2438 kind: "turn_started".into(),
2439 payload: json!({"type": "turn_started"}),
2440 };
2441 let first_user = first.project(&user);
2442 let second_user = second.project(&user);
2443 assert_eq!(
2444 first_user[0]["properties"]["info"]["id"],
2445 second_user[0]["properties"]["info"]["id"]
2446 );
2447 let first_assistant = first.project(&started);
2448 let second_assistant = second.project(&started);
2449 assert_eq!(
2450 first_assistant[1]["properties"]["info"]["id"],
2451 second_assistant[1]["properties"]["info"]["id"]
2452 );
2453 assert_eq!(
2454 first_assistant[2]["properties"]["part"]["id"],
2455 second_assistant[2]["properties"]["part"]["id"]
2456 );
2457 }
2458
2459 #[tokio::test]
2460 async fn traced_mutations_submit_and_answer_only_exact_typed_requests() {
2461 let runtime = FixtureRuntime::new();
2462 let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
2463 let session_id = adapter.session_id();
2464 let submitted = adapter
2465 .handle(
2466 OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
2467 json!({
2468 "messageID": "msg_stock", "agent": "build",
2469 "model": {"providerID": "openrouter", "modelID": "glm-5.2"},
2470 "parts": [{"id": "prt_stock", "type": "text", "text": "continue through GLM"}]
2471 }),
2472 ),
2473 )
2474 .await;
2475 assert_eq!(submitted.status, 200);
2476 assert_eq!(
2477 *runtime
2478 .submissions
2479 .lock()
2480 .unwrap_or_else(std::sync::PoisonError::into_inner),
2481 vec!["continue through GLM"]
2482 );
2483 let history = json_body(
2484 adapter
2485 .handle(OpenCodeRequest::new(
2486 "GET",
2487 format!("/session/{session_id}/message?limit=100"),
2488 ))
2489 .await,
2490 );
2491 let submitted_user = history
2492 .as_array()
2493 .unwrap()
2494 .iter()
2495 .find(|message| message["parts"][0]["text"] == "continue through GLM")
2496 .expect("submitted user message in history");
2497 assert_eq!(submitted_user["info"]["id"], "msg_stock");
2498 assert_eq!(submitted_user["parts"][0]["id"], "prt_stock");
2499
2500 let answered = adapter
2501 .handle(
2502 OpenCodeRequest::new("POST", "/permission/per_000000000000002a/reply")
2503 .with_body(json!({"reply": "once"})),
2504 )
2505 .await;
2506 assert_eq!(answered.status, 200);
2507 assert_eq!(
2508 *runtime
2509 .responses
2510 .lock()
2511 .unwrap_or_else(std::sync::PoisonError::into_inner),
2512 vec![FrontendResponse::Approval {
2513 request_id: 42,
2514 decision: FrontendApprovalDecision::Allow
2515 }]
2516 );
2517
2518 let extra = adapter
2519 .handle(
2520 OpenCodeRequest::new("POST", format!("/session/{session_id}/message")).with_body(
2521 json!({
2522 "messageID": "msg", "agent": "build", "model": {}, "parts": [],
2523 "untraced": true
2524 }),
2525 ),
2526 )
2527 .await;
2528 assert_eq!(extra.status, 400);
2529 assert_eq!(
2530 runtime
2531 .submissions
2532 .lock()
2533 .unwrap_or_else(std::sync::PoisonError::into_inner)
2534 .len(),
2535 1
2536 );
2537 for target in [
2538 format!("/session/{session_id}/message?untraced=1"),
2539 "/permission/per_000000000000002a/reply?untraced=1".into(),
2540 ] {
2541 assert_eq!(
2542 adapter
2543 .handle(
2544 OpenCodeRequest::new("POST", target).with_body(json!({"reply": "once"})),
2545 )
2546 .await
2547 .status,
2548 404
2549 );
2550 }
2551 }
2552
2553 #[tokio::test]
2554 async fn busy_steer_ignores_the_previous_turn_terminal_replay() {
2555 let runtime = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
2556 let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
2557 let response = adapter
2558 .handle(
2559 OpenCodeRequest::new("POST", format!("/session/{}/message", adapter.session_id()))
2560 .with_body(json!({
2561 "messageID": "msg_stock", "agent": "build",
2562 "model": {"providerID": "openrouter", "modelID": "glm-5.2"},
2563 "parts": [{"id": "prt_stock", "type": "text", "text": "change direction"}]
2564 })),
2565 )
2566 .await;
2567 assert_eq!(response.status, 200);
2568 let response = match response.body {
2569 ResponseBody::Json(body) => body,
2570 ResponseBody::EventStream(_) => panic!("expected steer response"),
2571 };
2572 assert_eq!(response["parts"][0]["text"], "steered:change direction");
2573 assert_ne!(response["parts"][0]["text"], "continued through GLM");
2574 assert!(runtime
2575 .submissions
2576 .lock()
2577 .unwrap_or_else(std::sync::PoisonError::into_inner)
2578 .is_empty());
2579 assert_eq!(
2580 *runtime
2581 .steers
2582 .lock()
2583 .unwrap_or_else(std::sync::PoisonError::into_inner),
2584 vec!["change direction"]
2585 );
2586 }
2587
2588 #[tokio::test]
2589 async fn busy_descriptor_race_propagates_atomic_steer_rejection() {
2590 let runtime = FixtureRuntime::with_turn_state(FrontendTurnState::Busy);
2591 runtime.steer_accepting.store(false, Ordering::SeqCst);
2592 let adapter = OpenCodeAdapter::new(runtime.clone(), "runtime-1", "/workspace");
2593 let response = adapter
2594 .handle(
2595 OpenCodeRequest::new("POST", format!("/session/{}/message", adapter.session_id()))
2596 .with_body(json!({
2597 "messageID": "msg_stock", "agent": "build",
2598 "model": {"providerID": "openrouter", "modelID": "glm-5.2"},
2599 "parts": [{"id": "prt_stock", "type": "text", "text": "too late"}]
2600 })),
2601 )
2602 .await;
2603 assert_eq!(response.status, 409);
2604 assert!(runtime
2605 .steers
2606 .lock()
2607 .unwrap_or_else(std::sync::PoisonError::into_inner)
2608 .is_empty());
2609 }
2610
2611 #[test]
2612 fn delayed_sse_reuses_identities_already_projected_by_history() {
2613 let descriptor = FixtureRuntime::new().descriptor.clone();
2614 let identities = Arc::new(Mutex::new(HistoryIdentityState::default()));
2615 let submitted = SubmittedMessage {
2616 prompt: "continue".into(),
2617 image_urls: Vec::new(),
2618 message_id: "msg_history_first".into(),
2619 part_id: "prt_history_first".into(),
2620 };
2621 let baseline = vec![
2622 ChatMessage::user("continue"),
2623 ChatMessage::assistant("prior identical prompt"),
2624 ];
2625 let history = vec![
2626 baseline[0].clone(),
2627 baseline[1].clone(),
2628 ChatMessage::user("continue"),
2629 ChatMessage::assistant("tool request"),
2630 ChatMessage::tool_result("call-1", "read", "tool result"),
2631 ChatMessage::assistant("history won the race"),
2632 ];
2633 let (history_user, history_assistant) = {
2634 let mut state = identities
2635 .lock()
2636 .unwrap_or_else(std::sync::PoisonError::into_inner);
2637 state.project(&baseline, 4);
2638 state.register_client_user(&submitted);
2639 let ordinals = state.project(&history, 10);
2640 (
2641 state.identity("ses_test", ordinals[2]),
2642 state.identity("ses_test", ordinals[5]),
2643 )
2644 };
2645
2646 let mut future = OpenCodeEventProjection::new(
2647 "ses_test".into(),
2648 PathBuf::from("/workspace"),
2649 &descriptor,
2650 identities.clone(),
2651 None,
2652 );
2653 let future_user = future.project(&SdkEvent {
2654 sequence: 11,
2655 kind: "user_message".into(),
2656 payload: json!({"type": "user_message", "text": "continue"}),
2657 });
2658 let future_assistant = future.project(&SdkEvent {
2659 sequence: 12,
2660 kind: "turn_started".into(),
2661 payload: json!({"type": "turn_started"}),
2662 });
2663 assert_ne!(future_user[0]["properties"]["info"]["id"], history_user.0);
2664 assert_ne!(
2665 future_assistant[1]["properties"]["info"]["id"],
2666 history_assistant.0
2667 );
2668
2669 let mut delayed = OpenCodeEventProjection::new(
2670 "ses_test".into(),
2671 PathBuf::from("/workspace"),
2672 &descriptor,
2673 identities,
2674 None,
2675 );
2676 let user = delayed.project(&SdkEvent {
2677 sequence: 6,
2678 kind: "user_message".into(),
2679 payload: json!({"type": "user_message", "text": "continue"}),
2680 });
2681 let assistant = delayed.project(&SdkEvent {
2682 sequence: 7,
2683 kind: "turn_started".into(),
2684 payload: json!({"type": "turn_started"}),
2685 });
2686
2687 assert_eq!(user[0]["properties"]["info"]["id"], history_user.0);
2688 assert_eq!(user[1]["properties"]["part"]["id"], history_user.1);
2689 assert_eq!(
2690 assistant[1]["properties"]["info"]["id"],
2691 history_assistant.0
2692 );
2693 assert_eq!(
2694 assistant[2]["properties"]["part"]["id"],
2695 history_assistant.1
2696 );
2697 }
2698
2699 #[test]
2700 fn delayed_sse_reuses_history_ids_for_another_frontend_or_scheduler_turn() {
2701 let descriptor = FixtureRuntime::new().descriptor.clone();
2702 let identities = Arc::new(Mutex::new(HistoryIdentityState::default()));
2703 let baseline = vec![
2704 ChatMessage::user("existing"),
2705 ChatMessage::assistant("existing reply"),
2706 ];
2707 let current = vec![
2708 baseline[0].clone(),
2709 baseline[1].clone(),
2710 ChatMessage::user("scheduled wakeup"),
2711 ChatMessage::assistant("scheduler reply"),
2712 ];
2713 let (history_user, history_assistant) = {
2714 let mut state = identities
2715 .lock()
2716 .unwrap_or_else(std::sync::PoisonError::into_inner);
2717 state.project(&baseline, 4);
2718 let ordinals = state.project(¤t, 8);
2719 (
2720 state.identity("ses_test", ordinals[2]),
2721 state.identity("ses_test", ordinals[3]),
2722 )
2723 };
2724 let mut delayed = OpenCodeEventProjection::new(
2725 "ses_test".into(),
2726 PathBuf::from("/workspace"),
2727 &descriptor,
2728 identities,
2729 None,
2730 );
2731 let user = delayed.project(&SdkEvent {
2732 sequence: 5,
2733 kind: "user_message".into(),
2734 payload: json!({"type": "user_message", "text": "scheduled wakeup"}),
2735 });
2736 let assistant = delayed.project(&SdkEvent {
2737 sequence: 6,
2738 kind: "turn_started".into(),
2739 payload: json!({"type": "turn_started"}),
2740 });
2741 assert_eq!(user[0]["properties"]["info"]["id"], history_user.0);
2742 assert_eq!(user[1]["properties"]["part"]["id"], history_user.1);
2743 assert_eq!(
2744 assistant[1]["properties"]["info"]["id"],
2745 history_assistant.0
2746 );
2747 assert_eq!(
2748 assistant[2]["properties"]["part"]["id"],
2749 history_assistant.1
2750 );
2751 }
2752
2753 #[test]
2754 fn older_attachment_projection_cannot_rewind_newer_identity_state() {
2755 let mut state = HistoryIdentityState::default();
2756 let old = vec![
2757 ChatMessage::user("old"),
2758 ChatMessage::assistant("old reply"),
2759 ];
2760 let newer = vec![
2761 old[0].clone(),
2762 old[1].clone(),
2763 ChatMessage::user("new"),
2764 ChatMessage::assistant("new reply"),
2765 ];
2766 let newer_ids = state.project(&newer, 8);
2767 assert_eq!(state.project(&old, 4), newer_ids[..2]);
2768
2769 let newest = vec![
2770 newer[0].clone(),
2771 newer[1].clone(),
2772 newer[2].clone(),
2773 newer[3].clone(),
2774 ChatMessage::user("newest"),
2775 ChatMessage::assistant("newest reply"),
2776 ];
2777 let newest_ids = state.project(&newest, 12);
2778 assert_eq!(&newest_ids[..4], &newer_ids);
2779 assert_eq!(state.project(&newer, 8), newer_ids);
2780 }
2781
2782 #[test]
2783 fn canonical_events_project_to_stock_delta_and_reversible_permission_ids() {
2784 let descriptor = FixtureRuntime::new().descriptor.clone();
2785 let identities = Arc::new(Mutex::new(HistoryIdentityState::default()));
2786 identities
2787 .lock()
2788 .unwrap_or_else(std::sync::PoisonError::into_inner)
2789 .register_client_user(&SubmittedMessage {
2790 prompt: "continue".into(),
2791 image_urls: Vec::new(),
2792 message_id: "msg_stock_client".into(),
2793 part_id: "prt_stock_client".into(),
2794 });
2795 let mut projection = OpenCodeEventProjection::new(
2796 "ses_test".into(),
2797 PathBuf::from("C:\\runtime\\worktree"),
2798 &descriptor,
2799 identities.clone(),
2800 None,
2801 );
2802 let user = projection.project(&SdkEvent {
2803 sequence: 6,
2804 kind: "user_message".into(),
2805 payload: json!({"type": "user_message", "text": "continue"}),
2806 });
2807 assert_eq!(user[0]["properties"]["info"]["id"], "msg_stock_client");
2808 assert_eq!(user[1]["properties"]["part"]["id"], "prt_stock_client");
2809 let started = projection.project(&SdkEvent {
2810 sequence: 7,
2811 kind: "turn_started".into(),
2812 payload: json!({"type": "turn_started"}),
2813 });
2814 assert_eq!(started[0]["type"], "session.status");
2815 assert_eq!(started[1]["type"], "message.updated");
2816 assert_eq!(started[2]["type"], "message.part.updated");
2817 let delta = projection.project(&SdkEvent {
2818 sequence: 8,
2819 kind: "text_delta".into(),
2820 payload: json!({"type": "text_delta", "text": "hello"}),
2821 });
2822 assert_eq!(delta[0]["type"], "message.part.delta");
2823 assert_eq!(delta[0]["properties"]["delta"], "hello");
2824 assert_eq!(
2825 delta[0]["properties"]["messageID"],
2826 started[1]["properties"]["info"]["id"]
2827 );
2828 assert_eq!(
2829 delta[0]["properties"]["partID"],
2830 started[2]["properties"]["part"]["id"]
2831 );
2832 projection.project(&SdkEvent {
2833 sequence: 9,
2834 kind: "turn_succeeded".into(),
2835 payload: json!({"type": "turn_succeeded", "reply": "hello"}),
2836 });
2837
2838 let history = vec![
2839 ChatMessage::user("continue"),
2840 ChatMessage::assistant("hello"),
2841 ];
2842 let mut identity_state = identities
2843 .lock()
2844 .unwrap_or_else(std::sync::PoisonError::into_inner);
2845 let ordinals = identity_state.project(&history, 9);
2846 let projected_identities = ordinals
2847 .iter()
2848 .map(|ordinal| identity_state.identity("ses_test", *ordinal))
2849 .collect::<Vec<_>>();
2850 drop(identity_state);
2851 let reconnected = history_messages(
2852 &history,
2853 &ordinals,
2854 &projected_identities,
2855 "ses_test",
2856 Path::new("C:\\runtime\\worktree"),
2857 &descriptor,
2858 );
2859 assert_eq!(
2860 reconnected[0]["info"]["id"],
2861 user[0]["properties"]["info"]["id"]
2862 );
2863 assert_eq!(
2864 reconnected[1]["info"]["id"],
2865 started[1]["properties"]["info"]["id"]
2866 );
2867 assert_eq!(
2868 reconnected[1]["parts"][0]["id"],
2869 started[2]["properties"]["part"]["id"]
2870 );
2871 assert_eq!(reconnected[1]["info"]["providerID"], "openrouter");
2872 assert_eq!(reconnected[1]["info"]["modelID"], "glm-5.2");
2873
2874 let permission = projection.project(&SdkEvent {
2875 sequence: 10,
2876 kind: "request".into(),
2877 payload: json!({"type": "request", "request": {
2878 "id": 42, "kind": "approval",
2879 "payload": {"tool": "edit", "subject": "C:\\runtime\\worktree\\probe.txt"}
2880 }}),
2881 });
2882 assert_eq!(permission[0]["properties"]["id"], "per_000000000000002a");
2883 assert_eq!(
2884 permission_request_id("/permission/per_000000000000002a/reply"),
2885 Some(42)
2886 );
2887 assert_eq!(
2888 path_text(Path::new("C:\\runtime\\worktree")),
2889 "C:/runtime/worktree"
2890 );
2891 }
2892
2893 #[test]
2894 fn client_projection_types_never_enter_harness_source() {
2895 let harness = std::fs::read_to_string(
2896 Path::new(env!("CARGO_MANIFEST_DIR")).join("../harness/src/frontend.rs"),
2897 )
2898 .unwrap();
2899 for forbidden in ["OpenCodeAdapter", "opencode_http", "OPENCODE_CLI_VERSION"] {
2900 assert!(!harness.contains(forbidden), "harness contains {forbidden}");
2901 }
2902 }
2903}