Skip to main content

scv_server/
lib.rs

1//! SCV's authoritative stdio server.
2
3mod agents;
4pub mod components;
5mod config;
6
7use std::{
8    collections::{HashMap, VecDeque},
9    io::Read as _,
10    path::{Path, PathBuf},
11    sync::{
12        Arc,
13        atomic::{AtomicU64, AtomicUsize, Ordering},
14    },
15    time::Duration,
16};
17
18use anyhow::{Context, Result, anyhow};
19use async_trait::async_trait;
20use config::Config;
21pub use config::{ApprovalPolicy, ConfigOverrides, user_home_path};
22use sha2::{Digest, Sha256};
23pub fn init_user_config() -> anyhow::Result<std::path::PathBuf> {
24    config::Config::init_user_config()
25}
26pub fn update_index_url(workspace: &std::path::Path) -> anyhow::Result<Option<String>> {
27    Ok(config::Config::load(workspace, ConfigOverrides::default())?
28        .update
29        .index_url)
30}
31
32/// Build a command for a native agent's CLI with the same private home and
33/// cleaned environment the daemon's `agent_<name>` tool uses, so the agent's
34/// own sign-in stores credentials where delegated runs will find them.
35/// Project configuration cannot set `[agents]`, so none is read.
36pub fn agent_command(agent: &str) -> Result<std::process::Command> {
37    let config = Config::load_user(ConfigOverrides::default())?;
38    config.prepare_adapter_homes()?;
39    let adapter = config
40        .adapters()
41        .remove(&format!("agent_{agent}"))
42        .ok_or_else(|| anyhow!("unknown agent {agent}"))?;
43    let mut command = std::process::Command::new(&adapter.command);
44    command
45        .current_dir(config.instance_home.join("adapters").join(agent))
46        .envs(adapter.environment);
47    for variable in scv_tools::AGENT_REMOVED_ENVIRONMENT {
48        command.env_remove(variable);
49    }
50    Ok(command)
51}
52
53/// Copy the user's own Codex setup from `source` into SCV's private Codex
54/// adapter home: `config.toml`, and `auth.json` only when it holds an API key.
55/// Returns display lines that never contain secret values.
56pub fn import_codex(source: &Path) -> Result<Vec<String>> {
57    let config = Config::load_user(ConfigOverrides::default())?;
58    config.prepare_adapter_homes()?;
59    agents::import_codex(source, &config.instance_home.join("adapters").join("codex"))
60}
61
62/// Return the user service name for the selected SCV instance.
63pub fn service_name() -> anyhow::Result<String> {
64    if std::env::var_os("SCV_HOME").is_none() {
65        return Ok("scv.service".into());
66    }
67    let home = user_home_path().ok_or_else(|| anyhow!("cannot determine SCV instance home"))?;
68    let digest = Sha256::digest(home.to_string_lossy().as_bytes());
69    let suffix = digest[..8]
70        .iter()
71        .map(|byte| format!("{byte:02x}"))
72        .collect::<String>();
73    Ok(format!("scv-{suffix}.service"))
74}
75
76pub fn service_unit_path() -> anyhow::Result<std::path::PathBuf> {
77    let config =
78        dirs::config_dir().ok_or_else(|| anyhow!("cannot determine XDG config directory"))?;
79    Ok(config.join("systemd/user").join(service_name()?))
80}
81use scv_core::{
82    AgentError, AgentRuntime, ApprovalGate, ApprovalRequest, BudgetContextPolicy, CoreEvent,
83    EventSink, Message, ToolRegistry, ToolRisk,
84};
85use scv_protocol::{ClientMessage, PROTOCOL_VERSION, PeerInfo, QueueEntry, ServerEvent, Usage};
86use scv_provider_openai::OpenAiProvider;
87use scv_tools::{SkillMap, builtin_registry};
88use tokio::{
89    io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader},
90    net::{UnixListener, UnixStream},
91    sync::{Mutex, OwnedSemaphorePermit, Semaphore, mpsc, oneshot},
92    task::JoinHandle,
93};
94use tokio_util::sync::CancellationToken;
95use tokio_util::task::TaskTracker;
96use uuid::Uuid;
97
98const PROMPT_LIMIT_BYTES: usize = 256 * 1024;
99const OUTPUT_QUEUE_CAPACITY: usize = 256;
100const OUTPUT_QUEUE_MIN_BYTES: usize = 16 * 1024 * 1024;
101const SHUTDOWN_GRACE: Duration = Duration::from_secs(3);
102const MAX_QUEUE_ITEMS: usize = 64;
103const MAX_QUEUE_BYTES: usize = 4 * 1024 * 1024;
104
105pub async fn run_stdio(overrides: ConfigOverrides) -> Result<()> {
106    let stdin = tokio::io::stdin();
107    let stdout = tokio::io::stdout();
108    let tasks = TaskTracker::new();
109    let result = run_managed(
110        stdin,
111        stdout,
112        overrides,
113        None,
114        CancellationToken::new(),
115        tasks.clone(),
116    )
117    .await;
118    tasks.close();
119    tasks.wait().await;
120    result
121}
122
123/// Return the local Unix socket used by the SCV daemon and TUI.
124pub fn default_socket_path() -> Result<PathBuf> {
125    scv_client::default_socket_path()
126}
127
128/// Run the authoritative server on the local Unix socket.
129pub async fn run_socket(path: &Path, overrides: ConfigOverrides) -> Result<()> {
130    if let Some(parent) = path.parent() {
131        tokio::fs::create_dir_all(parent)
132            .await
133            .context("create SCV socket directory")?;
134    }
135    let _lock = SocketLock::acquire(path)?;
136    if path.exists() {
137        if UnixStream::connect(path).await.is_ok() {
138            return Err(anyhow!(
139                "SCV server is already running at {}",
140                path.display()
141            ));
142        }
143        use std::os::unix::fs::FileTypeExt;
144        if !std::fs::symlink_metadata(path)?.file_type().is_socket() {
145            return Err(anyhow!(
146                "refusing to remove a non-socket at SCV socket path"
147            ));
148        }
149        tokio::fs::remove_file(path)
150            .await
151            .with_context(|| format!("remove stale SCV socket {}", path.display()))?;
152    }
153    let listener = UnixListener::bind(path)
154        .with_context(|| format!("bind SCV server socket {}", path.display()))?;
155    #[cfg(unix)]
156    {
157        use std::os::unix::fs::PermissionsExt;
158        std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
159            .context("secure SCV socket")?;
160    }
161    let components = Arc::new(Mutex::new(components::Components::new(
162        path.to_owned(),
163        std::env::current_dir()?,
164    )));
165    let cancellation = CancellationToken::new();
166    let tasks = TaskTracker::new();
167    let mut clients = tokio::task::JoinSet::new();
168    let refresh_components = components.clone();
169    let refresh_cancel = cancellation.clone();
170    let mut refresh_task = tokio::spawn(async move {
171        let mut refresh = tokio::time::interval(Duration::from_secs(2));
172        loop {
173            tokio::select! {
174                biased;
175                _ = refresh_cancel.cancelled() => break,
176                _ = refresh.tick() => {
177                    tokio::select! {
178                        biased;
179                        _ = refresh_cancel.cancelled() => break,
180                        result = async { refresh_components.lock().await.reconcile().await } => {
181                            if result.is_err() { tracing::warn!("Component account discovery failed"); }
182                        }
183                    }
184                }
185            }
186        }
187    });
188    let _refresh_abort = AbortGuard(refresh_task.abort_handle());
189    let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())?;
190    let result = loop {
191        tokio::select! {
192            accepted = listener.accept() => {
193                let (stream, _) = match accepted { Ok(value) => value, Err(error) => break Err(error.into()) };
194                let child_overrides = overrides.clone();
195                let components = components.clone();
196                let cancellation = cancellation.clone();
197                let tasks = tasks.clone();
198                clients.spawn(async move {
199                    let (reader, writer) = stream.into_split();
200                    if run_managed(reader, writer, child_overrides, Some(components), cancellation, tasks).await.is_err() {
201                        tracing::warn!("SCV socket client stopped");
202                    }
203                });
204            }
205            _ = clients.join_next(), if !clients.is_empty() => {},
206            _ = tokio::signal::ctrl_c() => break Ok(()),
207            _ = terminate.recv() => break Ok(()),
208        }
209    };
210    drop(listener);
211    cancellation.cancel();
212    let _ = (&mut refresh_task).await;
213    components.lock().await.shutdown().await;
214    if tokio::time::timeout(Duration::from_secs(8), async {
215        while clients.join_next().await.is_some() {}
216    })
217    .await
218    .is_err()
219    {
220        clients.abort_all();
221        while clients.join_next().await.is_some() {}
222    }
223    tasks.close();
224    tasks.wait().await;
225    let _ = tokio::fs::remove_file(path).await;
226    result
227}
228
229async fn run_managed<R, W>(
230    reader: R,
231    writer: W,
232    overrides: ConfigOverrides,
233    components: Option<Arc<Mutex<components::Components>>>,
234    cancellation: CancellationToken,
235    tasks: TaskTracker,
236) -> Result<()>
237where
238    R: tokio::io::AsyncRead + Unpin,
239    W: tokio::io::AsyncWrite + Unpin + Send + 'static,
240{
241    let initial_output_bytes =
242        output_queue_bytes(Config::default().protocol.max_server_frame_bytes)?;
243    let (output_tx, mut output_rx) = outbound_channel(initial_output_bytes);
244    let mut writer_task = tasks.spawn(async move {
245        let mut writer = writer;
246        while let Some(frame) = output_rx.recv().await {
247            writer.write_all(&frame.bytes).await?;
248            writer.write_all(b"\n").await?;
249            writer.flush().await?;
250        }
251        Ok::<(), std::io::Error>(())
252    });
253    let _writer_abort = AbortGuard(writer_task.abort_handle());
254    let (done_tx, mut done_rx) = mpsc::channel::<TurnDone>(4);
255    let approvals = Arc::new(ApprovalBroker::default());
256    let mut reader = BufReader::new(reader);
257    let mut frames = FrameBuffer::default();
258    let mut initialized = false;
259    let mut session: Option<Session> = None;
260    let mut active: Option<ActiveTurn> = None;
261    let mut fatal = false;
262    let mut writer_finished = false;
263
264    let loop_result: Result<()> = async {
265        loop {
266        let frame_limit = session.as_ref().map_or_else(
267            || Config::default().protocol.max_client_frame_bytes,
268            |value| value.config.protocol.max_client_frame_bytes,
269        );
270        tokio::select! {
271            _ = cancellation.cancelled() => break,
272            read = frames.read(&mut reader, frame_limit) => {
273                let frame = match read.context("read protocol input")? {
274                    FrameRead::Eof => {
275                        if let Some(active) = &active { active.cancellation.cancel(); }
276                        break;
277                    }
278                    FrameRead::TooLarge => {
279                        send_error(&output_tx, "", "invalid_request", "client frame exceeds configured limit", false, server_frame_limit(&session)).await?;
280                        continue;
281                    }
282                    FrameRead::Frame(frame) => frame,
283                };
284                if frame.is_empty() {
285                    send_error(&output_tx, "", "invalid_json", "protocol frame is empty", false, server_frame_limit(&session)).await?;
286                    continue;
287                }
288                let message = match serde_json::from_slice::<ClientMessage>(&frame) {
289                    Ok(message) => message,
290                    Err(error) => {
291                        send_error(&output_tx, "", "invalid_json", &format!("invalid protocol JSON: {error}"), false, server_frame_limit(&session)).await?;
292                        continue;
293                    }
294                };
295                match message {
296                    ClientMessage::Initialize { request_id, protocol_version, .. } => {
297                        if initialized {
298                            send_error(&output_tx, &request_id, "invalid_request", "connection is already initialized", false, server_frame_limit(&session)).await?;
299                            continue;
300                        }
301                        if protocol_version != PROTOCOL_VERSION {
302                            send_error(&output_tx, &request_id, "version_mismatch", &format!("server supports protocol {PROTOCOL_VERSION}"), true, server_frame_limit(&session)).await?;
303                            fatal = true;
304                            break;
305                        }
306                        initialized = true;
307                        send_event(&output_tx, ServerEvent::Initialized {
308                            request_id,
309                            protocol_version: PROTOCOL_VERSION,
310                            server: PeerInfo { name: "scv-server".into(), version: env!("CARGO_PKG_VERSION").into() },
311                        }, Config::default().protocol.max_server_frame_bytes).await?;
312                    }
313                    other if !initialized => {
314                        send_error(&output_tx, other.request_id(), "not_initialized", "initialize must be the first message", false, server_frame_limit(&session)).await?;
315                    }
316                    ClientMessage::DaemonControl { request_id, command } => {
317                        if let Some(components) = &components {
318                            let result = tokio::select! {
319                                biased;
320                                _ = cancellation.cancelled() => break,
321                                result = async { components.lock().await.control(command).await } => result,
322                            };
323                            match result {
324                                Ok(status) => send_event(&output_tx, ServerEvent::DaemonStatus { request_id, status }, server_frame_limit(&session)).await?,
325                                Err(_) => send_error(&output_tx, &request_id, "component_error", "Component operation failed; check account credentials, private file permissions and absolute workspace", false, server_frame_limit(&session)).await?,
326                            }
327                        } else {
328                            send_error(&output_tx, &request_id, "unsupported", "Component management requires the daemon socket", false, server_frame_limit(&session)).await?;
329                        }
330                    }
331                    ClientMessage::SessionStart { request_id, cwd, provider, model, base_url, no_tools } => {
332                        if session.is_some() {
333                            send_error(&output_tx, &request_id, "invalid_request", "this connection already has a session", false, server_frame_limit(&session)).await?;
334                            continue;
335                        }
336                        let session_overrides = ConfigOverrides {
337                            provider: provider.or_else(|| overrides.provider.clone()),
338                            model: model.or_else(|| overrides.model.clone()),
339                            base_url: base_url.or_else(|| overrides.base_url.clone()),
340                            approval_policy: overrides.approval_policy,
341                            no_tools: no_tools.unwrap_or(overrides.no_tools),
342                        };
343                        match build_session(&cwd, session_overrides).await {
344                            Ok(new_session) => {
345                                output_tx.ensure_capacity(output_queue_bytes(
346                                    new_session.config.protocol.max_server_frame_bytes,
347                                )?)?;
348                                let event = ServerEvent::SessionStarted {
349                                    request_id,
350                                    session_id: new_session.id.clone(),
351                                    cwd: new_session.workspace.display().to_string(),
352                                    model: new_session.runtime.model().to_owned(),
353                                    context_max_tokens: new_session.config.context.max_tokens,
354                                    max_server_frame_bytes: new_session.config.protocol.max_server_frame_bytes,
355                                    max_transcript_bytes: new_session.config.tui.max_transcript_bytes,
356                                    max_transcript_items: new_session.config.tui.max_transcript_items,
357                                    max_prompt_history_bytes: new_session.config.tui.max_prompt_history_bytes,
358                                    max_prompt_history_items: new_session.config.tui.max_prompt_history_items,
359                                };
360                                send_event(&output_tx, event, new_session.config.protocol.max_server_frame_bytes).await?;
361                                send_event(&output_tx, ServerEvent::QueueSnapshot {
362                                    request_id: None,
363                                    session_id: new_session.id.clone(),
364                                    seq: next_seq(&new_session.seq),
365                                    entries: new_session.queue.lock().await.iter().cloned().collect(),
366                                    paused: new_session.paused.load(Ordering::Acquire),
367                                }, new_session.config.protocol.max_server_frame_bytes).await?;
368                                session = Some(new_session);
369                            }
370                            Err(error) => {
371                                send_error(&output_tx, &request_id, "invalid_request", &error.to_string(), false, server_frame_limit(&session)).await?;
372                            }
373                        }
374                    }
375                    ClientMessage::SessionAttach { request_id, .. } => {
376                        send_error(&output_tx, &request_id, "unsupported", "session attach requires the shared socket server", false, server_frame_limit(&session)).await?;
377                    }
378                    ClientMessage::TurnStart { request_id, session_id, prompt } => {
379                        let Some(current) = session.as_ref() else {
380                            send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?;
381                            continue;
382                        };
383                        if current.id != session_id {
384                            send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
385                            continue;
386                        }
387                        if prompt.trim().is_empty() || prompt.len() > PROMPT_LIMIT_BYTES {
388                            send_error(&output_tx, &request_id, "invalid_request", "prompt must be non-empty and no larger than 256 KiB", false, server_frame_limit(&session)).await?;
389                            continue;
390                        }
391                        if active.is_some() {
392                            let entry = match current.enqueue(prompt, request_id.clone()).await {
393                                Ok(entry) => entry,
394                                Err(code) => { send_error(&output_tx, &request_id, code, "session queue limit reached", false, server_frame_limit(&session)).await?; continue; }
395                            };
396                            let position = current.queue.lock().await.len().saturating_sub(1);
397                            send_event(&output_tx, ServerEvent::QueueEnqueued {
398                                request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), entry, position,
399                            }, current.config.protocol.max_server_frame_bytes).await?;
400                            continue;
401                        }
402                        let turn_id = Uuid::new_v4().to_string();
403                        let cancellation = cancellation.child_token();
404                        let meta = TurnMeta {
405                            request_id: request_id.clone(),
406                            session_id: current.id.clone(),
407                            turn_id: turn_id.clone(),
408                            seq: Arc::clone(&current.seq),
409                            max_server_frame: current.config.protocol.max_server_frame_bytes,
410                        };
411                        send_event(&output_tx, ServerEvent::TurnStarted {
412                            request_id: request_id.clone(),
413                            session_id: current.id.clone(),
414                            turn_id: turn_id.clone(),
415                            seq: next_seq(&current.seq),
416                        }, current.config.protocol.max_server_frame_bytes).await?;
417                        let runtime = Arc::clone(&current.runtime);
418                        let history = Arc::clone(&current.history);
419                        let sink: Arc<dyn EventSink> = Arc::new(ProtocolSink {
420                            meta: meta.clone(),
421                            output: output_tx.clone(),
422                            cancellation: cancellation.clone(),
423                        });
424                        let gate: Arc<dyn ApprovalGate> = Arc::new(ProtocolApprovalGate {
425                            policy: current.config.tools.approval_policy,
426                            broker: Arc::clone(&approvals),
427                            meta,
428                            output: output_tx.clone(),
429                        });
430                        let task_cancel = cancellation.clone();
431                        let task_done = done_tx.clone();
432                        let task_request = request_id.clone();
433                        let task_session = current.id.clone();
434                        let task_turn = turn_id.clone();
435                        let task = tasks.spawn(async move {
436                            let mut history = history.lock().await;
437                            let result = runtime.run_turn(&mut history, prompt, sink, gate, task_cancel).await;
438                            let _ = task_done.send(TurnDone {
439                                request_id: task_request,
440                                session_id: task_session,
441                                turn_id: task_turn,
442                                result,
443                            }).await;
444                        });
445                        active = Some(ActiveTurn { turn_id, cancellation, task });
446                    }
447                    ClientMessage::QueueUpdate { request_id, session_id, queue_id, revision, prompt } => {
448                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
449                        if current.id != session_id { send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?; continue; }
450                        if prompt.trim().is_empty() || prompt.len() > PROMPT_LIMIT_BYTES { send_error(&output_tx, &request_id, "invalid_request", "prompt must be non-empty and no larger than 256 KiB", false, server_frame_limit(&session)).await?; continue; }
451                        match current.update_queue(&queue_id, revision, prompt).await {
452                            Ok(entry) => send_event(&output_tx, ServerEvent::QueueUpdated { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), entry }, current.config.protocol.max_server_frame_bytes).await?,
453                            Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
454                        }
455                    }
456                    ClientMessage::QueueMove { request_id, session_id, queue_id, revision, before_queue_id } => {
457                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
458                        match current.move_queue(&session_id, &queue_id, revision, before_queue_id).await {
459                            Ok((id, rev, pos)) => send_event(&output_tx, ServerEvent::QueueMoved { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), queue_id: id, position: pos, revision: rev }, current.config.protocol.max_server_frame_bytes).await?,
460                            Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
461                        }
462                    }
463                    ClientMessage::QueueRemove { request_id, session_id, queue_id, revision } => {
464                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
465                        match current.remove_queue(&session_id, &queue_id, revision).await {
466                            Ok((id, rev)) => send_event(&output_tx, ServerEvent::QueueRemoved { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), queue_id: id, revision: rev }, current.config.protocol.max_server_frame_bytes).await?,
467                            Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
468                        }
469                    }
470                    ClientMessage::SessionPause { request_id, session_id, paused } => {
471                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
472                        if current.id != session_id { send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?; continue; }
473                        current.paused.store(paused, Ordering::Release);
474                        send_event(&output_tx, ServerEvent::SessionPaused { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), paused }, current.config.protocol.max_server_frame_bytes).await?;
475                    }
476                    ClientMessage::TurnCancel { request_id, session_id, turn_id } => {
477                        match (&session, &active) {
478                            (Some(current), Some(running)) if current.id == session_id && running.turn_id == turn_id => running.cancellation.cancel(),
479                            _ => send_error(&output_tx, &request_id, "turn_not_found", "active turn was not found", false, server_frame_limit(&session)).await?,
480                        }
481                    }
482                    ClientMessage::ApprovalResolve { request_id, session_id, approval_id, approved } => {
483                        if session.as_ref().is_none_or(|current| current.id != session_id) {
484                            send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
485                        } else if !approvals.resolve(&approval_id, approved).await {
486                            send_error(&output_tx, &request_id, "approval_not_found", "approval was not found or already resolved", false, server_frame_limit(&session)).await?;
487                        }
488                    }
489                    ClientMessage::SessionClear { request_id, session_id } => {
490                        let Some(current) = session.as_ref() else {
491                            send_error(&output_tx, &request_id, "session_not_found", "session was not found", false, server_frame_limit(&session)).await?;
492                            continue;
493                        };
494                        if current.id != session_id {
495                            send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
496                        } else if active.is_some() {
497                            send_error(&output_tx, &request_id, "turn_active", "cancel the active turn before clearing", false, server_frame_limit(&session)).await?;
498                        } else {
499                            current.history.lock().await.clear();
500                            current.queue.lock().await.clear();
501                            send_event(&output_tx, ServerEvent::SessionCleared {
502                                request_id,
503                                session_id: current.id.clone(),
504                                seq: next_seq(&current.seq),
505                            }, current.config.protocol.max_server_frame_bytes).await?;
506                            send_event(&output_tx, ServerEvent::QueueSnapshot { request_id: None, session_id: current.id.clone(), seq: next_seq(&current.seq), entries: Vec::new(), paused: current.paused.load(Ordering::Acquire) }, current.config.protocol.max_server_frame_bytes).await?;
507                        }
508                    }
509                }
510            }
511            writer = &mut writer_task => {
512                writer_finished = true;
513                writer.context("join protocol writer")??;
514                break;
515            }
516            done = done_rx.recv(), if active.is_some() => {
517                if let Some(done) = done {
518                    if let Some(current) = session.as_ref() {
519                        let seq = next_seq(&current.seq);
520                        let event = match done.result {
521                            Ok(outcome) => ServerEvent::TurnCompleted {
522                                request_id: done.request_id,
523                                session_id: done.session_id,
524                                turn_id: done.turn_id,
525                                seq,
526                                steps: outcome.steps,
527                                usage: Usage { input_tokens: outcome.usage.input_tokens, output_tokens: outcome.usage.output_tokens },
528                            },
529                            Err(AgentError::Cancelled) => ServerEvent::TurnCancelled {
530                                request_id: done.request_id,
531                                session_id: done.session_id,
532                                turn_id: done.turn_id,
533                                seq,
534                            },
535                            Err(error) => ServerEvent::TurnFailed {
536                                request_id: done.request_id,
537                                session_id: done.session_id,
538                                turn_id: done.turn_id,
539                                seq,
540                                code: error.code().into(),
541                                message: error.to_string(),
542                            },
543                        };
544                        send_event(&output_tx, event, current.config.protocol.max_server_frame_bytes).await?;
545                    }
546                    if let Some(mut active) = active.take() {
547                        let _ = (&mut active.task).await;
548                    }
549                    if let Some(current) = session.as_ref()
550                        && !current.paused.load(Ordering::Acquire)
551                        && let Some(entry) = current.queue.lock().await.pop_front()
552                    {
553                        let turn_id = Uuid::new_v4().to_string();
554                        let cancellation = cancellation.child_token();
555                        send_event(&output_tx, ServerEvent::QueueDequeued {
556                            request_id: entry.submitter.clone(),
557                            session_id: current.id.clone(),
558                            seq: next_seq(&current.seq),
559                            queue_id: entry.queue_id,
560                            turn_id: turn_id.clone(),
561                        }, current.config.protocol.max_server_frame_bytes).await?;
562                        send_event(&output_tx, ServerEvent::TurnStarted {
563                            request_id: entry.submitter.clone(),
564                            session_id: current.id.clone(),
565                            turn_id: turn_id.clone(),
566                            seq: next_seq(&current.seq),
567                        }, current.config.protocol.max_server_frame_bytes).await?;
568                        let meta = TurnMeta {
569                            request_id: entry.submitter.clone(), session_id: current.id.clone(), turn_id: turn_id.clone(),
570                            seq: Arc::clone(&current.seq), max_server_frame: current.config.protocol.max_server_frame_bytes,
571                        };
572                        let sink: Arc<dyn EventSink> = Arc::new(ProtocolSink { meta: meta.clone(), output: output_tx.clone(), cancellation: cancellation.clone() });
573                        let gate: Arc<dyn ApprovalGate> = Arc::new(ProtocolApprovalGate { policy: current.config.tools.approval_policy, broker: Arc::clone(&approvals), meta, output: output_tx.clone() });
574                        let runtime = Arc::clone(&current.runtime);
575                        let history = Arc::clone(&current.history);
576                        let task_done = done_tx.clone();
577                        let task_request = entry.submitter;
578                        let task_session = current.id.clone();
579                        let task_turn = turn_id.clone();
580                        let task_cancel = cancellation.clone();
581                        let task = tasks.spawn(async move {
582                            let mut history = history.lock().await;
583                            let result = runtime.run_turn(&mut history, entry.prompt, sink, gate, task_cancel).await;
584                            let _ = task_done.send(TurnDone { request_id: task_request, session_id: task_session, turn_id: task_turn, result }).await;
585                        });
586                        active = Some(ActiveTurn { turn_id, cancellation, task });
587                    }
588                }
589            }
590        }
591        }
592        Ok(())
593    }
594    .await;
595
596    if let Some(active) = active.take() {
597        shutdown_active_turn(active, SHUTDOWN_GRACE).await;
598    }
599    drop(output_tx);
600    let writer_result = if writer_finished {
601        Ok(())
602    } else {
603        shutdown_writer(writer_task, SHUTDOWN_GRACE).await
604    };
605    loop_result?;
606    writer_result?;
607    if fatal {
608        return Err(anyhow!("protocol version mismatch"));
609    }
610    Ok(())
611}
612
613enum FrameRead {
614    Eof,
615    Frame(Vec<u8>),
616    TooLarge,
617}
618
619struct OutboundFrame {
620    bytes: Vec<u8>,
621    _byte_permit: OwnedSemaphorePermit,
622}
623
624#[derive(Clone)]
625struct OutboundSender {
626    frames: mpsc::Sender<OutboundFrame>,
627    budget: Arc<Semaphore>,
628    capacity: Arc<AtomicUsize>,
629}
630
631#[derive(Debug, PartialEq, Eq)]
632enum OutboundSendError {
633    Cancelled,
634    Closed,
635    TimedOut,
636    FrameExceedsQueue { frame_bytes: usize, capacity: usize },
637}
638
639impl std::fmt::Display for OutboundSendError {
640    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
641        match self {
642            Self::Cancelled => formatter.write_str("outbound send cancelled"),
643            Self::Closed => formatter.write_str("protocol client disconnected"),
644            Self::TimedOut => formatter.write_str("outbound send timed out under backpressure"),
645            Self::FrameExceedsQueue {
646                frame_bytes,
647                capacity,
648            } => write!(
649                formatter,
650                "outbound frame uses {frame_bytes} bytes but queue capacity is {capacity} bytes"
651            ),
652        }
653    }
654}
655
656impl std::error::Error for OutboundSendError {}
657
658fn outbound_channel(capacity: usize) -> (OutboundSender, mpsc::Receiver<OutboundFrame>) {
659    let (frames, receiver) = mpsc::channel(OUTPUT_QUEUE_CAPACITY);
660    (
661        OutboundSender {
662            frames,
663            budget: Arc::new(Semaphore::new(capacity)),
664            capacity: Arc::new(AtomicUsize::new(capacity)),
665        },
666        receiver,
667    )
668}
669
670impl OutboundSender {
671    fn ensure_capacity(&self, required: usize) -> Result<()> {
672        if required > Semaphore::MAX_PERMITS {
673            return Err(anyhow!(
674                "outbound queue capacity {required} exceeds runtime limit {}",
675                Semaphore::MAX_PERMITS
676            ));
677        }
678        let current = self.capacity.load(Ordering::Acquire);
679        if required > current {
680            self.budget.add_permits(required - current);
681            self.capacity.store(required, Ordering::Release);
682        }
683        Ok(())
684    }
685
686    async fn send(
687        &self,
688        bytes: Vec<u8>,
689        cancellation: Option<&CancellationToken>,
690    ) -> std::result::Result<(), OutboundSendError> {
691        self.send_with_timeout(bytes, cancellation, SHUTDOWN_GRACE)
692            .await
693    }
694
695    async fn send_with_timeout(
696        &self,
697        bytes: Vec<u8>,
698        cancellation: Option<&CancellationToken>,
699        control_timeout: Duration,
700    ) -> std::result::Result<(), OutboundSendError> {
701        let frame_bytes =
702            bytes
703                .len()
704                .checked_add(1)
705                .ok_or(OutboundSendError::FrameExceedsQueue {
706                    frame_bytes: usize::MAX,
707                    capacity: self.capacity.load(Ordering::Acquire),
708                })?;
709        let capacity = self.capacity.load(Ordering::Acquire);
710        let permits =
711            u32::try_from(frame_bytes).map_err(|_| OutboundSendError::FrameExceedsQueue {
712                frame_bytes,
713                capacity,
714            })?;
715        if frame_bytes > capacity {
716            return Err(OutboundSendError::FrameExceedsQueue {
717                frame_bytes,
718                capacity,
719            });
720        }
721
722        let control_deadline = tokio::time::Instant::now() + control_timeout;
723        let acquire = Arc::clone(&self.budget).acquire_many_owned(permits);
724        let permit = if let Some(cancellation) = cancellation {
725            tokio::select! {
726                biased;
727                _ = cancellation.cancelled() => return Err(OutboundSendError::Cancelled),
728                permit = acquire => permit.map_err(|_| OutboundSendError::Closed)?,
729            }
730        } else {
731            tokio::time::timeout_at(control_deadline, acquire)
732                .await
733                .map_err(|_| OutboundSendError::TimedOut)?
734                .map_err(|_| OutboundSendError::Closed)?
735        };
736        let frame = OutboundFrame {
737            bytes,
738            _byte_permit: permit,
739        };
740        if let Some(cancellation) = cancellation {
741            tokio::select! {
742                biased;
743                _ = cancellation.cancelled() => Err(OutboundSendError::Cancelled),
744                result = self.frames.send(frame) => result.map_err(|_| OutboundSendError::Closed),
745            }
746        } else {
747            tokio::time::timeout_at(control_deadline, self.frames.send(frame))
748                .await
749                .map_err(|_| OutboundSendError::TimedOut)?
750                .map_err(|_| OutboundSendError::Closed)
751        }
752    }
753}
754
755fn output_queue_bytes(max_frame_bytes: usize) -> Result<usize> {
756    let required = max_frame_bytes
757        .checked_add(1)
758        .and_then(|bytes| bytes.checked_mul(2))
759        .ok_or_else(|| anyhow!("configured server frame limit is too large"))?
760        .max(OUTPUT_QUEUE_MIN_BYTES);
761    if required > Semaphore::MAX_PERMITS {
762        return Err(anyhow!(
763            "configured server frame limit requires an outbound queue larger than the runtime supports"
764        ));
765    }
766    Ok(required)
767}
768
769/// A persistent decoder keeps consumed partial bytes across select cancellation.
770#[derive(Default)]
771struct FrameBuffer {
772    bytes: Vec<u8>,
773    oversized: bool,
774}
775
776impl FrameBuffer {
777    async fn read<R>(&mut self, reader: &mut R, max_bytes: usize) -> std::io::Result<FrameRead>
778    where
779        R: AsyncBufRead + Unpin,
780    {
781        loop {
782            let available = reader.fill_buf().await?;
783            let eof = available.is_empty();
784            let end = available.iter().position(|b| *b == b'\n');
785            let take = end.map_or(available.len(), |n| n + 1);
786            if !self.oversized {
787                if self.bytes.len().saturating_add(take) > max_bytes.saturating_add(2) {
788                    self.oversized = true;
789                    self.bytes.clear();
790                } else {
791                    self.bytes.extend_from_slice(&available[..take]);
792                }
793            }
794            reader.consume(take);
795            if end.is_some() || eof {
796                if std::mem::take(&mut self.oversized) {
797                    return Ok(FrameRead::TooLarge);
798                }
799                if eof && self.bytes.is_empty() {
800                    return Ok(FrameRead::Eof);
801                }
802                let mut bytes = std::mem::take(&mut self.bytes);
803                while matches!(bytes.last(), Some(b'\n' | b'\r')) {
804                    bytes.pop();
805                }
806                return Ok(if bytes.len() > max_bytes {
807                    FrameRead::TooLarge
808                } else {
809                    FrameRead::Frame(bytes)
810                });
811            }
812        }
813    }
814}
815
816#[cfg(test)]
817async fn read_bounded_frame<R: AsyncBufRead + Unpin>(
818    reader: &mut R,
819    max_bytes: usize,
820) -> std::io::Result<FrameRead> {
821    FrameBuffer::default().read(reader, max_bytes).await
822}
823
824/// A persistent advisory lock closes the stale-socket unlink/bind race.
825struct SocketLock(std::fs::File);
826impl SocketLock {
827    fn acquire(socket: &Path) -> Result<Self> {
828        use std::os::unix::{fs::OpenOptionsExt, io::AsRawFd};
829        let file = std::fs::OpenOptions::new()
830            .read(true)
831            .write(true)
832            .create(true)
833            .truncate(false)
834            .mode(0o600)
835            .custom_flags(libc::O_NOFOLLOW)
836            .open(socket.with_extension("lock"))?;
837        // SAFETY: flock operates on this owned, live file descriptor.
838        if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } != 0 {
839            return Err(anyhow!("SCV daemon already owns this socket"));
840        }
841        Ok(Self(file))
842    }
843}
844impl Drop for SocketLock {
845    fn drop(&mut self) {
846        use std::os::unix::io::AsRawFd;
847        // SAFETY: the descriptor remains live until this drop returns.
848        unsafe {
849            libc::flock(self.0.as_raw_fd(), libc::LOCK_UN);
850        }
851    }
852}
853
854fn server_frame_limit(session: &Option<Session>) -> usize {
855    session.as_ref().map_or_else(
856        || Config::default().protocol.max_server_frame_bytes,
857        |value| value.config.protocol.max_server_frame_bytes,
858    )
859}
860
861struct Session {
862    id: String,
863    workspace: PathBuf,
864    config: Config,
865    runtime: Arc<AgentRuntime>,
866    history: Arc<Mutex<Vec<Message>>>,
867    seq: Arc<AtomicU64>,
868    queue: Arc<Mutex<VecDeque<QueueEntry>>>,
869    paused: Arc<std::sync::atomic::AtomicBool>,
870}
871
872impl Session {
873    async fn enqueue(
874        &self,
875        prompt: String,
876        submitter: String,
877    ) -> std::result::Result<QueueEntry, &'static str> {
878        let entry = QueueEntry {
879            queue_id: Uuid::new_v4().to_string(),
880            revision: 1,
881            prompt,
882            submitter,
883        };
884        let mut queue = self.queue.lock().await;
885        let bytes: usize = queue.iter().map(|item| item.prompt.len()).sum();
886        if queue.len() >= MAX_QUEUE_ITEMS
887            || bytes.saturating_add(entry.prompt.len()) > MAX_QUEUE_BYTES
888        {
889            return Err("queue_limit");
890        }
891        queue.push_back(entry.clone());
892        Ok(entry)
893    }
894
895    async fn update_queue(
896        &self,
897        id: &str,
898        revision: u64,
899        prompt: String,
900    ) -> std::result::Result<QueueEntry, &'static str> {
901        let mut queue = self.queue.lock().await;
902        let bytes: usize = queue.iter().map(|item| item.prompt.len()).sum();
903        let entry = queue
904            .iter_mut()
905            .find(|entry| entry.queue_id == id)
906            .ok_or("queue_not_found")?;
907        if entry.revision != revision {
908            return Err("queue_conflict");
909        }
910        if bytes
911            .saturating_sub(entry.prompt.len())
912            .saturating_add(prompt.len())
913            > MAX_QUEUE_BYTES
914        {
915            return Err("queue_limit");
916        }
917        entry.prompt = prompt;
918        entry.revision += 1;
919        Ok(entry.clone())
920    }
921
922    async fn move_queue(
923        &self,
924        session_id: &str,
925        id: &str,
926        revision: u64,
927        before: Option<String>,
928    ) -> std::result::Result<(String, u64, usize), &'static str> {
929        if self.id != session_id {
930            return Err("session_not_found");
931        }
932        let mut queue = self.queue.lock().await;
933        let index = queue
934            .iter()
935            .position(|entry| entry.queue_id == id)
936            .ok_or("queue_not_found")?;
937        if queue[index].revision != revision {
938            return Err("queue_conflict");
939        }
940        // Validate the destination while the source is still present. This keeps
941        // the operation atomic and handles a self move as a no-op reorder.
942        let target_index = match before.as_deref() {
943            Some(target) if target == id => return Ok((id.to_string(), revision, index)),
944            Some(target) => Some(
945                queue
946                    .iter()
947                    .position(|item| item.queue_id == target)
948                    .ok_or("queue_not_found")?,
949            ),
950            None => None,
951        };
952        let mut entry = queue.remove(index).expect("queue index exists");
953        let target = target_index.map_or(queue.len(), |target| {
954            target.saturating_sub(usize::from(target > index))
955        });
956        let pos = target.min(queue.len());
957        let id = entry.queue_id.clone();
958        let rev = entry.revision + 1;
959        entry.revision = rev;
960        queue.insert(pos, entry);
961        Ok((id, rev, pos))
962    }
963
964    async fn remove_queue(
965        &self,
966        session_id: &str,
967        id: &str,
968        revision: u64,
969    ) -> std::result::Result<(String, u64), &'static str> {
970        if self.id != session_id {
971            return Err("session_not_found");
972        }
973        let mut queue = self.queue.lock().await;
974        let index = queue
975            .iter()
976            .position(|entry| entry.queue_id == id)
977            .ok_or("queue_not_found")?;
978        if queue[index].revision != revision {
979            return Err("queue_conflict");
980        }
981        let entry = queue.remove(index).expect("queue index exists");
982        Ok((entry.queue_id, entry.revision))
983    }
984}
985
986struct ActiveTurn {
987    turn_id: String,
988    cancellation: CancellationToken,
989    task: JoinHandle<()>,
990}
991
992impl Drop for ActiveTurn {
993    fn drop(&mut self) {
994        self.cancellation.cancel();
995        self.task.abort();
996    }
997}
998
999struct AbortGuard(tokio::task::AbortHandle);
1000impl Drop for AbortGuard {
1001    fn drop(&mut self) {
1002        self.0.abort();
1003    }
1004}
1005
1006async fn shutdown_active_turn(mut active: ActiveTurn, grace: Duration) -> bool {
1007    active.cancellation.cancel();
1008    if tokio::time::timeout(grace, &mut active.task).await.is_ok() {
1009        true
1010    } else {
1011        active.task.abort();
1012        let _ = (&mut active.task).await;
1013        false
1014    }
1015}
1016
1017async fn shutdown_writer(
1018    mut writer: JoinHandle<std::io::Result<()>>,
1019    grace: Duration,
1020) -> Result<()> {
1021    match tokio::time::timeout(grace, &mut writer).await {
1022        Ok(result) => {
1023            result.context("join protocol writer")??;
1024            Ok(())
1025        }
1026        Err(_) => {
1027            writer.abort();
1028            let _ = writer.await;
1029            Err(anyhow!("protocol writer shutdown timed out"))
1030        }
1031    }
1032}
1033
1034struct TurnDone {
1035    request_id: String,
1036    session_id: String,
1037    turn_id: String,
1038    result: Result<scv_core::TurnOutcome, AgentError>,
1039}
1040
1041#[derive(Clone)]
1042struct TurnMeta {
1043    request_id: String,
1044    session_id: String,
1045    turn_id: String,
1046    seq: Arc<AtomicU64>,
1047    max_server_frame: usize,
1048}
1049
1050fn next_seq(sequence: &AtomicU64) -> u64 {
1051    sequence.fetch_add(1, Ordering::Relaxed) + 1
1052}
1053
1054async fn build_session(cwd: &str, overrides: ConfigOverrides) -> Result<Session> {
1055    let workspace = std::fs::canonicalize(cwd).with_context(|| format!("resolve cwd {cwd}"))?;
1056    if !workspace.is_dir() {
1057        return Err(anyhow!("cwd is not a directory"));
1058    }
1059    let no_tools = overrides.no_tools;
1060    let config = Config::load(&workspace, overrides)?;
1061    if !no_tools {
1062        config.prepare_adapter_homes()?;
1063    }
1064    let provider_config = config.provider.clone();
1065    let api_key = provider_config.api_key.clone().or_else(|| {
1066        provider_config.api_key_env.as_deref().and_then(|name| std::env::var(name).ok())
1067    }).filter(|key| !key.trim().is_empty()).ok_or_else(|| anyhow!("provider credential is not configured; set provider.api_key or provider.api_key_env"))?;
1068    let (skills, skill_roots, skill_prompt) = discover_skills(&workspace, &config)?;
1069    let system_prompt = build_system_prompt(&workspace, &config, &skill_prompt)?;
1070    let provider = Arc::new(OpenAiProvider::new(
1071        provider_config.model.clone(),
1072        provider_config.base_url.clone(),
1073        api_key,
1074        Duration::from_secs(provider_config.timeout_seconds),
1075        config.provider_limits(),
1076        provider_config.headers.clone(),
1077    )?);
1078    let tools = if no_tools {
1079        Arc::new(ToolRegistry::default())
1080    } else {
1081        Arc::new(builtin_registry(
1082            config.tools(),
1083            skills,
1084            skill_roots,
1085            config.skills.max_skill_bytes,
1086            config.adapters(),
1087        )?)
1088    };
1089    let context = Arc::new(BudgetContextPolicy::new((&config.context).into())?);
1090    let runtime = Arc::new(AgentRuntime::new(
1091        provider,
1092        tools,
1093        context,
1094        config.core_agent(system_prompt),
1095        workspace.clone(),
1096    ));
1097    Ok(Session {
1098        id: Uuid::new_v4().to_string(),
1099        workspace,
1100        config,
1101        runtime,
1102        history: Arc::new(Mutex::new(Vec::new())),
1103        seq: Arc::new(AtomicU64::new(0)),
1104        queue: Arc::new(Mutex::new(VecDeque::new())),
1105        paused: Arc::new(std::sync::atomic::AtomicBool::new(false)),
1106    })
1107}
1108
1109fn build_system_prompt(workspace: &Path, config: &Config, skills: &str) -> Result<String> {
1110    let mut prompt = config.agent.system_prompt.clone();
1111    prompt.push_str(&format!(
1112        "\nCurrent working directory: {}\n",
1113        workspace.display()
1114    ));
1115    let agents_path = workspace.join("AGENTS.md");
1116    if agents_path.is_file() {
1117        let canonical = std::fs::canonicalize(&agents_path).context("resolve project AGENTS.md")?;
1118        if !canonical.starts_with(workspace) {
1119            return Err(anyhow!("project AGENTS.md escaped workspace"));
1120        }
1121        let (bytes, truncated) = read_prefix(&canonical, config.tools.max_read_bytes)
1122            .context("read project AGENTS.md")?;
1123        let instructions = std::str::from_utf8(&bytes).context("project AGENTS.md is not UTF-8")?;
1124        prompt.push_str("\n# Project instructions\n");
1125        prompt.push_str(instructions);
1126        if truncated {
1127            prompt.push_str("\n[AGENTS.md truncated by configured read limit]\n");
1128        }
1129    }
1130    if !skills.is_empty() {
1131        prompt.push_str("\n# Available skills\n");
1132        prompt.push_str(skills);
1133        prompt.push_str("\nUse read_skill with a skill name when its workflow applies.\n");
1134    }
1135    Ok(prompt)
1136}
1137
1138fn discover_skills(workspace: &Path, config: &Config) -> Result<(SkillMap, Vec<PathBuf>, String)> {
1139    let mut skills = SkillMap::new();
1140    let mut roots = Vec::new();
1141    let project_root = workspace.join(&config.skills.project_dir);
1142    for (root, must_be_workspace) in [(&project_root, true), (&config.skills.user_dir, false)] {
1143        if !root.is_dir() {
1144            continue;
1145        }
1146        let canonical = std::fs::canonicalize(root)
1147            .with_context(|| format!("resolve skill root {}", root.display()))?;
1148        if must_be_workspace && !canonical.starts_with(workspace) {
1149            return Err(anyhow!("project skill root escaped workspace"));
1150        }
1151        roots.push(canonical.clone());
1152        let mut entries: Vec<_> = std::fs::read_dir(&canonical)
1153            .with_context(|| format!("read skill root {}", canonical.display()))?
1154            .filter_map(Result::ok)
1155            .collect();
1156        entries.sort_by_key(|entry| entry.file_name());
1157        for entry in entries {
1158            if skills.len() >= config.skills.max_skills {
1159                break;
1160            }
1161            let path = entry.path().join("SKILL.md");
1162            if !path.is_file() {
1163                continue;
1164            }
1165            let canonical_file = std::fs::canonicalize(&path)
1166                .with_context(|| format!("resolve skill {}", path.display()))?;
1167            if !canonical_file.starts_with(&canonical) {
1168                continue;
1169            }
1170            let name = entry.file_name().to_string_lossy().to_string();
1171            skills.entry(name).or_insert(canonical_file);
1172        }
1173    }
1174    let mut names: Vec<_> = skills.keys().cloned().collect();
1175    names.sort();
1176    let mut listing = String::new();
1177    for name in names {
1178        let path = &skills[&name];
1179        let bytes = read_prefix(path, config.skills.max_skill_bytes)
1180            .map(|(bytes, _)| bytes)
1181            .unwrap_or_default();
1182        let content = String::from_utf8_lossy(&bytes);
1183        let description = skill_description(&content);
1184        listing.push_str(&format!("- {name}: {description}\n"));
1185    }
1186    Ok((skills, roots, listing))
1187}
1188
1189fn read_prefix(path: &Path, max_bytes: usize) -> std::io::Result<(Vec<u8>, bool)> {
1190    let file = std::fs::File::open(path)?;
1191    let mut bytes = Vec::with_capacity(max_bytes.min(8192));
1192    file.take(
1193        u64::try_from(max_bytes)
1194            .unwrap_or(u64::MAX)
1195            .saturating_add(1),
1196    )
1197    .read_to_end(&mut bytes)?;
1198    let truncated = bytes.len() > max_bytes;
1199    bytes.truncate(max_bytes);
1200    Ok((bytes, truncated))
1201}
1202
1203fn skill_description(content: &str) -> String {
1204    if let Some(frontmatter) = content.strip_prefix("---\n")
1205        && let Some((header, _)) = frontmatter.split_once("\n---")
1206    {
1207        for line in header.lines() {
1208            if let Some(description) = line.strip_prefix("description:") {
1209                return description.trim().trim_matches('"').to_owned();
1210            }
1211        }
1212    }
1213    content
1214        .lines()
1215        .map(str::trim)
1216        .find(|line| !line.is_empty() && !line.starts_with('#'))
1217        .unwrap_or("No description provided")
1218        .chars()
1219        .take(240)
1220        .collect()
1221}
1222
1223struct ProtocolSink {
1224    meta: TurnMeta,
1225    output: OutboundSender,
1226    cancellation: CancellationToken,
1227}
1228
1229#[async_trait]
1230impl EventSink for ProtocolSink {
1231    async fn emit(&self, event: CoreEvent) -> Result<(), AgentError> {
1232        let seq = next_seq(&self.meta.seq);
1233        let event = match event {
1234            CoreEvent::AssistantDelta { content } => ServerEvent::AssistantDelta {
1235                request_id: self.meta.request_id.clone(),
1236                session_id: self.meta.session_id.clone(),
1237                turn_id: self.meta.turn_id.clone(),
1238                seq,
1239                content,
1240            },
1241            CoreEvent::AssistantCompleted { content } => ServerEvent::AssistantCompleted {
1242                request_id: self.meta.request_id.clone(),
1243                session_id: self.meta.session_id.clone(),
1244                turn_id: self.meta.turn_id.clone(),
1245                seq,
1246                content,
1247            },
1248            CoreEvent::ToolProposed {
1249                call_id,
1250                name,
1251                arguments,
1252            } => ServerEvent::ToolProposed {
1253                request_id: self.meta.request_id.clone(),
1254                session_id: self.meta.session_id.clone(),
1255                turn_id: self.meta.turn_id.clone(),
1256                seq,
1257                call_id,
1258                name,
1259                arguments,
1260            },
1261            CoreEvent::ToolStarted { call_id, name } => ServerEvent::ToolStarted {
1262                request_id: self.meta.request_id.clone(),
1263                session_id: self.meta.session_id.clone(),
1264                turn_id: self.meta.turn_id.clone(),
1265                seq,
1266                call_id,
1267                name,
1268            },
1269            CoreEvent::ToolCompleted {
1270                call_id,
1271                name,
1272                output,
1273            } => ServerEvent::ToolCompleted {
1274                request_id: self.meta.request_id.clone(),
1275                session_id: self.meta.session_id.clone(),
1276                turn_id: self.meta.turn_id.clone(),
1277                seq,
1278                call_id,
1279                name,
1280                success: !output.is_error,
1281                output: output.content,
1282                truncated: output.truncated,
1283            },
1284            CoreEvent::ContextCompacted {
1285                before_tokens,
1286                after_tokens,
1287                removed_messages,
1288            } => ServerEvent::ContextCompacted {
1289                request_id: self.meta.request_id.clone(),
1290                session_id: self.meta.session_id.clone(),
1291                turn_id: self.meta.turn_id.clone(),
1292                seq,
1293                before_tokens,
1294                after_tokens,
1295                removed_messages,
1296            },
1297            CoreEvent::SessionTrimmed {
1298                removed_messages,
1299                history_bytes,
1300            } => ServerEvent::SessionTrimmed {
1301                request_id: self.meta.request_id.clone(),
1302                session_id: self.meta.session_id.clone(),
1303                seq,
1304                removed_messages,
1305                history_bytes,
1306            },
1307        };
1308        send_turn_event(
1309            &self.output,
1310            event,
1311            self.meta.max_server_frame,
1312            &self.cancellation,
1313        )
1314        .await
1315    }
1316}
1317
1318#[derive(Default)]
1319struct ApprovalBroker {
1320    pending: Mutex<HashMap<String, oneshot::Sender<bool>>>,
1321}
1322
1323impl ApprovalBroker {
1324    async fn insert(&self, id: String, sender: oneshot::Sender<bool>) {
1325        self.pending.lock().await.insert(id, sender);
1326    }
1327
1328    async fn remove(&self, id: &str) {
1329        self.pending.lock().await.remove(id);
1330    }
1331
1332    async fn resolve(&self, id: &str, approved: bool) -> bool {
1333        let sender = self.pending.lock().await.remove(id);
1334        sender.is_some_and(|sender| sender.send(approved).is_ok())
1335    }
1336}
1337
1338struct ProtocolApprovalGate {
1339    policy: ApprovalPolicy,
1340    broker: Arc<ApprovalBroker>,
1341    meta: TurnMeta,
1342    output: OutboundSender,
1343}
1344
1345#[async_trait]
1346impl ApprovalGate for ProtocolApprovalGate {
1347    async fn approve(
1348        &self,
1349        request: ApprovalRequest,
1350        cancellation: CancellationToken,
1351    ) -> Result<bool, AgentError> {
1352        match self.policy {
1353            ApprovalPolicy::OnRisk if request.risk == ToolRisk::ReadOnly => return Ok(true),
1354            ApprovalPolicy::Never => return Ok(request.risk == ToolRisk::ReadOnly),
1355            ApprovalPolicy::Always | ApprovalPolicy::OnRisk => {}
1356        }
1357        let approval_id = Uuid::new_v4().to_string();
1358        let (sender, receiver) = oneshot::channel();
1359        self.broker.insert(approval_id.clone(), sender).await;
1360        let event = ServerEvent::ApprovalRequested {
1361            request_id: self.meta.request_id.clone(),
1362            session_id: self.meta.session_id.clone(),
1363            turn_id: self.meta.turn_id.clone(),
1364            seq: next_seq(&self.meta.seq),
1365            approval_id: approval_id.clone(),
1366            call_id: request.call_id,
1367            name: request.name,
1368            risk: request.risk.as_str().into(),
1369            cwd: request.cwd.display().to_string(),
1370            summary: request.summary,
1371        };
1372        if let Err(error) = send_turn_event(
1373            &self.output,
1374            event,
1375            self.meta.max_server_frame,
1376            &cancellation,
1377        )
1378        .await
1379        {
1380            self.broker.remove(&approval_id).await;
1381            return Err(error);
1382        }
1383        tokio::select! {
1384            result = receiver => result.map_err(|_| AgentError::Cancelled),
1385            _ = cancellation.cancelled() => {
1386                self.broker.remove(&approval_id).await;
1387                Err(AgentError::Cancelled)
1388            }
1389        }
1390    }
1391}
1392
1393async fn send_event(output: &OutboundSender, event: ServerEvent, max_bytes: usize) -> Result<()> {
1394    let bytes = encode_event(&event, max_bytes)?;
1395    output.send(bytes, None).await.map_err(anyhow::Error::new)
1396}
1397
1398async fn send_turn_event(
1399    output: &OutboundSender,
1400    event: ServerEvent,
1401    max_bytes: usize,
1402    cancellation: &CancellationToken,
1403) -> Result<(), AgentError> {
1404    let bytes = encode_event(&event, max_bytes)
1405        .map_err(|error| AgentError::ResponseLimit(error.to_string()))?;
1406    match output.send(bytes, Some(cancellation)).await {
1407        Ok(()) => Ok(()),
1408        Err(OutboundSendError::Cancelled) => Err(AgentError::Cancelled),
1409        Err(error) => Err(AgentError::Internal(error.to_string())),
1410    }
1411}
1412
1413fn encode_event(event: &ServerEvent, max_bytes: usize) -> Result<Vec<u8>> {
1414    let bytes = serde_json::to_vec(event).context("serialize protocol event")?;
1415    if bytes.len() > max_bytes {
1416        return Err(anyhow!("server event exceeds configured frame limit"));
1417    }
1418    Ok(bytes)
1419}
1420
1421async fn send_error(
1422    output: &OutboundSender,
1423    request_id: &str,
1424    code: &str,
1425    message: &str,
1426    fatal: bool,
1427    max_bytes: usize,
1428) -> Result<()> {
1429    send_event(
1430        output,
1431        ServerEvent::Error {
1432            request_id: (!request_id.is_empty()).then(|| request_id.to_owned()),
1433            code: code.into(),
1434            message: message.into(),
1435            fatal,
1436        },
1437        max_bytes,
1438    )
1439    .await
1440}
1441
1442#[cfg(test)]
1443mod tests {
1444    use std::{
1445        future::pending,
1446        io::Cursor,
1447        sync::atomic::{AtomicBool, Ordering},
1448    };
1449
1450    use super::*;
1451
1452    struct DropSignal(Arc<AtomicBool>);
1453
1454    impl Drop for DropSignal {
1455        fn drop(&mut self) {
1456            self.0.store(true, Ordering::Release);
1457        }
1458    }
1459
1460    #[tokio::test]
1461    async fn nonreading_management_client_does_not_hold_component_lock() {
1462        let (mut input, server_input) = tokio::io::duplex(65536);
1463        let (server_output, _blocked_output) = tokio::io::duplex(1);
1464        let tasks = TaskTracker::new();
1465        let components = Arc::new(Mutex::new(components::Components::new(
1466            PathBuf::from("/unused.sock"),
1467            PathBuf::from("/"),
1468        )));
1469        let cancel = CancellationToken::new();
1470        let handler = tokio::spawn(run_managed(
1471            server_input,
1472            server_output,
1473            ConfigOverrides::default(),
1474            Some(components.clone()),
1475            cancel.clone(),
1476            tasks.clone(),
1477        ));
1478        input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
1479        for _ in 0..300 {
1480            input.write_all(b"{\"type\":\"daemon.control\",\"request_id\":\"s\",\"command\":{\"action\":\"status\"}}\n").await.unwrap();
1481        }
1482        tokio::time::sleep(Duration::from_millis(50)).await;
1483        let status = tokio::time::timeout(Duration::from_millis(100), async {
1484            components.lock().await.status()
1485        })
1486        .await
1487        .unwrap();
1488        assert_eq!(status.pid, std::process::id());
1489        cancel.cancel();
1490        handler.abort();
1491        let _ = handler.await;
1492        tasks.close();
1493        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1494            .await
1495            .unwrap();
1496    }
1497
1498    #[tokio::test]
1499    async fn forced_connection_abort_drops_and_joins_writer_descendants() {
1500        let (mut input, server_input) = tokio::io::duplex(512);
1501        let (server_output, _blocked_output) = tokio::io::duplex(1);
1502        let tasks = TaskTracker::new();
1503        let handler = tokio::spawn(run_managed(
1504            server_input,
1505            server_output,
1506            ConfigOverrides::default(),
1507            None,
1508            CancellationToken::new(),
1509            tasks.clone(),
1510        ));
1511        input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
1512        tokio::time::timeout(Duration::from_secs(1), async {
1513            while tasks.is_empty() {
1514                tokio::task::yield_now().await;
1515            }
1516        })
1517        .await
1518        .unwrap();
1519        handler.abort();
1520        let _ = handler.await;
1521        tasks.close();
1522        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1523            .await
1524            .unwrap();
1525        assert!(tasks.is_empty());
1526    }
1527
1528    #[tokio::test]
1529    async fn forced_handler_abort_cancels_and_joins_active_turn() {
1530        let tasks = TaskTracker::new();
1531        let cancellation = CancellationToken::new();
1532        let child_cancel = cancellation.child_token();
1533        let observed_cancel = child_cancel.clone();
1534        let dropped = Arc::new(AtomicBool::new(false));
1535        let (ready_tx, ready_rx) = oneshot::channel();
1536        let task = tasks.spawn({
1537            let dropped = dropped.clone();
1538            async move {
1539                let _guard = DropSignal(dropped);
1540                let _ = ready_tx.send(());
1541                pending::<()>().await;
1542            }
1543        });
1544        ready_rx.await.unwrap();
1545        let (owned_tx, owned_rx) = oneshot::channel();
1546        let handler = tokio::spawn(async move {
1547            let _active = ActiveTurn {
1548                turn_id: "test".into(),
1549                cancellation: child_cancel,
1550                task,
1551            };
1552            let _ = owned_tx.send(());
1553            pending::<()>().await;
1554        });
1555        owned_rx.await.unwrap();
1556        handler.abort();
1557        let _ = handler.await;
1558        tasks.close();
1559        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1560            .await
1561            .unwrap();
1562        assert!(observed_cancel.is_cancelled());
1563        assert!(dropped.load(Ordering::Acquire));
1564    }
1565
1566    #[tokio::test]
1567    async fn frame_buffer_preserves_partial_and_discard_state_across_cancellation() {
1568        let (mut input, output) = tokio::io::duplex(64);
1569        let mut reader = BufReader::new(output);
1570        let mut frames = FrameBuffer::default();
1571        input.write_all(b"12").await.unwrap();
1572        assert!(
1573            tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
1574                .await
1575                .is_err()
1576        );
1577        input.write_all(b"34\n").await.unwrap();
1578        assert!(
1579            matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"1234")
1580        );
1581        input.write_all(b"123456789").await.unwrap();
1582        assert!(
1583            tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
1584                .await
1585                .is_err()
1586        );
1587        input.write_all(b"\n{}\n").await.unwrap();
1588        assert!(matches!(
1589            frames.read(&mut reader, 4).await.unwrap(),
1590            FrameRead::TooLarge
1591        ));
1592        assert!(
1593            matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"{}")
1594        );
1595    }
1596
1597    #[tokio::test]
1598    async fn bounded_reader_discards_an_oversized_line() {
1599        let input = format!("{}\n{{}}\n", "x".repeat(10));
1600        let mut reader = BufReader::new(Cursor::new(input.into_bytes()));
1601        assert!(matches!(
1602            read_bounded_frame(&mut reader, 4).await.unwrap(),
1603            FrameRead::TooLarge
1604        ));
1605        match read_bounded_frame(&mut reader, 4).await.unwrap() {
1606            FrameRead::Frame(frame) => assert_eq!(frame, b"{}"),
1607            _ => panic!("expected the frame following the oversized line"),
1608        }
1609    }
1610
1611    #[tokio::test]
1612    async fn bounded_reader_accepts_exact_crlf_limit() {
1613        let mut reader = BufReader::new(Cursor::new(b"1234\r\n".to_vec()));
1614        match read_bounded_frame(&mut reader, 4).await.unwrap() {
1615            FrameRead::Frame(frame) => assert_eq!(frame, b"1234"),
1616            _ => panic!("expected an exact-limit frame"),
1617        }
1618    }
1619
1620    #[tokio::test]
1621    async fn outbound_byte_backpressure_is_cancellation_aware() {
1622        let (output, mut receiver) = outbound_channel(5);
1623        output.send(vec![0; 4], None).await.unwrap();
1624
1625        let cancellation = CancellationToken::new();
1626        let blocked = tokio::spawn({
1627            let output = output.clone();
1628            let cancellation = cancellation.clone();
1629            async move { output.send(vec![1; 4], Some(&cancellation)).await }
1630        });
1631        tokio::task::yield_now().await;
1632        assert!(!blocked.is_finished());
1633
1634        cancellation.cancel();
1635        assert_eq!(blocked.await.unwrap(), Err(OutboundSendError::Cancelled));
1636
1637        drop(receiver.recv().await.unwrap());
1638        output.send(vec![2; 4], None).await.unwrap();
1639    }
1640
1641    #[tokio::test]
1642    async fn outbound_control_send_times_out_under_byte_backpressure() {
1643        let (output, _receiver) = outbound_channel(5);
1644        output.send(vec![0; 4], None).await.unwrap();
1645        let result = output
1646            .send_with_timeout(vec![1; 4], None, Duration::from_millis(10))
1647            .await;
1648        assert_eq!(result, Err(OutboundSendError::TimedOut));
1649    }
1650
1651    #[tokio::test]
1652    async fn active_turn_shutdown_aborts_after_grace_period() {
1653        let cancellation = CancellationToken::new();
1654        let dropped = Arc::new(AtomicBool::new(false));
1655        let (started_tx, started_rx) = oneshot::channel();
1656        let task = tokio::spawn({
1657            let dropped = Arc::clone(&dropped);
1658            async move {
1659                let _signal = DropSignal(dropped);
1660                let _ = started_tx.send(());
1661                pending::<()>().await;
1662            }
1663        });
1664        started_rx.await.unwrap();
1665
1666        let graceful = shutdown_active_turn(
1667            ActiveTurn {
1668                turn_id: "turn".into(),
1669                cancellation,
1670                task,
1671            },
1672            Duration::from_millis(10),
1673        )
1674        .await;
1675
1676        assert!(!graceful);
1677        assert!(dropped.load(Ordering::Acquire));
1678    }
1679
1680    #[tokio::test]
1681    async fn writer_shutdown_aborts_after_grace_period() {
1682        let dropped = Arc::new(AtomicBool::new(false));
1683        let (started_tx, started_rx) = oneshot::channel();
1684        let writer = tokio::spawn({
1685            let dropped = Arc::clone(&dropped);
1686            async move {
1687                let _signal = DropSignal(dropped);
1688                let _ = started_tx.send(());
1689                pending::<std::io::Result<()>>().await
1690            }
1691        });
1692        started_rx.await.unwrap();
1693
1694        let result = shutdown_writer(writer, Duration::from_millis(10)).await;
1695
1696        assert!(result.is_err());
1697        assert!(dropped.load(Ordering::Acquire));
1698    }
1699}