返回 CodeWhale
parity_state.rs
根目录 / crates / state / tests / parity_state.rs
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, &current, &owner, &receipt)
541 .unwrap();
542 store
543 .publish_canonical_runtime_link("legacy-import", None, &current, &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
703 lines RUST