返回 CodeWhale
thread_history.rs
根目录 / crates / tui / src / runtime_api / thread_history.rs
1 //! Full graph migration into the held canonical Runtime. The old SQLite file
2 //! remains recovery evidence; this adapter never becomes another transcript writer.
3 use std::collections::HashMap;
4 use std::path::Path;
5 use std::sync::Arc;
6
7 use anyhow::{Context, Result, anyhow, bail, ensure};
8 use axum::{Json, extract::State};
9 use chrono::{DateTime, Utc};
10 use codewhale_models::{ContentBlock, Message, Role};
11 use codewhale_protocol::{
12 CanonicalHistoryImportRequest, CanonicalHistoryOptions, CanonicalHistorySource,
13 CanonicalThreadMutation, CanonicalThreadMutationRequest, CanonicalThreadReceipt,
14 LegacyThreadHistory, MAX_CANONICAL_HISTORY_BYTES, MAX_CANONICAL_HISTORY_ENTRIES,
15 };
16
17 use super::{ApiError, RuntimeApiState};
18 use crate::runtime_threads::{
19 CreateThreadRequest, RuntimeHistoryOperation, RuntimeHistoryWitness, RuntimeThreadManager,
20 ThreadListFilter,
21 };
22 use crate::session_manager::{SavedSession, SessionManager};
23 use crate::session_tree::{SessionEntry, SessionEntryKind, SessionImportContainer, SessionJournal};
24
25 pub(crate) async fn history_owner_work<T: Send + 'static>(
26 work: impl FnOnce() -> Result<T> + Send + 'static,
27 ) -> Result<T> {
28 #[cfg(test)]
29 let ticket = crate::test_support::env_scope_ticket();
30 codewhale_app_server::daemon_socket::owner_work(move || {
31 #[cfg(test)]
32 let _membership = crate::test_support::join_env_scope(ticket);
33 work()
34 })
35 .await
36 }
37
38 /// Validate every branch before traversing the active projection or creating a
39 /// canonical record. The existing journal validator alone does not reject cycles
40 /// or duplicate IDs, so this boundary must establish those stronger invariants.
41 fn legacy_journal(history: &LegacyThreadHistory) -> Result<SessionJournal> {
42 ensure!(
43 history.version == 1,
44 "unsupported legacy history schema; source retained for recovery"
45 );
46 ensure!(
47 history.state_store_id.len() == 64
48 && history
49 .state_store_id
50 .bytes()
51 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()),
52 "invalid source database identity"
53 );
54 ensure!(
55 !history.thread_id.is_empty() && history.thread_id.len() <= 128,
56 "invalid source thread identity"
57 );
58 ensure!(
59 history.messages.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
60 "legacy history entry bound exceeded"
61 );
62 let mut parents = HashMap::with_capacity(history.messages.len());
63 for row in &history.messages {
64 ensure!(
65 row.id > 0
66 && row.thread_id == history.thread_id
67 && parents.insert(row.id, row.parent_entry_id).is_none(),
68 "duplicate or foreign legacy history entry"
69 );
70 }
71 for parent in parents.values().flatten() {
72 ensure!(
73 parents.contains_key(parent),
74 "dangling legacy history parent; source retained for recovery"
75 );
76 }
77 match history.current_leaf_id {
78 Some(leaf) => ensure!(parents.contains_key(&leaf), "legacy history leaf is absent"),
79 None => ensure!(
80 parents.is_empty(),
81 "nonempty legacy history has no selected leaf"
82 ),
83 }
84 if let Some(goal) = history.goal.as_ref() {
85 ensure!(
86 goal.thread_id == history.thread_id,
87 "legacy goal belongs to another source thread"
88 );
89 goal.validate_stall_state().map_err(anyhow::Error::msg)?;
90 }
91 // Three-color iterative traversal bounds work and stack even for a deeply
92 // nested imported branch; no recursive call uses untrusted graph depth.
93 let mut colors = HashMap::with_capacity(parents.len());
94 for id in parents.keys().copied() {
95 let mut cursor = Some(id);
96 let mut path = Vec::new();
97 while let Some(id) = cursor {
98 match colors.get(&id) {
99 Some(1) => {
100 return Err(anyhow!(
101 "cyclic legacy history; source retained for recovery"
102 ));
103 }
104 Some(2) => break,
105 _ => {
106 colors.insert(id, 1);
107 path.push(id);
108 cursor = parents[&id];
109 }
110 }
111 }
112 for id in path {
113 colors.insert(id, 2);
114 }
115 }
116 let id = |id: i64| format!("legacy:{}:{id}", history.state_store_id);
117 let mut entries = Vec::with_capacity(history.messages.len());
118 for row in &history.messages {
119 let message = match (&row.item, row.role.as_str()) {
120 (Some(item), "history") => {
121 ensure!(
122 serde_json::from_str::<serde_json::Value>(&row.content)
123 .ok()
124 .as_ref()
125 == Some(item),
126 "legacy history payload and content disagree; source retained for recovery"
127 );
128 let message = serde_json::from_value::<Message>(item.clone()).context("opaque legacy item cannot be imported as model content; source retained for recovery")?;
129 ensure!(
130 serde_json::to_value(&message)? == *item,
131 "legacy message has unsupported fields; source retained for recovery"
132 );
133 message
134 }
135 (None, role @ ("user" | "assistant" | "system")) => Message {
136 role: match role {
137 "user" => Role::User,
138 "assistant" => Role::Assistant,
139 _ => Role::System,
140 },
141 content: vec![ContentBlock::Text {
142 text: row.content.clone(),
143 cache_control: None,
144 }],
145 },
146 _ => {
147 return Err(anyhow!(
148 "unsupported legacy message representation; source retained for recovery"
149 ));
150 }
151 };
152 if message.role == Role::User {
153 crate::image_attach::runtime_images_from_blocks(&message.content)
154 .map_err(|error| anyhow!("invalid imported image: {error}"))?;
155 }
156 entries.push(SessionEntry {
157 id: id(row.id),
158 parent_id: row.parent_entry_id.map(id),
159 kind: SessionEntryKind::Message { message },
160 created_at: DateTime::<Utc>::from_timestamp(row.created_at, 0)
161 .ok_or_else(|| anyhow!("invalid legacy history timestamp"))?,
162 spawn_depth: 0,
163 });
164 }
165 let journal = SessionJournal {
166 entries,
167 leaf_id: history.current_leaf_id.map(id),
168 schema_version: 1,
169 spawn_depth: 0,
170 };
171 journal.validate().map_err(anyhow::Error::msg)?;
172 Ok(journal)
173 }
174
175 pub(super) async fn import_thread_history(
176 State(state): State<RuntimeApiState>,
177 Json(request): Json<CanonicalHistoryImportRequest>,
178 ) -> Result<Json<CanonicalThreadReceipt>, ApiError> {
179 ensure_request_scope(&state.runtime_threads, &request)
180 .map_err(|error| ApiError::bad_request(error.to_string()))?;
181 // Immediate refusal avoids an unbounded queue retaining whole histories.
182 let admission = state
183 .runtime_threads
184 .try_history_import_guard()
185 .map_err(|error| ApiError::conflict(error.to_string()))?;
186 // The actual owner job keeps the admission and request after the HTTP waiter
187 // disappears. A retry observes its durable identity, never a fresh creation.
188 let job = tokio::spawn(async move {
189 let _admission = admission;
190 import_into_owner(state, request).await
191 });
192 job.await
193 .map_err(|_| {
194 ApiError::internal("canonical history worker did not settle; retry the same operation")
195 })?
196 .map(Json)
197 .map_err(|error| ApiError::conflict(error.to_string()))
198 }
199
200 fn ensure_request_scope(
201 runtime: &RuntimeThreadManager,
202 request: &CanonicalHistoryImportRequest,
203 ) -> Result<()> {
204 ensure!(
205 request.version == 1,
206 "unsupported canonical history request schema"
207 );
208 let binding = runtime.session_store_binding();
209 ensure!(
210 request.expected_data_dir == binding.data_dir
211 && request.expected_execution_scope == binding.execution_scope,
212 "canonical history request belongs to a different held owner store"
213 );
214 ensure!(
215 !request.operation_key.is_empty()
216 && request.operation_key.len() <= 128
217 && !request.operation_key.chars().any(char::is_control),
218 "invalid history operation key"
219 );
220 ensure!(
221 serde_json::to_vec(request)?.len() <= MAX_CANONICAL_HISTORY_BYTES,
222 "encoded history request exceeds canonical import bound"
223 );
224 Ok(())
225 }
226
227 async fn import_into_owner(
228 state: RuntimeApiState,
229 request: CanonicalHistoryImportRequest,
230 ) -> Result<CanonicalThreadReceipt> {
231 let runtime = Arc::clone(&state.runtime_threads);
232 // Validate and fingerprint the bounded graph on the existing filesystem/
233 // owner worker boundary before retaining or mutating canonical records.
234 let (request, journal, request_digest, history_digest) = history_owner_work(move || {
235 let journal = legacy_journal(&request.history)?;
236 let request_digest = crate::hashing::sha256_hex(
237 crate::client::canonical_json(&serde_json::to_value(&request)?).as_bytes(),
238 );
239 let history_digest = crate::hashing::sha256_hex(
240 crate::client::canonical_json(&serde_json::to_value(&journal)?).as_bytes(),
241 );
242 Ok((request, journal, request_digest, history_digest))
243 })
244 .await?;
245 let target = if let Some(id) = request.target_runtime_thread_id.as_deref() {
246 let thread = runtime.get_thread(id).await?;
247 ensure!(
248 request.workspace == thread.workspace,
249 "existing canonical target belongs to a different selected workspace"
250 );
251 let session = thread.session_id.ok_or_else(|| {
252 anyhow!("existing canonical target has no bound full document; recovery required")
253 })?;
254 Some((id.to_string(), session))
255 } else {
256 None
257 };
258 let target_for_reservation = target.clone();
259 let reservation_runtime = Arc::clone(&runtime);
260 let key = request.operation_key.clone();
261 let operation = history_owner_work(move || match target_for_reservation.as_ref() {
262 Some((thread, session)) => reservation_runtime.reserve_history_operation_for_target(
263 &key,
264 &request_digest,
265 &history_digest,
266 Some((thread.as_str(), session.as_str())),
267 ),
268 None => {
269 reservation_runtime.reserve_history_operation(&key, &request_digest, &history_digest)
270 }
271 })
272 .await?;
273 if operation.committed {
274 // A missing committed target is lost history, never a reason to remint.
275 let thread = runtime
276 .get_thread(&operation.receipt.runtime_thread_id)
277 .await?;
278 ensure!(
279 thread.session_id.as_deref() == Some(operation.receipt.session_id.as_str()),
280 "committed history target changed; recovery required"
281 );
282 let sessions_dir = state.sessions_dir.clone();
283 let binding = runtime.session_store_binding();
284 let expected = journal.clone();
285 history_owner_work(move || {
286 verify_import_checkpoint(
287 &SessionManager::new(sessions_dir)?,
288 &thread,
289 &binding,
290 &expected,
291 )
292 })
293 .await?;
294 let goal_runtime = Arc::clone(&runtime);
295 let original_thread = request.history.thread_id.clone();
296 let goal = request.history.goal.clone();
297 let expected = operation.clone();
298 history_owner_work(move || {
299 goal_runtime.adopt_history_goal(&expected, &original_thread, goal)
300 })
301 .await?;
302 return Ok(operation.receipt);
303 }
304 let lookup_runtime = Arc::clone(&runtime);
305 let thread_id = operation.receipt.runtime_thread_id.clone();
306 let existing =
307 history_owner_work(move || lookup_runtime.history_operation_thread(&thread_id)).await?;
308 let thread = match existing {
309 Some(thread) => thread,
310 None => {
311 runtime
312 .create_thread_with_reserved_id(
313 CreateThreadRequest {
314 workspace: Some(request.workspace),
315 model: request.model,
316 ..Default::default()
317 },
318 state.config_path.as_deref(),
319 state.config_profile.as_deref(),
320 Some(operation.receipt.runtime_thread_id.clone()),
321 )
322 .await?
323 }
324 };
325 let mut session = SavedSession::import_foreign(
326 SessionImportContainer::new("legacy-state-sqlite".into(), &journal, None),
327 thread.workspace.clone(),
328 thread.model.clone(),
329 )
330 .map_err(anyhow::Error::msg)?;
331 session.metadata.id = operation.receipt.session_id.clone();
332 session.system_prompt.clone_from(&thread.system_prompt);
333 session.metadata.created_at = operation.created_at;
334 session.metadata.updated_at = operation.created_at;
335 session.metadata.runtime_store = Some(runtime.session_store_binding());
336 session.metadata.set_model_provider_route(
337 thread
338 .model_provider
339 .as_deref()
340 .ok_or_else(|| anyhow!("canonical thread has no captured provider"))?,
341 thread.model_provider_id.as_deref(),
342 );
343 let mut captured_target_document_digest = None;
344 if target.is_some() {
345 let detail = runtime.get_thread_detail(&thread.id).await?;
346 ensure!(
347 !super::sessions::thread_detail_has_live_work(&detail),
348 "canonical target has live work; retry the same import after settlement"
349 );
350 let full = snapshot_in_owner(state.clone(), thread.id.clone())
351 .await
352 .map_err(|error| anyhow!(error.message))?;
353 captured_target_document_digest = full.saved_document_digest;
354 ensure!(
355 captured_target_document_digest.is_some(),
356 "canonical target full document witness is absent"
357 );
358 let mut observed: SavedSession = serde_json::from_value(full.session)?;
359 ensure!(
360 observed.metadata.id == operation.receipt.session_id,
361 "canonical target session changed before source merge"
362 );
363 let imported = session
364 .journal
365 .take()
366 .ok_or_else(|| anyhow!("source journal missing"))?;
367 let destination = observed
368 .journal
369 .as_mut()
370 .ok_or_else(|| anyhow!("target full journal missing"))?;
371 let mut positions: HashMap<String, usize> = destination
372 .entries
373 .iter()
374 .enumerate()
375 .map(|(index, entry)| (entry.id.clone(), index))
376 .collect();
377 for entry in imported.entries {
378 if let Some(index) = positions.get(&entry.id).copied() {
379 ensure!(
380 destination.entries[index] == entry,
381 "imported entry identity conflicts with canonical history"
382 );
383 } else {
384 ensure!(
385 destination.entries.len() < MAX_CANONICAL_HISTORY_ENTRIES,
386 "merged full history exceeds entry transport bound; both sources retained"
387 );
388 positions.insert(entry.id.clone(), destination.entries.len());
389 destination.entries.push(entry);
390 }
391 }
392 if destination.leaf_id.is_none() {
393 destination.leaf_id = imported.leaf_id;
394 observed.leaf_id = destination.leaf_id.clone();
395 observed.messages = destination.to_messages();
396 }
397 ensure!(
398 destination.entries.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
399 "merged full history exceeds entry transport bound; both sources retained"
400 );
401 let mut bound = HistorySizeBound(MAX_CANONICAL_HISTORY_BYTES);
402 serde_json::to_writer(&mut bound, &observed)?;
403 session = observed;
404 }
405 let target_merge = target.is_some();
406 let sessions_dir = state.sessions_dir.clone();
407 let (session, _target_lease) = history_owner_work(move || {
408 let manager = SessionManager::new(sessions_dir)?;
409 let lease = manager.reserve_session_for_external_write(&session.metadata.id)?;
410 match manager
411 .load_session_snapshot_bounded(&session.metadata.id, MAX_CANONICAL_HISTORY_BYTES)
412 {
413 Ok(observed) if target_merge => {
414 ensure!(
415 Some(saved_document_digest(&observed)?) == captured_target_document_digest,
416 "canonical full document changed before source merge; retry the same operation"
417 );
418 ensure!(
419 observed.metadata.runtime_store.as_ref().is_none_or(
420 |binding| Some(binding) == session.metadata.runtime_store.as_ref()
421 ) && observed.metadata.workspace == session.metadata.workspace,
422 "canonical target store binding changed before merge"
423 );
424 manager.save_session(&session)?;
425 return Ok((session, lease));
426 }
427 Ok(observed) => {
428 ensure!(
429 observed.journal == session.journal
430 && observed.metadata.runtime_store == session.metadata.runtime_store
431 && observed.metadata.workspace == session.metadata.workspace,
432 "reserved imported session diverged; recovery required"
433 );
434 return Ok((observed, lease));
435 }
436 Err(error) if error.kind() == std::io::ErrorKind::NotFound && !target_merge => {}
437 Err(error) => return Err(error.into()),
438 }
439 manager.save_session(&session)?;
440 Ok((session, lease))
441 })
442 .await?;
443 let detail = runtime.get_thread_detail(&thread.id).await?;
444 ensure!(
445 !super::sessions::thread_detail_has_live_work(&detail),
446 "reserved imported thread has active work; recovery required"
447 );
448 if detail.turns.is_empty() && detail.items.is_empty() {
449 runtime
450 .seed_thread_from_messages_with_history_operation(
451 &thread.id,
452 &session.messages,
453 Some(&operation),
454 &session.messages,
455 )
456 .await?;
457 } else if !target_merge {
458 let sessions_dir = state.sessions_dir.clone();
459 let binding = runtime.session_store_binding();
460 let expected = journal.clone();
461 let observed_thread = runtime.get_thread(&thread.id).await?;
462 history_owner_work(move || {
463 verify_import_checkpoint_under_lease(
464 &SessionManager::new(sessions_dir)?,
465 &observed_thread,
466 &binding,
467 &expected,
468 )
469 })
470 .await?;
471 }
472 runtime
473 .set_thread_session_checkpoint(&thread.id, &session)
474 .await?;
475 let goal_runtime = Arc::clone(&runtime);
476 let original_thread = request.history.thread_id.clone();
477 let goal = request.history.goal.clone();
478 let expected = operation.clone();
479 history_owner_work(move || goal_runtime.adopt_history_goal(&expected, &original_thread, goal))
480 .await?;
481 let commit_runtime = Arc::clone(&runtime);
482 history_owner_work(move || commit_runtime.commit_history_operation(operation)).await
483 }
484
485 fn verify_import_checkpoint(
486 manager: &SessionManager,
487 thread: &crate::runtime_threads::ThreadRecord,
488 binding: &crate::runtime_threads::RuntimeStoreBinding,
489 expected: &SessionJournal,
490 ) -> Result<()> {
491 let id = thread
492 .session_id
493 .as_deref()
494 .ok_or_else(|| anyhow!("committed history target has no session; recovery required"))?;
495 let _lease = manager.reserve_session_for_external_write(id)?;
496 verify_import_checkpoint_under_lease(manager, thread, binding, expected)
497 }
498
499 fn verify_import_checkpoint_under_lease(
500 manager: &SessionManager,
501 thread: &crate::runtime_threads::ThreadRecord,
502 binding: &crate::runtime_threads::RuntimeStoreBinding,
503 expected: &SessionJournal,
504 ) -> Result<()> {
505 let id = thread
506 .session_id
507 .as_deref()
508 .ok_or_else(|| anyhow!("canonical target session is absent"))?;
509 let observed = manager.load_session_snapshot_bounded(id, MAX_CANONICAL_HISTORY_BYTES)?;
510 ensure!(
511 observed.metadata.runtime_store.as_ref() == Some(binding)
512 && observed.metadata.workspace == thread.workspace,
513 "committed history session binding changed; recovery required"
514 );
515 let checkpoint = thread
516 .saved_session_checkpoint
517 .as_ref()
518 .ok_or_else(|| anyhow!("committed history has no verifiable checkpoint"))?;
519 ensure!(
520 crate::runtime_threads::checkpoint_prefix_len(checkpoint, &observed.messages)?.is_some(),
521 "committed history checkpoint differs from its document; recovery required"
522 );
523 let journal = observed
524 .journal
525 .as_ref()
526 .ok_or_else(|| anyhow!("committed full history journal is missing"))?;
527 let entries: HashMap<_, _> = journal
528 .entries
529 .iter()
530 .map(|entry| (entry.id.as_str(), entry))
531 .collect();
532 ensure!(
533 entries.len() == journal.entries.len()
534 && expected.entries.iter().all(|entry| entries
535 .get(entry.id.as_str())
536 .is_some_and(|observed| *observed == entry)),
537 "committed full history branches changed; recovery required"
538 );
539 Ok(())
540 }
541
542 /// Full graph read through the existing owner Engine and session checkpoint.
543 /// It never publishes an alias or becomes a second journal writer.
544 pub(super) async fn snapshot_thread_history(
545 State(state): State<RuntimeApiState>,
546 axum::extract::Path(id): axum::extract::Path<String>,
547 ) -> Result<Json<codewhale_protocol::CanonicalThreadSnapshot>, ApiError> {
548 let _admission = state.runtime_threads.session_checkpoint_guard().await;
549 snapshot_in_owner(state, id).await.map(Json)
550 }
551
552 async fn snapshot_in_owner(
553 state: RuntimeApiState,
554 id: String,
555 ) -> Result<codewhale_protocol::CanonicalThreadSnapshot, ApiError> {
556 snapshot_in_runtime(&state.runtime_threads, &state.sessions_dir, id).await
557 }
558
559 async fn snapshot_in_runtime(
560 runtime: &Arc<RuntimeThreadManager>,
561 sessions_dir: &Path,
562 id: String,
563 ) -> Result<codewhale_protocol::CanonicalThreadSnapshot, ApiError> {
564 let runtime = Arc::clone(runtime);
565 let thread = runtime
566 .get_thread(&id)
567 .await
568 .map_err(super::map_thread_err)?;
569 let bound_id = thread.session_id.clone();
570 let expected_binding = runtime.session_store_binding();
571 let document = if let Some(session_id) = bound_id.clone() {
572 let sessions_dir = sessions_dir.to_path_buf();
573 let checkpoint = thread.saved_session_checkpoint.clone();
574 let binding = expected_binding.clone();
575 let workspace = thread.workspace.clone();
576 Some(
577 history_owner_work(move || {
578 let manager = SessionManager::new(sessions_dir)?;
579 let _lease = manager.reserve_session_for_external_write(&session_id)?;
580 let document = manager
581 .load_session_snapshot_bounded(&session_id, MAX_CANONICAL_HISTORY_BYTES)?;
582 ensure!(
583 document.metadata.workspace == workspace,
584 "saved history workspace differs from its held thread"
585 );
586 if let Some(saved_binding) = document.metadata.runtime_store.as_ref() {
587 ensure!(
588 saved_binding == &binding,
589 "saved history belongs to a different Runtime store"
590 );
591 }
592 let checkpoint = checkpoint
593 .ok_or_else(|| anyhow!("bound history has no verifiable checkpoint"))?;
594 ensure!(
595 crate::runtime_threads::checkpoint_prefix_len(&checkpoint, &document.messages)?
596 .is_some(),
597 "bound full history checkpoint changed; recovery required"
598 );
599 let digest = saved_document_digest(&document)?;
600 Ok((document, digest, manager.load_session_goal(&session_id)?))
601 })
602 .await
603 .map_err(|error| ApiError::conflict(error.to_string()))?,
604 )
605 } else {
606 None
607 };
608 // The snapshot is taken on the one existing Engine mailbox. A saved prefix
609 // can lag a live turn, so merge that actual projection in memory while
610 // retaining all prior branches. No read writes or repairs the saved file.
611 let engine = runtime
612 .get_engine(&id)
613 .await
614 .map_err(super::map_thread_err)?;
615 let snapshot = engine
616 .get_session_snapshot()
617 .await
618 .map_err(|error| ApiError::internal(error.to_string()))?;
619 let response = history_owner_work(move || {
620 let persisted_entry_count = document
621 .as_ref()
622 .and_then(|(saved, _, _)| saved.journal.as_ref())
623 .map_or(0, |journal| journal.entries.len());
624 let previous_updated_at = document
625 .as_ref()
626 .map(|(saved, _, _)| saved.metadata.updated_at);
627 let has_projection_delta = document
628 .as_ref()
629 .is_none_or(|(saved, _, _)| saved.messages != snapshot.messages);
630 let source_goal_digest =
631 session_goal_digest(&document.as_ref().and_then(|(_, _, goal)| goal.clone()))?;
632 let persisted_document_digest = document.as_ref().map(|(_, digest, _)| digest.clone());
633 let mut session = match document {
634 Some((existing, _, _)) => crate::session_manager::update_session(
635 existing,
636 &snapshot.messages,
637 snapshot.total_tokens,
638 snapshot.system_prompt.as_ref(),
639 ),
640 None => crate::session_manager::create_saved_session_with_id_and_mode(
641 snapshot.session_id.clone(),
642 &snapshot.messages,
643 &snapshot.model,
644 &snapshot.workspace,
645 snapshot.total_tokens,
646 snapshot.system_prompt.as_ref(),
647 Some(&snapshot.mode),
648 ),
649 };
650 // The read projection uses the actual thread boundary, rather than
651 // fresh UUIDs and the reader's wall clock. Re-reading unchanged Engine
652 // content must not fabricate a changed full-history digest.
653 stable_projection_entries(&mut session, persisted_entry_count, &thread)?;
654 if previous_updated_at.is_none() {
655 session.metadata.created_at = thread.created_at;
656 }
657 session.metadata.updated_at = if has_projection_delta {
658 thread.updated_at
659 } else {
660 previous_updated_at.unwrap_or(thread.updated_at)
661 };
662 session.metadata.model = snapshot.model;
663 session.metadata.set_model_provider_route(
664 &snapshot.model_provider,
665 snapshot.model_provider_id.as_deref(),
666 );
667 session.metadata.runtime_store = Some(expected_binding.clone());
668 let graph = session
669 .journal
670 .as_ref()
671 .ok_or_else(|| anyhow!("full journal is absent"))?;
672 ensure!(
673 graph.entries.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
674 "full history exceeds entry transport bound; source retained"
675 );
676 // Measure before constructing the Value allocation. Refuse the complete
677 // result rather than truncating branches or returning a partial graph.
678 let mut size = HistorySizeBound(MAX_CANONICAL_HISTORY_BYTES);
679 serde_json::to_writer(&mut size, &session)?;
680 let response = codewhale_protocol::CanonicalThreadSnapshot {
681 version: 1,
682 data_dir: expected_binding.data_dir,
683 execution_scope: expected_binding.execution_scope,
684 runtime_thread_id: id,
685 saved_session_id: bound_id,
686 saved_document_digest: persisted_document_digest,
687 document_digest: history_source_digest(&session, &source_goal_digest)?,
688 session_goal_digest: source_goal_digest,
689 session: serde_json::to_value(session)?,
690 };
691 let mut size = HistorySizeBound(MAX_CANONICAL_HISTORY_BYTES);
692 serde_json::to_writer(&mut size, &response)?;
693 Ok(response)
694 })
695 .await
696 .map_err(|error| ApiError::conflict(error.to_string()))?;
697 Ok(response)
698 }
699
700 fn stable_projection_entries(
701 session: &mut SavedSession,
702 retained_count: usize,
703 thread: &crate::runtime_threads::ThreadRecord,
704 ) -> Result<()> {
705 let journal = session
706 .journal
707 .as_mut()
708 .ok_or_else(|| anyhow!("full journal is absent"))?;
709 let mut ids = HashMap::new();
710 for (index, entry) in journal.entries.iter_mut().enumerate().skip(retained_count) {
711 let old = entry.id.clone();
712 let digest = crate::hashing::sha256_hex(
713 crate::client::canonical_json(&serde_json::to_value(&entry.kind)?).as_bytes(),
714 );
715 entry.id = format!(
716 "projection:{}:{}:{index}:{digest}",
717 thread.id,
718 thread.latest_turn_id.as_deref().unwrap_or("initial")
719 );
720 entry.created_at = thread.updated_at;
721 ids.insert(old, entry.id.clone());
722 }
723 for entry in journal.entries.iter_mut().skip(retained_count) {
724 if let Some(id) = entry.parent_id.as_mut()
725 && let Some(replacement) = ids.get(id)
726 {
727 *id = replacement.clone();
728 }
729 }
730 if let Some(leaf) = journal.leaf_id.as_mut()
731 && let Some(replacement) = ids.get(leaf)
732 {
733 *leaf = replacement.clone();
734 }
735 session.leaf_id = journal.leaf_id.clone();
736 Ok(())
737 }
738
739 pub(crate) fn saved_document_digest(session: &SavedSession) -> Result<String> {
740 let mut size = HistorySizeBound(MAX_CANONICAL_HISTORY_BYTES);
741 serde_json::to_writer(&mut size, session)?;
742 Ok(crate::hashing::sha256_hex(
743 crate::client::canonical_json(&serde_json::to_value(session)?).as_bytes(),
744 ))
745 }
746
747 pub(super) async fn mutate_thread_history(
748 State(state): State<RuntimeApiState>,
749 Json(request): Json<CanonicalThreadMutationRequest>,
750 ) -> Result<Json<CanonicalThreadReceipt>, ApiError> {
751 mutate_thread_history_in_runtime(
752 &state.runtime_threads,
753 &state.sessions_dir,
754 state.config_path.as_deref(),
755 state.config_profile.as_deref(),
756 request,
757 )
758 .await
759 .map(Json)
760 .map_err(|error| ApiError::conflict(error.to_string()))
761 }
762
763 /// The HTTP frontend and the inactive local mounted handoff invoke this same
764 /// actual held manager. It creates no transport, guest manager or policy owner.
765 pub(crate) async fn mutate_thread_history_in_runtime(
766 runtime: &Arc<RuntimeThreadManager>,
767 sessions_dir: &Path,
768 config_path: Option<&Path>,
769 config_profile: Option<&str>,
770 request: CanonicalThreadMutationRequest,
771 ) -> Result<CanonicalThreadReceipt> {
772 let binding = runtime.session_store_binding();
773 ensure!(
774 request.version == 1
775 && request.expected_data_dir == binding.data_dir
776 && request.expected_execution_scope == binding.execution_scope,
777 "canonical mutation belongs to a different held owner store"
778 );
779 ensure!(
780 !request.operation_key.is_empty()
781 && request.operation_key.len() <= 128
782 && !request.operation_key.chars().any(char::is_control),
783 "invalid canonical operation key"
784 );
785 let mut size = HistorySizeBound(MAX_CANONICAL_HISTORY_BYTES);
786 serde_json::to_writer(&mut size, &request)?;
787 let admission = runtime.try_history_import_guard()?;
788 let runtime = Arc::clone(runtime);
789 let sessions_dir = sessions_dir.to_path_buf();
790 let config_path = config_path.map(Path::to_path_buf);
791 let config_profile = config_profile.map(str::to_string);
792 // Cancellation of a socket/HTTP waiter cannot release admission or abandon
793 // the reserved write. The caller retries only the same captured intent.
794 tokio::spawn(async move {
795 let _admission = admission;
796 mutate_in_owner(runtime, sessions_dir, config_path, config_profile, request).await
797 })
798 .await
799 .context("canonical operation outcome is uncertain; retain and retry the same operation key")?
800 }
801
802 fn journal_digest(journal: &SessionJournal) -> Result<String> {
803 Ok(crate::hashing::sha256_hex(
804 crate::client::canonical_json(&serde_json::to_value(journal)?).as_bytes(),
805 ))
806 }
807
808 pub(crate) fn session_goal_digest(
809 goal: &Option<crate::session_manager::SessionGoalState>,
810 ) -> Result<String> {
811 Ok(crate::hashing::sha256_hex(
812 crate::client::canonical_json(&serde_json::to_value(goal)?).as_bytes(),
813 ))
814 }
815
816 fn history_source_digest(session: &SavedSession, goal_digest: &str) -> Result<String> {
817 Ok(crate::hashing::sha256_hex(
818 format!(
819 "codewhale:full-history-source:v1\0{}\0{goal_digest}",
820 saved_document_digest(session)?
821 )
822 .as_bytes(),
823 ))
824 }
825
826 fn forked_session_goal(
827 mut goal: Option<crate::session_manager::SessionGoalState>,
828 ) -> Option<crate::session_manager::SessionGoalState> {
829 if let Some(goal) = goal.as_mut()
830 && goal.status == crate::session_manager::SessionGoalStatus::Active
831 {
832 goal.status = crate::session_manager::SessionGoalStatus::Paused;
833 goal.pause_reason = None;
834 }
835 goal
836 }
837
838 fn operation_witness(journal: &SessionJournal, workspace: &Path) -> RuntimeHistoryWitness {
839 RuntimeHistoryWitness {
840 entries_len: journal.entries.len(),
841 leaf_id: journal.leaf_id.clone(),
842 schema_version: journal.schema_version,
843 spawn_depth: journal.spawn_depth,
844 seed_from_message_index: 0,
845 workspace: workspace.to_path_buf(),
846 }
847 }
848
849 fn verify_operation_checkpoint(
850 manager: &SessionManager,
851 thread: &crate::runtime_threads::ThreadRecord,
852 binding: &crate::runtime_threads::RuntimeStoreBinding,
853 operation: &RuntimeHistoryOperation,
854 ) -> Result<()> {
855 let id = thread
856 .session_id
857 .as_deref()
858 .ok_or_else(|| anyhow!("canonical operation target lost its session"))?;
859 ensure!(
860 id == operation.receipt.session_id,
861 "canonical operation session changed"
862 );
863 let _lease = manager.reserve_session_for_external_write(id)?;
864 let observed = manager.load_session_snapshot_bounded(id, MAX_CANONICAL_HISTORY_BYTES)?;
865 ensure!(
866 observed.metadata.runtime_store.as_ref() == Some(binding)
867 && observed.metadata.workspace == thread.workspace,
868 "canonical operation document binding changed"
869 );
870 let checkpoint = thread
871 .saved_session_checkpoint
872 .as_ref()
873 .ok_or_else(|| anyhow!("canonical operation checkpoint is absent"))?;
874 ensure!(
875 crate::runtime_threads::checkpoint_prefix_len(checkpoint, &observed.messages)?.is_some(),
876 "canonical operation checkpoint differs from its document"
877 );
878 verify_operation_journal(&observed, operation)
879 }
880
881 fn verify_operation_journal(
882 observed: &SavedSession,
883 operation: &RuntimeHistoryOperation,
884 ) -> Result<()> {
885 let witness = operation
886 .journal_witness
887 .as_ref()
888 .ok_or_else(|| anyhow!("canonical operation witness missing; recovery required"))?;
889 let journal = observed
890 .journal
891 .as_ref()
892 .ok_or_else(|| anyhow!("canonical full graph missing"))?;
893 ensure!(
894 journal.entries.len() >= witness.entries_len
895 && journal.schema_version == witness.schema_version
896 && journal.spawn_depth == witness.spawn_depth,
897 "canonical operation full graph shrank or changed schema"
898 );
899 let original = SessionJournal {
900 entries: journal.entries[..witness.entries_len].to_vec(),
901 leaf_id: witness.leaf_id.clone(),
902 schema_version: witness.schema_version,
903 spawn_depth: witness.spawn_depth,
904 };
905 ensure!(
906 journal_digest(&original)? == operation.receipt.history_digest,
907 "canonical operation branches changed; recovery required"
908 );
909 Ok(())
910 }
911
912 /// Complete only the target whose bytes and process checkpoint were prepared by
913 /// this retained intent. Missing preparation never authorizes source recreation.
914 async fn settle_prepared_history_operation(
915 runtime: &Arc<RuntimeThreadManager>,
916 sessions_dir: &Path,
917 operation: &RuntimeHistoryOperation,
918 workspace: &Path,
919 ) -> Result<Option<CanonicalThreadReceipt>> {
920 let witness = operation
921 .journal_witness
922 .as_ref()
923 .ok_or_else(|| anyhow!("prepared operation lacks its full graph witness"))?;
924 ensure!(
925 witness.workspace == workspace,
926 "operation selected workspace changed"
927 );
928 let captured = Arc::clone(runtime);
929 let id = operation.receipt.runtime_thread_id.clone();
930 let Some(thread) = history_owner_work(move || captured.history_operation_thread(&id)).await?
931 else {
932 return Ok(None);
933 };
934
935 ensure!(
936 thread.workspace == workspace,
937 "reserved target workspace changed"
938 );
939 let detail = runtime.get_thread_detail(&thread.id).await?;
940 ensure!(
941 !super::sessions::thread_detail_has_live_work(&detail),
942 "reserved target has pending work"
943 );
944 let dir = sessions_dir.to_path_buf();
945 let binding = runtime.session_store_binding();
946 let expected = operation.clone();
947 let selected_workspace = workspace.to_path_buf();
948 let saved = history_owner_work(move || {
949 let manager = SessionManager::new(dir)?;
950 let lease = manager.reserve_session_for_external_write(&expected.receipt.session_id)?;
951 let saved = match manager.load_session_snapshot_bounded(
952 &expected.receipt.session_id,
953 MAX_CANONICAL_HISTORY_BYTES,
954 ) {
955 Ok(saved) => saved,
956 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
957 Err(error) => return Err(error.into()),
958 };
959 // An old source document is not the new prepared publication.
960 // The original source CAS below still decides whether it can be retried.
961 if expected.target_document_digest.as_ref() != Some(&saved_document_digest(&saved)?) {
962 return Ok(None);
963 }
964 ensure!(
965 saved.metadata.runtime_store.as_ref() == Some(&binding)
966 && saved.metadata.workspace == selected_workspace,
967 "reserved saved document binding changed"
968 );
969 if let Some(digest) = expected.session_goal_target_digest.as_ref() {
970 ensure!(
971 session_goal_digest(&manager.load_session_goal(&expected.receipt.session_id)?)?
972 == *digest,
973 "prepared local goal sidecar changed; no recovery write admitted"
974 );
975 }
976 verify_operation_journal(&saved, &expected)?;
977 let witness = expected
978 .journal_witness
979 .as_ref()
980 .ok_or_else(|| anyhow!("operation witness absent"))?;
981 let journal = saved
982 .journal
983 .as_ref()
984 .ok_or_else(|| anyhow!("prepared graph is absent"))?;
985 ensure!(
986 journal.entries.len() == witness.entries_len && journal.leaf_id == witness.leaf_id,
987 "uncommitted target acquired a successor; no replay is safe"
988 );
989 Ok(Some((saved, lease)))
990 })
991 .await?;
992 if let Some((saved, _target_lease)) = saved {
993 let from = witness.seed_from_message_index;
994 ensure!(
995 from <= saved.messages.len(),
996 "history suffix boundary is invalid"
997 );
998 let checkpoint = thread
999 .saved_session_checkpoint
1000 .as_ref()
1001 .map(|checkpoint| {
1002 crate::runtime_threads::checkpoint_prefix_len(checkpoint, &saved.messages)
1003 })
1004 .transpose()?
1005 .flatten();
1006 if checkpoint != Some(saved.messages.len()) {
1007 if from == 0 {
1008 ensure!(
1009 detail.turns.is_empty() && detail.items.is_empty(),
1010 "reserved target process history is unproven; recovery required"
1011 );
1012 } else {
1013 ensure!(
1014 checkpoint == Some(from)
1015 && thread.session_id.as_deref() == Some(saved.metadata.id.as_str()),
1016 "reserved Resume prefix changed; recovery required"
1017 );
1018 }
1019 runtime
1020 .seed_thread_from_messages_with_history_operation(
1021 &thread.id,
1022 &saved.messages[from..],
1023 Some(operation),
1024 &saved.messages,
1025 )
1026 .await?;
1027 }
1028 runtime
1029 .set_thread_session_checkpoint(&thread.id, &saved)
1030 .await?;
1031 runtime
1032 .synchronize_history_operation(&thread.id, &saved)
1033 .await?;
1034 let manager = Arc::clone(runtime);
1035 let operation = operation.clone();
1036 return history_owner_work(move || manager.commit_history_operation(operation))
1037 .await
1038 .map(Some);
1039 }
1040 Ok(None)
1041 }
1042
1043 async fn mutate_in_owner(
1044 runtime: Arc<RuntimeThreadManager>,
1045 sessions_dir: std::path::PathBuf,
1046 config_path: Option<std::path::PathBuf>,
1047 config_profile: Option<String>,
1048 request: CanonicalThreadMutationRequest,
1049 ) -> Result<CanonicalThreadReceipt> {
1050 let request_digest = crate::hashing::sha256_hex(
1051 crate::client::canonical_json(&serde_json::to_value(&request)?).as_bytes(),
1052 );
1053 let lookup_runtime = Arc::clone(&runtime);
1054 let key = request.operation_key.clone();
1055 let digest = request_digest.clone();
1056 let reserved =
1057 history_owner_work(move || lookup_runtime.lookup_history_operation(&key, &digest)).await?;
1058 if let Some(operation) = reserved
1059 .as_ref()
1060 .filter(|operation| operation.journal_witness.is_some())
1061 {
1062 let witness = operation.journal_witness.as_ref().ok_or_else(|| {
1063 anyhow!("operation scope witness absent; recover the retained intent")
1064 })?;
1065 ensure!(
1066 witness.workspace == request.workspace,
1067 "operation selected workspace changed"
1068 );
1069 let thread_id = operation.receipt.runtime_thread_id.clone();
1070 let captured = Arc::clone(&runtime);
1071 let thread =
1072 history_owner_work(move || captured.history_operation_thread(&thread_id)).await?;
1073 if operation.committed {
1074 let thread = thread.ok_or_else(|| anyhow!("committed canonical target is missing"))?;
1075 ensure!(
1076 thread.workspace == request.workspace,
1077 "canonical operation workspace changed"
1078 );
1079 let manager_dir = sessions_dir.clone();
1080 let binding = runtime.session_store_binding();
1081 let expected = operation.clone();
1082 history_owner_work(move || {
1083 verify_operation_checkpoint(
1084 &SessionManager::new(manager_dir)?,
1085 &thread,
1086 &binding,
1087 &expected,
1088 )
1089 })
1090 .await?;
1091 return Ok(operation.receipt.clone());
1092 }
1093 if let Some(receipt) = settle_prepared_history_operation(
1094 &runtime,
1095 &sessions_dir,
1096 operation,
1097 &request.workspace,
1098 )
1099 .await?
1100 {
1101 return Ok(receipt);
1102 }
1103 }
1104 let is_resume = matches!(&request.mutation, CanonicalThreadMutation::Resume { .. });
1105 let mut create = CreateThreadRequest::default();
1106 let mut source_thread = None;
1107 let mut source_disk_digest = None;
1108 let mut source_goal = None;
1109 let mut source = match &request.mutation {
1110 CanonicalThreadMutation::Create { config } => {
1111 create = serde_json::from_value(config.clone())
1112 .context("invalid create-thread configuration")?;
1113 // The accepted existing shape has no credentials. Unknown fields
1114 // cannot masquerade as an ignored policy or operation parameter.
1115 let encoded = serde_json::to_value(&create)?;
1116 ensure!(
1117 config.as_object().is_some_and(|fields| fields
1118 .iter()
1119 .all(|(key, value)| encoded.get(key) == Some(value))),
1120 "unsupported create-thread field"
1121 );
1122 ensure!(
1123 create
1124 .workspace
1125 .as_ref()
1126 .is_none_or(|workspace| workspace == &request.workspace),
1127 "create configuration workspace differs from selected admission"
1128 );
1129 None
1130 }
1131 CanonicalThreadMutation::Resume { source, .. }
1132 | CanonicalThreadMutation::Fork { source, .. } => Some(match source {
1133 CanonicalHistorySource::Thread {
1134 runtime_thread_id,
1135 expected_document_digest,
1136 } => {
1137 let thread = runtime.get_thread(runtime_thread_id).await?;
1138 ensure!(
1139 !is_resume || thread.workspace == request.workspace,
1140 "Resume source workspace differs from selected admission"
1141 );
1142 let detail = runtime.get_thread_detail(runtime_thread_id).await?;
1143 ensure!(
1144 !super::sessions::thread_detail_has_live_work(&detail),
1145 "selected history has pending work; wait before resume or fork"
1146 );
1147 let snapshot =
1148 snapshot_in_runtime(&runtime, &sessions_dir, runtime_thread_id.clone())
1149 .await
1150 .map_err(|error| anyhow!(error.message))?;
1151 ensure!(
1152 &snapshot.document_digest == expected_document_digest,
1153 "selected full history changed before canonical admission"
1154 );
1155 source_disk_digest = snapshot.saved_document_digest;
1156 let saved: SavedSession = serde_json::from_value(snapshot.session)?;
1157 let id = saved.metadata.id.clone();
1158 let dir = sessions_dir.clone();
1159 let expected = snapshot.session_goal_digest;
1160 source_goal = history_owner_work(move || {
1161 let manager = SessionManager::new(dir)?;
1162 let _lease = manager.reserve_session_for_external_write(&id)?;
1163 let goal = manager.load_session_goal(&id)?;
1164 ensure!(
1165 session_goal_digest(&goal)? == expected,
1166 "source local goal changed before admission"
1167 );
1168 Ok(goal)
1169 })
1170 .await?;
1171 source_thread = Some(thread);
1172 saved
1173 }
1174 CanonicalHistorySource::SavedSession {
1175 session,
1176 expected_document_digest,
1177 } => {
1178 let proposed: SavedSession = serde_json::from_value(session.clone())?;
1179 ensure!(
1180 serde_json::to_value(&proposed)? == *session,
1181 "saved history contains unsupported fields; source retained"
1182 );
1183 let id = proposed.metadata.id.clone();
1184 let dir = sessions_dir.clone();
1185 let expected = expected_document_digest.clone();
1186 let binding = runtime.session_store_binding();
1187 let selected_workspace = request.workspace.clone();
1188 let (observed, goal) = history_owner_work(move || {
1189 let manager = SessionManager::new(dir)?;
1190 let _lease = manager.reserve_session_for_external_write(&id)?;
1191 let observed =
1192 manager.load_session_snapshot_bounded(&id, MAX_CANONICAL_HISTORY_BYTES)?;
1193 ensure!(
1194 saved_document_digest(&observed)? == expected,
1195 "protected saved document changed before admission"
1196 );
1197 ensure!(
1198 (!is_resume || observed.metadata.workspace == selected_workspace)
1199 && observed
1200 .metadata
1201 .runtime_store
1202 .as_ref()
1203 .is_none_or(|saved| saved == &binding),
1204 "saved history belongs to another workspace or owner store"
1205 );
1206 Ok((observed, manager.load_session_goal(&id)?))
1207 })
1208 .await?;
1209 ensure!(
1210 saved_document_digest(&proposed)? == saved_document_digest(&observed)?,
1211 "proposed history differs from protected saved document"
1212 );
1213 source_goal = goal;
1214 source_disk_digest = Some(expected_document_digest.clone());
1215 let candidates = runtime
1216 .list_threads(ThreadListFilter::IncludeArchived, None)
1217 .await?
1218 .into_iter()
1219 .filter(|thread| {
1220 thread.session_id.as_deref() == Some(observed.metadata.id.as_str())
1221 })
1222 .collect::<Vec<_>>();
1223 ensure!(
1224 candidates.len() <= 1,
1225 "saved history has multiple canonical holders; explicit target required"
1226 );
1227 if let Some(thread) = candidates.into_iter().next() {
1228 ensure!(
1229 thread.workspace == observed.metadata.workspace,
1230 "saved history holder source workspace changed"
1231 );
1232 let detail = runtime.get_thread_detail(&thread.id).await?;
1233 ensure!(
1234 !super::sessions::thread_detail_has_live_work(&detail),
1235 "saved history holder has pending work"
1236 );
1237 let checkpoint = thread
1238 .saved_session_checkpoint
1239 .as_ref()
1240 .ok_or_else(|| anyhow!("saved holder checkpoint is missing"))?;
1241 // The checkpoint records the history the holder was
1242 // seeded with; turns the holder ran since then extend the
1243 // document past it. Like the other checkpoint guards here,
1244 // require it to cover a prefix, and let the exact
1245 // live == document check below prove that everything
1246 // after that prefix came from this holder. Requiring full
1247 // coverage refused every second resume of a continued
1248 // session.
1249 ensure!(
1250 crate::runtime_threads::checkpoint_prefix_len(
1251 checkpoint,
1252 &observed.messages
1253 )?
1254 .is_some(),
1255 "saved holder checkpoint differs from the saved history; refresh history"
1256 );
1257 let full = snapshot_in_runtime(&runtime, &sessions_dir, thread.id.clone())
1258 .await
1259 .map_err(|error| anyhow!(error.message))?;
1260 let live: SavedSession = serde_json::from_value(full.session)?;
1261 if live.messages == observed.messages && live.journal == observed.journal {
1262 source_thread = Some(thread);
1263 } else if live.messages.len() < observed.messages.len()
1264 && observed.messages.starts_with(&live.messages)
1265 {
1266 // After a resume the TUI runs its turns in its own
1267 // engine and saves them to the document; this holder
1268 // stayed at the history it was seeded with. It is
1269 // strictly behind the document and has no live work
1270 // (checked above), so it holds nothing the document
1271 // lacks: release its binding (receipt kept, turns
1272 // kept) and mount the document fresh. Refusing here
1273 // made every continued session impossible to resume
1274 // a second time.
1275 runtime.release_stale_session_holder(
1276 &thread,
1277 "the saved document advanced past this thread after a resume",
1278 )?;
1279 } else {
1280 bail!(
1281 "saved holder has a successor outside the selected document; refresh history"
1282 );
1283 }
1284 }
1285 observed
1286 }
1287 }),
1288 };
1289 let options = match &request.mutation {
1290 CanonicalThreadMutation::Resume { options, .. }
1291 | CanonicalThreadMutation::Fork { options, .. } => options.clone(),
1292 CanonicalThreadMutation::Create { .. } => CanonicalHistoryOptions::default(),
1293 };
1294 if matches!(
1295 &request.mutation,
1296 CanonicalThreadMutation::Resume {
1297 source: CanonicalHistorySource::SavedSession { .. },
1298 ..
1299 } | CanonicalThreadMutation::Fork {
1300 source: CanonicalHistorySource::SavedSession { .. },
1301 ..
1302 }
1303 ) {
1304 ensure!(
1305 options
1306 .expected_session_goal_digest
1307 .as_ref()
1308 .is_some_and(|expected| session_goal_digest(&source_goal)
1309 .is_ok_and(|observed| &observed == expected))
1310 || (options.expected_session_goal_digest.is_none() && source_goal.is_none()),
1311 "saved source local goal witness is absent or changed; source retained"
1312 );
1313 }
1314 let captured_source_goal_digest = session_goal_digest(&source_goal)?;
1315 let target_goal = if is_resume {
1316 source_goal.clone()
1317 } else {
1318 forked_session_goal(source_goal.clone())
1319 };
1320 let prepared_target_goal_digest = session_goal_digest(&target_goal)?;
1321 let base_message_count = source.as_ref().map_or(0, |session| session.messages.len());
1322 if let Some(path) = options.source_path.as_ref() {
1323 let session = source
1324 .as_ref()
1325 .ok_or_else(|| anyhow!("source path requires protected saved history"))?;
1326 ensure!(
1327 *path == sessions_dir.join(format!("{}.json", session.metadata.id)),
1328 "source path does not name the protected selected session; opaque legacy source retained"
1329 );
1330 }
1331 let updates = decode_history_overrides(
1332 &options.overrides,
1333 &runtime.read_config(),
1334 &request.workspace,
1335 )?;
1336 if let Some(session) = source.as_mut() {
1337 append_offered_history(session, &options.offered_history, &request.operation_key)?;
1338 } else {
1339 ensure!(
1340 options.offered_history.is_empty(),
1341 "offered history requires a full source"
1342 );
1343 }
1344 if let Some(session) = source.as_mut() {
1345 let journal = session
1346 .journal
1347 .as_ref()
1348 .ok_or_else(|| anyhow!("saved full journal is absent; source retained"))?;
1349 if let CanonicalThreadMutation::Fork {
1350 selected_entry_id, ..
1351 } = &request.mutation
1352 {
1353 let fork = match journal.fork_from(selected_entry_id.as_deref()) {
1354 Ok(fork) => fork,
1355 Err(error) => {
1356 // An entry a bounded TUI save archived (#6842) is not in
1357 // the document; fail closed and say where it went.
1358 let dir = sessions_dir.clone();
1359 let id = session.metadata.id.clone();
1360 let entry = selected_entry_id.clone();
1361 let archived = history_owner_work(move || {
1362 let Some(entry) = entry else {
1363 return Ok(false);
1364 };
1365 Ok(crate::session_manager::load_journal_archive(&dir, &id)?
1366 .iter()
1367 .any(|archived| archived.id == entry))
1368 })
1369 .await
1370 .unwrap_or(false);
1371 if archived {
1372 bail!(
1373 "{error}: the entry was moved to this session's journal archive to keep \
1374 the saved document bounded; restore it with `/branch {}` first. \
1375 Source retained",
1376 selected_entry_id.as_deref().unwrap_or_default()
1377 );
1378 }
1379 return Err(anyhow::Error::msg(error));
1380 }
1381 };
1382 session.metadata.mark_forked_from(&session.metadata.clone());
1383 session.metadata.spawn_depth = fork.spawn_depth;
1384 session.leaf_id = fork.leaf_id.clone();
1385 session.messages = fork.to_messages();
1386 session.journal = Some(fork);
1387 session.ensure_journal();
1388 }
1389 }
1390 let journal = source
1391 .as_ref()
1392 .and_then(|session| session.journal.clone())
1393 .unwrap_or_else(SessionJournal::new);
1394 ensure!(
1395 journal.entries.len() <= MAX_CANONICAL_HISTORY_ENTRIES,
1396 "complete graph exceeds entry bound"
1397 );
1398 // Import's strong validator is reused without changing ordinary module
1399 // limits or creating a new parser. It rejects duplicate/cycle/unknown graph.
1400 SavedSession::import_foreign(
1401 SessionImportContainer::new("canonical-operation-validation".into(), &journal, None),
1402 request.workspace.clone(),
1403 "validation".into(),
1404 )
1405 .map_err(anyhow::Error::msg)?;
1406 let history_digest = journal_digest(&journal)?;
1407 let mut witness = operation_witness(&journal, &request.workspace);
1408 if is_resume && source_thread.is_some() {
1409 witness.seed_from_message_index = base_message_count;
1410 }
1411 let target = if is_resume {
1412 match source_thread.as_ref() {
1413 Some(thread) => Some((
1414 thread.id.clone(),
1415 thread
1416 .session_id
1417 .clone()
1418 .ok_or_else(|| anyhow!("resume holder has no durable session"))?,
1419 )),
1420 None => source.as_ref().map(|session| {
1421 (
1422 format!("thr_{}", uuid::Uuid::new_v4()),
1423 session.metadata.id.clone(),
1424 )
1425 }),
1426 }
1427 } else {
1428 None
1429 };
1430 let association = codewhale_protocol::CanonicalThreadOperationAssociation {
1431 kind: match &request.mutation {
1432 CanonicalThreadMutation::Create { .. } => {
1433 codewhale_protocol::CanonicalThreadOperationKind::Create
1434 }
1435 CanonicalThreadMutation::Resume { .. } => {
1436 codewhale_protocol::CanonicalThreadOperationKind::Resume
1437 }
1438 CanonicalThreadMutation::Fork { .. } => {
1439 codewhale_protocol::CanonicalThreadOperationKind::Fork
1440 }
1441 },
1442 source_runtime_thread_id: source_thread.as_ref().map(|thread| thread.id.clone()),
1443 source_session_id: source.as_ref().map(|session| session.metadata.id.clone()),
1444 };
1445 let reserve_runtime = Arc::clone(&runtime);
1446 let key = request.operation_key.clone();
1447 let operation = history_owner_work(move || {
1448 let operation = reserve_runtime.reserve_history_operation_for_target(
1449 &key,
1450 &request_digest,
1451 &history_digest,
1452 target
1453 .as_ref()
1454 .map(|(thread, session)| (thread.as_str(), session.as_str())),
1455 )?;
1456 reserve_runtime.bind_history_operation_witness(operation, witness, association)
1457 })
1458 .await?;
1459 create.workspace = Some(request.workspace.clone());
1460 if source_thread.is_none()
1461 && let Some(saved) = source.as_ref()
1462 {
1463 create.model = Some(saved.metadata.model.clone());
1464 create.model_provider = Some(saved.metadata.model_provider.clone());
1465 create.model_provider_id = saved.metadata.model_provider_id.clone();
1466 create.system_prompt = saved.system_prompt.clone();
1467 }
1468 if let Some(thread) = source_thread.as_ref() {
1469 create.model = Some(thread.model.clone());
1470 create.model_provider = thread.model_provider.clone();
1471 create.model_provider_id = thread.model_provider_id.clone();
1472 create.reasoning_effort = thread.reasoning_effort.clone();
1473 create.allowed_tools = thread.allowed_tools.clone();
1474 create.mode = Some(thread.mode.clone());
1475 create.permission_posture = thread.permission_posture.clone();
1476 create.allow_shell = Some(thread.allow_shell);
1477 create.trust_mode = Some(thread.trust_mode);
1478 create.auto_approve = Some(thread.auto_approve);
1479 create.system_prompt = thread.system_prompt.clone();
1480 }
1481 if let Some(updates) = updates.as_ref() {
1482 apply_history_create_overrides(&mut create, updates);
1483 }
1484 let lookup_runtime = Arc::clone(&runtime);
1485 let id = operation.receipt.runtime_thread_id.clone();
1486 let mut thread =
1487 match history_owner_work(move || lookup_runtime.history_operation_thread(&id)).await? {
1488 Some(thread) => {
1489 ensure!(
1490 thread.workspace == request.workspace,
1491 "reserved target workspace changed"
1492 );
1493 thread
1494 }
1495 None => {
1496 runtime
1497 .create_thread_with_reserved_id(
1498 create,
1499 config_path.as_deref(),
1500 config_profile.as_deref(),
1501 Some(operation.receipt.runtime_thread_id.clone()),
1502 )
1503 .await?
1504 }
1505 };
1506 if is_resume && let Some(updates) = updates {
1507 thread = runtime
1508 .update_thread_under_history_guard(
1509 &thread.id,
1510 updates,
1511 config_path.as_deref(),
1512 config_profile.as_deref(),
1513 )
1514 .await?;
1515 }
1516 let source_id = source.as_ref().map(|session| session.metadata.id.clone());
1517 let mut session = source.unwrap_or_else(|| {
1518 crate::session_manager::create_saved_session_with_id_and_mode(
1519 operation.receipt.session_id.clone(),
1520 &[],
1521 &thread.model,
1522 &thread.workspace,
1523 0,
1524 None,
1525 Some(&thread.mode),
1526 )
1527 });
1528 session.metadata.id = operation.receipt.session_id.clone();
1529 if !is_resume {
1530 session.metadata.created_at = operation.created_at;
1531 }
1532 session.metadata.updated_at = operation.created_at;
1533 session.system_prompt.clone_from(&thread.system_prompt);
1534 session.metadata.runtime_store = Some(runtime.session_store_binding());
1535 session.metadata.workspace = thread.workspace.clone();
1536 session.metadata.model = thread.model.clone();
1537 session.metadata.mode = Some(thread.mode.clone());
1538 session.metadata.set_model_provider_route(
1539 thread
1540 .model_provider
1541 .as_deref()
1542 .ok_or_else(|| anyhow!("canonical thread has no provider identity"))?,
1543 thread.model_provider_id.as_deref(),
1544 );
1545 let document_digest = saved_document_digest(&session)?;
1546 let prepare_runtime = Arc::clone(&runtime);
1547 let operation = history_owner_work(move || {
1548 let operation =
1549 prepare_runtime.bind_history_operation_document(operation, &document_digest)?;
1550 prepare_runtime.bind_history_operation_session_goal(
1551 operation,
1552 &captured_source_goal_digest,
1553 &prepared_target_goal_digest,
1554 )
1555 })
1556 .await?;
1557 // Revalidate the exact source file immediately before target publication.
1558 // Both existing session leases are held until the target write settles.
1559 let dir = sessions_dir.clone();
1560 let expected_journal = journal.clone();
1561 let source_goal_digest = session_goal_digest(&source_goal)?;
1562 let target_goal_digest = session_goal_digest(&target_goal)?;
1563 let target_id = session.metadata.id.clone();
1564 let (session, _source_lease, _target_lease) = history_owner_work(move || {
1565 let manager = SessionManager::new(dir)?;
1566 let source_lease = source_id
1567 .as_ref()
1568 .filter(|id| **id != target_id)
1569 .map(|id| manager.reserve_session_for_external_write(id))
1570 .transpose()?;
1571 let target_lease = manager.reserve_session_for_external_write(&target_id)?;
1572 if let (Some(id), Some(digest)) = (source_id.as_ref(), source_disk_digest.as_ref()) {
1573 let observed =
1574 manager.load_session_snapshot_bounded(id, MAX_CANONICAL_HISTORY_BYTES)?;
1575 ensure!(
1576 session_goal_digest(&manager.load_session_goal(id)?)? == source_goal_digest,
1577 "source local goal changed before publication"
1578 );
1579 ensure!(
1580 saved_document_digest(&observed)? == *digest,
1581 "source full document changed before publication"
1582 );
1583 }
1584 if source_id.as_ref() != Some(&target_id) {
1585 let previous_goal = manager.load_session_goal(&target_id)?;
1586 ensure!(
1587 previous_goal.is_none()
1588 || session_goal_digest(&previous_goal)? == target_goal_digest,
1589 "reserved target has a conflicting local goal sidecar; source retained"
1590 );
1591 }
1592 match manager.load_session_snapshot_bounded(&target_id, MAX_CANONICAL_HISTORY_BYTES) {
1593 Ok(observed) if source_id.as_ref() != Some(&target_id) => {
1594 ensure!(
1595 observed.journal.as_ref() == Some(&expected_journal)
1596 && observed.metadata.runtime_store == session.metadata.runtime_store,
1597 "reserved canonical target changed; recovery required"
1598 );
1599 ensure!(
1600 session_goal_digest(&manager.load_session_goal(&target_id)?)?
1601 == target_goal_digest,
1602 "reserved target local goal changed; recovery required"
1603 );
1604 Ok((observed, source_lease, target_lease))
1605 }
1606 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1607 manager.save_session(&session)?;
1608 manager.save_session_goal(&target_id, target_goal.as_ref())?;
1609 Ok((session, source_lease, target_lease))
1610 }
1611 Ok(_) => {
1612 manager.save_session(&session)?;
1613 manager.save_session_goal(&target_id, target_goal.as_ref())?;
1614 Ok((session, source_lease, target_lease))
1615 }
1616 Err(error) => Err(error.into()),
1617 }
1618 })
1619 .await?;
1620 let detail = runtime.get_thread_detail(&thread.id).await?;
1621 ensure!(
1622 !super::sessions::thread_detail_has_live_work(&detail),
1623 "canonical target has pending work"
1624 );
1625 if detail.turns.is_empty() && detail.items.is_empty() {
1626 runtime
1627 .seed_thread_from_messages_with_history_operation(
1628 &thread.id,
1629 &session.messages,
1630 Some(&operation),
1631 &session.messages,
1632 )
1633 .await?;
1634 } else {
1635 ensure!(
1636 is_resume,
1637 "uncommitted new target has unproven history; recovery required"
1638 );
1639 let from = operation
1640 .journal_witness
1641 .as_ref()
1642 .ok_or_else(|| anyhow!("history witness absent"))?
1643 .seed_from_message_index;
1644 ensure!(
1645 from <= session.messages.len(),
1646 "history suffix boundary is invalid"
1647 );
1648 if from < session.messages.len() {
1649 runtime
1650 .seed_thread_from_messages_with_history_operation(
1651 &thread.id,
1652 &session.messages[from..],
1653 Some(&operation),
1654 &session.messages,
1655 )
1656 .await?;
1657 }
1658 }
1659 runtime
1660 .set_thread_session_checkpoint(&thread.id, &session)
1661 .await?;
1662 runtime
1663 .synchronize_history_operation(&thread.id, &session)
1664 .await?;
1665 let commit_runtime = Arc::clone(&runtime);
1666 history_owner_work(move || commit_runtime.commit_history_operation(operation)).await
1667 }
1668
1669 pub(super) async fn lookup_thread_history_operation(
1670 State(state): State<RuntimeApiState>,
1671 Json(request): Json<codewhale_protocol::CanonicalThreadOperationLookup>,
1672 ) -> Result<Json<codewhale_protocol::CanonicalThreadOperationStatus>, ApiError> {
1673 lookup_thread_history_operation_in_runtime(&state.runtime_threads, &state.sessions_dir, request)
1674 .await
1675 .map(Json)
1676 .map_err(|error| ApiError::conflict(error.to_string()))
1677 }
1678
1679 /// A read-only exact-key lookup never resumes an Engine, mutates a thread, or
1680 /// guesses absence from a damaged intent/document. It authenticates current
1681 /// binding separately from the historical receipt.
1682 pub(crate) async fn lookup_thread_history_operation_in_runtime(
1683 runtime: &Arc<RuntimeThreadManager>,
1684 sessions_dir: &Path,
1685 request: codewhale_protocol::CanonicalThreadOperationLookup,
1686 ) -> Result<codewhale_protocol::CanonicalThreadOperationStatus> {
1687 let binding = runtime.session_store_binding();
1688 ensure!(
1689 request.version == 1
1690 && request.expected_data_dir == binding.data_dir
1691 && request.expected_execution_scope == binding.execution_scope,
1692 "operation lookup belongs to a different held owner"
1693 );
1694 let manager = Arc::clone(runtime);
1695 let key = request.operation_key.clone();
1696 let operation =
1697 history_owner_work(move || manager.lookup_history_operation_by_key(&key)).await?;
1698 let Some(operation) = operation else {
1699 return Ok(codewhale_protocol::CanonicalThreadOperationStatus::Absent);
1700 };
1701 ensure!(
1702 operation
1703 .journal_witness
1704 .as_ref()
1705 .is_some_and(|witness| witness.workspace == request.workspace),
1706 "operation belongs to another acknowledged workspace"
1707 );
1708 let association = operation.association.clone().ok_or_else(|| {
1709 anyhow!("legacy operation lacks a typed action association; recovery required")
1710 })?;
1711 if !operation.committed {
1712 return Ok(
1713 codewhale_protocol::CanonicalThreadOperationStatus::Pending {
1714 receipt: operation.receipt,
1715 association,
1716 },
1717 );
1718 }
1719 let admission = runtime.try_history_import_guard()?;
1720 let manager = Arc::clone(runtime);
1721 let dir = sessions_dir.to_path_buf();
1722 history_owner_work(move || {
1723 let _admission = admission;
1724 let thread = manager
1725 .history_operation_thread(&operation.receipt.runtime_thread_id)?
1726 .ok_or_else(|| anyhow!("committed operation target is absent; recovery required"))?;
1727 ensure!(
1728 thread.workspace == request.workspace,
1729 "committed target belongs to another selected workspace"
1730 );
1731 verify_operation_checkpoint(&SessionManager::new(dir)?, &thread, &binding, &operation)?;
1732 Ok(
1733 codewhale_protocol::CanonicalThreadOperationStatus::Committed {
1734 receipt: operation.receipt,
1735 association,
1736 },
1737 )
1738 })
1739 .await
1740 }
1741
1742 pub(super) async fn recover_thread_history_operation(
1743 State(state): State<RuntimeApiState>,
1744 Json(request): Json<codewhale_protocol::CanonicalThreadOperationRecovery>,
1745 ) -> Result<Json<codewhale_protocol::CanonicalThreadOperationStatus>, ApiError> {
1746 recover_thread_history_operation_in_runtime(
1747 &state.runtime_threads,
1748 &state.sessions_dir,
1749 request,
1750 )
1751 .await
1752 .map(Json)
1753 .map_err(|error| ApiError::conflict(error.to_string()))
1754 }
1755
1756 /// Explicit exact-key recovery settles already-published preparation; it never
1757 /// reconstructs an intent from a changed source. Lookup remains read-only.
1758 pub(crate) async fn recover_thread_history_operation_in_runtime(
1759 runtime: &Arc<RuntimeThreadManager>,
1760 sessions_dir: &Path,
1761 request: codewhale_protocol::CanonicalThreadOperationRecovery,
1762 ) -> Result<codewhale_protocol::CanonicalThreadOperationStatus> {
1763 let binding = runtime.session_store_binding();
1764 ensure!(
1765 request.operation.version == 1
1766 && request.operation.expected_data_dir == binding.data_dir
1767 && request.operation.expected_execution_scope == binding.execution_scope,
1768 "operation recovery belongs to a different held owner"
1769 );
1770 let admission = runtime.try_history_import_guard()?;
1771 let runtime = Arc::clone(runtime);
1772 let dir = sessions_dir.to_path_buf();
1773 tokio::spawn(async move {
1774 let _admission = admission;
1775 let captured = Arc::clone(&runtime);
1776 let key = request.operation.operation_key.clone();
1777 let Some(operation) =
1778 history_owner_work(move || captured.lookup_history_operation_by_key(&key)).await?
1779 else {
1780 return Ok(codewhale_protocol::CanonicalThreadOperationStatus::Absent);
1781 };
1782 ensure!(
1783 operation.association.as_ref() == Some(&request.association),
1784 "retained operation action/source association differs; no recovery write admitted"
1785 );
1786 ensure!(
1787 operation
1788 .journal_witness
1789 .as_ref()
1790 .is_some_and(|witness| witness.workspace == request.operation.workspace),
1791 "operation belongs to another acknowledged workspace"
1792 );
1793 let committed = if operation.committed {
1794 let captured = Arc::clone(&runtime);
1795 let expected = operation.clone();
1796 let workspace = request.operation.workspace.clone();
1797 history_owner_work(move || {
1798 let thread = captured
1799 .history_operation_thread(&expected.receipt.runtime_thread_id)?
1800 .ok_or_else(|| anyhow!("committed operation target is absent"))?;
1801 ensure!(
1802 thread.workspace == workspace,
1803 "operation target workspace changed"
1804 );
1805 verify_operation_checkpoint(
1806 &SessionManager::new(dir)?,
1807 &thread,
1808 &binding,
1809 &expected,
1810 )?;
1811 Ok(expected.receipt)
1812 })
1813 .await?
1814 } else {
1815 let Some(receipt) = settle_prepared_history_operation(
1816 &runtime,
1817 &dir,
1818 &operation,
1819 &request.operation.workspace,
1820 )
1821 .await?
1822 else {
1823 return Ok(
1824 codewhale_protocol::CanonicalThreadOperationStatus::Pending {
1825 receipt: operation.receipt,
1826 association: request.association,
1827 },
1828 );
1829 };
1830 receipt
1831 };
1832 Ok(
1833 codewhale_protocol::CanonicalThreadOperationStatus::Committed {
1834 receipt: committed,
1835 association: request.association,
1836 },
1837 )
1838 })
1839 .await
1840 .context("operation recovery outcome is uncertain; retain the same operation key")?
1841 }
1842
1843 fn append_offered_history(
1844 session: &mut SavedSession,
1845 values: &[serde_json::Value],
1846 key: &str,
1847 ) -> Result<()> {
1848 let offered: Vec<Message> = values
1849 .iter()
1850 .map(|value| {
1851 let message: Message = serde_json::from_value(value.clone()).context(
1852 "opaque offered history retained; only complete model messages are admitted",
1853 )?;
1854 ensure!(
1855 serde_json::to_value(&message)? == *value,
1856 "offered message has unsupported fields"
1857 );
1858 if message.role == Role::User {
1859 crate::image_attach::runtime_images_from_blocks(&message.content)
1860 .map_err(anyhow::Error::msg)?;
1861 }
1862 Ok(message)
1863 })
1864 .collect::<Result<_>>()?;
1865 let overlap = codewhale_core::persisted_overlap(&session.messages, &offered);
1866 let journal = session
1867 .journal
1868 .as_mut()
1869 .ok_or_else(|| anyhow!("full journal missing"))?;
1870 ensure!(
1871 journal
1872 .entries
1873 .len()
1874 .checked_add(offered.len() - overlap)
1875 .is_some_and(|count| count <= MAX_CANONICAL_HISTORY_ENTRIES),
1876 "complete offered history exceeds entry bound"
1877 );
1878 for (index, message) in offered.into_iter().enumerate().skip(overlap) {
1879 journal.append_stamped(
1880 SessionEntryKind::Message { message },
1881 session.metadata.updated_at,
1882 );
1883 let entry = journal
1884 .entries
1885 .last_mut()
1886 .ok_or_else(|| anyhow!("offered journal append missing"))?;
1887 entry.id = format!(
1888 "offered:{}:{index}",
1889 crate::hashing::sha256_hex(key.as_bytes())
1890 );
1891 journal.leaf_id = Some(entry.id.clone());
1892 }
1893 session.leaf_id = journal.leaf_id.clone();
1894 // Publish the same active branch/count that the protected reader restores,
1895 // before retaining the full-document recovery witness.
1896 session.ensure_journal();
1897 Ok(())
1898 }
1899
1900 fn decode_history_overrides(
1901 value: &serde_json::Value,
1902 config: &crate::config::Config,
1903 workspace: &Path,
1904 ) -> Result<Option<crate::runtime_threads::UpdateThreadRequest>> {
1905 if value.is_null() {
1906 return Ok(None);
1907 }
1908 let fields = value
1909 .as_object()
1910 .ok_or_else(|| anyhow!("history overrides must be an object"))?;
1911 let mut wire = serde_json::Map::new();
1912 let mut prompt_parts = Vec::new();
1913 for (key, value) in fields {
1914 match key.as_str() {
1915 "model" | "model_provider" | "model_provider_id" | "mode" | "permission_posture"
1916 | "allow_shell" | "trust_mode" | "auto_approve" | "system_prompt" => {
1917 ensure!(
1918 wire.insert(key.clone(), value.clone()).is_none(),
1919 "conflicting history configuration fields"
1920 );
1921 }
1922 "cwd" | "workspace" => {
1923 let selected: std::path::PathBuf = serde_json::from_value(value.clone())?;
1924 ensure!(
1925 selected == workspace,
1926 "history workspace override differs from selected owner scope"
1927 );
1928 }
1929 "approval_policy" => {
1930 let policy = value
1931 .as_str()
1932 .ok_or_else(|| anyhow!("approval policy must be a string"))?;
1933 let posture = match policy.trim().to_ascii_lowercase().as_str() {
1934 "ask" | "suggest" | "on-request" | "untrusted" => "ask",
1935 "auto" | "auto-review" | "auto_review" => "auto_review",
1936 "full" | "full-access" | "full_access" | "bypass" | "never" => "full_access",
1937 _ => {
1938 return Err(anyhow!(
1939 "unsupported legacy approval policy; owner policy retained"
1940 ));
1941 }
1942 };
1943 ensure!(
1944 wire.insert("permission_posture".into(), serde_json::json!(posture))
1945 .is_none(),
1946 "conflicting history approval policy fields"
1947 );
1948 }
1949 "sandbox" => ensure!(
1950 value.as_str() == config.sandbox_mode.as_deref(),
1951 "legacy sandbox override conflicts with the captured owner sandbox ceiling"
1952 ),
1953 "config" => {
1954 let nested = value
1955 .as_object()
1956 .ok_or_else(|| anyhow!("history config must be a typed override object"))?;
1957 for (name, setting) in nested {
1958 ensure!(
1959 matches!(
1960 name.as_str(),
1961 "model"
1962 | "model_provider"
1963 | "model_provider_id"
1964 | "mode"
1965 | "permission_posture"
1966 | "allow_shell"
1967 | "trust_mode"
1968 | "auto_approve"
1969 | "system_prompt"
1970 ),
1971 "unsupported history config field; no credentials or global policy are imported"
1972 );
1973 ensure!(
1974 wire.insert(name.clone(), setting.clone()).is_none(),
1975 "conflicting history configuration fields"
1976 );
1977 }
1978 }
1979 "base_instructions" | "developer_instructions" | "personality" => {
1980 let text = value
1981 .as_str()
1982 .ok_or_else(|| anyhow!("history instruction must be a string"))?;
1983 prompt_parts.push(format!("{key}:\n{text}"));
1984 }
1985 _ => {
1986 return Err(anyhow!(
1987 "unsupported history override {key}; source retained"
1988 ));
1989 }
1990 }
1991 }
1992 if !prompt_parts.is_empty() {
1993 ensure!(
1994 !wire.contains_key("system_prompt"),
1995 "conflicting history instruction fields"
1996 );
1997 wire.insert(
1998 "system_prompt".into(),
1999 serde_json::json!(prompt_parts.join("\n\n")),
2000 );
2001 }
2002 if wire.is_empty() {
2003 return Ok(None);
2004 }
2005 Ok(Some(serde_json::from_value(serde_json::Value::Object(
2006 wire,
2007 ))?))
2008 }
2009
2010 fn apply_history_create_overrides(
2011 create: &mut CreateThreadRequest,
2012 update: &crate::runtime_threads::UpdateThreadRequest,
2013 ) {
2014 if update.model.is_some() {
2015 create.model.clone_from(&update.model);
2016 }
2017 if update.model_provider.is_some() {
2018 create.model_provider.clone_from(&update.model_provider);
2019 }
2020 if update.model_provider_id.is_some() {
2021 create
2022 .model_provider_id
2023 .clone_from(&update.model_provider_id);
2024 }
2025 if update.mode.is_some() {
2026 create.mode.clone_from(&update.mode);
2027 }
2028 if update.permission_posture.is_some() {
2029 create
2030 .permission_posture
2031 .clone_from(&update.permission_posture);
2032 }
2033 if update.allow_shell.is_some() {
2034 create.allow_shell = update.allow_shell;
2035 }
2036 if update.trust_mode.is_some() {
2037 create.trust_mode = update.trust_mode;
2038 }
2039 if update.auto_approve.is_some() {
2040 create.auto_approve = update.auto_approve;
2041 }
2042 if update.system_prompt.is_some() {
2043 create.system_prompt.clone_from(&update.system_prompt);
2044 }
2045 }
2046
2047 pub(crate) struct HistorySizeBound(pub(crate) usize);
2048 impl std::io::Write for HistorySizeBound {
2049 fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
2050 if bytes.len() > self.0 {
2051 return Err(std::io::Error::new(
2052 std::io::ErrorKind::InvalidData,
2053 "full history exceeds byte transport bound; source retained",
2054 ));
2055 }
2056 self.0 -= bytes.len();
2057 Ok(bytes.len())
2058 }
2059 fn flush(&mut self) -> std::io::Result<()> {
2060 Ok(())
2061 }
2062 }
2063
2064 #[cfg(test)]
2065 mod graph_tests {
2066 use super::*;
2067 use codewhale_protocol::MessageRecord;
2068
2069 fn history() -> LegacyThreadHistory {
2070 LegacyThreadHistory {
2071 goal: None,
2072 version: 1,
2073 state_store_id: "a".repeat(64),
2074 thread_id: "legacy".into(),
2075 current_leaf_id: Some(3),
2076 messages: [
2077 (1, None, "system"),
2078 (2, Some(1), "user"),
2079 (3, Some(2), "assistant"),
2080 (4, Some(1), "user"),
2081 ]
2082 .into_iter()
2083 .map(|(id, parent_entry_id, role)| MessageRecord {
2084 id,
2085 thread_id: "legacy".into(),
2086 role: role.into(),
2087 content: format!("content-{id}"),
2088 item: None,
2089 created_at: 1_700_000_000 + id,
2090 parent_entry_id,
2091 })
2092 .collect(),
2093 }
2094 }
2095
2096 #[test]
2097 fn history_override_conflicts_refuse_every_insert_and_preserve_full_access_posture() {
2098 let config = crate::config::Config::default();
2099 let workspace = std::path::PathBuf::from("/checked-workspace");
2100 for value in [
2101 serde_json::json!({"config":{"model":"first"},"model":"second"}),
2102 serde_json::json!({"approval_policy":"on-request","permission_posture":"full_access"}),
2103 serde_json::json!({"config":{"system_prompt":"first"},"system_prompt":"second"}),
2104 serde_json::json!({"base_instructions":"first","system_prompt":"second"}),
2105 ] {
2106 assert!(decode_history_overrides(&value, &config, &workspace).is_err());
2107 }
2108 for policy in ["full-access", "never"] {
2109 let update = decode_history_overrides(
2110 &serde_json::json!({"approval_policy":policy}),
2111 &config,
2112 &workspace,
2113 )
2114 .unwrap()
2115 .unwrap();
2116 assert_eq!(update.permission_posture.as_deref(), Some("full_access"));
2117 }
2118 assert!(
2119 decode_history_overrides(
2120 &serde_json::json!({"approval_policy":"never","permission_posture":"ask"}),
2121 &config,
2122 &workspace,
2123 )
2124 .is_err()
2125 );
2126 }
2127
2128 #[test]
2129 fn full_history_measurement_refuses_before_allocating_a_truncated_value() {
2130 let mut size = HistorySizeBound(32);
2131 assert!(serde_json::to_writer(&mut size, &"x".repeat(64)).is_err());
2132 assert!(size.0 <= 32);
2133 let mut size = HistorySizeBound(4);
2134 serde_json::to_writer(&mut size, &"ok").unwrap();
2135 assert_eq!(size.0, 0);
2136 }
2137
2138 #[test]
2139 fn import_graph_preserves_all_branches_stamps_and_selected_system_prefix() {
2140 let source = history();
2141 let before = serde_json::to_value(&source).unwrap();
2142 let graph = legacy_journal(&source).unwrap();
2143 assert_eq!(graph.entries.len(), 4);
2144 assert_eq!(graph.to_messages().len(), 3);
2145 assert_eq!(graph.to_messages()[0].role, Role::System);
2146 assert_eq!(graph.entries[3].parent_id, graph.entries[1].parent_id);
2147 assert_eq!(graph.entries[3].created_at.timestamp(), 1_700_000_004);
2148 assert_eq!(before, serde_json::to_value(source).unwrap());
2149 }
2150
2151 #[test]
2152 fn import_graph_refuses_cycles_duplicates_dangling_foreign_and_opaque_rows() {
2153 let original = history();
2154 let mut cases = Vec::new();
2155 let mut graph = original.clone();
2156 graph.messages[0].parent_entry_id = Some(3);
2157 cases.push(graph);
2158 let mut graph = original.clone();
2159 graph.messages[3].id = 2;
2160 cases.push(graph);
2161 let mut graph = original.clone();
2162 graph.messages[3].parent_entry_id = Some(99);
2163 cases.push(graph);
2164 let mut graph = original.clone();
2165 graph.messages[3].thread_id = "foreign".into();
2166 cases.push(graph);
2167 let mut graph = original.clone();
2168 graph.messages[3].item = Some(serde_json::json!({"unknown":"opaque"}));
2169 cases.push(graph);
2170 let mut graph = original.clone();
2171 graph.current_leaf_id = None;
2172 cases.push(graph);
2173 let mut graph = original.clone();
2174 graph.version = 2;
2175 cases.push(graph);
2176 for source in cases {
2177 let before = serde_json::to_value(&source).unwrap();
2178 assert!(legacy_journal(&source).is_err());
2179 assert_eq!(
2180 before,
2181 serde_json::to_value(source).unwrap(),
2182 "refusal preserves source"
2183 );
2184 }
2185 }
2186 }
2187
2187 lines RUST