1mod 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
32pub use agents::{Endpoint, PI_PROVIDER, WireApi, read_secret};
33pub use scv_tools::{adapters, delegation};
34
35pub fn agent_command(agent: &str) -> Result<std::process::Command> {
40 let config = Config::load_user(ConfigOverrides::default())?;
41 config.prepare_adapter_homes()?;
42 let adapter = config
43 .adapters()
44 .remove(&format!("agent_{agent}"))
45 .ok_or_else(|| anyhow!("unknown agent {agent}"))?;
46 let executable =
47 scv_tools::adapters::resolve_agent_executable(&adapter.command, &adapter.search_dirs)
48 .ok_or_else(|| {
49 anyhow!(
50 "{agent} is not installed: {:?} was not found on PATH or in ~/.local/bin",
51 adapter.command
52 )
53 })?;
54 let mut command = std::process::Command::new(executable);
55 command.current_dir(config.instance_home.join("adapters").join(agent));
56 scv_tools::apply_agent_environment(&mut command, &adapter.environment);
57 Ok(command)
58}
59
60pub fn agent_executable(agent: &str) -> Result<Option<PathBuf>> {
62 let config = Config::load_user(ConfigOverrides::default())?;
63 let adapter = config
64 .adapters()
65 .remove(&format!("agent_{agent}"))
66 .ok_or_else(|| anyhow!("unknown agent {agent}"))?;
67 Ok(scv_tools::adapters::resolve_agent_executable(
68 &adapter.command,
69 &adapter.search_dirs,
70 ))
71}
72
73pub fn agent_home(agent: &str) -> Result<PathBuf> {
75 let config = Config::load_user(ConfigOverrides::default())?;
76 config.prepare_adapter_homes()?;
77 let home = config.instance_home.join("adapters").join(agent);
78 if !home.is_dir() {
79 return Err(anyhow!("unknown agent {agent}"));
80 }
81 Ok(home)
82}
83
84fn key_store_home(agent: &str) -> Result<PathBuf> {
85 scv_tools::adapters::adapter(agent).ok_or_else(|| anyhow!("unknown agent {agent}"))?;
86 agent_home(agent)
87}
88
89pub fn store_agent_key(agent: &str, store: adapters::KeyStore, key: &str) -> Result<Vec<String>> {
91 agents::store_key(store, &key_store_home(agent)?, key)
92}
93
94pub fn agent_stored_status(agent: &str, store: adapters::KeyStore) -> Result<(bool, Vec<String>)> {
97 agents::stored_status(store, &key_store_home(agent)?)
98}
99
100pub fn remove_agent_credentials(agent: &str, store: adapters::KeyStore) -> Result<Vec<String>> {
102 agents::remove_stored(store, &key_store_home(agent)?)
103}
104
105pub fn configure_pi_endpoint(endpoint: &Endpoint, key: &str) -> Result<Vec<String>> {
107 agents::configure_pi_endpoint(&pi_agent_dir()?, endpoint, key)
108}
109
110pub fn import_pi_from_scv_provider() -> Result<Vec<String>> {
114 let config = Config::load_user(ConfigOverrides::default())?;
115 let provider = &config.provider;
116 if provider.kind != "openai-compatible" {
117 return Err(anyhow!("SCV's provider is not openai-compatible"));
118 }
119 let key = match (&provider.api_key, &provider.api_key_env) {
120 (Some(key), _) if !key.trim().is_empty() => key.trim().to_owned(),
121 (_, Some(variable)) => std::env::var(variable)
122 .ok()
123 .filter(|key| !key.trim().is_empty())
124 .ok_or_else(|| {
125 anyhow!("SCV's provider reads its key from ${variable}, which is not set here")
126 })?,
127 _ => return Err(anyhow!("SCV's provider has no API key configured")),
128 };
129 let endpoint = Endpoint {
130 base_url: provider.base_url.clone(),
131 api: WireApi::Responses,
132 model: provider.model.clone(),
133 };
134 let mut notes = agents::configure_pi_endpoint(&pi_agent_dir()?, &endpoint, &key)?;
135 if !provider.headers.is_empty() {
136 notes.push(
137 "Note: SCV's provider sends extra headers, which were not copied; add them to \
138 pi's models.json if the endpoint needs them"
139 .into(),
140 );
141 }
142 Ok(notes)
143}
144
145fn pi_agent_dir() -> Result<PathBuf> {
146 let descriptor =
147 scv_tools::adapters::adapter("pi").ok_or_else(|| anyhow!("unknown agent pi"))?;
148 let scv_tools::adapters::Status::Stored(scv_tools::adapters::KeyStore::Pi { dir }) =
149 descriptor.status
150 else {
151 return Err(anyhow!("pi has no SCV-managed store"));
152 };
153 Ok(agent_home("pi")?.join(dir))
154}
155
156pub fn import_codex(source: &Path) -> Result<Vec<String>> {
160 let config = Config::load_user(ConfigOverrides::default())?;
161 config.prepare_adapter_homes()?;
162 agents::import_codex(source, &config.instance_home.join("adapters").join("codex"))
163}
164
165pub fn service_name() -> anyhow::Result<String> {
167 if std::env::var_os("SCV_HOME").is_none() {
168 return Ok("scv.service".into());
169 }
170 let home = user_home_path().ok_or_else(|| anyhow!("cannot determine SCV instance home"))?;
171 let digest = Sha256::digest(home.to_string_lossy().as_bytes());
172 let suffix = digest[..8]
173 .iter()
174 .map(|byte| format!("{byte:02x}"))
175 .collect::<String>();
176 Ok(format!("scv-{suffix}.service"))
177}
178
179pub fn service_unit_path() -> anyhow::Result<std::path::PathBuf> {
180 let config =
181 dirs::config_dir().ok_or_else(|| anyhow!("cannot determine XDG config directory"))?;
182 Ok(config.join("systemd/user").join(service_name()?))
183}
184use scv_core::{
185 AgentError, AgentRuntime, ApprovalGate, ApprovalRequest, BudgetContextPolicy, CoreEvent,
186 EventSink, Message, ToolRegistry, ToolRisk,
187};
188use scv_protocol::{
189 ClientMessage, DaemonCommand, DaemonStatus, DelegationInfo, DelegationSummary,
190 PROTOCOL_VERSION, PeerInfo, QueueEntry, ServerEvent, Usage,
191};
192use scv_provider_openai::OpenAiProvider;
193use scv_tools::{
194 DelegationContext, SkillMap, builtin_registry,
195 delegation::{self as delegations, DelegationRegistry},
196};
197use tokio::{
198 io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader},
199 net::{UnixListener, UnixStream},
200 sync::{Mutex, OwnedSemaphorePermit, Semaphore, mpsc, oneshot},
201 task::JoinHandle,
202};
203use tokio_util::sync::CancellationToken;
204use tokio_util::task::TaskTracker;
205use uuid::Uuid;
206
207const PROMPT_LIMIT_BYTES: usize = 256 * 1024;
208const OUTPUT_QUEUE_CAPACITY: usize = 256;
209const OUTPUT_QUEUE_MIN_BYTES: usize = 16 * 1024 * 1024;
210const SHUTDOWN_GRACE: Duration = Duration::from_secs(3);
211const MAX_QUEUE_ITEMS: usize = 64;
212const MAX_QUEUE_BYTES: usize = 4 * 1024 * 1024;
213
214pub async fn run_stdio(overrides: ConfigOverrides) -> Result<()> {
215 let stdin = tokio::io::stdin();
216 let stdout = tokio::io::stdout();
217 let tasks = TaskTracker::new();
218 let registry = instance_delegations()?;
219 tokio::spawn(reconcile_delegations(Arc::clone(®istry)));
222 let result = run_managed(
223 stdin,
224 stdout,
225 overrides,
226 None,
227 registry,
228 CancellationToken::new(),
229 tasks.clone(),
230 )
231 .await;
232 tasks.close();
233 tasks.wait().await;
234 result
235}
236
237pub fn default_socket_path() -> Result<PathBuf> {
239 scv_client::default_socket_path()
240}
241
242pub async fn run_socket(path: &Path, overrides: ConfigOverrides) -> Result<()> {
244 if let Some(parent) = path.parent() {
245 tokio::fs::create_dir_all(parent)
246 .await
247 .context("create SCV socket directory")?;
248 }
249 let _lock = SocketLock::acquire(path)?;
250 if path.exists() {
251 if UnixStream::connect(path).await.is_ok() {
252 return Err(anyhow!(
253 "SCV server is already running at {}",
254 path.display()
255 ));
256 }
257 use std::os::unix::fs::FileTypeExt;
258 if !std::fs::symlink_metadata(path)?.file_type().is_socket() {
259 return Err(anyhow!(
260 "refusing to remove a non-socket at SCV socket path"
261 ));
262 }
263 tokio::fs::remove_file(path)
264 .await
265 .with_context(|| format!("remove stale SCV socket {}", path.display()))?;
266 }
267 let listener = UnixListener::bind(path)
268 .with_context(|| format!("bind SCV server socket {}", path.display()))?;
269 #[cfg(unix)]
270 {
271 use std::os::unix::fs::PermissionsExt;
272 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
273 .context("secure SCV socket")?;
274 }
275 let components = Arc::new(Mutex::new(components::Components::new(
276 path.to_owned(),
277 std::env::current_dir()?,
278 )));
279 let registry = instance_delegations()?;
280 if !delegations::become_child_subreaper() {
282 tracing::debug!("SCV daemon is not a child subreaper on this platform");
283 }
284 let cancellation = CancellationToken::new();
285 let delegation_registry = Arc::clone(®istry);
286 let delegation_cancel = cancellation.clone();
287 let mut delegation_task = tokio::spawn(async move {
288 let mut interval = tokio::time::interval(DELEGATION_RECONCILE_INTERVAL);
290 loop {
291 tokio::select! {
292 biased;
293 _ = delegation_cancel.cancelled() => break,
294 _ = interval.tick() => {
295 reconcile_delegations(Arc::clone(&delegation_registry)).await;
296 let zombies = delegations::reap_orphaned_zombies();
297 if zombies > 0 {
298 tracing::debug!("Reaped {zombies} exited orphan processes");
299 }
300 }
301 }
302 }
303 });
304 let _delegation_abort = AbortGuard(delegation_task.abort_handle());
305 let tasks = TaskTracker::new();
306 let mut clients = tokio::task::JoinSet::new();
307 let refresh_components = components.clone();
308 let refresh_cancel = cancellation.clone();
309 let mut refresh_task = tokio::spawn(async move {
310 let mut refresh = tokio::time::interval(Duration::from_secs(2));
311 loop {
312 tokio::select! {
313 biased;
314 _ = refresh_cancel.cancelled() => break,
315 _ = refresh.tick() => {
316 tokio::select! {
317 biased;
318 _ = refresh_cancel.cancelled() => break,
319 result = async { refresh_components.lock().await.reconcile().await } => {
320 if result.is_err() { tracing::warn!("Component account discovery failed"); }
321 }
322 }
323 }
324 }
325 }
326 });
327 let _refresh_abort = AbortGuard(refresh_task.abort_handle());
328 let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())?;
329 let result = loop {
330 tokio::select! {
331 accepted = listener.accept() => {
332 let (stream, _) = match accepted { Ok(value) => value, Err(error) => break Err(error.into()) };
333 let child_overrides = overrides.clone();
334 let components = components.clone();
335 let registry = Arc::clone(®istry);
336 let cancellation = cancellation.clone();
337 let tasks = tasks.clone();
338 clients.spawn(async move {
339 let (reader, writer) = stream.into_split();
340 if run_managed(reader, writer, child_overrides, Some(components), registry, cancellation, tasks).await.is_err() {
341 tracing::warn!("SCV socket client stopped");
342 }
343 });
344 }
345 _ = clients.join_next(), if !clients.is_empty() => {},
346 _ = tokio::signal::ctrl_c() => break Ok(()),
347 _ = terminate.recv() => break Ok(()),
348 }
349 };
350 drop(listener);
351 cancellation.cancel();
352 let _ = (&mut refresh_task).await;
353 let _ = (&mut delegation_task).await;
354 components.lock().await.shutdown().await;
355 if tokio::time::timeout(Duration::from_secs(8), async {
356 while clients.join_next().await.is_some() {}
357 })
358 .await
359 .is_err()
360 {
361 clients.abort_all();
362 while clients.join_next().await.is_some() {}
363 }
364 tasks.close();
365 tasks.wait().await;
366 let _ = tokio::fs::remove_file(path).await;
367 result
368}
369
370const DELEGATION_RECONCILE_INTERVAL: Duration = Duration::from_secs(60);
371
372fn instance_delegations() -> Result<Arc<DelegationRegistry>> {
374 let home =
375 config::user_home_path().ok_or_else(|| anyhow!("cannot determine SCV instance home"))?;
376 Ok(Arc::new(DelegationRegistry::new(&home)))
377}
378
379async fn reconcile_delegations(registry: Arc<DelegationRegistry>) {
381 let report = registry.reconcile().await;
382 if !report.reaped.is_empty() {
383 tracing::info!(
384 "Reaped {} orphaned delegations: {}",
385 report.reaped.len(),
386 report.reaped.join(", ")
387 );
388 }
389 if report.removed > 0 {
390 tracing::debug!(
391 "Removed {} delegation records whose processes had exited",
392 report.removed
393 );
394 }
395}
396
397enum ControlFailure {
399 Delegation(String),
401 Component,
402}
403
404async fn daemon_control(
406 components: &Arc<Mutex<components::Components>>,
407 registry: &DelegationRegistry,
408 command: DaemonCommand,
409) -> std::result::Result<DaemonStatus, ControlFailure> {
410 let mut killed = Vec::new();
411 let listing = match &command {
412 DaemonCommand::Delegations { all } => Some(*all),
413 DaemonCommand::DelegationKill { handle, orphans } => {
414 if handle.is_none() && !orphans {
415 return Err(ControlFailure::Delegation(
416 "name a delegation handle or ask for orphans".into(),
417 ));
418 }
419 if *orphans {
420 let report = registry.reconcile().await;
421 killed.extend(report.reaped);
422 }
423 if let Some(handle) = handle {
424 registry
425 .kill(handle)
426 .await
427 .map_err(ControlFailure::Delegation)?;
428 killed.push(handle.clone());
429 }
430 Some(true)
431 }
432 _ => None,
433 };
434 let mut status = components
435 .lock()
436 .await
437 .control(command)
438 .await
439 .map_err(|_| ControlFailure::Component)?;
440 let running = registry.list(false);
441 status.delegations = DelegationSummary {
442 active: running.len() as u64,
443 reaped: registry.reaped_total(),
444 entries: match listing {
445 Some(true) => registry.list(true),
446 Some(false) => running,
447 None => Vec::new(),
448 }
449 .into_iter()
450 .map(|entry| DelegationInfo {
451 handle: entry.record.handle,
452 agent: entry.record.agent,
453 session: entry.record.session,
454 depth: entry.record.depth,
455 pid: entry.record.process.pid,
456 owner_pid: entry.record.owner.pid,
457 processes: u32::try_from(entry.processes).unwrap_or(u32::MAX),
458 cwd: entry.record.cwd.display().to_string(),
459 started_unix_seconds: entry.record.started_unix,
460 orphaned: entry.orphaned,
461 })
462 .collect(),
463 killed,
464 };
465 Ok(status)
466}
467
468async fn run_managed<R, W>(
469 reader: R,
470 writer: W,
471 overrides: ConfigOverrides,
472 components: Option<Arc<Mutex<components::Components>>>,
473 registry: Arc<DelegationRegistry>,
474 cancellation: CancellationToken,
475 tasks: TaskTracker,
476) -> Result<()>
477where
478 R: tokio::io::AsyncRead + Unpin,
479 W: tokio::io::AsyncWrite + Unpin + Send + 'static,
480{
481 let initial_output_bytes =
482 output_queue_bytes(Config::default().protocol.max_server_frame_bytes)?;
483 let (output_tx, mut output_rx) = outbound_channel(initial_output_bytes);
484 let mut writer_task = tasks.spawn(async move {
485 let mut writer = writer;
486 while let Some(frame) = output_rx.recv().await {
487 writer.write_all(&frame.bytes).await?;
488 writer.write_all(b"\n").await?;
489 writer.flush().await?;
490 }
491 Ok::<(), std::io::Error>(())
492 });
493 let _writer_abort = AbortGuard(writer_task.abort_handle());
494 let (done_tx, mut done_rx) = mpsc::channel::<TurnDone>(4);
495 let approvals = Arc::new(ApprovalBroker::default());
496 let mut reader = BufReader::new(reader);
497 let mut frames = FrameBuffer::default();
498 let mut initialized = false;
499 let mut session: Option<Session> = None;
500 let mut active: Option<ActiveTurn> = None;
501 let mut fatal = false;
502 let mut writer_finished = false;
503
504 let loop_result: Result<()> = async {
505 loop {
506 let frame_limit = session.as_ref().map_or_else(
507 || Config::default().protocol.max_client_frame_bytes,
508 |value| value.config.protocol.max_client_frame_bytes,
509 );
510 tokio::select! {
511 _ = cancellation.cancelled() => break,
512 read = frames.read(&mut reader, frame_limit) => {
513 let frame = match read.context("read protocol input")? {
514 FrameRead::Eof => {
515 if let Some(active) = &active { active.cancellation.cancel(); }
516 break;
517 }
518 FrameRead::TooLarge => {
519 send_error(&output_tx, "", "invalid_request", "client frame exceeds configured limit", false, server_frame_limit(&session)).await?;
520 continue;
521 }
522 FrameRead::Frame(frame) => frame,
523 };
524 if frame.is_empty() {
525 send_error(&output_tx, "", "invalid_json", "protocol frame is empty", false, server_frame_limit(&session)).await?;
526 continue;
527 }
528 let message = match serde_json::from_slice::<ClientMessage>(&frame) {
529 Ok(message) => message,
530 Err(error) => {
531 send_error(&output_tx, "", "invalid_json", &format!("invalid protocol JSON: {error}"), false, server_frame_limit(&session)).await?;
532 continue;
533 }
534 };
535 match message {
536 ClientMessage::Initialize { request_id, protocol_version, .. } => {
537 if initialized {
538 send_error(&output_tx, &request_id, "invalid_request", "connection is already initialized", false, server_frame_limit(&session)).await?;
539 continue;
540 }
541 if protocol_version != PROTOCOL_VERSION {
542 send_error(&output_tx, &request_id, "version_mismatch", &format!("server supports protocol {PROTOCOL_VERSION}"), true, server_frame_limit(&session)).await?;
543 fatal = true;
544 break;
545 }
546 initialized = true;
547 send_event(&output_tx, ServerEvent::Initialized {
548 request_id,
549 protocol_version: PROTOCOL_VERSION,
550 server: PeerInfo { name: "scv-server".into(), version: env!("CARGO_PKG_VERSION").into() },
551 }, Config::default().protocol.max_server_frame_bytes).await?;
552 }
553 other if !initialized => {
554 send_error(&output_tx, other.request_id(), "not_initialized", "initialize must be the first message", false, server_frame_limit(&session)).await?;
555 }
556 ClientMessage::DaemonControl { request_id, command } => {
557 if let Some(components) = &components {
558 let result = tokio::select! {
559 biased;
560 _ = cancellation.cancelled() => break,
561 result = daemon_control(components, ®istry, command) => result,
562 };
563 match result {
564 Ok(status) => send_event(&output_tx, ServerEvent::DaemonStatus { request_id, status }, server_frame_limit(&session)).await?,
565 Err(ControlFailure::Delegation(message)) => send_error(&output_tx, &request_id, "delegation_error", &message, false, server_frame_limit(&session)).await?,
566 Err(ControlFailure::Component) => 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?,
567 }
568 } else {
569 send_error(&output_tx, &request_id, "unsupported", "Component management requires the daemon socket", false, server_frame_limit(&session)).await?;
570 }
571 }
572 ClientMessage::SessionStart { request_id, cwd, provider, model, base_url, no_tools } => {
573 if session.is_some() {
574 send_error(&output_tx, &request_id, "invalid_request", "this connection already has a session", false, server_frame_limit(&session)).await?;
575 continue;
576 }
577 let session_overrides = ConfigOverrides {
578 provider: provider.or_else(|| overrides.provider.clone()),
579 model: model.or_else(|| overrides.model.clone()),
580 base_url: base_url.or_else(|| overrides.base_url.clone()),
581 approval_policy: overrides.approval_policy,
582 no_tools: no_tools.unwrap_or(overrides.no_tools),
583 };
584 match build_session(&cwd, session_overrides, ®istry).await {
585 Ok(new_session) => {
586 output_tx.ensure_capacity(output_queue_bytes(
587 new_session.config.protocol.max_server_frame_bytes,
588 )?)?;
589 let event = ServerEvent::SessionStarted {
590 request_id,
591 session_id: new_session.id.clone(),
592 cwd: new_session.workspace.display().to_string(),
593 model: new_session.runtime.model().to_owned(),
594 context_max_tokens: new_session.config.context.max_tokens,
595 max_server_frame_bytes: new_session.config.protocol.max_server_frame_bytes,
596 max_transcript_bytes: new_session.config.tui.max_transcript_bytes,
597 max_transcript_items: new_session.config.tui.max_transcript_items,
598 max_prompt_history_bytes: new_session.config.tui.max_prompt_history_bytes,
599 max_prompt_history_items: new_session.config.tui.max_prompt_history_items,
600 };
601 send_event(&output_tx, event, new_session.config.protocol.max_server_frame_bytes).await?;
602 send_event(&output_tx, ServerEvent::QueueSnapshot {
603 request_id: None,
604 session_id: new_session.id.clone(),
605 seq: next_seq(&new_session.seq),
606 entries: new_session.queue.lock().await.iter().cloned().collect(),
607 paused: new_session.paused.load(Ordering::Acquire),
608 }, new_session.config.protocol.max_server_frame_bytes).await?;
609 session = Some(new_session);
610 }
611 Err(error) => {
612 send_error(&output_tx, &request_id, "invalid_request", &error.to_string(), false, server_frame_limit(&session)).await?;
613 }
614 }
615 }
616 ClientMessage::SessionAttach { request_id, .. } => {
617 send_error(&output_tx, &request_id, "unsupported", "session attach requires the shared socket server", false, server_frame_limit(&session)).await?;
618 }
619 ClientMessage::TurnStart { request_id, session_id, prompt } => {
620 let Some(current) = session.as_ref() else {
621 send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?;
622 continue;
623 };
624 if current.id != session_id {
625 send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
626 continue;
627 }
628 if prompt.trim().is_empty() || prompt.len() > PROMPT_LIMIT_BYTES {
629 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?;
630 continue;
631 }
632 if active.is_some() {
633 let entry = match current.enqueue(prompt, request_id.clone()).await {
634 Ok(entry) => entry,
635 Err(code) => { send_error(&output_tx, &request_id, code, "session queue limit reached", false, server_frame_limit(&session)).await?; continue; }
636 };
637 let position = current.queue.lock().await.len().saturating_sub(1);
638 send_event(&output_tx, ServerEvent::QueueEnqueued {
639 request_id, session_id: current.id.clone(), seq: next_seq(¤t.seq), entry, position,
640 }, current.config.protocol.max_server_frame_bytes).await?;
641 continue;
642 }
643 let turn_id = Uuid::new_v4().to_string();
644 let cancellation = cancellation.child_token();
645 let meta = TurnMeta {
646 request_id: request_id.clone(),
647 session_id: current.id.clone(),
648 turn_id: turn_id.clone(),
649 seq: Arc::clone(¤t.seq),
650 max_server_frame: current.config.protocol.max_server_frame_bytes,
651 };
652 send_event(&output_tx, ServerEvent::TurnStarted {
653 request_id: request_id.clone(),
654 session_id: current.id.clone(),
655 turn_id: turn_id.clone(),
656 seq: next_seq(¤t.seq),
657 }, current.config.protocol.max_server_frame_bytes).await?;
658 let runtime = Arc::clone(¤t.runtime);
659 let history = Arc::clone(¤t.history);
660 let sink: Arc<dyn EventSink> = Arc::new(ProtocolSink {
661 meta: meta.clone(),
662 output: output_tx.clone(),
663 cancellation: cancellation.clone(),
664 });
665 let gate: Arc<dyn ApprovalGate> = Arc::new(ProtocolApprovalGate {
666 policy: current.config.tools.approval_policy,
667 broker: Arc::clone(&approvals),
668 meta,
669 output: output_tx.clone(),
670 });
671 let task_cancel = cancellation.clone();
672 let task_done = done_tx.clone();
673 let task_request = request_id.clone();
674 let task_session = current.id.clone();
675 let task_turn = turn_id.clone();
676 let task = tasks.spawn(async move {
677 let mut history = history.lock().await;
678 let result = runtime.run_turn(&mut history, prompt, sink, gate, task_cancel).await;
679 let _ = task_done.send(TurnDone {
680 request_id: task_request,
681 session_id: task_session,
682 turn_id: task_turn,
683 result,
684 }).await;
685 });
686 active = Some(ActiveTurn { turn_id, cancellation, task });
687 }
688 ClientMessage::QueueUpdate { request_id, session_id, queue_id, revision, prompt } => {
689 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; };
690 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; }
691 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; }
692 match current.update_queue(&queue_id, revision, prompt).await {
693 Ok(entry) => send_event(&output_tx, ServerEvent::QueueUpdated { request_id, session_id: current.id.clone(), seq: next_seq(¤t.seq), entry }, current.config.protocol.max_server_frame_bytes).await?,
694 Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
695 }
696 }
697 ClientMessage::QueueMove { request_id, session_id, queue_id, revision, before_queue_id } => {
698 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; };
699 match current.move_queue(&session_id, &queue_id, revision, before_queue_id).await {
700 Ok((id, rev, pos)) => send_event(&output_tx, ServerEvent::QueueMoved { request_id, session_id: current.id.clone(), seq: next_seq(¤t.seq), queue_id: id, position: pos, revision: rev }, current.config.protocol.max_server_frame_bytes).await?,
701 Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
702 }
703 }
704 ClientMessage::QueueRemove { request_id, session_id, queue_id, revision } => {
705 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; };
706 match current.remove_queue(&session_id, &queue_id, revision).await {
707 Ok((id, rev)) => send_event(&output_tx, ServerEvent::QueueRemoved { request_id, session_id: current.id.clone(), seq: next_seq(¤t.seq), queue_id: id, revision: rev }, current.config.protocol.max_server_frame_bytes).await?,
708 Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
709 }
710 }
711 ClientMessage::SessionPause { request_id, session_id, paused } => {
712 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; };
713 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; }
714 current.paused.store(paused, Ordering::Release);
715 send_event(&output_tx, ServerEvent::SessionPaused { request_id, session_id: current.id.clone(), seq: next_seq(¤t.seq), paused }, current.config.protocol.max_server_frame_bytes).await?;
716 }
717 ClientMessage::TurnCancel { request_id, session_id, turn_id } => {
718 match (&session, &active) {
719 (Some(current), Some(running)) if current.id == session_id && running.turn_id == turn_id => running.cancellation.cancel(),
720 _ => send_error(&output_tx, &request_id, "turn_not_found", "active turn was not found", false, server_frame_limit(&session)).await?,
721 }
722 }
723 ClientMessage::ApprovalResolve { request_id, session_id, approval_id, approved } => {
724 if session.as_ref().is_none_or(|current| current.id != session_id) {
725 send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
726 } else if !approvals.resolve(&approval_id, approved).await {
727 send_error(&output_tx, &request_id, "approval_not_found", "approval was not found or already resolved", false, server_frame_limit(&session)).await?;
728 }
729 }
730 ClientMessage::SessionClear { request_id, session_id } => {
731 let Some(current) = session.as_ref() else {
732 send_error(&output_tx, &request_id, "session_not_found", "session was not found", false, server_frame_limit(&session)).await?;
733 continue;
734 };
735 if current.id != session_id {
736 send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
737 } else if active.is_some() {
738 send_error(&output_tx, &request_id, "turn_active", "cancel the active turn before clearing", false, server_frame_limit(&session)).await?;
739 } else {
740 current.history.lock().await.clear();
741 current.queue.lock().await.clear();
742 send_event(&output_tx, ServerEvent::SessionCleared {
743 request_id,
744 session_id: current.id.clone(),
745 seq: next_seq(¤t.seq),
746 }, current.config.protocol.max_server_frame_bytes).await?;
747 send_event(&output_tx, ServerEvent::QueueSnapshot { request_id: None, session_id: current.id.clone(), seq: next_seq(¤t.seq), entries: Vec::new(), paused: current.paused.load(Ordering::Acquire) }, current.config.protocol.max_server_frame_bytes).await?;
748 }
749 }
750 }
751 }
752 writer = &mut writer_task => {
753 writer_finished = true;
754 writer.context("join protocol writer")??;
755 break;
756 }
757 done = done_rx.recv(), if active.is_some() => {
758 if let Some(done) = done {
759 if let Some(current) = session.as_ref() {
760 let seq = next_seq(¤t.seq);
761 let event = match done.result {
762 Ok(outcome) => ServerEvent::TurnCompleted {
763 request_id: done.request_id,
764 session_id: done.session_id,
765 turn_id: done.turn_id,
766 seq,
767 steps: outcome.steps,
768 usage: Usage { input_tokens: outcome.usage.input_tokens, output_tokens: outcome.usage.output_tokens },
769 },
770 Err(AgentError::Cancelled) => ServerEvent::TurnCancelled {
771 request_id: done.request_id,
772 session_id: done.session_id,
773 turn_id: done.turn_id,
774 seq,
775 },
776 Err(error) => ServerEvent::TurnFailed {
777 request_id: done.request_id,
778 session_id: done.session_id,
779 turn_id: done.turn_id,
780 seq,
781 code: error.code().into(),
782 message: error.to_string(),
783 },
784 };
785 send_event(&output_tx, event, current.config.protocol.max_server_frame_bytes).await?;
786 }
787 if let Some(mut active) = active.take() {
788 let _ = (&mut active.task).await;
789 }
790 if let Some(current) = session.as_ref()
791 && !current.paused.load(Ordering::Acquire)
792 && let Some(entry) = current.queue.lock().await.pop_front()
793 {
794 let turn_id = Uuid::new_v4().to_string();
795 let cancellation = cancellation.child_token();
796 send_event(&output_tx, ServerEvent::QueueDequeued {
797 request_id: entry.submitter.clone(),
798 session_id: current.id.clone(),
799 seq: next_seq(¤t.seq),
800 queue_id: entry.queue_id,
801 turn_id: turn_id.clone(),
802 }, current.config.protocol.max_server_frame_bytes).await?;
803 send_event(&output_tx, ServerEvent::TurnStarted {
804 request_id: entry.submitter.clone(),
805 session_id: current.id.clone(),
806 turn_id: turn_id.clone(),
807 seq: next_seq(¤t.seq),
808 }, current.config.protocol.max_server_frame_bytes).await?;
809 let meta = TurnMeta {
810 request_id: entry.submitter.clone(), session_id: current.id.clone(), turn_id: turn_id.clone(),
811 seq: Arc::clone(¤t.seq), max_server_frame: current.config.protocol.max_server_frame_bytes,
812 };
813 let sink: Arc<dyn EventSink> = Arc::new(ProtocolSink { meta: meta.clone(), output: output_tx.clone(), cancellation: cancellation.clone() });
814 let gate: Arc<dyn ApprovalGate> = Arc::new(ProtocolApprovalGate { policy: current.config.tools.approval_policy, broker: Arc::clone(&approvals), meta, output: output_tx.clone() });
815 let runtime = Arc::clone(¤t.runtime);
816 let history = Arc::clone(¤t.history);
817 let task_done = done_tx.clone();
818 let task_request = entry.submitter;
819 let task_session = current.id.clone();
820 let task_turn = turn_id.clone();
821 let task_cancel = cancellation.clone();
822 let task = tasks.spawn(async move {
823 let mut history = history.lock().await;
824 let result = runtime.run_turn(&mut history, entry.prompt, sink, gate, task_cancel).await;
825 let _ = task_done.send(TurnDone { request_id: task_request, session_id: task_session, turn_id: task_turn, result }).await;
826 });
827 active = Some(ActiveTurn { turn_id, cancellation, task });
828 }
829 }
830 }
831 }
832 }
833 Ok(())
834 }
835 .await;
836
837 if let Some(active) = active.take() {
838 shutdown_active_turn(active, SHUTDOWN_GRACE).await;
839 }
840 drop(output_tx);
841 let writer_result = if writer_finished {
842 Ok(())
843 } else {
844 shutdown_writer(writer_task, SHUTDOWN_GRACE).await
845 };
846 loop_result?;
847 writer_result?;
848 if fatal {
849 return Err(anyhow!("protocol version mismatch"));
850 }
851 Ok(())
852}
853
854enum FrameRead {
855 Eof,
856 Frame(Vec<u8>),
857 TooLarge,
858}
859
860struct OutboundFrame {
861 bytes: Vec<u8>,
862 _byte_permit: OwnedSemaphorePermit,
863}
864
865#[derive(Clone)]
866struct OutboundSender {
867 frames: mpsc::Sender<OutboundFrame>,
868 budget: Arc<Semaphore>,
869 capacity: Arc<AtomicUsize>,
870}
871
872#[derive(Debug, PartialEq, Eq)]
873enum OutboundSendError {
874 Cancelled,
875 Closed,
876 TimedOut,
877 FrameExceedsQueue { frame_bytes: usize, capacity: usize },
878}
879
880impl std::fmt::Display for OutboundSendError {
881 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
882 match self {
883 Self::Cancelled => formatter.write_str("outbound send cancelled"),
884 Self::Closed => formatter.write_str("protocol client disconnected"),
885 Self::TimedOut => formatter.write_str("outbound send timed out under backpressure"),
886 Self::FrameExceedsQueue {
887 frame_bytes,
888 capacity,
889 } => write!(
890 formatter,
891 "outbound frame uses {frame_bytes} bytes but queue capacity is {capacity} bytes"
892 ),
893 }
894 }
895}
896
897impl std::error::Error for OutboundSendError {}
898
899fn outbound_channel(capacity: usize) -> (OutboundSender, mpsc::Receiver<OutboundFrame>) {
900 let (frames, receiver) = mpsc::channel(OUTPUT_QUEUE_CAPACITY);
901 (
902 OutboundSender {
903 frames,
904 budget: Arc::new(Semaphore::new(capacity)),
905 capacity: Arc::new(AtomicUsize::new(capacity)),
906 },
907 receiver,
908 )
909}
910
911impl OutboundSender {
912 fn ensure_capacity(&self, required: usize) -> Result<()> {
913 if required > Semaphore::MAX_PERMITS {
914 return Err(anyhow!(
915 "outbound queue capacity {required} exceeds runtime limit {}",
916 Semaphore::MAX_PERMITS
917 ));
918 }
919 let current = self.capacity.load(Ordering::Acquire);
920 if required > current {
921 self.budget.add_permits(required - current);
922 self.capacity.store(required, Ordering::Release);
923 }
924 Ok(())
925 }
926
927 async fn send(
928 &self,
929 bytes: Vec<u8>,
930 cancellation: Option<&CancellationToken>,
931 ) -> std::result::Result<(), OutboundSendError> {
932 self.send_with_timeout(bytes, cancellation, SHUTDOWN_GRACE)
933 .await
934 }
935
936 async fn send_with_timeout(
937 &self,
938 bytes: Vec<u8>,
939 cancellation: Option<&CancellationToken>,
940 control_timeout: Duration,
941 ) -> std::result::Result<(), OutboundSendError> {
942 let frame_bytes =
943 bytes
944 .len()
945 .checked_add(1)
946 .ok_or(OutboundSendError::FrameExceedsQueue {
947 frame_bytes: usize::MAX,
948 capacity: self.capacity.load(Ordering::Acquire),
949 })?;
950 let capacity = self.capacity.load(Ordering::Acquire);
951 let permits =
952 u32::try_from(frame_bytes).map_err(|_| OutboundSendError::FrameExceedsQueue {
953 frame_bytes,
954 capacity,
955 })?;
956 if frame_bytes > capacity {
957 return Err(OutboundSendError::FrameExceedsQueue {
958 frame_bytes,
959 capacity,
960 });
961 }
962
963 let control_deadline = tokio::time::Instant::now() + control_timeout;
964 let acquire = Arc::clone(&self.budget).acquire_many_owned(permits);
965 let permit = if let Some(cancellation) = cancellation {
966 tokio::select! {
967 biased;
968 _ = cancellation.cancelled() => return Err(OutboundSendError::Cancelled),
969 permit = acquire => permit.map_err(|_| OutboundSendError::Closed)?,
970 }
971 } else {
972 tokio::time::timeout_at(control_deadline, acquire)
973 .await
974 .map_err(|_| OutboundSendError::TimedOut)?
975 .map_err(|_| OutboundSendError::Closed)?
976 };
977 let frame = OutboundFrame {
978 bytes,
979 _byte_permit: permit,
980 };
981 if let Some(cancellation) = cancellation {
982 tokio::select! {
983 biased;
984 _ = cancellation.cancelled() => Err(OutboundSendError::Cancelled),
985 result = self.frames.send(frame) => result.map_err(|_| OutboundSendError::Closed),
986 }
987 } else {
988 tokio::time::timeout_at(control_deadline, self.frames.send(frame))
989 .await
990 .map_err(|_| OutboundSendError::TimedOut)?
991 .map_err(|_| OutboundSendError::Closed)
992 }
993 }
994}
995
996fn output_queue_bytes(max_frame_bytes: usize) -> Result<usize> {
997 let required = max_frame_bytes
998 .checked_add(1)
999 .and_then(|bytes| bytes.checked_mul(2))
1000 .ok_or_else(|| anyhow!("configured server frame limit is too large"))?
1001 .max(OUTPUT_QUEUE_MIN_BYTES);
1002 if required > Semaphore::MAX_PERMITS {
1003 return Err(anyhow!(
1004 "configured server frame limit requires an outbound queue larger than the runtime supports"
1005 ));
1006 }
1007 Ok(required)
1008}
1009
1010#[derive(Default)]
1012struct FrameBuffer {
1013 bytes: Vec<u8>,
1014 oversized: bool,
1015}
1016
1017impl FrameBuffer {
1018 async fn read<R>(&mut self, reader: &mut R, max_bytes: usize) -> std::io::Result<FrameRead>
1019 where
1020 R: AsyncBufRead + Unpin,
1021 {
1022 loop {
1023 let available = reader.fill_buf().await?;
1024 let eof = available.is_empty();
1025 let end = available.iter().position(|b| *b == b'\n');
1026 let take = end.map_or(available.len(), |n| n + 1);
1027 if !self.oversized {
1028 if self.bytes.len().saturating_add(take) > max_bytes.saturating_add(2) {
1029 self.oversized = true;
1030 self.bytes.clear();
1031 } else {
1032 self.bytes.extend_from_slice(&available[..take]);
1033 }
1034 }
1035 reader.consume(take);
1036 if end.is_some() || eof {
1037 if std::mem::take(&mut self.oversized) {
1038 return Ok(FrameRead::TooLarge);
1039 }
1040 if eof && self.bytes.is_empty() {
1041 return Ok(FrameRead::Eof);
1042 }
1043 let mut bytes = std::mem::take(&mut self.bytes);
1044 while matches!(bytes.last(), Some(b'\n' | b'\r')) {
1045 bytes.pop();
1046 }
1047 return Ok(if bytes.len() > max_bytes {
1048 FrameRead::TooLarge
1049 } else {
1050 FrameRead::Frame(bytes)
1051 });
1052 }
1053 }
1054 }
1055}
1056
1057#[cfg(test)]
1058async fn read_bounded_frame<R: AsyncBufRead + Unpin>(
1059 reader: &mut R,
1060 max_bytes: usize,
1061) -> std::io::Result<FrameRead> {
1062 FrameBuffer::default().read(reader, max_bytes).await
1063}
1064
1065struct SocketLock(std::fs::File);
1067impl SocketLock {
1068 fn acquire(socket: &Path) -> Result<Self> {
1069 use std::os::unix::{fs::OpenOptionsExt, io::AsRawFd};
1070 let file = std::fs::OpenOptions::new()
1071 .read(true)
1072 .write(true)
1073 .create(true)
1074 .truncate(false)
1075 .mode(0o600)
1076 .custom_flags(libc::O_NOFOLLOW)
1077 .open(socket.with_extension("lock"))?;
1078 if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } != 0 {
1080 return Err(anyhow!("SCV daemon already owns this socket"));
1081 }
1082 Ok(Self(file))
1083 }
1084}
1085impl Drop for SocketLock {
1086 fn drop(&mut self) {
1087 use std::os::unix::io::AsRawFd;
1088 unsafe {
1090 libc::flock(self.0.as_raw_fd(), libc::LOCK_UN);
1091 }
1092 }
1093}
1094
1095fn server_frame_limit(session: &Option<Session>) -> usize {
1096 session.as_ref().map_or_else(
1097 || Config::default().protocol.max_server_frame_bytes,
1098 |value| value.config.protocol.max_server_frame_bytes,
1099 )
1100}
1101
1102struct Session {
1103 id: String,
1104 workspace: PathBuf,
1105 config: Config,
1106 runtime: Arc<AgentRuntime>,
1107 history: Arc<Mutex<Vec<Message>>>,
1108 seq: Arc<AtomicU64>,
1109 queue: Arc<Mutex<VecDeque<QueueEntry>>>,
1110 paused: Arc<std::sync::atomic::AtomicBool>,
1111}
1112
1113impl Session {
1114 async fn enqueue(
1115 &self,
1116 prompt: String,
1117 submitter: String,
1118 ) -> std::result::Result<QueueEntry, &'static str> {
1119 let entry = QueueEntry {
1120 queue_id: Uuid::new_v4().to_string(),
1121 revision: 1,
1122 prompt,
1123 submitter,
1124 };
1125 let mut queue = self.queue.lock().await;
1126 let bytes: usize = queue.iter().map(|item| item.prompt.len()).sum();
1127 if queue.len() >= MAX_QUEUE_ITEMS
1128 || bytes.saturating_add(entry.prompt.len()) > MAX_QUEUE_BYTES
1129 {
1130 return Err("queue_limit");
1131 }
1132 queue.push_back(entry.clone());
1133 Ok(entry)
1134 }
1135
1136 async fn update_queue(
1137 &self,
1138 id: &str,
1139 revision: u64,
1140 prompt: String,
1141 ) -> std::result::Result<QueueEntry, &'static str> {
1142 let mut queue = self.queue.lock().await;
1143 let bytes: usize = queue.iter().map(|item| item.prompt.len()).sum();
1144 let entry = queue
1145 .iter_mut()
1146 .find(|entry| entry.queue_id == id)
1147 .ok_or("queue_not_found")?;
1148 if entry.revision != revision {
1149 return Err("queue_conflict");
1150 }
1151 if bytes
1152 .saturating_sub(entry.prompt.len())
1153 .saturating_add(prompt.len())
1154 > MAX_QUEUE_BYTES
1155 {
1156 return Err("queue_limit");
1157 }
1158 entry.prompt = prompt;
1159 entry.revision += 1;
1160 Ok(entry.clone())
1161 }
1162
1163 async fn move_queue(
1164 &self,
1165 session_id: &str,
1166 id: &str,
1167 revision: u64,
1168 before: Option<String>,
1169 ) -> std::result::Result<(String, u64, usize), &'static str> {
1170 if self.id != session_id {
1171 return Err("session_not_found");
1172 }
1173 let mut queue = self.queue.lock().await;
1174 let index = queue
1175 .iter()
1176 .position(|entry| entry.queue_id == id)
1177 .ok_or("queue_not_found")?;
1178 if queue[index].revision != revision {
1179 return Err("queue_conflict");
1180 }
1181 let target_index = match before.as_deref() {
1184 Some(target) if target == id => return Ok((id.to_string(), revision, index)),
1185 Some(target) => Some(
1186 queue
1187 .iter()
1188 .position(|item| item.queue_id == target)
1189 .ok_or("queue_not_found")?,
1190 ),
1191 None => None,
1192 };
1193 let mut entry = queue.remove(index).expect("queue index exists");
1194 let target = target_index.map_or(queue.len(), |target| {
1195 target.saturating_sub(usize::from(target > index))
1196 });
1197 let pos = target.min(queue.len());
1198 let id = entry.queue_id.clone();
1199 let rev = entry.revision + 1;
1200 entry.revision = rev;
1201 queue.insert(pos, entry);
1202 Ok((id, rev, pos))
1203 }
1204
1205 async fn remove_queue(
1206 &self,
1207 session_id: &str,
1208 id: &str,
1209 revision: u64,
1210 ) -> std::result::Result<(String, u64), &'static str> {
1211 if self.id != session_id {
1212 return Err("session_not_found");
1213 }
1214 let mut queue = self.queue.lock().await;
1215 let index = queue
1216 .iter()
1217 .position(|entry| entry.queue_id == id)
1218 .ok_or("queue_not_found")?;
1219 if queue[index].revision != revision {
1220 return Err("queue_conflict");
1221 }
1222 let entry = queue.remove(index).expect("queue index exists");
1223 Ok((entry.queue_id, entry.revision))
1224 }
1225}
1226
1227struct ActiveTurn {
1228 turn_id: String,
1229 cancellation: CancellationToken,
1230 task: JoinHandle<()>,
1231}
1232
1233impl Drop for ActiveTurn {
1234 fn drop(&mut self) {
1235 self.cancellation.cancel();
1236 self.task.abort();
1237 }
1238}
1239
1240struct AbortGuard(tokio::task::AbortHandle);
1241impl Drop for AbortGuard {
1242 fn drop(&mut self) {
1243 self.0.abort();
1244 }
1245}
1246
1247async fn shutdown_active_turn(mut active: ActiveTurn, grace: Duration) -> bool {
1248 active.cancellation.cancel();
1249 if tokio::time::timeout(grace, &mut active.task).await.is_ok() {
1250 true
1251 } else {
1252 active.task.abort();
1253 let _ = (&mut active.task).await;
1254 false
1255 }
1256}
1257
1258async fn shutdown_writer(
1259 mut writer: JoinHandle<std::io::Result<()>>,
1260 grace: Duration,
1261) -> Result<()> {
1262 match tokio::time::timeout(grace, &mut writer).await {
1263 Ok(result) => {
1264 result.context("join protocol writer")??;
1265 Ok(())
1266 }
1267 Err(_) => {
1268 writer.abort();
1269 let _ = writer.await;
1270 Err(anyhow!("protocol writer shutdown timed out"))
1271 }
1272 }
1273}
1274
1275struct TurnDone {
1276 request_id: String,
1277 session_id: String,
1278 turn_id: String,
1279 result: Result<scv_core::TurnOutcome, AgentError>,
1280}
1281
1282#[derive(Clone)]
1283struct TurnMeta {
1284 request_id: String,
1285 session_id: String,
1286 turn_id: String,
1287 seq: Arc<AtomicU64>,
1288 max_server_frame: usize,
1289}
1290
1291fn next_seq(sequence: &AtomicU64) -> u64 {
1292 sequence.fetch_add(1, Ordering::Relaxed) + 1
1293}
1294
1295async fn build_session(
1296 cwd: &str,
1297 overrides: ConfigOverrides,
1298 registry: &Arc<DelegationRegistry>,
1299) -> Result<Session> {
1300 let id = Uuid::new_v4().to_string();
1301 let workspace = std::fs::canonicalize(cwd).with_context(|| format!("resolve cwd {cwd}"))?;
1302 if !workspace.is_dir() {
1303 return Err(anyhow!("cwd is not a directory"));
1304 }
1305 let no_tools = overrides.no_tools;
1306 let config = Config::load(&workspace, overrides)?;
1307 if !no_tools {
1308 config.prepare_adapter_homes()?;
1309 }
1310 let provider_config = config.provider.clone();
1311 let api_key = provider_config.api_key.clone().or_else(|| {
1312 provider_config.api_key_env.as_deref().and_then(|name| std::env::var(name).ok())
1313 }).filter(|key| !key.trim().is_empty()).ok_or_else(|| anyhow!("provider credential is not configured; set provider.api_key or provider.api_key_env"))?;
1314 let skills = discover_skills(&workspace, &config, !no_tools)?;
1315 let system_prompt = build_system_prompt(&workspace, &config, &skills)?;
1316 let mut provider = OpenAiProvider::new(
1317 provider_config.model.clone(),
1318 provider_config.base_url.clone(),
1319 api_key,
1320 Duration::from_secs(provider_config.timeout_seconds),
1321 config.provider_limits(),
1322 provider_config.headers.clone(),
1323 )?;
1324 if !no_tools && config.hosted_web_search() {
1325 provider = provider.with_web_search();
1326 }
1327 let provider = Arc::new(provider);
1328 let tools = if no_tools {
1329 Arc::new(ToolRegistry::default())
1330 } else {
1331 let mut tools = config.tools();
1332 tools.delegation = Some(DelegationContext {
1333 registry: Arc::clone(registry),
1334 session: id.clone(),
1335 });
1336 let mut registry = builtin_registry(
1337 tools,
1338 skills.map,
1339 skills.roots,
1340 config.skills.max_skill_bytes,
1341 config.adapters(),
1342 )?;
1343 if let Some(web) = config.web_tools() {
1344 scv_tools::web::register(&mut registry, web)?;
1345 }
1346 Arc::new(registry)
1347 };
1348 let context = Arc::new(BudgetContextPolicy::new((&config.context).into())?);
1349 let runtime = Arc::new(AgentRuntime::new(
1350 provider,
1351 tools,
1352 context,
1353 config.core_agent(system_prompt),
1354 workspace.clone(),
1355 ));
1356 Ok(Session {
1357 id,
1358 workspace,
1359 config,
1360 runtime,
1361 history: Arc::new(Mutex::new(Vec::new())),
1362 seq: Arc::new(AtomicU64::new(0)),
1363 queue: Arc::new(Mutex::new(VecDeque::new())),
1364 paused: Arc::new(std::sync::atomic::AtomicBool::new(false)),
1365 })
1366}
1367
1368fn build_system_prompt(
1369 workspace: &Path,
1370 config: &Config,
1371 skills: &DiscoveredSkills,
1372) -> Result<String> {
1373 let mut prompt = config.agent.system_prompt.clone();
1374 prompt.push_str(&format!(
1375 "\nCurrent working directory: {}\n",
1376 workspace.display()
1377 ));
1378 let agents_path = workspace.join("AGENTS.md");
1379 if agents_path.is_file() {
1380 let canonical = std::fs::canonicalize(&agents_path).context("resolve project AGENTS.md")?;
1381 if !canonical.starts_with(workspace) {
1382 return Err(anyhow!("project AGENTS.md escaped workspace"));
1383 }
1384 let (bytes, truncated) = read_prefix(&canonical, config.tools.max_read_bytes)
1385 .context("read project AGENTS.md")?;
1386 let instructions = std::str::from_utf8(&bytes).context("project AGENTS.md is not UTF-8")?;
1387 prompt.push_str("\n# Project instructions\n");
1388 prompt.push_str(instructions);
1389 if truncated {
1390 prompt.push_str("\n[AGENTS.md truncated by configured read limit]\n");
1391 }
1392 }
1393 if !skills.listing.is_empty() {
1394 prompt.push_str("\n# Available skills\n");
1395 prompt.push_str(&skills.listing);
1396 prompt.push_str("\nUse read_skill with a skill name when its workflow applies.\n");
1397 }
1398 if !skills.project_listing.is_empty() {
1399 prompt.push_str("\n# Project skills\n");
1400 prompt.push_str(
1401 "Projects in this workspace provide these skills to agents working in them:\n",
1402 );
1403 prompt.push_str(&skills.project_listing);
1404 prompt.push_str(
1405 "\nTo use one, delegate with an agent_* tool such as agent_codex or agent_claude, \
1406 set its cwd to the skill's project, and name the skill in the prompt: that agent \
1407 then loads the project's instructions and skills itself. read_skill loads a \
1408 skill for reference.\n",
1409 );
1410 }
1411 Ok(prompt)
1412}
1413
1414struct DiscoveredSkills {
1417 map: SkillMap,
1418 roots: Vec<PathBuf>,
1419 listing: String,
1420 project_listing: String,
1421}
1422
1423const PROJECT_SKILL_DIRS: [&str; 2] = [".agents/skills", ".claude/skills"];
1426const MAX_WORKSPACE_ENTRIES: usize = 4096;
1429const MAX_SKILL_PROJECTS: usize = 256;
1430const PROJECT_SKILL_HEADER_BYTES: usize = 16 * 1024;
1432const MAX_PROJECT_SKILL_DESCRIPTION: usize = 400;
1433
1434fn discover_skills(workspace: &Path, config: &Config, tools: bool) -> Result<DiscoveredSkills> {
1435 let mut skills = SkillMap::new();
1436 let mut roots = Vec::new();
1437 let project_root = workspace.join(&config.skills.project_dir);
1438 for (root, must_be_workspace) in [(&project_root, true), (&config.skills.user_dir, false)] {
1439 if !root.is_dir() {
1440 continue;
1441 }
1442 let canonical = std::fs::canonicalize(root)
1443 .with_context(|| format!("resolve skill root {}", root.display()))?;
1444 if must_be_workspace && !canonical.starts_with(workspace) {
1445 return Err(anyhow!("project skill root escaped workspace"));
1446 }
1447 roots.push(canonical.clone());
1448 let mut entries: Vec<_> = std::fs::read_dir(&canonical)
1449 .with_context(|| format!("read skill root {}", canonical.display()))?
1450 .filter_map(Result::ok)
1451 .collect();
1452 entries.sort_by_key(|entry| entry.file_name());
1453 for entry in entries {
1454 if skills.len() >= config.skills.max_skills {
1455 break;
1456 }
1457 let path = entry.path().join("SKILL.md");
1458 if !path.is_file() {
1459 continue;
1460 }
1461 let canonical_file = std::fs::canonicalize(&path)
1462 .with_context(|| format!("resolve skill {}", path.display()))?;
1463 if !canonical_file.starts_with(&canonical) {
1464 continue;
1465 }
1466 let name = entry.file_name().to_string_lossy().to_string();
1467 skills.entry(name).or_insert(canonical_file);
1468 }
1469 }
1470 let mut names: Vec<_> = skills.keys().cloned().collect();
1471 names.sort();
1472 let mut listing = String::new();
1473 for name in names {
1474 let path = &skills[&name];
1475 let bytes = read_prefix(path, config.skills.max_skill_bytes)
1476 .map(|(bytes, _)| bytes)
1477 .unwrap_or_default();
1478 let content = String::from_utf8_lossy(&bytes);
1479 let description = skill_description(&content);
1480 listing.push_str(&format!("- {name}: {description}\n"));
1481 }
1482 let project_listing = if tools && config.skills.scan_projects {
1485 discover_project_skills(workspace, config, &mut skills, &mut roots)
1486 } else {
1487 String::new()
1488 };
1489 Ok(DiscoveredSkills {
1490 map: skills,
1491 roots,
1492 listing,
1493 project_listing,
1494 })
1495}
1496
1497fn discover_project_skills(
1503 workspace: &Path,
1504 config: &Config,
1505 skills: &mut SkillMap,
1506 roots: &mut Vec<PathBuf>,
1507) -> String {
1508 let mut projects = vec![(None, workspace.to_path_buf())];
1509 let mut names: Vec<_> = std::fs::read_dir(workspace)
1510 .into_iter()
1511 .flatten()
1512 .filter_map(Result::ok)
1513 .take(MAX_WORKSPACE_ENTRIES)
1514 .map(|entry| entry.file_name().to_string_lossy().into_owned())
1515 .filter(|name| !name.starts_with('.'))
1516 .collect();
1517 names.sort();
1518 for name in names {
1519 if projects.len() > MAX_SKILL_PROJECTS {
1520 break;
1521 }
1522 let Ok(directory) = std::fs::canonicalize(workspace.join(&name)) else {
1523 continue;
1524 };
1525 if directory.is_dir()
1526 && directory.starts_with(workspace)
1527 && !projects.iter().any(|(_, seen)| seen == &directory)
1528 {
1529 projects.push((Some(name), directory));
1530 }
1531 }
1532 let mut seen_files = std::collections::HashSet::new();
1533 let mut listing = String::new();
1534 'projects: for (project, directory) in projects {
1535 let project_roots: Vec<PathBuf> = PROJECT_SKILL_DIRS
1536 .iter()
1537 .filter_map(|relative| std::fs::canonicalize(directory.join(relative)).ok())
1538 .filter(|root| root.is_dir() && root.starts_with(workspace))
1539 .collect();
1540 for root in &project_roots {
1541 if !roots.contains(root) {
1542 roots.push(root.clone());
1543 }
1544 }
1545 for root in &project_roots {
1546 let mut entries: Vec<_> = std::fs::read_dir(root)
1547 .into_iter()
1548 .flatten()
1549 .filter_map(Result::ok)
1550 .take(MAX_WORKSPACE_ENTRIES)
1551 .collect();
1552 entries.sort_by_key(|entry| entry.file_name());
1553 for entry in entries {
1554 if skills.len() >= config.skills.max_skills {
1555 break 'projects;
1556 }
1557 let Ok(file) = std::fs::canonicalize(entry.path().join("SKILL.md")) else {
1558 continue;
1559 };
1560 if !file.is_file()
1561 || !project_roots.iter().any(|root| file.starts_with(root))
1562 || !seen_files.insert(file.clone())
1563 {
1564 continue;
1565 }
1566 let skill = entry.file_name().to_string_lossy().into_owned();
1567 let (name, location) = match &project {
1568 Some(project) => (format!("{project}:{skill}"), format!("project {project}")),
1569 None => (skill, "workspace root".to_owned()),
1570 };
1571 if skills.contains_key(&name) {
1573 continue;
1574 }
1575 let header = read_prefix(
1576 &file,
1577 config
1578 .skills
1579 .max_skill_bytes
1580 .min(PROJECT_SKILL_HEADER_BYTES),
1581 )
1582 .map(|(bytes, _)| bytes)
1583 .unwrap_or_default();
1584 let description: String = skill_description(&String::from_utf8_lossy(&header))
1585 .chars()
1586 .take(MAX_PROJECT_SKILL_DESCRIPTION)
1587 .collect();
1588 listing.push_str(&format!("- {name} ({location}): {description}\n"));
1589 skills.insert(name, file);
1590 }
1591 }
1592 }
1593 listing
1594}
1595
1596fn read_prefix(path: &Path, max_bytes: usize) -> std::io::Result<(Vec<u8>, bool)> {
1597 let file = std::fs::File::open(path)?;
1598 let mut bytes = Vec::with_capacity(max_bytes.min(8192));
1599 file.take(
1600 u64::try_from(max_bytes)
1601 .unwrap_or(u64::MAX)
1602 .saturating_add(1),
1603 )
1604 .read_to_end(&mut bytes)?;
1605 let truncated = bytes.len() > max_bytes;
1606 bytes.truncate(max_bytes);
1607 Ok((bytes, truncated))
1608}
1609
1610fn skill_description(content: &str) -> String {
1611 if let Some(frontmatter) = content.strip_prefix("---\n")
1612 && let Some((header, _)) = frontmatter.split_once("\n---")
1613 {
1614 for line in header.lines() {
1615 if let Some(description) = line.strip_prefix("description:") {
1616 return description.trim().trim_matches('"').to_owned();
1617 }
1618 }
1619 }
1620 content
1621 .lines()
1622 .map(str::trim)
1623 .find(|line| !line.is_empty() && !line.starts_with('#'))
1624 .unwrap_or("No description provided")
1625 .chars()
1626 .take(240)
1627 .collect()
1628}
1629
1630struct ProtocolSink {
1631 meta: TurnMeta,
1632 output: OutboundSender,
1633 cancellation: CancellationToken,
1634}
1635
1636#[async_trait]
1637impl EventSink for ProtocolSink {
1638 async fn emit(&self, event: CoreEvent) -> Result<(), AgentError> {
1639 let seq = next_seq(&self.meta.seq);
1640 let event = match event {
1641 CoreEvent::AssistantDelta { content } => ServerEvent::AssistantDelta {
1642 request_id: self.meta.request_id.clone(),
1643 session_id: self.meta.session_id.clone(),
1644 turn_id: self.meta.turn_id.clone(),
1645 seq,
1646 content,
1647 },
1648 CoreEvent::AssistantCompleted { content } => ServerEvent::AssistantCompleted {
1649 request_id: self.meta.request_id.clone(),
1650 session_id: self.meta.session_id.clone(),
1651 turn_id: self.meta.turn_id.clone(),
1652 seq,
1653 content,
1654 },
1655 CoreEvent::ToolProposed {
1656 call_id,
1657 name,
1658 arguments,
1659 } => ServerEvent::ToolProposed {
1660 request_id: self.meta.request_id.clone(),
1661 session_id: self.meta.session_id.clone(),
1662 turn_id: self.meta.turn_id.clone(),
1663 seq,
1664 call_id,
1665 name,
1666 arguments,
1667 },
1668 CoreEvent::ToolStarted { call_id, name } => ServerEvent::ToolStarted {
1669 request_id: self.meta.request_id.clone(),
1670 session_id: self.meta.session_id.clone(),
1671 turn_id: self.meta.turn_id.clone(),
1672 seq,
1673 call_id,
1674 name,
1675 },
1676 CoreEvent::ToolCompleted {
1677 call_id,
1678 name,
1679 output,
1680 } => ServerEvent::ToolCompleted {
1681 request_id: self.meta.request_id.clone(),
1682 session_id: self.meta.session_id.clone(),
1683 turn_id: self.meta.turn_id.clone(),
1684 seq,
1685 call_id,
1686 name,
1687 success: !output.is_error,
1688 output: output.content,
1689 truncated: output.truncated,
1690 },
1691 CoreEvent::ContextCompacted {
1692 before_tokens,
1693 after_tokens,
1694 removed_messages,
1695 } => ServerEvent::ContextCompacted {
1696 request_id: self.meta.request_id.clone(),
1697 session_id: self.meta.session_id.clone(),
1698 turn_id: self.meta.turn_id.clone(),
1699 seq,
1700 before_tokens,
1701 after_tokens,
1702 removed_messages,
1703 },
1704 CoreEvent::SessionTrimmed {
1705 removed_messages,
1706 history_bytes,
1707 } => ServerEvent::SessionTrimmed {
1708 request_id: self.meta.request_id.clone(),
1709 session_id: self.meta.session_id.clone(),
1710 seq,
1711 removed_messages,
1712 history_bytes,
1713 },
1714 };
1715 send_turn_event(
1716 &self.output,
1717 event,
1718 self.meta.max_server_frame,
1719 &self.cancellation,
1720 )
1721 .await
1722 }
1723}
1724
1725#[derive(Default)]
1726struct ApprovalBroker {
1727 pending: Mutex<HashMap<String, oneshot::Sender<bool>>>,
1728}
1729
1730impl ApprovalBroker {
1731 async fn insert(&self, id: String, sender: oneshot::Sender<bool>) {
1732 self.pending.lock().await.insert(id, sender);
1733 }
1734
1735 async fn remove(&self, id: &str) {
1736 self.pending.lock().await.remove(id);
1737 }
1738
1739 async fn resolve(&self, id: &str, approved: bool) -> bool {
1740 let sender = self.pending.lock().await.remove(id);
1741 sender.is_some_and(|sender| sender.send(approved).is_ok())
1742 }
1743}
1744
1745struct ProtocolApprovalGate {
1746 policy: ApprovalPolicy,
1747 broker: Arc<ApprovalBroker>,
1748 meta: TurnMeta,
1749 output: OutboundSender,
1750}
1751
1752#[async_trait]
1753impl ApprovalGate for ProtocolApprovalGate {
1754 async fn approve(
1755 &self,
1756 request: ApprovalRequest,
1757 cancellation: CancellationToken,
1758 ) -> Result<bool, AgentError> {
1759 match self.policy {
1760 ApprovalPolicy::OnRisk if request.risk == ToolRisk::ReadOnly => return Ok(true),
1761 ApprovalPolicy::Never => return Ok(request.risk == ToolRisk::ReadOnly),
1762 ApprovalPolicy::Always | ApprovalPolicy::OnRisk => {}
1763 }
1764 let approval_id = Uuid::new_v4().to_string();
1765 let (sender, receiver) = oneshot::channel();
1766 self.broker.insert(approval_id.clone(), sender).await;
1767 let event = ServerEvent::ApprovalRequested {
1768 request_id: self.meta.request_id.clone(),
1769 session_id: self.meta.session_id.clone(),
1770 turn_id: self.meta.turn_id.clone(),
1771 seq: next_seq(&self.meta.seq),
1772 approval_id: approval_id.clone(),
1773 call_id: request.call_id,
1774 name: request.name,
1775 risk: request.risk.as_str().into(),
1776 cwd: request.cwd.display().to_string(),
1777 summary: request.summary,
1778 };
1779 if let Err(error) = send_turn_event(
1780 &self.output,
1781 event,
1782 self.meta.max_server_frame,
1783 &cancellation,
1784 )
1785 .await
1786 {
1787 self.broker.remove(&approval_id).await;
1788 return Err(error);
1789 }
1790 tokio::select! {
1791 result = receiver => result.map_err(|_| AgentError::Cancelled),
1792 _ = cancellation.cancelled() => {
1793 self.broker.remove(&approval_id).await;
1794 Err(AgentError::Cancelled)
1795 }
1796 }
1797 }
1798}
1799
1800async fn send_event(output: &OutboundSender, event: ServerEvent, max_bytes: usize) -> Result<()> {
1801 let bytes = encode_event(&event, max_bytes)?;
1802 output.send(bytes, None).await.map_err(anyhow::Error::new)
1803}
1804
1805async fn send_turn_event(
1806 output: &OutboundSender,
1807 event: ServerEvent,
1808 max_bytes: usize,
1809 cancellation: &CancellationToken,
1810) -> Result<(), AgentError> {
1811 let bytes = encode_event(&event, max_bytes)
1812 .map_err(|error| AgentError::ResponseLimit(error.to_string()))?;
1813 match output.send(bytes, Some(cancellation)).await {
1814 Ok(()) => Ok(()),
1815 Err(OutboundSendError::Cancelled) => Err(AgentError::Cancelled),
1816 Err(error) => Err(AgentError::Internal(error.to_string())),
1817 }
1818}
1819
1820fn encode_event(event: &ServerEvent, max_bytes: usize) -> Result<Vec<u8>> {
1821 let bytes = serde_json::to_vec(event).context("serialize protocol event")?;
1822 if bytes.len() > max_bytes {
1823 return Err(anyhow!("server event exceeds configured frame limit"));
1824 }
1825 Ok(bytes)
1826}
1827
1828async fn send_error(
1829 output: &OutboundSender,
1830 request_id: &str,
1831 code: &str,
1832 message: &str,
1833 fatal: bool,
1834 max_bytes: usize,
1835) -> Result<()> {
1836 send_event(
1837 output,
1838 ServerEvent::Error {
1839 request_id: (!request_id.is_empty()).then(|| request_id.to_owned()),
1840 code: code.into(),
1841 message: message.into(),
1842 fatal,
1843 },
1844 max_bytes,
1845 )
1846 .await
1847}
1848
1849#[cfg(test)]
1850mod tests {
1851 use std::{
1852 future::pending,
1853 io::Cursor,
1854 sync::atomic::{AtomicBool, Ordering},
1855 };
1856
1857 fn test_registry() -> Arc<DelegationRegistry> {
1859 let home = tempfile::tempdir().unwrap().keep();
1860 Arc::new(DelegationRegistry::new(&home))
1861 }
1862
1863 use super::*;
1864
1865 struct DropSignal(Arc<AtomicBool>);
1866
1867 impl Drop for DropSignal {
1868 fn drop(&mut self) {
1869 self.0.store(true, Ordering::Release);
1870 }
1871 }
1872
1873 #[tokio::test]
1874 async fn nonreading_management_client_does_not_hold_component_lock() {
1875 let (mut input, server_input) = tokio::io::duplex(65536);
1876 let (server_output, _blocked_output) = tokio::io::duplex(1);
1877 let tasks = TaskTracker::new();
1878 let components = Arc::new(Mutex::new(components::Components::new(
1879 PathBuf::from("/unused.sock"),
1880 PathBuf::from("/"),
1881 )));
1882 let cancel = CancellationToken::new();
1883 let handler = tokio::spawn(run_managed(
1884 server_input,
1885 server_output,
1886 ConfigOverrides::default(),
1887 Some(components.clone()),
1888 test_registry(),
1889 cancel.clone(),
1890 tasks.clone(),
1891 ));
1892 input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
1893 for _ in 0..300 {
1894 input.write_all(b"{\"type\":\"daemon.control\",\"request_id\":\"s\",\"command\":{\"action\":\"status\"}}\n").await.unwrap();
1895 }
1896 tokio::time::sleep(Duration::from_millis(50)).await;
1897 let status = tokio::time::timeout(Duration::from_millis(100), async {
1898 components.lock().await.status()
1899 })
1900 .await
1901 .unwrap();
1902 assert_eq!(status.pid, std::process::id());
1903 cancel.cancel();
1904 handler.abort();
1905 let _ = handler.await;
1906 tasks.close();
1907 tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1908 .await
1909 .unwrap();
1910 }
1911
1912 #[tokio::test]
1913 async fn forced_connection_abort_drops_and_joins_writer_descendants() {
1914 let (mut input, server_input) = tokio::io::duplex(512);
1915 let (server_output, _blocked_output) = tokio::io::duplex(1);
1916 let tasks = TaskTracker::new();
1917 let handler = tokio::spawn(run_managed(
1918 server_input,
1919 server_output,
1920 ConfigOverrides::default(),
1921 None,
1922 test_registry(),
1923 CancellationToken::new(),
1924 tasks.clone(),
1925 ));
1926 input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
1927 tokio::time::timeout(Duration::from_secs(1), async {
1928 while tasks.is_empty() {
1929 tokio::task::yield_now().await;
1930 }
1931 })
1932 .await
1933 .unwrap();
1934 handler.abort();
1935 let _ = handler.await;
1936 tasks.close();
1937 tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1938 .await
1939 .unwrap();
1940 assert!(tasks.is_empty());
1941 }
1942
1943 #[tokio::test]
1944 async fn forced_handler_abort_cancels_and_joins_active_turn() {
1945 let tasks = TaskTracker::new();
1946 let cancellation = CancellationToken::new();
1947 let child_cancel = cancellation.child_token();
1948 let observed_cancel = child_cancel.clone();
1949 let dropped = Arc::new(AtomicBool::new(false));
1950 let (ready_tx, ready_rx) = oneshot::channel();
1951 let task = tasks.spawn({
1952 let dropped = dropped.clone();
1953 async move {
1954 let _guard = DropSignal(dropped);
1955 let _ = ready_tx.send(());
1956 pending::<()>().await;
1957 }
1958 });
1959 ready_rx.await.unwrap();
1960 let (owned_tx, owned_rx) = oneshot::channel();
1961 let handler = tokio::spawn(async move {
1962 let _active = ActiveTurn {
1963 turn_id: "test".into(),
1964 cancellation: child_cancel,
1965 task,
1966 };
1967 let _ = owned_tx.send(());
1968 pending::<()>().await;
1969 });
1970 owned_rx.await.unwrap();
1971 handler.abort();
1972 let _ = handler.await;
1973 tasks.close();
1974 tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1975 .await
1976 .unwrap();
1977 assert!(observed_cancel.is_cancelled());
1978 assert!(dropped.load(Ordering::Acquire));
1979 }
1980
1981 #[tokio::test]
1982 async fn frame_buffer_preserves_partial_and_discard_state_across_cancellation() {
1983 let (mut input, output) = tokio::io::duplex(64);
1984 let mut reader = BufReader::new(output);
1985 let mut frames = FrameBuffer::default();
1986 input.write_all(b"12").await.unwrap();
1987 assert!(
1988 tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
1989 .await
1990 .is_err()
1991 );
1992 input.write_all(b"34\n").await.unwrap();
1993 assert!(
1994 matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"1234")
1995 );
1996 input.write_all(b"123456789").await.unwrap();
1997 assert!(
1998 tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
1999 .await
2000 .is_err()
2001 );
2002 input.write_all(b"\n{}\n").await.unwrap();
2003 assert!(matches!(
2004 frames.read(&mut reader, 4).await.unwrap(),
2005 FrameRead::TooLarge
2006 ));
2007 assert!(
2008 matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"{}")
2009 );
2010 }
2011
2012 #[tokio::test]
2013 async fn bounded_reader_discards_an_oversized_line() {
2014 let input = format!("{}\n{{}}\n", "x".repeat(10));
2015 let mut reader = BufReader::new(Cursor::new(input.into_bytes()));
2016 assert!(matches!(
2017 read_bounded_frame(&mut reader, 4).await.unwrap(),
2018 FrameRead::TooLarge
2019 ));
2020 match read_bounded_frame(&mut reader, 4).await.unwrap() {
2021 FrameRead::Frame(frame) => assert_eq!(frame, b"{}"),
2022 _ => panic!("expected the frame following the oversized line"),
2023 }
2024 }
2025
2026 #[tokio::test]
2027 async fn bounded_reader_accepts_exact_crlf_limit() {
2028 let mut reader = BufReader::new(Cursor::new(b"1234\r\n".to_vec()));
2029 match read_bounded_frame(&mut reader, 4).await.unwrap() {
2030 FrameRead::Frame(frame) => assert_eq!(frame, b"1234"),
2031 _ => panic!("expected an exact-limit frame"),
2032 }
2033 }
2034
2035 #[tokio::test]
2036 async fn outbound_byte_backpressure_is_cancellation_aware() {
2037 let (output, mut receiver) = outbound_channel(5);
2038 output.send(vec![0; 4], None).await.unwrap();
2039
2040 let cancellation = CancellationToken::new();
2041 let blocked = tokio::spawn({
2042 let output = output.clone();
2043 let cancellation = cancellation.clone();
2044 async move { output.send(vec![1; 4], Some(&cancellation)).await }
2045 });
2046 tokio::task::yield_now().await;
2047 assert!(!blocked.is_finished());
2048
2049 cancellation.cancel();
2050 assert_eq!(blocked.await.unwrap(), Err(OutboundSendError::Cancelled));
2051
2052 drop(receiver.recv().await.unwrap());
2053 output.send(vec![2; 4], None).await.unwrap();
2054 }
2055
2056 #[tokio::test]
2057 async fn outbound_control_send_times_out_under_byte_backpressure() {
2058 let (output, _receiver) = outbound_channel(5);
2059 output.send(vec![0; 4], None).await.unwrap();
2060 let result = output
2061 .send_with_timeout(vec![1; 4], None, Duration::from_millis(10))
2062 .await;
2063 assert_eq!(result, Err(OutboundSendError::TimedOut));
2064 }
2065
2066 #[tokio::test]
2067 async fn active_turn_shutdown_aborts_after_grace_period() {
2068 let cancellation = CancellationToken::new();
2069 let dropped = Arc::new(AtomicBool::new(false));
2070 let (started_tx, started_rx) = oneshot::channel();
2071 let task = tokio::spawn({
2072 let dropped = Arc::clone(&dropped);
2073 async move {
2074 let _signal = DropSignal(dropped);
2075 let _ = started_tx.send(());
2076 pending::<()>().await;
2077 }
2078 });
2079 started_rx.await.unwrap();
2080
2081 let graceful = shutdown_active_turn(
2082 ActiveTurn {
2083 turn_id: "turn".into(),
2084 cancellation,
2085 task,
2086 },
2087 Duration::from_millis(10),
2088 )
2089 .await;
2090
2091 assert!(!graceful);
2092 assert!(dropped.load(Ordering::Acquire));
2093 }
2094
2095 #[tokio::test]
2096 async fn writer_shutdown_aborts_after_grace_period() {
2097 let dropped = Arc::new(AtomicBool::new(false));
2098 let (started_tx, started_rx) = oneshot::channel();
2099 let writer = tokio::spawn({
2100 let dropped = Arc::clone(&dropped);
2101 async move {
2102 let _signal = DropSignal(dropped);
2103 let _ = started_tx.send(());
2104 pending::<std::io::Result<()>>().await
2105 }
2106 });
2107 started_rx.await.unwrap();
2108
2109 let result = shutdown_writer(writer, Duration::from_millis(10)).await;
2110
2111 assert!(result.is_err());
2112 assert!(dropped.load(Ordering::Acquire));
2113 }
2114
2115 #[tokio::test]
2116 async fn workspace_projects_list_their_agent_skills_for_delegation() {
2117 use std::os::unix::fs::symlink;
2118 let temporary = tempfile::tempdir().unwrap();
2119 let workspace = temporary.path().canonicalize().unwrap();
2120 let outside = tempfile::tempdir().unwrap();
2121 let outside = outside.path().canonicalize().unwrap();
2122 let write_skill = |directory: &Path, description: &str| {
2123 std::fs::create_dir_all(directory).unwrap();
2124 std::fs::write(
2125 directory.join("SKILL.md"),
2126 format!(
2127 "---\nname: skill\ndescription: {description}\n---\nBody of {description}\n"
2128 ),
2129 )
2130 .unwrap();
2131 };
2132 write_skill(
2133 &workspace.join("scv/.agents/skills/feature-flow"),
2134 "Land SCV",
2135 );
2136 std::fs::create_dir_all(workspace.join("scv/.claude/skills")).unwrap();
2137 symlink(
2138 "../../.agents/skills/feature-flow",
2139 workspace.join("scv/.claude/skills/feature-flow"),
2140 )
2141 .unwrap();
2142 write_skill(
2143 &workspace.join("web/.claude/skills/deploy"),
2144 "Deploy the site",
2145 );
2146 write_skill(&workspace.join(".agents/skills/triage"), "Root triage");
2147 write_skill(&workspace.join(".agents/skills/notes"), "Root notes");
2148 write_skill(&workspace.join(".scv/skills/triage"), "SCV triage");
2149 write_skill(&workspace.join(".hidden/.agents/skills/secret"), "Hidden");
2150 write_skill(&outside.join(".agents/skills/evil"), "Outside");
2151 symlink(&outside, workspace.join("escape")).unwrap();
2152 std::fs::create_dir_all(workspace.join("rogue/.agents")).unwrap();
2153 symlink(
2154 outside.join(".agents/skills"),
2155 workspace.join("rogue/.agents/skills"),
2156 )
2157 .unwrap();
2158 std::fs::create_dir_all(workspace.join("sneaky/.agents/skills/leak")).unwrap();
2159 symlink(
2160 outside.join(".agents/skills/evil/SKILL.md"),
2161 workspace.join("sneaky/.agents/skills/leak/SKILL.md"),
2162 )
2163 .unwrap();
2164 std::fs::write(workspace.join("file"), "not a project").unwrap();
2165 let mut config = Config::default();
2166 config.skills.user_dir = workspace.join("no-user-skills");
2167
2168 let skills = discover_skills(&workspace, &config, true).unwrap();
2169 let mut names: Vec<_> = skills.map.keys().cloned().collect();
2170 names.sort();
2171 assert_eq!(names, ["notes", "scv:feature-flow", "triage", "web:deploy"]);
2172 assert_eq!(
2173 skills.map["triage"],
2174 workspace.join(".scv/skills/triage/SKILL.md")
2175 );
2176 assert_eq!(
2177 skills.project_listing,
2178 "- notes (workspace root): Root notes\n\
2179 - scv:feature-flow (project scv): Land SCV\n\
2180 - web:deploy (project web): Deploy the site\n"
2181 );
2182 let prompt = build_system_prompt(&workspace, &config, &skills).unwrap();
2183 assert!(prompt.contains("# Project skills"));
2184 assert!(prompt.contains("set its cwd to the skill's project"));
2185
2186 let registry = builtin_registry(
2187 config.tools(),
2188 skills.map,
2189 skills.roots,
2190 config.skills.max_skill_bytes,
2191 HashMap::new(),
2192 )
2193 .unwrap();
2194 let read_skill = registry.get("read_skill").unwrap();
2195 let loaded = read_skill
2196 .execute(
2197 serde_json::json!({"name":"scv:feature-flow"}),
2198 scv_core::ToolContext {
2199 workspace: workspace.clone(),
2200 cancellation: CancellationToken::new(),
2201 },
2202 )
2203 .await
2204 .unwrap();
2205 assert!(loaded.content.contains("Body of Land SCV"));
2206
2207 let tool_free = discover_skills(&workspace, &config, false).unwrap();
2208 assert!(tool_free.project_listing.is_empty());
2209 assert!(!tool_free.map.contains_key("scv:feature-flow"));
2210 config.skills.scan_projects = false;
2211 let disabled = discover_skills(&workspace, &config, true).unwrap();
2212 assert!(disabled.project_listing.is_empty());
2213 config.skills.scan_projects = true;
2214 config.skills.max_skills = 3;
2215 let capped = discover_skills(&workspace, &config, true).unwrap();
2216 assert_eq!(capped.map.len(), 3);
2217 assert!(capped.map.contains_key("scv:feature-flow"));
2218 assert!(!capped.map.contains_key("web:deploy"));
2219 }
2220
2221 #[tokio::test]
2222 async fn daemon_control_lists_and_stops_delegations() {
2223 use std::os::unix::process::CommandExt as _;
2224 let home = tempfile::tempdir().unwrap();
2225 let registry = DelegationRegistry::new(home.path());
2226 let components = Arc::new(Mutex::new(components::Components::new(
2227 PathBuf::from("/unused.sock"),
2228 PathBuf::from("/"),
2229 )));
2230 let mut owner = std::process::Command::new("sleep")
2232 .arg("30")
2233 .spawn()
2234 .unwrap();
2235 let mut agent = std::process::Command::new("sleep")
2236 .arg("30")
2237 .process_group(0)
2238 .spawn()
2239 .unwrap();
2240 let identity = |pid| delegations::ProcessIdentity::of(pid).unwrap();
2241 let record = delegations::DelegationRecord {
2242 handle: "codex-a1b2c3".into(),
2243 agent: "codex".into(),
2244 instance: registry.instance().into(),
2245 session: "session".into(),
2246 owner: identity(owner.id()),
2247 process: identity(agent.id()),
2248 pgid: agent.id(),
2249 cwd: "/work/project\u{7}".into(),
2250 started_unix: 1,
2251 depth: 1,
2252 };
2253 std::fs::create_dir_all(registry.record_dir()).unwrap();
2254 std::fs::write(
2255 registry.record_dir().join("codex-a1b2c3.json"),
2256 serde_json::to_vec(&record).unwrap(),
2257 )
2258 .unwrap();
2259 let control = |command| daemon_control(&components, ®istry, command);
2260 let Ok(status) = control(DaemonCommand::Status).await else {
2261 panic!("status failed");
2262 };
2263 assert_eq!(status.delegations.active, 1);
2264 assert!(status.delegations.entries.is_empty());
2265 let Ok(status) = control(DaemonCommand::Delegations { all: false }).await else {
2266 panic!("listing failed");
2267 };
2268 let [entry] = status.delegations.entries.as_slice() else {
2269 panic!("{:?}", status.delegations);
2270 };
2271 assert_eq!(entry.handle, "codex-a1b2c3");
2272 assert_eq!(entry.pid, agent.id());
2273 assert_eq!(entry.owner_pid, owner.id());
2274 assert!(!entry.orphaned);
2275 assert_eq!(entry.processes, 1);
2276 for command in [
2277 DaemonCommand::DelegationKill {
2278 handle: Some("codex-nosuch".into()),
2279 orphans: false,
2280 },
2281 DaemonCommand::DelegationKill {
2282 handle: None,
2283 orphans: false,
2284 },
2285 ] {
2286 assert!(matches!(
2287 control(command).await,
2288 Err(ControlFailure::Delegation(_))
2289 ));
2290 }
2291 let Ok(status) = control(DaemonCommand::DelegationKill {
2292 handle: Some("codex-a1b2c3".into()),
2293 orphans: false,
2294 })
2295 .await
2296 else {
2297 panic!("kill failed");
2298 };
2299 assert_eq!(status.delegations.killed, ["codex-a1b2c3"]);
2300 assert!(agent.wait().unwrap().code().is_none());
2301 owner.kill().unwrap();
2304 owner.wait().unwrap();
2305 let Ok(status) = control(DaemonCommand::Delegations { all: true }).await else {
2306 panic!("listing failed");
2307 };
2308 assert!(status.delegations.entries[0].orphaned);
2309 assert_eq!(status.delegations.active, 0);
2310 let Ok(_) = control(DaemonCommand::DelegationKill {
2311 handle: None,
2312 orphans: true,
2313 })
2314 .await
2315 else {
2316 panic!("orphan sweep failed");
2317 };
2318 assert!(registry.list(true).is_empty());
2319 }
2320}