返回 CodeWhale
lib.rs
根目录 / crates / state / src / lib.rs
1 //! Historical thread recovery, canonical owner receipts, metadata and jobs.
2 //!
3 //! [`StateStore`] retains legacy conversation graphs as read-only archives.
4 //! A bounded absent-only restore can admit a complete historical graph in one
5 //! SQLite transaction; live transcript writes belong to Runtime's Engine and
6 //! SessionManager. Canonical aliases and durable operations use this same
7 //! SQLite connection. Thread metadata, jobs, dynamic tools and the session
8 //! index retain their existing ownership.
9
10 mod runtime_aliases;
11
12 use std::collections::HashMap;
13 use std::fs::{self, OpenOptions};
14 use std::io::{BufRead, BufReader, Write};
15 use std::path::{Path, PathBuf};
16 use std::sync::{Arc, LazyLock, Mutex, MutexGuard};
17
18 /// Serializes all `session_index.jsonl` read/append/compact/rename operations so
19 /// concurrent `StateStore` clones cannot interleave an append with a compaction
20 /// rename and silently drop entries.
21 static SESSION_INDEX_LOCK: LazyLock<Mutex<()>> = LazyLock::new(|| Mutex::new(()));
22
23 use anyhow::{Context, Result};
24 use chrono::Utc;
25 use codewhale_paths::{CODEWHALE_APP_DIR, LEGACY_APP_DIR, codewhale_home_override};
26 use rusqlite::{Connection, OptionalExtension, params};
27 use serde::{Deserialize, Serialize};
28 use serde_json::Value;
29
30 // Re-export protocol's ThreadStatus so callers in the state crate and
31 // external consumers (e.g. core) can reference a single canonical definition.
32 pub use codewhale_protocol::ThreadStatus;
33
34 /// Indicates how a session was initiated.
35 ///
36 /// Serialized as lowercase snake_case strings.
37 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
38 #[serde(rename_all = "snake_case")]
39 pub enum SessionSource {
40 /// Started by a user interacting with the CLI.
41 Interactive,
42 /// Resumed from a previously persisted session.
43 Resume,
44 /// Created by forking an existing conversation at a specific message.
45 Fork,
46 /// Initiated programmatically via the API.
47 Api,
48 /// Source is unknown or unspecified.
49 Unknown,
50 }
51
52 /// Metadata for a persisted conversation thread.
53 ///
54 /// Each thread represents a single conversation session and stores its
55 /// configuration, git context, and current status.
56 #[derive(Debug, Clone, Serialize, Deserialize)]
57 pub struct ThreadMetadata {
58 /// Unique identifier for this thread.
59 pub id: String,
60 /// Optional filesystem path to the rollout (JSONL transcript) file.
61 pub rollout_path: Option<PathBuf>,
62 /// Short preview or summary of the thread content.
63 pub preview: String,
64 /// Whether this thread is ephemeral (not persisted long-term).
65 pub ephemeral: bool,
66 /// Identifier of the model provider used for this thread (e.g. `"openai"`).
67 pub model_provider: String,
68 /// Unix timestamp (seconds) when the thread was created.
69 pub created_at: i64,
70 /// Unix timestamp (seconds) of the most recent update to the thread.
71 pub updated_at: i64,
72 /// Current lifecycle status of the thread.
73 pub status: ThreadStatus,
74 /// Optional filesystem path associated with the thread working context.
75 pub path: Option<PathBuf>,
76 /// Working directory that was active when the thread was created.
77 pub cwd: PathBuf,
78 /// Version of the CLI that created this thread.
79 pub cli_version: String,
80 /// How this session was initiated.
81 pub source: SessionSource,
82 /// User-assigned display name for the thread.
83 pub name: Option<String>,
84 /// Serialized sandbox policy applied to this thread, if any.
85 pub sandbox_policy: Option<String>,
86 /// Approval mode configured for tool calls in this thread.
87 pub approval_mode: Option<String>,
88 /// Whether the thread has been archived.
89 pub archived: bool,
90 /// Unix timestamp (seconds) when the thread was archived, or `None` if not archived.
91 pub archived_at: Option<i64>,
92 /// Git commit SHA of the working tree when the thread was created.
93 pub git_sha: Option<String>,
94 /// Git branch checked out when the thread was created.
95 pub git_branch: Option<String>,
96 /// URL of the git remote origin, if available.
97 pub git_origin_url: Option<String>,
98 /// Memory mode configured for this thread (e.g. `"local"`, `"remote"`).
99 pub memory_mode: Option<String>,
100 /// ID of the current leaf message in the conversation tree.
101 pub current_leaf_id: Option<i64>,
102 }
103
104 /// A complete historical SQLite thread, restored once into an absent target.
105 ///
106 /// This is recovery/import data, never a live transcript writer. Entry IDs,
107 /// timestamps, all branches and the exact selected leaf are retained.
108 #[derive(Debug, Clone, Serialize, Deserialize)]
109 #[serde(deny_unknown_fields)]
110 pub struct LegacyThreadArchive {
111 pub thread: ThreadMetadata,
112 pub messages: Vec<MessageRecord>,
113 #[serde(default)]
114 pub goal: Option<ThreadGoalRecord>,
115 #[serde(default)]
116 pub checkpoints: Vec<CheckpointRecord>,
117 }
118
119 /// A dynamically registered tool associated with a thread.
120 #[derive(Debug, Clone, Serialize, Deserialize)]
121 pub struct DynamicToolRecord {
122 /// Ordinal position of this tool in the thread tool list.
123 pub position: i64,
124 /// Unique name identifying the tool.
125 pub name: String,
126 /// Human-readable description of what the tool does.
127 pub description: Option<String>,
128 /// JSON Schema describing the tool input parameters.
129 pub input_schema: Value,
130 }
131
132 pub use codewhale_protocol::MessageRecord;
133
134 /// Private historical database fixture for State unit tests.
135 #[cfg(test)]
136 #[derive(Debug, Clone)]
137 struct NewMessage {
138 /// Role of the message sender (e.g. `"user"`, `"history"`).
139 pub role: String,
140 /// Text content of the message.
141 pub content: String,
142 /// Optional structured item payload.
143 pub item: Option<Value>,
144 }
145
146 /// A named checkpoint capturing the state of a thread at a point in time.
147 #[derive(Debug, Clone, Serialize, Deserialize)]
148 pub struct CheckpointRecord {
149 /// ID of the thread this checkpoint belongs to.
150 pub thread_id: String,
151 /// Unique identifier for this checkpoint within its thread.
152 pub checkpoint_id: String,
153 /// Serialized state snapshot stored as a JSON value.
154 pub state: Value,
155 /// Unix timestamp (seconds) when the checkpoint was created or last updated.
156 pub created_at: i64,
157 }
158
159 /// Status of a background job.
160 ///
161 /// Serialized as lowercase snake_case strings.
162 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
163 #[serde(rename_all = "snake_case")]
164 pub enum JobStateStatus {
165 /// Job is waiting to be executed.
166 Queued,
167 /// Job is currently executing.
168 Running,
169 /// Job has been temporarily paused.
170 Paused,
171 /// Job has finished successfully.
172 Completed,
173 /// Job has failed with an error.
174 Failed,
175 /// Job was cancelled before completion.
176 Cancelled,
177 }
178
179 /// Persisted state of a background job.
180 #[derive(Debug, Clone, Serialize, Deserialize)]
181 pub struct JobStateRecord {
182 /// Unique identifier for the job.
183 pub id: String,
184 /// Human-readable name describing the job.
185 pub name: String,
186 /// Current lifecycle status of the job.
187 pub status: JobStateStatus,
188 /// Completion progress as a percentage (0--100), if available.
189 pub progress: Option<u8>,
190 /// Optional detail message providing additional status information.
191 pub detail: Option<String>,
192 /// Unix timestamp (seconds) when the job was created.
193 pub created_at: i64,
194 /// Unix timestamp (seconds) of the most recent status update.
195 pub updated_at: i64,
196 }
197
198 /// Persisted lifecycle status for a thread goal.
199 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
200 #[serde(rename_all = "snake_case")]
201 pub enum ThreadGoalStatus {
202 /// Goal is active and should continue receiving work.
203 Active,
204 /// Goal is paused by the user.
205 Paused,
206 /// Goal is blocked and cannot make meaningful progress.
207 Blocked,
208 /// Goal stopped because account/service usage limits were reached.
209 UsageLimited,
210 /// Goal stopped because its explicit token budget was reached.
211 BudgetLimited,
212 /// Goal has been completed.
213 Complete,
214 }
215
216 /// Persisted goal state attached to a thread.
217 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
218 pub struct ThreadGoalRecord {
219 /// Thread this goal belongs to.
220 pub thread_id: String,
221 /// Stable identifier for this goal revision.
222 pub goal_id: String,
223 /// User-visible objective.
224 pub objective: String,
225 /// Current lifecycle status.
226 pub status: ThreadGoalStatus,
227 /// Optional token budget requested by the user.
228 pub token_budget: Option<i64>,
229 /// Tokens consumed while pursuing the goal.
230 pub tokens_used: i64,
231 /// Elapsed wall-clock work time in seconds.
232 pub time_used_seconds: i64,
233 /// Durable continuation passes dispatched for this objective.
234 pub continuation_count: i64,
235 /// Unix timestamp (seconds) when the goal was created.
236 pub created_at: i64,
237 /// Unix timestamp (seconds) when the goal was last updated.
238 pub updated_at: i64,
239 #[serde(default, skip_serializing_if = "Option::is_none")]
240 pub last_gap_fingerprint: Option<String>,
241 #[serde(default)]
242 pub repeated_gap_count: u32,
243 #[serde(default, skip_serializing_if = "Option::is_none")]
244 pub last_gap_pass: Option<u32>,
245 #[serde(default, skip_serializing_if = "Option::is_none")]
246 pub pause_reason: Option<codewhale_protocol::GoalPauseReason>,
247 }
248
249 /// Filters for listing conversation threads.
250 #[derive(Debug, Clone)]
251 pub struct ThreadListFilters {
252 /// Whether to include archived threads in the results.
253 pub include_archived: bool,
254 /// Maximum number of threads to return. Defaults to 50.
255 pub limit: Option<usize>,
256 }
257
258 impl Default for ThreadListFilters {
259 fn default() -> Self {
260 Self {
261 include_archived: false,
262 limit: Some(50),
263 }
264 }
265 }
266
267 #[derive(Debug, Clone, Serialize, Deserialize)]
268 struct SessionIndexEntry {
269 thread_id: String,
270 thread_name: Option<String>,
271 updated_at: i64,
272 rollout_path: Option<PathBuf>,
273 }
274
275 /// Rewrite the session index once the append-only log grows large enough that
276 /// full-file scans become costly. Lookups already dedupe by thread id, so
277 /// compaction keeps only the latest entry per thread.
278 fn session_index_compact_line_threshold() -> usize {
279 if cfg!(test) { 5 } else { 5_000 }
280 }
281
282 /// The `user_version` the last step of `StateStore::migrate_schema` writes.
283 const SCHEMA_VERSION: u32 = 7;
284
285 /// Persistent storage for conversation threads, messages, checkpoints, and jobs.
286 ///
287 /// Backed by a SQLite database and an append-only JSONL session index file.
288 /// The database schema is automatically initialized and migrated on [`open`](Self::open).
289 #[derive(Debug, Clone)]
290 pub struct StateStore {
291 db_path: PathBuf,
292 session_index_path: PathBuf,
293 // Single long-lived connection shared by all clones. SQLite pragmas are
294 // per-connection, so opening once in `open` and applying them there keeps
295 // every operation consistent without re-opening the database per call.
296 conn: Arc<Mutex<Connection>>,
297 }
298
299 impl StateStore {
300 /// Open (or create) a state store at the given database path.
301 ///
302 /// If `path` is `None`, the default location (`~/.codewhale/state.db`, with
303 /// `~/.deepseek/state.db` as a legacy fallback) is used.
304 /// The database schema is created automatically if it does not exist.
305 pub fn open(path: Option<PathBuf>) -> Result<Self> {
306 let db_path = path.unwrap_or_else(default_state_db_path);
307 let session_index_path = db_path
308 .parent()
309 .unwrap_or_else(|| Path::new("."))
310 .join("session_index.jsonl");
311 if let Some(parent) = db_path.parent() {
312 fs::create_dir_all(parent).with_context(|| {
313 format!("failed to create state directory {}", parent.display())
314 })?;
315 }
316 let conn = Connection::open(&db_path)
317 .with_context(|| format!("failed to open state db {}", db_path.display()))?;
318 Self::configure_connection(&conn, &db_path)?;
319 Self::init_schema(&conn)?;
320 Ok(Self {
321 db_path,
322 session_index_path,
323 conn: Arc::new(Mutex::new(conn)),
324 })
325 }
326
327 /// Apply connection-level SQLite settings that must hold for every open.
328 ///
329 /// Enables WAL so readers and writers from concurrent CodeWhale processes
330 /// do not block each other as aggressively as the default rollback journal,
331 /// and sets a multi-second busy timeout so a second process retries on
332 /// `SQLITE_BUSY` instead of failing immediately (issue #4734).
333 fn configure_connection(conn: &Connection, db_path: &Path) -> Result<()> {
334 // Install our wait policy before touching database-level settings or
335 // schema. Connection::open may currently provide a dependency default,
336 // but StateStore must not rely on that incidental behavior.
337 conn.busy_timeout(std::time::Duration::from_secs(5))
338 .with_context(|| format!("failed to set busy_timeout for {}", db_path.display()))?;
339 conn.pragma_update(None, "foreign_keys", "ON")
340 .with_context(|| format!("failed to enable foreign keys for {}", db_path.display()))?;
341
342 // WAL persists in the database header, so established stores need no
343 // write-like journal transition on every process start. Fresh or
344 // explicitly downgraded stores still transition once, and we verify
345 // SQLite accepted the requested mode instead of silently retaining the
346 // previous one (for example on an unsupported VFS).
347 let journal_mode: String = conn
348 .pragma_query_value(None, "journal_mode", |row| row.get(0))
349 .with_context(|| format!("failed to read journal mode for {}", db_path.display()))?;
350 if !journal_mode.eq_ignore_ascii_case("wal") {
351 let configured_mode: String = conn
352 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
353 .with_context(|| format!("failed to enable WAL for {}", db_path.display()))?;
354 if !configured_mode.eq_ignore_ascii_case("wal") {
355 anyhow::bail!(
356 "failed to enable WAL for {}: SQLite retained journal mode {configured_mode}",
357 db_path.display()
358 );
359 }
360 }
361 Ok(())
362 }
363
364 /// Returns the filesystem path of the underlying SQLite database.
365 pub fn db_path(&self) -> &Path {
366 &self.db_path
367 }
368
369 fn conn(&self) -> Result<MutexGuard<'_, Connection>> {
370 // Poisoning means a panic mid-operation; any open transaction was
371 // rolled back when it dropped, but surface the condition rather than
372 // silently continuing on a connection whose state we can't vouch for.
373 self.conn
374 .lock()
375 .map_err(|_| anyhow::anyhow!("state db connection mutex poisoned"))
376 }
377
378 fn init_schema(conn: &Connection) -> Result<()> {
379 let user_version: u32 = conn.query_row("PRAGMA user_version;", [], |row| row.get(0))?;
380 if user_version >= SCHEMA_VERSION {
381 return Ok(());
382 }
383 // Take the write lock *before* reading the version or probing
384 // columns. Two processes opening a fresh database used to both read
385 // v0 and "column missing" outside any transaction; the loser then
386 // waited out the winner's commit and re-ran its `ADD COLUMN`, failing
387 // the open with "duplicate column name". Under `BEGIN IMMEDIATE` the
388 // loser waits (busy_timeout) and then decides from the committed
389 // schema. The version read above is only a fast path for the
390 // already-current database, which never takes the write lock.
391 conn.execute_batch("BEGIN IMMEDIATE;")
392 .context("failed to lock the state db for its schema migration")?;
393 let migrated = Self::migrate_schema(conn).and_then(|()| {
394 conn.execute_batch("COMMIT;")
395 .context("failed to commit the state db schema migration")
396 });
397 if migrated.is_err() {
398 let _ = conn.execute_batch("ROLLBACK;");
399 }
400 migrated
401 }
402
403 /// Every schema step, run inside the one transaction [`Self::init_schema`]
404 /// holds; the steps themselves never begin or commit.
405 fn migrate_schema(conn: &Connection) -> Result<()> {
406 let mut user_version: u32 = conn.query_row("PRAGMA user_version;", [], |row| row.get(0))?;
407 if user_version == 0 {
408 // Guard each ALTER: a database restored with a v0 header (or
409 // stamped by a racing process that crashed before setting
410 // user_version) can already carry these columns, and an
411 // unguarded ADD COLUMN aborts the whole open with
412 // "duplicate column name".
413 let add_parent_entry_id = if column_exists(conn, "messages", "parent_entry_id")? {
414 ""
415 } else {
416 "ALTER TABLE messages ADD COLUMN parent_entry_id INTEGER NULL;"
417 };
418 let add_current_leaf_id = if column_exists(conn, "threads", "current_leaf_id")? {
419 ""
420 } else {
421 "ALTER TABLE threads ADD COLUMN current_leaf_id INTEGER NULL;"
422 };
423 conn.execute_batch(&format!(
424 r#"
425 CREATE TABLE IF NOT EXISTS threads (
426 id TEXT PRIMARY KEY,
427 rollout_path TEXT,
428 preview TEXT NOT NULL,
429 ephemeral INTEGER NOT NULL,
430 model_provider TEXT NOT NULL,
431 created_at INTEGER NOT NULL,
432 updated_at INTEGER NOT NULL,
433 status TEXT NOT NULL,
434 path TEXT,
435 cwd TEXT NOT NULL,
436 cli_version TEXT NOT NULL,
437 source TEXT NOT NULL,
438 title TEXT,
439 sandbox_policy TEXT,
440 approval_mode TEXT,
441 archived INTEGER NOT NULL DEFAULT 0,
442 archived_at INTEGER,
443 git_sha TEXT,
444 git_branch TEXT,
445 git_origin_url TEXT,
446 memory_mode TEXT
447 );
448 CREATE INDEX IF NOT EXISTS idx_threads_updated_at ON threads(updated_at DESC);
449 CREATE INDEX IF NOT EXISTS idx_threads_archived_at ON threads(archived_at DESC);
450 CREATE INDEX IF NOT EXISTS idx_threads_archived_updated ON threads(archived, updated_at DESC);
451
452 CREATE TABLE IF NOT EXISTS thread_dynamic_tools (
453 thread_id TEXT NOT NULL,
454 position INTEGER NOT NULL,
455 name TEXT NOT NULL,
456 description TEXT,
457 input_schema TEXT NOT NULL,
458 PRIMARY KEY (thread_id, position),
459 FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE
460 );
461
462 CREATE TABLE IF NOT EXISTS messages (
463 id INTEGER PRIMARY KEY AUTOINCREMENT,
464 thread_id TEXT NOT NULL,
465 role TEXT NOT NULL,
466 content TEXT NOT NULL,
467 item_json TEXT,
468 created_at INTEGER NOT NULL,
469 FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE
470 );
471 CREATE INDEX IF NOT EXISTS idx_messages_thread_created_at ON messages(thread_id, created_at ASC);
472
473 CREATE TABLE IF NOT EXISTS checkpoints (
474 thread_id TEXT NOT NULL,
475 checkpoint_id TEXT NOT NULL,
476 state_json TEXT NOT NULL,
477 created_at INTEGER NOT NULL,
478 PRIMARY KEY(thread_id, checkpoint_id),
479 FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE
480 );
481 CREATE INDEX IF NOT EXISTS idx_checkpoints_thread_created_at ON checkpoints(thread_id, created_at DESC);
482
483 CREATE TABLE IF NOT EXISTS jobs (
484 id TEXT PRIMARY KEY,
485 name TEXT NOT NULL,
486 status TEXT NOT NULL,
487 progress INTEGER,
488 detail TEXT,
489 created_at INTEGER NOT NULL,
490 updated_at INTEGER NOT NULL
491 );
492 CREATE INDEX IF NOT EXISTS idx_jobs_updated_at ON jobs(updated_at DESC);
493
494 -- Add parent_entry_id column, and set to last message before current message
495 {add_parent_entry_id}
496 UPDATE messages
497 SET parent_entry_id = (
498 SELECT m2.id
499 FROM messages m2
500 WHERE m2.thread_id = messages.thread_id
501 AND (
502 m2.created_at < messages.created_at
503 OR (
504 m2.created_at = messages.created_at
505 AND m2.id < messages.id
506 )
507 )
508 ORDER BY m2.created_at DESC, m2.id DESC
509 LIMIT 1
510 );
511 CREATE INDEX IF NOT EXISTS idx_messages_parent_entry_id ON messages(parent_entry_id);
512
513 -- Add current_leaf_id column, and set to last message in thread
514 {add_current_leaf_id}
515 UPDATE threads
516 SET current_leaf_id = (
517 SELECT m.id
518 FROM messages m
519 WHERE m.thread_id = threads.id
520 ORDER BY m.id DESC
521 LIMIT 1
522 );
523
524 PRAGMA user_version = 1;
525 "#
526 ))
527 .context("failed to initialize thread schema")?;
528 user_version = 1;
529 }
530 if user_version < 2 {
531 conn.execute_batch(
532 r#"
533 CREATE TABLE IF NOT EXISTS workflow_runs (
534 id TEXT PRIMARY KEY,
535 workflow_id TEXT NOT NULL,
536 goal TEXT NOT NULL,
537 status TEXT NOT NULL,
538 input_hash TEXT,
539 started_at INTEGER NOT NULL,
540 completed_at INTEGER,
541 metadata_json TEXT NOT NULL DEFAULT '{}'
542 );
543 CREATE INDEX IF NOT EXISTS idx_workflow_runs_status_started_at
544 ON workflow_runs(status, started_at DESC);
545 CREATE INDEX IF NOT EXISTS idx_workflow_runs_workflow_started_at
546 ON workflow_runs(workflow_id, started_at DESC);
547
548 CREATE TABLE IF NOT EXISTS branch_runs (
549 id TEXT PRIMARY KEY,
550 workflow_run_id TEXT NOT NULL,
551 branch_id TEXT NOT NULL,
552 node_id TEXT NOT NULL,
553 status TEXT NOT NULL,
554 started_at INTEGER NOT NULL,
555 completed_at INTEGER,
556 result_json TEXT NOT NULL DEFAULT '{}',
557 FOREIGN KEY(workflow_run_id) REFERENCES workflow_runs(id) ON DELETE CASCADE
558 );
559 CREATE INDEX IF NOT EXISTS idx_branch_runs_workflow_run_id
560 ON branch_runs(workflow_run_id);
561 CREATE INDEX IF NOT EXISTS idx_branch_runs_branch_id
562 ON branch_runs(branch_id);
563
564 CREATE TABLE IF NOT EXISTS leaf_runs (
565 id TEXT PRIMARY KEY,
566 workflow_run_id TEXT NOT NULL,
567 branch_run_id TEXT,
568 leaf_id TEXT NOT NULL,
569 task_id TEXT NOT NULL,
570 input_hash TEXT,
571 status TEXT NOT NULL,
572 output_json TEXT NOT NULL DEFAULT '{}',
573 artifacts_json TEXT NOT NULL DEFAULT '[]',
574 started_at INTEGER NOT NULL,
575 completed_at INTEGER,
576 FOREIGN KEY(workflow_run_id) REFERENCES workflow_runs(id) ON DELETE CASCADE,
577 FOREIGN KEY(branch_run_id) REFERENCES branch_runs(id) ON DELETE SET NULL
578 );
579 CREATE INDEX IF NOT EXISTS idx_leaf_runs_workflow_run_id
580 ON leaf_runs(workflow_run_id);
581 CREATE INDEX IF NOT EXISTS idx_leaf_runs_replay_lookup
582 ON leaf_runs(workflow_run_id, leaf_id, input_hash);
583
584 CREATE TABLE IF NOT EXISTS control_node_runs (
585 id TEXT PRIMARY KEY,
586 workflow_run_id TEXT NOT NULL,
587 node_id TEXT NOT NULL,
588 kind TEXT NOT NULL,
589 status TEXT NOT NULL,
590 selected_children_json TEXT NOT NULL DEFAULT '[]',
591 result_json TEXT NOT NULL DEFAULT '{}',
592 started_at INTEGER NOT NULL,
593 completed_at INTEGER,
594 FOREIGN KEY(workflow_run_id) REFERENCES workflow_runs(id) ON DELETE CASCADE
595 );
596 CREATE INDEX IF NOT EXISTS idx_control_node_runs_workflow_run_id
597 ON control_node_runs(workflow_run_id);
598 CREATE INDEX IF NOT EXISTS idx_control_node_runs_node_id
599 ON control_node_runs(node_id);
600
601 CREATE TABLE IF NOT EXISTS teacher_candidates (
602 id TEXT PRIMARY KEY,
603 workflow_run_id TEXT NOT NULL,
604 control_node_run_id TEXT NOT NULL,
605 candidate_id TEXT NOT NULL,
606 branch_run_id TEXT,
607 score REAL,
608 passed INTEGER,
609 rationale_json TEXT NOT NULL DEFAULT '{}',
610 created_at INTEGER NOT NULL,
611 FOREIGN KEY(workflow_run_id) REFERENCES workflow_runs(id) ON DELETE CASCADE,
612 FOREIGN KEY(control_node_run_id) REFERENCES control_node_runs(id) ON DELETE CASCADE,
613 FOREIGN KEY(branch_run_id) REFERENCES branch_runs(id) ON DELETE SET NULL
614 );
615 CREATE INDEX IF NOT EXISTS idx_teacher_candidates_workflow_run_id
616 ON teacher_candidates(workflow_run_id);
617 CREATE INDEX IF NOT EXISTS idx_teacher_candidates_control_node_run_id
618 ON teacher_candidates(control_node_run_id);
619
620 PRAGMA user_version = 2;
621 "#,
622 )
623 .context("failed to initialize workflow trace schema")?;
624 user_version = 2;
625 }
626 if user_version < 3 {
627 conn.execute_batch(
628 r#"
629 CREATE TABLE IF NOT EXISTS thread_goals (
630 thread_id TEXT PRIMARY KEY NOT NULL,
631 goal_id TEXT NOT NULL,
632 objective TEXT NOT NULL,
633 status TEXT NOT NULL CHECK(status IN (
634 'active',
635 'paused',
636 'blocked',
637 'usage_limited',
638 'budget_limited',
639 'complete'
640 )),
641 token_budget INTEGER,
642 tokens_used INTEGER NOT NULL DEFAULT 0,
643 time_used_seconds INTEGER NOT NULL DEFAULT 0,
644 created_at INTEGER NOT NULL,
645 updated_at INTEGER NOT NULL,
646 FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE
647 );
648
649 PRAGMA user_version = 3;
650 "#,
651 )
652 .context("failed to initialize thread goal schema")?;
653 user_version = 3;
654 }
655 if user_version < 4 {
656 // Same restore/race guard as the v0 block: the column may
657 // already exist even though the header predates version 4.
658 let add_continuation_count = if column_exists(
659 conn,
660 "thread_goals",
661 "continuation_count",
662 )? {
663 ""
664 } else {
665 "ALTER TABLE thread_goals\n ADD COLUMN continuation_count INTEGER NOT NULL DEFAULT 0;"
666 };
667 conn.execute_batch(&format!(
668 r#"
669 {add_continuation_count}
670
671 PRAGMA user_version = 4;
672 "#
673 ))
674 .context("failed to initialize thread goal continuation schema")?;
675 user_version = 4;
676 }
677 if user_version < 5 {
678 let mut additions = String::new();
679 for (column, definition) in [
680 ("last_gap_fingerprint", "TEXT"),
681 ("repeated_gap_count", "INTEGER NOT NULL DEFAULT 0"),
682 ("last_gap_pass", "INTEGER"),
683 ("pause_reason", "TEXT"),
684 ] {
685 if !column_exists(conn, "thread_goals", column)? {
686 additions.push_str(&format!(
687 "ALTER TABLE thread_goals ADD COLUMN {column} {definition};\n"
688 ));
689 }
690 }
691 conn.execute_batch(&format!("{additions} PRAGMA user_version = 5;"))
692 .context("failed to initialize durable goal stall schema")?;
693 user_version = 5;
694 }
695 if user_version < 6 {
696 conn.execute_batch(
697 r#"
698 CREATE TABLE IF NOT EXISTS thread_runtime_links (
699 thread_id TEXT PRIMARY KEY NOT NULL,
700 runtime_thread_id TEXT NOT NULL,
701 created_at INTEGER NOT NULL,
702 FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE
703 );
704
705 PRAGMA user_version = 6;
706 "#,
707 )
708 .context("failed to initialize thread runtime link schema")?;
709 user_version = 6;
710 }
711 if user_version < 7 {
712 conn.execute_batch(
713 "CREATE TABLE IF NOT EXISTS state_store_identity (singleton INTEGER PRIMARY KEY CHECK(singleton = 1), identity TEXT NOT NULL);
714 INSERT OR IGNORE INTO state_store_identity VALUES(1, lower(hex(randomblob(32))));
715 CREATE TABLE IF NOT EXISTS thread_runtime_receipts (thread_id TEXT PRIMARY KEY REFERENCES threads(id) ON DELETE CASCADE, receipt_json TEXT NOT NULL);
716 CREATE TABLE IF NOT EXISTS thread_runtime_operations (operation_key TEXT PRIMARY KEY, request_digest TEXT NOT NULL, receipt_json TEXT NOT NULL);
717 PRAGMA user_version = 7;",
718 ).context("failed to initialize bound canonical Runtime receipts")?;
719 }
720 Ok(())
721 }
722
723 /// The runtime (turn engine) thread a client-facing thread runs on, if
724 /// one was linked. Links outlive the app-server process, so a thread
725 /// keeps its conversation across a daemon restart.
726 pub fn get_runtime_thread_link(&self, thread_id: &str) -> Result<Option<String>> {
727 let conn = self.conn()?;
728 conn.query_row(
729 "SELECT runtime_thread_id FROM thread_runtime_links WHERE thread_id = ?1",
730 params![thread_id],
731 |row| row.get(0),
732 )
733 .optional()
734 .with_context(|| format!("failed to read runtime link for thread {thread_id}"))
735 }
736
737 /// Record the runtime thread for `thread_id`. The thread must exist; the
738 /// link is removed with it.
739 #[cfg(test)]
740 fn set_runtime_thread_link(&self, thread_id: &str, runtime_thread_id: &str) -> Result<()> {
741 let mut conn = self.conn()?;
742 let tx = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
743 let bound: bool = tx.query_row(
744 "SELECT EXISTS(SELECT 1 FROM thread_runtime_receipts WHERE thread_id = ?1)",
745 params![thread_id],
746 |row| row.get(0),
747 )?;
748 anyhow::ensure!(
749 !bound,
750 "bound canonical aliases cannot be changed by the legacy link writer"
751 );
752 tx.execute(
753 r#"
754 INSERT INTO thread_runtime_links(thread_id, runtime_thread_id, created_at)
755 VALUES (?1, ?2, ?3)
756 ON CONFLICT(thread_id) DO UPDATE SET
757 runtime_thread_id = excluded.runtime_thread_id,
758 created_at = excluded.created_at
759 "#,
760 params![thread_id, runtime_thread_id, Utc::now().timestamp()],
761 )
762 .with_context(|| format!("failed to save runtime link for thread {thread_id}"))?;
763 tx.commit()?;
764 Ok(())
765 }
766
767 /// Insert or update thread metadata.
768 ///
769 /// This does **not** update legacy history or its selected leaf. Restoring an
770 /// old transcript requires an absent target and a complete immutable archive.
771 pub fn upsert_thread(&self, thread: &ThreadMetadata) -> Result<()> {
772 let conn = self.conn()?;
773 write_thread_metadata_on(&conn, thread)?;
774
775 self.append_thread_name(
776 &thread.id,
777 thread.name.clone(),
778 thread.updated_at,
779 thread.rollout_path.clone(),
780 )?;
781 Ok(())
782 }
783
784 /// Retrieve a single thread by its ID.
785 ///
786 /// Returns `None` if no thread with the given ID exists.
787 pub fn get_thread(&self, id: &str) -> Result<Option<ThreadMetadata>> {
788 let conn = self.conn()?;
789 conn.query_row(
790 r#"
791 SELECT id, rollout_path, preview, ephemeral, model_provider, created_at, updated_at, status, path, cwd,
792 cli_version, source, title, sandbox_policy, approval_mode, archived, archived_at,
793 git_sha, git_branch, git_origin_url, memory_mode, current_leaf_id
794 FROM threads
795 WHERE id = ?1
796 "#,
797 params![id],
798 row_to_thread,
799 )
800 .optional()
801 .context("failed to read thread")
802 }
803
804 /// List threads ordered by most recently updated.
805 ///
806 /// Use [`ThreadListFilters`] to control whether archived threads are included
807 /// and the maximum number of results returned.
808 pub fn list_threads(&self, filters: ThreadListFilters) -> Result<Vec<ThreadMetadata>> {
809 let conn = self.conn()?;
810 let sql = if filters.include_archived {
811 "SELECT id, rollout_path, preview, ephemeral, model_provider, created_at, updated_at, status, path, cwd, cli_version, source, title, sandbox_policy, approval_mode, archived, archived_at, git_sha, git_branch, git_origin_url, memory_mode, current_leaf_id FROM threads ORDER BY updated_at DESC LIMIT ?1"
812 } else {
813 "SELECT id, rollout_path, preview, ephemeral, model_provider, created_at, updated_at, status, path, cwd, cli_version, source, title, sandbox_policy, approval_mode, archived, archived_at, git_sha, git_branch, git_origin_url, memory_mode, current_leaf_id FROM threads WHERE archived = 0 ORDER BY updated_at DESC LIMIT ?1"
814 };
815
816 let mut stmt = conn.prepare(sql).context("failed to prepare list query")?;
817 let limit = i64::try_from(filters.limit.unwrap_or(50)).unwrap_or(50);
818 let mut rows = stmt
819 .query(params![limit])
820 .context("failed to query threads")?;
821 let mut out = Vec::new();
822 while let Some(row) = rows.next().context("failed to iterate thread rows")? {
823 out.push(row_to_thread(row)?);
824 }
825 Ok(out)
826 }
827
828 /// Archive a thread, setting its status to [`ThreadStatus::Archived`] and
829 /// recording the current timestamp.
830 pub fn mark_archived(&self, id: &str) -> Result<()> {
831 let conn = self.conn()?;
832 conn.execute(
833 "UPDATE threads SET archived = 1, archived_at = ?2, status = ?3 WHERE id = ?1",
834 params![
835 id,
836 Utc::now().timestamp(),
837 thread_status_to_str(&ThreadStatus::Archived)
838 ],
839 )
840 .context("failed to archive thread")?;
841 Ok(())
842 }
843
844 /// Unarchive a thread, removing the archived flag and clearing `archived_at`.
845 pub fn mark_unarchived(&self, id: &str) -> Result<()> {
846 let conn = self.conn()?;
847 conn.execute(
848 "UPDATE threads SET archived = 0, archived_at = NULL, status = CASE WHEN status = ?2 THEN ?3 ELSE status END WHERE id = ?1",
849 params![
850 id,
851 thread_status_to_str(&ThreadStatus::Archived),
852 thread_status_to_str(&ThreadStatus::Idle),
853 ],
854 )
855 .context("failed to unarchive thread")?;
856 Ok(())
857 }
858
859 /// Permanently delete a thread and all of its associated data
860 /// (messages, checkpoints, dynamic tools) via cascading foreign keys, and
861 /// drop it from the session index.
862 ///
863 /// The index is append-only and name lookups read it without the
864 /// database, so a deleted thread left there kept answering for its name —
865 /// shadowing a live thread of the same name when it was newer.
866 pub fn delete_thread(&self, id: &str) -> Result<()> {
867 {
868 let conn = self.conn()?;
869 conn.execute("DELETE FROM threads WHERE id = ?1", params![id])
870 .context("failed to delete thread")?;
871 }
872 let _guard = SESSION_INDEX_LOCK.lock().unwrap();
873 self.with_session_index_lock(|| {
874 let mut latest = self.session_index_map()?;
875 if latest.remove(id).is_some() {
876 self.rewrite_session_index_locked(&latest)?;
877 }
878 Ok(())
879 })
880 }
881
882 /// Insert or replace the persisted goal for a thread.
883 #[cfg(test)]
884 fn upsert_thread_goal(&self, goal: &ThreadGoalRecord) -> Result<()> {
885 let conn = self.conn()?;
886 write_thread_goal_on(&conn, goal)
887 }
888
889 /// Accrue additional token and wall-clock usage onto a thread's persisted goal.
890 ///
891 /// This is the durable, additive accounting path for the persistent goal loop: it
892 /// increments `tokens_used` and `time_used_seconds` in a single atomic SQL `UPDATE`
893 /// (`col = col + ?`) so concurrent accruals do not race a read-modify-write. The
894 /// goal's `updated_at` is advanced to the larger of its current value and `now`,
895 /// keeping the timestamp monotonic even if a stale `now` is supplied.
896 ///
897 /// `token_delta` and `time_delta_seconds` are added on the database side; callers
898 /// should pass non-negative deltas (negative values are accepted and will decrement,
899 /// which is intentionally left to the caller's discretion).
900 ///
901 /// Returns the updated [`ThreadGoalRecord`], or `Ok(None)` if the thread has no
902 /// persisted goal. Unlike the historical fixture upsert this never
903 /// creates a goal row; it only accumulates onto an existing one.
904 #[cfg(test)]
905 fn record_thread_goal_usage(
906 &self,
907 thread_id: &str,
908 token_delta: i64,
909 time_delta_seconds: i64,
910 now: i64,
911 ) -> Result<Option<ThreadGoalRecord>> {
912 let conn = self.conn()?;
913 let changed = conn
914 .execute(
915 r#"
916 UPDATE thread_goals
917 SET tokens_used = tokens_used + ?2,
918 time_used_seconds = time_used_seconds + ?3,
919 updated_at = MAX(updated_at, ?4)
920 WHERE thread_id = ?1
921 "#,
922 params![thread_id, token_delta, time_delta_seconds, now],
923 )
924 .context("failed to record thread goal usage")?;
925 if changed == 0 {
926 return Ok(None);
927 }
928 Self::read_thread_goal(&conn, thread_id)
929 }
930
931 /// Increment the durable cross-turn continuation counter for a thread goal.
932 ///
933 /// The older TUI continuation guard is scoped to one engine turn. This
934 /// counter is intentionally persisted so a resumed goal loop can feed
935 /// `goal_loop::decide_continuation` with the true cross-turn count.
936 #[cfg(test)]
937 fn record_thread_goal_continuation(
938 &self,
939 thread_id: &str,
940 now: i64,
941 ) -> Result<Option<ThreadGoalRecord>> {
942 let conn = self.conn()?;
943 let changed = conn
944 .execute(
945 r#"
946 UPDATE thread_goals
947 SET continuation_count = continuation_count + 1,
948 updated_at = MAX(updated_at, ?2)
949 WHERE thread_id = ?1
950 "#,
951 params![thread_id, now],
952 )
953 .context("failed to record thread goal continuation")?;
954 if changed == 0 {
955 return Ok(None);
956 }
957 Self::read_thread_goal(&conn, thread_id)
958 }
959
960 /// Retrieve the persisted goal for a thread.
961 pub fn get_thread_goal(&self, thread_id: &str) -> Result<Option<ThreadGoalRecord>> {
962 let conn = self.conn()?;
963 Self::read_thread_goal(&conn, thread_id)
964 }
965
966 /// Read a goal on an already-held connection. The `record_*` mutators call
967 /// this instead of the public goal reader, which would re-lock the
968 /// connection mutex and self-deadlock.
969 fn read_thread_goal(conn: &Connection, thread_id: &str) -> Result<Option<ThreadGoalRecord>> {
970 let mut goal = Self::read_thread_goal_snapshot(conn, thread_id)?;
971 if let Some(goal) = goal.as_mut() {
972 goal.normalize_restored_stall_state();
973 }
974 Ok(goal)
975 }
976
977 /// Immutable migration/CAS snapshot keeps the original persisted status.
978 /// Restore normalization belongs to the reader/target, never source CAS.
979 fn read_thread_goal_snapshot(
980 conn: &Connection,
981 thread_id: &str,
982 ) -> Result<Option<ThreadGoalRecord>> {
983 conn.query_row(
984 r#"
985 SELECT thread_id, goal_id, objective, status, token_budget, tokens_used,
986 time_used_seconds, continuation_count, created_at, updated_at,
987 last_gap_fingerprint, repeated_gap_count, last_gap_pass, pause_reason
988 FROM thread_goals
989 WHERE thread_id = ?1
990 "#,
991 params![thread_id],
992 row_to_thread_goal,
993 )
994 .optional()
995 .context("failed to read thread goal")
996 }
997
998 /// Delete the persisted goal for a thread.
999 #[cfg(test)]
1000 fn delete_thread_goal(&self, thread_id: &str) -> Result<bool> {
1001 let conn = self.conn()?;
1002 let changed = conn
1003 .execute(
1004 "DELETE FROM thread_goals WHERE thread_id = ?1",
1005 params![thread_id],
1006 )
1007 .context("failed to delete thread goal")?;
1008 Ok(changed > 0)
1009 }
1010
1011 /// List all leaf messages in a thread.
1012 ///
1013 /// A leaf message is one that has no other message referencing it as a parent.
1014 /// In a branching conversation tree, there may be multiple leaf messages.
1015 pub fn list_leaf_messages(&self, thread_id: &str) -> Result<Vec<MessageRecord>> {
1016 let conn = self.conn()?;
1017 let mut stmt = conn
1018 .prepare(
1019 r#"
1020 SELECT m1.id, m1.thread_id, m1.role, m1.content, m1.item_json, m1.created_at, m1.parent_entry_id
1021 FROM messages m1
1022 LEFT JOIN messages m2 ON m1.id = m2.parent_entry_id
1023 WHERE m1.thread_id = ?1 AND m2.id IS NULL
1024 "#,
1025 )
1026 .context("failed to prepare message listing query")?;
1027 let mut rows = stmt
1028 .query(params![thread_id])
1029 .with_context(|| format!("failed to list leaf messages for thread {thread_id}"))?;
1030 let mut out = Vec::new();
1031 while let Some(row) = rows.next().context("failed to iterate message rows")? {
1032 let item_json: Option<String> = row.get(4).context("failed to read item json")?;
1033 let item = item_json
1034 .as_deref()
1035 .map(serde_json::from_str)
1036 .transpose()
1037 .with_context(|| {
1038 format!("failed to parse message item json in thread {thread_id}")
1039 })?;
1040 out.push(MessageRecord {
1041 id: row.get(0).context("failed to read message id")?,
1042 thread_id: row.get(1).context("failed to read message thread id")?,
1043 role: row.get(2).context("failed to read message role")?,
1044 content: row.get(3).context("failed to read message content")?,
1045 item,
1046 created_at: row.get(5).context("failed to read message timestamp")?,
1047 parent_entry_id: row.get(6).context("failed to read parent entry id")?,
1048 });
1049 }
1050 Ok(out)
1051 }
1052
1053 /// Replace the dynamic tools for a thread.
1054 ///
1055 /// All existing dynamic tools for the thread are deleted and replaced with the
1056 /// provided list. The operation is performed within a transaction.
1057 pub fn persist_dynamic_tools(
1058 &self,
1059 thread_id: &str,
1060 tools: &[DynamicToolRecord],
1061 ) -> Result<()> {
1062 let mut conn = self.conn()?;
1063 let tx = conn
1064 .transaction()
1065 .context("failed to begin dynamic tools transaction")?;
1066 tx.execute(
1067 "DELETE FROM thread_dynamic_tools WHERE thread_id = ?1",
1068 params![thread_id],
1069 )
1070 .context("failed to clear dynamic tools")?;
1071 for tool in tools {
1072 tx.execute(
1073 "INSERT INTO thread_dynamic_tools(thread_id, position, name, description, input_schema) VALUES (?1, ?2, ?3, ?4, ?5)",
1074 params![
1075 thread_id,
1076 tool.position,
1077 tool.name,
1078 tool.description,
1079 tool.input_schema.to_string()
1080 ],
1081 )
1082 .with_context(|| format!("failed to persist dynamic tool {}", tool.name))?;
1083 }
1084 tx.commit().context("failed to commit dynamic tools")?;
1085 Ok(())
1086 }
1087
1088 /// Retrieve all dynamic tools registered for a thread, ordered by position.
1089 pub fn get_dynamic_tools(&self, thread_id: &str) -> Result<Vec<DynamicToolRecord>> {
1090 let conn = self.conn()?;
1091 let mut stmt = conn
1092 .prepare(
1093 "SELECT position, name, description, input_schema FROM thread_dynamic_tools WHERE thread_id = ?1 ORDER BY position ASC",
1094 )
1095 .context("failed to prepare get dynamic tools query")?;
1096 let mut rows = stmt
1097 .query(params![thread_id])
1098 .context("failed to query dynamic tools")?;
1099 let mut out = Vec::new();
1100 while let Some(row) = rows.next().context("failed to iterate dynamic tools")? {
1101 let input_schema_raw: String =
1102 row.get(3).context("failed to read tool input schema")?;
1103 let input_schema: Value =
1104 serde_json::from_str(&input_schema_raw).with_context(|| {
1105 format!("failed to parse input schema for dynamic tool in thread {thread_id}")
1106 })?;
1107 out.push(DynamicToolRecord {
1108 position: row.get(0).context("failed to read tool position")?,
1109 name: row.get(1).context("failed to read tool name")?,
1110 description: row.get(2).context("failed to read tool description")?,
1111 input_schema,
1112 });
1113 }
1114 Ok(out)
1115 }
1116
1117 /// Append a new message to a thread.
1118 ///
1119 /// The message is linked to the thread's current leaf as its parent, and the
1120 /// thread's `current_leaf_id` is updated to the new message. Returns the ID
1121 /// of the newly created message.
1122 #[cfg(test)]
1123 fn append_message(
1124 &self,
1125 thread_id: &str,
1126 role: &str,
1127 content: &str,
1128 item: Option<Value>,
1129 ) -> Result<i64> {
1130 let ids = self.append_messages(
1131 thread_id,
1132 &[NewMessage {
1133 role: role.to_string(),
1134 content: content.to_string(),
1135 item,
1136 }],
1137 )?;
1138 ids.first()
1139 .copied()
1140 .context("append message transaction returned no id")
1141 }
1142
1143 /// Append `messages` to a thread's current branch in one transaction.
1144 ///
1145 /// Each message is linked to the one before it (the first to the thread's
1146 /// current leaf), and the leaf ends on the last. All or nothing: a failure
1147 /// part-way leaves the thread exactly as it was, never a partial history.
1148 /// Returns the new message ids in order.
1149 #[cfg(test)]
1150 fn append_messages(&self, thread_id: &str, messages: &[NewMessage]) -> Result<Vec<i64>> {
1151 if messages.is_empty() {
1152 return Ok(Vec::new());
1153 }
1154 let encoded = messages
1155 .iter()
1156 .map(|message| {
1157 message
1158 .item
1159 .as_ref()
1160 .map(serde_json::to_string)
1161 .transpose()
1162 .context("failed to serialize message item payload")
1163 })
1164 .collect::<Result<Vec<_>>>()?;
1165 let mut conn = self.conn()?;
1166 let created_at = Utc::now().timestamp();
1167
1168 let tx = conn
1169 .transaction()
1170 .context("failed to begin append message transaction")?;
1171
1172 let mut leaf_id: Option<i64> = tx
1173 .query_row(
1174 "SELECT current_leaf_id FROM threads WHERE id = ?1",
1175 params![thread_id],
1176 |row| row.get(0),
1177 )
1178 .with_context(|| {
1179 format!("failed to query thread current leaf id for thread {thread_id}")
1180 })?;
1181
1182 let mut ids = Vec::with_capacity(messages.len());
1183 for (message, item_json) in messages.iter().zip(encoded) {
1184 let next_leaf_id: i64 = tx.query_row(
1185 r#"
1186 INSERT INTO messages(thread_id, role, content, item_json, created_at, parent_entry_id)
1187 SELECT ?1, ?2, ?3, ?4, ?5, ?6
1188 RETURNING id
1189 "#, params![thread_id, message.role, message.content, item_json, created_at, leaf_id], |row| row.get(0)
1190 ).with_context(|| format!("failed to append message for thread {thread_id}"))?;
1191 leaf_id = Some(next_leaf_id);
1192 ids.push(next_leaf_id);
1193 }
1194
1195 tx.execute(
1196 r#"
1197 UPDATE threads
1198 SET current_leaf_id = ?1
1199 WHERE id = ?2;
1200 "#,
1201 params![leaf_id, thread_id],
1202 )
1203 .with_context(|| {
1204 format!("failed to update thread current leaf id for thread {thread_id}")
1205 })?;
1206
1207 tx.commit()
1208 .context("failed to commit append message transaction")?;
1209
1210 Ok(ids)
1211 }
1212
1213 /// List messages in the current conversation branch, walking backwards from
1214 /// the thread's `current_leaf_id`.
1215 ///
1216 /// Messages are returned in chronological order (oldest first). The `limit`
1217 /// parameter caps how many ancestor messages are traversed; it defaults to 500.
1218 pub fn list_messages(
1219 &self,
1220 thread_id: &str,
1221 limit: Option<usize>,
1222 ) -> Result<Vec<MessageRecord>> {
1223 let conn = self.conn()?;
1224 let limit = i64::try_from(limit.unwrap_or(500)).unwrap_or(500);
1225 let mut stmt = conn
1226 .prepare(
1227 r#"
1228 WITH RECURSIVE
1229 leaf_id AS (
1230 SELECT current_leaf_id FROM threads WHERE id = ?1
1231 ),
1232 ancestors AS (
1233 SELECT id, thread_id, role, content, item_json, created_at, parent_entry_id, 0 AS depth
1234 FROM messages
1235 WHERE id = (SELECT current_leaf_id FROM leaf_id)
1236
1237 UNION ALL
1238
1239 SELECT m.id, m.thread_id, m.role, m.content, m.item_json, m.created_at, m.parent_entry_id, a.depth + 1
1240 FROM messages m
1241 JOIN ancestors a ON m.id = a.parent_entry_id
1242 WHERE a.depth < ?2
1243 )
1244 SELECT id, thread_id, role, content, item_json, created_at, parent_entry_id FROM ancestors
1245 ORDER BY depth DESC
1246 "#
1247 )
1248 .context("failed to prepare message listing query")?;
1249 let mut rows = stmt
1250 .query(params![thread_id, limit - 1])
1251 .with_context(|| format!("failed to list messages for thread {thread_id}"))?;
1252 let mut out = Vec::new();
1253 while let Some(row) = rows.next().context("failed to iterate message rows")? {
1254 let item_json: Option<String> = row.get(4).context("failed to read item json")?;
1255 let item = item_json
1256 .as_deref()
1257 .map(serde_json::from_str)
1258 .transpose()
1259 .with_context(|| {
1260 format!("failed to parse message item json in thread {thread_id}")
1261 })?;
1262 out.push(MessageRecord {
1263 id: row.get(0).context("failed to read message id")?,
1264 thread_id: row.get(1).context("failed to read message thread id")?,
1265 role: row.get(2).context("failed to read message role")?,
1266 content: row.get(3).context("failed to read message content")?,
1267 item,
1268 created_at: row.get(5).context("failed to read message timestamp")?,
1269 parent_entry_id: row.get(6).context("failed to read parent entry id")?,
1270 });
1271 }
1272 Ok(out)
1273 }
1274
1275 /// Delete all messages belonging to a thread and reset its `current_leaf_id`.
1276 ///
1277 /// Returns the number of messages deleted.
1278 #[cfg(test)]
1279 fn clear_messages(&self, thread_id: &str) -> Result<usize> {
1280 let mut conn = self.conn()?;
1281 let tx = conn
1282 .transaction()
1283 .context("failed to begin clear messages transaction")?;
1284
1285 tx.execute(
1286 r#"
1287 UPDATE threads
1288 SET current_leaf_id = NULL
1289 WHERE id = ?1;
1290 "#,
1291 params![thread_id],
1292 )
1293 .with_context(|| format!("failed to clear messages for thread {thread_id}"))?;
1294 let result = tx
1295 .execute(
1296 r#"
1297 DELETE FROM messages WHERE thread_id = ?1
1298 "#,
1299 params![thread_id],
1300 )
1301 .with_context(|| format!("failed to clear messages for thread {thread_id}"))?;
1302 tx.commit()
1303 .context("failed to commit clear messages transaction")?;
1304
1305 Ok(result)
1306 }
1307
1308 /// Save (or update) a named checkpoint for a thread.
1309 ///
1310 /// If a checkpoint with the same `thread_id` and `checkpoint_id` already exists,
1311 /// its state and timestamp are overwritten.
1312 #[cfg(test)]
1313 fn save_checkpoint(&self, thread_id: &str, checkpoint_id: &str, state: &Value) -> Result<()> {
1314 let conn = self.conn()?;
1315 let state_json =
1316 serde_json::to_string(state).context("failed to encode checkpoint state")?;
1317 conn.execute(
1318 r#"
1319 INSERT INTO checkpoints(thread_id, checkpoint_id, state_json, created_at)
1320 VALUES (?1, ?2, ?3, ?4)
1321 ON CONFLICT(thread_id, checkpoint_id) DO UPDATE SET
1322 state_json = excluded.state_json,
1323 created_at = excluded.created_at
1324 "#,
1325 params![thread_id, checkpoint_id, state_json, Utc::now().timestamp()],
1326 )
1327 .with_context(|| {
1328 format!("failed to save checkpoint {checkpoint_id} for thread {thread_id}")
1329 })?;
1330 Ok(())
1331 }
1332
1333 /// Load a checkpoint for a thread.
1334 ///
1335 /// If `checkpoint_id` is provided, loads that specific checkpoint. Otherwise,
1336 /// loads the most recently created checkpoint for the thread. Returns `None`
1337 /// if no matching checkpoint exists.
1338 pub fn load_checkpoint(
1339 &self,
1340 thread_id: &str,
1341 checkpoint_id: Option<&str>,
1342 ) -> Result<Option<CheckpointRecord>> {
1343 let conn = self.conn()?;
1344 if let Some(checkpoint_id) = checkpoint_id {
1345 let row = conn
1346 .query_row(
1347 "SELECT thread_id, checkpoint_id, state_json, created_at FROM checkpoints WHERE thread_id = ?1 AND checkpoint_id = ?2",
1348 params![thread_id, checkpoint_id],
1349 |row| {
1350 Ok((
1351 row.get::<_, String>(0)?,
1352 row.get::<_, String>(1)?,
1353 row.get::<_, String>(2)?,
1354 row.get::<_, i64>(3)?,
1355 ))
1356 },
1357 )
1358 .optional()
1359 .with_context(|| {
1360 format!("failed to load checkpoint {checkpoint_id} for thread {thread_id}")
1361 })?;
1362 if let Some((thread_id, checkpoint_id, state_json, created_at)) = row {
1363 let state = parse_checkpoint_state(&state_json)?;
1364 return Ok(Some(CheckpointRecord {
1365 thread_id,
1366 checkpoint_id,
1367 state,
1368 created_at,
1369 }));
1370 }
1371 return Ok(None);
1372 }
1373
1374 let row = conn
1375 .query_row(
1376 "SELECT thread_id, checkpoint_id, state_json, created_at FROM checkpoints WHERE thread_id = ?1 ORDER BY created_at DESC LIMIT 1",
1377 params![thread_id],
1378 |row| {
1379 Ok((
1380 row.get::<_, String>(0)?,
1381 row.get::<_, String>(1)?,
1382 row.get::<_, String>(2)?,
1383 row.get::<_, i64>(3)?,
1384 ))
1385 },
1386 )
1387 .optional()
1388 .with_context(|| format!("failed to load latest checkpoint for thread {thread_id}"))?;
1389 if let Some((thread_id, checkpoint_id, state_json, created_at)) = row {
1390 let state = parse_checkpoint_state(&state_json)?;
1391 return Ok(Some(CheckpointRecord {
1392 thread_id,
1393 checkpoint_id,
1394 state,
1395 created_at,
1396 }));
1397 }
1398 Ok(None)
1399 }
1400
1401 /// List checkpoints for a thread, ordered by creation time (newest first).
1402 ///
1403 /// The `limit` parameter caps the number of results and defaults to 100.
1404 pub fn list_checkpoints(
1405 &self,
1406 thread_id: &str,
1407 limit: Option<usize>,
1408 ) -> Result<Vec<CheckpointRecord>> {
1409 let conn = self.conn()?;
1410 let limit = i64::try_from(limit.unwrap_or(100)).unwrap_or(100);
1411 let mut stmt = conn
1412 .prepare(
1413 "SELECT thread_id, checkpoint_id, state_json, created_at FROM checkpoints WHERE thread_id = ?1 ORDER BY created_at DESC LIMIT ?2",
1414 )
1415 .context("failed to prepare checkpoint list query")?;
1416 let mut rows = stmt
1417 .query(params![thread_id, limit])
1418 .with_context(|| format!("failed to list checkpoints for thread {thread_id}"))?;
1419
1420 let mut out = Vec::new();
1421 while let Some(row) = rows.next().context("failed to iterate checkpoint rows")? {
1422 let state_json: String = row.get(2).context("failed to read checkpoint state json")?;
1423 let state = parse_checkpoint_state(&state_json)?;
1424 out.push(CheckpointRecord {
1425 thread_id: row.get(0).context("failed to read checkpoint thread id")?,
1426 checkpoint_id: row.get(1).context("failed to read checkpoint id")?,
1427 state,
1428 created_at: row.get(3).context("failed to read checkpoint timestamp")?,
1429 });
1430 }
1431 Ok(out)
1432 }
1433
1434 /// Insert or update a background job record.
1435 pub fn upsert_job(&self, job: &JobStateRecord) -> Result<()> {
1436 let conn = self.conn()?;
1437 conn.execute(
1438 r#"
1439 INSERT INTO jobs(id, name, status, progress, detail, created_at, updated_at)
1440 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
1441 ON CONFLICT(id) DO UPDATE SET
1442 name = excluded.name,
1443 status = excluded.status,
1444 progress = excluded.progress,
1445 detail = excluded.detail,
1446 created_at = excluded.created_at,
1447 updated_at = excluded.updated_at
1448 "#,
1449 params![
1450 job.id,
1451 job.name,
1452 job_state_status_to_str(&job.status),
1453 job.progress.map(i64::from),
1454 job.detail,
1455 job.created_at,
1456 job.updated_at
1457 ],
1458 )
1459 .with_context(|| format!("failed to upsert job {}", job.id))?;
1460 Ok(())
1461 }
1462
1463 /// List jobs ordered by most recently updated.
1464 ///
1465 /// The `limit` parameter caps the number of results and defaults to 100.
1466 pub fn list_jobs(&self, limit: Option<usize>) -> Result<Vec<JobStateRecord>> {
1467 let conn = self.conn()?;
1468 let limit = i64::try_from(limit.unwrap_or(100)).unwrap_or(100);
1469 let mut stmt = conn
1470 .prepare(
1471 "SELECT id, name, status, progress, detail, created_at, updated_at FROM jobs ORDER BY updated_at DESC LIMIT ?1",
1472 )
1473 .context("failed to prepare job list query")?;
1474 let mut rows = stmt
1475 .query(params![limit])
1476 .context("failed to query persisted jobs")?;
1477 let mut out = Vec::new();
1478 while let Some(row) = rows.next().context("failed to iterate persisted jobs")? {
1479 let status_raw: String = row.get(2).context("failed to read job status")?;
1480 let progress: Option<i64> = row.get(3).context("failed to read job progress")?;
1481 out.push(JobStateRecord {
1482 id: row.get(0).context("failed to read job id")?,
1483 name: row.get(1).context("failed to read job name")?,
1484 status: job_state_status_from_str(&status_raw),
1485 progress: progress.and_then(|v| u8::try_from(v).ok()),
1486 detail: row.get(4).context("failed to read job detail")?,
1487 created_at: row.get(5).context("failed to read job created_at")?,
1488 updated_at: row.get(6).context("failed to read job updated_at")?,
1489 });
1490 }
1491 Ok(out)
1492 }
1493
1494 /// Look up the rollout file path for a thread by its ID.
1495 pub fn find_rollout_path_by_id(&self, id: &str) -> Result<Option<PathBuf>> {
1496 let conn = self.conn()?;
1497 conn.query_row(
1498 "SELECT rollout_path FROM threads WHERE id = ?1",
1499 params![id],
1500 |row| row.get::<_, Option<String>>(0),
1501 )
1502 .optional()
1503 .context("failed to lookup rollout path")
1504 .map(|opt| opt.flatten().map(PathBuf::from))
1505 }
1506
1507 /// Append an entry to the JSONL session index file.
1508 ///
1509 /// The session index is an append-only log that maps thread IDs to their names,
1510 /// update timestamps, and rollout paths. It is used for fast name-based lookups
1511 /// without opening the SQLite database.
1512 pub fn append_thread_name(
1513 &self,
1514 thread_id: &str,
1515 thread_name: Option<String>,
1516 updated_at: i64,
1517 rollout_path: Option<PathBuf>,
1518 ) -> Result<()> {
1519 // Hold the index lock for the entire append + compaction so a concurrent
1520 // `StateStore` clone cannot rename the index while we are appending.
1521 let _guard = SESSION_INDEX_LOCK.lock().unwrap();
1522 if let Some(parent) = self.session_index_path.parent() {
1523 fs::create_dir_all(parent).with_context(|| {
1524 format!(
1525 "failed to create session index directory {}",
1526 parent.display()
1527 )
1528 })?;
1529 }
1530 let entry = SessionIndexEntry {
1531 thread_id: thread_id.to_string(),
1532 thread_name,
1533 updated_at,
1534 rollout_path,
1535 };
1536 let encoded =
1537 serde_json::to_string(&entry).context("failed to serialize session index entry")?;
1538 // Append and compaction share one lock. Compaction rewrites the file
1539 // from a snapshot and renames over it, so an append landing between
1540 // that snapshot and the rename would be discarded — silently, since
1541 // the append already returned success to its caller.
1542 self.with_session_index_lock(|| {
1543 let mut file = OpenOptions::new()
1544 .create(true)
1545 .append(true)
1546 .open(&self.session_index_path)
1547 .with_context(|| {
1548 format!(
1549 "failed to open session index {}",
1550 self.session_index_path.display()
1551 )
1552 })?;
1553 writeln!(file, "{encoded}").context("failed to append session index entry")?;
1554 // Durability: without this a crash mid-write can leave a torn
1555 // final line. Reads tolerate one (see `session_index_map`), but
1556 // not losing the entry beats recovering from having lost it.
1557 file.sync_data()
1558 .context("failed to flush session index entry")?;
1559 drop(file);
1560 self.compact_session_index_locked()
1561 })
1562 }
1563
1564 /// Run `operation` holding the exclusive session-index lock.
1565 ///
1566 /// The lock is an adjacent `.lock` file rather than the index itself, so
1567 /// compaction's rename cannot pull the lock out from under a waiter. This
1568 /// mirrors the discipline `codewhale-config` uses for `config.toml`.
1569 fn with_session_index_lock<T>(&self, operation: impl FnOnce() -> Result<T>) -> Result<T> {
1570 if let Some(parent) = self.session_index_path.parent() {
1571 fs::create_dir_all(parent).with_context(|| {
1572 format!(
1573 "failed to create session index directory {}",
1574 parent.display()
1575 )
1576 })?;
1577 }
1578 let lock_path = self.session_index_path.with_extension("jsonl.lock");
1579 let lock_file = OpenOptions::new()
1580 .create(true)
1581 .read(true)
1582 .write(true)
1583 // The file is only a lock handle; its contents are never read and
1584 // truncating it would race other holders for no benefit.
1585 .truncate(false)
1586 .open(&lock_path)
1587 .with_context(|| {
1588 format!("failed to open session index lock {}", lock_path.display())
1589 })?;
1590 #[cfg(unix)]
1591 {
1592 use std::os::unix::fs::PermissionsExt as _;
1593 lock_file
1594 .set_permissions(fs::Permissions::from_mode(0o600))
1595 .with_context(|| {
1596 format!(
1597 "failed to secure session index lock {}",
1598 lock_path.display()
1599 )
1600 })?;
1601 }
1602 let mut lock = fd_lock::RwLock::new(lock_file);
1603 let _guard = lock
1604 .write()
1605 .with_context(|| format!("failed to lock session index {}", lock_path.display()))?;
1606 operation()
1607 }
1608
1609 /// Find the display name for a thread by its ID, using the session index.
1610 ///
1611 /// Returns `None` if the thread is not in the index or has no name.
1612 pub fn find_thread_name_by_id(&self, thread_id: &str) -> Result<Option<String>> {
1613 let map = self.session_index_map()?;
1614 Ok(map
1615 .get(thread_id)
1616 .and_then(|entry| entry.thread_name.clone()))
1617 }
1618
1619 /// Look up display names for multiple thread IDs at once.
1620 ///
1621 /// Returns a map from thread ID to its name (which may be `None`).
1622 pub fn find_thread_names_by_ids(
1623 &self,
1624 ids: &[String],
1625 ) -> Result<HashMap<String, Option<String>>> {
1626 let map = self.session_index_map()?;
1627 let mut out = HashMap::new();
1628 for id in ids {
1629 let name = map.get(id).and_then(|entry| entry.thread_name.clone());
1630 out.insert(id.clone(), name);
1631 }
1632 Ok(out)
1633 }
1634
1635 /// Find the rollout path for a thread by its display name (case-insensitive).
1636 ///
1637 /// If multiple threads share the same name, the most recently updated one is returned.
1638 /// Returns `None` if no matching thread is found.
1639 pub fn find_thread_path_by_name_str(&self, name: &str) -> Result<Option<PathBuf>> {
1640 let map = self.session_index_map()?;
1641 let matched = map
1642 .values()
1643 .filter(|entry| {
1644 entry
1645 .thread_name
1646 .as_deref()
1647 .is_some_and(|n| n.eq_ignore_ascii_case(name))
1648 })
1649 .max_by_key(|entry| entry.updated_at);
1650 Ok(matched.and_then(|entry| entry.rollout_path.clone()))
1651 }
1652
1653 /// Compact the session index. The caller must already hold the lock from
1654 /// [`Self::with_session_index_lock`]: this reads a snapshot and renames a
1655 /// rewritten file over the live one, and an append interleaved between
1656 /// those two steps is lost.
1657 fn compact_session_index_locked(&self) -> Result<()> {
1658 if !self.session_index_path.exists() {
1659 return Ok(());
1660 }
1661 let line_count = BufReader::new(
1662 OpenOptions::new()
1663 .read(true)
1664 .open(&self.session_index_path)
1665 .with_context(|| {
1666 format!(
1667 "failed to read session index {}",
1668 self.session_index_path.display()
1669 )
1670 })?,
1671 )
1672 .lines()
1673 .filter(|line| {
1674 line.as_ref()
1675 .map(|value| !value.trim().is_empty())
1676 .unwrap_or(false)
1677 })
1678 .count();
1679 if line_count <= session_index_compact_line_threshold() {
1680 return Ok(());
1681 }
1682
1683 let latest = self.session_index_map()?;
1684 self.rewrite_session_index_locked(&latest)
1685 }
1686
1687 /// Replace the session index with exactly `latest`, one line per thread.
1688 /// The caller holds the lock from [`Self::with_session_index_lock`].
1689 fn rewrite_session_index_locked(
1690 &self,
1691 latest: &HashMap<String, SessionIndexEntry>,
1692 ) -> Result<()> {
1693 let compact_path = self.session_index_path.with_extension("jsonl.compact");
1694 {
1695 let mut file = OpenOptions::new()
1696 .create(true)
1697 .write(true)
1698 .truncate(true)
1699 .open(&compact_path)
1700 .with_context(|| {
1701 format!(
1702 "failed to open compact session index {}",
1703 compact_path.display()
1704 )
1705 })?;
1706 for entry in latest.values() {
1707 let encoded = serde_json::to_string(entry)
1708 .context("failed to serialize compact session index entry")?;
1709 writeln!(file, "{encoded}")
1710 .context("failed to write compact session index entry")?;
1711 }
1712 }
1713 // The snapshot is written but the live file is still the old one:
1714 // this is the window an unsynchronized appender would write into and
1715 // lose. Tests widen it deliberately to prove the lock closes it.
1716 #[cfg(test)]
1717 tests::compaction_midpoint(&self.session_index_path);
1718 fs::rename(&compact_path, &self.session_index_path).with_context(|| {
1719 format!(
1720 "failed to replace session index {}",
1721 self.session_index_path.display()
1722 )
1723 })?;
1724 Ok(())
1725 }
1726
1727 #[cfg(test)]
1728 fn session_index_line_count(&self) -> Result<usize> {
1729 if !self.session_index_path.exists() {
1730 return Ok(0);
1731 }
1732 Ok(BufReader::new(
1733 OpenOptions::new()
1734 .read(true)
1735 .open(&self.session_index_path)
1736 .with_context(|| {
1737 format!(
1738 "failed to read session index {}",
1739 self.session_index_path.display()
1740 )
1741 })?,
1742 )
1743 .lines()
1744 .filter(|line| {
1745 line.as_ref()
1746 .map(|value| !value.trim().is_empty())
1747 .unwrap_or(false)
1748 })
1749 .count())
1750 }
1751
1752 fn session_index_map(&self) -> Result<HashMap<String, SessionIndexEntry>> {
1753 if !self.session_index_path.exists() {
1754 return Ok(HashMap::new());
1755 }
1756 let file = OpenOptions::new()
1757 .read(true)
1758 .open(&self.session_index_path)
1759 .with_context(|| {
1760 format!(
1761 "failed to read session index {}",
1762 self.session_index_path.display()
1763 )
1764 })?;
1765 let reader = BufReader::new(file);
1766 let mut latest = HashMap::<String, SessionIndexEntry>::new();
1767 for line in reader.lines() {
1768 let line = line.context("failed to read session index line")?;
1769 if line.trim().is_empty() {
1770 continue;
1771 }
1772 // Skip a line we can't parse instead of failing the whole read.
1773 // An append that was interrupted mid-write leaves a torn final
1774 // line; aborting here broke every thread-name lookup, and because
1775 // compaction reads through this same function, the index could
1776 // never repair itself either — the file stayed broken until
1777 // someone deleted it by hand.
1778 match serde_json::from_str::<SessionIndexEntry>(&line) {
1779 Ok(parsed) => {
1780 latest.insert(parsed.thread_id.clone(), parsed);
1781 }
1782 Err(err) => {
1783 tracing::warn!(
1784 "skipping unparseable session index entry in {}: {err}",
1785 self.session_index_path.display()
1786 );
1787 }
1788 }
1789 }
1790 Ok(latest)
1791 }
1792 }
1793
1794 /// Resolve the default SQLite state path without opening or creating it.
1795 ///
1796 /// An explicit `CODEWHALE_HOME` always yields `<override>/state.db` and blocks
1797 /// ambient legacy fallback. Without an override, an existing legacy database
1798 /// remains readable until it is migrated.
1799 #[must_use]
1800 pub fn default_state_db_path() -> PathBuf {
1801 // $CODEWHALE_HOME is a hard override of the base data directory
1802 // (docs/CONFIGURATION.md): when set, the state DB lives under it and we do
1803 // NOT fall back to the legacy ~/.deepseek path — silent fallback would
1804 // defeat the isolation the override promises (CI, containers, multi-project,
1805 // test harnesses). Legacy ~/.deepseek migration only applies to the default
1806 // home location.
1807 if let Some(overridden) = codewhale_home_override().ok().flatten() {
1808 return overridden.join("state.db");
1809 }
1810 let home = codewhale_paths::user_home().unwrap_or_else(|| PathBuf::from("."));
1811 // Prefer the CodeWhale directory, falling back to legacy DeepSeek path
1812 // so existing installs don't lose their session history.
1813 let primary = home.join(CODEWHALE_APP_DIR).join("state.db");
1814 if primary.exists() || !home.join(LEGACY_APP_DIR).join("state.db").exists() {
1815 primary
1816 } else {
1817 home.join(LEGACY_APP_DIR).join("state.db")
1818 }
1819 }
1820
1821 fn write_thread_metadata_on(conn: &Connection, thread: &ThreadMetadata) -> Result<()> {
1822 conn.execute(
1823 r#"
1824 INSERT INTO threads (
1825 id, rollout_path, preview, ephemeral, model_provider, created_at, updated_at, status, path, cwd,
1826 cli_version, source, title, sandbox_policy, approval_mode, archived, archived_at,
1827 git_sha, git_branch, git_origin_url, memory_mode
1828 ) VALUES (
1829 ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10,
1830 ?11, ?12, ?13, ?14, ?15, ?16, ?17,
1831 ?18, ?19, ?20, ?21
1832 )
1833 ON CONFLICT(id) DO UPDATE SET
1834 rollout_path=excluded.rollout_path,
1835 preview=excluded.preview,
1836 ephemeral=excluded.ephemeral,
1837 model_provider=excluded.model_provider,
1838 created_at=excluded.created_at,
1839 updated_at=excluded.updated_at,
1840 status=excluded.status,
1841 path=excluded.path,
1842 cwd=excluded.cwd,
1843 cli_version=excluded.cli_version,
1844 source=excluded.source,
1845 title=excluded.title,
1846 sandbox_policy=excluded.sandbox_policy,
1847 approval_mode=excluded.approval_mode,
1848 archived=excluded.archived,
1849 archived_at=excluded.archived_at,
1850 git_sha=excluded.git_sha,
1851 git_branch=excluded.git_branch,
1852 git_origin_url=excluded.git_origin_url,
1853 memory_mode=excluded.memory_mode
1854 "#,
1855 params![
1856 thread.id,
1857 path_to_opt_string(thread.rollout_path.as_deref()),
1858 thread.preview,
1859 bool_to_i64(thread.ephemeral),
1860 thread.model_provider,
1861 thread.created_at,
1862 thread.updated_at,
1863 thread_status_to_str(&thread.status),
1864 path_to_opt_string(thread.path.as_deref()),
1865 thread.cwd.display().to_string(),
1866 thread.cli_version,
1867 session_source_to_str(&thread.source),
1868 thread.name,
1869 thread.sandbox_policy,
1870 thread.approval_mode,
1871 bool_to_i64(thread.archived),
1872 thread.archived_at,
1873 thread.git_sha,
1874 thread.git_branch,
1875 thread.git_origin_url,
1876 thread.memory_mode,
1877 ],
1878 )
1879 .context("failed to upsert thread metadata")?;
1880 Ok(())
1881 }
1882
1883 fn write_thread_goal_on(conn: &Connection, goal: &ThreadGoalRecord) -> Result<()> {
1884 codewhale_protocol::validate_goal_stall_state(
1885 goal.last_gap_fingerprint.as_deref(),
1886 goal.repeated_gap_count,
1887 goal.last_gap_pass,
1888 u32::try_from(goal.continuation_count.max(0)).unwrap_or(u32::MAX),
1889 )
1890 .map_err(anyhow::Error::msg)?;
1891 let pause_reason = goal
1892 .pause_reason
1893 .map(|reason| serde_json::to_string(&reason))
1894 .transpose()?;
1895 let exists: Option<i64> = conn
1896 .query_row(
1897 "SELECT 1 FROM threads WHERE id = ?1",
1898 params![goal.thread_id],
1899 |row| row.get(0),
1900 )
1901 .optional()
1902 .context("failed to verify thread before saving goal")?;
1903 if exists.is_none() {
1904 anyhow::bail!("thread {} not found", goal.thread_id);
1905 }
1906
1907 conn.execute(
1908 r#"
1909 INSERT INTO thread_goals (
1910 thread_id, goal_id, objective, status, token_budget, tokens_used,
1911 time_used_seconds, continuation_count, created_at, updated_at,
1912 last_gap_fingerprint, repeated_gap_count, last_gap_pass, pause_reason
1913 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)
1914 ON CONFLICT(thread_id) DO UPDATE SET
1915 goal_id=excluded.goal_id,
1916 objective=excluded.objective,
1917 status=excluded.status,
1918 token_budget=excluded.token_budget,
1919 tokens_used=excluded.tokens_used,
1920 time_used_seconds=excluded.time_used_seconds,
1921 continuation_count=excluded.continuation_count,
1922 created_at=excluded.created_at,
1923 updated_at=excluded.updated_at,
1924 last_gap_fingerprint=excluded.last_gap_fingerprint,
1925 repeated_gap_count=excluded.repeated_gap_count,
1926 last_gap_pass=excluded.last_gap_pass,
1927 pause_reason=excluded.pause_reason
1928 "#,
1929 params![
1930 goal.thread_id,
1931 goal.goal_id,
1932 goal.objective,
1933 thread_goal_status_to_str(&goal.status),
1934 goal.token_budget,
1935 goal.tokens_used,
1936 goal.time_used_seconds,
1937 goal.continuation_count,
1938 goal.created_at,
1939 goal.updated_at,
1940 goal.last_gap_fingerprint,
1941 goal.repeated_gap_count,
1942 goal.last_gap_pass,
1943 pause_reason,
1944 ],
1945 )
1946 .context("failed to upsert thread goal")?;
1947 Ok(())
1948 }
1949
1950 fn bool_to_i64(value: bool) -> i64 {
1951 if value { 1 } else { 0 }
1952 }
1953
1954 /// Whether `table` currently has a column named `column`.
1955 ///
1956 /// Used to guard `ALTER TABLE ... ADD COLUMN` migrations so they are
1957 /// idempotent. Both identifiers are compile-time literals at every call
1958 /// site, never user input. A missing table reports `false`, matching the
1959 /// fresh-database case where the migration must still run.
1960 fn column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool> {
1961 let mut stmt = conn.prepare(&format!("PRAGMA table_info({table})"))?;
1962 let names = stmt.query_map([], |row| row.get::<_, String>(1))?;
1963 for name in names {
1964 if name? == column {
1965 return Ok(true);
1966 }
1967 }
1968 Ok(false)
1969 }
1970
1971 fn i64_to_bool(value: i64) -> bool {
1972 value != 0
1973 }
1974
1975 fn thread_status_to_str(status: &ThreadStatus) -> &'static str {
1976 match status {
1977 ThreadStatus::Running => "running",
1978 ThreadStatus::Idle => "idle",
1979 ThreadStatus::Completed => "completed",
1980 ThreadStatus::Failed => "failed",
1981 ThreadStatus::Paused => "paused",
1982 ThreadStatus::Archived => "archived",
1983 }
1984 }
1985
1986 fn thread_status_from_str(value: &str) -> ThreadStatus {
1987 match value {
1988 "running" => ThreadStatus::Running,
1989 "idle" => ThreadStatus::Idle,
1990 "completed" => ThreadStatus::Completed,
1991 "failed" => ThreadStatus::Failed,
1992 "paused" => ThreadStatus::Paused,
1993 "archived" => ThreadStatus::Archived,
1994 _ => ThreadStatus::Idle,
1995 }
1996 }
1997
1998 fn session_source_to_str(source: &SessionSource) -> &'static str {
1999 match source {
2000 SessionSource::Interactive => "interactive",
2001 SessionSource::Resume => "resume",
2002 SessionSource::Fork => "fork",
2003 SessionSource::Api => "api",
2004 SessionSource::Unknown => "unknown",
2005 }
2006 }
2007
2008 fn session_source_from_str(value: &str) -> SessionSource {
2009 match value {
2010 "interactive" => SessionSource::Interactive,
2011 "resume" => SessionSource::Resume,
2012 "fork" => SessionSource::Fork,
2013 "api" => SessionSource::Api,
2014 _ => SessionSource::Unknown,
2015 }
2016 }
2017
2018 fn path_to_opt_string(path: Option<&Path>) -> Option<String> {
2019 path.map(|p| p.display().to_string())
2020 }
2021
2022 fn parse_checkpoint_state(state_json: &str) -> Result<Value> {
2023 serde_json::from_str(state_json).context("failed to parse checkpoint state json")
2024 }
2025
2026 fn job_state_status_to_str(status: &JobStateStatus) -> &'static str {
2027 match status {
2028 JobStateStatus::Queued => "queued",
2029 JobStateStatus::Running => "running",
2030 JobStateStatus::Paused => "paused",
2031 JobStateStatus::Completed => "completed",
2032 JobStateStatus::Failed => "failed",
2033 JobStateStatus::Cancelled => "cancelled",
2034 }
2035 }
2036
2037 fn job_state_status_from_str(value: &str) -> JobStateStatus {
2038 match value {
2039 "queued" => JobStateStatus::Queued,
2040 "running" => JobStateStatus::Running,
2041 "paused" => JobStateStatus::Paused,
2042 "completed" => JobStateStatus::Completed,
2043 "failed" => JobStateStatus::Failed,
2044 "cancelled" => JobStateStatus::Cancelled,
2045 _ => JobStateStatus::Queued,
2046 }
2047 }
2048
2049 fn thread_goal_status_to_str(status: &ThreadGoalStatus) -> &'static str {
2050 match status {
2051 ThreadGoalStatus::Active => "active",
2052 ThreadGoalStatus::Paused => "paused",
2053 ThreadGoalStatus::Blocked => "blocked",
2054 ThreadGoalStatus::UsageLimited => "usage_limited",
2055 ThreadGoalStatus::BudgetLimited => "budget_limited",
2056 ThreadGoalStatus::Complete => "complete",
2057 }
2058 }
2059
2060 fn thread_goal_status_from_str(value: &str) -> ThreadGoalStatus {
2061 match value {
2062 "active" => ThreadGoalStatus::Active,
2063 "paused" => ThreadGoalStatus::Paused,
2064 "blocked" => ThreadGoalStatus::Blocked,
2065 "usage_limited" => ThreadGoalStatus::UsageLimited,
2066 "budget_limited" => ThreadGoalStatus::BudgetLimited,
2067 "complete" => ThreadGoalStatus::Complete,
2068 // Fail closed: an unknown or corrupted persisted value must never
2069 // resurrect a self-driving goal. The user can inspect and explicitly
2070 // resume a paused goal after repairing or replacing the record.
2071 _ => ThreadGoalStatus::Paused,
2072 }
2073 }
2074
2075 fn row_to_thread(row: &rusqlite::Row<'_>) -> rusqlite::Result<ThreadMetadata> {
2076 let status_raw: String = row.get(7)?;
2077 let source_raw: String = row.get(11)?;
2078 let rollout_path: Option<String> = row.get(1)?;
2079 let path: Option<String> = row.get(8)?;
2080 Ok(ThreadMetadata {
2081 id: row.get(0)?,
2082 rollout_path: rollout_path.map(PathBuf::from),
2083 preview: row.get(2)?,
2084 ephemeral: i64_to_bool(row.get(3)?),
2085 model_provider: row.get(4)?,
2086 created_at: row.get(5)?,
2087 updated_at: row.get(6)?,
2088 status: thread_status_from_str(&status_raw),
2089 path: path.map(PathBuf::from),
2090 cwd: PathBuf::from(row.get::<_, String>(9)?),
2091 cli_version: row.get(10)?,
2092 source: session_source_from_str(&source_raw),
2093 name: row.get(12)?,
2094 sandbox_policy: row.get(13)?,
2095 approval_mode: row.get(14)?,
2096 archived: i64_to_bool(row.get(15)?),
2097 archived_at: row.get(16)?,
2098 git_sha: row.get(17)?,
2099 git_branch: row.get(18)?,
2100 git_origin_url: row.get(19)?,
2101 memory_mode: row.get(20)?,
2102 current_leaf_id: row.get(21)?,
2103 })
2104 }
2105
2106 fn row_to_thread_goal(row: &rusqlite::Row<'_>) -> rusqlite::Result<ThreadGoalRecord> {
2107 let status_raw: String = row.get(3)?;
2108 let pause_reason: Option<String> = row.get(13)?;
2109 let goal = ThreadGoalRecord {
2110 thread_id: row.get(0)?,
2111 goal_id: row.get(1)?,
2112 objective: row.get(2)?,
2113 status: thread_goal_status_from_str(&status_raw),
2114 token_budget: row.get(4)?,
2115 tokens_used: row.get(5)?,
2116 time_used_seconds: row.get(6)?,
2117 continuation_count: row.get(7)?,
2118 created_at: row.get(8)?,
2119 updated_at: row.get(9)?,
2120 last_gap_fingerprint: row.get(10)?,
2121 repeated_gap_count: row.get(11)?,
2122 last_gap_pass: row.get(12)?,
2123 pause_reason: pause_reason
2124 .map(|value| serde_json::from_str(&value))
2125 .transpose()
2126 .map_err(|error| {
2127 rusqlite::Error::FromSqlConversionFailure(
2128 13,
2129 rusqlite::types::Type::Text,
2130 Box::new(error),
2131 )
2132 })?,
2133 };
2134 codewhale_protocol::validate_goal_stall_state(
2135 goal.last_gap_fingerprint.as_deref(),
2136 goal.repeated_gap_count,
2137 goal.last_gap_pass,
2138 u32::try_from(goal.continuation_count.max(0)).unwrap_or(u32::MAX),
2139 )
2140 .map_err(|error| {
2141 rusqlite::Error::FromSqlConversionFailure(
2142 10,
2143 rusqlite::types::Type::Text,
2144 Box::new(std::io::Error::new(std::io::ErrorKind::InvalidData, error)),
2145 )
2146 })?;
2147 Ok(goal)
2148 }
2149
2150 impl ThreadGoalRecord {
2151 /// Same restore rule as `codewhale_protocol::ThreadGoal`: an Active
2152 /// record with an exhausted stall window is corrupt and reads as paused.
2153 fn normalize_restored_stall_state(&mut self) {
2154 if matches!(self.status, ThreadGoalStatus::Active)
2155 && self.repeated_gap_count >= codewhale_protocol::MAX_REPEATED_GAP_COUNT
2156 {
2157 self.status = ThreadGoalStatus::Paused;
2158 self.pause_reason = Some(codewhale_protocol::GoalPauseReason::NoProgress);
2159 }
2160 }
2161 }
2162
2163 #[cfg(test)]
2164 mod tests {
2165 use super::*;
2166 use serde_json::json;
2167 use std::sync::{Arc, Barrier, Mutex, mpsc};
2168 use std::thread;
2169 use std::time::{Duration, SystemTime, UNIX_EPOCH};
2170
2171 fn temp_state_dir(name: &str) -> PathBuf {
2172 let suffix = SystemTime::now()
2173 .duration_since(UNIX_EPOCH)
2174 .expect("system time")
2175 .as_nanos();
2176 let dir = std::env::temp_dir().join(format!(
2177 "codewhale-state-{name}-{}-{suffix}",
2178 std::process::id()
2179 ));
2180 fs::create_dir_all(&dir).expect("create temp state dir");
2181 dir
2182 }
2183
2184 fn temp_state_store(name: &str) -> StateStore {
2185 let dir = temp_state_dir(name);
2186 StateStore::open(Some(dir.join("state.db"))).expect("open state store")
2187 }
2188
2189 fn test_thread(id: &str) -> ThreadMetadata {
2190 ThreadMetadata {
2191 id: id.to_string(),
2192 rollout_path: None,
2193 preview: "test thread".to_string(),
2194 ephemeral: false,
2195 model_provider: "deepseek".to_string(),
2196 created_at: 10,
2197 updated_at: 10,
2198 status: ThreadStatus::Running,
2199 path: None,
2200 cwd: PathBuf::from("/tmp/codewhale"),
2201 cli_version: "0.0.0-test".to_string(),
2202 source: SessionSource::Interactive,
2203 name: None,
2204 sandbox_policy: None,
2205 approval_mode: None,
2206 archived: false,
2207 archived_at: None,
2208 git_sha: None,
2209 git_branch: None,
2210 git_origin_url: None,
2211 memory_mode: None,
2212 current_leaf_id: None,
2213 }
2214 }
2215
2216 fn test_goal(thread_id: &str, objective: &str) -> ThreadGoalRecord {
2217 ThreadGoalRecord {
2218 thread_id: thread_id.to_string(),
2219 goal_id: "goal-1".to_string(),
2220 objective: objective.to_string(),
2221 status: ThreadGoalStatus::Active,
2222 token_budget: Some(123),
2223 tokens_used: 7,
2224 time_used_seconds: 11,
2225 continuation_count: 0,
2226 last_gap_fingerprint: None,
2227 repeated_gap_count: 0,
2228 last_gap_pass: None,
2229 pause_reason: None,
2230 created_at: 100,
2231 updated_at: 101,
2232 }
2233 }
2234
2235 #[test]
2236 fn unknown_persisted_goal_status_fails_closed() {
2237 assert_eq!(
2238 thread_goal_status_from_str("future_or_corrupt_status"),
2239 ThreadGoalStatus::Paused
2240 );
2241 }
2242
2243 #[test]
2244 fn thread_goal_crud_round_trips_and_replaces() {
2245 let store = temp_state_store("thread-goal-crud");
2246 store
2247 .upsert_thread(&test_thread("thread-1"))
2248 .expect("upsert thread");
2249
2250 let goal = test_goal("thread-1", "Ship v0.8.59");
2251 store.upsert_thread_goal(&goal).expect("upsert goal");
2252 assert_eq!(
2253 store
2254 .get_thread_goal("thread-1")
2255 .expect("read goal")
2256 .as_ref(),
2257 Some(&goal)
2258 );
2259
2260 let mut replacement = test_goal("thread-1", "Ship v0.8.59 safely");
2261 replacement.goal_id = "goal-2".to_string();
2262 replacement.status = ThreadGoalStatus::BudgetLimited;
2263 replacement.token_budget = None;
2264 replacement.updated_at = 202;
2265 store
2266 .upsert_thread_goal(&replacement)
2267 .expect("replace goal");
2268 assert_eq!(
2269 store.get_thread_goal("thread-1").expect("read replacement"),
2270 Some(replacement)
2271 );
2272
2273 assert!(store.delete_thread_goal("thread-1").expect("delete goal"));
2274 assert!(
2275 store
2276 .get_thread_goal("thread-1")
2277 .expect("read empty")
2278 .is_none()
2279 );
2280 assert!(!store.delete_thread_goal("thread-1").expect("delete empty"));
2281 }
2282
2283 #[test]
2284 fn thread_goal_requires_existing_thread() {
2285 let store = temp_state_store("thread-goal-missing-thread");
2286 let err = store
2287 .upsert_thread_goal(&test_goal("missing-thread", "nope"))
2288 .expect_err("goal without a thread should fail");
2289 assert!(err.to_string().contains("thread missing-thread not found"));
2290 }
2291
2292 /// Audit R03-07: two first opens of one fresh database. The loser must
2293 /// decide from the winner's committed schema, not from a probe taken
2294 /// before it could write. The old code read "column missing" outside its
2295 /// transaction, waited out the winner, then failed its own `ADD COLUMN`
2296 /// with "duplicate column name".
2297 #[test]
2298 fn a_first_open_that_loses_the_migration_race_still_opens() {
2299 let dir = temp_state_dir("migration-race");
2300 let path = dir.join("state.db");
2301 let winner = Connection::open(&path).expect("open winner");
2302 let mode: String = winner
2303 .query_row("PRAGMA journal_mode=WAL;", [], |row| row.get(0))
2304 .expect("wal");
2305 assert_eq!(mode, "wal");
2306 // The winner is mid-migration: the v1 columns exist inside its
2307 // uncommitted write transaction, user_version is still 0.
2308 winner
2309 .execute_batch(
2310 r#"
2311 BEGIN IMMEDIATE;
2312 CREATE TABLE threads (
2313 id TEXT PRIMARY KEY, rollout_path TEXT, preview TEXT NOT NULL,
2314 ephemeral INTEGER NOT NULL, model_provider TEXT NOT NULL,
2315 created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL,
2316 status TEXT NOT NULL, path TEXT, cwd TEXT NOT NULL,
2317 cli_version TEXT NOT NULL, source TEXT NOT NULL, title TEXT,
2318 sandbox_policy TEXT, approval_mode TEXT,
2319 archived INTEGER NOT NULL DEFAULT 0, archived_at INTEGER,
2320 git_sha TEXT, git_branch TEXT, git_origin_url TEXT,
2321 memory_mode TEXT, current_leaf_id INTEGER NULL
2322 );
2323 CREATE TABLE messages (
2324 id INTEGER PRIMARY KEY AUTOINCREMENT, thread_id TEXT NOT NULL,
2325 role TEXT NOT NULL, content TEXT NOT NULL, item_json TEXT,
2326 created_at INTEGER NOT NULL, parent_entry_id INTEGER NULL,
2327 FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE
2328 );
2329 "#,
2330 )
2331 .expect("winner begins its migration");
2332
2333 let loser = {
2334 let path = path.clone();
2335 thread::spawn(move || StateStore::open(Some(path)).map(|_| ()))
2336 };
2337 // Long enough for the loser to probe the schema; well inside its
2338 // five-second busy timeout.
2339 thread::sleep(Duration::from_millis(500));
2340 winner.execute_batch("COMMIT;").expect("winner commits");
2341
2342 loser
2343 .join()
2344 .expect("loser thread")
2345 .expect("the losing first open must succeed");
2346 }
2347
2348 /// Audit R05-05: a batch append is all or nothing. The pattern it
2349 /// replaces (one `append_message` transaction per item) leaves a
2350 /// partial chain when a later item fails.
2351 #[test]
2352 fn append_messages_is_all_or_nothing() {
2353 let store = temp_state_store("append-batch-atomic");
2354 store
2355 .upsert_thread(&test_thread("thread-1"))
2356 .expect("upsert thread");
2357 store
2358 .conn()
2359 .expect("conn")
2360 .execute_batch(
2361 "CREATE TRIGGER refuse_boom BEFORE INSERT ON messages \
2362 WHEN NEW.content = 'boom' BEGIN SELECT RAISE(ABORT, 'boom'); END;",
2363 )
2364 .expect("install failing trigger");
2365 let batch = ["one", "boom"].map(|content| NewMessage {
2366 role: "history".to_string(),
2367 content: content.to_string(),
2368 item: None,
2369 });
2370
2371 // The per-item pattern leaves "one" behind.
2372 let per_item: Result<Vec<i64>> = batch
2373 .iter()
2374 .map(|message| store.append_message("thread-1", &message.role, &message.content, None))
2375 .collect();
2376 assert!(per_item.is_err());
2377 let partial = store.list_messages("thread-1", None).expect("list");
2378 assert_eq!(partial.len(), 1, "per-item appends leave a partial chain");
2379 store.clear_messages("thread-1").expect("reset");
2380
2381 assert!(store.append_messages("thread-1", &batch).is_err());
2382 assert!(
2383 store
2384 .list_messages("thread-1", None)
2385 .expect("list")
2386 .is_empty(),
2387 "a failed batch must leave no messages"
2388 );
2389 let leaf: Option<i64> = store
2390 .conn()
2391 .expect("conn")
2392 .query_row(
2393 "SELECT current_leaf_id FROM threads WHERE id = 'thread-1'",
2394 [],
2395 |row| row.get(0),
2396 )
2397 .expect("leaf");
2398 assert_eq!(leaf, None, "a failed batch must not move the leaf");
2399
2400 let ok = ["a", "b", "c"].map(|content| NewMessage {
2401 role: "history".to_string(),
2402 content: content.to_string(),
2403 item: None,
2404 });
2405 let ids = store
2406 .append_messages("thread-1", &ok)
2407 .expect("append batch");
2408 let chain = store.list_messages("thread-1", None).expect("list");
2409 assert_eq!(chain.iter().map(|m| m.id).collect::<Vec<_>>(), ids);
2410 assert_eq!(chain[1].parent_entry_id, Some(ids[0]));
2411 assert_eq!(chain[2].parent_entry_id, Some(ids[1]));
2412 }
2413
2414 /// Audit R03-09: a deleted thread must not keep answering for its name
2415 /// in the session index, shadowing a live thread of the same name.
2416 #[test]
2417 fn deleted_thread_leaves_the_session_index() {
2418 let store = temp_state_store("delete-index");
2419 store
2420 .append_thread_name(
2421 "live",
2422 Some("Build".to_string()),
2423 100,
2424 Some(PathBuf::from("/live")),
2425 )
2426 .expect("index live");
2427 store
2428 .append_thread_name(
2429 "gone",
2430 Some("build".to_string()),
2431 200,
2432 Some(PathBuf::from("/gone")),
2433 )
2434 .expect("index gone");
2435 assert_eq!(
2436 store.find_thread_path_by_name_str("build").expect("find"),
2437 Some(PathBuf::from("/gone"))
2438 );
2439
2440 store.delete_thread("gone").expect("delete");
2441
2442 assert_eq!(
2443 store.find_thread_path_by_name_str("build").expect("find"),
2444 Some(PathBuf::from("/live"))
2445 );
2446 }
2447
2448 #[test]
2449 fn delete_thread_cascades_child_rows() {
2450 let store = temp_state_store("thread-delete-cascade");
2451 store
2452 .upsert_thread(&test_thread("thread-1"))
2453 .expect("upsert thread");
2454 store
2455 .append_message("thread-1", "user", "hello", None)
2456 .expect("append message");
2457 store
2458 .save_checkpoint("thread-1", "checkpoint-1", &serde_json::json!({"ok": true}))
2459 .expect("save checkpoint");
2460 store
2461 .persist_dynamic_tools(
2462 "thread-1",
2463 &[DynamicToolRecord {
2464 position: 0,
2465 name: "test_tool".to_string(),
2466 description: Some("test".to_string()),
2467 input_schema: serde_json::json!({"type": "object"}),
2468 }],
2469 )
2470 .expect("persist dynamic tools");
2471 store
2472 .upsert_thread_goal(&test_goal("thread-1", "Ship v0.8.67"))
2473 .expect("upsert goal");
2474
2475 store.delete_thread("thread-1").expect("delete thread");
2476
2477 let conn = store.conn().expect("conn");
2478 for table in [
2479 "messages",
2480 "checkpoints",
2481 "thread_dynamic_tools",
2482 "thread_goals",
2483 ] {
2484 let sql = format!("SELECT COUNT(*) FROM {table} WHERE thread_id = ?1");
2485 let count: i64 = conn
2486 .query_row(&sql, params!["thread-1"], |row| row.get(0))
2487 .expect("count child rows");
2488 assert_eq!(count, 0, "{table} row survived thread deletion");
2489 }
2490 }
2491
2492 #[test]
2493 fn state_store_reuses_one_connection_across_operations_and_clones() {
2494 let store = temp_state_store("conn-reuse");
2495 {
2496 let conn = store.conn().expect("conn");
2497 conn.execute_batch("CREATE TEMP TABLE conn_reuse_probe(id INTEGER);")
2498 .expect("create temp table");
2499 }
2500 // TEMP tables are visible only on the connection that created them, so
2501 // seeing the probe again — through a clone, after real operations ran —
2502 // proves the store holds one long-lived connection instead of
2503 // reopening the database (and reapplying pragmas) per call.
2504 let clone = store.clone();
2505 clone
2506 .upsert_thread(&test_thread("thread-conn-reuse"))
2507 .expect("upsert thread");
2508 let conn = clone.conn().expect("conn");
2509 let probe_count: i64 = conn
2510 .query_row(
2511 "SELECT COUNT(*) FROM sqlite_temp_master WHERE name = 'conn_reuse_probe'",
2512 [],
2513 |row| row.get(0),
2514 )
2515 .expect("query temp master");
2516 assert_eq!(
2517 probe_count, 1,
2518 "temp table not visible: a fresh connection was opened"
2519 );
2520 // The pragma applied once at open still governs the shared connection.
2521 let foreign_keys: i64 = conn
2522 .query_row("PRAGMA foreign_keys;", [], |row| row.get(0))
2523 .expect("read foreign_keys pragma");
2524 assert_eq!(foreign_keys, 1);
2525 let journal_mode: String = conn
2526 .query_row("PRAGMA journal_mode;", [], |row| row.get(0))
2527 .expect("read journal_mode pragma");
2528 assert_eq!(
2529 journal_mode.to_ascii_lowercase(),
2530 "wal",
2531 "open should enable WAL for multi-process readers/writers"
2532 );
2533 }
2534
2535 #[test]
2536 fn connection_setup_waits_for_database_lock_before_enabling_wal() {
2537 let dir = temp_state_dir("locked-open");
2538 let db_path = dir.join("state.db");
2539
2540 let candidate = Connection::open(&db_path).expect("open candidate connection");
2541 // Do not let rusqlite's current default mask StateStore's own setup
2542 // contract: configure_connection must install the wait policy before
2543 // it performs any operation that can need a database lock.
2544 candidate
2545 .busy_timeout(Duration::ZERO)
2546 .expect("disable dependency default timeout");
2547 let blocker = Connection::open(&db_path).expect("open blocking connection");
2548 let (locked_tx, locked_rx) = mpsc::sync_channel(0);
2549 let blocker_thread = thread::spawn(move || {
2550 blocker
2551 .execute_batch("BEGIN EXCLUSIVE;")
2552 .expect("acquire exclusive database lock");
2553 locked_tx.send(()).expect("announce database lock");
2554 thread::sleep(Duration::from_millis(200));
2555 blocker
2556 .execute_batch("COMMIT;")
2557 .expect("release exclusive database lock");
2558 });
2559
2560 locked_rx.recv().expect("wait for database lock");
2561 StateStore::configure_connection(&candidate, &db_path)
2562 .expect("connection setup should wait for the brief database lock");
2563 blocker_thread.join().expect("blocking thread panicked");
2564
2565 let journal_mode: String = candidate
2566 .query_row("PRAGMA journal_mode;", [], |row| row.get(0))
2567 .expect("read journal_mode");
2568 assert_eq!(journal_mode.to_ascii_lowercase(), "wal");
2569
2570 drop(candidate);
2571 let _ = fs::remove_dir_all(dir);
2572 }
2573
2574 /// A second process must wait for a brief active writer instead of
2575 /// surfacing SQLITE_BUSY (#4734).
2576 ///
2577 /// The lock handoff is explicit: unlike a race between many autocommit
2578 /// writes, this proves the busy timeout while keeping the contention
2579 /// duration below its documented five-second bound on every platform.
2580 #[test]
2581 fn second_connection_waits_for_active_writer() {
2582 let dir = temp_state_dir("concurrent-write");
2583 let db_path = dir.join("state.db");
2584
2585 let store_a = StateStore::open(Some(db_path.clone())).expect("open store a");
2586 let store_b = StateStore::open(Some(db_path.clone())).expect("open store b");
2587 let (locked_tx, locked_rx) = mpsc::sync_channel(0);
2588 let (release_tx, release_rx) = mpsc::sync_channel(0);
2589
2590 let writer_a = thread::spawn(move || {
2591 let conn = store_a.conn().expect("connection a");
2592 conn.execute_batch(
2593 r#"
2594 BEGIN IMMEDIATE;
2595 INSERT INTO jobs(id, name, status, created_at, updated_at)
2596 VALUES ('job-a', 'writer-a', 'running', 0, 0);
2597 "#,
2598 )
2599 .expect("writer a should acquire the database write lock");
2600 locked_tx.send(()).expect("announce active writer");
2601 release_rx.recv().expect("wait to release active writer");
2602 conn.execute_batch("COMMIT;")
2603 .expect("writer a should commit");
2604 });
2605
2606 locked_rx.recv().expect("wait for active writer");
2607 let (attempting_tx, attempting_rx) = mpsc::sync_channel(0);
2608 let writer_b = thread::spawn(move || {
2609 attempting_tx.send(()).expect("announce second write");
2610 store_b.upsert_job(&JobStateRecord {
2611 id: "job-b".to_string(),
2612 name: "writer-b".to_string(),
2613 status: JobStateStatus::Running,
2614 progress: None,
2615 detail: Some("waited for writer a".to_string()),
2616 created_at: 1,
2617 updated_at: 1,
2618 })
2619 });
2620
2621 attempting_rx.recv().expect("wait for second write attempt");
2622 thread::sleep(Duration::from_millis(100));
2623 assert!(
2624 !writer_b.is_finished(),
2625 "second writer should still be waiting while the first holds the lock"
2626 );
2627 release_tx.send(()).expect("release active writer");
2628 writer_a.join().expect("writer a panicked");
2629 writer_b
2630 .join()
2631 .expect("writer b panicked")
2632 .expect("writer b should succeed after the lock is released");
2633
2634 let store = StateStore::open(Some(db_path)).expect("reopen for verify");
2635 let listed = store.list_jobs(Some(2)).expect("list jobs");
2636 assert_eq!(listed.len(), 2, "both writers should persist their jobs");
2637
2638 let _ = fs::remove_dir_all(dir);
2639 }
2640
2641 #[test]
2642 fn runtime_thread_link_round_trips_and_goes_with_its_thread() {
2643 let store = temp_state_store("runtime-thread-link");
2644 store
2645 .upsert_thread(&test_thread("thread-linked"))
2646 .expect("upsert thread");
2647 assert_eq!(
2648 store
2649 .get_runtime_thread_link("thread-linked")
2650 .expect("read"),
2651 None
2652 );
2653 store
2654 .set_runtime_thread_link("thread-linked", "thr_runtime_1")
2655 .expect("link");
2656 store
2657 .set_runtime_thread_link("thread-linked", "thr_runtime_2")
2658 .expect("relink");
2659 assert_eq!(
2660 store
2661 .get_runtime_thread_link("thread-linked")
2662 .expect("read")
2663 .as_deref(),
2664 Some("thr_runtime_2")
2665 );
2666 assert!(
2667 store
2668 .set_runtime_thread_link("no-such-thread", "thr_runtime_3")
2669 .is_err(),
2670 "a link needs an existing thread"
2671 );
2672 store.delete_thread("thread-linked").expect("delete thread");
2673 assert_eq!(
2674 store
2675 .get_runtime_thread_link("thread-linked")
2676 .expect("read"),
2677 None
2678 );
2679 }
2680
2681 #[test]
2682 fn migration_runs_cleanly_when_schema_predates_user_version_header() {
2683 // Simulate a restore (or a racing process that crashed before
2684 // stamping user_version): the on-disk schema is fully migrated but
2685 // the header still says 0. The v0 block used to re-run unconditional
2686 // ADD COLUMN statements and abort the open with
2687 // "duplicate column name".
2688 let dir = temp_state_dir("migration-v0-idempotent");
2689 let db_path = dir.join("state.db");
2690 drop(StateStore::open(Some(db_path.clone())).expect("initial open"));
2691 {
2692 let conn = Connection::open(&db_path).expect("raw connection");
2693 conn.pragma_update(None, "user_version", 0)
2694 .expect("reset user_version");
2695 }
2696
2697 let store = StateStore::open(Some(db_path.clone())).expect("reopen with v0 header");
2698 store
2699 .upsert_thread(&test_thread("thread-migrated"))
2700 .expect("write after guarded migration");
2701
2702 // Reopening again (now stamped at the current version) still works.
2703 drop(store);
2704 let store = StateStore::open(Some(db_path)).expect("third open");
2705 let persisted = store
2706 .get_thread("thread-migrated")
2707 .expect("read after reopen");
2708 assert!(persisted.is_some());
2709
2710 let _ = fs::remove_dir_all(dir);
2711 }
2712
2713 #[test]
2714 fn record_thread_goal_usage_accumulates_tokens_and_time() {
2715 let store = temp_state_store("thread-goal-usage");
2716 store
2717 .upsert_thread(&test_thread("thread-1"))
2718 .expect("upsert thread");
2719
2720 // Mirror the runtime, which creates goals with zeroed accounting.
2721 let mut goal = test_goal("thread-1", "Ship the persistent goal loop");
2722 goal.tokens_used = 0;
2723 goal.time_used_seconds = 0;
2724 goal.updated_at = 100;
2725 store.upsert_thread_goal(&goal).expect("upsert goal");
2726
2727 // First accrual lands the deltas and advances updated_at.
2728 let after_first = store
2729 .record_thread_goal_usage("thread-1", 250, 12, 150)
2730 .expect("record usage")
2731 .expect("goal exists");
2732 assert_eq!(after_first.tokens_used, 250);
2733 assert_eq!(after_first.time_used_seconds, 12);
2734 assert_eq!(after_first.updated_at, 150);
2735 // Identity fields are preserved across accrual.
2736 assert_eq!(after_first.goal_id, goal.goal_id);
2737 assert_eq!(after_first.objective, goal.objective);
2738 assert_eq!(after_first.status, goal.status);
2739 assert_eq!(after_first.token_budget, goal.token_budget);
2740 assert_eq!(after_first.created_at, goal.created_at);
2741 assert_eq!(after_first.continuation_count, 0);
2742
2743 // Second accrual adds on top of the first (additive, not replacing).
2744 let after_second = store
2745 .record_thread_goal_usage("thread-1", 75, 8, 200)
2746 .expect("record usage")
2747 .expect("goal exists");
2748 assert_eq!(after_second.tokens_used, 325);
2749 assert_eq!(after_second.time_used_seconds, 20);
2750 assert_eq!(after_second.updated_at, 200);
2751
2752 // A stale `now` must not move updated_at backwards.
2753 let after_stale = store
2754 .record_thread_goal_usage("thread-1", 5, 1, 1)
2755 .expect("record usage")
2756 .expect("goal exists");
2757 assert_eq!(after_stale.tokens_used, 330);
2758 assert_eq!(after_stale.time_used_seconds, 21);
2759 assert_eq!(after_stale.updated_at, 200);
2760
2761 // Read back through the normal getter to confirm durability.
2762 let persisted = store
2763 .get_thread_goal("thread-1")
2764 .expect("read goal")
2765 .expect("goal exists");
2766 assert_eq!(persisted.tokens_used, 330);
2767 assert_eq!(persisted.time_used_seconds, 21);
2768 }
2769
2770 #[test]
2771 fn record_thread_goal_usage_returns_none_without_goal() {
2772 let store = temp_state_store("thread-goal-usage-missing");
2773 store
2774 .upsert_thread(&test_thread("thread-1"))
2775 .expect("upsert thread");
2776 // Thread exists but has no goal row yet: accrual is a no-op, not an error,
2777 // and must not create a goal.
2778 let result = store
2779 .record_thread_goal_usage("thread-1", 100, 5, 999)
2780 .expect("record usage on goalless thread");
2781 assert!(result.is_none());
2782 assert!(
2783 store
2784 .get_thread_goal("thread-1")
2785 .expect("read goal")
2786 .is_none()
2787 );
2788 }
2789
2790 #[test]
2791 fn record_thread_goal_continuation_accumulates_durably() {
2792 let store = temp_state_store("thread-goal-continuation");
2793 store
2794 .upsert_thread(&test_thread("thread-1"))
2795 .expect("upsert thread");
2796
2797 let mut goal = test_goal("thread-1", "Keep working across turns");
2798 goal.updated_at = 100;
2799 store.upsert_thread_goal(&goal).expect("upsert goal");
2800
2801 let after_first = store
2802 .record_thread_goal_continuation("thread-1", 120)
2803 .expect("record continuation")
2804 .expect("goal exists");
2805 assert_eq!(after_first.continuation_count, 1);
2806 assert_eq!(after_first.tokens_used, goal.tokens_used);
2807 assert_eq!(after_first.time_used_seconds, goal.time_used_seconds);
2808 assert_eq!(after_first.updated_at, 120);
2809
2810 let after_second = store
2811 .record_thread_goal_continuation("thread-1", 110)
2812 .expect("record second continuation")
2813 .expect("goal exists");
2814 assert_eq!(after_second.continuation_count, 2);
2815 assert_eq!(after_second.updated_at, 120);
2816
2817 let persisted = store
2818 .get_thread_goal("thread-1")
2819 .expect("read goal")
2820 .expect("goal exists");
2821 assert_eq!(persisted.continuation_count, 2);
2822 }
2823
2824 // ── $CODEWHALE_HOME override tests ──────────────────────────────
2825 //
2826 // These touch a process-global env var, so they serialize against each
2827 // other (and restore the prior value) to stay hermetic under parallel test
2828 // runs — the same concern AGENTS.md flags for config_command_allow_shell_*.
2829
2830 static CODEWHALE_HOME_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
2831
2832 struct CodeWhaleHomeGuard {
2833 prior: Option<std::ffi::OsString>,
2834 }
2835 impl CodeWhaleHomeGuard {
2836 fn set(value: &str) -> Self {
2837 let prior = std::env::var_os("CODEWHALE_HOME");
2838 // SAFETY: serialised by CODEWHALE_HOME_TEST_LOCK.
2839 unsafe { std::env::set_var("CODEWHALE_HOME", value) };
2840 Self { prior }
2841 }
2842 fn remove() -> Self {
2843 let prior = std::env::var_os("CODEWHALE_HOME");
2844 // SAFETY: serialised by CODEWHALE_HOME_TEST_LOCK.
2845 unsafe { std::env::remove_var("CODEWHALE_HOME") };
2846 Self { prior }
2847 }
2848 }
2849 impl Drop for CodeWhaleHomeGuard {
2850 fn drop(&mut self) {
2851 // SAFETY: serialised by CODEWHALE_HOME_TEST_LOCK.
2852 unsafe {
2853 match &self.prior {
2854 Some(value) => std::env::set_var("CODEWHALE_HOME", value),
2855 None => std::env::remove_var("CODEWHALE_HOME"),
2856 }
2857 }
2858 }
2859 }
2860
2861 #[test]
2862 fn codewhale_home_override_returns_the_env_value_verbatim() {
2863 let _lock = CODEWHALE_HOME_TEST_LOCK.lock().unwrap();
2864 let override_path = std::env::temp_dir().join("cw-isolated-state");
2865 let _g = CodeWhaleHomeGuard::set(override_path.to_str().unwrap());
2866 // The env var IS the home dir — no ".codewhale" appended. This matches
2867 // codewhale_home() in config ($CODEWHALE_HOME=/x means home is /x).
2868 assert_eq!(
2869 codewhale_home_override().unwrap().as_deref(),
2870 Some(override_path.as_path())
2871 );
2872 }
2873
2874 #[test]
2875 fn codewhale_home_override_none_when_unset() {
2876 let _lock = CODEWHALE_HOME_TEST_LOCK.lock().unwrap();
2877 let _g = CodeWhaleHomeGuard::remove();
2878 assert!(codewhale_home_override().unwrap().is_none());
2879 }
2880
2881 #[test]
2882 fn codewhale_home_override_none_when_whitespace_only() {
2883 let _lock = CODEWHALE_HOME_TEST_LOCK.lock().unwrap();
2884 let _g = CodeWhaleHomeGuard::set(" ");
2885 assert!(
2886 codewhale_home_override().unwrap().is_none(),
2887 "whitespace-only CODEWHALE_HOME must not establish isolation"
2888 );
2889 }
2890
2891 #[test]
2892 fn default_state_db_path_uses_codewhale_home_when_set() {
2893 let _lock = CODEWHALE_HOME_TEST_LOCK.lock().unwrap();
2894 let dir = std::env::temp_dir().join(format!(
2895 "cw-home-state-{}-{}",
2896 std::process::id(),
2897 std::time::SystemTime::now()
2898 .duration_since(std::time::UNIX_EPOCH)
2899 .unwrap()
2900 .as_nanos()
2901 ));
2902 let _g = CodeWhaleHomeGuard::set(dir.to_str().unwrap());
2903 // Hard override: the DB is <CODEWHALE_HOME>/state.db, NOT
2904 // <CODEWHALE_HOME>/.codewhale/state.db, and the legacy ~/.deepseek
2905 // fallback is bypassed entirely.
2906 assert_eq!(default_state_db_path(), dir.join("state.db"));
2907 }
2908
2909 #[test]
2910 fn load_checkpoint_propagates_invalid_state_json() {
2911 let store = temp_state_store("checkpoint-parse-error");
2912 store
2913 .upsert_thread(&test_thread("thread-1"))
2914 .expect("upsert thread");
2915 store
2916 .save_checkpoint("thread-1", "broken", &json!({"ok": true}))
2917 .expect("save checkpoint");
2918
2919 {
2920 let conn = store.conn().expect("conn");
2921 conn.execute(
2922 "UPDATE checkpoints SET state_json = ?1 WHERE thread_id = ?2 AND checkpoint_id = ?3",
2923 params!["not-json", "thread-1", "broken"],
2924 )
2925 .expect("corrupt checkpoint");
2926 }
2927
2928 let err = store
2929 .load_checkpoint("thread-1", Some("broken"))
2930 .expect_err("invalid checkpoint json should fail");
2931 assert!(
2932 err.to_string()
2933 .contains("failed to parse checkpoint state json")
2934 );
2935 }
2936
2937 #[test]
2938 fn session_index_compacts_after_threshold() {
2939 let store = temp_state_store("session-index-compact");
2940 for idx in 0..6 {
2941 store
2942 .append_thread_name("thread-1", Some(format!("name-{idx}")), idx, None)
2943 .expect("append session index entry");
2944 }
2945
2946 let line_count = store
2947 .session_index_line_count()
2948 .expect("count session index lines");
2949 assert_eq!(line_count, 1);
2950
2951 let name = store
2952 .find_thread_name_by_id("thread-1")
2953 .expect("lookup thread name");
2954 assert_eq!(name.as_deref(), Some("name-5"));
2955 }
2956
2957 #[test]
2958 fn session_index_read_skips_a_torn_line() {
2959 // #4735: a crash mid-append leaves a truncated final line. Failing the
2960 // whole read broke every thread-name lookup at once, and compaction
2961 // reads through the same path, so the index could not repair itself.
2962 let store = temp_state_store("session-index-torn");
2963 store
2964 .append_thread_name("thread-1", Some("first".to_string()), 1, None)
2965 .expect("append first entry");
2966
2967 {
2968 let mut file = OpenOptions::new()
2969 .append(true)
2970 .open(&store.session_index_path)
2971 .expect("open session index");
2972 // A write cut off mid-JSON, exactly as a crash would leave it.
2973 writeln!(file, "{{\"thread_id\":\"thread-2\",\"thread_na").expect("write torn line");
2974 }
2975
2976 store
2977 .append_thread_name("thread-3", Some("third".to_string()), 3, None)
2978 .expect("append third entry");
2979
2980 assert_eq!(
2981 store
2982 .find_thread_name_by_id("thread-1")
2983 .expect("lookup thread-1")
2984 .as_deref(),
2985 Some("first"),
2986 );
2987 assert_eq!(
2988 store
2989 .find_thread_name_by_id("thread-3")
2990 .expect("lookup thread-3")
2991 .as_deref(),
2992 Some("third"),
2993 );
2994 }
2995
2996 /// Hook fired by compaction between writing the snapshot and renaming it
2997 /// over the live index — the window a concurrent append can be lost in.
2998 /// Only the store whose path a test registered is affected, so tests
2999 /// running in parallel don't disturb each other.
3000 type MidpointHook = Box<dyn Fn() + Send + Sync>;
3001 static COMPACTION_MIDPOINT: Mutex<Option<(PathBuf, MidpointHook)>> = Mutex::new(None);
3002
3003 /// Fires at most once: the racing append compacts too, and a hook that
3004 /// fired twice would re-enter the test's one-shot handshake.
3005 pub(super) fn compaction_midpoint(index_path: &Path) {
3006 let mut hook = COMPACTION_MIDPOINT.lock().expect("midpoint hook lock");
3007 let registered_for_this_store = match hook.as_ref() {
3008 Some((registered, _)) => registered == index_path,
3009 None => return,
3010 };
3011 if !registered_for_this_store {
3012 return;
3013 }
3014 let (_, callback) = hook.take().expect("presence checked above");
3015 drop(hook);
3016 callback();
3017 }
3018
3019 #[test]
3020 fn session_index_compaction_does_not_drop_a_concurrent_append() {
3021 // #4736: compaction snapshots the file, rewrites it, and renames over
3022 // the live one. An append landing between those two steps used to
3023 // vanish — silently, since it had already returned success to its
3024 // caller. The shared lock serializes the two.
3025 //
3026 // The race is real but narrow, so the test drives it deterministically:
3027 // a hook at the compaction midpoint releases the appender and then
3028 // waits. Without the lock the appender writes into the doomed file and
3029 // the rename discards it; with the lock it blocks until compaction
3030 // finishes, and its entry survives.
3031 let store = Arc::new(temp_state_store("session-index-race"));
3032 let threshold = session_index_compact_line_threshold();
3033
3034 // One line short of the threshold, so the next append compacts.
3035 for idx in 0..threshold {
3036 store
3037 .append_thread_name(
3038 &format!("thread-{idx}"),
3039 Some(format!("name-{idx}")),
3040 1,
3041 None,
3042 )
3043 .expect("append filler entry");
3044 }
3045
3046 let appender_released = Arc::new(Barrier::new(2));
3047 {
3048 let released = Arc::clone(&appender_released);
3049 *COMPACTION_MIDPOINT.lock().expect("midpoint hook lock") = Some((
3050 store.session_index_path.clone(),
3051 Box::new(move || {
3052 released.wait();
3053 // Give the appender time to complete its write into the
3054 // window. Under the fix it is blocked on the lock instead.
3055 thread::sleep(Duration::from_millis(300));
3056 }),
3057 ));
3058 }
3059
3060 let appender = {
3061 let store = Arc::clone(&store);
3062 let released = Arc::clone(&appender_released);
3063 thread::spawn(move || {
3064 released.wait();
3065 store
3066 .append_thread_name("racer", Some("racer-name".to_string()), 2, None)
3067 .expect("append racing entry");
3068 })
3069 };
3070
3071 store
3072 .append_thread_name("trigger", Some("trigger-name".to_string()), 1, None)
3073 .expect("append entry that triggers compaction");
3074 appender.join().expect("appender thread");
3075 *COMPACTION_MIDPOINT.lock().expect("midpoint hook lock") = None;
3076
3077 assert_eq!(
3078 store
3079 .find_thread_name_by_id("racer")
3080 .expect("lookup racer")
3081 .as_deref(),
3082 Some("racer-name"),
3083 "append was dropped by a concurrent compaction",
3084 );
3085 assert_eq!(
3086 store
3087 .find_thread_name_by_id("trigger")
3088 .expect("lookup trigger")
3089 .as_deref(),
3090 Some("trigger-name"),
3091 );
3092 }
3093 fn legacy_archive_fixture(id: &str) -> LegacyThreadArchive {
3094 let mut thread = test_thread(id);
3095 thread.current_leaf_id = Some(3);
3096 LegacyThreadArchive {
3097 thread,
3098 messages: [
3099 (1, None, "root"),
3100 (2, Some(1), "inactive"),
3101 (3, Some(1), "active"),
3102 ]
3103 .into_iter()
3104 .map(|(entry_id, parent_entry_id, content)| MessageRecord {
3105 id: entry_id,
3106 thread_id: id.into(),
3107 role: "user".into(),
3108 content: content.into(),
3109 item: Some(serde_json::json!({"content":content})),
3110 created_at: 100 + entry_id,
3111 parent_entry_id,
3112 })
3113 .collect(),
3114 goal: Some(test_goal(id, "historical goal")),
3115 checkpoints: vec![CheckpointRecord {
3116 thread_id: id.into(),
3117 checkpoint_id: "retained".into(),
3118 state: serde_json::json!({"receipt":"historical"}),
3119 created_at: 123,
3120 }],
3121 }
3122 }
3123
3124 #[test]
3125 fn legacy_archive_restore_keeps_full_graph_goal_checkpoint_and_refuses_replacement() {
3126 let store = temp_state_store("archive-once");
3127 let archive = legacy_archive_fixture("legacy");
3128 store.restore_legacy_thread_archive(&archive).unwrap();
3129 let snapshot = store.snapshot_legacy_thread_history("legacy").unwrap();
3130 assert_eq!(
3131 serde_json::to_value(&snapshot.messages).unwrap(),
3132 serde_json::to_value(&archive.messages).unwrap()
3133 );
3134 assert_eq!(snapshot.current_leaf_id, Some(3));
3135 assert_eq!(store.list_messages("legacy", None).unwrap().len(), 2);
3136 let goal = snapshot.goal.unwrap();
3137 assert_eq!(goal.status, codewhale_protocol::ThreadGoalStatus::Active);
3138 assert_eq!(goal.tokens_used, archive.goal.as_ref().unwrap().tokens_used);
3139 assert_eq!(goal.updated_at, archive.goal.as_ref().unwrap().updated_at);
3140 let checkpoint = store
3141 .load_checkpoint("legacy", Some("retained"))
3142 .unwrap()
3143 .unwrap();
3144 assert_eq!(checkpoint.created_at, 123);
3145 assert_eq!(checkpoint.state, archive.checkpoints[0].state);
3146 let mut replacement = archive.clone();
3147 replacement.messages[1].content = "changed inactive history".into();
3148 assert!(store.restore_legacy_thread_archive(&replacement).is_err());
3149 assert_eq!(
3150 store
3151 .snapshot_legacy_thread_history("legacy")
3152 .unwrap()
3153 .messages[1]
3154 .content,
3155 "inactive"
3156 );
3157 }
3158
3159 #[test]
3160 fn legacy_archive_restore_refuses_cycles_duplicates_foreign_rows_and_bounds_without_publication()
3161 {
3162 let store = temp_state_store("archive-refusals");
3163 let valid = legacy_archive_fixture("legacy");
3164 let mut invalid = Vec::new();
3165 let mut archive = valid.clone();
3166 archive.messages[0].parent_entry_id = Some(2);
3167 invalid.push(archive);
3168 let mut archive = valid.clone();
3169 archive.messages[1].id = 1;
3170 invalid.push(archive);
3171 let mut archive = valid.clone();
3172 archive.messages[1].parent_entry_id = Some(99);
3173 invalid.push(archive);
3174 let mut archive = valid.clone();
3175 archive.messages[1].thread_id = "foreign".into();
3176 invalid.push(archive);
3177 let mut archive = valid.clone();
3178 archive.thread.current_leaf_id = Some(99);
3179 invalid.push(archive);
3180 let mut archive = valid.clone();
3181 archive.goal.as_mut().unwrap().thread_id = "foreign".into();
3182 invalid.push(archive);
3183 let mut archive = valid.clone();
3184 archive.checkpoints[0].thread_id = "foreign".into();
3185 invalid.push(archive);
3186 let mut archive = valid.clone();
3187 archive.checkpoints.push(archive.checkpoints[0].clone());
3188 invalid.push(archive);
3189 let mut archive = valid.clone();
3190 archive.messages[0].content =
3191 "x".repeat(codewhale_protocol::MAX_CANONICAL_HISTORY_BYTES + 1);
3192 invalid.push(archive);
3193 for archive in invalid {
3194 assert!(store.restore_legacy_thread_archive(&archive).is_err());
3195 assert!(store.get_thread("legacy").unwrap().is_none());
3196 }
3197 store.restore_legacy_thread_archive(&valid).unwrap();
3198 }
3199
3200 #[test]
3201 fn legacy_archive_restore_rolls_back_metadata_messages_goal_on_late_checkpoint_failure() {
3202 let store = temp_state_store("archive-transaction");
3203 store.conn().unwrap().execute_batch("CREATE TRIGGER refuse_archive_checkpoint BEFORE INSERT ON checkpoints BEGIN SELECT RAISE(ABORT, 'fixture checkpoint failure'); END;").unwrap();
3204 let archive = legacy_archive_fixture("legacy");
3205 assert!(store.restore_legacy_thread_archive(&archive).is_err());
3206 let conn = store.conn().unwrap();
3207 for table in ["threads", "messages", "thread_goals", "checkpoints"] {
3208 let count: i64 = conn
3209 .query_row(&format!("SELECT count(*) FROM {table}"), [], |row| {
3210 row.get(0)
3211 })
3212 .unwrap();
3213 assert_eq!(
3214 count, 0,
3215 "{table} must roll back with the immutable archive"
3216 );
3217 }
3218 conn.execute_batch("DROP TRIGGER refuse_archive_checkpoint")
3219 .unwrap();
3220 drop(conn);
3221 store.restore_legacy_thread_archive(&archive).unwrap();
3222 }
3223 }
3224
3224 lines RUST