返回 CodeWhale
lib.rs
根目录 / crates / core / src / lib.rs
1 pub mod context_reference;
2 pub mod fragments;
3 pub mod ids;
4 pub mod journal;
5 pub mod prefix_cache;
6 pub mod request;
7 pub mod role;
8 pub mod secret_eq;
9 pub mod session;
10 pub mod tool_parser;
11
12 pub use context_reference::{
13 ContextReference, ContextReferenceKind, ContextReferenceSource, MediaAttachmentReference,
14 media_attachment_references,
15 };
16
17 use std::collections::HashMap;
18
19 use anyhow::Result;
20 use codewhale_config::{ConfigToml, ProviderKind};
21 use codewhale_hooks::HookDispatcher;
22 use codewhale_protocol::{AppResponse, EventFrame, ResponseChannel, Status};
23 use codewhale_state::{JobStateRecord, JobStateStatus, StateStore};
24 use serde_json::{Value, json};
25 use uuid::Uuid;
26
27 /// Status of a background job.
28 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
29 pub enum JobStatus {
30 /// Waiting to be picked up.
31 Queued,
32 /// Currently executing.
33 Running,
34 /// Temporarily paused.
35 Paused,
36 /// Finished successfully.
37 Completed,
38 /// Finished with an error.
39 Failed,
40 /// Cancelled by the user.
41 Cancelled,
42 }
43
44 impl Status for JobStatus {
45 fn is_terminal(&self) -> bool {
46 matches!(self, Self::Completed | Self::Failed | Self::Cancelled)
47 }
48 fn is_active(&self) -> bool {
49 matches!(self, Self::Queued | Self::Running)
50 }
51 fn is_paused(&self) -> bool {
52 matches!(self, Self::Paused)
53 }
54 }
55
56 const JOB_DETAIL_SCHEMA_VERSION: u8 = 1;
57 const DEFAULT_JOB_MAX_ATTEMPTS: u32 = 3;
58 const DEFAULT_JOB_BACKOFF_BASE_MS: u64 = 500;
59 const MAX_JOB_HISTORY_ENTRIES: usize = 64;
60
61 /// Retry state for a job that failed and may be retried.
62 #[derive(Debug, Clone)]
63 pub struct JobRetryMetadata {
64 /// Current attempt number (0 = not yet retried).
65 pub attempt: u32,
66 /// Maximum number of retry attempts before giving up.
67 pub max_attempts: u32,
68 /// Base delay in milliseconds for exponential backoff.
69 pub backoff_base_ms: u64,
70 /// Computed delay in milliseconds until the next retry.
71 pub next_backoff_ms: u64,
72 /// Timestamp when the next retry should be attempted.
73 pub next_retry_at: Option<i64>,
74 }
75
76 impl Default for JobRetryMetadata {
77 fn default() -> Self {
78 Self {
79 attempt: 0,
80 max_attempts: DEFAULT_JOB_MAX_ATTEMPTS,
81 backoff_base_ms: DEFAULT_JOB_BACKOFF_BASE_MS,
82 next_backoff_ms: 0,
83 next_retry_at: None,
84 }
85 }
86 }
87
88 /// A single entry in a job's history log.
89 #[derive(Debug, Clone)]
90 pub struct JobHistoryEntry {
91 /// Timestamp when this entry was recorded.
92 pub at: i64,
93 /// Phase name (e.g., "created", "running", "failed").
94 pub phase: String,
95 /// Job status at this point in time.
96 pub status: JobStatus,
97 /// Progress percentage at this point, if available.
98 pub progress: Option<u8>,
99 /// Human-readable detail message.
100 pub detail: Option<String>,
101 /// Retry state snapshot at this point.
102 pub retry: JobRetryMetadata,
103 }
104
105 #[derive(Debug, Clone)]
106 struct PersistedJobDetail {
107 pub status: JobStatus,
108 pub detail: Option<String>,
109 pub retry: JobRetryMetadata,
110 pub history: Vec<JobHistoryEntry>,
111 }
112
113 /// A complete job record with all metadata and history.
114 #[derive(Debug, Clone)]
115 pub struct JobRecord {
116 /// Unique job identifier.
117 pub id: String,
118 /// Human-readable job name.
119 pub name: String,
120 /// Current job status.
121 pub status: JobStatus,
122 /// Current progress percentage (0-100).
123 pub progress: Option<u8>,
124 /// Human-readable detail about the current state.
125 pub detail: Option<String>,
126 /// Retry state for failed jobs.
127 pub retry: JobRetryMetadata,
128 /// Chronological history of state transitions.
129 pub history: Vec<JobHistoryEntry>,
130 /// Timestamp when the job was created.
131 pub created_at: i64,
132 /// Timestamp of the last state change.
133 pub updated_at: i64,
134 }
135
136 /// Map a durable [`JobRecord`] to the dependency-neutral run read model.
137 ///
138 /// Pure projection of the record as persisted: unknown budgets stay unset and
139 /// nothing is fabricated. `updated_at` (epoch seconds) provides the terminal
140 /// timestamp because the job manager records no separate end time. The
141 /// free-form job detail is intentionally omitted because this owner does not
142 /// classify it as safe for a cross-surface read model.
143 #[must_use]
144 pub fn job_record_to_agent_run(
145 record: &JobRecord,
146 ) -> codewhale_protocol::agent_run::AgentRunSnapshot {
147 use codewhale_protocol::agent_run::{
148 AgentRunSnapshot, BudgetSummary, RunSource, RunState, TerminalOutcome, TerminalSummary,
149 };
150
151 let (state, terminal) = match record.status {
152 JobStatus::Queued => (RunState::Queued, None),
153 JobStatus::Running => (RunState::Running, None),
154 JobStatus::Paused => (RunState::Paused, None),
155 JobStatus::Completed | JobStatus::Failed | JobStatus::Cancelled => {
156 let outcome = match record.status {
157 JobStatus::Completed => TerminalOutcome::Completed,
158 JobStatus::Failed => TerminalOutcome::Failed,
159 _ => TerminalOutcome::Cancelled,
160 };
161 (
162 RunState::Terminal,
163 Some(TerminalSummary {
164 outcome,
165 ended_at_ms: record.updated_at.checked_mul(1000),
166 detail: None,
167 }),
168 )
169 }
170 };
171
172 AgentRunSnapshot {
173 run_id: record.id.clone(),
174 parent: None,
175 source: RunSource::CoreJob,
176 state,
177 budget: BudgetSummary::default(),
178 terminal,
179 refs: Vec::new(),
180 }
181 }
182
183 /// Manages background jobs with retry logic and persistence.
184 #[derive(Debug, Default)]
185 pub struct JobManager {
186 jobs: HashMap<String, JobRecord>,
187 }
188
189 impl JobManager {
190 fn now_ts() -> i64 {
191 chrono::Utc::now().timestamp()
192 }
193
194 fn deterministic_backoff_ms(retry: &JobRetryMetadata) -> u64 {
195 if retry.attempt == 0 {
196 return 0;
197 }
198 let exponent = retry.attempt.saturating_sub(1).min(20);
199 let multiplier = 1u64.checked_shl(exponent).unwrap_or(u64::MAX);
200 retry.backoff_base_ms.saturating_mul(multiplier)
201 }
202
203 fn clear_retry_schedule(retry: &mut JobRetryMetadata) {
204 retry.next_backoff_ms = 0;
205 retry.next_retry_at = None;
206 }
207
208 fn push_history(job: &mut JobRecord, phase: &str) {
209 job.history.push(JobHistoryEntry {
210 at: job.updated_at,
211 phase: phase.to_string(),
212 status: job.status,
213 progress: job.progress,
214 detail: job.detail.clone(),
215 retry: job.retry.clone(),
216 });
217 if job.history.len() > MAX_JOB_HISTORY_ENTRIES {
218 let to_drain = job.history.len() - MAX_JOB_HISTORY_ENTRIES;
219 job.history.drain(0..to_drain);
220 }
221 }
222
223 fn parse_persisted_detail(raw: Option<&str>) -> Option<PersistedJobDetail> {
224 let raw = raw?;
225 let parsed: Value = serde_json::from_str(raw).ok()?;
226 let status = parsed
227 .get("status")
228 .and_then(Value::as_str)
229 .and_then(job_status_from_str)?;
230 let detail = parsed.get("detail").and_then(json_optional_string);
231 let retry = parse_retry_metadata(parsed.get("retry"));
232 let history = parsed
233 .get("history")
234 .and_then(Value::as_array)
235 .map(|items| {
236 items
237 .iter()
238 .filter_map(parse_history_entry)
239 .collect::<Vec<_>>()
240 })
241 .unwrap_or_default();
242 Some(PersistedJobDetail {
243 status,
244 detail,
245 retry,
246 history,
247 })
248 }
249
250 fn encode_persisted_detail(job: &JobRecord) -> Result<Option<String>> {
251 let encoded = json!({
252 "schema_version": JOB_DETAIL_SCHEMA_VERSION,
253 "status": job_status_to_str(job.status),
254 "detail": job.detail.clone(),
255 "retry": job_retry_to_value(&job.retry),
256 "history": job.history.iter().map(job_history_to_value).collect::<Vec<_>>()
257 })
258 .to_string();
259 Ok(Some(encoded))
260 }
261
262 /// Enqueues a new job and returns its record.
263 pub fn enqueue(&mut self, name: impl Into<String>) -> JobRecord {
264 let now = Self::now_ts();
265 let id = format!("job-{}", Uuid::new_v4());
266 let mut job = JobRecord {
267 id: id.clone(),
268 name: name.into(),
269 status: JobStatus::Queued,
270 progress: Some(0),
271 detail: None,
272 retry: JobRetryMetadata::default(),
273 history: Vec::new(),
274 created_at: now,
275 updated_at: now,
276 };
277 Self::push_history(&mut job, "created");
278 self.jobs.insert(id, job.clone());
279 job
280 }
281
282 /// Transitions a job to running and clears its retry schedule.
283 pub fn set_running(&mut self, id: &str) {
284 if let Some(job) = self.jobs.get_mut(id) {
285 job.status = JobStatus::Running;
286 Self::clear_retry_schedule(&mut job.retry);
287 job.updated_at = Self::now_ts();
288 Self::push_history(job, "running");
289 }
290 }
291
292 /// Updates a job's progress (clamped to 100) and optional detail message.
293 pub fn update_progress(&mut self, id: &str, progress: u8, detail: Option<String>) {
294 if let Some(job) = self.jobs.get_mut(id) {
295 job.progress = Some(progress.min(100));
296 job.detail = detail;
297 job.updated_at = Self::now_ts();
298 Self::push_history(job, "progress_updated");
299 }
300 }
301
302 /// Marks a job as completed with 100% progress and clears its retry schedule.
303 pub fn complete(&mut self, id: &str) {
304 if let Some(job) = self.jobs.get_mut(id) {
305 job.status = JobStatus::Completed;
306 job.progress = Some(100);
307 Self::clear_retry_schedule(&mut job.retry);
308 job.updated_at = Self::now_ts();
309 Self::push_history(job, "completed");
310 }
311 }
312
313 /// Marks a job as failed and schedules a retry if attempts remain.
314 pub fn fail(&mut self, id: &str, detail: impl Into<String>) {
315 if let Some(job) = self.jobs.get_mut(id) {
316 let now = Self::now_ts();
317 job.status = JobStatus::Failed;
318 job.detail = Some(detail.into());
319 if job.retry.attempt < job.retry.max_attempts {
320 job.retry.attempt += 1;
321 job.retry.next_backoff_ms = Self::deterministic_backoff_ms(&job.retry);
322 let delay_secs = ((job.retry.next_backoff_ms.saturating_add(999)) / 1000)
323 .min(i64::MAX as u64) as i64;
324 job.retry.next_retry_at = Some(now.saturating_add(delay_secs));
325 } else {
326 Self::clear_retry_schedule(&mut job.retry);
327 }
328 job.updated_at = now;
329 Self::push_history(job, "failed");
330 }
331 }
332
333 /// Cancels a job and clears any pending retry schedule.
334 pub fn cancel(&mut self, id: &str) {
335 if let Some(job) = self.jobs.get_mut(id) {
336 job.status = JobStatus::Cancelled;
337 Self::clear_retry_schedule(&mut job.retry);
338 job.updated_at = Self::now_ts();
339 Self::push_history(job, "cancelled");
340 }
341 }
342
343 /// Pauses a job, optionally updating its detail message.
344 pub fn pause(&mut self, id: &str, detail: Option<String>) {
345 if let Some(job) = self.jobs.get_mut(id) {
346 job.status = JobStatus::Paused;
347 if detail.is_some() {
348 job.detail = detail;
349 }
350 job.updated_at = Self::now_ts();
351 Self::push_history(job, "paused");
352 }
353 }
354
355 /// Resumes a paused or failed job back to running status.
356 pub fn resume(&mut self, id: &str, detail: Option<String>) {
357 if let Some(job) = self.jobs.get_mut(id) {
358 job.status = JobStatus::Running;
359 if detail.is_some() {
360 job.detail = detail;
361 }
362 Self::clear_retry_schedule(&mut job.retry);
363 job.updated_at = Self::now_ts();
364 Self::push_history(job, "resumed");
365 }
366 }
367
368 /// Returns all jobs sorted by most recently updated first.
369 pub fn list(&self) -> Vec<JobRecord> {
370 let mut out = self.jobs.values().cloned().collect::<Vec<_>>();
371 out.sort_by_key(|job| std::cmp::Reverse(job.updated_at));
372 out
373 }
374
375 /// Returns the history entries for a job, or an empty vec if not found.
376 pub fn history(&self, id: &str) -> Vec<JobHistoryEntry> {
377 self.jobs
378 .get(id)
379 .map(|job| job.history.clone())
380 .unwrap_or_default()
381 }
382
383 /// Resets queued or running jobs back to queued on application resume.
384 pub fn resume_pending(&mut self) -> Vec<JobRecord> {
385 let mut resumed = Vec::new();
386 for job in self.jobs.values_mut() {
387 if matches!(job.status, JobStatus::Queued | JobStatus::Running) {
388 job.status = JobStatus::Queued;
389 job.updated_at = Self::now_ts();
390 Self::push_history(job, "queued_after_resume");
391 resumed.push(job.clone());
392 }
393 }
394 resumed
395 }
396
397 /// Loads jobs from the state store, deserializing extended detail when available.
398 pub fn load_from_store(&mut self, store: &StateStore) -> Result<()> {
399 let persisted = store.list_jobs(Some(500))?;
400 for job in persisted {
401 let fallback_status = job_state_status_to_runtime(job.status);
402 let parsed = Self::parse_persisted_detail(job.detail.as_deref());
403 let (status, detail, retry, history) = if let Some(detail_state) = parsed {
404 (
405 detail_state.status,
406 detail_state.detail,
407 detail_state.retry,
408 detail_state.history,
409 )
410 } else {
411 (
412 fallback_status,
413 job.detail,
414 JobRetryMetadata::default(),
415 Vec::new(),
416 )
417 };
418 self.jobs.insert(
419 job.id.clone(),
420 JobRecord {
421 id: job.id,
422 name: job.name,
423 status,
424 progress: job.progress,
425 detail,
426 retry,
427 history,
428 created_at: job.created_at,
429 updated_at: job.updated_at,
430 },
431 );
432 }
433 Ok(())
434 }
435
436 /// Persists a single job's current state to the state store.
437 pub fn persist_job(&self, store: &StateStore, id: &str) -> Result<()> {
438 let Some(job) = self.jobs.get(id) else {
439 return Ok(());
440 };
441 let encoded_detail = Self::encode_persisted_detail(job)?;
442 store.upsert_job(&JobStateRecord {
443 id: job.id.clone(),
444 name: job.name.clone(),
445 status: runtime_status_to_job_state(job.status),
446 progress: job.progress,
447 detail: encoded_detail,
448 created_at: job.created_at,
449 updated_at: job.updated_at,
450 })
451 }
452
453 /// Persists all in-memory jobs to the state store.
454 pub fn persist_all(&self, store: &StateStore) -> Result<()> {
455 for id in self.jobs.keys() {
456 self.persist_job(store, id)?;
457 }
458 Ok(())
459 }
460 }
461
462 /// Compatibility configuration, hooks and job bookkeeping.
463 ///
464 /// Conversation execution and durable history belong to the held canonical
465 /// RuntimeThreadManager/Engine. Legacy State data is retained for read-only
466 /// recovery and bound alias publication, never as a second transcript writer.
467 pub struct Runtime {
468 pub config: ConfigToml,
469 state: StateStore,
470 pub hooks: HookDispatcher,
471 pub jobs: JobManager,
472 }
473
474 impl Runtime {
475 pub fn new(config: ConfigToml, state: StateStore, hooks: HookDispatcher) -> Self {
476 let mut jobs = JobManager::default();
477 if let Err(e) = jobs.load_from_store(&state) {
478 tracing::warn!("Failed to load job store, starting with empty job list: {e}");
479 }
480 Self {
481 config,
482 state,
483 hooks,
484 jobs,
485 }
486 }
487
488 /// Existing compatibility data and receipts; no transcript controller.
489 pub fn state_store(&self) -> &StateStore {
490 &self.state
491 }
492
493 pub fn update_config(&mut self, config: ConfigToml) {
494 self.config = config;
495 }
496
497 /// Returns the current application status including all jobs and their history.
498 pub fn app_status(&self) -> AppResponse {
499 let jobs = self.jobs.list();
500 let events = jobs
501 .iter()
502 .flat_map(|job| {
503 job.history.iter().map(|entry| EventFrame::ResponseDelta {
504 response_id: job.id.clone(),
505 delta: json!({
506 "kind": "job_transition",
507 "job_id": job.id.clone(),
508 "phase": entry.phase.clone(),
509 "status": job_status_to_str(entry.status),
510 "progress": entry.progress,
511 "detail": entry.detail.clone(),
512 "retry": job_retry_to_value(&entry.retry),
513 "at": entry.at
514 })
515 .to_string(),
516 channel: ResponseChannel::Text,
517 })
518 })
519 .collect::<Vec<_>>();
520 AppResponse {
521 ok: true,
522 data: json!({
523 "jobs": jobs.into_iter().map(|job| {
524 json!({
525 "id": job.id,
526 "name": job.name,
527 "status": job_status_to_str(job.status),
528 "progress": job.progress,
529 "detail": job.detail,
530 "retry": job_retry_to_value(&job.retry),
531 "history": job.history.iter().map(job_history_to_value).collect::<Vec<_>>()
532 })
533 }).collect::<Vec<_>>()
534 }),
535 events,
536 }
537 }
538
539 /// Returns the default model provider from the resolved configuration.
540 pub fn provider_default(&self) -> ProviderKind {
541 self.config.provider
542 }
543 }
544
545 /// How many leading entries of `offered` repeat the end of `persisted`: the
546 /// longest suffix of `persisted` that is also a prefix of `offered`.
547 /// Longest persisted suffix repeated by an offered history prefix. Linear
548 /// work preserves genuine repeated messages without an unbounded nested scan.
549 pub fn persisted_overlap<T: PartialEq>(persisted: &[T], offered: &[T]) -> usize {
550 if offered.is_empty() {
551 return 0;
552 }
553 let mut failure = vec![0; offered.len()];
554 for index in 1..offered.len() {
555 let mut matched = failure[index - 1];
556 while matched > 0 && offered[index] != offered[matched] {
557 matched = failure[matched - 1];
558 }
559 if offered[index] == offered[matched] {
560 matched += 1;
561 }
562 failure[index] = matched;
563 }
564 let mut matched = 0;
565 for item in persisted {
566 if matched == offered.len() {
567 matched = failure[matched - 1];
568 }
569 while matched > 0 && *item != offered[matched] {
570 matched = failure[matched - 1];
571 }
572 if *item == offered[matched] {
573 matched += 1;
574 }
575 }
576 matched
577 }
578
579 fn json_optional_string(value: &Value) -> Option<String> {
580 if value.is_null() {
581 None
582 } else {
583 value.as_str().map(ToString::to_string)
584 }
585 }
586
587 fn parse_retry_metadata(value: Option<&Value>) -> JobRetryMetadata {
588 let Some(value) = value else {
589 return JobRetryMetadata::default();
590 };
591 JobRetryMetadata {
592 attempt: value
593 .get("attempt")
594 .and_then(Value::as_u64)
595 .unwrap_or(0)
596 .min(u32::MAX as u64) as u32,
597 max_attempts: value
598 .get("max_attempts")
599 .and_then(Value::as_u64)
600 .unwrap_or(DEFAULT_JOB_MAX_ATTEMPTS as u64)
601 .min(u32::MAX as u64) as u32,
602 backoff_base_ms: value
603 .get("backoff_base_ms")
604 .and_then(Value::as_u64)
605 .unwrap_or(DEFAULT_JOB_BACKOFF_BASE_MS),
606 next_backoff_ms: value
607 .get("next_backoff_ms")
608 .and_then(Value::as_u64)
609 .unwrap_or(0),
610 next_retry_at: value.get("next_retry_at").and_then(Value::as_i64),
611 }
612 }
613
614 fn parse_history_entry(value: &Value) -> Option<JobHistoryEntry> {
615 let status = value
616 .get("status")
617 .and_then(Value::as_str)
618 .and_then(job_status_from_str)?;
619 Some(JobHistoryEntry {
620 at: value.get("at").and_then(Value::as_i64).unwrap_or(0),
621 phase: value
622 .get("phase")
623 .and_then(Value::as_str)
624 .unwrap_or("unknown")
625 .to_string(),
626 status,
627 progress: value
628 .get("progress")
629 .and_then(Value::as_u64)
630 .map(|v| v.min(u8::MAX as u64) as u8),
631 detail: value.get("detail").and_then(json_optional_string),
632 retry: parse_retry_metadata(value.get("retry")),
633 })
634 }
635
636 fn job_status_to_str(status: JobStatus) -> &'static str {
637 match status {
638 JobStatus::Queued => "queued",
639 JobStatus::Running => "running",
640 JobStatus::Paused => "paused",
641 JobStatus::Completed => "completed",
642 JobStatus::Failed => "failed",
643 JobStatus::Cancelled => "cancelled",
644 }
645 }
646
647 fn job_status_from_str(value: &str) -> Option<JobStatus> {
648 match value {
649 "queued" => Some(JobStatus::Queued),
650 "running" => Some(JobStatus::Running),
651 "paused" => Some(JobStatus::Paused),
652 "completed" => Some(JobStatus::Completed),
653 "failed" => Some(JobStatus::Failed),
654 "cancelled" => Some(JobStatus::Cancelled),
655 _ => None,
656 }
657 }
658
659 fn job_retry_to_value(retry: &JobRetryMetadata) -> Value {
660 json!({
661 "attempt": retry.attempt,
662 "max_attempts": retry.max_attempts,
663 "backoff_base_ms": retry.backoff_base_ms,
664 "next_backoff_ms": retry.next_backoff_ms,
665 "next_retry_at": retry.next_retry_at
666 })
667 }
668
669 fn job_history_to_value(entry: &JobHistoryEntry) -> Value {
670 json!({
671 "at": entry.at,
672 "phase": entry.phase.clone(),
673 "status": job_status_to_str(entry.status),
674 "progress": entry.progress,
675 "detail": entry.detail.clone(),
676 "retry": job_retry_to_value(&entry.retry)
677 })
678 }
679
680 fn runtime_status_to_job_state(status: JobStatus) -> JobStateStatus {
681 match status {
682 JobStatus::Queued => JobStateStatus::Queued,
683 JobStatus::Running => JobStateStatus::Running,
684 JobStatus::Paused => JobStateStatus::Paused,
685 JobStatus::Completed => JobStateStatus::Completed,
686 JobStatus::Failed => JobStateStatus::Failed,
687 JobStatus::Cancelled => JobStateStatus::Cancelled,
688 }
689 }
690
691 fn job_state_status_to_runtime(status: JobStateStatus) -> JobStatus {
692 match status {
693 JobStateStatus::Queued => JobStatus::Queued,
694 JobStateStatus::Running => JobStatus::Running,
695 JobStateStatus::Paused => JobStatus::Paused,
696 JobStateStatus::Completed => JobStatus::Completed,
697 JobStateStatus::Failed => JobStatus::Failed,
698 JobStateStatus::Cancelled => JobStatus::Cancelled,
699 }
700 }
701
702 #[cfg(test)]
703 mod tests {
704 use super::*;
705
706 fn temp_core_state(name: &str) -> StateStore {
707 let dir =
708 std::env::temp_dir().join(format!("codewhale-core-{name}-{}", Uuid::new_v4().simple()));
709 std::fs::create_dir_all(&dir).expect("create temp state dir");
710 StateStore::open(Some(dir.join("state.db"))).expect("open state store")
711 }
712
713 // ── JobManager: lifecycle ──────────────────────────────────────────
714
715 #[test]
716 fn enqueue_creates_queued_job_with_zero_progress() {
717 let mut jm = JobManager::default();
718 let job = jm.enqueue("build");
719 assert_eq!(job.name, "build");
720 assert_eq!(job.status, JobStatus::Queued);
721 assert_eq!(job.progress, Some(0));
722 assert!(job.detail.is_none());
723 assert_eq!(job.history.len(), 1);
724 assert_eq!(job.history[0].phase, "created");
725 }
726
727 #[test]
728 fn set_running_transitions_from_queued() {
729 let mut jm = JobManager::default();
730 let job = jm.enqueue("deploy");
731 let id = job.id.clone();
732 jm.set_running(&id);
733 let jobs = jm.list();
734 let updated = jobs.iter().find(|j| j.id == id).unwrap();
735 assert_eq!(updated.status, JobStatus::Running);
736 assert_eq!(updated.history.last().unwrap().phase, "running");
737 }
738
739 #[test]
740 fn update_progress_clamps_to_100() {
741 let mut jm = JobManager::default();
742 let job = jm.enqueue("task");
743 let id = job.id.clone();
744 jm.update_progress(&id, 150, Some("over".to_string()));
745 let jobs = jm.list();
746 let updated = jobs.iter().find(|j| j.id == id).unwrap();
747 assert_eq!(updated.progress, Some(100));
748 }
749
750 #[test]
751 fn complete_sets_progress_to_100() {
752 let mut jm = JobManager::default();
753 let job = jm.enqueue("task");
754 let id = job.id.clone();
755 jm.set_running(&id);
756 jm.complete(&id);
757 let jobs = jm.list();
758 let updated = jobs.iter().find(|j| j.id == id).unwrap();
759 assert_eq!(updated.status, JobStatus::Completed);
760 assert_eq!(updated.progress, Some(100));
761 }
762
763 #[test]
764 fn fail_increments_attempt_and_sets_backoff() {
765 let mut jm = JobManager::default();
766 let job = jm.enqueue("fragile");
767 let id = job.id.clone();
768 jm.set_running(&id);
769 jm.fail(&id, "crashed");
770 let jobs = jm.list();
771 let updated = jobs.iter().find(|j| j.id == id).unwrap();
772 assert_eq!(updated.status, JobStatus::Failed);
773 assert_eq!(updated.retry.attempt, 1);
774 assert!(updated.retry.next_backoff_ms > 0);
775 assert!(updated.retry.next_retry_at.is_some());
776 assert_eq!(updated.detail.as_deref(), Some("crashed"));
777 }
778
779 #[test]
780 fn fail_clears_retry_after_max_attempts() {
781 let mut jm = JobManager::default();
782 let job = jm.enqueue("fragile");
783 let id = job.id.clone();
784 for _ in 0..=DEFAULT_JOB_MAX_ATTEMPTS {
785 jm.set_running(&id);
786 jm.fail(&id, "boom");
787 }
788 let jobs = jm.list();
789 let updated = jobs.iter().find(|j| j.id == id).unwrap();
790 assert_eq!(updated.retry.attempt, DEFAULT_JOB_MAX_ATTEMPTS);
791 assert_eq!(updated.retry.next_backoff_ms, 0);
792 assert!(updated.retry.next_retry_at.is_none());
793 }
794
795 #[test]
796 fn cancel_sets_status_and_clears_retry() {
797 let mut jm = JobManager::default();
798 let job = jm.enqueue("task");
799 let id = job.id.clone();
800 jm.cancel(&id);
801 let jobs = jm.list();
802 let updated = jobs.iter().find(|j| j.id == id).unwrap();
803 assert_eq!(updated.status, JobStatus::Cancelled);
804 assert_eq!(updated.retry.next_backoff_ms, 0);
805 }
806
807 #[test]
808 fn pause_and_resume_round_trip() {
809 let mut jm = JobManager::default();
810 let job = jm.enqueue("task");
811 let id = job.id.clone();
812 jm.set_running(&id);
813 jm.pause(&id, Some("waiting".to_string()));
814 let jobs = jm.list();
815 let paused = jobs.iter().find(|j| j.id == id).unwrap();
816 assert_eq!(paused.status, JobStatus::Paused);
817 assert_eq!(paused.detail.as_deref(), Some("waiting"));
818
819 jm.resume(&id, None);
820 let jobs = jm.list();
821 let resumed = jobs.iter().find(|j| j.id == id).unwrap();
822 assert_eq!(resumed.status, JobStatus::Running);
823 assert_eq!(resumed.history.last().unwrap().phase, "resumed");
824 }
825
826 #[test]
827 fn list_returns_jobs_sorted_by_updated_at_desc() {
828 let mut jm = JobManager::default();
829 jm.enqueue("first");
830 jm.enqueue("second");
831 jm.enqueue("third");
832 let jobs = jm.list();
833 assert_eq!(jobs.len(), 3);
834 for window in jobs.windows(2) {
835 assert!(window[0].updated_at >= window[1].updated_at);
836 }
837 }
838
839 #[test]
840 fn history_returns_entries_for_existing_job() {
841 let mut jm = JobManager::default();
842 let job = jm.enqueue("task");
843 let id = job.id.clone();
844 jm.set_running(&id);
845 jm.complete(&id);
846 let history = jm.history(&id);
847 assert_eq!(history.len(), 3); // created, running, completed
848 assert_eq!(history[0].phase, "created");
849 assert_eq!(history[1].phase, "running");
850 assert_eq!(history[2].phase, "completed");
851 }
852
853 #[test]
854 fn history_returns_empty_for_unknown_job() {
855 let jm = JobManager::default();
856 assert!(jm.history("nonexistent").is_empty());
857 }
858
859 #[test]
860 fn resume_pending_requeues_running_and_queued() {
861 let mut jm = JobManager::default();
862 let _j1 = jm.enqueue("queued_task");
863 let j2 = jm.enqueue("running_task");
864 let j3 = jm.enqueue("completed_task");
865 let id2 = j2.id.clone();
866 let id3 = j3.id.clone();
867 jm.set_running(&id2);
868 jm.set_running(&id3);
869 jm.complete(&id3);
870
871 let resumed = jm.resume_pending();
872 assert_eq!(resumed.len(), 2);
873 for job in &resumed {
874 assert_eq!(job.status, JobStatus::Queued);
875 }
876 }
877
878 // ── JobManager: backoff ────────────────────────────────────────────
879
880 #[test]
881 fn deterministic_backoff_zero_on_first_attempt() {
882 let retry = JobRetryMetadata {
883 attempt: 0,
884 ..Default::default()
885 };
886 assert_eq!(JobManager::deterministic_backoff_ms(&retry), 0);
887 }
888
889 #[test]
890 fn deterministic_backoff_exponential_growth() {
891 let base = DEFAULT_JOB_BACKOFF_BASE_MS;
892 for attempt in 1..=5 {
893 let retry = JobRetryMetadata {
894 attempt,
895 backoff_base_ms: base,
896 ..Default::default()
897 };
898 let expected = base * 2u64.pow(attempt.saturating_sub(1).min(20));
899 assert_eq!(
900 JobManager::deterministic_backoff_ms(&retry),
901 expected,
902 "attempt {attempt}"
903 );
904 }
905 }
906
907 #[test]
908 fn deterministic_backoff_saturates_at_high_exponent() {
909 let retry = JobRetryMetadata {
910 attempt: 63,
911 backoff_base_ms: 1000,
912 ..Default::default()
913 };
914 // Should not panic; result saturates
915 let _ = JobManager::deterministic_backoff_ms(&retry);
916 }
917
918 // ── JobManager: history truncation ─────────────────────────────────
919
920 #[test]
921 fn push_history_truncates_beyond_max() {
922 let mut jm = JobManager::default();
923 let job = jm.enqueue("task");
924 let id = job.id.clone();
925 // Generate more history entries than the limit
926 for i in 0..(MAX_JOB_HISTORY_ENTRIES + 20) {
927 jm.update_progress(&id, (i % 100) as u8, Some(format!("step {i}")));
928 }
929 let history = jm.history(&id);
930 assert_eq!(history.len(), MAX_JOB_HISTORY_ENTRIES);
931 }
932
933 // ── JobManager: persistence encoding/parsing ───────────────────────
934
935 #[test]
936 fn encode_and_parse_persisted_detail_round_trip() {
937 let mut jm = JobManager::default();
938 let job = jm.enqueue("task");
939 let id = job.id.clone();
940 jm.set_running(&id);
941 jm.fail(&id, "oops");
942 let job = jm.list().into_iter().find(|j| j.id == id).unwrap();
943
944 let encoded = JobManager::encode_persisted_detail(&job).unwrap().unwrap();
945 let parsed = JobManager::parse_persisted_detail(Some(&encoded)).unwrap();
946
947 assert_eq!(parsed.status, job.status);
948 assert_eq!(parsed.detail, job.detail);
949 assert_eq!(parsed.retry.attempt, job.retry.attempt);
950 assert_eq!(parsed.history.len(), job.history.len());
951 }
952
953 #[test]
954 fn parse_persisted_detail_returns_none_for_none_input() {
955 assert!(JobManager::parse_persisted_detail(None).is_none());
956 }
957
958 #[test]
959 fn parse_persisted_detail_returns_none_for_invalid_json() {
960 assert!(JobManager::parse_persisted_detail(Some("not json")).is_none());
961 }
962
963 // ── Helper functions ───────────────────────────────────────────────
964
965 #[test]
966 fn job_status_round_trip_str() {
967 let statuses = [
968 JobStatus::Queued,
969 JobStatus::Running,
970 JobStatus::Paused,
971 JobStatus::Completed,
972 JobStatus::Failed,
973 JobStatus::Cancelled,
974 ];
975 for status in &statuses {
976 let s = job_status_to_str(*status);
977 let parsed = job_status_from_str(s);
978 assert_eq!(parsed, Some(*status), "round-trip failed for {s:?}");
979 }
980 }
981
982 #[test]
983 fn job_status_from_str_returns_none_for_unknown() {
984 assert_eq!(job_status_from_str("unknown"), None);
985 assert_eq!(job_status_from_str(""), None);
986 }
987
988 #[test]
989 fn runtime_status_to_job_state_maps_correctly() {
990 assert_eq!(
991 runtime_status_to_job_state(JobStatus::Queued),
992 JobStateStatus::Queued
993 );
994 assert_eq!(
995 runtime_status_to_job_state(JobStatus::Running),
996 JobStateStatus::Running
997 );
998 assert_eq!(
999 runtime_status_to_job_state(JobStatus::Paused),
1000 JobStateStatus::Paused
1001 );
1002 assert_eq!(
1003 runtime_status_to_job_state(JobStatus::Completed),
1004 JobStateStatus::Completed
1005 );
1006 assert_eq!(
1007 runtime_status_to_job_state(JobStatus::Failed),
1008 JobStateStatus::Failed
1009 );
1010 assert_eq!(
1011 runtime_status_to_job_state(JobStatus::Cancelled),
1012 JobStateStatus::Cancelled
1013 );
1014 }
1015
1016 #[test]
1017 fn job_state_status_to_runtime_maps_correctly() {
1018 assert_eq!(
1019 job_state_status_to_runtime(JobStateStatus::Queued),
1020 JobStatus::Queued
1021 );
1022 assert_eq!(
1023 job_state_status_to_runtime(JobStateStatus::Running),
1024 JobStatus::Running
1025 );
1026 assert_eq!(
1027 job_state_status_to_runtime(JobStateStatus::Paused),
1028 JobStatus::Paused
1029 );
1030 assert_eq!(
1031 job_state_status_to_runtime(JobStateStatus::Completed),
1032 JobStatus::Completed
1033 );
1034 assert_eq!(
1035 job_state_status_to_runtime(JobStateStatus::Failed),
1036 JobStatus::Failed
1037 );
1038 assert_eq!(
1039 job_state_status_to_runtime(JobStateStatus::Cancelled),
1040 JobStatus::Cancelled
1041 );
1042 }
1043
1044 #[test]
1045 fn json_optional_string_handles_null() {
1046 assert!(json_optional_string(&Value::Null).is_none());
1047 }
1048
1049 #[test]
1050 fn json_optional_string_handles_string() {
1051 assert_eq!(
1052 json_optional_string(&Value::String("hello".to_string())),
1053 Some("hello".to_string())
1054 );
1055 }
1056
1057 #[test]
1058 fn json_optional_string_handles_non_string() {
1059 assert!(json_optional_string(&json!(42)).is_none());
1060 }
1061
1062 #[test]
1063 fn parse_retry_metadata_returns_default_for_none() {
1064 let retry = parse_retry_metadata(None);
1065 assert_eq!(retry.attempt, 0);
1066 assert_eq!(retry.max_attempts, DEFAULT_JOB_MAX_ATTEMPTS);
1067 assert_eq!(retry.backoff_base_ms, DEFAULT_JOB_BACKOFF_BASE_MS);
1068 }
1069
1070 #[test]
1071 fn parse_retry_metadata_parses_fields() {
1072 let value = json!({
1073 "attempt": 2,
1074 "max_attempts": 5,
1075 "backoff_base_ms": 1000,
1076 "next_backoff_ms": 2000,
1077 "next_retry_at": 1234567890i64
1078 });
1079 let retry = parse_retry_metadata(Some(&value));
1080 assert_eq!(retry.attempt, 2);
1081 assert_eq!(retry.max_attempts, 5);
1082 assert_eq!(retry.backoff_base_ms, 1000);
1083 assert_eq!(retry.next_backoff_ms, 2000);
1084 assert_eq!(retry.next_retry_at, Some(1234567890));
1085 }
1086
1087 #[test]
1088 fn parse_history_entry_returns_none_without_status() {
1089 let value = json!({"at": 1, "phase": "test"});
1090 assert!(parse_history_entry(&value).is_none());
1091 }
1092
1093 #[test]
1094 fn parse_history_entry_parses_valid_entry() {
1095 let value = json!({
1096 "at": 100,
1097 "phase": "running",
1098 "status": "running",
1099 "progress": 50,
1100 "detail": "working",
1101 "retry": {"attempt": 0, "max_attempts": 3, "backoff_base_ms": 500}
1102 });
1103 let entry = parse_history_entry(&value).unwrap();
1104 assert_eq!(entry.at, 100);
1105 assert_eq!(entry.phase, "running");
1106 assert_eq!(entry.status, JobStatus::Running);
1107 assert_eq!(entry.progress, Some(50));
1108 assert_eq!(entry.detail.as_deref(), Some("working"));
1109 }
1110
1111 #[test]
1112 fn paused_job_persists_as_paused_not_running() {
1113 let store = temp_core_state("paused-persist");
1114 let mut jm = JobManager::default();
1115 let job = jm.enqueue("task");
1116 let id = job.id.clone();
1117 jm.set_running(&id);
1118 jm.pause(&id, Some("waiting".to_string()));
1119 jm.persist_job(&store, &id).expect("persist paused job");
1120
1121 let persisted = store.list_jobs(Some(10)).expect("list jobs");
1122 let record = persisted.iter().find(|job| job.id == id).unwrap();
1123 assert_eq!(record.status, JobStateStatus::Paused);
1124
1125 let mut reloaded = JobManager::default();
1126 reloaded.load_from_store(&store).expect("reload jobs");
1127 let jobs = reloaded.list();
1128 let reloaded_job = jobs.iter().find(|job| job.id == id).unwrap();
1129 assert_eq!(reloaded_job.status, JobStatus::Paused);
1130 }
1131
1132 // ── O1: JobRecord → AgentRunSnapshot adapter ────────────────────────
1133
1134 fn sample_job_record(status: JobStatus, detail: Option<&str>) -> JobRecord {
1135 JobRecord {
1136 id: "job-o1-1".to_string(),
1137 name: "sample".to_string(),
1138 status,
1139 progress: None,
1140 detail: detail.map(str::to_string),
1141 retry: JobRetryMetadata {
1142 attempt: 0,
1143 max_attempts: DEFAULT_JOB_MAX_ATTEMPTS,
1144 backoff_base_ms: DEFAULT_JOB_BACKOFF_BASE_MS,
1145 next_backoff_ms: 0,
1146 next_retry_at: None,
1147 },
1148 history: Vec::new(),
1149 created_at: 1_700_000_000,
1150 updated_at: 1_700_000_042,
1151 }
1152 }
1153
1154 #[test]
1155 fn job_record_to_agent_run_maps_non_terminal_states() {
1156 use codewhale_protocol::agent_run::RunState;
1157
1158 for (status, expected) in [
1159 (JobStatus::Queued, RunState::Queued),
1160 (JobStatus::Running, RunState::Running),
1161 (JobStatus::Paused, RunState::Paused),
1162 ] {
1163 let snapshot = job_record_to_agent_run(&sample_job_record(status, None));
1164 assert!(snapshot.is_coherent());
1165 assert_eq!(snapshot.run_id, "job-o1-1");
1166 assert_eq!(snapshot.parent, None);
1167 assert_eq!(
1168 snapshot.source,
1169 codewhale_protocol::agent_run::RunSource::CoreJob
1170 );
1171 assert_eq!(snapshot.state, expected);
1172 assert!(snapshot.terminal.is_none());
1173 assert!(snapshot.refs.is_empty());
1174 assert_eq!(
1175 snapshot.budget,
1176 codewhale_protocol::agent_run::BudgetSummary::default()
1177 );
1178 }
1179 }
1180
1181 #[test]
1182 fn job_record_to_agent_run_maps_terminal_states_without_fabricating_fields() {
1183 use codewhale_protocol::agent_run::{RunState, TerminalOutcome};
1184
1185 let cases = [
1186 (
1187 JobStatus::Completed,
1188 TerminalOutcome::Completed,
1189 Some("done"),
1190 ),
1191 (JobStatus::Failed, TerminalOutcome::Failed, Some("boom")),
1192 (JobStatus::Cancelled, TerminalOutcome::Cancelled, None),
1193 ];
1194
1195 for (status, outcome, detail) in cases {
1196 let snapshot = job_record_to_agent_run(&sample_job_record(status, detail));
1197 assert!(snapshot.is_coherent());
1198 assert_eq!(snapshot.state, RunState::Terminal);
1199 let terminal = snapshot.terminal.expect("terminal summary");
1200 assert_eq!(terminal.outcome, outcome);
1201 assert_eq!(terminal.ended_at_ms, Some(1_700_000_042_000));
1202 assert_eq!(terminal.detail, None);
1203 assert_eq!(
1204 snapshot.budget,
1205 codewhale_protocol::agent_run::BudgetSummary::default()
1206 );
1207 assert!(snapshot.refs.is_empty());
1208 assert_eq!(snapshot.parent, None);
1209 }
1210 }
1211
1212 #[test]
1213 fn job_record_to_agent_run_does_not_export_unclassified_detail() {
1214 let record = sample_job_record(JobStatus::Failed, Some("owner-private diagnostic"));
1215 let snapshot = job_record_to_agent_run(&record);
1216 let terminal = snapshot.terminal.as_ref().expect("terminal summary");
1217 assert_eq!(terminal.detail, None);
1218 let serialized = serde_json::to_string(&snapshot).expect("serialize snapshot");
1219 assert!(!serialized.contains("owner-private diagnostic"));
1220 }
1221
1222 #[test]
1223 fn job_record_to_agent_run_omits_ended_at_on_updated_at_overflow() {
1224 let mut record = sample_job_record(JobStatus::Completed, Some("ok"));
1225 record.updated_at = i64::MAX;
1226 let snapshot = job_record_to_agent_run(&record);
1227 assert!(snapshot.is_coherent());
1228 let terminal = snapshot.terminal.expect("terminal summary");
1229 assert_eq!(terminal.ended_at_ms, None);
1230 }
1231
1232 #[test]
1233 fn runtime_bookkeeping_retains_jobs_and_read_only_legacy_store() {
1234 let store = temp_core_state("read-only-runtime");
1235 let before = store.list_threads(Default::default()).unwrap();
1236 let mut runtime = Runtime::new(ConfigToml::default(), store, HookDispatcher::default());
1237 let job = runtime.jobs.enqueue("retained job");
1238 runtime
1239 .jobs
1240 .persist_job(runtime.state_store(), &job.id)
1241 .unwrap();
1242 runtime.update_config(ConfigToml::default());
1243 assert_eq!(runtime.app_status().data["jobs"][0]["id"], job.id);
1244 assert_eq!(
1245 runtime
1246 .state_store()
1247 .list_threads(Default::default())
1248 .unwrap()
1249 .len(),
1250 before.len()
1251 );
1252 let reopened = Runtime::new(
1253 runtime.config.clone(),
1254 runtime.state_store().clone(),
1255 HookDispatcher::default(),
1256 );
1257 assert_eq!(reopened.jobs.list().len(), 1);
1258 assert!(
1259 reopened
1260 .state_store()
1261 .list_threads(Default::default())
1262 .unwrap()
1263 .is_empty()
1264 );
1265 }
1266 }
1267
1267 lines RUST