| 1 | use std::path::PathBuf; |
| 2 | |
| 3 | use codewhale_state::{SessionSource, StateStore, ThreadListFilters, ThreadMetadata, ThreadStatus}; |
| 4 | use rusqlite::Connection; |
| 5 | |
| 6 | fn temp_state_path(label: &str) -> PathBuf { |
| 7 | std::env::temp_dir().join(format!( |
| 8 | "deepseek_state_test_{}_{}_{}.db", |
| 9 | label, |
| 10 | std::process::id(), |
| 11 | chrono::Utc::now().timestamp_nanos_opt().unwrap_or(0) |
| 12 | )) |
| 13 | } |
| 14 | |
| 15 | fn assert_workflow_trace_schema(conn: &Connection) { |
| 16 | let user_version: u32 = conn |
| 17 | .query_row("PRAGMA user_version;", [], |row| row.get(0)) |
| 18 | .expect("read user_version"); |
| 19 | // v5 (goal stall-history migration) adds `thread_goals.last_gap_fingerprint`, |
| 20 | // `repeated_gap_count`, `last_gap_pass` and `pause_reason` on top of the v4 |
| 21 | // continuation-count column and the v3 workflow-trace + thread_goals tables. |
| 22 | // v6 adds `thread_runtime_links` (client thread -> runtime thread). |
| 23 | // v7 adds bound canonical alias/operation receipts on this same connection. |
| 24 | assert_eq!(user_version, 7); |
| 25 | |
| 26 | for table in [ |
| 27 | "workflow_runs", |
| 28 | "branch_runs", |
| 29 | "leaf_runs", |
| 30 | "control_node_runs", |
| 31 | "teacher_candidates", |
| 32 | "thread_goals", |
| 33 | "thread_runtime_links", |
| 34 | "state_store_identity", |
| 35 | "thread_runtime_receipts", |
| 36 | "thread_runtime_operations", |
| 37 | ] { |
| 38 | let exists: bool = conn |
| 39 | .query_row( |
| 40 | "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1)", |
| 41 | [table], |
| 42 | |row| row.get(0), |
| 43 | ) |
| 44 | .unwrap_or_else(|err| panic!("read sqlite_master for {table}: {err}")); |
| 45 | assert!(exists, "missing workflow trace table {table}"); |
| 46 | } |
| 47 | } |
| 48 | |
| 49 | #[test] |
| 50 | fn upsert_and_resume_thread_metadata() { |
| 51 | let path = temp_state_path("upsert_resume"); |
| 52 | let store = StateStore::open(Some(path.clone())).expect("open state store"); |
| 53 | let now = chrono::Utc::now().timestamp(); |
| 54 | let thread = ThreadMetadata { |
| 55 | id: "thread-test-1".to_string(), |
| 56 | rollout_path: Some(PathBuf::from("/tmp/rollout.jsonl")), |
| 57 | preview: "hello".to_string(), |
| 58 | ephemeral: false, |
| 59 | model_provider: "deepseek".to_string(), |
| 60 | created_at: now, |
| 61 | updated_at: now, |
| 62 | status: ThreadStatus::Running, |
| 63 | path: Some(PathBuf::from("/tmp/project")), |
| 64 | cwd: PathBuf::from("/tmp/project"), |
| 65 | cli_version: "0.0.0-test".to_string(), |
| 66 | source: SessionSource::Interactive, |
| 67 | name: Some("Test Thread".to_string()), |
| 68 | sandbox_policy: Some("workspace-write".to_string()), |
| 69 | approval_mode: Some("on-request".to_string()), |
| 70 | archived: false, |
| 71 | archived_at: None, |
| 72 | git_sha: None, |
| 73 | git_branch: None, |
| 74 | git_origin_url: None, |
| 75 | memory_mode: Some("extended".to_string()), |
| 76 | current_leaf_id: None, |
| 77 | }; |
| 78 | store.upsert_thread(&thread).expect("upsert thread"); |
| 79 | |
| 80 | let loaded = store |
| 81 | .get_thread("thread-test-1") |
| 82 | .expect("read thread") |
| 83 | .expect("thread must exist"); |
| 84 | assert_eq!(loaded.id, "thread-test-1"); |
| 85 | assert_eq!(loaded.name.as_deref(), Some("Test Thread")); |
| 86 | assert_eq!(loaded.memory_mode.as_deref(), Some("extended")); |
| 87 | assert_eq!( |
| 88 | loaded.rollout_path, |
| 89 | Some(PathBuf::from("/tmp/rollout.jsonl")) |
| 90 | ); |
| 91 | |
| 92 | store |
| 93 | .mark_archived("thread-test-1") |
| 94 | .expect("archive thread"); |
| 95 | let archived = store |
| 96 | .get_thread("thread-test-1") |
| 97 | .expect("read archived thread") |
| 98 | .expect("thread exists after archive"); |
| 99 | assert!(archived.archived); |
| 100 | |
| 101 | let listed = store |
| 102 | .list_threads(ThreadListFilters { |
| 103 | include_archived: true, |
| 104 | limit: Some(10), |
| 105 | }) |
| 106 | .expect("list threads"); |
| 107 | assert!(!listed.is_empty()); |
| 108 | } |
| 109 | |
| 110 | #[test] |
| 111 | fn init_schema_migration() { |
| 112 | let path = temp_state_path("init_schema_migration"); |
| 113 | let conn = Connection::open(&path).expect("open state db"); |
| 114 | conn.execute_batch( |
| 115 | r#" |
| 116 | CREATE TABLE IF NOT EXISTS threads ( |
| 117 | id TEXT PRIMARY KEY, |
| 118 | rollout_path TEXT, |
| 119 | preview TEXT NOT NULL, |
| 120 | ephemeral INTEGER NOT NULL, |
| 121 | model_provider TEXT NOT NULL, |
| 122 | created_at INTEGER NOT NULL, |
| 123 | updated_at INTEGER NOT NULL, |
| 124 | status TEXT NOT NULL, |
| 125 | path TEXT, |
| 126 | cwd TEXT NOT NULL, |
| 127 | cli_version TEXT NOT NULL, |
| 128 | source TEXT NOT NULL, |
| 129 | title TEXT, |
| 130 | sandbox_policy TEXT, |
| 131 | approval_mode TEXT, |
| 132 | archived INTEGER NOT NULL DEFAULT 0, |
| 133 | archived_at INTEGER, |
| 134 | git_sha TEXT, |
| 135 | git_branch TEXT, |
| 136 | git_origin_url TEXT, |
| 137 | memory_mode TEXT |
| 138 | ); |
| 139 | CREATE TABLE IF NOT EXISTS messages ( |
| 140 | id INTEGER PRIMARY KEY AUTOINCREMENT, |
| 141 | thread_id TEXT NOT NULL, |
| 142 | role TEXT NOT NULL, |
| 143 | content TEXT NOT NULL, |
| 144 | item_json TEXT, |
| 145 | created_at INTEGER NOT NULL, |
| 146 | FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE |
| 147 | ); |
| 148 | INSERT INTO threads ( |
| 149 | id, preview, ephemeral, model_provider, created_at, updated_at, status, cwd, cli_version, source, archived |
| 150 | ) |
| 151 | VALUES ( |
| 152 | 'thread-test-1', 'hello', false, 'deepseek', 0, 0, 'running', '/tmp/project', '0.0.0-test', 'interactive', false |
| 153 | ); |
| 154 | INSERT INTO messages (thread_id, role, content, created_at) VALUES |
| 155 | ('thread-test-1', 'foo0', 'bar0', 0), |
| 156 | ('thread-test-1', 'foo1', 'bar1', 1), |
| 157 | ('thread-test-1', 'foo2', 'bar2', 2); |
| 158 | "#, |
| 159 | ) |
| 160 | .expect("init schema migration"); |
| 161 | |
| 162 | let store = StateStore::open(Some(path.clone())).expect("open state store"); |
| 163 | let thread = store |
| 164 | .get_thread("thread-test-1") |
| 165 | .expect("read thread") |
| 166 | .unwrap(); |
| 167 | assert_eq!(thread.id, "thread-test-1"); |
| 168 | assert_eq!(thread.preview, "hello"); |
| 169 | assert!(!thread.ephemeral); |
| 170 | assert_eq!(thread.model_provider, "deepseek"); |
| 171 | assert_eq!(thread.created_at, 0); |
| 172 | assert_eq!(thread.updated_at, 0); |
| 173 | assert_eq!(thread.status, ThreadStatus::Running); |
| 174 | assert_eq!(thread.cwd, PathBuf::from("/tmp/project")); |
| 175 | assert_eq!(thread.cli_version, "0.0.0-test"); |
| 176 | assert_eq!(thread.source, SessionSource::Interactive); |
| 177 | assert!(thread.current_leaf_id.is_some()); |
| 178 | |
| 179 | let messages = store |
| 180 | .list_messages("thread-test-1", None) |
| 181 | .expect("list messages"); |
| 182 | assert_eq!(messages.len(), 3); |
| 183 | for (i, message) in messages.iter().enumerate() { |
| 184 | assert_eq!(message.thread_id, "thread-test-1"); |
| 185 | assert_eq!(message.role, format!("foo{i}")); |
| 186 | assert_eq!(message.content, format!("bar{i}")); |
| 187 | assert_eq!(message.created_at, i as i64); |
| 188 | } |
| 189 | |
| 190 | // Test idempotent |
| 191 | StateStore::open(Some(path.clone())).expect("open state store"); |
| 192 | } |
| 193 | |
| 194 | #[test] |
| 195 | fn fresh_schema_includes_workflow_trace_tables() { |
| 196 | let path = temp_state_path("fresh_schema_includes_workflow_trace_tables"); |
| 197 | |
| 198 | StateStore::open(Some(path.clone())).expect("open state store"); |
| 199 | |
| 200 | let conn = Connection::open(&path).expect("open state db"); |
| 201 | assert_workflow_trace_schema(&conn); |
| 202 | } |
| 203 | |
| 204 | #[test] |
| 205 | fn v1_schema_migrates_workflow_trace_tables() { |
| 206 | let path = temp_state_path("v1_schema_migrates_workflow_trace_tables"); |
| 207 | let conn = Connection::open(&path).expect("open state db"); |
| 208 | conn.execute_batch( |
| 209 | r#" |
| 210 | CREATE TABLE threads ( |
| 211 | id TEXT PRIMARY KEY, |
| 212 | rollout_path TEXT, |
| 213 | preview TEXT NOT NULL, |
| 214 | ephemeral INTEGER NOT NULL, |
| 215 | model_provider TEXT NOT NULL, |
| 216 | created_at INTEGER NOT NULL, |
| 217 | updated_at INTEGER NOT NULL, |
| 218 | status TEXT NOT NULL, |
| 219 | path TEXT, |
| 220 | cwd TEXT NOT NULL, |
| 221 | cli_version TEXT NOT NULL, |
| 222 | source TEXT NOT NULL, |
| 223 | title TEXT, |
| 224 | sandbox_policy TEXT, |
| 225 | approval_mode TEXT, |
| 226 | archived INTEGER NOT NULL DEFAULT 0, |
| 227 | archived_at INTEGER, |
| 228 | git_sha TEXT, |
| 229 | git_branch TEXT, |
| 230 | git_origin_url TEXT, |
| 231 | memory_mode TEXT, |
| 232 | current_leaf_id INTEGER |
| 233 | ); |
| 234 | CREATE TABLE messages ( |
| 235 | id INTEGER PRIMARY KEY AUTOINCREMENT, |
| 236 | thread_id TEXT NOT NULL, |
| 237 | role TEXT NOT NULL, |
| 238 | content TEXT NOT NULL, |
| 239 | item_json TEXT, |
| 240 | created_at INTEGER NOT NULL, |
| 241 | parent_entry_id INTEGER |
| 242 | ); |
| 243 | CREATE TABLE checkpoints ( |
| 244 | thread_id TEXT NOT NULL, |
| 245 | checkpoint_id TEXT NOT NULL, |
| 246 | state_json TEXT NOT NULL, |
| 247 | created_at INTEGER NOT NULL, |
| 248 | PRIMARY KEY(thread_id, checkpoint_id) |
| 249 | ); |
| 250 | CREATE TABLE jobs ( |
| 251 | id TEXT PRIMARY KEY, |
| 252 | name TEXT NOT NULL, |
| 253 | status TEXT NOT NULL, |
| 254 | progress INTEGER, |
| 255 | detail TEXT, |
| 256 | created_at INTEGER NOT NULL, |
| 257 | updated_at INTEGER NOT NULL |
| 258 | ); |
| 259 | CREATE TABLE thread_dynamic_tools ( |
| 260 | thread_id TEXT NOT NULL, |
| 261 | position INTEGER NOT NULL, |
| 262 | name TEXT NOT NULL, |
| 263 | description TEXT, |
| 264 | input_schema TEXT NOT NULL, |
| 265 | PRIMARY KEY (thread_id, position) |
| 266 | ); |
| 267 | INSERT INTO threads ( |
| 268 | id, preview, ephemeral, model_provider, created_at, updated_at, status, cwd, cli_version, source, archived |
| 269 | ) |
| 270 | VALUES ( |
| 271 | 'thread-test-1', 'hello', false, 'deepseek', 0, 0, 'running', '/tmp/project', '0.0.0-test', 'interactive', false |
| 272 | ); |
| 273 | PRAGMA user_version = 1; |
| 274 | "#, |
| 275 | ) |
| 276 | .expect("create v1 schema"); |
| 277 | drop(conn); |
| 278 | |
| 279 | let store = StateStore::open(Some(path.clone())).expect("open state store"); |
| 280 | let thread = store |
| 281 | .get_thread("thread-test-1") |
| 282 | .expect("read thread") |
| 283 | .expect("thread survives migration"); |
| 284 | assert_eq!(thread.preview, "hello"); |
| 285 | |
| 286 | let conn = Connection::open(&path).expect("open state db"); |
| 287 | assert_workflow_trace_schema(&conn); |
| 288 | } |
| 289 | |
| 290 | #[test] |
| 291 | fn init_schema_migration_same_second_messages() { |
| 292 | let path = temp_state_path("init_schema_migration_same_second_messages"); |
| 293 | let conn = Connection::open(&path).expect("open state db"); |
| 294 | conn.execute_batch( |
| 295 | r#" |
| 296 | CREATE TABLE IF NOT EXISTS threads ( |
| 297 | id TEXT PRIMARY KEY, |
| 298 | rollout_path TEXT, |
| 299 | preview TEXT NOT NULL, |
| 300 | ephemeral INTEGER NOT NULL, |
| 301 | model_provider TEXT NOT NULL, |
| 302 | created_at INTEGER NOT NULL, |
| 303 | updated_at INTEGER NOT NULL, |
| 304 | status TEXT NOT NULL, |
| 305 | path TEXT, |
| 306 | cwd TEXT NOT NULL, |
| 307 | cli_version TEXT NOT NULL, |
| 308 | source TEXT NOT NULL, |
| 309 | title TEXT, |
| 310 | sandbox_policy TEXT, |
| 311 | approval_mode TEXT, |
| 312 | archived INTEGER NOT NULL DEFAULT 0, |
| 313 | archived_at INTEGER, |
| 314 | git_sha TEXT, |
| 315 | git_branch TEXT, |
| 316 | git_origin_url TEXT, |
| 317 | memory_mode TEXT |
| 318 | ); |
| 319 | CREATE TABLE IF NOT EXISTS messages ( |
| 320 | id INTEGER PRIMARY KEY AUTOINCREMENT, |
| 321 | thread_id TEXT NOT NULL, |
| 322 | role TEXT NOT NULL, |
| 323 | content TEXT NOT NULL, |
| 324 | item_json TEXT, |
| 325 | created_at INTEGER NOT NULL, |
| 326 | FOREIGN KEY(thread_id) REFERENCES threads(id) ON DELETE CASCADE |
| 327 | ); |
| 328 | INSERT INTO threads ( |
| 329 | id, preview, ephemeral, model_provider, created_at, updated_at, status, cwd, cli_version, source, archived |
| 330 | ) |
| 331 | VALUES ( |
| 332 | 'thread-test-2', 'hello', false, 'deepseek', 0, 0, 'running', '/tmp/project', '0.0.0-test', 'interactive', false |
| 333 | ); |
| 334 | INSERT INTO messages (thread_id, role, content, created_at) VALUES |
| 335 | ('thread-test-2', 'foo0', 'bar0', 123), |
| 336 | ('thread-test-2', 'foo1', 'bar1', 123), |
| 337 | ('thread-test-2', 'foo2', 'bar2', 123), |
| 338 | ('thread-test-2', 'foo3', 'bar3', 123); |
| 339 | "#, |
| 340 | ) |
| 341 | .expect("init schema migration"); |
| 342 | |
| 343 | let store = StateStore::open(Some(path.clone())).expect("open state store"); |
| 344 | let messages = store |
| 345 | .list_messages("thread-test-2", None) |
| 346 | .expect("list messages"); |
| 347 | assert_eq!(messages.len(), 4); |
| 348 | for (i, message) in messages.iter().enumerate() { |
| 349 | assert_eq!(message.thread_id, "thread-test-2"); |
| 350 | assert_eq!(message.role, format!("foo{i}")); |
| 351 | assert_eq!(message.content, format!("bar{i}")); |
| 352 | assert_eq!(message.created_at, 123); |
| 353 | } |
| 354 | assert_eq!(messages[0].parent_entry_id, None); |
| 355 | assert_eq!(messages[1].parent_entry_id, Some(messages[0].id)); |
| 356 | assert_eq!(messages[2].parent_entry_id, Some(messages[1].id)); |
| 357 | assert_eq!(messages[3].parent_entry_id, Some(messages[2].id)); |
| 358 | |
| 359 | // Test idempotent reopen after same-second parent links are migrated. |
| 360 | StateStore::open(Some(path.clone())).expect("open state store - idempotent"); |
| 361 | } |
| 362 | |
| 363 | #[test] |
| 364 | fn test_fork() { |
| 365 | // Historical branching remains readable in every selected projection; live |
| 366 | // fork/mutation now belongs to the canonical Runtime owner. |
| 367 | for (leaf, expected) in [ |
| 368 | (5, vec![0, 1, 2, 3, 4]), |
| 369 | (6, vec![0, 1, 2, 5]), |
| 370 | (7, vec![0, 1, 2, 3, 4, 6]), |
| 371 | ] { |
| 372 | let (_path, store, _, _) = canonical_alias_fixture( |
| 373 | &format!("archived-fork-{leaf}"), |
| 374 | (0..7) |
| 375 | .map(|index| { |
| 376 | let parent = match index { |
| 377 | 0 => None, |
| 378 | 5 => Some(3), |
| 379 | 6 => Some(5), |
| 380 | _ => Some(index), |
| 381 | }; |
| 382 | legacy_entry( |
| 383 | index + 1, |
| 384 | parent, |
| 385 | &format!("foo{index}"), |
| 386 | &format!("bar{index}"), |
| 387 | ) |
| 388 | }) |
| 389 | .collect(), |
| 390 | Some(leaf), |
| 391 | None, |
| 392 | ); |
| 393 | let messages = store.list_messages("legacy-import", None).unwrap(); |
| 394 | assert_eq!(messages.len(), expected.len()); |
| 395 | for (message, index) in messages.iter().zip(expected) { |
| 396 | assert_eq!(message.role, format!("foo{index}")); |
| 397 | assert_eq!(message.content, format!("bar{index}")); |
| 398 | assert_eq!(message.thread_id, "legacy-import"); |
| 399 | } |
| 400 | assert_eq!(store.list_leaf_messages("legacy-import").unwrap().len(), 2); |
| 401 | let archive = store |
| 402 | .snapshot_legacy_thread_history("legacy-import") |
| 403 | .unwrap(); |
| 404 | assert_eq!( |
| 405 | archive.messages.len(), |
| 406 | 7, |
| 407 | "all inactive branches are retained" |
| 408 | ); |
| 409 | assert_eq!(archive.current_leaf_id, Some(leaf)); |
| 410 | } |
| 411 | } |
| 412 | |
| 413 | fn legacy_entry( |
| 414 | id: i64, |
| 415 | parent_entry_id: Option<i64>, |
| 416 | role: &str, |
| 417 | content: &str, |
| 418 | ) -> codewhale_state::MessageRecord { |
| 419 | codewhale_state::MessageRecord { |
| 420 | id, |
| 421 | thread_id: "legacy-import".into(), |
| 422 | role: role.into(), |
| 423 | content: content.into(), |
| 424 | item: None, |
| 425 | created_at: 1_700_000_000 + id, |
| 426 | parent_entry_id, |
| 427 | } |
| 428 | } |
| 429 | |
| 430 | fn canonical_alias_fixture( |
| 431 | label: &str, |
| 432 | messages: Vec<codewhale_state::MessageRecord>, |
| 433 | current_leaf_id: Option<i64>, |
| 434 | goal: Option<codewhale_state::ThreadGoalRecord>, |
| 435 | ) -> ( |
| 436 | PathBuf, |
| 437 | StateStore, |
| 438 | codewhale_protocol::RuntimeOwnerReceipt, |
| 439 | codewhale_protocol::CanonicalThreadReceipt, |
| 440 | ) { |
| 441 | let path = temp_state_path(label); |
| 442 | let store = StateStore::open(Some(path.clone())).expect("store"); |
| 443 | let now = chrono::Utc::now().timestamp(); |
| 444 | let thread = ThreadMetadata { |
| 445 | id: "legacy-import".into(), |
| 446 | rollout_path: None, |
| 447 | preview: "import".into(), |
| 448 | ephemeral: false, |
| 449 | model_provider: "deepseek".into(), |
| 450 | created_at: now, |
| 451 | updated_at: now, |
| 452 | status: ThreadStatus::Idle, |
| 453 | path: None, |
| 454 | cwd: std::env::temp_dir(), |
| 455 | cli_version: "test".into(), |
| 456 | source: SessionSource::Resume, |
| 457 | name: None, |
| 458 | sandbox_policy: None, |
| 459 | approval_mode: None, |
| 460 | archived: false, |
| 461 | archived_at: None, |
| 462 | git_sha: None, |
| 463 | git_branch: None, |
| 464 | git_origin_url: None, |
| 465 | memory_mode: None, |
| 466 | current_leaf_id, |
| 467 | }; |
| 468 | store |
| 469 | .restore_legacy_thread_archive(&codewhale_state::LegacyThreadArchive { |
| 470 | thread, |
| 471 | messages, |
| 472 | goal, |
| 473 | checkpoints: Vec::new(), |
| 474 | }) |
| 475 | .expect("immutable legacy archive"); |
| 476 | // This is a database-only test of captured receipt comparison, not proof of |
| 477 | // peer authentication; the actual owner transport is tested in Runtime. |
| 478 | let owner = codewhale_protocol::RuntimeOwnerReceipt { |
| 479 | version: 1, |
| 480 | data_dir: std::env::temp_dir().join("canonical-alias-fixture"), |
| 481 | execution_scope: "held-scope".into(), |
| 482 | lease_generation: "held-generation".into(), |
| 483 | pid: std::process::id(), |
| 484 | process_start: "fixture-start".into(), |
| 485 | principal: "fixture-user".into(), |
| 486 | socket_path: std::env::temp_dir().join("fixture.sock"), |
| 487 | config_path: None, |
| 488 | }; |
| 489 | let receipt = codewhale_protocol::CanonicalThreadReceipt { |
| 490 | version: 1, |
| 491 | data_dir: owner.data_dir.clone(), |
| 492 | execution_scope: owner.execution_scope.clone(), |
| 493 | operation_key: format!("import-{label}"), |
| 494 | request_digest: "a".repeat(64), |
| 495 | history_digest: "b".repeat(64), |
| 496 | runtime_thread_id: "thr_reserved".into(), |
| 497 | session_id: "session_reserved".into(), |
| 498 | }; |
| 499 | (path, store, owner, receipt) |
| 500 | } |
| 501 | |
| 502 | #[test] |
| 503 | fn canonical_alias_publication_rejects_changed_full_graph_and_leaf() { |
| 504 | let (path, store, owner, receipt) = canonical_alias_fixture( |
| 505 | "source_cas", |
| 506 | vec![ |
| 507 | legacy_entry(1, None, "system", "system context"), |
| 508 | legacy_entry(2, Some(1), "user", "current branch"), |
| 509 | ], |
| 510 | Some(2), |
| 511 | None, |
| 512 | ); |
| 513 | let captured = store |
| 514 | .snapshot_legacy_thread_history("legacy-import") |
| 515 | .unwrap(); |
| 516 | // A historical external process can still tamper with SQLite. The canonical |
| 517 | // alias CAS must refuse it; no product append/leaf writer is exposed. |
| 518 | let conn = Connection::open(&path).unwrap(); |
| 519 | conn.execute_batch("INSERT INTO messages(id,thread_id,role,content,created_at,parent_entry_id) VALUES(3,'legacy-import','user','concurrent alternate branch',1700000003,1); UPDATE threads SET current_leaf_id=3 WHERE id='legacy-import';").unwrap(); |
| 520 | assert!( |
| 521 | store |
| 522 | .publish_canonical_runtime_link("legacy-import", None, &captured, &owner, &receipt) |
| 523 | .is_err() |
| 524 | ); |
| 525 | assert!( |
| 526 | store |
| 527 | .get_canonical_runtime_link("legacy-import", &owner) |
| 528 | .unwrap() |
| 529 | .is_none() |
| 530 | ); |
| 531 | let current = store |
| 532 | .snapshot_legacy_thread_history("legacy-import") |
| 533 | .unwrap(); |
| 534 | assert_eq!( |
| 535 | current.messages.len(), |
| 536 | 3, |
| 537 | "both source branches survive refusal" |
| 538 | ); |
| 539 | store |
| 540 | .publish_canonical_runtime_link("legacy-import", None, ¤t, &owner, &receipt) |
| 541 | .unwrap(); |
| 542 | store |
| 543 | .publish_canonical_runtime_link("legacy-import", None, ¤t, &owner, &receipt) |
| 544 | .unwrap(); |
| 545 | assert_eq!( |
| 546 | store |
| 547 | .get_canonical_runtime_link("legacy-import", &owner) |
| 548 | .unwrap(), |
| 549 | Some(receipt) |
| 550 | ); |
| 551 | assert!( |
| 552 | store |
| 553 | .restore_legacy_thread_archive(&codewhale_state::LegacyThreadArchive { |
| 554 | thread: store.get_thread("legacy-import").unwrap().unwrap(), |
| 555 | messages: current.messages, |
| 556 | goal: None, |
| 557 | checkpoints: Vec::new(), |
| 558 | }) |
| 559 | .is_err(), |
| 560 | "a bound alias cannot be replaced with an archive" |
| 561 | ); |
| 562 | } |
| 563 | |
| 564 | #[test] |
| 565 | fn canonical_alias_transaction_rolls_back_and_restart_keeps_exact_owner_binding() { |
| 566 | let (path, store, owner, receipt) = canonical_alias_fixture( |
| 567 | "atomic_restart", |
| 568 | vec![legacy_entry(1, None, "user", "retained original")], |
| 569 | Some(1), |
| 570 | None, |
| 571 | ); |
| 572 | let captured = store |
| 573 | .snapshot_legacy_thread_history("legacy-import") |
| 574 | .unwrap(); |
| 575 | let conn = Connection::open(&path).unwrap(); |
| 576 | conn.execute_batch("CREATE TRIGGER fail_canonical_receipt BEFORE INSERT ON thread_runtime_receipts BEGIN SELECT RAISE(ABORT, 'injected publication failure'); END;").unwrap(); |
| 577 | assert!( |
| 578 | store |
| 579 | .publish_canonical_runtime_link("legacy-import", None, &captured, &owner, &receipt) |
| 580 | .is_err() |
| 581 | ); |
| 582 | for table in [ |
| 583 | "thread_runtime_links", |
| 584 | "thread_runtime_operations", |
| 585 | "thread_runtime_receipts", |
| 586 | ] { |
| 587 | let count: i64 = conn |
| 588 | .query_row(&format!("SELECT count(*) FROM {table}"), [], |row| { |
| 589 | row.get(0) |
| 590 | }) |
| 591 | .unwrap(); |
| 592 | assert_eq!(count, 0, "{table} rolls back with the failed transaction"); |
| 593 | } |
| 594 | conn.execute_batch("DROP TRIGGER fail_canonical_receipt") |
| 595 | .unwrap(); |
| 596 | store |
| 597 | .publish_canonical_runtime_link("legacy-import", None, &captured, &owner, &receipt) |
| 598 | .unwrap(); |
| 599 | drop(store); |
| 600 | let reopened = StateStore::open(Some(path)).unwrap(); |
| 601 | assert_eq!( |
| 602 | reopened |
| 603 | .get_canonical_runtime_link("legacy-import", &owner) |
| 604 | .unwrap(), |
| 605 | Some(receipt) |
| 606 | ); |
| 607 | let mut foreign = owner.clone(); |
| 608 | foreign.execution_scope = "another-store-scope".into(); |
| 609 | assert!( |
| 610 | reopened |
| 611 | .get_canonical_runtime_link("legacy-import", &foreign) |
| 612 | .is_err() |
| 613 | ); |
| 614 | } |
| 615 | |
| 616 | #[test] |
| 617 | fn full_legacy_snapshot_keeps_branches_beyond_active_listing_limit() { |
| 618 | let mut messages = vec![legacy_entry(1, None, "user", "root")]; |
| 619 | messages |
| 620 | .extend((0..550).map(|index| { |
| 621 | legacy_entry(index + 2, Some(index + 1), "assistant", &index.to_string()) |
| 622 | })); |
| 623 | let alternate = 552; |
| 624 | messages.push(legacy_entry(alternate, Some(1), "user", "alternate")); |
| 625 | let (_path, store, _, _) = |
| 626 | canonical_alias_fixture("full_graph", messages, Some(alternate), None); |
| 627 | let snapshot = store |
| 628 | .snapshot_legacy_thread_history("legacy-import") |
| 629 | .unwrap(); |
| 630 | assert_eq!(snapshot.messages.len(), 552); |
| 631 | assert_eq!(snapshot.current_leaf_id, Some(alternate)); |
| 632 | assert_eq!(store.list_messages("legacy-import", None).unwrap().len(), 2); |
| 633 | } |
| 634 | |
| 635 | #[test] |
| 636 | fn canonical_goal_snapshot_and_publication_cas_keep_original_source_status_and_timestamps() { |
| 637 | let goal = codewhale_state::ThreadGoalRecord { |
| 638 | thread_id: "legacy-import".into(), |
| 639 | goal_id: "goal-source-revision".into(), |
| 640 | objective: "original goal".into(), |
| 641 | status: codewhale_state::ThreadGoalStatus::Active, |
| 642 | token_budget: Some(1000), |
| 643 | tokens_used: 41, |
| 644 | time_used_seconds: 9, |
| 645 | continuation_count: 2, |
| 646 | created_at: 1_700_000_000, |
| 647 | updated_at: 1_700_000_010, |
| 648 | last_gap_fingerprint: Some("a".repeat(64)), |
| 649 | repeated_gap_count: 3, |
| 650 | last_gap_pass: Some(2), |
| 651 | pause_reason: None, |
| 652 | }; |
| 653 | let (path, store, owner, receipt) = |
| 654 | canonical_alias_fixture("goal_source_cas", Vec::new(), None, Some(goal.clone())); |
| 655 | let captured = store |
| 656 | .snapshot_legacy_thread_history("legacy-import") |
| 657 | .unwrap(); |
| 658 | let saved = captured.goal.as_ref().unwrap(); |
| 659 | assert_eq!( |
| 660 | saved.status, |
| 661 | codewhale_protocol::ThreadGoalStatus::Active, |
| 662 | "source CAS keeps original persisted status rather than target normalization" |
| 663 | ); |
| 664 | assert_eq!(saved.updated_at, goal.updated_at); |
| 665 | assert_eq!( |
| 666 | store |
| 667 | .get_thread_goal("legacy-import") |
| 668 | .unwrap() |
| 669 | .unwrap() |
| 670 | .status, |
| 671 | codewhale_state::ThreadGoalStatus::Paused, |
| 672 | "ordinary restored-reader behavior is unchanged" |
| 673 | ); |
| 674 | let mut changed = goal; |
| 675 | changed.objective = "goal changed during owner await".into(); |
| 676 | changed.updated_at += 1; |
| 677 | let conn = Connection::open(&path).unwrap(); |
| 678 | conn.execute( |
| 679 | "UPDATE thread_goals SET objective=?1, updated_at=?2 WHERE thread_id='legacy-import'", |
| 680 | rusqlite::params![changed.objective, changed.updated_at], |
| 681 | ) |
| 682 | .unwrap(); |
| 683 | assert!( |
| 684 | store |
| 685 | .publish_canonical_runtime_link("legacy-import", None, &captured, &owner, &receipt) |
| 686 | .is_err() |
| 687 | ); |
| 688 | assert!( |
| 689 | store |
| 690 | .get_canonical_runtime_link("legacy-import", &owner) |
| 691 | .unwrap() |
| 692 | .is_none() |
| 693 | ); |
| 694 | assert_eq!( |
| 695 | store |
| 696 | .get_thread_goal("legacy-import") |
| 697 | .unwrap() |
| 698 | .unwrap() |
| 699 | .objective, |
| 700 | changed.objective |
| 701 | ); |
| 702 | } |
| 703 |