Skip to main content

supercode_harness/runtime/
adapters.rs

1//! Harness-specific live-runtime adapters built on the primitive contracts.
2
3use std::collections::{BTreeMap, VecDeque};
4use std::net::TcpListener;
5use std::path::{Path, PathBuf};
6use std::process::Stdio;
7use std::sync::Arc;
8use std::time::{Duration, SystemTime, UNIX_EPOCH};
9
10use async_trait::async_trait;
11use futures::StreamExt;
12use serde_json::{json, Value};
13use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
14use tokio::process::{Child, ChildStdin, Command};
15use tokio::sync::{mpsc, Mutex};
16
17use super::{
18    BearerToken, HarnessEvent, JsonLineClient, McpServerLaunch, RuntimeAttachRequest,
19    RuntimeBackend, RuntimeCapabilities, RuntimeConnection, RuntimeEndpoint, RuntimeHandle,
20    RuntimeInput, RuntimeLaunch, RuntimeStartRequest,
21};
22use crate::{Error, HarnessId, Result};
23
24/// Pi live-runtime backend using `pi --mode rpc` JSONL.
25#[derive(Debug, Clone)]
26pub struct PiRuntimeBackend {
27    launch: RuntimeLaunch,
28}
29
30impl Default for PiRuntimeBackend {
31    fn default() -> Self {
32        Self::new()
33    }
34}
35
36impl PiRuntimeBackend {
37    /// Use `pi --mode rpc` from `PATH`.
38    pub fn new() -> Self {
39        Self {
40            launch: RuntimeLaunch {
41                program: "pi".into(),
42                arguments: vec!["--mode".into(), "rpc".into()],
43                env: BTreeMap::new(),
44            },
45        }
46    }
47
48    /// Use an explicit Pi RPC command prefix.
49    pub fn with_launch(launch: RuntimeLaunch) -> Self {
50        Self { launch }
51    }
52
53    async fn open(
54        &self,
55        cwd: &Path,
56        runtime_id: String,
57        launch: Option<RuntimeLaunch>,
58        mcp_servers: &[McpServerLaunch],
59        resume: bool,
60    ) -> Result<Box<dyn RuntimeConnection>> {
61        let mut launch = launch.unwrap_or_else(|| self.launch.clone());
62        // Pi has no MCP client of its own; `pi-mcp-adapter` is the extension
63        // that gives it one, loaded for this process from npm by Pi's own
64        // package door and pointed at the servers written for this runtime,
65        // as Claude Code's backend points `--mcp-config` at the same file shape.
66        if !mcp_servers.is_empty() {
67            let file = mcp_config_file("pi", &runtime_id, mcp_servers).await?;
68            launch.arguments.extend([
69                "--extension".into(),
70                std::env::var(PI_MCP_ADAPTER_ENV)
71                    .ok()
72                    .filter(|value| !value.is_empty())
73                    .unwrap_or_else(|| PI_MCP_ADAPTER.into()),
74                "--mcp-config".into(),
75                file.to_string_lossy().into_owned(),
76            ]);
77        }
78        if resume {
79            launch
80                .arguments
81                .extend(["--session".into(), runtime_id.clone()]);
82        } else {
83            launch
84                .arguments
85                .extend(["--session-id".into(), runtime_id.clone()]);
86        }
87        let transport = RawLineTransport::spawn(&launch, Some(cwd), "pi-rpc-jsonl").await?;
88        let handle = RuntimeHandle {
89            harness: HarnessId::from(HarnessId::PI),
90            runtime_id,
91            endpoint: transport.endpoint.clone(),
92        };
93        Ok(Box::new(PiRuntimeConnection {
94            handle,
95            transport,
96            next_request: 1,
97        }))
98    }
99}
100
101#[async_trait]
102impl RuntimeBackend for PiRuntimeBackend {
103    fn harness(&self) -> HarnessId {
104        HarnessId::from(HarnessId::PI)
105    }
106
107    fn capabilities(&self) -> RuntimeCapabilities {
108        RuntimeCapabilities {
109            start_session: true,
110            resume_session: true,
111            attach_existing_process: false,
112            send_input: true,
113            stream_events: true,
114            interrupt: true,
115            steer: false,
116            respond_to_requests: true,
117        }
118    }
119
120    async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
121        self.open(
122            &request.cwd,
123            generated_session_id(),
124            request.launch,
125            &request.mcp_servers,
126            false,
127        )
128        .await
129    }
130
131    async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
132        let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
133        self.open(
134            &cwd,
135            request.runtime_id,
136            request.launch,
137            &request.mcp_servers,
138            true,
139        )
140        .await
141    }
142}
143
144/// The Pi extension that mounts MCP servers, pinned: Pi fetches an `npm:`
145/// extension into its own package directory on first use.
146const PI_MCP_ADAPTER: &str = "npm:pi-mcp-adapter@2.37.0";
147/// Where the adapter already is, for a Pi that cannot fetch one: a path or a
148/// Pi package source in place of `PI_MCP_ADAPTER` (a browser tab's Pi pack
149/// carries it, and the tab's npm resolves no `npm:` source).
150const PI_MCP_ADAPTER_ENV: &str = "SUPERCODE_PI_MCP_ADAPTER";
151
152struct PiRuntimeConnection {
153    handle: RuntimeHandle,
154    transport: RawLineTransport,
155    next_request: u64,
156}
157
158#[async_trait]
159impl RuntimeConnection for PiRuntimeConnection {
160    fn handle(&self) -> &RuntimeHandle {
161        &self.handle
162    }
163
164    async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
165        if !input.image_urls.is_empty() {
166            return Err(Error::Other(
167                "Pi RPC image input is not verified by the installed protocol contract".into(),
168            ));
169        }
170        let id = format!("supercode-{}", self.next_request);
171        self.next_request += 1;
172        self.transport
173            .write(json!({"id": id, "type": "prompt", "message": input.text}))
174            .await?;
175        Ok(Some(id))
176    }
177
178    async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
179        raw_next_event(&mut self.transport.receiver).await
180    }
181
182    async fn interrupt(&mut self) -> Result<()> {
183        self.transport.write(json!({"type": "abort"})).await
184    }
185
186    async fn respond(&mut self, request_id: Value, mut response: Value) -> Result<()> {
187        if let Value::Object(object) = &mut response {
188            object.entry("id").or_insert(request_id);
189            self.transport.write(response).await
190        } else {
191            self.transport
192                .write(json!({"id": request_id, "response": response}))
193                .await
194        }
195    }
196
197    async fn close(&mut self) -> Result<()> {
198        self.transport.close().await
199    }
200}
201
202/// Claude Code live-runtime backend using bidirectional stream-json print
203/// mode with supercode registered as the CLI's permission handler.
204///
205/// It can create/resume sessions, cancel the running turn through the
206/// stream-json control channel, and answer the permission requests the CLI
207/// raises while a tool call is blocked.
208#[derive(Debug, Clone)]
209pub struct ClaudeCodeRuntimeBackend {
210    launch: RuntimeLaunch,
211    permission_timeout: Duration,
212}
213
214impl Default for ClaudeCodeRuntimeBackend {
215    fn default() -> Self {
216        Self::new()
217    }
218}
219
220impl ClaudeCodeRuntimeBackend {
221    /// Use `claude` from `PATH` in bidirectional stream-json mode, with
222    /// supercode registered as the permission handler.
223    pub fn new() -> Self {
224        Self {
225            launch: RuntimeLaunch {
226                program: "claude".into(),
227                arguments: vec![
228                    "--print".into(),
229                    "--input-format".into(),
230                    "stream-json".into(),
231                    "--output-format".into(),
232                    "stream-json".into(),
233                    "--verbose".into(),
234                    // Register this adapter as the permission handler. Measured
235                    // against claude 2.1.258: without it a tool call that needs
236                    // an answer is refused outright with
237                    // `{"type":"system","subtype":"permission_denied"}`; with
238                    // it the CLI raises a `can_use_tool` control request on
239                    // stdout and blocks the turn until a `control_response`
240                    // arrives. See
241                    // `docs/interop/research/orc2-claude-respond-receipt-2026-09-04.json`.
242                    "--permission-prompt-tool".into(),
243                    "stdio".into(),
244                ],
245                env: BTreeMap::new(),
246            },
247            permission_timeout: CLAUDE_PERMISSION_RESPONSE_TIMEOUT,
248        }
249    }
250
251    /// The command prefix this backend launches: what the registry publishes
252    /// as `default_launch`, so a caller that hands the published launch back
253    /// (the orchestrator does) gets stream-json mode and not the interactive
254    /// TUI blocked on a pipe.
255    pub fn launch(&self) -> &RuntimeLaunch {
256        &self.launch
257    }
258
259    /// Use an explicit Claude Code stream-json command prefix.
260    pub fn with_launch(launch: RuntimeLaunch) -> Self {
261        Self {
262            launch,
263            permission_timeout: CLAUDE_PERMISSION_RESPONSE_TIMEOUT,
264        }
265    }
266
267    /// Override how long an unanswered permission request is held before this
268    /// adapter denies it on the caller's behalf.
269    ///
270    /// See [`CLAUDE_PERMISSION_RESPONSE_TIMEOUT`] for the default and what it
271    /// was measured against.
272    pub fn with_permission_timeout(mut self, timeout: Duration) -> Self {
273        self.permission_timeout = timeout;
274        self
275    }
276
277    async fn open(
278        &self,
279        cwd: &Path,
280        runtime_id: String,
281        launch: Option<RuntimeLaunch>,
282        mcp_servers: &[McpServerLaunch],
283        resume: bool,
284    ) -> Result<Box<dyn RuntimeConnection>> {
285        // The prefix WITHOUT session arguments is kept: it is what reopens
286        // this same session later, and appending `--session-id` to a launch
287        // that already carries one is what the CLI refuses.
288        let mut prefix = launch.unwrap_or_else(|| self.launch.clone());
289        // MCP servers ride the prefix, not the session: Claude Code mounts them
290        // from `--mcp-config <file>` on every process, so the file written for
291        // this runtime is what every reopen of the session hands back too.
292        if !mcp_servers.is_empty() {
293            let file = mcp_config_file("claude", &runtime_id, mcp_servers).await?;
294            prefix
295                .arguments
296                .extend(["--mcp-config".into(), file.to_string_lossy().into_owned()]);
297        }
298        let mut launch = prefix.clone();
299        launch.arguments.extend(if resume {
300            vec!["--resume".into(), runtime_id.clone()]
301        } else {
302            vec!["--session-id".into(), runtime_id.clone()]
303        });
304        let transport = RawLineTransport::spawn(&launch, Some(cwd), "claude-stream-json").await?;
305        Ok(Box::new(ClaudeRuntimeConnection {
306            handle: RuntimeHandle {
307                harness: HarnessId::from(HarnessId::CLAUDE_CODE),
308                runtime_id,
309                endpoint: transport.endpoint.clone(),
310            },
311            transport,
312            prefix,
313            cwd: cwd.to_path_buf(),
314            spoke: false,
315            buffered_events: VecDeque::new(),
316            next_control_request: 1,
317            control_timeout: CLAUDE_CONTROL_RESPONSE_TIMEOUT,
318            pending_permissions: Vec::new(),
319            permission_timeout: self.permission_timeout,
320            mounted_mcp_servers: mcp_servers
321                .iter()
322                .map(|server| server.name.clone())
323                .collect(),
324        }))
325    }
326}
327
328/// How long `interrupt` waits for the CLI's matching `control_response` before
329/// returning a structured error instead of hanging the caller.
330///
331/// Measured against claude 2.1.224: an interrupt issued while a turn is in
332/// flight is acknowledged in ~1 ms, but one issued during process startup —
333/// before the CLI has emitted `system/init` — is queued behind session-start
334/// hooks and took 1.15 s to acknowledge on a warm box. The bound is set well
335/// above the slow case so a legitimately busy startup is never reported as a
336/// protocol failure.
337const CLAUDE_CONTROL_RESPONSE_TIMEOUT: Duration = Duration::from_secs(10);
338
339/// How long an unanswered `can_use_tool` request is held before this adapter
340/// denies it on the caller's behalf.
341///
342/// Measured against claude 2.1.258: a `can_use_tool` control request the host
343/// never answers blocks the turn indefinitely — the probe watched one sit for
344/// 20 s with no further frame and no tool result, and the CLI has no bound of
345/// its own (receipt step `unanswered_request_blocks_the_turn` in
346/// `docs/interop/research/orc2-claude-respond-receipt-2026-09-04.json`). The
347/// bound therefore has to live here. Five minutes is long enough for an
348/// operator to read `supercode approvals list` and answer, and short enough
349/// that a forgotten prompt does not wedge a driven session forever. Override
350/// per backend with [`ClaudeCodeRuntimeBackend::with_permission_timeout`].
351pub const CLAUDE_PERMISSION_RESPONSE_TIMEOUT: Duration = Duration::from_secs(300);
352
353/// The `behavior` values Claude Code's permission handler protocol accepts.
354///
355/// Measured against claude 2.1.258: anything else — including a `deny` with no
356/// `message` — is refused by the CLI with "The canUseTool callback returned an
357/// invalid permission result. Expected {behavior: 'allow', updatedInput?:
358/// object} or {behavior: 'deny', message: string}."
359const CLAUDE_PERMISSION_BEHAVIORS: [&str; 2] = ["allow", "deny"];
360
361/// What this adapter says when it denies a request nobody answered in time.
362const CLAUDE_PERMISSION_TIMEOUT_MESSAGE: &str =
363    "Volter Harness denied this permission request: no answer arrived before the adapter's \
364     permission timeout elapsed";
365
366#[async_trait]
367impl RuntimeBackend for ClaudeCodeRuntimeBackend {
368    fn harness(&self) -> HarnessId {
369        HarnessId::from(HarnessId::CLAUDE_CODE)
370    }
371
372    fn capabilities(&self) -> RuntimeCapabilities {
373        RuntimeCapabilities {
374            start_session: true,
375            resume_session: true,
376            attach_existing_process: false,
377            send_input: true,
378            stream_events: true,
379            interrupt: true,
380            steer: true,
381            respond_to_requests: true,
382        }
383    }
384
385    async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
386        let session_id = generated_session_id();
387        // The agent the launch named is recorded in the session's config home before its transcript exists, so
388        // whoever discovers the session finds its agent with it (launch_agent.rs).
389        let env = request
390            .launch
391            .as_ref()
392            .map(|launch| launch.env.clone())
393            .unwrap_or_default();
394        let home = env.get("CLAUDE_CONFIG_DIR").map(PathBuf::from).or_else(|| {
395            crate::HarnessHomes::default()
396                .claude_code
397                .parent()
398                .map(Path::to_path_buf)
399        });
400        if let Some(home) = home {
401            if let Err(error) = crate::launch_agent::record(&home, &session_id, &env) {
402                tracing::warn!(
403                    "could not record the launch agent of Claude Code session {session_id}: {error}"
404                );
405            }
406        }
407        // A model the caller named is Claude Code's own `--model`; it rides the launch, so the session
408        // reopens on the same model.
409        let launch = match &request.model {
410            Some(model) => {
411                let mut launch = request.launch.unwrap_or_else(|| self.launch.clone());
412                launch.arguments.extend(["--model".into(), model.clone()]);
413                Some(launch)
414            }
415            None => request.launch,
416        };
417        self.open(
418            &request.cwd,
419            session_id,
420            launch,
421            &request.mcp_servers,
422            false,
423        )
424        .await
425    }
426
427    async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
428        let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
429        self.open(
430            &cwd,
431            request.runtime_id,
432            request.launch,
433            &request.mcp_servers,
434            true,
435        )
436        .await
437    }
438}
439
440/// Claude Code's own MCP configuration shape for the servers a start request
441/// mounts, written once per runtime under the OS temp directory (mode 0600 on
442/// Unix: a server's env may carry a token meant for that server alone).
443async fn mcp_config_file(
444    harness: &str,
445    runtime_id: &str,
446    servers: &[McpServerLaunch],
447) -> Result<PathBuf> {
448    let mut entries = serde_json::Map::new();
449    for server in servers {
450        entries.insert(
451            server.name.clone(),
452            json!({
453                "type": "stdio",
454                "command": server.command,
455                "args": server.arguments,
456                "env": server.env,
457            }),
458        );
459    }
460    let safe: String = runtime_id
461        .chars()
462        .filter(|c| c.is_ascii_alphanumeric() || *c == '-' || *c == '_')
463        .collect();
464    let path = std::env::temp_dir().join(format!("supercode-{harness}-mcp-{safe}.json"));
465    let body = serde_json::to_vec_pretty(&json!({ "mcpServers": entries })).map_err(|error| {
466        Error::Other(format!(
467            "{harness} mcp config could not be encoded: {error}"
468        ))
469    })?;
470    tokio::fs::write(&path, body).await.map_err(|error| {
471        Error::Other(format!(
472            "{harness} mcp config could not be written: {error}"
473        ))
474    })?;
475    #[cfg(unix)]
476    {
477        use std::os::unix::fs::PermissionsExt;
478        tokio::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600))
479            .await
480            .map_err(|error| {
481                Error::Other(format!(
482                    "{harness} mcp config could not be protected: {error}"
483                ))
484            })?;
485    }
486    Ok(path)
487}
488
489struct ClaudeRuntimeConnection {
490    handle: RuntimeHandle,
491    transport: RawLineTransport,
492    /// The launch prefix this session was opened with, carrying no session
493    /// argument of its own: half of what reopens the session after its
494    /// process has exited.
495    prefix: RuntimeLaunch,
496    /// The workspace the session belongs to. `--resume` resolves a session id
497    /// against the project the process runs in, so a reopen has to run in the
498    /// same directory the first process did.
499    cwd: PathBuf,
500    /// Whether the process now behind this connection has produced a frame.
501    /// A process that exited without ever speaking never opened a session, so
502    /// its close is a launch failure and is reported; one that spoke left a
503    /// session `--resume` reopens.
504    spoke: bool,
505    /// Native events read off the transport while `interrupt` was waiting for
506    /// its `control_response`. They are handed to `next_event` in arrival order
507    /// so cancelling a turn never costs the consumer an event.
508    buffered_events: VecDeque<Value>,
509    next_control_request: u64,
510    control_timeout: Duration,
511    /// Permission requests the CLI raised and is still blocked on, oldest
512    /// first. A request is recorded when this adapter hands it to the event
513    /// consumer and dropped when it is answered or denied on timeout.
514    pending_permissions: Vec<PendingPermission>,
515    permission_timeout: Duration,
516    /// Names of the MCP servers the START REQUEST mounted. A caller that mounts
517    /// a server intends its tools: a `can_use_tool` for `mcp__<name>__…` is
518    /// allowed here at once, never held for the caller to answer.
519    mounted_mcp_servers: Vec<String>,
520}
521
522/// One `can_use_tool` control request Claude Code is blocked on.
523struct PendingPermission {
524    /// The CLI's own `request_id`, echoed verbatim in the answer.
525    request_id: String,
526    /// When this adapter denies it if nobody has answered.
527    deadline: tokio::time::Instant,
528}
529
530impl ClaudeRuntimeConnection {
531    fn is_relay(&self) -> bool {
532        self.prefix
533            .env
534            .get("SUPERCODE_CLAUDE_RELAY")
535            .is_some_and(|value| value == "1")
536    }
537
538    /// Whether the `claude` process behind this connection has exited.
539    ///
540    /// A stream-json `claude --print` process is the session's CURRENT
541    /// speaker, not the session: it ends when it has nothing left to do (its
542    /// own exit, a crash, a reaped process group), while the session it wrote
543    /// stays on disk under the same id.
544    async fn process_ended(&self) -> bool {
545        matches!(self.transport.child.lock().await.try_wait(), Ok(Some(_)))
546    }
547
548    /// Put a live `claude` process back behind this connection by resuming the
549    /// session it already owns.
550    ///
551    /// The runtime id, the connection and the endpoint's protocol are
552    /// unchanged — only the process is new — so a caller of `harness.v1`
553    /// never learns that one ended. Verified against claude 2.1.258:
554    /// `--resume <id>` in bidirectional stream-json print mode reopens the
555    /// SAME session id (its `system/init` and every `result` echo it) with the
556    /// conversation intact, and goes on reading turns from stdin.
557    async fn reopen(&mut self) -> Result<()> {
558        let mut launch = self.prefix.clone();
559        launch
560            .arguments
561            .extend(["--resume".into(), self.handle.runtime_id.clone()]);
562        let transport = RawLineTransport::spawn(&launch, Some(&self.cwd), "claude-stream-json")
563            .await
564            .map_err(|error| {
565                Error::Other(format!(
566                    "could not resume Claude Code session `{}` after its process exited: {error}",
567                    self.handle.runtime_id
568                ))
569            })?;
570        self.handle.endpoint = transport.endpoint.clone();
571        self.transport = transport;
572        // Whatever the dead process was blocked on died with it: those ids
573        // name nothing the new process can be told about.
574        self.pending_permissions.clear();
575        self.spoke = false;
576        Ok(())
577    }
578
579    /// Write one frame to the CLI, resuming the session first when the process
580    /// that was reading stdin is gone.
581    ///
582    /// Both orders are covered because both happen: the exit can be visible
583    /// before the write (the child is reaped) or only in the write itself (a
584    /// grandchild — a hook, an MCP server — still holds the pipe, so nothing
585    /// upstream has noticed the exit and only `EPIPE` reveals it).
586    async fn write_turn(&mut self, frame: Value) -> Result<()> {
587        if self.process_ended().await {
588            if self.is_relay() {
589                return Err(Error::Other("Claude relay process exited".into()));
590            }
591            self.reopen().await?;
592        }
593        match self.transport.write(frame.clone()).await {
594            Ok(()) => Ok(()),
595            Err(error) if broken_pipe(&error) && !self.is_relay() => {
596                self.reopen().await?;
597                self.transport.write(frame).await
598            }
599            Err(error) => Err(error),
600        }
601    }
602
603    /// The control channel is adapter-private plumbing: a `control_response` is
604    /// the reply to a frame this adapter sent, never a harness event, so it is
605    /// dropped rather than forwarded to event consumers. A late reply that
606    /// arrives after `interrupt` gave up is dropped here too.
607    fn is_control_response(value: &Value) -> bool {
608        value.get("type").and_then(Value::as_str) == Some("control_response")
609    }
610
611    /// Match one `control_response` envelope against an outstanding request id.
612    ///
613    /// Ground truth (claude 2.1.224, verified live): the CLI answers
614    /// `{"type":"control_request","request_id":ID,"request":{"subtype":"interrupt"}}`
615    /// with
616    /// `{"type":"control_response","response":{"subtype":"success","request_id":ID,"response":{"still_queued":[]}}}`,
617    /// or with `{"subtype":"error","request_id":ID,"error":"…"}` on failure.
618    fn control_result(value: &Value, request_id: &str) -> Option<Result<()>> {
619        let response = value.get("response")?;
620        if response.get("request_id").and_then(Value::as_str) != Some(request_id) {
621            return None;
622        }
623        match response.get("subtype").and_then(Value::as_str) {
624            Some("success") => Some(Ok(())),
625            other => Some(Err(Error::Other(format!(
626                "Claude Code rejected the interrupt control request: {}",
627                response
628                    .get("error")
629                    .and_then(Value::as_str)
630                    .map(str::to_string)
631                    .unwrap_or_else(|| format!(
632                        "control_response subtype {}",
633                        other.unwrap_or("(missing)")
634                    ))
635            )))),
636        }
637    }
638
639    /// The CLI's own `request_id` when this frame is a permission request.
640    ///
641    /// Ground truth (claude 2.1.258, recorded live): a tool call that needs an
642    /// answer from the registered permission handler arrives as
643    /// `{"type":"control_request","request_id":"<uuid>","request":{"subtype":"can_use_tool","tool_name":…,"input":{…},"permission_suggestions":[…],"tool_use_id":…}}`.
644    /// No other inbound `control_request` subtype is answered here.
645    fn permission_request_id(value: &Value) -> Option<&str> {
646        if value.get("type").and_then(Value::as_str)? != "control_request" {
647            return None;
648        }
649        let request = value.get("request")?;
650        if request.get("subtype").and_then(Value::as_str)? != "can_use_tool" {
651            return None;
652        }
653        value.get("request_id").and_then(Value::as_str)
654    }
655
656    /// The request id of a `can_use_tool` for a tool of an MCP server this
657    /// connection's start request mounted (`mcp__<server>__<tool>`), if that
658    /// is what the frame is.
659    fn mounted_tool_permission_request(&self, payload: &Value) -> Option<String> {
660        let request_id = Self::permission_request_id(payload)?;
661        let tool = payload
662            .get("request")?
663            .get("tool_name")
664            .and_then(Value::as_str)?;
665        let mounted = self.mounted_mcp_servers.iter().any(|name| {
666            tool.strip_prefix("mcp__")
667                .and_then(|rest| rest.strip_prefix(name.as_str()))
668                .is_some_and(|rest| rest.starts_with("__"))
669        });
670        mounted.then(|| request_id.to_string())
671    }
672
673    /// Start the clock on a permission request about to reach the consumer.
674    fn note_permission_request(&mut self, payload: &Value) {
675        let Some(request_id) = Self::permission_request_id(payload) else {
676            return;
677        };
678        if self
679            .pending_permissions
680            .iter()
681            .any(|pending| pending.request_id == request_id)
682        {
683            return;
684        }
685        self.pending_permissions.push(PendingPermission {
686            request_id: request_id.to_string(),
687            deadline: tokio::time::Instant::now() + self.permission_timeout,
688        });
689    }
690
691    /// Write one `control_response` for a request the CLI is blocked on.
692    async fn write_permission_response(&mut self, request_id: &str, body: Value) -> Result<()> {
693        self.transport
694            .write(json!({
695                "type": "control_response",
696                "response": {
697                    "subtype": "success",
698                    "request_id": request_id,
699                    "response": body,
700                },
701            }))
702            .await
703    }
704
705    /// Deny every held request whose deadline has passed.
706    ///
707    /// Claude Code blocks the turn forever on an unanswered request, so an
708    /// unanswered request is denied rather than left to wedge the session.
709    async fn deny_expired_permissions(&mut self) -> Result<()> {
710        let now = tokio::time::Instant::now();
711        let expired = self
712            .pending_permissions
713            .iter()
714            .filter(|pending| pending.deadline <= now)
715            .map(|pending| pending.request_id.clone())
716            .collect::<Vec<_>>();
717        self.pending_permissions
718            .retain(|pending| pending.deadline > now);
719        for request_id in expired {
720            self.write_permission_response(
721                &request_id,
722                json!({"behavior": "deny", "message": CLAUDE_PERMISSION_TIMEOUT_MESSAGE}),
723            )
724            .await?;
725        }
726        Ok(())
727    }
728
729    /// What `next_event` reports when the CLI's stdout ended.
730    ///
731    /// For every other harness a closed transport is a closed runtime, and
732    /// reporting it is right. Claude Code is the one whose process is not the
733    /// session: `claude --print` exits when it has answered, and the session
734    /// it wrote is reopened by the next `send_input`. Reporting that exit as
735    /// `transport_closed` is what made this harness the only one whose
736    /// connection died between turns — the service drops a closed runtime, so
737    /// the caller's next input is refused on a connection the harness itself
738    /// considers perfectly resumable.
739    ///
740    /// So a session that has spoken PARKS here instead: the connection stays,
741    /// and it has nothing to say until someone speaks to it again — which is
742    /// exactly the state an idle runtime of any other harness is in. A process
743    /// that exited without ever speaking opened no session to park on, and its
744    /// close is reported as before.
745    async fn transport_ended(&mut self) -> Result<Option<HarnessEvent>> {
746        // A relay is a persistent mailbox speaker. Its death must close its host;
747        // parking it would leave a busy host steering messages into a dead pipe.
748        if self.is_relay() || !self.spoke {
749            return Ok(None);
750        }
751        std::future::pending().await
752    }
753
754    /// How long until the oldest held request must be denied.
755    fn next_permission_deadline(&self) -> Option<Duration> {
756        let now = tokio::time::Instant::now();
757        self.pending_permissions
758            .iter()
759            .map(|pending| pending.deadline.saturating_duration_since(now))
760            .min()
761    }
762}
763
764/// Validate one caller-supplied answer against the permission-result shape
765/// Claude Code accepts, filling in the parts the CLI requires.
766///
767/// The uniform door (`harness.v1.approvals.resolve`) sends
768/// `{"behavior":"allow"}` or `{"behavior":"deny","message":…}`; a caller using
769/// `harness.v1.runtimes.respond` directly may add the CLI's optional fields
770/// (`updatedInput`, `updatedPermissions`) and they are passed through
771/// untouched. Anything without a recognized `behavior` is refused by name
772/// rather than guessed at, because the CLI itself refuses it.
773fn claude_permission_result(response: Value) -> Result<Value> {
774    let Value::Object(mut body) = response else {
775        return Err(claude_permission_shape_error(&response));
776    };
777    match body.get("behavior").and_then(Value::as_str) {
778        Some("allow") => {}
779        Some("deny") => {
780            // Measured: the CLI rejects a `deny` with no `message`.
781            let empty = body
782                .get("message")
783                .and_then(Value::as_str)
784                .is_none_or(str::is_empty);
785            if empty {
786                body.insert(
787                    "message".into(),
788                    Value::String("Volter Harness denied this permission request".into()),
789                );
790            }
791        }
792        _ => return Err(claude_permission_shape_error(&Value::Object(body))),
793    }
794    Ok(Value::Object(body))
795}
796
797fn claude_permission_shape_error(response: &Value) -> Error {
798    Error::Other(format!(
799        "Claude Code permission answers must carry a `behavior` of {}; got {response}",
800        CLAUDE_PERMISSION_BEHAVIORS
801            .map(|behavior| format!("`{behavior}`"))
802            .join(" or "),
803    ))
804}
805
806#[async_trait]
807impl RuntimeConnection for ClaudeRuntimeConnection {
808    fn handle(&self) -> &RuntimeHandle {
809        &self.handle
810    }
811
812    async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
813        let content = if input.image_urls.is_empty() {
814            Value::String(input.text)
815        } else {
816            let mut parts = Vec::new();
817            if !input.text.is_empty() {
818                parts.push(json!({"type":"text", "text":input.text}));
819            }
820            for url in input.image_urls {
821                parts.push(claude_image_part(&url)?);
822            }
823            Value::Array(parts)
824        };
825        self.write_turn(json!({
826            "type": "user",
827            "session_id": self.handle.runtime_id,
828            "message": {"role": "user", "content": content},
829        }))
830        .await?;
831        Ok(None)
832    }
833
834    async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
835        if let Some(payload) = self.buffered_events.pop_front() {
836            self.spoke = true;
837            return Ok(Some(harness_event(payload)));
838        }
839        loop {
840            // A held permission request has a deadline, so the wait is bounded
841            // by the nearest one: an unanswered request is denied here rather
842            // than left blocking the CLI's turn forever.
843            self.deny_expired_permissions().await?;
844            if self.is_relay() && self.process_ended().await {
845                return Ok(None);
846            }
847            let payload = match self
848                .next_permission_deadline()
849                .map(|remaining| remaining.min(Duration::from_millis(200)))
850            {
851                Some(remaining) => {
852                    match tokio::time::timeout(remaining, self.transport.receiver.recv()).await {
853                        Err(_) => continue,
854                        Ok(None) => return self.transport_ended().await,
855                        Ok(Some(payload)) => payload,
856                    }
857                }
858                None => match tokio::time::timeout(
859                    Duration::from_millis(200),
860                    self.transport.receiver.recv(),
861                )
862                .await
863                {
864                    Err(_) => continue,
865                    Ok(None) => return self.transport_ended().await,
866                    Ok(Some(payload)) => payload,
867                },
868            };
869            self.spoke = true;
870            if Self::is_control_response(&payload) {
871                continue;
872            }
873            if let Some(request_id) = self.mounted_tool_permission_request(&payload) {
874                self.write_permission_response(&request_id, json!({"behavior": "allow"}))
875                    .await?;
876                continue;
877            }
878            self.note_permission_request(&payload);
879            return Ok(Some(harness_event(payload)));
880        }
881    }
882
883    /// Cancel the running turn through the stream-json control channel and wait
884    /// for the CLI's acknowledgement.
885    ///
886    /// Interrupting with no turn in flight is safe and succeeds: claude 2.1.224
887    /// acknowledges the request with `subtype: "success"` and an empty
888    /// `still_queued` list rather than erroring, and the session keeps
889    /// accepting input. The adapter reports what the harness reports instead of
890    /// inventing a turn-state gate of its own.
891    async fn interrupt(&mut self) -> Result<()> {
892        // A session whose process has exited has no turn in flight, and the
893        // adapter does not start one just to cancel nothing. This is the same
894        // answer the CLI gives an interrupt with nothing running.
895        if self.process_ended().await {
896            return Ok(());
897        }
898        let request_id = format!(
899            "supercode-{}-interrupt-{}",
900            self.handle.runtime_id, self.next_control_request
901        );
902        self.next_control_request += 1;
903        self.transport
904            .write(json!({
905                "type": "control_request",
906                "request_id": request_id,
907                "request": {"subtype": "interrupt"},
908            }))
909            .await?;
910
911        let deadline = tokio::time::Instant::now() + self.control_timeout;
912        loop {
913            let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
914            if remaining.is_zero() {
915                return Err(claude_interrupt_timeout(self.control_timeout));
916            }
917            match tokio::time::timeout(remaining, self.transport.receiver.recv()).await {
918                Err(_) => return Err(claude_interrupt_timeout(self.control_timeout)),
919                Ok(None) => return Err(Error::Other(
920                    "Claude Code stream-json transport closed before acknowledging the interrupt"
921                        .into(),
922                )),
923                Ok(Some(payload)) => {
924                    if Self::is_control_response(&payload) {
925                        if let Some(result) = Self::control_result(&payload, &request_id) {
926                            return result;
927                        }
928                        continue;
929                    }
930                    // A permission request raised while the interrupt was in
931                    // flight is still a request the CLI is blocked on: start
932                    // its clock now, not when the consumer drains the buffer.
933                    if let Some(request_id) = self.mounted_tool_permission_request(&payload) {
934                        self.write_permission_response(&request_id, json!({"behavior": "allow"}))
935                            .await?;
936                        continue;
937                    }
938                    self.note_permission_request(&payload);
939                    self.buffered_events.push_back(payload);
940                }
941            }
942        }
943    }
944
945    async fn steer(&mut self, text: String) -> Result<()> {
946        self.send_input(RuntimeInput {
947            text,
948            image_urls: Vec::new(),
949        })
950        .await
951        .map(|_| ())
952    }
953
954    /// Answer one `can_use_tool` permission request through the stream-json
955    /// control channel.
956    ///
957    /// Only a request this adapter has surfaced to the event consumer is
958    /// answerable: the CLI's `request_id` is the identity, the answer is sent
959    /// once, and a second answer for the same id is refused rather than
960    /// silently ignored (claude 2.1.258 logs and drops a duplicate
961    /// `control_response`).
962    async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
963        let Some(request_id) = request_id.as_str().map(str::to_string) else {
964            return Err(Error::Other(format!(
965                "Claude Code control requests are identified by a string `request_id`; got \
966                 {request_id}"
967            )));
968        };
969        let Some(index) = self
970            .pending_permissions
971            .iter()
972            .position(|pending| pending.request_id == request_id)
973        else {
974            return Err(Error::Other(format!(
975                "no Claude Code permission request `{request_id}` is waiting on this connection — \
976                 a `can_use_tool` request is answerable only while its turn is blocked on it, and \
977                 only until it is answered or denied on timeout"
978            )));
979        };
980        let body = claude_permission_result(response)?;
981        self.pending_permissions.remove(index);
982        self.write_permission_response(&request_id, body).await
983    }
984
985    async fn close(&mut self) -> Result<()> {
986        self.transport.close().await
987    }
988}
989
990/// Generic ACP v1 client backend for any ACP agent command.
991#[derive(Debug, Clone)]
992pub struct AcpRuntimeBackend {
993    harness: HarnessId,
994    launch: RuntimeLaunch,
995    resume_session: bool,
996}
997
998impl AcpRuntimeBackend {
999    /// Construct an ACP adapter for a named harness/agent command.
1000    pub fn new(harness: HarnessId, launch: RuntimeLaunch) -> Self {
1001        Self {
1002            harness,
1003            launch,
1004            resume_session: false,
1005        }
1006    }
1007
1008    /// Declare that this known ACP agent advertises `session/load` or
1009    /// `session/resume`. Attach still validates the capability negotiated by
1010    /// `initialize`, so a changed or incompatible agent fails honestly.
1011    pub fn with_resume_support(mut self, supported: bool) -> Self {
1012        self.resume_session = supported;
1013        self
1014    }
1015
1016    async fn connect(
1017        &self,
1018        cwd: &Path,
1019        launch: Option<RuntimeLaunch>,
1020    ) -> Result<(
1021        Arc<JsonLineClient>,
1022        mpsc::UnboundedReceiver<Value>,
1023        RuntimeEndpoint,
1024        Value,
1025    )> {
1026        let launch = launch.unwrap_or_else(|| self.launch.clone());
1027        let (client, receiver, endpoint) =
1028            JsonLineClient::spawn(&launch, Some(cwd), true, "acp-v1-jsonrpc").await?;
1029        let initialized = client
1030            .request(
1031                "initialize",
1032                json!({
1033                    "protocolVersion": 1,
1034                    "clientCapabilities": {},
1035                    "clientInfo": {
1036                        "name": "supercode",
1037                        "title": "Volter Harness",
1038                        "version": env!("CARGO_PKG_VERSION"),
1039                    },
1040                }),
1041            )
1042            .await?;
1043        if initialized.get("protocolVersion").and_then(Value::as_u64) != Some(1) {
1044            return Err(Error::Other(format!(
1045                "ACP agent negotiated unsupported protocol version: {}",
1046                initialized
1047                    .get("protocolVersion")
1048                    .cloned()
1049                    .unwrap_or(Value::Null)
1050            )));
1051        }
1052        Ok((client, receiver, endpoint, initialized))
1053    }
1054
1055    async fn session_request(
1056        &self,
1057        client: &JsonLineClient,
1058        initialized: &Value,
1059        method: &str,
1060        params: Value,
1061    ) -> Result<Value> {
1062        match client.request(method, params.clone()).await {
1063            Ok(response) => Ok(response),
1064            Err(error) if acp_auth_required(&error.to_string()) => {
1065                let cached = initialized
1066                    .get("authMethods")
1067                    .and_then(Value::as_array)
1068                    .and_then(|methods| {
1069                        methods.iter().find_map(|candidate| {
1070                            (candidate.get("id").and_then(Value::as_str) == Some("cached_token"))
1071                                .then_some("cached_token")
1072                        })
1073                    });
1074                let Some(method_id) = cached else {
1075                    return Err(Error::Other(
1076                        "ACP agent requires authentication but did not advertise the non-interactive `cached_token` method"
1077                            .into(),
1078                    ));
1079                };
1080                client
1081                    .request(
1082                        "authenticate",
1083                        json!({"methodId": method_id, "_meta": {"headless": true}}),
1084                    )
1085                    .await?;
1086                client.request(method, params).await
1087            }
1088            Err(error) => Err(error),
1089        }
1090    }
1091
1092    async fn connection(
1093        &self,
1094        cwd: &Path,
1095        runtime_id: Option<String>,
1096        launch: Option<RuntimeLaunch>,
1097        mcp_servers: Vec<McpServerLaunch>,
1098    ) -> Result<Box<dyn RuntimeConnection>> {
1099        let mcp_servers = acp_mcp_servers(&mcp_servers);
1100        // Supercode's ACP server serves exactly the session it was launched
1101        // with, and a resumed Supercode conversation continues under a new
1102        // id. So Supercode is resumed by launching the continuation
1103        // (`supercode resume <id> --acp`, with the default launch's flags) and
1104        // opening the session that process serves; asking a fresh `acp`
1105        // process to resume another id is always refused.
1106        let (launch, runtime_id) = match runtime_id {
1107            Some(session_id) if self.harness.as_str() == HarnessId::SUPERCODE => {
1108                let mut continuation = launch.unwrap_or_else(|| self.launch.clone());
1109                let flags = continuation
1110                    .arguments
1111                    .iter()
1112                    .skip_while(|argument| argument.as_str() != "acp")
1113                    .skip(1)
1114                    .cloned()
1115                    .collect::<Vec<_>>();
1116                continuation.arguments = ["resume", &session_id, "--harness", "supercode", "--acp"]
1117                    .into_iter()
1118                    .map(String::from)
1119                    .chain(flags)
1120                    .collect();
1121                (Some(continuation), None)
1122            }
1123            other => (launch, other),
1124        };
1125        let (client, mut receiver, endpoint, initialized) = self.connect(cwd, launch).await?;
1126        let session_id = if let Some(session_id) = runtime_id {
1127            let resume = initialized
1128                .pointer("/agentCapabilities/sessionCapabilities/resume")
1129                .is_some();
1130            let load = initialized
1131                .pointer("/agentCapabilities/loadSession")
1132                .and_then(Value::as_bool)
1133                .unwrap_or(false);
1134            let method = if resume {
1135                "session/resume"
1136            } else if load {
1137                "session/load"
1138            } else {
1139                return Err(Error::Other(
1140                    "ACP agent did not advertise session resume or load".into(),
1141                ));
1142            };
1143            self.session_request(
1144                client.as_ref(),
1145                &initialized,
1146                method,
1147                json!({"sessionId": session_id, "cwd": cwd, "mcpServers": mcp_servers}),
1148            )
1149            .await?;
1150            session_id
1151        } else {
1152            self.session_request(
1153                client.as_ref(),
1154                &initialized,
1155                "session/new",
1156                json!({"cwd": cwd, "mcpServers": mcp_servers}),
1157            )
1158            .await?
1159            .get("sessionId")
1160            .and_then(Value::as_str)
1161            .ok_or_else(|| Error::Other("ACP session/new omitted sessionId".into()))?
1162            .to_string()
1163        };
1164        // `session/load` is allowed to replay the persisted conversation as
1165        // `session/update` notifications before returning its response. Those
1166        // are bootstrap data, not output from a newly submitted prompt. If
1167        // they escape through the live runtime stream, clients fabricate an
1168        // assistant delta and a turn that can never complete because no
1169        // `session/prompt` request exists. The persisted transcript already
1170        // supplies this history, so discard every notification queued by the
1171        // completed new/load handshake before exposing the connection.
1172        while receiver.try_recv().is_ok() {}
1173        Ok(Box::new(AcpRuntimeConnection {
1174            handle: RuntimeHandle {
1175                harness: self.harness.clone(),
1176                runtime_id: session_id,
1177                endpoint,
1178            },
1179            client,
1180            receiver,
1181            active_prompt: None,
1182        }))
1183    }
1184}
1185
1186/// The uniform [`McpServerLaunch`] list in ACP's own `session/new` shape:
1187/// `{name, command, args, env: [{name, value}]}` (Agent Client Protocol v1
1188/// `McpServer`, the stdio form). An empty list stays `[]`, which is what this
1189/// backend has always sent.
1190fn acp_mcp_servers(servers: &[McpServerLaunch]) -> Value {
1191    Value::Array(
1192        servers
1193            .iter()
1194            .map(|server| {
1195                json!({
1196                    "name": server.name,
1197                    "command": server.command,
1198                    "args": server.arguments,
1199                    "env": server
1200                        .env
1201                        .iter()
1202                        .map(|(name, value)| json!({"name": name, "value": value}))
1203                        .collect::<Vec<_>>(),
1204                })
1205            })
1206            .collect::<Vec<_>>(),
1207    )
1208}
1209
1210fn acp_auth_required(message: &str) -> bool {
1211    let message = message.to_ascii_lowercase();
1212    [
1213        "auth",
1214        "login",
1215        "sign in",
1216        "sign-in",
1217        "unauthorized",
1218        "forbidden",
1219        "credential",
1220    ]
1221    .iter()
1222    .any(|needle| message.contains(needle))
1223}
1224
1225#[async_trait]
1226impl RuntimeBackend for AcpRuntimeBackend {
1227    fn harness(&self) -> HarnessId {
1228        self.harness.clone()
1229    }
1230
1231    fn capabilities(&self) -> RuntimeCapabilities {
1232        RuntimeCapabilities {
1233            start_session: true,
1234            // Optional in ACP v1. Known agents may declare it here; attach
1235            // still checks the actual initialize response before use.
1236            resume_session: self.resume_session,
1237            attach_existing_process: false,
1238            send_input: true,
1239            stream_events: true,
1240            interrupt: true,
1241            steer: false,
1242            respond_to_requests: true,
1243        }
1244    }
1245
1246    async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
1247        self.connection(&request.cwd, None, request.launch, request.mcp_servers)
1248            .await
1249    }
1250
1251    async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
1252        let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1253        self.connection(
1254            &cwd,
1255            Some(request.runtime_id),
1256            request.launch,
1257            request.mcp_servers,
1258        )
1259        .await
1260    }
1261}
1262
1263struct AcpRuntimeConnection {
1264    handle: RuntimeHandle,
1265    client: Arc<JsonLineClient>,
1266    receiver: mpsc::UnboundedReceiver<Value>,
1267    active_prompt: Option<u64>,
1268}
1269
1270#[async_trait]
1271impl RuntimeConnection for AcpRuntimeConnection {
1272    fn handle(&self) -> &RuntimeHandle {
1273        &self.handle
1274    }
1275
1276    async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
1277        let mut prompt = Vec::new();
1278        if !input.text.is_empty() {
1279            prompt.push(json!({"type": "text", "text": input.text}));
1280        }
1281        for url in input.image_urls {
1282            let (mime_type, data) = data_image_parts(&url).ok_or_else(|| {
1283                Error::Other("ACP image prompts require base64 image data URLs".into())
1284            })?;
1285            prompt.push(json!({"type":"image", "mimeType":mime_type, "data":data}));
1286        }
1287        let (id, response) = self
1288            .client
1289            .begin_request(
1290                "session/prompt",
1291                json!({
1292                    "sessionId": self.handle.runtime_id,
1293                    "prompt": prompt,
1294                }),
1295            )
1296            .await?;
1297        self.active_prompt = Some(id);
1298        let client = self.client.clone();
1299        tokio::spawn(async move {
1300            let result = match response.await {
1301                Ok(Ok(result)) => json!({"id": id, "result": result}),
1302                Ok(Err(error)) => json!({"id": id, "error": error}),
1303                Err(_) => json!({"id": id, "error": "response channel closed"}),
1304            };
1305            client.emit(json!({
1306                "jsonrpc": "2.0",
1307                "method": "supercode/acp_request_completed",
1308                "params": result,
1309            }));
1310        });
1311        Ok(Some(id.to_string()))
1312    }
1313
1314    async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
1315        let Some(payload) = self.receiver.recv().await else {
1316            return Ok(None);
1317        };
1318        let kind = payload
1319            .get("method")
1320            .and_then(Value::as_str)
1321            .or_else(|| payload.get("type").and_then(Value::as_str))
1322            .unwrap_or("protocol")
1323            .to_string();
1324        if kind == "supercode/acp_request_completed" {
1325            self.active_prompt = None;
1326        }
1327        Ok(Some(HarnessEvent {
1328            sequence: None,
1329            kind,
1330            payload,
1331        }))
1332    }
1333
1334    async fn interrupt(&mut self) -> Result<()> {
1335        self.client
1336            .notify(
1337                "session/cancel",
1338                json!({"sessionId": self.handle.runtime_id}),
1339            )
1340            .await
1341    }
1342
1343    async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
1344        self.client.respond(request_id, response).await
1345    }
1346
1347    async fn close(&mut self) -> Result<()> {
1348        self.client.close().await
1349    }
1350}
1351
1352/// OpenCode live-runtime backend using its official HTTP API and SSE event
1353/// stream. [`OpenCodeRuntimeBackend::connect`] can join the server embedded in
1354/// an already-running TUI when that TUI was launched with a known host/port.
1355#[derive(Debug, Clone)]
1356pub struct OpenCodeRuntimeBackend {
1357    launch: RuntimeLaunch,
1358    base_url: Option<String>,
1359    bearer: Option<BearerToken>,
1360}
1361
1362impl Default for OpenCodeRuntimeBackend {
1363    fn default() -> Self {
1364        Self::new()
1365    }
1366}
1367
1368impl OpenCodeRuntimeBackend {
1369    /// Launch a fresh `opencode serve` process for each connection.
1370    pub fn new() -> Self {
1371        Self {
1372            launch: RuntimeLaunch {
1373                program: "opencode".into(),
1374                arguments: vec!["serve".into()],
1375                env: BTreeMap::new(),
1376            },
1377            base_url: None,
1378            bearer: None,
1379        }
1380    }
1381
1382    /// Connect to an existing OpenCode server, including a TUI's server when
1383    /// it was launched on a known address.
1384    pub fn connect(base_url: impl Into<String>) -> Self {
1385        Self {
1386            base_url: Some(base_url.into().trim_end_matches('/').to_string()),
1387            ..Self::new()
1388        }
1389    }
1390
1391    /// Override the command used when launching a new OpenCode server.
1392    pub fn with_launch(mut self, launch: RuntimeLaunch) -> Self {
1393        self.launch = launch;
1394        self
1395    }
1396
1397    /// Send the resolved connect-mode credential as a bearer Authorization
1398    /// header on every request to the joined server.
1399    pub fn with_bearer(mut self, token: BearerToken) -> Self {
1400        self.bearer = Some(token);
1401        self
1402    }
1403
1404    fn http_client(&self) -> Result<reqwest::Client> {
1405        let Some(token) = &self.bearer else {
1406            return Ok(reqwest::Client::new());
1407        };
1408        let mut headers = reqwest::header::HeaderMap::new();
1409        let mut value =
1410            reqwest::header::HeaderValue::from_str(&format!("Bearer {}", token.secret())).map_err(
1411                |_| Error::Other("connect-mode bearer token is not a valid header value".into()),
1412            )?;
1413        value.set_sensitive(true);
1414        headers.insert(reqwest::header::AUTHORIZATION, value);
1415        reqwest::Client::builder()
1416            .default_headers(headers)
1417            .build()
1418            .map_err(|error| Error::Other(format!("could not build HTTP client: {error}")))
1419    }
1420
1421    async fn service(
1422        &self,
1423        client: &reqwest::Client,
1424        launch: Option<RuntimeLaunch>,
1425    ) -> Result<(String, Option<super::GroupLeader>)> {
1426        if let Some(base_url) = &self.base_url {
1427            wait_for_health(client, base_url).await?;
1428            return Ok((base_url.clone(), None));
1429        }
1430        let port = TcpListener::bind(("127.0.0.1", 0))?.local_addr()?.port();
1431        let mut launch = launch.unwrap_or_else(|| self.launch.clone());
1432        launch.arguments.extend([
1433            "--hostname".into(),
1434            "127.0.0.1".into(),
1435            "--port".into(),
1436            port.to_string(),
1437        ]);
1438        let mut command = Command::new(&launch.program);
1439        command
1440            .args(&launch.arguments)
1441            .envs(&launch.env)
1442            .stdin(Stdio::null())
1443            .stdout(Stdio::null())
1444            .stderr(Stdio::inherit())
1445            .kill_on_drop(true);
1446        // OpenCode's launcher may replace itself with or spawn a native
1447        // worker. Give the runtime its own group so close can reap the whole
1448        // server tree instead of orphaning the worker and its inherited FDs.
1449        #[cfg(unix)]
1450        command.process_group(0);
1451        // `kill_on_drop` only targets the launcher, and a startup that blows
1452        // its caller's deadline is DROPPED mid-wait. `GroupLeader` turns that
1453        // drop into a group signal; the explicit teardown below covers a
1454        // startup that fails while this future is still running.
1455        let mut child = super::GroupLeader(command.spawn().map_err(|error| {
1456            Error::Other(format!("could not launch {}: {error}", launch.program))
1457        })?);
1458        let base_url = format!("http://127.0.0.1:{port}");
1459        if let Err(error) = wait_for_health(client, &base_url).await {
1460            let _ = terminate_opencode_server(&mut child).await;
1461            return Err(error);
1462        }
1463        Ok((base_url, Some(child)))
1464    }
1465
1466    async fn open(
1467        &self,
1468        cwd: &Path,
1469        runtime_id: Option<String>,
1470        launch: Option<RuntimeLaunch>,
1471    ) -> Result<Box<dyn RuntimeConnection>> {
1472        let client = self.http_client()?;
1473        let (base_url, child) = self.service(&client, launch).await?;
1474        let cwd_string = cwd.to_string_lossy().to_string();
1475        let runtime_id = match runtime_id {
1476            Some(id) => {
1477                http_ok(
1478                    client
1479                        .get(format!("{base_url}/session/{id}"))
1480                        .query(&[("directory", &cwd_string)])
1481                        .send()
1482                        .await,
1483                )
1484                .await?;
1485                id
1486            }
1487            None => {
1488                let response = http_ok(
1489                    client
1490                        .post(format!("{base_url}/session"))
1491                        .query(&[("directory", &cwd_string)])
1492                        .json(&json!({}))
1493                        .send()
1494                        .await,
1495                )
1496                .await?;
1497                response
1498                    .json::<Value>()
1499                    .await
1500                    .map_err(http_error)?
1501                    .get("id")
1502                    .and_then(Value::as_str)
1503                    .ok_or_else(|| Error::Other("OpenCode create session omitted id".into()))?
1504                    .to_string()
1505            }
1506        };
1507        let receiver = spawn_sse(
1508            client.clone(),
1509            format!("{base_url}/event"),
1510            cwd_string.clone(),
1511        );
1512        Ok(Box::new(OpenCodeRuntimeConnection {
1513            handle: RuntimeHandle {
1514                harness: HarnessId::from(HarnessId::OPENCODE),
1515                runtime_id,
1516                endpoint: RuntimeEndpoint::Http {
1517                    base_url: base_url.clone(),
1518                    protocol: "opencode-http-sse".into(),
1519                },
1520            },
1521            base_url,
1522            cwd: cwd_string,
1523            client,
1524            receiver,
1525            child,
1526        }))
1527    }
1528}
1529
1530#[async_trait]
1531impl RuntimeBackend for OpenCodeRuntimeBackend {
1532    fn harness(&self) -> HarnessId {
1533        HarnessId::from(HarnessId::OPENCODE)
1534    }
1535
1536    fn capabilities(&self) -> RuntimeCapabilities {
1537        RuntimeCapabilities {
1538            start_session: true,
1539            resume_session: true,
1540            attach_existing_process: self.base_url.is_some(),
1541            send_input: true,
1542            stream_events: true,
1543            interrupt: true,
1544            steer: false,
1545            respond_to_requests: true,
1546        }
1547    }
1548
1549    async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
1550        self.open(&request.cwd, None, request.launch).await
1551    }
1552
1553    async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
1554        let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1555        self.open(&cwd, Some(request.runtime_id), request.launch)
1556            .await
1557    }
1558
1559    async fn attach_existing(
1560        &self,
1561        request: RuntimeAttachRequest,
1562    ) -> Result<Box<dyn RuntimeConnection>> {
1563        if self.base_url.is_none() {
1564            return Err(Error::Other(
1565                "OpenCode live attach requires the existing server's `base_url`".into(),
1566            ));
1567        }
1568        let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1569        self.open(&cwd, Some(request.runtime_id), request.launch)
1570            .await
1571    }
1572}
1573
1574struct OpenCodeRuntimeConnection {
1575    handle: RuntimeHandle,
1576    base_url: String,
1577    cwd: String,
1578    client: reqwest::Client,
1579    receiver: mpsc::UnboundedReceiver<Value>,
1580    child: Option<super::GroupLeader>,
1581}
1582
1583#[async_trait]
1584impl RuntimeConnection for OpenCodeRuntimeConnection {
1585    fn handle(&self) -> &RuntimeHandle {
1586        &self.handle
1587    }
1588
1589    async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
1590        let mut parts = Vec::new();
1591        if !input.text.is_empty() {
1592            parts.push(json!({"type": "text", "text": input.text}));
1593        }
1594        for url in input.image_urls {
1595            let mime = image_mime_type(&url).ok_or_else(|| {
1596                Error::Other("OpenCode image prompts require a recognizable image MIME type".into())
1597            })?;
1598            parts.push(json!({"type":"file", "mime":mime, "url":url}));
1599        }
1600        http_ok(
1601            self.client
1602                .post(format!(
1603                    "{}/session/{}/prompt_async",
1604                    self.base_url, self.handle.runtime_id
1605                ))
1606                .query(&[("directory", &self.cwd)])
1607                .json(&json!({"parts": parts}))
1608                .send()
1609                .await,
1610        )
1611        .await?;
1612        Ok(None)
1613    }
1614
1615    async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
1616        loop {
1617            let Some(payload) = self.receiver.recv().await else {
1618                return Ok(None);
1619            };
1620            if opencode_event_session_id(&payload)
1621                .is_some_and(|session_id| session_id != self.handle.runtime_id)
1622            {
1623                continue;
1624            }
1625            let kind = payload
1626                .get("type")
1627                .and_then(Value::as_str)
1628                .unwrap_or("event")
1629                .to_string();
1630            return Ok(Some(HarnessEvent {
1631                sequence: None,
1632                kind,
1633                payload,
1634            }));
1635        }
1636    }
1637
1638    async fn interrupt(&mut self) -> Result<()> {
1639        http_ok(
1640            self.client
1641                .post(format!(
1642                    "{}/session/{}/abort",
1643                    self.base_url, self.handle.runtime_id
1644                ))
1645                .query(&[("directory", &self.cwd)])
1646                .send()
1647                .await,
1648        )
1649        .await?;
1650        Ok(())
1651    }
1652
1653    async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
1654        let permission = request_id.as_str().ok_or_else(|| {
1655            Error::Other("OpenCode permission request id must be a string".into())
1656        })?;
1657        http_ok(
1658            self.client
1659                .post(format!(
1660                    "{}/session/{}/permissions/{permission}",
1661                    self.base_url, self.handle.runtime_id
1662                ))
1663                .query(&[("directory", &self.cwd)])
1664                .json(&response)
1665                .send()
1666                .await,
1667        )
1668        .await?;
1669        Ok(())
1670    }
1671
1672    async fn close(&mut self) -> Result<()> {
1673        if let Some(child) = &mut self.child {
1674            terminate_opencode_server(child).await?;
1675        }
1676        Ok(())
1677    }
1678}
1679
1680fn data_image_parts(url: &str) -> Option<(&str, &str)> {
1681    let rest = url.strip_prefix("data:")?;
1682    let (mime_type, data) = rest.split_once(";base64,")?;
1683    mime_type.starts_with("image/").then_some((mime_type, data))
1684}
1685
1686fn image_mime_type(url: &str) -> Option<&str> {
1687    if let Some((mime_type, _)) = data_image_parts(url) {
1688        return Some(mime_type);
1689    }
1690    let path = url.split(['?', '#']).next()?.to_ascii_lowercase();
1691    if path.ends_with(".png") {
1692        Some("image/png")
1693    } else if path.ends_with(".jpg") || path.ends_with(".jpeg") {
1694        Some("image/jpeg")
1695    } else if path.ends_with(".gif") {
1696        Some("image/gif")
1697    } else if path.ends_with(".webp") {
1698        Some("image/webp")
1699    } else {
1700        None
1701    }
1702}
1703
1704fn claude_image_part(url: &str) -> Result<Value> {
1705    if let Some((media_type, data)) = data_image_parts(url) {
1706        return Ok(json!({
1707            "type":"image",
1708            "source":{"type":"base64", "media_type":media_type, "data":data}
1709        }));
1710    }
1711    if url.starts_with("https://") || url.starts_with("http://") {
1712        return Ok(json!({"type":"image", "source":{"type":"url", "url":url}}));
1713    }
1714    Err(Error::Other(
1715        "Claude image prompts require image data URLs or HTTP(S) URLs".into(),
1716    ))
1717}
1718
1719fn opencode_event_session_id(payload: &Value) -> Option<&str> {
1720    let properties = payload.get("properties").unwrap_or(payload);
1721    properties
1722        .get("sessionID")
1723        .and_then(Value::as_str)
1724        .or_else(|| {
1725            properties
1726                .get("part")
1727                .and_then(|part| part.get("sessionID"))
1728                .and_then(Value::as_str)
1729        })
1730        .or_else(|| {
1731            properties
1732                .get("info")
1733                .and_then(|info| info.get("sessionID"))
1734                .and_then(Value::as_str)
1735        })
1736}
1737
1738async fn terminate_opencode_server(child: &mut Child) -> Result<()> {
1739    #[cfg(unix)]
1740    let process_group = child.id();
1741    let leader_exited = child.try_wait()?.is_some();
1742    if leader_exited {
1743        #[cfg(unix)]
1744        if let Some(pid) = process_group.filter(|pid| process_group_exists(*pid)) {
1745            crate::lsp::kill_process_group(pid);
1746            wait_for_process_group_exit(pid, Duration::from_secs(3)).await?;
1747        }
1748        return Ok(());
1749    }
1750    // `Child::kill().await` waits for process reaping and can block forever
1751    // when a launcher leaves its native worker and inherited handles alive.
1752    // Terminate the isolated group while its leader can still reap workers;
1753    // killing leader and workers simultaneously can leave transient orphan
1754    // zombies and made close observably race process cleanup on Linux.
1755    #[cfg(unix)]
1756    if let Some(pid) = process_group {
1757        unsafe {
1758            libc::kill(-(pid as libc::pid_t), libc::SIGTERM);
1759        }
1760        let mut leader_reaped = false;
1761        if let Ok(status) = tokio::time::timeout(Duration::from_millis(500), child.wait()).await {
1762            status?;
1763            leader_reaped = true;
1764            if !process_group_exists(pid) {
1765                return Ok(());
1766            }
1767        }
1768        // The leader may exit while a detached worker ignores SIGTERM. Do
1769        // not mistake a reaped launcher for a stopped server tree.
1770        crate::lsp::kill_process_group(pid);
1771        if leader_reaped {
1772            return wait_for_process_group_exit(pid, Duration::from_secs(3)).await;
1773        }
1774    }
1775    #[cfg(not(unix))]
1776    child.start_kill()?;
1777    tokio::time::timeout(Duration::from_secs(3), child.wait())
1778        .await
1779        .map_err(|_| Error::Other("timed out reaping the OpenCode server".into()))??;
1780    #[cfg(unix)]
1781    if let Some(pid) = process_group {
1782        wait_for_process_group_exit(pid, Duration::from_secs(3)).await?;
1783    }
1784    Ok(())
1785}
1786
1787#[cfg(unix)]
1788fn process_group_exists(pid: u32) -> bool {
1789    let result = unsafe { libc::kill(-(pid as libc::pid_t), 0) };
1790    result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
1791}
1792
1793#[cfg(unix)]
1794async fn wait_for_process_group_exit(pid: u32, timeout: Duration) -> Result<()> {
1795    let deadline = tokio::time::Instant::now() + timeout;
1796    while process_group_exists(pid) {
1797        if tokio::time::Instant::now() >= deadline {
1798            return Err(Error::Other(format!(
1799                "timed out stopping OpenCode process group {pid}"
1800            )));
1801        }
1802        tokio::time::sleep(Duration::from_millis(10)).await;
1803    }
1804    Ok(())
1805}
1806
1807struct RawLineTransport {
1808    stdin: Mutex<ChildStdin>,
1809    child: Mutex<super::GroupLeader>,
1810    receiver: mpsc::UnboundedReceiver<Value>,
1811    endpoint: RuntimeEndpoint,
1812}
1813
1814impl RawLineTransport {
1815    async fn spawn(launch: &RuntimeLaunch, cwd: Option<&Path>, protocol: &str) -> Result<Self> {
1816        let mut command = Command::new(&launch.program);
1817        command
1818            .args(&launch.arguments)
1819            .envs(&launch.env)
1820            .stdin(Stdio::piped())
1821            .stdout(Stdio::piped())
1822            .stderr(Stdio::inherit())
1823            .kill_on_drop(true);
1824        // A raw-protocol launcher spawns its worker the same way a JSON-line
1825        // one does. Give it its own group so close — and a dropped deadline —
1826        // reap the worker instead of orphaning it holding this stdout.
1827        #[cfg(unix)]
1828        command.process_group(0);
1829        if let Some(cwd) = cwd {
1830            command.current_dir(cwd);
1831        }
1832        let mut child = command.spawn().map_err(|error| {
1833            Error::Other(format!("could not launch {}: {error}", launch.program))
1834        })?;
1835        let pid = child.id();
1836        let stdin = child
1837            .stdin
1838            .take()
1839            .ok_or_else(|| Error::Other("runtime child has no stdin".into()))?;
1840        let stdout = child
1841            .stdout
1842            .take()
1843            .ok_or_else(|| Error::Other("runtime child has no stdout".into()))?;
1844        let (sender, receiver) = mpsc::unbounded_channel();
1845        tokio::spawn(async move {
1846            let mut lines = BufReader::new(stdout).lines();
1847            while let Ok(Some(line)) = lines.next_line().await {
1848                let value = serde_json::from_str(&line)
1849                    .unwrap_or_else(|_| json!({"type": "malformed_output", "line": line}));
1850                let _ = sender.send(value);
1851            }
1852        });
1853        Ok(Self {
1854            stdin: Mutex::new(stdin),
1855            child: Mutex::new(super::GroupLeader(child)),
1856            receiver,
1857            endpoint: RuntimeEndpoint::LocalProcess {
1858                pid,
1859                command: std::iter::once(launch.program.clone())
1860                    .chain(launch.arguments.iter().cloned())
1861                    .collect(),
1862                protocol: protocol.into(),
1863            },
1864        })
1865    }
1866
1867    async fn write(&self, value: Value) -> Result<()> {
1868        let mut stdin = self.stdin.lock().await;
1869        stdin.write_all(value.to_string().as_bytes()).await?;
1870        stdin.write_all(b"\n").await?;
1871        stdin.flush().await?;
1872        Ok(())
1873    }
1874
1875    async fn close(&self) -> Result<()> {
1876        let mut child = self.child.lock().await;
1877        if child.try_wait()?.is_some() {
1878            return Ok(());
1879        }
1880        // `Child::kill` reaches the launcher only. Signal the whole group so
1881        // a shim's worker cannot outlive the connection that owns it.
1882        #[cfg(unix)]
1883        if let Some(pid) = child.id() {
1884            crate::lsp::kill_process_group(pid);
1885            tokio::time::timeout(Duration::from_secs(3), child.wait())
1886                .await
1887                .map_err(|_| Error::Other("timed out reaping runtime process group".into()))??;
1888            return Ok(());
1889        }
1890        #[cfg(not(unix))]
1891        child.kill().await?;
1892        Ok(())
1893    }
1894}
1895
1896async fn raw_next_event(
1897    receiver: &mut mpsc::UnboundedReceiver<Value>,
1898) -> Result<Option<HarnessEvent>> {
1899    let Some(payload) = receiver.recv().await else {
1900        return Ok(None);
1901    };
1902    Ok(Some(harness_event(payload)))
1903}
1904
1905fn harness_event(payload: Value) -> HarnessEvent {
1906    let kind = payload
1907        .get("type")
1908        .and_then(Value::as_str)
1909        .unwrap_or("event")
1910        .to_string();
1911    HarnessEvent {
1912        sequence: None,
1913        kind,
1914        payload,
1915    }
1916}
1917
1918/// Whether this error is the write end of a pipe whose reader is gone.
1919///
1920/// `EPIPE` is how a process exit reaches a writer that had no other way to
1921/// see it: the child's stdout can still be held open by a grandchild it
1922/// spawned, so nothing upstream has noticed, and the failed write is the
1923/// first evidence.
1924fn broken_pipe(error: &Error) -> bool {
1925    matches!(error, Error::Io(io) if io.kind() == std::io::ErrorKind::BrokenPipe)
1926}
1927
1928fn claude_interrupt_timeout(bound: Duration) -> Error {
1929    Error::Other(format!(
1930        "Claude Code did not acknowledge the interrupt control request within {}s",
1931        bound.as_secs_f32()
1932    ))
1933}
1934
1935pub(crate) fn generated_session_id() -> String {
1936    let mut bytes = [0_u8; 16];
1937    if getrandom::getrandom(&mut bytes).is_err() {
1938        let nanos = SystemTime::now()
1939            .duration_since(UNIX_EPOCH)
1940            .unwrap_or_default()
1941            .as_nanos()
1942            .to_le_bytes();
1943        bytes.copy_from_slice(&nanos);
1944    }
1945    bytes[6] = (bytes[6] & 0x0f) | 0x40;
1946    bytes[8] = (bytes[8] & 0x3f) | 0x80;
1947    format!(
1948        "{:02x}{:02x}{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}{:02x}{:02x}{:02x}{:02x}",
1949        bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5], bytes[6], bytes[7],
1950        bytes[8], bytes[9], bytes[10], bytes[11], bytes[12], bytes[13], bytes[14], bytes[15]
1951    )
1952}
1953
1954async fn wait_for_health(client: &reqwest::Client, base_url: &str) -> Result<()> {
1955    wait_for_health_for(client, base_url, Duration::from_secs(10)).await
1956}
1957
1958async fn wait_for_health_for(
1959    client: &reqwest::Client,
1960    base_url: &str,
1961    total_timeout: Duration,
1962) -> Result<()> {
1963    let url = format!("{base_url}/global/health");
1964    let mut last = None;
1965    let deadline = tokio::time::Instant::now() + total_timeout;
1966    // Local package-manager shims can take longer than five seconds to start
1967    // under build or indexing load. Ten seconds avoids false unavailability
1968    // without permitting an unbounded launch; inventory handshakes retain
1969    // their separate 30-second bound around the complete startup.
1970    loop {
1971        let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
1972        if remaining.is_zero() {
1973            break;
1974        }
1975        let request_timeout = remaining.min(Duration::from_millis(500));
1976        match tokio::time::timeout(request_timeout, client.get(&url).send()).await {
1977            Ok(Ok(response)) if response.status().is_success() => return Ok(()),
1978            Ok(Ok(response)) => last = Some(format!("HTTP {}", response.status())),
1979            Ok(Err(error)) => last = Some(error.to_string()),
1980            Err(_) => last = Some("health request timed out".into()),
1981        }
1982        let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
1983        if !remaining.is_zero() {
1984            tokio::time::sleep(remaining.min(Duration::from_millis(100))).await;
1985        }
1986    }
1987    Err(Error::Other(format!(
1988        "OpenCode server at {base_url} did not become healthy: {}",
1989        last.unwrap_or_else(|| "no response".into())
1990    )))
1991}
1992
1993async fn http_ok(
1994    response: std::result::Result<reqwest::Response, reqwest::Error>,
1995) -> Result<reqwest::Response> {
1996    response
1997        .map_err(http_error)?
1998        .error_for_status()
1999        .map_err(http_error)
2000}
2001
2002fn http_error(error: reqwest::Error) -> Error {
2003    Error::Other(format!("runtime HTTP request failed: {error}"))
2004}
2005
2006fn spawn_sse(
2007    client: reqwest::Client,
2008    url: String,
2009    directory: String,
2010) -> mpsc::UnboundedReceiver<Value> {
2011    let (sender, receiver) = mpsc::unbounded_channel();
2012    tokio::spawn(async move {
2013        let response = client
2014            .get(url)
2015            .query(&[("directory", directory)])
2016            .send()
2017            .await;
2018        let Ok(response) = response.and_then(reqwest::Response::error_for_status) else {
2019            let _ = sender.send(
2020                json!({"type": "stream_error", "message": "could not open OpenCode SSE stream"}),
2021            );
2022            return;
2023        };
2024        let mut stream = response.bytes_stream();
2025        let mut buffer = String::new();
2026        while let Some(chunk) = stream.next().await {
2027            let Ok(chunk) = chunk else {
2028                break;
2029            };
2030            buffer.push_str(&String::from_utf8_lossy(&chunk));
2031            while let Some(newline) = buffer.find('\n') {
2032                let line = buffer[..newline].trim_end_matches('\r').to_string();
2033                buffer.drain(..=newline);
2034                if let Some(data) = line.strip_prefix("data:") {
2035                    let data = data.trim();
2036                    if let Ok(value) = serde_json::from_str(data) {
2037                        let _ = sender.send(value);
2038                    }
2039                }
2040            }
2041        }
2042    });
2043    receiver
2044}