Skip to main content

codewhale_core/
lib.rs

1use std::collections::{HashMap, HashSet};
2use std::path::{Path, PathBuf};
3use std::sync::Arc;
4
5use std::time::Duration;
6
7use anyhow::Result;
8use codewhale_agent::ModelRegistry;
9use codewhale_config::{CliRuntimeOverrides, ConfigToml, ProviderKind};
10use codewhale_execpolicy::{
11    AskForApproval, ExecApprovalRequirement, ExecPolicyContext, ExecPolicyDecision,
12    ExecPolicyEngine,
13};
14use codewhale_hooks::{HookDispatcher, HookEvent};
15use codewhale_mcp::{
16    McpManager, McpStartupCompleteEvent, McpStartupStatus as McpManagerStartupStatus,
17};
18use codewhale_protocol::{
19    AppResponse, EventFrame, ExecApprovalRequestEvent, PromptRequest, PromptResponse,
20    ResponseChannel, ReviewDecision, Status, Thread, ThreadForkParams, ThreadGoal,
21    ThreadGoalClearParams, ThreadGoalGetParams, ThreadGoalProgressParams, ThreadGoalSetParams,
22    ThreadGoalStatus, ThreadListParams, ThreadReadParams, ThreadRequest, ThreadResponse,
23    ThreadResumeParams, ThreadSetNameParams, ThreadStatus, ToolPayload, UserInputRequestEvent,
24};
25use codewhale_state::{
26    JobStateRecord, JobStateStatus, SessionSource, StateStore, ThreadGoalRecord,
27    ThreadGoalStatus as PersistedThreadGoalStatus, ThreadListFilters, ThreadMetadata,
28    ThreadStatus as PersistedThreadStatus,
29};
30use codewhale_tools::{ToolCall, ToolRegistry};
31use serde_json::{Value, json};
32use tokio::time;
33use uuid::Uuid;
34
35/// Per-tool dispatch budget for the headless runtime. Matches the generous
36/// subagent default so long-running tools are not cut off prematurely.
37fn tool_dispatch_timeout() -> Duration {
38    if cfg!(test) {
39        Duration::from_millis(50)
40    } else {
41        Duration::from_secs(300)
42    }
43}
44
45/// How a new thread's conversation history is initialized.
46#[derive(Debug, Clone)]
47pub enum InitialHistory {
48    /// Start with an empty conversation.
49    New,
50    /// Forked from an existing thread with the given history items.
51    Forked(Vec<Value>),
52    /// Resumed from a persisted thread with its full history.
53    Resumed {
54        conversation_id: String,
55        history: Vec<Value>,
56        rollout_path: PathBuf,
57    },
58}
59
60/// Result of spawning or resuming a thread.
61#[derive(Debug, Clone)]
62pub struct NewThread {
63    /// The thread metadata.
64    pub thread: Thread,
65    /// Resolved model identifier.
66    pub model: String,
67    /// Provider that serves the model.
68    pub model_provider: String,
69    /// Working directory for the thread.
70    pub cwd: PathBuf,
71    /// Approval policy override, if any.
72    pub approval_policy: Option<String>,
73    /// Sandbox mode override, if any.
74    pub sandbox: Option<String>,
75}
76
77/// Status of a background job.
78#[derive(Debug, Clone, Copy, PartialEq, Eq)]
79pub enum JobStatus {
80    /// Waiting to be picked up.
81    Queued,
82    /// Currently executing.
83    Running,
84    /// Temporarily paused.
85    Paused,
86    /// Finished successfully.
87    Completed,
88    /// Finished with an error.
89    Failed,
90    /// Cancelled by the user.
91    Cancelled,
92}
93
94impl Status for JobStatus {
95    fn is_terminal(&self) -> bool {
96        matches!(self, Self::Completed | Self::Failed | Self::Cancelled)
97    }
98    fn is_active(&self) -> bool {
99        matches!(self, Self::Queued | Self::Running)
100    }
101    fn is_paused(&self) -> bool {
102        matches!(self, Self::Paused)
103    }
104}
105
106const JOB_DETAIL_SCHEMA_VERSION: u8 = 1;
107const DEFAULT_JOB_MAX_ATTEMPTS: u32 = 3;
108const DEFAULT_JOB_BACKOFF_BASE_MS: u64 = 500;
109const MAX_JOB_HISTORY_ENTRIES: usize = 64;
110
111/// Retry state for a job that failed and may be retried.
112#[derive(Debug, Clone)]
113pub struct JobRetryMetadata {
114    /// Current attempt number (0 = not yet retried).
115    pub attempt: u32,
116    /// Maximum number of retry attempts before giving up.
117    pub max_attempts: u32,
118    /// Base delay in milliseconds for exponential backoff.
119    pub backoff_base_ms: u64,
120    /// Computed delay in milliseconds until the next retry.
121    pub next_backoff_ms: u64,
122    /// Timestamp when the next retry should be attempted.
123    pub next_retry_at: Option<i64>,
124}
125
126impl Default for JobRetryMetadata {
127    fn default() -> Self {
128        Self {
129            attempt: 0,
130            max_attempts: DEFAULT_JOB_MAX_ATTEMPTS,
131            backoff_base_ms: DEFAULT_JOB_BACKOFF_BASE_MS,
132            next_backoff_ms: 0,
133            next_retry_at: None,
134        }
135    }
136}
137
138/// A single entry in a job's history log.
139#[derive(Debug, Clone)]
140pub struct JobHistoryEntry {
141    /// Timestamp when this entry was recorded.
142    pub at: i64,
143    /// Phase name (e.g., "created", "running", "failed").
144    pub phase: String,
145    /// Job status at this point in time.
146    pub status: JobStatus,
147    /// Progress percentage at this point, if available.
148    pub progress: Option<u8>,
149    /// Human-readable detail message.
150    pub detail: Option<String>,
151    /// Retry state snapshot at this point.
152    pub retry: JobRetryMetadata,
153}
154
155#[derive(Debug, Clone)]
156struct PersistedJobDetail {
157    pub status: JobStatus,
158    pub detail: Option<String>,
159    pub retry: JobRetryMetadata,
160    pub history: Vec<JobHistoryEntry>,
161}
162
163/// A complete job record with all metadata and history.
164#[derive(Debug, Clone)]
165pub struct JobRecord {
166    /// Unique job identifier.
167    pub id: String,
168    /// Human-readable job name.
169    pub name: String,
170    /// Current job status.
171    pub status: JobStatus,
172    /// Current progress percentage (0-100).
173    pub progress: Option<u8>,
174    /// Human-readable detail about the current state.
175    pub detail: Option<String>,
176    /// Retry state for failed jobs.
177    pub retry: JobRetryMetadata,
178    /// Chronological history of state transitions.
179    pub history: Vec<JobHistoryEntry>,
180    /// Timestamp when the job was created.
181    pub created_at: i64,
182    /// Timestamp of the last state change.
183    pub updated_at: i64,
184}
185
186/// Map a durable [`JobRecord`] to the dependency-neutral run read model.
187///
188/// Pure projection of the record as persisted: unknown budgets stay unset and
189/// nothing is fabricated. `updated_at` (epoch seconds) provides the terminal
190/// timestamp because the job manager records no separate end time. The
191/// free-form job detail is intentionally omitted because this owner does not
192/// classify it as safe for a cross-surface read model.
193#[must_use]
194pub fn job_record_to_agent_run(
195    record: &JobRecord,
196) -> codewhale_protocol::agent_run::AgentRunSnapshot {
197    use codewhale_protocol::agent_run::{
198        AgentRunSnapshot, BudgetSummary, RunSource, RunState, TerminalOutcome, TerminalSummary,
199    };
200
201    let (state, terminal) = match record.status {
202        JobStatus::Queued => (RunState::Queued, None),
203        JobStatus::Running => (RunState::Running, None),
204        JobStatus::Paused => (RunState::Paused, None),
205        JobStatus::Completed | JobStatus::Failed | JobStatus::Cancelled => {
206            let outcome = match record.status {
207                JobStatus::Completed => TerminalOutcome::Completed,
208                JobStatus::Failed => TerminalOutcome::Failed,
209                _ => TerminalOutcome::Cancelled,
210            };
211            (
212                RunState::Terminal,
213                Some(TerminalSummary {
214                    outcome,
215                    ended_at_ms: record.updated_at.checked_mul(1000),
216                    detail: None,
217                }),
218            )
219        }
220    };
221
222    AgentRunSnapshot {
223        run_id: record.id.clone(),
224        parent: None,
225        source: RunSource::CoreJob,
226        state,
227        budget: BudgetSummary::default(),
228        terminal,
229        refs: Vec::new(),
230    }
231}
232
233/// Manages background jobs with retry logic and persistence.
234#[derive(Debug, Default)]
235pub struct JobManager {
236    jobs: HashMap<String, JobRecord>,
237}
238
239impl JobManager {
240    fn now_ts() -> i64 {
241        chrono::Utc::now().timestamp()
242    }
243
244    fn deterministic_backoff_ms(retry: &JobRetryMetadata) -> u64 {
245        if retry.attempt == 0 {
246            return 0;
247        }
248        let exponent = retry.attempt.saturating_sub(1).min(20);
249        let multiplier = 1u64.checked_shl(exponent).unwrap_or(u64::MAX);
250        retry.backoff_base_ms.saturating_mul(multiplier)
251    }
252
253    fn clear_retry_schedule(retry: &mut JobRetryMetadata) {
254        retry.next_backoff_ms = 0;
255        retry.next_retry_at = None;
256    }
257
258    fn push_history(job: &mut JobRecord, phase: &str) {
259        job.history.push(JobHistoryEntry {
260            at: job.updated_at,
261            phase: phase.to_string(),
262            status: job.status,
263            progress: job.progress,
264            detail: job.detail.clone(),
265            retry: job.retry.clone(),
266        });
267        if job.history.len() > MAX_JOB_HISTORY_ENTRIES {
268            let to_drain = job.history.len() - MAX_JOB_HISTORY_ENTRIES;
269            job.history.drain(0..to_drain);
270        }
271    }
272
273    fn parse_persisted_detail(raw: Option<&str>) -> Option<PersistedJobDetail> {
274        let raw = raw?;
275        let parsed: Value = serde_json::from_str(raw).ok()?;
276        let status = parsed
277            .get("status")
278            .and_then(Value::as_str)
279            .and_then(job_status_from_str)?;
280        let detail = parsed.get("detail").and_then(json_optional_string);
281        let retry = parse_retry_metadata(parsed.get("retry"));
282        let history = parsed
283            .get("history")
284            .and_then(Value::as_array)
285            .map(|items| {
286                items
287                    .iter()
288                    .filter_map(parse_history_entry)
289                    .collect::<Vec<_>>()
290            })
291            .unwrap_or_default();
292        Some(PersistedJobDetail {
293            status,
294            detail,
295            retry,
296            history,
297        })
298    }
299
300    fn encode_persisted_detail(job: &JobRecord) -> Result<Option<String>> {
301        let encoded = json!({
302            "schema_version": JOB_DETAIL_SCHEMA_VERSION,
303            "status": job_status_to_str(job.status),
304            "detail": job.detail.clone(),
305            "retry": job_retry_to_value(&job.retry),
306            "history": job.history.iter().map(job_history_to_value).collect::<Vec<_>>()
307        })
308        .to_string();
309        Ok(Some(encoded))
310    }
311
312    /// Enqueues a new job and returns its record.
313    pub fn enqueue(&mut self, name: impl Into<String>) -> JobRecord {
314        let now = Self::now_ts();
315        let id = format!("job-{}", Uuid::new_v4());
316        let mut job = JobRecord {
317            id: id.clone(),
318            name: name.into(),
319            status: JobStatus::Queued,
320            progress: Some(0),
321            detail: None,
322            retry: JobRetryMetadata::default(),
323            history: Vec::new(),
324            created_at: now,
325            updated_at: now,
326        };
327        Self::push_history(&mut job, "created");
328        self.jobs.insert(id, job.clone());
329        job
330    }
331
332    /// Transitions a job to running and clears its retry schedule.
333    pub fn set_running(&mut self, id: &str) {
334        if let Some(job) = self.jobs.get_mut(id) {
335            job.status = JobStatus::Running;
336            Self::clear_retry_schedule(&mut job.retry);
337            job.updated_at = Self::now_ts();
338            Self::push_history(job, "running");
339        }
340    }
341
342    /// Updates a job's progress (clamped to 100) and optional detail message.
343    pub fn update_progress(&mut self, id: &str, progress: u8, detail: Option<String>) {
344        if let Some(job) = self.jobs.get_mut(id) {
345            job.progress = Some(progress.min(100));
346            job.detail = detail;
347            job.updated_at = Self::now_ts();
348            Self::push_history(job, "progress_updated");
349        }
350    }
351
352    /// Marks a job as completed with 100% progress and clears its retry schedule.
353    pub fn complete(&mut self, id: &str) {
354        if let Some(job) = self.jobs.get_mut(id) {
355            job.status = JobStatus::Completed;
356            job.progress = Some(100);
357            Self::clear_retry_schedule(&mut job.retry);
358            job.updated_at = Self::now_ts();
359            Self::push_history(job, "completed");
360        }
361    }
362
363    /// Marks a job as failed and schedules a retry if attempts remain.
364    pub fn fail(&mut self, id: &str, detail: impl Into<String>) {
365        if let Some(job) = self.jobs.get_mut(id) {
366            let now = Self::now_ts();
367            job.status = JobStatus::Failed;
368            job.detail = Some(detail.into());
369            if job.retry.attempt < job.retry.max_attempts {
370                job.retry.attempt += 1;
371                job.retry.next_backoff_ms = Self::deterministic_backoff_ms(&job.retry);
372                let delay_secs = ((job.retry.next_backoff_ms.saturating_add(999)) / 1000)
373                    .min(i64::MAX as u64) as i64;
374                job.retry.next_retry_at = Some(now.saturating_add(delay_secs));
375            } else {
376                Self::clear_retry_schedule(&mut job.retry);
377            }
378            job.updated_at = now;
379            Self::push_history(job, "failed");
380        }
381    }
382
383    /// Cancels a job and clears any pending retry schedule.
384    pub fn cancel(&mut self, id: &str) {
385        if let Some(job) = self.jobs.get_mut(id) {
386            job.status = JobStatus::Cancelled;
387            Self::clear_retry_schedule(&mut job.retry);
388            job.updated_at = Self::now_ts();
389            Self::push_history(job, "cancelled");
390        }
391    }
392
393    /// Pauses a job, optionally updating its detail message.
394    pub fn pause(&mut self, id: &str, detail: Option<String>) {
395        if let Some(job) = self.jobs.get_mut(id) {
396            job.status = JobStatus::Paused;
397            if detail.is_some() {
398                job.detail = detail;
399            }
400            job.updated_at = Self::now_ts();
401            Self::push_history(job, "paused");
402        }
403    }
404
405    /// Resumes a paused or failed job back to running status.
406    pub fn resume(&mut self, id: &str, detail: Option<String>) {
407        if let Some(job) = self.jobs.get_mut(id) {
408            job.status = JobStatus::Running;
409            if detail.is_some() {
410                job.detail = detail;
411            }
412            Self::clear_retry_schedule(&mut job.retry);
413            job.updated_at = Self::now_ts();
414            Self::push_history(job, "resumed");
415        }
416    }
417
418    /// Returns all jobs sorted by most recently updated first.
419    pub fn list(&self) -> Vec<JobRecord> {
420        let mut out = self.jobs.values().cloned().collect::<Vec<_>>();
421        out.sort_by_key(|job| std::cmp::Reverse(job.updated_at));
422        out
423    }
424
425    /// Returns the history entries for a job, or an empty vec if not found.
426    pub fn history(&self, id: &str) -> Vec<JobHistoryEntry> {
427        self.jobs
428            .get(id)
429            .map(|job| job.history.clone())
430            .unwrap_or_default()
431    }
432
433    /// Resets queued or running jobs back to queued on application resume.
434    pub fn resume_pending(&mut self) -> Vec<JobRecord> {
435        let mut resumed = Vec::new();
436        for job in self.jobs.values_mut() {
437            if matches!(job.status, JobStatus::Queued | JobStatus::Running) {
438                job.status = JobStatus::Queued;
439                job.updated_at = Self::now_ts();
440                Self::push_history(job, "queued_after_resume");
441                resumed.push(job.clone());
442            }
443        }
444        resumed
445    }
446
447    /// Loads jobs from the state store, deserializing extended detail when available.
448    pub fn load_from_store(&mut self, store: &StateStore) -> Result<()> {
449        let persisted = store.list_jobs(Some(500))?;
450        for job in persisted {
451            let fallback_status = job_state_status_to_runtime(job.status);
452            let parsed = Self::parse_persisted_detail(job.detail.as_deref());
453            let (status, detail, retry, history) = if let Some(detail_state) = parsed {
454                (
455                    detail_state.status,
456                    detail_state.detail,
457                    detail_state.retry,
458                    detail_state.history,
459                )
460            } else {
461                (
462                    fallback_status,
463                    job.detail,
464                    JobRetryMetadata::default(),
465                    Vec::new(),
466                )
467            };
468            self.jobs.insert(
469                job.id.clone(),
470                JobRecord {
471                    id: job.id,
472                    name: job.name,
473                    status,
474                    progress: job.progress,
475                    detail,
476                    retry,
477                    history,
478                    created_at: job.created_at,
479                    updated_at: job.updated_at,
480                },
481            );
482        }
483        Ok(())
484    }
485
486    /// Persists a single job's current state to the state store.
487    pub fn persist_job(&self, store: &StateStore, id: &str) -> Result<()> {
488        let Some(job) = self.jobs.get(id) else {
489            return Ok(());
490        };
491        let encoded_detail = Self::encode_persisted_detail(job)?;
492        store.upsert_job(&JobStateRecord {
493            id: job.id.clone(),
494            name: job.name.clone(),
495            status: runtime_status_to_job_state(job.status),
496            progress: job.progress,
497            detail: encoded_detail,
498            created_at: job.created_at,
499            updated_at: job.updated_at,
500        })
501    }
502
503    /// Persists all in-memory jobs to the state store.
504    pub fn persist_all(&self, store: &StateStore) -> Result<()> {
505        for id in self.jobs.keys() {
506            self.persist_job(store, id)?;
507        }
508        Ok(())
509    }
510}
511
512/// Manages thread lifecycle: spawn, resume, fork, archive, and persistence.
513pub struct ThreadManager {
514    store: StateStore,
515    running_threads: HashMap<String, Thread>,
516    cli_version: String,
517}
518
519impl ThreadManager {
520    /// Creates a new `ThreadManager` backed by the given state store.
521    pub fn new(store: StateStore) -> Self {
522        Self {
523            store,
524            running_threads: HashMap::new(),
525            cli_version: env!("CARGO_PKG_VERSION").to_string(),
526        }
527    }
528
529    /// Returns a reference to the underlying state store.
530    pub fn state_store(&self) -> &StateStore {
531        &self.store
532    }
533
534    /// Spawns a new thread with the given initial history and persists it.
535    pub fn spawn_thread_with_history(
536        &mut self,
537        model_provider: String,
538        cwd: PathBuf,
539        initial_history: InitialHistory,
540        persist_extended_history: bool,
541    ) -> Result<NewThread> {
542        let id = format!("thread-{}", Uuid::new_v4());
543        let now = chrono::Utc::now().timestamp();
544        let preview = preview_from_initial_history(&initial_history);
545        let source = match initial_history {
546            InitialHistory::New => SessionSource::Interactive,
547            InitialHistory::Forked(_) => SessionSource::Fork,
548            InitialHistory::Resumed { .. } => SessionSource::Resume,
549        };
550        let thread = Thread {
551            id: id.clone(),
552            preview,
553            ephemeral: !persist_extended_history,
554            model_provider: model_provider.clone(),
555            created_at: now,
556            updated_at: now,
557            status: ThreadStatus::Running,
558            path: None,
559            cwd: cwd.clone(),
560            cli_version: self.cli_version.clone(),
561            source: match source {
562                SessionSource::Interactive => codewhale_protocol::SessionSource::Interactive,
563                SessionSource::Resume => codewhale_protocol::SessionSource::Resume,
564                SessionSource::Fork => codewhale_protocol::SessionSource::Fork,
565                SessionSource::Api => codewhale_protocol::SessionSource::Api,
566                SessionSource::Unknown => codewhale_protocol::SessionSource::Unknown,
567            },
568            name: None,
569        };
570        self.persist_thread(&thread, None)?;
571        match &initial_history {
572            InitialHistory::Forked(items) => {
573                for item in items {
574                    self.store.append_message(
575                        &thread.id,
576                        "history",
577                        &item.to_string(),
578                        Some(item.clone()),
579                    )?;
580                }
581            }
582            InitialHistory::Resumed { history, .. } => {
583                for item in history {
584                    self.store.append_message(
585                        &thread.id,
586                        "history",
587                        &item.to_string(),
588                        Some(item.clone()),
589                    )?;
590                }
591            }
592            InitialHistory::New => {}
593        }
594        self.running_threads
595            .insert(thread.id.clone(), thread.clone());
596        Ok(NewThread {
597            thread,
598            model: "auto".to_string(),
599            model_provider,
600            cwd,
601            approval_policy: None,
602            sandbox: None,
603        })
604    }
605
606    /// Resumes an existing thread, returning `None` if not found.
607    pub fn resume_thread_with_history(
608        &mut self,
609        params: &ThreadResumeParams,
610        fallback_cwd: &Path,
611        model_provider: String,
612    ) -> Result<Option<NewThread>> {
613        if params.history.is_none()
614            && let Some(thread) = self.running_threads.get(&params.thread_id).cloned()
615        {
616            return Ok(Some(NewThread {
617                model: params.model.clone().unwrap_or_else(|| "auto".to_string()),
618                model_provider: params.model_provider.clone().unwrap_or(model_provider),
619                cwd: params.cwd.clone().unwrap_or_else(|| thread.cwd.clone()),
620                approval_policy: params.approval_policy.clone(),
621                sandbox: params.sandbox.clone(),
622                thread,
623            }));
624        }
625
626        let persisted = self.store.get_thread(&params.thread_id)?;
627        let Some(metadata) = persisted else {
628            return Ok(None);
629        };
630        let mut thread = to_protocol_thread(metadata);
631        thread.status = ThreadStatus::Running;
632        thread.updated_at = chrono::Utc::now().timestamp();
633        thread.cwd = params
634            .cwd
635            .clone()
636            .unwrap_or_else(|| fallback_cwd.to_path_buf());
637        self.persist_thread(&thread, None)?;
638        self.running_threads
639            .insert(thread.id.clone(), thread.clone());
640        if let Some(history) = params.history.as_ref() {
641            // A read→resume flow hands back items that are already on the
642            // persisted chain; appending them again would double the
643            // conversation on every resume, compounding. Dedup by content
644            // fingerprint (the item's JSON, matching what append_message
645            // stores as content) against the persisted chain and against
646            // items already appended in this loop.
647            let mut seen: HashSet<String> = self
648                .store
649                .list_messages(&thread.id, None)?
650                .into_iter()
651                .map(|message| {
652                    message
653                        .item
654                        .as_ref()
655                        .map_or(message.content.clone(), |item| item.to_string())
656                })
657                .collect();
658            for item in history {
659                let fingerprint = item.to_string();
660                if !seen.insert(fingerprint.clone()) {
661                    continue;
662                }
663                self.store.append_message(
664                    &thread.id,
665                    "history",
666                    &fingerprint,
667                    Some(item.clone()),
668                )?;
669            }
670        }
671
672        Ok(Some(NewThread {
673            model: params.model.clone().unwrap_or_else(|| "auto".to_string()),
674            model_provider: params.model_provider.clone().unwrap_or(model_provider),
675            cwd: thread.cwd.clone(),
676            approval_policy: params.approval_policy.clone(),
677            sandbox: params.sandbox.clone(),
678            thread,
679        }))
680    }
681
682    /// Forks an existing thread into a new one, inheriting the parent's provider.
683    pub fn fork_thread(
684        &mut self,
685        params: &ThreadForkParams,
686        fallback_cwd: &Path,
687    ) -> Result<Option<NewThread>> {
688        let parent = self.store.get_thread(&params.thread_id)?;
689        let Some(parent) = parent else {
690            return Ok(None);
691        };
692        let parent_thread = to_protocol_thread(parent);
693        let new = self.spawn_thread_with_history(
694            params
695                .model_provider
696                .clone()
697                .unwrap_or_else(|| parent_thread.model_provider.clone()),
698            params
699                .cwd
700                .clone()
701                .unwrap_or_else(|| fallback_cwd.to_path_buf()),
702            InitialHistory::Forked(vec![json!({
703                "type": "fork",
704                "from_thread_id": parent_thread.id
705            })]),
706            params.persist_extended_history,
707        )?;
708        Ok(Some(new))
709    }
710
711    /// Lists threads matching the given filter parameters.
712    pub fn list_threads(&self, params: &ThreadListParams) -> Result<Vec<Thread>> {
713        let list = self.store.list_threads(ThreadListFilters {
714            include_archived: params.include_archived,
715            limit: params.limit,
716        })?;
717        Ok(list.into_iter().map(to_protocol_thread).collect())
718    }
719
720    /// Reads a single thread by id, or `None` if not found.
721    pub fn read_thread(&self, params: &ThreadReadParams) -> Result<Option<Thread>> {
722        Ok(self
723            .store
724            .get_thread(&params.thread_id)?
725            .map(to_protocol_thread))
726    }
727
728    /// Sets the display name for a thread, returning the updated thread or `None`.
729    pub fn set_thread_name(&mut self, params: &ThreadSetNameParams) -> Result<Option<Thread>> {
730        let Some(mut metadata) = self.store.get_thread(&params.thread_id)? else {
731            return Ok(None);
732        };
733        metadata.name = Some(params.name.clone());
734        metadata.updated_at = chrono::Utc::now().timestamp();
735        self.store.upsert_thread(&metadata)?;
736        let updated = to_protocol_thread(metadata);
737        self.running_threads
738            .insert(updated.id.clone(), updated.clone());
739        Ok(Some(updated))
740    }
741
742    /// Sets or replaces the persisted goal for a thread.
743    pub fn set_thread_goal(&mut self, params: &ThreadGoalSetParams) -> Result<Option<ThreadGoal>> {
744        if self.store.get_thread(&params.thread_id)?.is_none() {
745            return Ok(None);
746        }
747        let now = chrono::Utc::now().timestamp();
748        let goal = ThreadGoalRecord {
749            thread_id: params.thread_id.clone(),
750            goal_id: format!("goal-{}", Uuid::new_v4()),
751            objective: params.objective.clone(),
752            status: PersistedThreadGoalStatus::Active,
753            token_budget: params.token_budget,
754            tokens_used: 0,
755            time_used_seconds: 0,
756            continuation_count: 0,
757            created_at: now,
758            updated_at: now,
759        };
760        self.store.upsert_thread_goal(&goal)?;
761        Ok(Some(to_protocol_goal(goal)))
762    }
763
764    /// Reads the persisted goal for a thread.
765    pub fn get_thread_goal(&self, params: &ThreadGoalGetParams) -> Result<Option<ThreadGoal>> {
766        Ok(self
767            .store
768            .get_thread_goal(&params.thread_id)?
769            .map(to_protocol_goal))
770    }
771
772    /// Accrues durable per-goal usage and/or a continuation pass for a thread.
773    pub fn record_thread_goal_progress(
774        &mut self,
775        params: &ThreadGoalProgressParams,
776    ) -> Result<Option<ThreadGoal>> {
777        if self.store.get_thread(&params.thread_id)?.is_none() {
778            return Ok(None);
779        }
780
781        let now = chrono::Utc::now().timestamp();
782        let mut goal = if params.token_delta != 0 || params.time_delta_seconds != 0 {
783            self.store.record_thread_goal_usage(
784                &params.thread_id,
785                params.token_delta,
786                params.time_delta_seconds,
787                now,
788            )?
789        } else {
790            self.store.get_thread_goal(&params.thread_id)?
791        };
792
793        if params.record_continuation {
794            goal = self
795                .store
796                .record_thread_goal_continuation(&params.thread_id, now)?;
797        }
798
799        Ok(goal.map(to_protocol_goal))
800    }
801
802    /// Clears the persisted goal for a thread, returning whether one existed.
803    pub fn clear_thread_goal(&mut self, params: &ThreadGoalClearParams) -> Result<bool> {
804        self.store.delete_thread_goal(&params.thread_id)
805    }
806
807    /// Archives a thread so it no longer appears in default listings.
808    pub fn archive_thread(&mut self, thread_id: &str) -> Result<()> {
809        self.store.mark_archived(thread_id)?;
810        if let Some(thread) = self.running_threads.get_mut(thread_id) {
811            thread.status = ThreadStatus::Archived;
812        }
813        Ok(())
814    }
815
816    /// Restores an archived thread to active status.
817    pub fn unarchive_thread(&mut self, thread_id: &str) -> Result<()> {
818        self.store.mark_unarchived(thread_id)?;
819        if let Some(metadata) = self.store.get_thread(thread_id)? {
820            let thread = to_protocol_thread(metadata);
821            if let Some(cached) = self.running_threads.get_mut(thread_id) {
822                *cached = thread;
823            }
824        }
825        Ok(())
826    }
827
828    /// Records a user message in a thread and updates its preview and timestamp.
829    pub fn touch_message(&mut self, thread_id: &str, input: &str) -> Result<()> {
830        let Some(mut metadata) = self.store.get_thread(thread_id)? else {
831            return Ok(());
832        };
833        metadata.updated_at = chrono::Utc::now().timestamp();
834        metadata.preview = truncate_preview(input);
835        metadata.status = PersistedThreadStatus::Running;
836        self.store.upsert_thread(&metadata)?;
837        if let Some(thread) = self.running_threads.get_mut(thread_id) {
838            thread.updated_at = metadata.updated_at;
839            thread.preview = metadata.preview;
840            thread.status = ThreadStatus::Running;
841        }
842        let message_id = self.store.append_message(thread_id, "user", input, None)?;
843        self.store.save_checkpoint(
844            thread_id,
845            "latest",
846            &json!({
847                "reason": "thread_message",
848                "message_id": message_id,
849                "role": "user",
850                "preview": truncate_preview(input),
851                "updated_at": metadata.updated_at
852            }),
853        )?;
854        Ok(())
855    }
856
857    fn persist_thread(&self, thread: &Thread, rollout_path: Option<PathBuf>) -> Result<()> {
858        // This update payload carries no per-thread policy, so preserve any
859        // policy already stored for the thread rather than erasing it with
860        // NULLs on every persist/resume.
861        let existing = self.store.get_thread(&thread.id)?;
862        self.store.upsert_thread(&ThreadMetadata {
863            id: thread.id.clone(),
864            rollout_path,
865            preview: thread.preview.clone(),
866            ephemeral: thread.ephemeral,
867            model_provider: thread.model_provider.clone(),
868            created_at: thread.created_at,
869            updated_at: thread.updated_at,
870            status: to_persisted_status(&thread.status),
871            path: thread.path.clone(),
872            cwd: thread.cwd.clone(),
873            cli_version: thread.cli_version.clone(),
874            source: to_persisted_source(&thread.source),
875            name: thread.name.clone(),
876            sandbox_policy: existing
877                .as_ref()
878                .and_then(|metadata| metadata.sandbox_policy.clone()),
879            approval_mode: existing
880                .as_ref()
881                .and_then(|metadata| metadata.approval_mode.clone()),
882            archived: matches!(thread.status, ThreadStatus::Archived),
883            archived_at: None,
884            git_sha: None,
885            git_branch: None,
886            git_origin_url: None,
887            memory_mode: None,
888            current_leaf_id: None,
889        })
890    }
891}
892
893/// Top-level runtime combining config, model registry, threads, tools, MCP, and hooks.
894pub struct Runtime {
895    /// Resolved application configuration.
896    pub config: ConfigToml,
897    /// Registry of available model providers.
898    pub model_registry: ModelRegistry,
899    /// Manages conversation thread lifecycle.
900    pub thread_manager: ThreadManager,
901    /// Registry of callable tools.
902    pub tool_registry: Arc<ToolRegistry>,
903    /// Manager for MCP server connections.
904    pub mcp_manager: Arc<McpManager>,
905    /// Engine for evaluating execution policy decisions.
906    pub exec_policy: ExecPolicyEngine,
907    /// Dispatcher for lifecycle hooks.
908    pub hooks: HookDispatcher,
909    /// Manager for background job lifecycle.
910    pub jobs: JobManager,
911}
912
913impl Runtime {
914    /// Constructs a new `Runtime`, loading existing jobs from the state store.
915    pub fn new(
916        config: ConfigToml,
917        model_registry: ModelRegistry,
918        state: StateStore,
919        tool_registry: Arc<ToolRegistry>,
920        mcp_manager: Arc<McpManager>,
921        exec_policy: ExecPolicyEngine,
922        hooks: HookDispatcher,
923    ) -> Self {
924        let mut jobs = JobManager::default();
925        if let Err(e) = jobs.load_from_store(&state) {
926            tracing::warn!("Failed to load job store, starting with empty job list: {e}");
927        }
928        Self {
929            config,
930            model_registry,
931            thread_manager: ThreadManager::new(state),
932            tool_registry,
933            mcp_manager,
934            exec_policy,
935            hooks,
936            jobs,
937        }
938    }
939
940    /// Update the live configuration in-place so the next turn picks up
941    /// changes without a restart.  Called by the app-server after
942    /// `ConfigSet` or `ConfigUnset`.
943    ///
944    /// Only `config.toml` is touched by those operations, so the sibling
945    /// `permissions.toml` (and therefore `exec_policy`) is left unchanged.
946    ///
947    /// Fields that the TUI caches on its `App` struct (`api_provider`,
948    /// `reasoning_effort`, `mcp_config_path`, `skills_dir`, …) are read
949    /// live from `self.config` here via `resolve_runtime_options`, so they
950    /// take effect on the next prompt turn without any extra plumbing.
951    pub fn update_config(&mut self, config: ConfigToml) {
952        self.config = config;
953    }
954
955    /// Reload the live configuration **and** the exec policy from a
956    /// freshly-loaded `ConfigStore`.  Used by the app-server's
957    /// `ConfigReload` request, which re-reads both `config.toml` and the
958    /// sibling `permissions.toml` from disk.
959    ///
960    /// Unlike `update_config`, this also refreshes `self.exec_policy` so
961    /// externally edited permission rules take effect without a restart.
962    ///
963    /// Mirrors the TUI `reload_runtime_config` codepath for everything
964    /// that is reachable from the headless `Runtime`. The TUI-only caches
965    /// (`last_effective_reasoning_effort`, `model_compaction_budget`,
966    /// `ui_locale`, …) do not exist on `Runtime` and need no work here.
967    ///
968    /// **Not** refreshed by this call:
969    /// * `mcp_manager` — MCP server connections are loaded once at
970    ///   startup from `mcp_config_path`. Changing `mcp_config_path` or the
971    ///   referenced `mcp.json` still requires a headless-runtime restart;
972    ///   the TUI owns a separate explicit `/mcp reload` operation.
973    /// * `tool_registry` — built once at startup.
974    /// * `model_registry` — static catalog.
975    pub fn reload_config_and_policy(&mut self, config: ConfigToml, exec_policy: ExecPolicyEngine) {
976        self.config = config;
977        self.exec_policy = exec_policy;
978    }
979
980    fn persisted_thread_data(&self, thread_id: &str) -> Result<Value> {
981        let history = self
982            .thread_manager
983            .state_store()
984            .list_messages(thread_id, Some(500))?
985            .into_iter()
986            .map(|message| {
987                json!({
988                    "id": message.id,
989                    "role": message.role,
990                    "content": message.content,
991                    "item": message.item,
992                    "created_at": message.created_at
993                })
994            })
995            .collect::<Vec<_>>();
996
997        let checkpoint = self
998            .thread_manager
999            .state_store()
1000            .load_checkpoint(thread_id, None)?
1001            .map(|record| {
1002                json!({
1003                    "checkpoint_id": record.checkpoint_id,
1004                    "state": record.state,
1005                    "created_at": record.created_at
1006                })
1007            });
1008
1009        let goal = self
1010            .thread_manager
1011            .state_store()
1012            .get_thread_goal(thread_id)?
1013            .map(to_protocol_goal);
1014
1015        Ok(json!({
1016            "history": history,
1017            "checkpoint": checkpoint,
1018            "goal": goal
1019        }))
1020    }
1021
1022    fn persist_latest_checkpoint(&self, thread_id: &str, reason: &str, state: Value) -> Result<()> {
1023        self.thread_manager.state_store().save_checkpoint(
1024            thread_id,
1025            "latest",
1026            &json!({
1027                "reason": reason,
1028                "saved_at": chrono::Utc::now().timestamp(),
1029                "state": state
1030            }),
1031        )
1032    }
1033
1034    /// Dispatches a thread request (create, start, resume, fork, list, read, etc.).
1035    pub async fn handle_thread(&mut self, req: ThreadRequest) -> Result<ThreadResponse> {
1036        match req {
1037            ThreadRequest::Create { .. } => {
1038                let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
1039                let new = self.thread_manager.spawn_thread_with_history(
1040                    "deepseek".to_string(),
1041                    cwd,
1042                    InitialHistory::New,
1043                    false,
1044                )?;
1045                let mut response = thread_response_from_new("created", new);
1046                response.data = self.persisted_thread_data(&response.thread_id)?;
1047                Ok(response)
1048            }
1049            ThreadRequest::Start(params) => {
1050                let cwd = params.cwd.clone().unwrap_or_else(|| {
1051                    std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))
1052                });
1053                let new = self.thread_manager.spawn_thread_with_history(
1054                    params
1055                        .model_provider
1056                        .clone()
1057                        .unwrap_or_else(|| "deepseek".to_string()),
1058                    cwd,
1059                    InitialHistory::New,
1060                    params.persist_extended_history,
1061                )?;
1062                let mut response = thread_response_from_new("started", new);
1063                response.data = self.persisted_thread_data(&response.thread_id)?;
1064                Ok(response)
1065            }
1066            ThreadRequest::Resume(params) => {
1067                let fallback_cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
1068                if let Some(new) = self.thread_manager.resume_thread_with_history(
1069                    &params,
1070                    &fallback_cwd,
1071                    "deepseek".to_string(),
1072                )? {
1073                    let mut response = thread_response_from_new("resumed", new);
1074                    response.data = self.persisted_thread_data(&response.thread_id)?;
1075                    Ok(response)
1076                } else {
1077                    Ok(ThreadResponse {
1078                        thread_id: params.thread_id,
1079                        status: "missing".to_string(),
1080                        thread: None,
1081                        threads: Vec::new(),
1082                        goal: None,
1083                        model: None,
1084                        model_provider: None,
1085                        cwd: None,
1086                        approval_policy: params.approval_policy,
1087                        sandbox: params.sandbox,
1088                        events: Vec::new(),
1089                        data: json!({"error":"thread not found"}),
1090                    })
1091                }
1092            }
1093            ThreadRequest::Fork(params) => {
1094                let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
1095                if let Some(new) = self.thread_manager.fork_thread(&params, &cwd)? {
1096                    let mut response = thread_response_from_new("forked", new);
1097                    response.data = self.persisted_thread_data(&response.thread_id)?;
1098                    Ok(response)
1099                } else {
1100                    Ok(ThreadResponse {
1101                        thread_id: params.thread_id,
1102                        status: "missing".to_string(),
1103                        thread: None,
1104                        threads: Vec::new(),
1105                        goal: None,
1106                        model: None,
1107                        model_provider: None,
1108                        cwd: None,
1109                        approval_policy: params.approval_policy,
1110                        sandbox: params.sandbox,
1111                        events: Vec::new(),
1112                        data: json!({"error":"thread not found"}),
1113                    })
1114                }
1115            }
1116            ThreadRequest::List(params) => Ok(ThreadResponse {
1117                thread_id: "list".to_string(),
1118                status: "ok".to_string(),
1119                thread: None,
1120                threads: self.thread_manager.list_threads(&params)?,
1121                goal: None,
1122                model: None,
1123                model_provider: None,
1124                cwd: None,
1125                approval_policy: None,
1126                sandbox: None,
1127                events: Vec::new(),
1128                data: json!({}),
1129            }),
1130            ThreadRequest::Read(params) => {
1131                let id = params.thread_id.clone();
1132                let data = self.persisted_thread_data(&id)?;
1133                Ok(ThreadResponse {
1134                    thread_id: id,
1135                    status: "ok".to_string(),
1136                    thread: self.thread_manager.read_thread(&params)?,
1137                    threads: Vec::new(),
1138                    goal: self.thread_manager.get_thread_goal(&ThreadGoalGetParams {
1139                        thread_id: params.thread_id,
1140                    })?,
1141                    model: None,
1142                    model_provider: None,
1143                    cwd: None,
1144                    approval_policy: None,
1145                    sandbox: None,
1146                    events: Vec::new(),
1147                    data,
1148                })
1149            }
1150            ThreadRequest::SetName(params) => Ok(ThreadResponse {
1151                thread_id: params.thread_id.clone(),
1152                status: "ok".to_string(),
1153                thread: self.thread_manager.set_thread_name(&params)?,
1154                threads: Vec::new(),
1155                goal: None,
1156                model: None,
1157                model_provider: None,
1158                cwd: None,
1159                approval_policy: None,
1160                sandbox: None,
1161                events: Vec::new(),
1162                data: json!({}),
1163            }),
1164            ThreadRequest::GoalSet(params) => {
1165                let thread_id = params.thread_id.clone();
1166                if let Some(goal) = self.thread_manager.set_thread_goal(&params)? {
1167                    Ok(ThreadResponse {
1168                        thread_id,
1169                        status: "ok".to_string(),
1170                        thread: None,
1171                        threads: Vec::new(),
1172                        goal: Some(goal.clone()),
1173                        model: None,
1174                        model_provider: None,
1175                        cwd: None,
1176                        approval_policy: None,
1177                        sandbox: None,
1178                        events: vec![EventFrame::ThreadGoalUpdated { goal: goal.clone() }],
1179                        data: json!({ "goal": goal }),
1180                    })
1181                } else {
1182                    Ok(ThreadResponse {
1183                        thread_id,
1184                        status: "missing".to_string(),
1185                        thread: None,
1186                        threads: Vec::new(),
1187                        goal: None,
1188                        model: None,
1189                        model_provider: None,
1190                        cwd: None,
1191                        approval_policy: None,
1192                        sandbox: None,
1193                        events: Vec::new(),
1194                        data: json!({"error":"thread not found"}),
1195                    })
1196                }
1197            }
1198            ThreadRequest::GoalGet(params) => {
1199                let goal = self.thread_manager.get_thread_goal(&params)?;
1200                Ok(ThreadResponse {
1201                    thread_id: params.thread_id,
1202                    status: "ok".to_string(),
1203                    thread: None,
1204                    threads: Vec::new(),
1205                    goal: goal.clone(),
1206                    model: None,
1207                    model_provider: None,
1208                    cwd: None,
1209                    approval_policy: None,
1210                    sandbox: None,
1211                    events: Vec::new(),
1212                    data: json!({ "goal": goal }),
1213                })
1214            }
1215            ThreadRequest::GoalClear(params) => {
1216                let thread_id = params.thread_id.clone();
1217                let cleared = self.thread_manager.clear_thread_goal(&params)?;
1218                Ok(ThreadResponse {
1219                    thread_id: thread_id.clone(),
1220                    status: if cleared { "cleared" } else { "empty" }.to_string(),
1221                    thread: None,
1222                    threads: Vec::new(),
1223                    goal: None,
1224                    model: None,
1225                    model_provider: None,
1226                    cwd: None,
1227                    approval_policy: None,
1228                    sandbox: None,
1229                    events: if cleared {
1230                        vec![EventFrame::ThreadGoalCleared { thread_id }]
1231                    } else {
1232                        Vec::new()
1233                    },
1234                    data: json!({ "cleared": cleared }),
1235                })
1236            }
1237            ThreadRequest::GoalRecordProgress(params) => {
1238                let thread_id = params.thread_id.clone();
1239                if let Some(goal) = self.thread_manager.record_thread_goal_progress(&params)? {
1240                    Ok(ThreadResponse {
1241                        thread_id,
1242                        status: "ok".to_string(),
1243                        thread: None,
1244                        threads: Vec::new(),
1245                        goal: Some(goal.clone()),
1246                        model: None,
1247                        model_provider: None,
1248                        cwd: None,
1249                        approval_policy: None,
1250                        sandbox: None,
1251                        events: vec![EventFrame::ThreadGoalUpdated { goal: goal.clone() }],
1252                        data: json!({ "goal": goal }),
1253                    })
1254                } else {
1255                    Ok(ThreadResponse {
1256                        thread_id,
1257                        status: "missing".to_string(),
1258                        thread: None,
1259                        threads: Vec::new(),
1260                        goal: None,
1261                        model: None,
1262                        model_provider: None,
1263                        cwd: None,
1264                        approval_policy: None,
1265                        sandbox: None,
1266                        events: Vec::new(),
1267                        data: json!({"error":"thread or goal not found"}),
1268                    })
1269                }
1270            }
1271            ThreadRequest::Archive { thread_id } => {
1272                self.thread_manager.archive_thread(&thread_id)?;
1273                Ok(ThreadResponse {
1274                    thread_id,
1275                    status: "archived".to_string(),
1276                    thread: None,
1277                    threads: Vec::new(),
1278                    goal: None,
1279                    model: None,
1280                    model_provider: None,
1281                    cwd: None,
1282                    approval_policy: None,
1283                    sandbox: None,
1284                    events: Vec::new(),
1285                    data: json!({}),
1286                })
1287            }
1288            ThreadRequest::Unarchive { thread_id } => {
1289                self.thread_manager.unarchive_thread(&thread_id)?;
1290                Ok(ThreadResponse {
1291                    thread_id,
1292                    status: "unarchived".to_string(),
1293                    thread: None,
1294                    threads: Vec::new(),
1295                    goal: None,
1296                    model: None,
1297                    model_provider: None,
1298                    cwd: None,
1299                    approval_policy: None,
1300                    sandbox: None,
1301                    events: Vec::new(),
1302                    data: json!({}),
1303                })
1304            }
1305            ThreadRequest::Message { thread_id, input } => {
1306                self.thread_manager.touch_message(&thread_id, &input)?;
1307                // Keyed by a fresh uuid, like handle_prompt: keying on
1308                // `{thread_id}:{input.len()}` made any two equal-length
1309                // messages share a response_id, breaking hook correlation.
1310                let response_id = format!("resp-{}", Uuid::new_v4());
1311                self.hooks
1312                    .emit(HookEvent::ResponseStart {
1313                        response_id: response_id.clone(),
1314                    })
1315                    .await;
1316                self.hooks
1317                    .emit(HookEvent::ResponseEnd {
1318                        response_id: response_id.clone(),
1319                    })
1320                    .await;
1321
1322                Ok(ThreadResponse {
1323                    thread_id,
1324                    status: "accepted".to_string(),
1325                    thread: None,
1326                    threads: Vec::new(),
1327                    goal: None,
1328                    model: None,
1329                    model_provider: None,
1330                    cwd: None,
1331                    approval_policy: None,
1332                    sandbox: None,
1333                    events: vec![
1334                        EventFrame::ResponseStart {
1335                            response_id: response_id.clone(),
1336                        },
1337                        EventFrame::ResponseDelta {
1338                            response_id: response_id.clone(),
1339                            delta: "queued".to_string(),
1340                            channel: ResponseChannel::Text,
1341                        },
1342                        EventFrame::ResponseEnd { response_id },
1343                    ],
1344                    data: json!({}),
1345                })
1346            }
1347        }
1348    }
1349
1350    /// Resolves the model for a prompt, records the message, and returns the response.
1351    pub async fn handle_prompt(
1352        &mut self,
1353        req: PromptRequest,
1354        cli_overrides: &CliRuntimeOverrides,
1355    ) -> Result<PromptResponse> {
1356        let resolved = self.config.resolve_runtime_options(cli_overrides);
1357        let requested_model = req.model.clone().unwrap_or_else(|| resolved.model.clone());
1358        let selection = self
1359            .model_registry
1360            .resolve(Some(&requested_model), Some(resolved.provider));
1361        let resolved_model = selection.resolved.id.clone();
1362        let response_id = format!("resp-{}", Uuid::new_v4());
1363
1364        self.hooks
1365            .emit(HookEvent::ResponseStart {
1366                response_id: response_id.clone(),
1367            })
1368            .await;
1369        self.hooks
1370            .emit(HookEvent::ResponseDelta {
1371                response_id: response_id.clone(),
1372                delta: "model-selected".to_string(),
1373            })
1374            .await;
1375        self.hooks
1376            .emit(HookEvent::ResponseEnd {
1377                response_id: response_id.clone(),
1378            })
1379            .await;
1380
1381        let payload = json!({
1382            "provider": resolved.provider.as_str(),
1383            "model": resolved_model.clone(),
1384            "prompt": req.prompt,
1385            "telemetry": resolved.telemetry,
1386            "base_url": resolved.base_url,
1387            "has_api_key": resolved.api_key.as_ref().is_some_and(|k| !k.trim().is_empty()),
1388            "approval_policy": resolved.approval_policy,
1389            "sandbox_mode": resolved.sandbox_mode
1390        });
1391        if let Some(thread_id) = req.thread_id.as_ref() {
1392            self.thread_manager.touch_message(thread_id, &req.prompt)?;
1393            let assistant_message_id = self.thread_manager.store.append_message(
1394                thread_id,
1395                "assistant",
1396                &payload.to_string(),
1397                Some(payload.clone()),
1398            )?;
1399            self.persist_latest_checkpoint(
1400                thread_id,
1401                "prompt_response",
1402                json!({
1403                    "response_id": response_id.clone(),
1404                    "model": resolved_model.clone(),
1405                    "provider": resolved.provider.as_str(),
1406                    "assistant_message_id": assistant_message_id
1407                }),
1408            )?;
1409        }
1410
1411        Ok(PromptResponse {
1412            output: payload.to_string(),
1413            model: resolved_model,
1414            events: vec![
1415                EventFrame::ResponseStart {
1416                    response_id: response_id.clone(),
1417                },
1418                EventFrame::ResponseDelta {
1419                    response_id: response_id.clone(),
1420                    delta: "model-selected".to_string(),
1421                    channel: ResponseChannel::Text,
1422                },
1423                EventFrame::ResponseEnd { response_id },
1424            ],
1425        })
1426    }
1427
1428    /// Evaluates execution policy and dispatches a tool call.
1429    pub async fn invoke_tool(
1430        &self,
1431        call: ToolCall,
1432        approval_mode: AskForApproval,
1433        cwd: &Path,
1434    ) -> Result<Value> {
1435        let fallback_cwd = cwd.display().to_string();
1436        let (command, policy_cwd, execution_kind) = call.execution_subject(&fallback_cwd);
1437        let policy_tool = match &call.payload {
1438            ToolPayload::LocalShell { .. } => "exec_shell",
1439            _ => call.name.as_str(),
1440        };
1441        let policy_path = permission_path_for_call(&call);
1442        let decision = self.exec_policy.check(ExecPolicyContext {
1443            command: &command,
1444            cwd: &policy_cwd,
1445            tool: Some(policy_tool),
1446            path: policy_path.as_deref(),
1447            ask_for_approval: approval_mode,
1448            sandbox_mode: None,
1449        })?;
1450        let precheck = policy_precheck_payload(&decision, &command, &policy_cwd, execution_kind);
1451        let response_id = format!("tool-{}", Uuid::new_v4());
1452        let call_id = call
1453            .raw_tool_call_id
1454            .clone()
1455            .unwrap_or_else(|| format!("tool-call-{}", Uuid::new_v4()));
1456        self.hooks
1457            .emit(HookEvent::ToolLifecycle {
1458                response_id: response_id.clone(),
1459                tool_name: call.name.clone(),
1460                phase: "precheck".to_string(),
1461                payload: precheck.clone(),
1462            })
1463            .await;
1464
1465        if !decision.allow {
1466            let reason = decision.reason().to_string();
1467            let approval_id = format!("approval-{}", Uuid::new_v4());
1468            let error_frame = EventFrame::Error {
1469                response_id: response_id.clone(),
1470                message: reason.clone(),
1471            };
1472            self.hooks
1473                .emit(HookEvent::ApprovalLifecycle {
1474                    approval_id,
1475                    phase: "denied".to_string(),
1476                    reason: Some(reason.clone()),
1477                })
1478                .await;
1479            self.hooks
1480                .emit(HookEvent::GenericEventFrame {
1481                    frame: Box::new(error_frame.clone()),
1482                })
1483                .await;
1484            return Ok(json!({
1485                "ok": false,
1486                "status": "denied",
1487                "execution_kind": execution_kind,
1488                "response_id": response_id,
1489                "precheck": precheck,
1490                "error": reason,
1491                "events": [event_frame_payload(&error_frame)],
1492            }));
1493        }
1494
1495        if decision.requires_approval {
1496            let approval_id = format!("approval-{}", Uuid::new_v4());
1497            let reason = decision.reason().to_string();
1498            let maybe_approval_frame = approval_request_frame(
1499                &decision.requirement,
1500                decision.matched_rule.as_deref(),
1501                call_id,
1502                approval_id.clone(),
1503                response_id.clone(),
1504                command.clone(),
1505                policy_cwd.clone(),
1506            );
1507            self.hooks
1508                .emit(HookEvent::ApprovalLifecycle {
1509                    approval_id: approval_id.clone(),
1510                    phase: "requested".to_string(),
1511                    reason: Some(reason.clone()),
1512                })
1513                .await;
1514            let mut events = Vec::new();
1515            if let Some(frame) = maybe_approval_frame {
1516                self.hooks
1517                    .emit(HookEvent::GenericEventFrame {
1518                        frame: Box::new(frame.clone()),
1519                    })
1520                    .await;
1521                events.push(event_frame_payload(&frame));
1522            }
1523            return Ok(json!({
1524                "ok": false,
1525                "status": "approval_required",
1526                "execution_kind": execution_kind,
1527                "response_id": response_id,
1528                "approval_id": approval_id,
1529                "precheck": precheck,
1530                "error": reason,
1531                "events": events,
1532            }));
1533        }
1534
1535        // Headless `request_user_input`: mirror the approval fire-and-return
1536        // branch (issue #3102). The TUI intercepts this tool by name before
1537        // dispatch and blocks on a reply channel; the headless runtime instead
1538        // emits a typed `UserInputRequest` frame and returns a
1539        // `user_input_required` status so the client can render the question
1540        // and POST answers back via `AppRequest::SubmitUserInput`. It does NOT
1541        // block — consistent with the headless approval model, which has no
1542        // resume channel either.
1543        if call.name == REQUEST_USER_INPUT_TOOL_NAME {
1544            let request_id = format!("user-input-{}", Uuid::new_v4());
1545            let arguments = match &call.payload {
1546                ToolPayload::Function { arguments } => arguments.as_str(),
1547                // Custom/Mcp/LocalShell can't carry a user_input payload; fall
1548                // through to the generic dispatch error below.
1549                _ => "",
1550            };
1551            let maybe_frame = user_input_request_frame(
1552                call_id.clone(),
1553                response_id.clone(),
1554                request_id.clone(),
1555                arguments,
1556            );
1557            let mut events = Vec::new();
1558            if let Some(frame) = maybe_frame {
1559                self.hooks
1560                    .emit(HookEvent::GenericEventFrame {
1561                        frame: Box::new(frame.clone()),
1562                    })
1563                    .await;
1564                events.push(event_frame_payload(&frame));
1565            }
1566            return Ok(json!({
1567                "ok": false,
1568                "status": "user_input_required",
1569                "execution_kind": execution_kind,
1570                "response_id": response_id,
1571                "request_id": request_id,
1572                "precheck": precheck,
1573                "events": events,
1574            }));
1575        }
1576
1577        let start_frame = EventFrame::ToolCallStart {
1578            response_id: response_id.clone(),
1579            tool_name: call.name.clone(),
1580            arguments: tool_payload_value(&call.payload),
1581        };
1582        self.hooks
1583            .emit(HookEvent::GenericEventFrame {
1584                frame: Box::new(start_frame.clone()),
1585            })
1586            .await;
1587        self.hooks
1588            .emit(HookEvent::ToolLifecycle {
1589                response_id: response_id.clone(),
1590                tool_name: call.name.clone(),
1591                phase: "dispatching".to_string(),
1592                payload: json!({
1593                    "call_id": call_id,
1594                    "execution_kind": execution_kind
1595                }),
1596            })
1597            .await;
1598
1599        match time::timeout(
1600            tool_dispatch_timeout(),
1601            self.tool_registry.dispatch(call.clone(), true),
1602        )
1603        .await
1604        {
1605            Ok(Ok(tool_output)) => {
1606                let success = tool_output.success();
1607                let status = if success { "completed" } else { "failed" };
1608                let result_frame = EventFrame::ToolCallResult {
1609                    response_id: response_id.clone(),
1610                    tool_name: call.name.clone(),
1611                    output: tool_output_value(&tool_output),
1612                };
1613                self.hooks
1614                    .emit(HookEvent::GenericEventFrame {
1615                        frame: Box::new(result_frame.clone()),
1616                    })
1617                    .await;
1618                self.hooks
1619                    .emit(HookEvent::ToolLifecycle {
1620                        response_id: response_id.clone(),
1621                        tool_name: call.name,
1622                        phase: status.to_string(),
1623                        payload: json!({ "ok": success }),
1624                    })
1625                    .await;
1626                Ok(json!({
1627                    "ok": success,
1628                    "status": status,
1629                    "execution_kind": execution_kind,
1630                    "response_id": response_id,
1631                    "precheck": precheck,
1632                    "output": tool_output,
1633                    "events": [
1634                        event_frame_payload(&start_frame),
1635                        event_frame_payload(&result_frame)
1636                    ]
1637                }))
1638            }
1639            Ok(Err(err)) => {
1640                let message = format!("{err:?}");
1641                let error_frame = EventFrame::Error {
1642                    response_id: response_id.clone(),
1643                    message: message.clone(),
1644                };
1645                self.hooks
1646                    .emit(HookEvent::GenericEventFrame {
1647                        frame: Box::new(error_frame.clone()),
1648                    })
1649                    .await;
1650                self.hooks
1651                    .emit(HookEvent::ToolLifecycle {
1652                        response_id: response_id.clone(),
1653                        tool_name: call.name,
1654                        phase: "failed".to_string(),
1655                        payload: json!({ "error": message.clone() }),
1656                    })
1657                    .await;
1658                Ok(json!({
1659                    "ok": false,
1660                    "status": "failed",
1661                    "execution_kind": execution_kind,
1662                    "response_id": response_id,
1663                    "precheck": precheck,
1664                    "error": message,
1665                    "events": [
1666                        event_frame_payload(&start_frame),
1667                        event_frame_payload(&error_frame)
1668                    ]
1669                }))
1670            }
1671            Err(_elapsed) => {
1672                let seconds = tool_dispatch_timeout().as_secs().max(1);
1673                let message = format!("Tool '{}' timed out after {seconds}s", call.name);
1674                let error_frame = EventFrame::Error {
1675                    response_id: response_id.clone(),
1676                    message: message.clone(),
1677                };
1678                self.hooks
1679                    .emit(HookEvent::GenericEventFrame {
1680                        frame: Box::new(error_frame.clone()),
1681                    })
1682                    .await;
1683                self.hooks
1684                    .emit(HookEvent::ToolLifecycle {
1685                        response_id: response_id.clone(),
1686                        tool_name: call.name,
1687                        phase: "failed".to_string(),
1688                        payload: json!({ "error": message.clone(), "timeout": true }),
1689                    })
1690                    .await;
1691                Ok(json!({
1692                    "ok": false,
1693                    "status": "timeout",
1694                    "execution_kind": execution_kind,
1695                    "response_id": response_id,
1696                    "precheck": precheck,
1697                    "error": message,
1698                    "events": [
1699                        event_frame_payload(&start_frame),
1700                        event_frame_payload(&error_frame)
1701                    ]
1702                }))
1703            }
1704        }
1705    }
1706
1707    /// Starts all configured MCP servers and emits startup events via hooks.
1708    pub async fn mcp_startup(&self) -> McpStartupCompleteEvent {
1709        let mut updates = Vec::new();
1710        let summary = self.mcp_manager.start_all(|update| {
1711            updates.push(update);
1712        });
1713        for update in updates {
1714            let status = match update.status {
1715                McpManagerStartupStatus::Starting => codewhale_protocol::McpStartupStatus::Starting,
1716                McpManagerStartupStatus::Ready => codewhale_protocol::McpStartupStatus::Ready,
1717                McpManagerStartupStatus::Failed { error } => {
1718                    codewhale_protocol::McpStartupStatus::Failed { error }
1719                }
1720                McpManagerStartupStatus::Cancelled => {
1721                    codewhale_protocol::McpStartupStatus::Cancelled
1722                }
1723            };
1724            self.hooks
1725                .emit(HookEvent::GenericEventFrame {
1726                    frame: Box::new(EventFrame::McpStartupUpdate {
1727                        update: codewhale_protocol::McpStartupUpdateEvent {
1728                            server_name: update.server_name,
1729                            status,
1730                        },
1731                    }),
1732                })
1733                .await;
1734        }
1735        self.hooks
1736            .emit(HookEvent::GenericEventFrame {
1737                frame: Box::new(EventFrame::McpStartupComplete {
1738                    summary: codewhale_protocol::McpStartupCompleteEvent {
1739                        ready: summary.ready.clone(),
1740                        failed: summary
1741                            .failed
1742                            .iter()
1743                            .map(|f| codewhale_protocol::McpStartupFailure {
1744                                server_name: f.server_name.clone(),
1745                                error: f.error.clone(),
1746                            })
1747                            .collect(),
1748                        cancelled: summary.cancelled.clone(),
1749                    },
1750                }),
1751            })
1752            .await;
1753        summary
1754    }
1755
1756    /// Returns the current application status including all jobs and their history.
1757    pub fn app_status(&self) -> AppResponse {
1758        let jobs = self.jobs.list();
1759        let events = jobs
1760            .iter()
1761            .flat_map(|job| {
1762                job.history.iter().map(|entry| EventFrame::ResponseDelta {
1763                    response_id: job.id.clone(),
1764                    delta: json!({
1765                        "kind": "job_transition",
1766                        "job_id": job.id.clone(),
1767                        "phase": entry.phase.clone(),
1768                        "status": job_status_to_str(entry.status),
1769                        "progress": entry.progress,
1770                        "detail": entry.detail.clone(),
1771                        "retry": job_retry_to_value(&entry.retry),
1772                        "at": entry.at
1773                    })
1774                    .to_string(),
1775                    channel: ResponseChannel::Text,
1776                })
1777            })
1778            .collect::<Vec<_>>();
1779        AppResponse {
1780            ok: true,
1781            data: json!({
1782                "jobs": jobs.into_iter().map(|job| {
1783                    json!({
1784                        "id": job.id,
1785                        "name": job.name,
1786                        "status": job_status_to_str(job.status),
1787                        "progress": job.progress,
1788                        "detail": job.detail,
1789                        "retry": job_retry_to_value(&job.retry),
1790                        "history": job.history.iter().map(job_history_to_value).collect::<Vec<_>>()
1791                    })
1792                }).collect::<Vec<_>>()
1793            }),
1794            events,
1795        }
1796    }
1797
1798    /// Returns the default model provider from the resolved configuration.
1799    pub fn provider_default(&self) -> ProviderKind {
1800        self.config.provider
1801    }
1802
1803    /// Saves a named checkpoint for a thread.
1804    pub fn save_thread_checkpoint(
1805        &self,
1806        thread_id: &str,
1807        checkpoint_id: &str,
1808        state: &Value,
1809    ) -> Result<()> {
1810        self.thread_manager
1811            .state_store()
1812            .save_checkpoint(thread_id, checkpoint_id, state)
1813    }
1814
1815    /// Loads a checkpoint for a thread. Pass `None` for the latest.
1816    pub fn load_thread_checkpoint(
1817        &self,
1818        thread_id: &str,
1819        checkpoint_id: Option<&str>,
1820    ) -> Result<Option<Value>> {
1821        Ok(self
1822            .thread_manager
1823            .state_store()
1824            .load_checkpoint(thread_id, checkpoint_id)?
1825            .map(|checkpoint| checkpoint.state))
1826    }
1827
1828    /// Enqueues a new background job and persists it immediately.
1829    pub fn enqueue_job(&mut self, name: impl Into<String>) -> Result<JobRecord> {
1830        let job = self.jobs.enqueue(name);
1831        self.jobs
1832            .persist_job(self.thread_manager.state_store(), &job.id)?;
1833        Ok(job)
1834    }
1835
1836    /// Transitions a job to running and persists the change.
1837    pub fn set_job_running(&mut self, job_id: &str) -> Result<()> {
1838        self.jobs.set_running(job_id);
1839        self.jobs
1840            .persist_job(self.thread_manager.state_store(), job_id)
1841    }
1842
1843    /// Updates a job's progress and persists the change.
1844    pub fn update_job_progress(
1845        &mut self,
1846        job_id: &str,
1847        progress: u8,
1848        detail: Option<String>,
1849    ) -> Result<()> {
1850        self.jobs.update_progress(job_id, progress, detail);
1851        self.jobs
1852            .persist_job(self.thread_manager.state_store(), job_id)
1853    }
1854
1855    /// Marks a job as completed and persists the change.
1856    pub fn complete_job(&mut self, job_id: &str) -> Result<()> {
1857        self.jobs.complete(job_id);
1858        self.jobs
1859            .persist_job(self.thread_manager.state_store(), job_id)
1860    }
1861
1862    /// Marks a job as failed and persists the change.
1863    pub fn fail_job(&mut self, job_id: &str, detail: impl Into<String>) -> Result<()> {
1864        self.jobs.fail(job_id, detail);
1865        self.jobs
1866            .persist_job(self.thread_manager.state_store(), job_id)
1867    }
1868
1869    /// Cancels a job and persists the change.
1870    pub fn cancel_job(&mut self, job_id: &str) -> Result<()> {
1871        self.jobs.cancel(job_id);
1872        self.jobs
1873            .persist_job(self.thread_manager.state_store(), job_id)
1874    }
1875
1876    /// Pauses a job and persists the change.
1877    pub fn pause_job(&mut self, job_id: &str, detail: Option<String>) -> Result<()> {
1878        self.jobs.pause(job_id, detail);
1879        self.jobs
1880            .persist_job(self.thread_manager.state_store(), job_id)
1881    }
1882
1883    /// Resumes a paused job and persists the change.
1884    pub fn resume_job(&mut self, job_id: &str, detail: Option<String>) -> Result<()> {
1885        self.jobs.resume(job_id, detail);
1886        self.jobs
1887            .persist_job(self.thread_manager.state_store(), job_id)
1888    }
1889
1890    /// Returns the state-transition history for a job.
1891    pub fn job_history(&self, job_id: &str) -> Vec<JobHistoryEntry> {
1892        self.jobs.history(job_id)
1893    }
1894}
1895
1896fn thread_response_from_new(status: &str, new: NewThread) -> ThreadResponse {
1897    ThreadResponse {
1898        thread_id: new.thread.id.clone(),
1899        status: status.to_string(),
1900        thread: Some(new.thread),
1901        threads: Vec::new(),
1902        goal: None,
1903        model: Some(new.model),
1904        model_provider: Some(new.model_provider),
1905        cwd: Some(new.cwd),
1906        approval_policy: new.approval_policy,
1907        sandbox: new.sandbox,
1908        events: Vec::new(),
1909        data: json!({}),
1910    }
1911}
1912
1913fn preview_from_initial_history(initial_history: &InitialHistory) -> String {
1914    match initial_history {
1915        InitialHistory::New => "New conversation".to_string(),
1916        InitialHistory::Forked(items) => truncate_preview(
1917            &items
1918                .first()
1919                .map(Value::to_string)
1920                .unwrap_or_else(|| "Forked conversation".to_string()),
1921        ),
1922        InitialHistory::Resumed { history, .. } => truncate_preview(
1923            &history
1924                .first()
1925                .map(Value::to_string)
1926                .unwrap_or_else(|| "Resumed conversation".to_string()),
1927        ),
1928    }
1929}
1930
1931fn permission_path_for_call(call: &ToolCall) -> Option<String> {
1932    match &call.payload {
1933        ToolPayload::Function { arguments } => serde_json::from_str::<Value>(arguments)
1934            .ok()
1935            .and_then(|value| {
1936                value
1937                    .get("path")
1938                    .and_then(Value::as_str)
1939                    .map(str::to_string)
1940            }),
1941        ToolPayload::Mcp { raw_arguments, .. } => raw_arguments
1942            .get("path")
1943            .and_then(Value::as_str)
1944            .map(str::to_string),
1945        ToolPayload::Custom { .. } | ToolPayload::LocalShell { .. } => None,
1946    }
1947}
1948
1949fn truncate_preview(value: &str) -> String {
1950    value.chars().take(120).collect()
1951}
1952
1953fn to_protocol_thread(thread: ThreadMetadata) -> Thread {
1954    Thread {
1955        id: thread.id,
1956        preview: thread.preview,
1957        ephemeral: thread.ephemeral,
1958        model_provider: thread.model_provider,
1959        created_at: thread.created_at,
1960        updated_at: thread.updated_at,
1961        status: match thread.status {
1962            PersistedThreadStatus::Running => ThreadStatus::Running,
1963            PersistedThreadStatus::Idle => ThreadStatus::Idle,
1964            PersistedThreadStatus::Completed => ThreadStatus::Completed,
1965            PersistedThreadStatus::Failed => ThreadStatus::Failed,
1966            PersistedThreadStatus::Paused => ThreadStatus::Paused,
1967            PersistedThreadStatus::Archived => ThreadStatus::Archived,
1968        },
1969        path: thread.path,
1970        cwd: thread.cwd,
1971        cli_version: thread.cli_version,
1972        source: match thread.source {
1973            SessionSource::Interactive => codewhale_protocol::SessionSource::Interactive,
1974            SessionSource::Resume => codewhale_protocol::SessionSource::Resume,
1975            SessionSource::Fork => codewhale_protocol::SessionSource::Fork,
1976            SessionSource::Api => codewhale_protocol::SessionSource::Api,
1977            SessionSource::Unknown => codewhale_protocol::SessionSource::Unknown,
1978        },
1979        name: thread.name,
1980    }
1981}
1982
1983fn to_protocol_goal(goal: ThreadGoalRecord) -> ThreadGoal {
1984    ThreadGoal {
1985        thread_id: goal.thread_id,
1986        goal_id: goal.goal_id,
1987        objective: goal.objective,
1988        status: to_protocol_goal_status(goal.status),
1989        token_budget: goal.token_budget,
1990        tokens_used: goal.tokens_used,
1991        time_used_seconds: goal.time_used_seconds,
1992        continuation_count: goal.continuation_count,
1993        created_at: goal.created_at,
1994        updated_at: goal.updated_at,
1995    }
1996}
1997
1998fn to_protocol_goal_status(status: PersistedThreadGoalStatus) -> ThreadGoalStatus {
1999    match status {
2000        PersistedThreadGoalStatus::Active => ThreadGoalStatus::Active,
2001        PersistedThreadGoalStatus::Paused => ThreadGoalStatus::Paused,
2002        PersistedThreadGoalStatus::Blocked => ThreadGoalStatus::Blocked,
2003        PersistedThreadGoalStatus::UsageLimited => ThreadGoalStatus::UsageLimited,
2004        PersistedThreadGoalStatus::BudgetLimited => ThreadGoalStatus::BudgetLimited,
2005        PersistedThreadGoalStatus::Complete => ThreadGoalStatus::Complete,
2006    }
2007}
2008
2009fn to_persisted_status(status: &ThreadStatus) -> PersistedThreadStatus {
2010    match status {
2011        ThreadStatus::Running => PersistedThreadStatus::Running,
2012        ThreadStatus::Idle => PersistedThreadStatus::Idle,
2013        ThreadStatus::Completed => PersistedThreadStatus::Completed,
2014        ThreadStatus::Failed => PersistedThreadStatus::Failed,
2015        ThreadStatus::Paused => PersistedThreadStatus::Paused,
2016        ThreadStatus::Archived => PersistedThreadStatus::Archived,
2017    }
2018}
2019
2020fn to_persisted_source(source: &codewhale_protocol::SessionSource) -> SessionSource {
2021    match source {
2022        codewhale_protocol::SessionSource::Interactive => SessionSource::Interactive,
2023        codewhale_protocol::SessionSource::Resume => SessionSource::Resume,
2024        codewhale_protocol::SessionSource::Fork => SessionSource::Fork,
2025        codewhale_protocol::SessionSource::Api => SessionSource::Api,
2026        codewhale_protocol::SessionSource::Unknown => SessionSource::Unknown,
2027    }
2028}
2029
2030fn approval_request_frame(
2031    requirement: &ExecApprovalRequirement,
2032    matched_rule: Option<&str>,
2033    call_id: String,
2034    approval_id: String,
2035    turn_id: String,
2036    command: String,
2037    cwd: String,
2038) -> Option<EventFrame> {
2039    let ExecApprovalRequirement::NeedsApproval {
2040        reason,
2041        proposed_execpolicy_amendment,
2042        proposed_network_policy_amendments,
2043    } = requirement
2044    else {
2045        return None;
2046    };
2047
2048    let mut available_decisions = vec![
2049        ReviewDecision::Approved,
2050        ReviewDecision::ApprovedForSession,
2051        ReviewDecision::Denied,
2052        ReviewDecision::Abort,
2053    ];
2054    if proposed_execpolicy_amendment
2055        .as_ref()
2056        .is_some_and(|amendment| !amendment.prefixes.is_empty())
2057    {
2058        available_decisions.push(ReviewDecision::ApprovedExecpolicyAmendment);
2059    }
2060    available_decisions.extend(proposed_network_policy_amendments.iter().cloned().map(
2061        |amendment| ReviewDecision::NetworkPolicyAmendment {
2062            host: amendment.host,
2063            action: amendment.action,
2064        },
2065    ));
2066
2067    Some(EventFrame::ExecApprovalRequest {
2068        request: ExecApprovalRequestEvent {
2069            call_id,
2070            approval_id,
2071            turn_id,
2072            command,
2073            cwd,
2074            reason: reason.clone(),
2075            matched_rule: matched_rule.map(|rule| rule.to_string().into_boxed_str()),
2076            network_approval_context: None,
2077            proposed_execpolicy_amendment: proposed_execpolicy_amendment
2078                .as_ref()
2079                .map(|amendment| amendment.prefixes.clone())
2080                .unwrap_or_default(),
2081            proposed_network_policy_amendments: proposed_network_policy_amendments.clone(),
2082            additional_permissions: Vec::new(),
2083            available_decisions,
2084        },
2085    })
2086}
2087
2088/// Build an [`EventFrame::UserInputRequest`] for a headless
2089/// `request_user_input` tool call, mirroring [`approval_request_frame`].
2090///
2091/// `arguments` is the raw JSON arguments string the model supplied to the
2092/// `request_user_input` tool (a `ToolPayload::Function` body). On parse
2093/// failure we return `None` so the caller falls through to the generic tool
2094/// error path rather than silently dropping the request.
2095fn user_input_request_frame(
2096    call_id: String,
2097    turn_id: String,
2098    request_id: String,
2099    arguments: &str,
2100) -> Option<EventFrame> {
2101    let parsed: Value = serde_json::from_str(arguments).ok()?;
2102    // Extract the `questions` array and lift it into the headless event
2103    // shape. We tolerate missing `allow_free_text`/`multi_select` (default
2104    // false) and extra fields, matching the lenient TUI `from_value` path.
2105    let questions = parsed.get("questions").cloned().filter(Value::is_array)?;
2106    let request = UserInputRequestEvent {
2107        call_id,
2108        turn_id,
2109        request_id,
2110        questions: serde_json::from_value(questions).ok()?,
2111    };
2112    Some(EventFrame::UserInputRequest { request })
2113}
2114
2115fn approval_requirement_payload(requirement: &ExecApprovalRequirement) -> Value {
2116    match requirement {
2117        ExecApprovalRequirement::Skip {
2118            bypass_sandbox,
2119            proposed_execpolicy_amendment,
2120        } => json!({
2121            "type": "skip",
2122            "bypass_sandbox": bypass_sandbox,
2123            "reason": requirement.reason(),
2124            "proposed_execpolicy_amendment": proposed_execpolicy_amendment
2125                .as_ref()
2126                .map(|amendment| amendment.prefixes.clone())
2127                .unwrap_or_default()
2128        }),
2129        ExecApprovalRequirement::NeedsApproval {
2130            reason,
2131            proposed_execpolicy_amendment,
2132            proposed_network_policy_amendments,
2133        } => json!({
2134            "type": "needs_approval",
2135            "reason": reason,
2136            "proposed_execpolicy_amendment": proposed_execpolicy_amendment
2137                .as_ref()
2138                .map(|amendment| amendment.prefixes.clone())
2139                .unwrap_or_default(),
2140            "proposed_network_policy_amendments": proposed_network_policy_amendments
2141        }),
2142        ExecApprovalRequirement::Forbidden { reason } => json!({
2143            "type": "forbidden",
2144            "reason": reason
2145        }),
2146    }
2147}
2148
2149fn policy_precheck_payload(
2150    decision: &ExecPolicyDecision,
2151    command: &str,
2152    cwd: &str,
2153    execution_kind: &str,
2154) -> Value {
2155    json!({
2156        "execution_kind": execution_kind,
2157        "command": command,
2158        "cwd": cwd,
2159        "allow": decision.allow,
2160        "requires_approval": decision.requires_approval,
2161        "matched_rule": decision.matched_rule.clone(),
2162        "phase": decision.requirement.phase(),
2163        "reason": decision.reason(),
2164        "requirement": approval_requirement_payload(&decision.requirement)
2165    })
2166}
2167
2168fn tool_payload_value(payload: &ToolPayload) -> Value {
2169    serde_json::to_value(payload).unwrap_or_else(
2170        |_| json!({"type":"serialization_error","message":"tool payload unavailable"}),
2171    )
2172}
2173
2174fn tool_output_value(output: &codewhale_protocol::ToolOutput) -> Value {
2175    serde_json::to_value(output).unwrap_or_else(
2176        |_| json!({"type":"serialization_error","message":"tool output unavailable"}),
2177    )
2178}
2179
2180fn event_frame_payload(frame: &EventFrame) -> Value {
2181    serde_json::to_value(frame)
2182        .unwrap_or_else(|_| json!({"event":"error","message":"failed to encode event frame"}))
2183}
2184
2185/// Tool name that triggers the headless clarification-question flow.
2186///
2187/// Mirrors the TUI's `REQUEST_USER_INPUT_NAME`
2188/// (`crates/tui/src/core/engine/tool_catalog.rs`); duplicated here rather than
2189/// depended on across crates so `core` stays free of `tui` imports.
2190const REQUEST_USER_INPUT_TOOL_NAME: &str = "request_user_input";
2191
2192fn json_optional_string(value: &Value) -> Option<String> {
2193    if value.is_null() {
2194        None
2195    } else {
2196        value.as_str().map(ToString::to_string)
2197    }
2198}
2199
2200fn parse_retry_metadata(value: Option<&Value>) -> JobRetryMetadata {
2201    let Some(value) = value else {
2202        return JobRetryMetadata::default();
2203    };
2204    JobRetryMetadata {
2205        attempt: value
2206            .get("attempt")
2207            .and_then(Value::as_u64)
2208            .unwrap_or(0)
2209            .min(u32::MAX as u64) as u32,
2210        max_attempts: value
2211            .get("max_attempts")
2212            .and_then(Value::as_u64)
2213            .unwrap_or(DEFAULT_JOB_MAX_ATTEMPTS as u64)
2214            .min(u32::MAX as u64) as u32,
2215        backoff_base_ms: value
2216            .get("backoff_base_ms")
2217            .and_then(Value::as_u64)
2218            .unwrap_or(DEFAULT_JOB_BACKOFF_BASE_MS),
2219        next_backoff_ms: value
2220            .get("next_backoff_ms")
2221            .and_then(Value::as_u64)
2222            .unwrap_or(0),
2223        next_retry_at: value.get("next_retry_at").and_then(Value::as_i64),
2224    }
2225}
2226
2227fn parse_history_entry(value: &Value) -> Option<JobHistoryEntry> {
2228    let status = value
2229        .get("status")
2230        .and_then(Value::as_str)
2231        .and_then(job_status_from_str)?;
2232    Some(JobHistoryEntry {
2233        at: value.get("at").and_then(Value::as_i64).unwrap_or(0),
2234        phase: value
2235            .get("phase")
2236            .and_then(Value::as_str)
2237            .unwrap_or("unknown")
2238            .to_string(),
2239        status,
2240        progress: value
2241            .get("progress")
2242            .and_then(Value::as_u64)
2243            .map(|v| v.min(u8::MAX as u64) as u8),
2244        detail: value.get("detail").and_then(json_optional_string),
2245        retry: parse_retry_metadata(value.get("retry")),
2246    })
2247}
2248
2249fn job_status_to_str(status: JobStatus) -> &'static str {
2250    match status {
2251        JobStatus::Queued => "queued",
2252        JobStatus::Running => "running",
2253        JobStatus::Paused => "paused",
2254        JobStatus::Completed => "completed",
2255        JobStatus::Failed => "failed",
2256        JobStatus::Cancelled => "cancelled",
2257    }
2258}
2259
2260fn job_status_from_str(value: &str) -> Option<JobStatus> {
2261    match value {
2262        "queued" => Some(JobStatus::Queued),
2263        "running" => Some(JobStatus::Running),
2264        "paused" => Some(JobStatus::Paused),
2265        "completed" => Some(JobStatus::Completed),
2266        "failed" => Some(JobStatus::Failed),
2267        "cancelled" => Some(JobStatus::Cancelled),
2268        _ => None,
2269    }
2270}
2271
2272fn job_retry_to_value(retry: &JobRetryMetadata) -> Value {
2273    json!({
2274        "attempt": retry.attempt,
2275        "max_attempts": retry.max_attempts,
2276        "backoff_base_ms": retry.backoff_base_ms,
2277        "next_backoff_ms": retry.next_backoff_ms,
2278        "next_retry_at": retry.next_retry_at
2279    })
2280}
2281
2282fn job_history_to_value(entry: &JobHistoryEntry) -> Value {
2283    json!({
2284        "at": entry.at,
2285        "phase": entry.phase.clone(),
2286        "status": job_status_to_str(entry.status),
2287        "progress": entry.progress,
2288        "detail": entry.detail.clone(),
2289        "retry": job_retry_to_value(&entry.retry)
2290    })
2291}
2292
2293fn runtime_status_to_job_state(status: JobStatus) -> JobStateStatus {
2294    match status {
2295        JobStatus::Queued => JobStateStatus::Queued,
2296        JobStatus::Running => JobStateStatus::Running,
2297        JobStatus::Paused => JobStateStatus::Paused,
2298        JobStatus::Completed => JobStateStatus::Completed,
2299        JobStatus::Failed => JobStateStatus::Failed,
2300        JobStatus::Cancelled => JobStateStatus::Cancelled,
2301    }
2302}
2303
2304fn job_state_status_to_runtime(status: JobStateStatus) -> JobStatus {
2305    match status {
2306        JobStateStatus::Queued => JobStatus::Queued,
2307        JobStateStatus::Running => JobStatus::Running,
2308        JobStateStatus::Paused => JobStatus::Paused,
2309        JobStateStatus::Completed => JobStatus::Completed,
2310        JobStateStatus::Failed => JobStatus::Failed,
2311        JobStateStatus::Cancelled => JobStatus::Cancelled,
2312    }
2313}
2314
2315#[cfg(test)]
2316mod tests {
2317    use super::*;
2318    use codewhale_protocol::ThreadResumeParams;
2319    use codewhale_tools::ToolCallSource;
2320
2321    fn temp_core_state(name: &str) -> StateStore {
2322        let dir =
2323            std::env::temp_dir().join(format!("codewhale-core-{name}-{}", Uuid::new_v4().simple()));
2324        std::fs::create_dir_all(&dir).expect("create temp state dir");
2325        StateStore::open(Some(dir.join("state.db"))).expect("open state store")
2326    }
2327
2328    fn test_thread_metadata(id: &str) -> ThreadMetadata {
2329        ThreadMetadata {
2330            id: id.to_string(),
2331            rollout_path: None,
2332            preview: "test thread".to_string(),
2333            ephemeral: false,
2334            model_provider: "deepseek".to_string(),
2335            created_at: 10,
2336            updated_at: 10,
2337            status: PersistedThreadStatus::Running,
2338            path: None,
2339            cwd: PathBuf::from("/tmp/codewhale"),
2340            cli_version: "0.0.0-test".to_string(),
2341            source: SessionSource::Interactive,
2342            name: None,
2343            sandbox_policy: None,
2344            approval_mode: None,
2345            archived: false,
2346            archived_at: None,
2347            git_sha: None,
2348            git_branch: None,
2349            git_origin_url: None,
2350            memory_mode: None,
2351            current_leaf_id: None,
2352        }
2353    }
2354
2355    // ── JobManager: lifecycle ──────────────────────────────────────────
2356
2357    #[test]
2358    fn permission_path_for_call_extracts_function_path_argument() {
2359        let call = ToolCall {
2360            name: "read_file".to_string(),
2361            payload: ToolPayload::Function {
2362                arguments: json!({ "path": "README.md" }).to_string(),
2363            },
2364            source: ToolCallSource::Direct,
2365            raw_tool_call_id: None,
2366        };
2367
2368        assert_eq!(
2369            permission_path_for_call(&call).as_deref(),
2370            Some("README.md")
2371        );
2372    }
2373
2374    #[test]
2375    fn permission_path_for_call_extracts_mcp_path_argument() {
2376        let call = ToolCall {
2377            name: "mcp_fs_read".to_string(),
2378            payload: ToolPayload::Mcp {
2379                server: "fs".to_string(),
2380                tool: "read".to_string(),
2381                raw_arguments: json!({ "path": "secrets/token.txt" }),
2382                raw_tool_call_id: None,
2383            },
2384            source: ToolCallSource::Direct,
2385            raw_tool_call_id: None,
2386        };
2387
2388        assert_eq!(
2389            permission_path_for_call(&call).as_deref(),
2390            Some("secrets/token.txt")
2391        );
2392    }
2393
2394    #[test]
2395    fn permission_path_for_call_ignores_shell_payload() {
2396        let call = ToolCall {
2397            name: "exec_shell".to_string(),
2398            payload: ToolPayload::LocalShell {
2399                params: codewhale_protocol::LocalShellParams {
2400                    command: "cargo test".to_string(),
2401                    cwd: None,
2402                    timeout_ms: None,
2403                },
2404            },
2405            source: ToolCallSource::Direct,
2406            raw_tool_call_id: None,
2407        };
2408
2409        assert_eq!(permission_path_for_call(&call), None);
2410    }
2411
2412    #[test]
2413    fn thread_goal_progress_accumulates_durable_accounting() {
2414        let store = temp_core_state("thread-goal-progress");
2415        store
2416            .upsert_thread(&test_thread_metadata("thread-1"))
2417            .expect("upsert thread");
2418        let mut manager = ThreadManager::new(store);
2419        manager
2420            .set_thread_goal(&ThreadGoalSetParams {
2421                thread_id: "thread-1".to_string(),
2422                objective: "Carry the goal across turns".to_string(),
2423                token_budget: Some(2_000),
2424            })
2425            .expect("set goal")
2426            .expect("goal exists");
2427
2428        let updated = manager
2429            .record_thread_goal_progress(&ThreadGoalProgressParams {
2430                thread_id: "thread-1".to_string(),
2431                token_delta: 750,
2432                time_delta_seconds: 12,
2433                record_continuation: true,
2434            })
2435            .expect("record progress")
2436            .expect("goal exists");
2437
2438        assert_eq!(updated.tokens_used, 750);
2439        assert_eq!(updated.time_used_seconds, 12);
2440        assert_eq!(updated.continuation_count, 1);
2441
2442        let persisted = manager
2443            .get_thread_goal(&ThreadGoalGetParams {
2444                thread_id: "thread-1".to_string(),
2445            })
2446            .expect("read goal")
2447            .expect("goal exists");
2448        assert_eq!(persisted.tokens_used, 750);
2449        assert_eq!(persisted.time_used_seconds, 12);
2450        assert_eq!(persisted.continuation_count, 1);
2451    }
2452
2453    #[test]
2454    fn approval_request_frame_includes_matched_rule() {
2455        let requirement = ExecApprovalRequirement::NeedsApproval {
2456            reason: "Typed ask rule 'tool=exec_shell command=cargo test' requires approval."
2457                .to_string(),
2458            proposed_execpolicy_amendment: None,
2459            proposed_network_policy_amendments: Vec::new(),
2460        };
2461
2462        let frame = approval_request_frame(
2463            &requirement,
2464            Some("tool=exec_shell command=cargo test"),
2465            "call-1".to_string(),
2466            "approval-1".to_string(),
2467            "turn-1".to_string(),
2468            "cargo test --workspace".to_string(),
2469            "/repo".to_string(),
2470        )
2471        .expect("approval frame");
2472
2473        let EventFrame::ExecApprovalRequest { request } = frame else {
2474            panic!("expected exec approval request frame");
2475        };
2476        assert_eq!(
2477            request.matched_rule.as_deref(),
2478            Some("tool=exec_shell command=cargo test")
2479        );
2480        assert_eq!(request.reason, requirement.reason());
2481    }
2482
2483    #[test]
2484    fn user_input_request_frame_lifts_questions_from_arguments() {
2485        // issue #3102: the headless frame constructor must parse the model's
2486        // `request_user_input` arguments and lift the questions into the
2487        // UserInputRequestEvent, defaulting the boolean flags when omitted.
2488        let arguments = r#"{"questions":[{"header":"Scope","id":"scope","question":"Which?","options":[{"label":"A","description":"a"},{"label":"B","description":"b"}],"allow_free_text":true}]}"#;
2489        let frame = user_input_request_frame(
2490            "call-1".to_string(),
2491            "turn-1".to_string(),
2492            "ui-1".to_string(),
2493            arguments,
2494        )
2495        .expect("user input frame");
2496
2497        let EventFrame::UserInputRequest { request } = frame else {
2498            panic!("expected user_input_request frame");
2499        };
2500        assert_eq!(request.call_id, "call-1");
2501        assert_eq!(request.turn_id, "turn-1");
2502        assert_eq!(request.request_id, "ui-1");
2503        assert_eq!(request.questions.len(), 1);
2504        assert_eq!(request.questions[0].id, "scope");
2505        assert!(request.questions[0].allow_free_text);
2506        // multi_select omitted in the payload → defaults to false.
2507        assert!(!request.questions[0].multi_select);
2508        assert_eq!(request.questions[0].options.len(), 2);
2509    }
2510
2511    #[test]
2512    fn user_input_request_frame_returns_none_on_invalid_arguments() {
2513        // On parse failure the constructor returns None so invoke_tool falls
2514        // through to the generic tool error path instead of silently dropping.
2515        let frame = user_input_request_frame(
2516            "call-1".to_string(),
2517            "turn-1".to_string(),
2518            "ui-1".to_string(),
2519            "not json",
2520        );
2521        assert!(frame.is_none());
2522
2523        // Valid JSON but missing the questions array is also rejected.
2524        let frame = user_input_request_frame(
2525            "call-1".to_string(),
2526            "turn-1".to_string(),
2527            "ui-1".to_string(),
2528            r#"{"foo":"bar"}"#,
2529        );
2530        assert!(frame.is_none());
2531    }
2532
2533    #[test]
2534    fn enqueue_creates_queued_job_with_zero_progress() {
2535        let mut jm = JobManager::default();
2536        let job = jm.enqueue("build");
2537        assert_eq!(job.name, "build");
2538        assert_eq!(job.status, JobStatus::Queued);
2539        assert_eq!(job.progress, Some(0));
2540        assert!(job.detail.is_none());
2541        assert_eq!(job.history.len(), 1);
2542        assert_eq!(job.history[0].phase, "created");
2543    }
2544
2545    #[test]
2546    fn set_running_transitions_from_queued() {
2547        let mut jm = JobManager::default();
2548        let job = jm.enqueue("deploy");
2549        let id = job.id.clone();
2550        jm.set_running(&id);
2551        let jobs = jm.list();
2552        let updated = jobs.iter().find(|j| j.id == id).unwrap();
2553        assert_eq!(updated.status, JobStatus::Running);
2554        assert_eq!(updated.history.last().unwrap().phase, "running");
2555    }
2556
2557    #[test]
2558    fn update_progress_clamps_to_100() {
2559        let mut jm = JobManager::default();
2560        let job = jm.enqueue("task");
2561        let id = job.id.clone();
2562        jm.update_progress(&id, 150, Some("over".to_string()));
2563        let jobs = jm.list();
2564        let updated = jobs.iter().find(|j| j.id == id).unwrap();
2565        assert_eq!(updated.progress, Some(100));
2566    }
2567
2568    #[test]
2569    fn complete_sets_progress_to_100() {
2570        let mut jm = JobManager::default();
2571        let job = jm.enqueue("task");
2572        let id = job.id.clone();
2573        jm.set_running(&id);
2574        jm.complete(&id);
2575        let jobs = jm.list();
2576        let updated = jobs.iter().find(|j| j.id == id).unwrap();
2577        assert_eq!(updated.status, JobStatus::Completed);
2578        assert_eq!(updated.progress, Some(100));
2579    }
2580
2581    #[test]
2582    fn fail_increments_attempt_and_sets_backoff() {
2583        let mut jm = JobManager::default();
2584        let job = jm.enqueue("fragile");
2585        let id = job.id.clone();
2586        jm.set_running(&id);
2587        jm.fail(&id, "crashed");
2588        let jobs = jm.list();
2589        let updated = jobs.iter().find(|j| j.id == id).unwrap();
2590        assert_eq!(updated.status, JobStatus::Failed);
2591        assert_eq!(updated.retry.attempt, 1);
2592        assert!(updated.retry.next_backoff_ms > 0);
2593        assert!(updated.retry.next_retry_at.is_some());
2594        assert_eq!(updated.detail.as_deref(), Some("crashed"));
2595    }
2596
2597    #[test]
2598    fn fail_clears_retry_after_max_attempts() {
2599        let mut jm = JobManager::default();
2600        let job = jm.enqueue("fragile");
2601        let id = job.id.clone();
2602        for _ in 0..=DEFAULT_JOB_MAX_ATTEMPTS {
2603            jm.set_running(&id);
2604            jm.fail(&id, "boom");
2605        }
2606        let jobs = jm.list();
2607        let updated = jobs.iter().find(|j| j.id == id).unwrap();
2608        assert_eq!(updated.retry.attempt, DEFAULT_JOB_MAX_ATTEMPTS);
2609        assert_eq!(updated.retry.next_backoff_ms, 0);
2610        assert!(updated.retry.next_retry_at.is_none());
2611    }
2612
2613    #[test]
2614    fn cancel_sets_status_and_clears_retry() {
2615        let mut jm = JobManager::default();
2616        let job = jm.enqueue("task");
2617        let id = job.id.clone();
2618        jm.cancel(&id);
2619        let jobs = jm.list();
2620        let updated = jobs.iter().find(|j| j.id == id).unwrap();
2621        assert_eq!(updated.status, JobStatus::Cancelled);
2622        assert_eq!(updated.retry.next_backoff_ms, 0);
2623    }
2624
2625    #[test]
2626    fn pause_and_resume_round_trip() {
2627        let mut jm = JobManager::default();
2628        let job = jm.enqueue("task");
2629        let id = job.id.clone();
2630        jm.set_running(&id);
2631        jm.pause(&id, Some("waiting".to_string()));
2632        let jobs = jm.list();
2633        let paused = jobs.iter().find(|j| j.id == id).unwrap();
2634        assert_eq!(paused.status, JobStatus::Paused);
2635        assert_eq!(paused.detail.as_deref(), Some("waiting"));
2636
2637        jm.resume(&id, None);
2638        let jobs = jm.list();
2639        let resumed = jobs.iter().find(|j| j.id == id).unwrap();
2640        assert_eq!(resumed.status, JobStatus::Running);
2641        assert_eq!(resumed.history.last().unwrap().phase, "resumed");
2642    }
2643
2644    #[test]
2645    fn list_returns_jobs_sorted_by_updated_at_desc() {
2646        let mut jm = JobManager::default();
2647        jm.enqueue("first");
2648        jm.enqueue("second");
2649        jm.enqueue("third");
2650        let jobs = jm.list();
2651        assert_eq!(jobs.len(), 3);
2652        for window in jobs.windows(2) {
2653            assert!(window[0].updated_at >= window[1].updated_at);
2654        }
2655    }
2656
2657    #[test]
2658    fn history_returns_entries_for_existing_job() {
2659        let mut jm = JobManager::default();
2660        let job = jm.enqueue("task");
2661        let id = job.id.clone();
2662        jm.set_running(&id);
2663        jm.complete(&id);
2664        let history = jm.history(&id);
2665        assert_eq!(history.len(), 3); // created, running, completed
2666        assert_eq!(history[0].phase, "created");
2667        assert_eq!(history[1].phase, "running");
2668        assert_eq!(history[2].phase, "completed");
2669    }
2670
2671    #[test]
2672    fn history_returns_empty_for_unknown_job() {
2673        let jm = JobManager::default();
2674        assert!(jm.history("nonexistent").is_empty());
2675    }
2676
2677    #[test]
2678    fn resume_pending_requeues_running_and_queued() {
2679        let mut jm = JobManager::default();
2680        let _j1 = jm.enqueue("queued_task");
2681        let j2 = jm.enqueue("running_task");
2682        let j3 = jm.enqueue("completed_task");
2683        let id2 = j2.id.clone();
2684        let id3 = j3.id.clone();
2685        jm.set_running(&id2);
2686        jm.set_running(&id3);
2687        jm.complete(&id3);
2688
2689        let resumed = jm.resume_pending();
2690        assert_eq!(resumed.len(), 2);
2691        for job in &resumed {
2692            assert_eq!(job.status, JobStatus::Queued);
2693        }
2694    }
2695
2696    // ── JobManager: backoff ────────────────────────────────────────────
2697
2698    #[test]
2699    fn deterministic_backoff_zero_on_first_attempt() {
2700        let retry = JobRetryMetadata {
2701            attempt: 0,
2702            ..Default::default()
2703        };
2704        assert_eq!(JobManager::deterministic_backoff_ms(&retry), 0);
2705    }
2706
2707    #[test]
2708    fn deterministic_backoff_exponential_growth() {
2709        let base = DEFAULT_JOB_BACKOFF_BASE_MS;
2710        for attempt in 1..=5 {
2711            let retry = JobRetryMetadata {
2712                attempt,
2713                backoff_base_ms: base,
2714                ..Default::default()
2715            };
2716            let expected = base * 2u64.pow(attempt.saturating_sub(1).min(20));
2717            assert_eq!(
2718                JobManager::deterministic_backoff_ms(&retry),
2719                expected,
2720                "attempt {attempt}"
2721            );
2722        }
2723    }
2724
2725    #[test]
2726    fn deterministic_backoff_saturates_at_high_exponent() {
2727        let retry = JobRetryMetadata {
2728            attempt: 63,
2729            backoff_base_ms: 1000,
2730            ..Default::default()
2731        };
2732        // Should not panic; result saturates
2733        let _ = JobManager::deterministic_backoff_ms(&retry);
2734    }
2735
2736    // ── JobManager: history truncation ─────────────────────────────────
2737
2738    #[test]
2739    fn push_history_truncates_beyond_max() {
2740        let mut jm = JobManager::default();
2741        let job = jm.enqueue("task");
2742        let id = job.id.clone();
2743        // Generate more history entries than the limit
2744        for i in 0..(MAX_JOB_HISTORY_ENTRIES + 20) {
2745            jm.update_progress(&id, (i % 100) as u8, Some(format!("step {i}")));
2746        }
2747        let history = jm.history(&id);
2748        assert_eq!(history.len(), MAX_JOB_HISTORY_ENTRIES);
2749    }
2750
2751    // ── JobManager: persistence encoding/parsing ───────────────────────
2752
2753    #[test]
2754    fn encode_and_parse_persisted_detail_round_trip() {
2755        let mut jm = JobManager::default();
2756        let job = jm.enqueue("task");
2757        let id = job.id.clone();
2758        jm.set_running(&id);
2759        jm.fail(&id, "oops");
2760        let job = jm.list().into_iter().find(|j| j.id == id).unwrap();
2761
2762        let encoded = JobManager::encode_persisted_detail(&job).unwrap().unwrap();
2763        let parsed = JobManager::parse_persisted_detail(Some(&encoded)).unwrap();
2764
2765        assert_eq!(parsed.status, job.status);
2766        assert_eq!(parsed.detail, job.detail);
2767        assert_eq!(parsed.retry.attempt, job.retry.attempt);
2768        assert_eq!(parsed.history.len(), job.history.len());
2769    }
2770
2771    #[test]
2772    fn parse_persisted_detail_returns_none_for_none_input() {
2773        assert!(JobManager::parse_persisted_detail(None).is_none());
2774    }
2775
2776    #[test]
2777    fn parse_persisted_detail_returns_none_for_invalid_json() {
2778        assert!(JobManager::parse_persisted_detail(Some("not json")).is_none());
2779    }
2780
2781    // ── Helper functions ───────────────────────────────────────────────
2782
2783    #[test]
2784    fn job_status_round_trip_str() {
2785        let statuses = [
2786            JobStatus::Queued,
2787            JobStatus::Running,
2788            JobStatus::Paused,
2789            JobStatus::Completed,
2790            JobStatus::Failed,
2791            JobStatus::Cancelled,
2792        ];
2793        for status in &statuses {
2794            let s = job_status_to_str(*status);
2795            let parsed = job_status_from_str(s);
2796            assert_eq!(parsed, Some(*status), "round-trip failed for {s:?}");
2797        }
2798    }
2799
2800    #[test]
2801    fn job_status_from_str_returns_none_for_unknown() {
2802        assert_eq!(job_status_from_str("unknown"), None);
2803        assert_eq!(job_status_from_str(""), None);
2804    }
2805
2806    #[test]
2807    fn truncate_preview_limits_to_120_chars() {
2808        let long = "a".repeat(200);
2809        let truncated = truncate_preview(&long);
2810        assert_eq!(truncated.len(), 120);
2811    }
2812
2813    #[test]
2814    fn truncate_preview_preserves_short_strings() {
2815        let short = "hello";
2816        assert_eq!(truncate_preview(short), "hello");
2817    }
2818
2819    #[test]
2820    fn runtime_status_to_job_state_maps_correctly() {
2821        assert_eq!(
2822            runtime_status_to_job_state(JobStatus::Queued),
2823            JobStateStatus::Queued
2824        );
2825        assert_eq!(
2826            runtime_status_to_job_state(JobStatus::Running),
2827            JobStateStatus::Running
2828        );
2829        assert_eq!(
2830            runtime_status_to_job_state(JobStatus::Paused),
2831            JobStateStatus::Paused
2832        );
2833        assert_eq!(
2834            runtime_status_to_job_state(JobStatus::Completed),
2835            JobStateStatus::Completed
2836        );
2837        assert_eq!(
2838            runtime_status_to_job_state(JobStatus::Failed),
2839            JobStateStatus::Failed
2840        );
2841        assert_eq!(
2842            runtime_status_to_job_state(JobStatus::Cancelled),
2843            JobStateStatus::Cancelled
2844        );
2845    }
2846
2847    #[test]
2848    fn job_state_status_to_runtime_maps_correctly() {
2849        assert_eq!(
2850            job_state_status_to_runtime(JobStateStatus::Queued),
2851            JobStatus::Queued
2852        );
2853        assert_eq!(
2854            job_state_status_to_runtime(JobStateStatus::Running),
2855            JobStatus::Running
2856        );
2857        assert_eq!(
2858            job_state_status_to_runtime(JobStateStatus::Paused),
2859            JobStatus::Paused
2860        );
2861        assert_eq!(
2862            job_state_status_to_runtime(JobStateStatus::Completed),
2863            JobStatus::Completed
2864        );
2865        assert_eq!(
2866            job_state_status_to_runtime(JobStateStatus::Failed),
2867            JobStatus::Failed
2868        );
2869        assert_eq!(
2870            job_state_status_to_runtime(JobStateStatus::Cancelled),
2871            JobStatus::Cancelled
2872        );
2873    }
2874
2875    #[test]
2876    fn preview_from_initial_history_new() {
2877        let preview = preview_from_initial_history(&InitialHistory::New);
2878        assert_eq!(preview, "New conversation");
2879    }
2880
2881    #[test]
2882    fn preview_from_initial_history_forked() {
2883        let preview = preview_from_initial_history(&InitialHistory::Forked(vec![json!("hello")]));
2884        assert!(preview.contains("hello"));
2885    }
2886
2887    #[test]
2888    fn preview_from_initial_history_resumed() {
2889        let preview = preview_from_initial_history(&InitialHistory::Resumed {
2890            conversation_id: "test".to_string(),
2891            history: vec![json!("world")],
2892            rollout_path: PathBuf::from("/tmp/test"),
2893        });
2894        assert!(preview.contains("world"));
2895    }
2896
2897    #[test]
2898    fn json_optional_string_handles_null() {
2899        assert!(json_optional_string(&Value::Null).is_none());
2900    }
2901
2902    #[test]
2903    fn json_optional_string_handles_string() {
2904        assert_eq!(
2905            json_optional_string(&Value::String("hello".to_string())),
2906            Some("hello".to_string())
2907        );
2908    }
2909
2910    #[test]
2911    fn json_optional_string_handles_non_string() {
2912        assert!(json_optional_string(&json!(42)).is_none());
2913    }
2914
2915    #[test]
2916    fn parse_retry_metadata_returns_default_for_none() {
2917        let retry = parse_retry_metadata(None);
2918        assert_eq!(retry.attempt, 0);
2919        assert_eq!(retry.max_attempts, DEFAULT_JOB_MAX_ATTEMPTS);
2920        assert_eq!(retry.backoff_base_ms, DEFAULT_JOB_BACKOFF_BASE_MS);
2921    }
2922
2923    #[test]
2924    fn parse_retry_metadata_parses_fields() {
2925        let value = json!({
2926            "attempt": 2,
2927            "max_attempts": 5,
2928            "backoff_base_ms": 1000,
2929            "next_backoff_ms": 2000,
2930            "next_retry_at": 1234567890i64
2931        });
2932        let retry = parse_retry_metadata(Some(&value));
2933        assert_eq!(retry.attempt, 2);
2934        assert_eq!(retry.max_attempts, 5);
2935        assert_eq!(retry.backoff_base_ms, 1000);
2936        assert_eq!(retry.next_backoff_ms, 2000);
2937        assert_eq!(retry.next_retry_at, Some(1234567890));
2938    }
2939
2940    #[test]
2941    fn parse_history_entry_returns_none_without_status() {
2942        let value = json!({"at": 1, "phase": "test"});
2943        assert!(parse_history_entry(&value).is_none());
2944    }
2945
2946    #[test]
2947    fn parse_history_entry_parses_valid_entry() {
2948        let value = json!({
2949            "at": 100,
2950            "phase": "running",
2951            "status": "running",
2952            "progress": 50,
2953            "detail": "working",
2954            "retry": {"attempt": 0, "max_attempts": 3, "backoff_base_ms": 500}
2955        });
2956        let entry = parse_history_entry(&value).unwrap();
2957        assert_eq!(entry.at, 100);
2958        assert_eq!(entry.phase, "running");
2959        assert_eq!(entry.status, JobStatus::Running);
2960        assert_eq!(entry.progress, Some(50));
2961        assert_eq!(entry.detail.as_deref(), Some("working"));
2962    }
2963
2964    #[test]
2965    fn paused_job_persists_as_paused_not_running() {
2966        let store = temp_core_state("paused-persist");
2967        let mut jm = JobManager::default();
2968        let job = jm.enqueue("task");
2969        let id = job.id.clone();
2970        jm.set_running(&id);
2971        jm.pause(&id, Some("waiting".to_string()));
2972        jm.persist_job(&store, &id).expect("persist paused job");
2973
2974        let persisted = store.list_jobs(Some(10)).expect("list jobs");
2975        let record = persisted.iter().find(|job| job.id == id).unwrap();
2976        assert_eq!(record.status, JobStateStatus::Paused);
2977
2978        let mut reloaded = JobManager::default();
2979        reloaded.load_from_store(&store).expect("reload jobs");
2980        let jobs = reloaded.list();
2981        let reloaded_job = jobs.iter().find(|job| job.id == id).unwrap();
2982        assert_eq!(reloaded_job.status, JobStatus::Paused);
2983    }
2984
2985    // ── O1: JobRecord → AgentRunSnapshot adapter ────────────────────────
2986
2987    fn sample_job_record(status: JobStatus, detail: Option<&str>) -> JobRecord {
2988        JobRecord {
2989            id: "job-o1-1".to_string(),
2990            name: "sample".to_string(),
2991            status,
2992            progress: None,
2993            detail: detail.map(str::to_string),
2994            retry: JobRetryMetadata {
2995                attempt: 0,
2996                max_attempts: DEFAULT_JOB_MAX_ATTEMPTS,
2997                backoff_base_ms: DEFAULT_JOB_BACKOFF_BASE_MS,
2998                next_backoff_ms: 0,
2999                next_retry_at: None,
3000            },
3001            history: Vec::new(),
3002            created_at: 1_700_000_000,
3003            updated_at: 1_700_000_042,
3004        }
3005    }
3006
3007    #[test]
3008    fn job_record_to_agent_run_maps_non_terminal_states() {
3009        use codewhale_protocol::agent_run::RunState;
3010
3011        for (status, expected) in [
3012            (JobStatus::Queued, RunState::Queued),
3013            (JobStatus::Running, RunState::Running),
3014            (JobStatus::Paused, RunState::Paused),
3015        ] {
3016            let snapshot = job_record_to_agent_run(&sample_job_record(status, None));
3017            assert!(snapshot.is_coherent());
3018            assert_eq!(snapshot.run_id, "job-o1-1");
3019            assert_eq!(snapshot.parent, None);
3020            assert_eq!(
3021                snapshot.source,
3022                codewhale_protocol::agent_run::RunSource::CoreJob
3023            );
3024            assert_eq!(snapshot.state, expected);
3025            assert!(snapshot.terminal.is_none());
3026            assert!(snapshot.refs.is_empty());
3027            assert_eq!(
3028                snapshot.budget,
3029                codewhale_protocol::agent_run::BudgetSummary::default()
3030            );
3031        }
3032    }
3033
3034    #[test]
3035    fn job_record_to_agent_run_maps_terminal_states_without_fabricating_fields() {
3036        use codewhale_protocol::agent_run::{RunState, TerminalOutcome};
3037
3038        let cases = [
3039            (
3040                JobStatus::Completed,
3041                TerminalOutcome::Completed,
3042                Some("done"),
3043            ),
3044            (JobStatus::Failed, TerminalOutcome::Failed, Some("boom")),
3045            (JobStatus::Cancelled, TerminalOutcome::Cancelled, None),
3046        ];
3047
3048        for (status, outcome, detail) in cases {
3049            let snapshot = job_record_to_agent_run(&sample_job_record(status, detail));
3050            assert!(snapshot.is_coherent());
3051            assert_eq!(snapshot.state, RunState::Terminal);
3052            let terminal = snapshot.terminal.expect("terminal summary");
3053            assert_eq!(terminal.outcome, outcome);
3054            assert_eq!(terminal.ended_at_ms, Some(1_700_000_042_000));
3055            assert_eq!(terminal.detail, None);
3056            assert_eq!(
3057                snapshot.budget,
3058                codewhale_protocol::agent_run::BudgetSummary::default()
3059            );
3060            assert!(snapshot.refs.is_empty());
3061            assert_eq!(snapshot.parent, None);
3062        }
3063    }
3064
3065    #[test]
3066    fn job_record_to_agent_run_does_not_export_unclassified_detail() {
3067        let record = sample_job_record(JobStatus::Failed, Some("owner-private diagnostic"));
3068        let snapshot = job_record_to_agent_run(&record);
3069        let terminal = snapshot.terminal.as_ref().expect("terminal summary");
3070        assert_eq!(terminal.detail, None);
3071        let serialized = serde_json::to_string(&snapshot).expect("serialize snapshot");
3072        assert!(!serialized.contains("owner-private diagnostic"));
3073    }
3074
3075    #[test]
3076    fn job_record_to_agent_run_omits_ended_at_on_updated_at_overflow() {
3077        let mut record = sample_job_record(JobStatus::Completed, Some("ok"));
3078        record.updated_at = i64::MAX;
3079        let snapshot = job_record_to_agent_run(&record);
3080        assert!(snapshot.is_coherent());
3081        let terminal = snapshot.terminal.expect("terminal summary");
3082        assert_eq!(terminal.ended_at_ms, None);
3083    }
3084
3085    #[test]
3086    fn unarchive_thread_updates_running_threads_cache() {
3087        let store = temp_core_state("unarchive-cache");
3088        let mut manager = ThreadManager::new(store);
3089        let spawned = manager
3090            .spawn_thread_with_history(
3091                "deepseek".to_string(),
3092                PathBuf::from("/tmp/codewhale"),
3093                InitialHistory::New,
3094                true,
3095            )
3096            .expect("spawn thread");
3097        let thread_id = spawned.thread.id.clone();
3098        let resume_params = ThreadResumeParams {
3099            thread_id: thread_id.clone(),
3100            history: None,
3101            path: None,
3102            model: None,
3103            model_provider: None,
3104            cwd: None,
3105            approval_policy: None,
3106            sandbox: None,
3107            config: None,
3108            base_instructions: None,
3109            developer_instructions: None,
3110            personality: None,
3111            persist_extended_history: false,
3112        };
3113
3114        manager.archive_thread(&thread_id).expect("archive thread");
3115        let archived = manager
3116            .resume_thread_with_history(
3117                &resume_params,
3118                Path::new("/tmp/codewhale"),
3119                "deepseek".to_string(),
3120            )
3121            .expect("resume archived thread")
3122            .expect("thread in cache");
3123        assert_eq!(archived.thread.status, ThreadStatus::Archived);
3124
3125        manager
3126            .unarchive_thread(&thread_id)
3127            .expect("unarchive thread");
3128        let restored = manager
3129            .resume_thread_with_history(
3130                &resume_params,
3131                Path::new("/tmp/codewhale"),
3132                "deepseek".to_string(),
3133            )
3134            .expect("resume unarchived thread")
3135            .expect("thread in cache");
3136        assert_eq!(restored.thread.status, ThreadStatus::Idle);
3137    }
3138
3139    #[test]
3140    fn resume_with_history_does_not_reappend_persisted_messages() {
3141        // A read→resume flow hands the thread's own history back to
3142        // `thread/resume`; appending it verbatim doubled the conversation on
3143        // every resume, compounding.
3144        let store = temp_core_state("resume-history-dedup");
3145        let mut manager = ThreadManager::new(store);
3146        let history = vec![
3147            json!({"type": "user_message", "message": "hello"}),
3148            json!({"type": "assistant_message", "message": "hi there"}),
3149        ];
3150        let spawned = manager
3151            .spawn_thread_with_history(
3152                "deepseek".to_string(),
3153                PathBuf::from("/tmp/codewhale"),
3154                InitialHistory::Forked(history.clone()),
3155                true,
3156            )
3157            .expect("spawn thread");
3158        let thread_id = spawned.thread.id.clone();
3159        let message_count = |manager: &ThreadManager| {
3160            manager
3161                .state_store()
3162                .list_messages(&thread_id, None)
3163                .expect("list messages")
3164                .len()
3165        };
3166        assert_eq!(message_count(&manager), 2);
3167
3168        let resume_params = ThreadResumeParams {
3169            thread_id: thread_id.clone(),
3170            history: Some(history.clone()),
3171            path: None,
3172            model: None,
3173            model_provider: None,
3174            cwd: None,
3175            approval_policy: None,
3176            sandbox: None,
3177            config: None,
3178            base_instructions: None,
3179            developer_instructions: None,
3180            personality: None,
3181            persist_extended_history: false,
3182        };
3183
3184        // Resuming twice with the same history must be idempotent.
3185        for _ in 0..2 {
3186            manager
3187                .resume_thread_with_history(
3188                    &resume_params,
3189                    Path::new("/tmp/codewhale"),
3190                    "deepseek".to_string(),
3191                )
3192                .expect("resume thread")
3193                .expect("thread found");
3194        }
3195        assert_eq!(
3196            message_count(&manager),
3197            2,
3198            "resume re-appended messages already on the persisted chain"
3199        );
3200
3201        // A genuinely new history item is still appended, exactly once.
3202        let mut extended = history.clone();
3203        extended.push(json!({"type": "user_message", "message": "something new"}));
3204        let resume_params = ThreadResumeParams {
3205            history: Some(extended),
3206            ..resume_params
3207        };
3208        manager
3209            .resume_thread_with_history(
3210                &resume_params,
3211                Path::new("/tmp/codewhale"),
3212                "deepseek".to_string(),
3213            )
3214            .expect("resume thread")
3215            .expect("thread found");
3216        assert_eq!(message_count(&manager), 3);
3217    }
3218
3219    #[test]
3220    fn persist_thread_preserves_stored_policy() {
3221        // persist_thread's update payload carries no per-thread policy;
3222        // writing NULLs unconditionally erased any policy stored earlier
3223        // (e.g. on every resume).
3224        let store = temp_core_state("persist-policy");
3225        let mut metadata = test_thread_metadata("thread-policy");
3226        metadata.sandbox_policy = Some("workspace-write".to_string());
3227        metadata.approval_mode = Some("on-request".to_string());
3228        store.upsert_thread(&metadata).expect("seed thread");
3229
3230        // A fresh manager has an empty running-thread cache, so resume goes
3231        // through the persisted path, which calls persist_thread.
3232        let mut manager = ThreadManager::new(store);
3233        let resume_params = ThreadResumeParams {
3234            thread_id: "thread-policy".to_string(),
3235            history: None,
3236            path: None,
3237            model: None,
3238            model_provider: None,
3239            cwd: None,
3240            approval_policy: None,
3241            sandbox: None,
3242            config: None,
3243            base_instructions: None,
3244            developer_instructions: None,
3245            personality: None,
3246            persist_extended_history: false,
3247        };
3248        manager
3249            .resume_thread_with_history(
3250                &resume_params,
3251                Path::new("/tmp/codewhale"),
3252                "deepseek".to_string(),
3253            )
3254            .expect("resume thread")
3255            .expect("thread found");
3256
3257        let persisted = manager
3258            .state_store()
3259            .get_thread("thread-policy")
3260            .expect("read thread")
3261            .expect("thread persisted");
3262        assert_eq!(persisted.sandbox_policy.as_deref(), Some("workspace-write"));
3263        assert_eq!(persisted.approval_mode.as_deref(), Some("on-request"));
3264    }
3265
3266    #[tokio::test]
3267    async fn invoke_tool_returns_timeout_status_for_slow_tools() {
3268        use async_trait::async_trait;
3269        use codewhale_agent::ModelRegistry;
3270        use codewhale_config::ConfigToml;
3271        use codewhale_execpolicy::{AskForApproval, ExecPolicyEngine};
3272        use codewhale_hooks::HookDispatcher;
3273        use codewhale_mcp::McpManager;
3274        use codewhale_protocol::{ToolKind, ToolOutput, ToolPayload};
3275        use codewhale_tools::{FunctionCallError, ToolDescriptor, ToolHandler, ToolInvocation};
3276
3277        struct SlowTool;
3278        #[async_trait]
3279        impl ToolHandler for SlowTool {
3280            fn kind(&self) -> ToolKind {
3281                ToolKind::Function
3282            }
3283
3284            async fn handle(
3285                &self,
3286                _invocation: ToolInvocation,
3287            ) -> std::result::Result<ToolOutput, FunctionCallError> {
3288                time::sleep(Duration::from_millis(200)).await;
3289                Ok(ToolOutput::Function {
3290                    body: Some(json!("late")),
3291                    success: true,
3292                })
3293            }
3294        }
3295
3296        let mut registry = ToolRegistry::default();
3297        registry
3298            .register(
3299                ToolDescriptor {
3300                    name: "slow_tool".to_string(),
3301                    input_schema: json!({"type":"object"}),
3302                    output_schema: json!({"type":"object"}),
3303                    supports_parallel_tool_calls: true,
3304                    timeout_ms: None,
3305                },
3306                Arc::new(SlowTool),
3307            )
3308            .expect("register slow tool");
3309
3310        let runtime = Runtime::new(
3311            ConfigToml::default(),
3312            ModelRegistry::default(),
3313            temp_core_state("invoke-tool-timeout"),
3314            Arc::new(registry),
3315            Arc::new(McpManager::default()),
3316            ExecPolicyEngine::new(vec![], vec![]),
3317            HookDispatcher::default(),
3318        );
3319
3320        let result = runtime
3321            .invoke_tool(
3322                ToolCall {
3323                    name: "slow_tool".to_string(),
3324                    payload: ToolPayload::Function {
3325                        arguments: "{}".to_string(),
3326                    },
3327                    source: ToolCallSource::Direct,
3328                    raw_tool_call_id: None,
3329                },
3330                AskForApproval::Never,
3331                Path::new("/tmp/codewhale"),
3332            )
3333            .await
3334            .expect("invoke tool");
3335
3336        assert_eq!(result["status"], "timeout");
3337        assert_eq!(result["ok"], false);
3338    }
3339
3340    #[tokio::test]
3341    async fn thread_message_response_ids_are_unique_for_equal_length_inputs() {
3342        // The Message arm used to key response_id as `{thread_id}:{input.len()}`,
3343        // so any two equal-length messages collided and hooks could not tell
3344        // their ResponseStart/ResponseEnd pairs apart.
3345        let mut runtime = Runtime::new(
3346            ConfigToml::default(),
3347            ModelRegistry::default(),
3348            temp_core_state("response-id-unique"),
3349            Arc::new(ToolRegistry::default()),
3350            Arc::new(McpManager::default()),
3351            ExecPolicyEngine::new(vec![], vec![]),
3352            HookDispatcher::default(),
3353        );
3354        let spawned = runtime
3355            .thread_manager
3356            .spawn_thread_with_history(
3357                "deepseek".to_string(),
3358                PathBuf::from("/tmp/codewhale"),
3359                InitialHistory::New,
3360                true,
3361            )
3362            .expect("spawn thread");
3363        let thread_id = spawned.thread.id.clone();
3364
3365        let mut response_ids = Vec::new();
3366        for input in ["aaaa", "bbbb"] {
3367            let response = runtime
3368                .handle_thread(ThreadRequest::Message {
3369                    thread_id: thread_id.clone(),
3370                    input: input.to_string(),
3371                })
3372                .await
3373                .expect("handle message");
3374            let response_id = response
3375                .events
3376                .iter()
3377                .find_map(|frame| match frame {
3378                    EventFrame::ResponseStart { response_id } => Some(response_id.clone()),
3379                    _ => None,
3380                })
3381                .expect("response start event");
3382            response_ids.push(response_id);
3383        }
3384
3385        assert_ne!(
3386            response_ids[0], response_ids[1],
3387            "equal-length inputs must not share a response_id"
3388        );
3389        assert!(
3390            response_ids.iter().all(|id| id.starts_with("resp-")),
3391            "response ids should use the resp-<uuid> shape: {response_ids:?}"
3392        );
3393    }
3394}