返回 CodeWhale
tests.rs
1 use super::*;
2 use axum::extract::Request;
3 use tokio::sync::Notify;
4
5 #[derive(Clone)]
6 struct FakeOwner {
7 record: Arc<Mutex<Value>>,
8 calls: Arc<Mutex<Vec<(String, String, Value)>>>,
9 owner: RuntimeOwnerReceipt,
10 block_import: bool,
11 import_seen: Arc<Notify>,
12 import_release: Arc<Notify>,
13 operations: Arc<Mutex<HashMap<String, (Value, CanonicalThreadReceipt)>>>,
14 fork_record: Arc<Mutex<Option<Value>>>,
15 missing_after_commit: Arc<std::sync::atomic::AtomicBool>,
16 goal: Arc<Mutex<Option<Value>>>,
17 turn_seq: Arc<std::sync::atomic::AtomicU64>,
18 pending_operation: Arc<std::sync::atomic::AtomicBool>,
19 recovery_prepared: Arc<std::sync::atomic::AtomicBool>,
20 recovery_refusal: Arc<std::sync::atomic::AtomicBool>,
21 document_digest: Arc<Mutex<String>>,
22 }
23
24 async fn owned_http(State(fake): State<FakeOwner>, request: Request) -> Response {
25 let method = request.method().clone();
26 let path = request.uri().path().to_owned();
27 let authorization = request
28 .headers()
29 .get(axum::http::header::AUTHORIZATION)
30 .cloned();
31 let bytes = axum::body::to_bytes(request.into_body(), MAX_CANONICAL_HISTORY_BYTES)
32 .await
33 .unwrap();
34 let body = if bytes.is_empty() {
35 Value::Null
36 } else {
37 serde_json::from_slice(&bytes).unwrap()
38 };
39 fake.calls
40 .lock()
41 .await
42 .push((method.to_string(), path.clone(), body.clone()));
43 if path == "/v1/config/reload" {
44 assert_eq!(method, Method::POST);
45 assert_eq!(
46 authorization.as_ref().map(|value| value.as_bytes()),
47 Some(b"Bearer private-compatibility-reload-fixture".as_slice())
48 );
49 return StatusCode::NO_CONTENT.into_response();
50 }
51 if path == "/v1/thread-history/operations/lookup"
52 || path == "/v1/thread-history/operations/recover"
53 {
54 let recovery = if path.ends_with("/recover") {
55 Some(serde_json::from_value::<CanonicalThreadOperationRecovery>(body.clone()).unwrap())
56 } else {
57 None
58 };
59 let input: CanonicalThreadOperationLookup = match recovery.as_ref() {
60 Some(recovery) => recovery.operation.clone(),
61 None => serde_json::from_value(body).unwrap(),
62 };
63 assert_eq!(input.expected_data_dir, fake.owner.data_dir);
64 assert_eq!(input.expected_execution_scope, fake.owner.execution_scope);
65 let operations = fake.operations.lock().await;
66 let Some((prior, receipt)) = operations.get(&input.operation_key) else {
67 return Json(CanonicalThreadOperationStatus::Absent).into_response();
68 };
69 let prior: CanonicalThreadMutationRequest = serde_json::from_value(prior.clone()).unwrap();
70 assert_eq!(input.workspace, prior.workspace);
71 let (kind, source) = match prior.mutation {
72 CanonicalThreadMutation::Create { .. } => (CanonicalThreadOperationKind::Create, None),
73 CanonicalThreadMutation::Resume { source, .. } => {
74 (CanonicalThreadOperationKind::Resume, Some(source))
75 }
76 CanonicalThreadMutation::Fork { source, .. } => {
77 (CanonicalThreadOperationKind::Fork, Some(source))
78 }
79 };
80 let source_runtime_thread_id = source.as_ref().and_then(|source| match source {
81 CanonicalHistorySource::Thread {
82 runtime_thread_id, ..
83 } => Some(runtime_thread_id.clone()),
84 CanonicalHistorySource::SavedSession { .. } => None,
85 });
86 let association = codewhale_protocol::CanonicalThreadOperationAssociation {
87 kind,
88 source_runtime_thread_id,
89 source_session_id: None,
90 };
91 if let Some(recovery) = recovery {
92 if recovery.association != association {
93 return (
94 StatusCode::CONFLICT,
95 Json(json!({"error":"retained association changed"})),
96 )
97 .into_response();
98 }
99 if fake
100 .recovery_refusal
101 .load(std::sync::atomic::Ordering::SeqCst)
102 {
103 return (
104 StatusCode::CONFLICT,
105 Json(json!({"error":"prepared target graph changed"})),
106 )
107 .into_response();
108 }
109 if fake
110 .recovery_prepared
111 .load(std::sync::atomic::Ordering::SeqCst)
112 {
113 fake.pending_operation
114 .store(false, std::sync::atomic::Ordering::SeqCst);
115 }
116 }
117 let status = if fake
118 .pending_operation
119 .load(std::sync::atomic::Ordering::SeqCst)
120 {
121 CanonicalThreadOperationStatus::Pending {
122 receipt: receipt.clone(),
123 association,
124 }
125 } else {
126 CanonicalThreadOperationStatus::Committed {
127 receipt: receipt.clone(),
128 association,
129 }
130 };
131 return Json(status).into_response();
132 }
133 if path == "/v1/thread-history/mutate" {
134 let input: CanonicalThreadMutationRequest = serde_json::from_value(body.clone()).unwrap();
135 let mut operations = fake.operations.lock().await;
136 if let Some((prior, receipt)) = operations.get(&input.operation_key) {
137 if prior != &body {
138 return (
139 StatusCode::CONFLICT,
140 Json(json!({"error":"intent changed"})),
141 )
142 .into_response();
143 }
144 return Json(receipt.clone()).into_response();
145 }
146 let mut record = fake.record.lock().await.clone();
147 let id = if matches!(input.mutation, CanonicalThreadMutation::Fork { .. }) {
148 record["id"] = json!("canonical-fork");
149 record["workspace"] = json!(input.workspace);
150 *fake.fork_record.lock().await = Some(record);
151 "canonical-fork"
152 } else {
153 "canonical-1"
154 };
155 let receipt = CanonicalThreadReceipt {
156 version: 1,
157 data_dir: fake.owner.data_dir,
158 execution_scope: fake.owner.execution_scope,
159 operation_key: input.operation_key.clone(),
160 request_digest: "1".repeat(64),
161 history_digest: "2".repeat(64),
162 runtime_thread_id: id.into(),
163 session_id: if id == "canonical-fork" {
164 "session-fork"
165 } else {
166 "session-1"
167 }
168 .into(),
169 };
170 operations.insert(input.operation_key, (body, receipt.clone()));
171 return Json(receipt).into_response();
172 }
173 if path == "/v1/thread-history/import" {
174 let input: CanonicalHistoryImportRequest = serde_json::from_value(body).unwrap();
175 if input
176 .target_runtime_thread_id
177 .as_deref()
178 .is_some_and(|id| id != "canonical-1")
179 {
180 return (
181 StatusCode::NOT_FOUND,
182 Json(json!({"error":"existing canonical target missing; no replacement"})),
183 )
184 .into_response();
185 }
186 fake.import_seen.notify_one();
187 if fake.block_import {
188 fake.import_release.notified().await;
189 }
190 if let Some(mut goal) = input.history.goal {
191 goal.thread_id = input
192 .target_runtime_thread_id
193 .clone()
194 .unwrap_or_else(|| "canonical-1".into());
195 if goal.status == codewhale_protocol::ThreadGoalStatus::Active {
196 goal.status = codewhale_protocol::ThreadGoalStatus::Paused;
197 goal.pause_reason = None;
198 }
199 let goal = serde_json::to_value(goal).unwrap();
200 let mut current = fake.goal.lock().await;
201 if current.as_ref().is_some_and(|prior| prior != &goal) {
202 return (
203 StatusCode::CONFLICT,
204 Json(json!({"error":"canonical goal conflict; source retained"})),
205 )
206 .into_response();
207 }
208 *current = Some(goal);
209 }
210 return Json(CanonicalThreadReceipt {
211 version: 1,
212 data_dir: fake.owner.data_dir,
213 execution_scope: fake.owner.execution_scope,
214 operation_key: input.operation_key,
215 request_digest: "1".repeat(64),
216 history_digest: "2".repeat(64),
217 runtime_thread_id: input
218 .target_runtime_thread_id
219 .unwrap_or_else(|| "canonical-1".into()),
220 session_id: "session-1".into(),
221 })
222 .into_response();
223 }
224 if path == "/v1/threads/running" {
225 return Json(json!([])).into_response();
226 }
227 if path == "/v1/threads/canonical-1/turns" && method == Method::POST {
228 let seq = fake
229 .turn_seq
230 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
231 + 1;
232 return Json(json!({"turn":{"id":format!("fixture-turn-{seq}")}})).into_response();
233 }
234 if path == "/v1/threads/canonical-1/events" {
235 let seq = fake.turn_seq.load(std::sync::atomic::Ordering::SeqCst);
236 let frame = json!({"seq":seq,"turn_id":format!("fixture-turn-{seq}"),
237 "payload":{"turn":{"status":"completed"}}});
238 return (
239 [(header::CONTENT_TYPE, "text/event-stream")],
240 format!("event: turn.completed\ndata: {frame}\n\n"),
241 )
242 .into_response();
243 }
244 if path == "/v1/threads/canonical-1/goal" {
245 if method == Method::PUT {
246 let goal = json!({"thread_id":"canonical-1","goal_id":"fixture-goal",
247 "objective":body["objective"],"token_budget":body["token_budget"],"status":"active",
248 "tokens_used":0,"time_used_seconds":0,"continuation_count":0,
249 "created_at":1,"updated_at":1,"repeated_gap_count":0});
250 *fake.goal.lock().await = Some(goal.clone());
251 return Json(goal).into_response();
252 }
253 let mut goal = fake.goal.lock().await;
254 if method == Method::DELETE && goal.take().is_some() {
255 return StatusCode::NO_CONTENT.into_response();
256 }
257 return match goal.clone() {
258 Some(goal) => Json(goal).into_response(),
259 None => (StatusCode::NOT_FOUND, Json(json!({"error":"no goal"}))).into_response(),
260 };
261 }
262 if path == "/v1/threads" {
263 return Json(json!([fake.record.lock().await.clone()])).into_response();
264 }
265 if path == "/v1/threads/canonical-1/history" {
266 return Json(json!({"version":1,"data_dir":fake.owner.data_dir,"execution_scope":fake.owner.execution_scope,
267 "runtime_thread_id":"canonical-1","saved_session_id":"session-1","saved_document_digest":"a".repeat(64),"session_goal_digest":"74234e98afe7498fb5daf1f36ac2d78acc339464f950703b8c019892f982b90b","document_digest":fake.document_digest.lock().await.clone(),
268 "session":{"metadata":{"id":"session-1"},"messages":[{"role":"user","content":"active"}],
269 "journal":{"entries":[{"id":"root"},{"id":"inactive-branch"},{"id":"active"}]},"leaf_id":"active"}})).into_response();
270 }
271 if path == "/v1/threads/canonical-fork" {
272 return Json(fake.fork_record.lock().await.clone().unwrap()).into_response();
273 }
274 if path != "/v1/threads/canonical-1" {
275 return (StatusCode::NOT_FOUND, Json(json!({"error":"missing"}))).into_response();
276 }
277 if fake
278 .missing_after_commit
279 .load(std::sync::atomic::Ordering::SeqCst)
280 && !fake.operations.lock().await.is_empty()
281 {
282 return (
283 StatusCode::NOT_FOUND,
284 Json(json!({"error":"metadata disappeared"})),
285 )
286 .into_response();
287 }
288 if method == Method::PATCH {
289 let mut record = fake.record.lock().await;
290 for (key, value) in body.as_object().unwrap() {
291 record[key] = if key == "title" && value.as_str().is_some_and(|s| s.trim().is_empty()) {
292 Value::Null
293 } else {
294 value.clone()
295 };
296 }
297 return Json(record.clone()).into_response();
298 }
299 Json(json!({"thread":fake.record.lock().await.clone(),"items":[],"turns":[]})).into_response()
300 }
301
302 fn fixture_state(block_import: bool) -> (AppState, tempfile::TempDir, FakeOwner) {
303 let temp = tempfile::tempdir().unwrap();
304 let config_path = temp.path().join("config.toml");
305 std::fs::write(&config_path, "").unwrap();
306 let mut state = build_state(Some(config_path.clone()), None).unwrap();
307 let owner = RuntimeOwnerReceipt {
308 version: 1,
309 data_dir: temp.path().to_owned(),
310 execution_scope: "fixture-owner".into(),
311 lease_generation: "fixture-generation".into(),
312 pid: std::process::id(),
313 process_start: "fixture".into(),
314 principal: "fixture".into(),
315 socket_path: temp.path().join("owner.sock"),
316 config_path: Some(config_path),
317 };
318 state.captured_owner = Some(owner.clone());
319 state.frontend_workspace = Some(temp.path().to_owned());
320 let fake = FakeOwner {
321 record: Arc::new(Mutex::new(
322 json!({"id":"canonical-1","created_at":"2026-10-02T00:00:00Z",
323 "updated_at":"2026-10-02T00:00:01Z","model":"fixture-model","model_provider":"custom","model_provider_id":"fixture-account",
324 "workspace":temp.path(),"archived":false,"title":"named"}),
325 )),
326 calls: Arc::default(),
327 owner,
328 block_import,
329 import_seen: Arc::new(Notify::new()),
330 import_release: Arc::new(Notify::new()),
331 operations: Arc::default(),
332 fork_record: Arc::default(),
333 missing_after_commit: Arc::default(),
334 goal: Arc::default(),
335 turn_seq: Arc::default(),
336 pending_operation: Arc::default(),
337 recovery_prepared: Arc::default(),
338 recovery_refusal: Arc::default(),
339 document_digest: Arc::new(Mutex::new("b".repeat(64))),
340 };
341 (state, temp, fake)
342 }
343
344 async fn fixture(
345 block_import: bool,
346 ) -> (
347 AppState,
348 tempfile::TempDir,
349 FakeOwner,
350 tokio::task::JoinHandle<()>,
351 ) {
352 let (state, temp, fake) = fixture_state(block_import);
353 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
354 let address = listener.local_addr().unwrap();
355 let app = Router::new().fallback(owned_http).with_state(fake.clone());
356 let server = tokio::spawn(async move {
357 axum::serve(listener, app).await.unwrap();
358 });
359 let mut bridge = RuntimeBridge::from_base_url_for_test(format!("http://{address}"));
360 bridge.auth_token = Some("private-compatibility-reload-fixture".into());
361 *state.runtime_bridge.lock().await = Some(Arc::new(Mutex::new(bridge)));
362 (state, temp, fake, server)
363 }
364
365 pub(super) fn compatibility_router() -> (AppState, tempfile::TempDir, Router) {
366 let (state, temp, fake) = fixture_state(false);
367 (
368 state,
369 temp,
370 Router::new().fallback(owned_http).with_state(fake),
371 )
372 }
373
374 pub(super) async fn compatibility_fixture()
375 -> (AppState, tempfile::TempDir, tokio::task::JoinHandle<()>) {
376 let (state, temp, _fake, server) = fixture(false).await;
377 (state, temp, server)
378 }
379
380 async fn seed(state: &AppState, id: &str, workspace: &Path) -> StateStore {
381 seed_archive(state, id, workspace, Vec::new(), None).await
382 }
383
384 async fn seed_archive(
385 state: &AppState,
386 id: &str,
387 workspace: &Path,
388 messages: Vec<codewhale_state::MessageRecord>,
389 goal: Option<codewhale_state::ThreadGoalRecord>,
390 ) -> StateStore {
391 let store = state.runtime.read().await.state_store().clone();
392 let thread = codewhale_state::ThreadMetadata {
393 id: id.into(),
394 rollout_path: None,
395 preview: "legacy preview".into(),
396 ephemeral: false,
397 model_provider: "legacy-provider".into(),
398 created_at: 1,
399 updated_at: 1,
400 status: codewhale_state::ThreadStatus::Idle,
401 path: None,
402 cwd: workspace.to_owned(),
403 cli_version: "old".into(),
404 source: codewhale_state::SessionSource::Api,
405 name: Some("old name".into()),
406 sandbox_policy: None,
407 approval_mode: None,
408 archived: false,
409 archived_at: None,
410 git_sha: None,
411 git_branch: None,
412 git_origin_url: None,
413 memory_mode: None,
414 current_leaf_id: messages.last().map(|message| message.id),
415 };
416 store
417 .restore_legacy_thread_archive(&codewhale_state::LegacyThreadArchive {
418 thread,
419 messages,
420 goal,
421 checkpoints: Vec::new(),
422 })
423 .unwrap();
424 store
425 }
426
427 fn archive_message(
428 id: i64,
429 text: &str,
430 parent_entry_id: Option<i64>,
431 ) -> codewhale_state::MessageRecord {
432 codewhale_state::MessageRecord {
433 id,
434 thread_id: "legacy".into(),
435 role: "user".into(),
436 content: text.into(),
437 item: None,
438 created_at: 1,
439 parent_entry_id,
440 }
441 }
442
443 // Model an external historical writer only in test-owned SQLite. Production
444 // controls cannot append/replace history or mint bare links through State APIs.
445 fn seed_bare_link(store: &StateStore, thread: &str, runtime: &str) {
446 rusqlite::Connection::open(store.db_path()).unwrap().execute(
447 "INSERT INTO thread_runtime_links(thread_id,runtime_thread_id,created_at) VALUES(?1,?2,1)",
448 rusqlite::params![thread,runtime],
449 ).unwrap();
450 }
451
452 #[tokio::test]
453 async fn legacy_resolution_imports_all_branches_and_publishes_the_same_operation_once() {
454 let (state, temp, fake, server) = fixture(false).await;
455 let store = seed_archive(
456 &state,
457 "legacy",
458 temp.path(),
459 vec![
460 archive_message(1, "root", None),
461 archive_message(2, "old branch", Some(1)),
462 archive_message(3, "new branch", Some(1)),
463 ],
464 None,
465 )
466 .await;
467 let before = store.snapshot_legacy_thread_history("legacy").unwrap();
468 let first = resolve(&state, "legacy", true).await.unwrap();
469 let second = resolve(&state, "legacy", true).await.unwrap();
470 assert_eq!(first, second);
471 let calls = fake.calls.lock().await;
472 let imports: Vec<_> = calls
473 .iter()
474 .filter(|(_, path, _)| path == "/v1/thread-history/import")
475 .collect();
476 assert_eq!(imports.len(), 1);
477 assert_eq!(
478 imports[0].2["history"],
479 serde_json::to_value(&before).unwrap()
480 );
481 assert_eq!(imports[0].2["target_runtime_thread_id"], Value::Null);
482 let receipt = store
483 .get_canonical_runtime_link("legacy", state.captured_owner.as_ref().unwrap())
484 .unwrap()
485 .unwrap();
486 assert_eq!(receipt.operation_key, migration_key(&before).unwrap());
487 assert_eq!(receipt.runtime_thread_id, "canonical-1");
488 server.abort();
489 }
490
491 #[tokio::test]
492 async fn existing_bare_link_adopts_exact_target_without_empty_creation() {
493 let (state, temp, fake, server) = fixture(false).await;
494 let store = seed_archive(
495 &state,
496 "legacy",
497 temp.path(),
498 vec![archive_message(1, "retained source", None)],
499 None,
500 )
501 .await;
502 seed_bare_link(&store, "legacy", "canonical-1");
503 resolve(&state, "legacy", true).await.unwrap();
504 let calls = fake.calls.lock().await;
505 let import = calls
506 .iter()
507 .find(|(_, p, _)| p == "/v1/thread-history/import")
508 .unwrap();
509 assert_eq!(import.2["target_runtime_thread_id"], "canonical-1");
510 assert_eq!(
511 import.2["history"]["messages"][0]["content"],
512 "retained source"
513 );
514 assert!(
515 !calls
516 .iter()
517 .any(|(m, p, _)| m == "POST" && p == "/v1/threads")
518 );
519 server.abort();
520 }
521
522 #[tokio::test]
523 async fn cancelled_waiter_cannot_drop_committed_alias_publication() {
524 let (state, temp, fake, server) = fixture(true).await;
525 let store = seed(&state, "legacy", temp.path()).await;
526 let worker_state = state.clone();
527 let waiter = tokio::spawn(async move { resolve(&worker_state, "legacy", true).await });
528 fake.import_seen.notified().await;
529 waiter.abort();
530 fake.import_release.notify_one();
531 tokio::time::timeout(Duration::from_secs(5), async {
532 loop {
533 if store
534 .get_canonical_runtime_link("legacy", state.captured_owner.as_ref().unwrap())
535 .unwrap()
536 .is_some()
537 {
538 break;
539 }
540 tokio::task::yield_now().await;
541 }
542 })
543 .await
544 .unwrap();
545 assert_eq!(
546 store.get_runtime_thread_link("legacy").unwrap().as_deref(),
547 Some("canonical-1")
548 );
549 server.abort();
550 }
551
552 #[tokio::test]
553 async fn source_graph_change_during_import_refuses_publication_and_preserves_both_histories() {
554 let (state, temp, fake, server) = fixture(true).await;
555 let store = seed_archive(
556 &state,
557 "legacy",
558 temp.path(),
559 vec![archive_message(1, "first", None)],
560 None,
561 )
562 .await;
563 let worker_state = state.clone();
564 let waiter = tokio::spawn(async move { resolve(&worker_state, "legacy", true).await });
565 fake.import_seen.notified().await;
566 let mut foreign = rusqlite::Connection::open(store.db_path()).unwrap();
567 let transaction = foreign
568 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
569 .unwrap();
570 transaction.execute("INSERT INTO messages(id,thread_id,role,content,item_json,created_at,parent_entry_id) VALUES(2,'legacy','user','changed while waiting',NULL,1,1)", []).unwrap();
571 transaction
572 .execute("UPDATE threads SET current_leaf_id=2 WHERE id='legacy'", [])
573 .unwrap();
574 transaction.commit().unwrap();
575 fake.import_release.notify_one();
576 let error = waiter.await.unwrap().unwrap_err();
577 assert!(format!("{error:#}").contains("legacy history changed"));
578 assert!(
579 store
580 .get_canonical_runtime_link("legacy", state.captured_owner.as_ref().unwrap())
581 .unwrap()
582 .is_none()
583 );
584 assert_eq!(
585 store
586 .snapshot_legacy_thread_history("legacy")
587 .unwrap()
588 .messages
589 .len(),
590 2
591 );
592 assert_eq!(fake.record.lock().await["id"], "canonical-1");
593 server.abort();
594 }
595
596 #[tokio::test]
597 async fn list_is_global_read_only_union_and_bound_alias_dedupes_actual_metadata() {
598 let (state, temp, fake, server) = fixture(false).await;
599 let store = seed(&state, "legacy", temp.path()).await;
600 resolve(&state, "legacy", true).await.unwrap();
601 let other = temp.path().join("other-workspace");
602 seed(&state, "other-legacy", &other).await;
603 fake.record.lock().await["title"] = json!("canonical title");
604 let before =
605 serde_json::to_value(store.snapshot_legacy_thread_history("legacy").unwrap()).unwrap();
606 let result = handle(
607 &state,
608 ThreadRequest::List(codewhale_protocol::ThreadListParams {
609 include_archived: false,
610 limit: None,
611 }),
612 )
613 .await
614 .unwrap();
615 assert_eq!(result.threads.len(), 2);
616 assert!(
617 result
618 .threads
619 .iter()
620 .any(|t| t.id == "other-legacy" && t.cwd == other)
621 );
622 let alias = result.threads.iter().find(|t| t.id == "legacy").unwrap();
623 assert_eq!(alias.name.as_deref(), Some("canonical title"));
624 assert!(!result.threads.iter().any(|t| t.id == "canonical-1"));
625 assert_eq!(
626 serde_json::to_value(store.snapshot_legacy_thread_history("legacy").unwrap()).unwrap(),
627 before
628 );
629 assert_eq!(
630 fake.calls
631 .lock()
632 .await
633 .iter()
634 .filter(|(_, p, _)| p == "/v1/thread-history/import")
635 .count(),
636 1
637 );
638 server.abort();
639 }
640
641 #[tokio::test]
642 async fn canonical_title_clear_and_archive_filter_ignore_stale_sqlite_names() {
643 let (state, temp, fake, server) = fixture(false).await;
644 let store = seed(&state, "legacy", temp.path()).await;
645 resolve(&state, "legacy", true).await.unwrap();
646 let clear = handle(
647 &state,
648 ThreadRequest::SetName(codewhale_protocol::ThreadSetNameParams {
649 thread_id: "legacy".into(),
650 name: String::new(),
651 }),
652 )
653 .await
654 .unwrap();
655 assert!(clear.thread.unwrap().name.is_none());
656 assert_eq!(fake.record.lock().await["title"], Value::Null);
657 assert_eq!(
658 store.get_thread("legacy").unwrap().unwrap().name.as_deref(),
659 Some("old name")
660 );
661 handle(
662 &state,
663 ThreadRequest::Archive {
664 thread_id: "legacy".into(),
665 },
666 )
667 .await
668 .unwrap();
669 let active = list(
670 &state,
671 codewhale_protocol::ThreadListParams {
672 include_archived: false,
673 limit: None,
674 },
675 )
676 .await
677 .unwrap();
678 assert!(active.threads.is_empty());
679 let all = list(
680 &state,
681 codewhale_protocol::ThreadListParams {
682 include_archived: true,
683 limit: None,
684 },
685 )
686 .await
687 .unwrap();
688 assert_eq!(all.threads.len(), 1);
689 assert_eq!(all.threads[0].status, ThreadStatus::Archived);
690 server.abort();
691 }
692
693 #[tokio::test]
694 async fn read_returns_full_journal_including_inactive_branch() {
695 let (state, _temp, _fake, server) = fixture(false).await;
696 let value = handle(
697 &state,
698 ThreadRequest::Read(codewhale_protocol::ThreadReadParams {
699 thread_id: "canonical-1".into(),
700 }),
701 )
702 .await
703 .unwrap();
704 assert_eq!(
705 value.data["history"]["session"]["journal"]["entries"]
706 .as_array()
707 .unwrap()
708 .len(),
709 3
710 );
711 assert_eq!(
712 value.data["history"]["session"]["journal"]["entries"][1]["id"],
713 "inactive-branch"
714 );
715 server.abort();
716 }
717
718 #[tokio::test]
719 async fn unknown_canonical_target_never_creates_or_changes_an_alias() {
720 let (state, _temp, fake, server) = fixture(false).await;
721 let error = resolve(&state, "missing", true).await.unwrap_err();
722 assert!(format!("{error:#}").contains("no replacement thread"));
723 assert!(state.runtime_thread_map.lock().await.is_empty());
724 assert!(!fake.calls.lock().await.iter().any(|(m, _, _)| m == "POST"));
725 server.abort();
726 }
727
728 #[tokio::test]
729 async fn unreceipted_goal_deltas_and_unqualified_controls_are_refused_before_io() {
730 let (state, _temp, fake, server) = fixture(false).await;
731 for request in [
732 ThreadRequest::GoalRecordProgress(codewhale_protocol::ThreadGoalProgressParams {
733 thread_id: "canonical-1".into(),
734 token_delta: 1,
735 time_delta_seconds: 1,
736 record_continuation: true,
737 }),
738 ThreadRequest::Create {
739 metadata: Value::Null,
740 },
741 ] {
742 assert!(handle(&state, request).await.is_err());
743 }
744 assert!(fake.calls.lock().await.is_empty());
745 server.abort();
746 }
747
748 #[tokio::test]
749 async fn keyed_start_carries_exact_captured_config_and_never_opens_a_legacy_writer() {
750 let (state, temp, fake, server) = fixture(false).await;
751 let request: ThreadRequest = serde_json::from_value(json!({"kind":"start",
752 "operation_key":"start-intent","model":"fixture-model",
753 "model_provider":"fixture-account","cwd":temp.path(),"persist_extended_history":true}))
754 .unwrap();
755 let result = handle(&state, request).await.unwrap();
756 assert_eq!(result.status, "started");
757 assert_eq!(result.thread_id, "canonical-1");
758 assert_eq!(result.data["receipt"]["operation_key"], "start-intent");
759 let operations = fake.operations.lock().await;
760 assert_eq!(operations.len(), 1);
761 let (body, _) = operations.get("start-intent").unwrap();
762 assert_eq!(
763 body["mutation"]["config"],
764 json!({"workspace":temp.path(),
765 "model":"fixture-model","model_provider":"fixture-account"})
766 );
767 assert_eq!(body["expected_data_dir"], json!(fake.owner.data_dir));
768 assert_eq!(body["expected_execution_scope"], fake.owner.execution_scope);
769 let store = state.runtime.read().await.state_store().clone();
770 assert!(
771 store
772 .list_threads(codewhale_state::ThreadListFilters {
773 include_archived: true,
774 limit: None
775 })
776 .unwrap()
777 .is_empty()
778 );
779 server.abort();
780 }
781
782 #[tokio::test]
783 async fn create_retry_keeps_exact_client_intent_and_returns_same_owner_receipt() {
784 let (state, temp, fake, server) = fixture(false).await;
785 let request = ThreadRequest::Create {
786 metadata: json!({"operation_key":"create-intent",
787 "workspace":temp.path(),"model":"fixture-model","allowed_tools":[],"allow_shell":false}),
788 };
789 let first = handle(&state, request.clone()).await.unwrap();
790 let second = handle(&state, request).await.unwrap();
791 assert_eq!(first.thread_id, second.thread_id);
792 assert_eq!(first.data["receipt"], second.data["receipt"]);
793 assert_eq!(fake.operations.lock().await.len(), 1);
794 let calls = fake.calls.lock().await;
795 let requests: Vec<_> = calls
796 .iter()
797 .filter(|(_, path, _)| path == "/v1/thread-history/mutate")
798 .map(|(_, _, body)| body)
799 .collect();
800 assert_eq!(requests.len(), 1);
801 assert_eq!(
802 calls
803 .iter()
804 .filter(|(_, path, _)| path.ends_with("/operations/lookup"))
805 .count(),
806 2
807 );
808 assert!(
809 requests[0]["mutation"]["config"]
810 .get("operation_key")
811 .is_none()
812 );
813 server.abort();
814 }
815
816 #[tokio::test]
817 async fn creation_without_acknowledged_scope_or_stable_key_refuses_before_io() {
818 let (mut state, temp, fake, server) = fixture(false).await;
819 for input in [
820 json!({"kind":"start"}),
821 json!({"kind":"create","metadata":{"operation_key":"wrong-scope","workspace":temp.path().join("other")}}),
822 ] {
823 assert!(
824 handle(&state, serde_json::from_value(input).unwrap())
825 .await
826 .is_err()
827 );
828 }
829 state.frontend_workspace = None;
830 let request =
831 serde_json::from_value(json!({"kind":"start","operation_key":"no-default"})).unwrap();
832 assert!(
833 handle(&state, request)
834 .await
835 .unwrap_err()
836 .message
837 .contains("acknowledged workspace")
838 );
839 assert!(fake.calls.lock().await.is_empty());
840 server.abort();
841 }
842
843 #[tokio::test]
844 async fn ordinary_resume_binds_full_document_digest_and_preserves_public_alias() {
845 let (state, temp, fake, server) = fixture(false).await;
846 let store = seed(&state, "legacy", temp.path()).await;
847 resolve(&state, "legacy", true).await.unwrap();
848 let before = store.snapshot_legacy_thread_history("legacy").unwrap();
849 let request = serde_json::from_value(json!({"kind":"resume","thread_id":"legacy",
850 "operation_key":"resume-intent","persist_extended_history":true}))
851 .unwrap();
852 let result = handle(&state, request).await.unwrap();
853 assert_eq!(result.thread_id, "legacy");
854 assert_eq!(result.status, "resumed");
855 let operations = fake.operations.lock().await;
856 assert_eq!(
857 operations["resume-intent"].0["mutation"]["source"],
858 json!({"kind":"thread",
859 "runtime_thread_id":"canonical-1","expected_document_digest":"b".repeat(64)})
860 );
861 assert_eq!(
862 serde_json::to_value(store.snapshot_legacy_thread_history("legacy").unwrap()).unwrap(),
863 serde_json::to_value(before).unwrap()
864 );
865 server.abort();
866 }
867
868 #[tokio::test]
869 async fn ordinary_fork_returns_canonical_identity_and_preserves_original_branches() {
870 let (state, temp, fake, server) = fixture(false).await;
871 let store = seed(&state, "legacy", temp.path()).await;
872 resolve(&state, "legacy", true).await.unwrap();
873 let before = store.snapshot_legacy_thread_history("legacy").unwrap();
874 let request = serde_json::from_value(json!({"kind":"fork","thread_id":"legacy",
875 "operation_key":"fork-intent"}))
876 .unwrap();
877 let result = handle(&state, request).await.unwrap();
878 assert_eq!(result.thread_id, "canonical-fork");
879 assert_eq!(result.status, "forked");
880 assert_eq!(result.data["receipt"]["session_id"], "session-fork");
881 assert_eq!(fake.record.lock().await["id"], "canonical-1");
882 assert_eq!(
883 store.get_runtime_thread_link("legacy").unwrap().as_deref(),
884 Some("canonical-1")
885 );
886 assert_eq!(
887 serde_json::to_value(store.snapshot_legacy_thread_history("legacy").unwrap()).unwrap(),
888 serde_json::to_value(before).unwrap()
889 );
890 server.abort();
891 }
892
893 #[tokio::test]
894 async fn confirmed_mutation_metadata_404_keeps_completed_effect_and_retry_identity() {
895 let (state, _temp, fake, server) = fixture(false).await;
896 fake.missing_after_commit
897 .store(true, std::sync::atomic::Ordering::SeqCst);
898 let request =
899 serde_json::from_value(json!({"kind":"start","operation_key":"completed-intent"})).unwrap();
900 let error = handle(&state, request).await.unwrap_err();
901 assert_eq!(fake.operations.lock().await.len(), 1);
902 assert!(error.message.contains("completed-intent"));
903 assert!(
904 error
905 .message
906 .contains("completed as thread canonical-1 / session session-1")
907 );
908 assert!(error.message.contains("without replay"));
909 assert_ne!(
910 error.code,
911 JsonRpcError::thread_not_found("canonical-1").code
912 );
913 server.abort();
914 }
915
916 #[tokio::test]
917 async fn encoded_oversize_creation_refuses_before_owner_effect() {
918 let (state, _temp, fake, server) = fixture(false).await;
919 let request = ThreadRequest::Create {
920 metadata: json!({"operation_key":"oversize-intent",
921 "system_prompt":"x".repeat(MAX_CANONICAL_HISTORY_BYTES)}),
922 };
923 let error = handle(&state, request).await.unwrap_err();
924 assert!(
925 error
926 .message
927 .contains("encoded canonical control exceeds bound")
928 );
929 assert!(fake.operations.lock().await.is_empty());
930 assert!(
931 fake.calls
932 .lock()
933 .await
934 .iter()
935 .all(|(_, path, _)| path.ends_with("/operations/lookup"))
936 );
937 server.abort();
938 }
939
940 #[tokio::test]
941 async fn declared_oversize_response_refuses_complete_document_without_truncation() {
942 crate::install_test_crypto_provider();
943 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
944 let address = listener.local_addr().unwrap();
945 let app = Router::new().route(
946 "/large",
947 get(|| async {
948 (
949 [(header::CONTENT_TYPE, "application/json")],
950 format!("\"{}\"", "x".repeat(MAX_CANONICAL_HISTORY_BYTES)),
951 )
952 }),
953 );
954 let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
955 let response = reqwest::get(format!("http://{address}/large"))
956 .await
957 .unwrap();
958 let error = read_json_response(response).await.unwrap_err();
959 assert!(format!("{error:#}").contains("complete-document bound"));
960 server.abort();
961 }
962
963 #[tokio::test]
964 async fn http_and_stdio_dispatch_share_the_same_full_history_owner() {
965 let (state, _temp, fake, server) = fixture(false).await;
966 let request = ThreadRequest::Read(codewhale_protocol::ThreadReadParams {
967 thread_id: "canonical-1".into(),
968 });
969 let http = thread_handler(State(state.clone()), Json(request.clone())).await;
970 assert_eq!(http.status(), StatusCode::OK);
971 let bytes = axum::body::to_bytes(http.into_body(), MAX_CANONICAL_HISTORY_BYTES)
972 .await
973 .unwrap();
974 let http: ThreadResponse = serde_json::from_slice(&bytes).unwrap();
975 let stdio = dispatch_stdio_request(
976 &state,
977 "thread/request",
978 serde_json::to_value(request).unwrap(),
979 )
980 .await
981 .unwrap();
982 let stdio: ThreadResponse = serde_json::from_value(stdio.result).unwrap();
983 assert_eq!(http.data["history"], stdio.data["history"]);
984 assert_eq!(
985 fake.calls
986 .lock()
987 .await
988 .iter()
989 .filter(|(_, path, _)| path.ends_with("/history"))
990 .count(),
991 2
992 );
993 server.abort();
994 }
995
996 #[tokio::test]
997 async fn completed_resume_retry_reads_original_operation_before_changed_history() {
998 let (state, _temp, fake, server) = fixture(false).await;
999 let request: ThreadRequest =
1000 serde_json::from_value(json!({"kind":"resume","thread_id":"canonical-1",
1001 "operation_key":"resume-recovery"}))
1002 .unwrap();
1003 let first = handle(&state, request.clone()).await.unwrap();
1004 *fake.document_digest.lock().await = "c".repeat(64);
1005 let count = fake.calls.lock().await.len();
1006 let second = handle(&state, request).await.unwrap();
1007 assert_eq!(first.data["receipt"], second.data["receipt"]);
1008 let calls = fake.calls.lock().await;
1009 assert!(
1010 calls[count..]
1011 .iter()
1012 .all(|(_, path, _)| path.ends_with("/operations/lookup")
1013 || path == "/v1/threads/canonical-1")
1014 );
1015 assert_eq!(
1016 calls
1017 .iter()
1018 .filter(|(_, path, _)| path == "/v1/thread-history/mutate")
1019 .count(),
1020 1
1021 );
1022 server.abort();
1023 }
1024
1025 #[tokio::test]
1026 async fn pending_retained_operation_never_reconstructs_or_replays_source() {
1027 let (state, _temp, fake, server) = fixture(false).await;
1028 let request: ThreadRequest =
1029 serde_json::from_value(json!({"kind":"resume","thread_id":"canonical-1",
1030 "operation_key":"pending-recovery"}))
1031 .unwrap();
1032 handle(&state, request.clone()).await.unwrap();
1033 fake.pending_operation
1034 .store(true, std::sync::atomic::Ordering::SeqCst);
1035 let count = fake.calls.lock().await.len();
1036 let error = handle(&state, request).await.unwrap_err();
1037 assert!(
1038 error.message.contains("pending-recovery")
1039 && error.message.contains("pending")
1040 && error.message.contains("no resume, replay or replacement")
1041 );
1042 let calls = fake.calls.lock().await;
1043 assert_eq!(calls.len(), count + 2);
1044 assert!(calls[count].1.ends_with("/operations/lookup"));
1045 assert!(calls[count + 1].1.ends_with("/operations/recover"));
1046 assert_eq!(calls[count + 1].2["operation"], calls[count].2);
1047 assert!(
1048 calls[count..]
1049 .iter()
1050 .all(|(_, path, _)| !path.ends_with("/history") && !path.ends_with("/mutate"))
1051 );
1052 server.abort();
1053 }
1054
1055 #[tokio::test]
1056 async fn prepared_pending_resume_settles_reserved_receipt_without_reconstructing_source() {
1057 let (state, _temp, fake, server) = fixture(false).await;
1058 let request: ThreadRequest =
1059 serde_json::from_value(json!({"kind":"resume","thread_id":"canonical-1",
1060 "operation_key":"prepared-recovery"}))
1061 .unwrap();
1062 let first = handle(&state, request.clone()).await.unwrap();
1063 fake.pending_operation
1064 .store(true, std::sync::atomic::Ordering::SeqCst);
1065 fake.recovery_prepared
1066 .store(true, std::sync::atomic::Ordering::SeqCst);
1067 *fake.document_digest.lock().await = "c".repeat(64);
1068 let count = fake.calls.lock().await.len();
1069 let second = handle(&state, request.clone()).await.unwrap();
1070 assert_eq!(first.data["receipt"], second.data["receipt"]);
1071 let calls = fake.calls.lock().await;
1072 assert_eq!(calls[count].1, "/v1/thread-history/operations/lookup");
1073 assert_eq!(calls[count + 1].1, "/v1/thread-history/operations/recover");
1074 assert_eq!(
1075 calls[count + 1].2["association"]["source_runtime_thread_id"],
1076 "canonical-1"
1077 );
1078 assert!(
1079 calls[count..]
1080 .iter()
1081 .all(|(_, path, _)| !path.ends_with("/history") && !path.ends_with("/mutate"))
1082 );
1083 drop(calls);
1084 let again = handle(&state, request).await.unwrap();
1085 assert_eq!(again.data["receipt"], first.data["receipt"]);
1086 assert_eq!(
1087 fake.calls
1088 .lock()
1089 .await
1090 .iter()
1091 .filter(|(_, path, _)| path.ends_with("/recover"))
1092 .count(),
1093 1
1094 );
1095 server.abort();
1096 }
1097
1098 #[tokio::test]
1099 async fn pending_wrong_action_refuses_before_any_prepared_recovery() {
1100 let (state, _temp, fake, server) = fixture(false).await;
1101 let create =
1102 serde_json::from_value(json!({"kind":"start","operation_key":"wrong-pending-action"}))
1103 .unwrap();
1104 handle(&state, create).await.unwrap();
1105 fake.pending_operation
1106 .store(true, std::sync::atomic::Ordering::SeqCst);
1107 fake.recovery_prepared
1108 .store(true, std::sync::atomic::Ordering::SeqCst);
1109 let count = fake.calls.lock().await.len();
1110 let request = serde_json::from_value(
1111 json!({"kind":"fork","thread_id":"canonical-1","operation_key":"wrong-pending-action"}),
1112 )
1113 .unwrap();
1114 let error = handle(&state, request).await.unwrap_err();
1115 assert!(
1116 error.message.contains("another action or source")
1117 && error.message.contains("wrong-pending-action")
1118 );
1119 let calls = fake.calls.lock().await;
1120 assert_eq!(calls.len(), count + 1);
1121 assert!(calls[count].1.ends_with("/lookup"));
1122 assert!(
1123 fake.pending_operation
1124 .load(std::sync::atomic::Ordering::SeqCst)
1125 );
1126 server.abort();
1127 }
1128
1129 #[tokio::test]
1130 async fn changed_prepared_target_refusal_keeps_reserved_operation_and_no_replay() {
1131 let (state, _temp, fake, server) = fixture(false).await;
1132 let request: ThreadRequest =
1133 serde_json::from_value(json!({"kind":"resume","thread_id":"canonical-1",
1134 "operation_key":"changed-prepared-target"}))
1135 .unwrap();
1136 handle(&state, request.clone()).await.unwrap();
1137 fake.pending_operation
1138 .store(true, std::sync::atomic::Ordering::SeqCst);
1139 fake.recovery_refusal
1140 .store(true, std::sync::atomic::Ordering::SeqCst);
1141 let count = fake.calls.lock().await.len();
1142 let error = handle(&state, request).await.unwrap_err();
1143 assert!(
1144 error.message.contains("changed-prepared-target")
1145 && error.message.contains("canonical-1 / session session-1")
1146 && error.message.contains("prepared target graph changed")
1147 );
1148 let calls = fake.calls.lock().await;
1149 assert_eq!(calls.len(), count + 2);
1150 assert!(
1151 calls[count..]
1152 .iter()
1153 .all(|(_, path, _)| path.ends_with("/lookup") || path.ends_with("/recover"))
1154 );
1155 assert_eq!(fake.operations.lock().await.len(), 1);
1156 server.abort();
1157 }
1158
1159 #[tokio::test]
1160 async fn retained_operation_rejects_another_action_or_source_with_completed_facts() {
1161 let (state, _temp, fake, server) = fixture(false).await;
1162 let create: ThreadRequest =
1163 serde_json::from_value(json!({"kind":"start","operation_key":"action-bound"})).unwrap();
1164 handle(&state, create).await.unwrap();
1165 let request = serde_json::from_value(json!({"kind":"fork","thread_id":"canonical-1",
1166 "operation_key":"action-bound"}))
1167 .unwrap();
1168 let error = handle(&state, request).await.unwrap_err();
1169 assert!(
1170 error
1171 .message
1172 .contains("completed as thread canonical-1 / session session-1")
1173 );
1174 assert!(error.message.contains("another action or source"));
1175 assert_eq!(fake.operations.lock().await.len(), 1);
1176 server.abort();
1177 }
1178
1179 #[tokio::test]
1180 async fn resume_options_are_forwarded_exactly_to_the_owning_decoder() {
1181 let (state, temp, fake, server) = fixture(false).await;
1182 let history = json!([{"role":"user","content":"offered"}]);
1183 let config = json!({"allow_shell":false,"mode":"ask","system_prompt":"captured override"});
1184 let source_path = temp.path().join("sessions/session-1.json");
1185 let request = serde_json::from_value(json!({"kind":"resume","thread_id":"canonical-1",
1186 "operation_key":"options-intent","history":history,"path":source_path,
1187 "model":"fixture-model","model_provider":"fixture-account","cwd":temp.path(),
1188 "approval_policy":"on-request","sandbox":"captured-ceiling","config":config,
1189 "base_instructions":"base","developer_instructions":"developer","personality":"brief",
1190 "persist_extended_history":false}))
1191 .unwrap();
1192 handle(&state, request).await.unwrap();
1193 let operations = fake.operations.lock().await;
1194 let options = &operations["options-intent"].0["mutation"]["options"];
1195 assert_eq!(options["offered_history"], history);
1196 assert_eq!(options["source_path"], json!(source_path));
1197 assert!(options.get("expected_session_goal_digest").is_none());
1198 assert_eq!(
1199 options["overrides"],
1200 json!({"model":"fixture-model","model_provider":"fixture-account",
1201 "cwd":temp.path(),"approval_policy":"on-request","sandbox":"captured-ceiling","config":config,
1202 "base_instructions":"base","developer_instructions":"developer","personality":"brief"})
1203 );
1204 server.abort();
1205 }
1206
1207 #[tokio::test]
1208 async fn fork_reads_other_workspace_source_but_uses_acknowledged_target_scope() {
1209 let (state, temp, fake, server) = fixture(false).await;
1210 let original = temp.path().join("original");
1211 fake.record.lock().await["workspace"] = json!(original);
1212 let request = serde_json::from_value(json!({"kind":"fork","thread_id":"canonical-1",
1213 "operation_key":"cross-workspace-fork","cwd":temp.path()}))
1214 .unwrap();
1215 let response = handle(&state, request).await.unwrap();
1216 assert_eq!(response.cwd.as_deref(), Some(temp.path()));
1217 assert_eq!(
1218 fake.operations.lock().await["cross-workspace-fork"].0["workspace"],
1219 json!(temp.path())
1220 );
1221 assert_eq!(fake.record.lock().await["workspace"], json!(original));
1222 let request = serde_json::from_value(json!({"kind":"resume","thread_id":"canonical-1",
1223 "operation_key":"cross-workspace-resume","cwd":temp.path()}))
1224 .unwrap();
1225 assert!(handle(&state, request).await.is_err());
1226 assert_eq!(fake.operations.lock().await.len(), 1);
1227 server.abort();
1228 }
1229
1230 #[tokio::test]
1231 async fn bridged_turn_sends_the_observed_owner_workspace_as_final_narrowing() {
1232 let (state, temp, fake, server) = fixture(false).await;
1233 run_http_thread_message(
1234 &state,
1235 "canonical-1".into(),
1236 "hello".into(),
1237 Vec::new(),
1238 None,
1239 )
1240 .await
1241 .unwrap();
1242 let calls = fake.calls.lock().await;
1243 let turn = calls
1244 .iter()
1245 .find(|(_, path, _)| path.ends_with("/turns"))
1246 .unwrap();
1247 assert_eq!(turn.2["expected_workspace"], json!(temp.path()));
1248 server.abort();
1249 }
1250
1251 fn legacy_goal_record(status: &str, objective: &str) -> codewhale_state::ThreadGoalRecord {
1252 serde_json::from_value(
1253 json!({"thread_id":"legacy","goal_id":"old-goal","objective":objective,
1254 "status":status,"token_budget":1000,"tokens_used":137,"time_used_seconds":23,
1255 "continuation_count":2,"created_at":1,"updated_at":2,"repeated_gap_count":0}),
1256 )
1257 .unwrap()
1258 }
1259
1260 #[tokio::test]
1261 async fn legacy_goal_remains_visible_and_cannot_silently_restart_or_be_replaced() {
1262 let (state, temp, fake, server) = fixture(false).await;
1263 let goal = legacy_goal_record("active", "retained work");
1264 let store = seed_archive(&state, "legacy", temp.path(), Vec::new(), Some(goal)).await;
1265 let before = store.snapshot_legacy_thread_history("legacy").unwrap();
1266 assert_eq!(
1267 before.goal.as_ref().unwrap().status,
1268 codewhale_protocol::ThreadGoalStatus::Active
1269 );
1270 let request = ThreadRequest::GoalGet(codewhale_protocol::ThreadGoalGetParams {
1271 thread_id: "legacy".into(),
1272 });
1273 let result = handle(&state, request.clone()).await.unwrap();
1274 assert_eq!(result.status, "ok");
1275 let recovered = result.goal.unwrap();
1276 assert_eq!(recovered.thread_id, "legacy");
1277 assert_eq!(recovered.goal_id, "old-goal");
1278 assert_eq!(recovered.tokens_used, 137);
1279 assert_eq!(recovered.time_used_seconds, 23);
1280 assert_eq!(recovered.continuation_count, 2);
1281 assert_eq!(recovered.created_at, 1);
1282 assert_eq!(recovered.updated_at, 2);
1283 assert_eq!(recovered.token_budget, Some(1000));
1284 assert_eq!(
1285 recovered.status,
1286 codewhale_protocol::ThreadGoalStatus::Paused
1287 );
1288 assert_eq!(recovered.pause_reason, None);
1289 assert_eq!(
1290 serde_json::to_value(store.snapshot_legacy_thread_history("legacy").unwrap()).unwrap(),
1291 serde_json::to_value(before).unwrap()
1292 );
1293 // Later owner progress is authoritative. A repeated alias read neither
1294 // imports the old goal again nor resets its real current progress.
1295 fake.goal.lock().await.as_mut().unwrap()["tokens_used"] = json!(151);
1296 let again = handle(&state, request).await.unwrap();
1297 assert_eq!(again.goal.unwrap().tokens_used, 151);
1298 let calls = fake.calls.lock().await;
1299 assert_eq!(
1300 calls
1301 .iter()
1302 .filter(|(_, p, _)| p.ends_with("/import"))
1303 .count(),
1304 1
1305 );
1306 assert!(
1307 calls
1308 .iter()
1309 .all(|(m, p, _)| !(p.ends_with("/turns") || p.ends_with("/goal") && m != "GET"))
1310 );
1311 server.abort();
1312 }
1313
1314 #[tokio::test]
1315 async fn changed_legacy_goal_during_import_refuses_alias_publication() {
1316 let (state, temp, fake, server) = fixture(true).await;
1317 let store = seed_archive(
1318 &state,
1319 "legacy",
1320 temp.path(),
1321 Vec::new(),
1322 Some(legacy_goal_record("paused", "original")),
1323 )
1324 .await;
1325 let worker_state = state.clone();
1326 let waiter = tokio::spawn(async move { resolve(&worker_state, "legacy", true).await });
1327 fake.import_seen.notified().await;
1328 rusqlite::Connection::open(store.db_path())
1329 .unwrap()
1330 .execute(
1331 "UPDATE thread_goals SET objective='changed during import' WHERE thread_id='legacy'",
1332 [],
1333 )
1334 .unwrap();
1335 fake.import_release.notify_one();
1336 let error = waiter.await.unwrap().unwrap_err();
1337 assert!(format!("{error:#}").contains("legacy history changed"));
1338 assert!(
1339 store
1340 .get_canonical_runtime_link("legacy", state.captured_owner.as_ref().unwrap())
1341 .unwrap()
1342 .is_none()
1343 );
1344 assert_eq!(
1345 store.get_thread_goal("legacy").unwrap().unwrap().objective,
1346 "changed during import"
1347 );
1348 assert_eq!(
1349 fake.goal.lock().await.as_ref().unwrap()["objective"],
1350 "original"
1351 );
1352 server.abort();
1353 }
1354
1355 #[tokio::test]
1356 async fn conflicting_existing_owner_goal_is_visible_refusal_without_empty_alias() {
1357 let (state, temp, fake, server) = fixture(false).await;
1358 let store = seed_archive(
1359 &state,
1360 "legacy",
1361 temp.path(),
1362 Vec::new(),
1363 Some(legacy_goal_record("paused", "legacy objective")),
1364 )
1365 .await;
1366 let mut existing =
1367 serde_json::to_value(legacy_goal_record("paused", "owned objective")).unwrap();
1368 existing["thread_id"] = json!("canonical-1");
1369 existing["goal_id"] = json!("owned-goal");
1370 *fake.goal.lock().await = Some(existing.clone());
1371 let request = ThreadRequest::GoalGet(codewhale_protocol::ThreadGoalGetParams {
1372 thread_id: "legacy".into(),
1373 });
1374 let error = handle(&state, request).await.unwrap_err();
1375 assert!(error.message.contains("canonical goal conflict"));
1376 assert!(
1377 store
1378 .get_canonical_runtime_link("legacy", state.captured_owner.as_ref().unwrap())
1379 .unwrap()
1380 .is_none()
1381 );
1382 assert_eq!(
1383 store.get_thread_goal("legacy").unwrap().unwrap().objective,
1384 "legacy objective"
1385 );
1386 assert_eq!(fake.goal.lock().await.as_ref().unwrap(), &existing);
1387 server.abort();
1388 }
1389
1389 lines RUST