Skip to main content

supercode_opencode_frontend/
opencode_http_v1_2_15.rs

1//! Compatibility namespace for the HTTP protocol spoken by OpenCode 1.2.15.
2
3use 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
15/// Exact protocol namespace pinned by the SUP-66 stock-client corpus.
16pub const PROTOCOL_NAMESPACE: &str = "opencode_http/v1_2_15";
17/// Only stock client release accepted by this compatibility namespace.
18pub const OPENCODE_CLI_VERSION: &str = "1.2.15";
19/// Maximum historical messages projected into a stock-client attachment.
20pub 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/// One transport-neutral HTTP request before a listener adds wire framing.
49#[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
71/// Adapter output. Event streams remain typed until the authenticated HTTP
72/// host binds them to one connection.
73pub enum ResponseBody {
74    Json(Value),
75    EventStream(Box<FrontendAttachment>),
76}
77
78/// One compatible HTTP response.
79pub 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/// Named adapter failures. The HTTP host preserves these stable compatible
101/// status codes and messages.
102#[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
141/// Version-pinned projection around one existing SDK-owned runtime.
142pub 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    /// Project one traced request without giving client protocol types to the
214    /// SDK. Unobserved routes fail closed.
215    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        // Attach before accepting the mutation so the response can be proven
370        // to belong to this request. Returning an arbitrary latest assistant
371        // makes a successful steer look complete before the active turn has
372        // consumed it and invites duplicate retries.
373        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            // A completed turn commits its history cursor before publishing
556            // its terminal event. A fresh Busy attachment may still replay
557            // the preceding turn's terminal, so accept only a terminal that
558            // follows a history boundary newer than the pre-steer snapshot.
559            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
963/// Stateful, client-protocol-only projection of canonical SDK events. It
964/// retains no agent or session state; only the currently rendered part IDs.
965pub 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    /// Preserve identities through the only canonical history transition:
1267    /// append, optionally dropping a bounded prefix. The largest old-suffix /
1268    /// new-prefix overlap also handles repeated messages deterministically.
1269    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            // Concurrent attachments may finish projection out of order.
1282            // Derive an older window against the monotonic current state,
1283            // but never replace that state or rewind identity allocation.
1284            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 &current[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                // HTTP completion can project canonical history before a
1311                // delayed SSE consumer handles the already-queued live
1312                // events. This applies to every canonical owner (another
1313                // frontend and the scheduler included), not only OpenCode
1314                // POSTs carrying explicit stock ids.
1315                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            // One stock assistant message spans a canonical tool loop, so its
1363            // identity is the final assistant history message for the turn.
1364            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(&current, 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}