返回 CodeWhale
sessions.rs
根目录 / crates / tui / src / runtime_api / sessions.rs
1 use std::collections::HashMap;
2 use std::path::PathBuf;
3
4 use axum::Json;
5 use axum::extract::{Path, Query, State};
6 use axum::http::StatusCode;
7 use serde::{Deserialize, Serialize};
8 use serde_json::{Value, json};
9
10 use crate::runtime_threads::{
11 CreateThreadRequest, RuntimeThreadManager, RuntimeTurnStatus, ThreadDetail, ThreadListFilter,
12 TurnItemLifecycleStatus,
13 };
14 use crate::session_manager::{
15 SavedSession, SessionListFilter, SessionManager, SessionMetadata, SessionMutator,
16 create_saved_session_with_id_and_mode,
17 };
18 use crate::session_peek::{MAX_PEEK_ENTRIES, SessionPeek, build_peek};
19 use crate::session_projection::{SessionQuery, SessionSortMode, SessionSummary, project_sessions};
20 use codewhale_models::{Message, Role};
21
22 use super::{ApiError, RuntimeApiState, map_thread_err, truncate_text};
23
24 #[derive(Debug, Serialize)]
25 pub(super) struct SessionsResponse {
26 sessions: Vec<SessionMetadata>,
27 }
28
29 #[derive(Debug, Serialize)]
30 pub(super) struct SessionDetailResponse {
31 pub(super) metadata: SessionMetadata,
32 pub(super) messages: Vec<Value>,
33 pub(super) system_prompt: Option<String>,
34 /// Turns that ended `Failed`, with the redacted reason the transcript
35 /// showed. Absent when none did.
36 #[serde(skip_serializing_if = "Vec::is_empty")]
37 pub(super) turn_outcomes: Vec<crate::session_manager::SavedTurnOutcome>,
38 }
39
40 #[derive(Debug, Deserialize)]
41 pub(super) struct CreateSessionRequest {
42 thread_id: String,
43 title: Option<String>,
44 }
45
46 #[derive(Debug, Serialize)]
47 pub(super) struct CreateSessionResponse {
48 session_id: String,
49 thread_id: String,
50 message_count: usize,
51 title: String,
52 }
53
54 #[derive(Debug, Deserialize)]
55 pub(crate) struct ResumeSessionRequest {
56 pub(crate) model: Option<String>,
57 pub(crate) mode: Option<String>,
58 }
59
60 #[derive(Debug, Serialize)]
61 pub(crate) struct ResumeSessionResponse {
62 pub(crate) thread_id: String,
63 pub(crate) session_id: String,
64 pub(crate) message_count: usize,
65 pub(crate) summary: String,
66 }
67
68 #[derive(Debug, Deserialize)]
69 pub(super) struct SessionsQuery {
70 limit: Option<usize>,
71 search: Option<String>,
72 /// Include archived sessions. Same name and meaning as the `/v1/threads`
73 /// query pair, so a client does not need two mental models (#4397).
74 #[serde(default)]
75 include_archived: Option<bool>,
76 /// Return archived sessions only. Overrides `include_archived`.
77 #[serde(default)]
78 archived_only: Option<bool>,
79 /// Restrict to sessions recorded against this workspace. Absent means
80 /// every workspace, matching the historical behaviour of this route.
81 #[serde(default)]
82 workspace: Option<PathBuf>,
83 /// `recent` (default), `name`, or `size`.
84 #[serde(default)]
85 sort: Option<String>,
86 }
87
88 /// `PATCH /v1/sessions/{id}` body. Both fields are optional; omitting one
89 /// leaves it untouched.
90 #[derive(Debug, Deserialize)]
91 pub(super) struct PatchSessionRequest {
92 #[serde(default)]
93 title: Option<String>,
94 #[serde(default)]
95 archived: Option<bool>,
96 }
97
98 /// Lifecycle receipt for a session mutation.
99 ///
100 /// Deliberately shaped like the thread patch receipt: the caller gets the
101 /// resulting record plus an explicit `changes` map of what actually moved, so
102 /// a no-op patch is distinguishable from an applied one without diffing.
103 #[derive(Debug, Serialize)]
104 pub(super) struct PatchSessionResponse {
105 session: SessionMetadata,
106 changes: HashMap<String, Value>,
107 }
108
109 #[derive(Debug, Deserialize)]
110 pub(crate) struct SaveSessionRequest {
111 /// Thread ID to save as a session. If omitted, saves the most recently
112 /// active thread.
113 #[serde(default)]
114 pub(crate) thread_id: Option<String>,
115 /// If provided, update the existing session with this ID instead of
116 /// creating a new one. This matches TUI's `build_session_snapshot`
117 /// behavior where it updates the current session in-place.
118 #[serde(default)]
119 pub(crate) session_id: Option<String>,
120 }
121
122 #[derive(Debug, Serialize)]
123 pub(crate) struct SaveSessionResponse {
124 pub(crate) session_id: String,
125 pub(super) session: SessionDetailResponse,
126 }
127
128 /// Turn a `SessionsQuery` into the shared projection query.
129 ///
130 /// The whole point of routing through [`SessionQuery`] is that the API's
131 /// filter/sort/search semantics are the *same code* the TUI picker and the
132 /// sidebar rail run, not a parallel reimplementation that drifts.
133 fn projection_query(query: &SessionsQuery) -> SessionQuery {
134 let mut projected = SessionQuery::default()
135 .with_filter(SessionListFilter::from_query(
136 query.include_archived,
137 query.archived_only,
138 ))
139 .with_sort(
140 query
141 .sort
142 .as_deref()
143 .map_or(SessionSortMode::Recent, SessionSortMode::from_str_or_recent),
144 )
145 .with_search(query.search.clone().unwrap_or_default())
146 .with_limit(query.limit.unwrap_or(50).clamp(1, 500));
147 if let Some(workspace) = query.workspace.as_deref() {
148 projected = projected.scoped_to(workspace);
149 }
150 projected
151 }
152
153 pub(super) async fn list_sessions(
154 State(state): State<RuntimeApiState>,
155 Query(query): Query<SessionsQuery>,
156 ) -> Result<Json<SessionsResponse>, ApiError> {
157 let manager = SessionManager::new(state.sessions_dir.clone())
158 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
159 let all = manager
160 .list_sessions()
161 .map_err(|e| ApiError::internal(format!("Failed to list sessions: {e}")))?;
162 // This route keeps returning full `SessionMetadata` for compatibility;
163 // `/v1/sessions/summary` is the projected shape. Membership *and* order
164 // come from the shared projection so the two routes never disagree.
165 let sessions: Vec<SessionMetadata> = project_sessions(&all, &projection_query(&query), None)
166 .into_iter()
167 .filter_map(|summary| all.iter().find(|m| m.id == summary.id).cloned())
168 .collect();
169 Ok(Json(SessionsResponse { sessions }))
170 }
171
172 /// `GET /v1/sessions/summary` — the projected row shape.
173 ///
174 /// Field-compatible with `/v1/threads/summary` so the embedded dashboard can
175 /// render a saved session and a live thread with one row renderer, which is
176 /// what "one projection" means in practice rather than as an aspiration.
177 pub(super) async fn list_sessions_summary(
178 State(state): State<RuntimeApiState>,
179 Query(query): Query<SessionsQuery>,
180 ) -> Result<Json<Vec<SessionSummary>>, ApiError> {
181 let manager = SessionManager::new(state.sessions_dir.clone())
182 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
183 let all = manager
184 .list_sessions()
185 .map_err(|e| ApiError::internal(format!("Failed to list sessions: {e}")))?;
186 Ok(Json(project_sessions(
187 &all,
188 &projection_query(&query),
189 None,
190 )))
191 }
192
193 /// `PATCH /v1/sessions/{id}` — rename and/or archive a saved session.
194 ///
195 /// Both mutations go through the manager's single writers
196 /// (`rename_session`, `set_session_archived`), which is what keeps the web
197 /// dashboard, the TUI picker, and `/sessions archive` from producing three
198 /// different notions of the same lifecycle state.
199 pub(super) async fn patch_session(
200 State(state): State<RuntimeApiState>,
201 Path(id): Path<String>,
202 Json(req): Json<PatchSessionRequest>,
203 ) -> Result<Json<PatchSessionResponse>, ApiError> {
204 if req.title.is_none() && req.archived.is_none() {
205 return Err(ApiError::bad_request(
206 "PATCH /v1/sessions/{id} requires at least one of `title` or `archived`",
207 ));
208 }
209 // The single writers hold the session's live lease across their load and
210 // save; taking it may retry with short sleeps. Keep all of it off the
211 // async worker (#6149).
212 tokio::task::spawn_blocking(move || patch_session_blocking(state.sessions_dir, &id, &req))
213 .await
214 .map_err(|_| ApiError::internal("session update failed"))?
215 .map(Json)
216 }
217
218 fn patch_session_blocking(
219 sessions_dir: PathBuf,
220 id: &str,
221 req: &PatchSessionRequest,
222 ) -> Result<PatchSessionResponse, ApiError> {
223 let manager = SessionManager::new(sessions_dir)
224 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
225
226 let before = manager
227 .load_session(id)
228 .map_err(|e| map_session_err(id, e, "read"))?
229 .metadata;
230 let mut metadata = before.clone();
231 let mut changes: HashMap<String, Value> = HashMap::new();
232
233 if let Some(title) = req.title.as_deref() {
234 // Validate the title before touching the store so a rejected title
235 // reports *why* it was rejected rather than the generic "invalid
236 // session id" that `map_session_err` produces for `InvalidInput`.
237 crate::session_manager::normalize_session_title(title)
238 .map_err(|e| ApiError::bad_request(e.to_string()))?;
239 metadata = manager
240 .rename_session(id, title, SessionMutator::External)
241 .map_err(|e| map_session_err(id, e, "rename"))?;
242 if metadata.title != before.title {
243 changes.insert("title".to_string(), json!(metadata.title));
244 }
245 }
246 if let Some(archived) = req.archived {
247 metadata = manager
248 .set_session_archived(id, archived, SessionMutator::External)
249 .map_err(|e| map_session_err(id, e, "archive"))?;
250 if metadata.archived != before.archived {
251 changes.insert("archived".to_string(), json!(metadata.archived));
252 }
253 }
254
255 Ok(PatchSessionResponse {
256 session: metadata,
257 changes,
258 })
259 }
260
261 /// Hold `id`'s live lease across an external load and save, refusing (409)
262 /// as rename, archive and delete do when an interactive session holds the
263 /// document open — its next autosave would revert the write — and rejecting
264 /// a malformed id (400). A released liveness probe cannot protect the write
265 /// that follows it (#6144). Taking the lease may retry with short sleeps, so
266 /// it runs off the async worker; drop the lease only after the save.
267 async fn reserve_external_session_write(
268 state: &RuntimeApiState,
269 id: &str,
270 action: &'static str,
271 ) -> Result<crate::session_manager::SessionLease, ApiError> {
272 reserve_session_write(&state.sessions_dir, id, action).await
273 }
274
275 async fn reserve_session_write(
276 sessions_dir: &std::path::Path,
277 id: &str,
278 action: &'static str,
279 ) -> Result<crate::session_manager::SessionLease, ApiError> {
280 let sessions_dir = sessions_dir.to_path_buf();
281 let id = id.to_string();
282 tokio::task::spawn_blocking(move || {
283 SessionManager::new(sessions_dir)
284 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?
285 .reserve_session_for_external_write(&id)
286 .map_err(|e| map_session_err(&id, e, action))
287 })
288 .await
289 .map_err(|_| ApiError::internal("session lease reservation failed"))?
290 }
291
292 /// `GET /v1/sessions/{id}` query options.
293 #[derive(Debug, Deserialize, Default)]
294 pub(super) struct SessionDetailQuery {
295 /// When true, return a bounded, redacted [`SessionPeek`] instead of the
296 /// full transcript. The dashboard always asks for this: shipping a
297 /// multi-megabyte transcript to a browser in order to show twelve lines is
298 /// both wasteful and a needless place to re-emit secrets.
299 #[serde(default)]
300 peek: Option<bool>,
301 /// Entry budget for the peek, clamped to [`MAX_PEEK_ENTRIES`].
302 #[serde(default)]
303 entries: Option<usize>,
304 }
305
306 /// Either the full session or a bounded peek, chosen by `?peek=true`.
307 #[derive(Debug, Serialize)]
308 #[serde(untagged)]
309 pub(super) enum SessionDetailOrPeek {
310 Peek(Box<SessionPeek>),
311 Detail(Box<SessionDetailResponse>),
312 }
313
314 pub(super) async fn get_session(
315 State(state): State<RuntimeApiState>,
316 Path(id): Path<String>,
317 Query(query): Query<SessionDetailQuery>,
318 ) -> Result<Json<SessionDetailOrPeek>, ApiError> {
319 let manager = SessionManager::new(state.sessions_dir.clone())
320 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
321 let session = manager
322 .load_session(&id)
323 .map_err(|e| map_session_err(&id, e, "read"))?;
324
325 if query.peek.unwrap_or(false) {
326 let entries = query.entries.unwrap_or(MAX_PEEK_ENTRIES);
327 return Ok(Json(SessionDetailOrPeek::Peek(Box::new(build_peek(
328 &session, entries,
329 )))));
330 }
331 Ok(Json(SessionDetailOrPeek::Detail(Box::new(
332 session_to_detail(session),
333 ))))
334 }
335
336 /// `POST /v1/sessions/{id}/resume-thread` — open a saved session as a live
337 /// thread.
338 ///
339 /// Idempotent for a conversation that is already open: when an active thread
340 /// already holds this session (and its checkpoint still describes the file),
341 /// that thread is returned with `200 OK` instead of minting a second one with
342 /// `201 Created`. Minting unconditionally is what made "continue this
343 /// conversation" grow the rail by a row per visit.
344 ///
345 /// `req.model` / `req.mode` apply only when a thread is created. An open thread
346 /// keeps the route it was opened with — a caller that needs a different route
347 /// is creating a conversation, not resuming one.
348 ///
349 /// `message_count` reports the *saved session's* count. A reused thread may hold
350 /// more than that: it keeps the turns it ran after the session's last save.
351 pub(super) async fn resume_session_thread(
352 State(state): State<RuntimeApiState>,
353 Path(id): Path<String>,
354 Json(req): Json<ResumeSessionRequest>,
355 ) -> Result<(StatusCode, Json<ResumeSessionResponse>), ApiError> {
356 resume_session_in_runtime(
357 &state.runtime_threads,
358 &state.sessions_dir,
359 &id,
360 req,
361 (
362 state.config_path.as_deref(),
363 state.config_profile.as_deref(),
364 ),
365 None,
366 )
367 .await
368 }
369
370 pub(crate) async fn resume_session_in_runtime(
371 runtime: &std::sync::Arc<RuntimeThreadManager>,
372 sessions_dir: &std::path::Path,
373 id: &str,
374 req: ResumeSessionRequest,
375 config_source: (Option<&std::path::Path>, Option<&str>),
376 allow_shell: Option<bool>,
377 ) -> Result<(StatusCode, Json<ResumeSessionResponse>), ApiError> {
378 let _checkpoint_admission = runtime.session_checkpoint_guard().await;
379 let manager = SessionManager::new(sessions_dir.to_path_buf())
380 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
381 let session = manager
382 .resume_session(id)
383 .map_err(|e| map_session_err(id, e, "read"))?
384 .session;
385
386 if runtime.is_acp_host()
387 && session
388 .metadata
389 .runtime_store
390 .as_ref()
391 .is_some_and(|binding| *binding != runtime.session_store_binding())
392 {
393 return Err(ApiError::conflict(
394 "This saved conversation belongs to another Runtime store; secure owner attachment is not qualified, so ACP will not copy its history",
395 ));
396 }
397
398 // Validate imported image bytes before allocating a Runtime thread. This
399 // retains local history's existing bounds; invalid content cannot leave an
400 // empty session, and no path or remote image reference is dereferenced.
401 for message in session
402 .messages
403 .iter()
404 .filter(|message| message.role == Role::User)
405 {
406 crate::image_attach::runtime_images_from_blocks(&message.content).map_err(|error| {
407 ApiError::bad_request(format!("Cannot restore session image: {error}"))
408 })?;
409 }
410
411 // The conversation may already be open. Answer with the thread that holds
412 // it rather than adding a second row for the same history (see
413 // `RuntimeThreadManager::thread_holding_session`).
414 if let Some(existing) = runtime.thread_holding_session(id, &session) {
415 let thread_id = existing.id;
416 let message_count = session.messages.len();
417 let summary = format!(
418 "Session '{}' is already open in thread {thread_id} ({message_count} messages)",
419 session.metadata.title
420 );
421 return Ok((
422 StatusCode::OK,
423 Json(ResumeSessionResponse {
424 thread_id,
425 session_id: id.to_string(),
426 message_count,
427 summary,
428 }),
429 ));
430 }
431
432 let model = req.model.unwrap_or_else(|| session.metadata.model.clone());
433 let mode = req.mode.unwrap_or_else(|| {
434 session
435 .metadata
436 .mode
437 .clone()
438 .unwrap_or_else(|| "agent".to_string())
439 });
440
441 let thread = runtime
442 .create_thread_with_shell_policy(
443 CreateThreadRequest {
444 model: Some(model),
445 model_provider: Some(session.metadata.model_provider.clone()),
446 model_provider_id: session.metadata.model_provider_id.clone(),
447 workspace: Some(session.metadata.workspace.clone()),
448 mode: Some(mode),
449 allow_shell,
450 trust_mode: None,
451 auto_approve: None,
452 archived: false,
453 system_prompt: session.system_prompt.clone(),
454 task_id: None,
455 ..Default::default()
456 },
457 config_source.0,
458 config_source.1,
459 )
460 .await
461 .map_err(map_resume_thread_create_err)?;
462
463 let msg_count = session.messages.len();
464 runtime
465 .seed_thread_from_messages(&thread.id, &session.messages)
466 .await
467 .map_err(|e| ApiError::internal(format!("Failed to seed thread history: {e}")))?;
468
469 // Link the session to the new thread so that `ensure_engine_loaded`
470 // can restore the full message history from the session file.
471 runtime
472 .set_thread_session_checkpoint(&thread.id, &session)
473 .await
474 .map_err(|e| {
475 ApiError::internal(format!(
476 "Saved session was read but its Runtime checkpoint could not be bound: {e}"
477 ))
478 })?;
479
480 let summary = format!(
481 "Resumed session '{}' ({} messages) into thread {}",
482 session.metadata.title, msg_count, thread.id
483 );
484
485 Ok((
486 StatusCode::CREATED,
487 Json(ResumeSessionResponse {
488 thread_id: thread.id,
489 session_id: id.to_string(),
490 message_count: msg_count,
491 summary,
492 }),
493 ))
494 }
495
496 pub(super) async fn create_session_from_thread(
497 State(state): State<RuntimeApiState>,
498 Json(req): Json<CreateSessionRequest>,
499 ) -> Result<(StatusCode, Json<CreateSessionResponse>), ApiError> {
500 let _checkpoint_admission = state.runtime_threads.session_checkpoint_guard().await;
501 let thread_id = req.thread_id.trim();
502 if thread_id.is_empty() {
503 return Err(ApiError::bad_request("thread_id is required"));
504 }
505
506 let detail = state
507 .runtime_threads
508 .get_thread_detail(thread_id)
509 .await
510 .map_err(map_thread_err)?;
511
512 if thread_detail_has_live_work(&detail) {
513 return Err(ApiError {
514 status: StatusCode::CONFLICT,
515 message: format!(
516 "Thread {thread_id} has a queued or active turn; wait for completion before saving as a session"
517 ),
518 code: None,
519 });
520 }
521
522 let messages = messages_from_thread_detail(&detail).map_err(|error| {
523 ApiError::internal(format!("Failed to reconstruct thread history: {error}"))
524 })?;
525 if messages.is_empty() {
526 return Err(ApiError::bad_request(format!(
527 "Thread {thread_id} has no user or assistant messages to save"
528 )));
529 }
530
531 let manager = SessionManager::new(state.sessions_dir.clone())
532 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
533 let total_tokens = total_tokens_from_thread_detail(&detail);
534 // Export is idempotent (#6144). Every POST used to mint a fresh document
535 // and rebind the thread to it, so each re-export left the previous
536 // document unreferenced, and a crash between the save and the bind below
537 // left the new one unreferenced too. Export writes only the document id
538 // derived from the thread — the one its engine already writes artifacts
539 // under — so a re-export or a retry after such a crash updates and binds
540 // the document it already wrote.
541 //
542 // The document the thread is currently bound to is deliberately *not*
543 // the target: a thread opened with `resume-thread` is bound to the
544 // original saved session, often a TUI conversation, and rewriting it from
545 // this thread's lossier projection would drop its images, tool work and
546 // system prompt. Export leaves that document untouched.
547 let session_handle = crate::runtime_threads::thread_session_id(&detail.thread.id);
548 let _lease = reserve_external_session_write(&state, &session_handle, "export").await?;
549 let existing = match manager.load_session(&session_handle) {
550 Ok(existing) => Some(existing),
551 Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
552 Err(error) => return Err(map_session_err(&session_handle, error, "read")),
553 };
554 let created = existing.is_none();
555 let mut session = match existing {
556 Some(existing) => {
557 let mut updated =
558 crate::session_manager::update_session(existing, &messages, total_tokens, None);
559 updated.metadata.model = detail.thread.model.clone();
560 updated.metadata.mode = Some(detail.thread.mode.clone());
561 updated
562 }
563 None => create_saved_session_with_id_and_mode(
564 session_handle.clone(),
565 &messages,
566 &detail.thread.model,
567 &detail.thread.workspace,
568 total_tokens,
569 None,
570 Some(&detail.thread.mode),
571 ),
572 };
573 {
574 let config = state.runtime_threads.read_config();
575 stamp_session_provider_from_thread(&config, &detail, &mut session.metadata).map_err(
576 |reason| {
577 ApiError::bad_request(format!(
578 "Thread {thread_id} provider route is unavailable; session export will not fall back: {reason}"
579 ))
580 },
581 )?;
582 }
583 session.system_prompt = detail.thread.system_prompt.clone();
584
585 if let Some(title) =
586 session_title_override(req.title.as_deref(), detail.thread.title.as_deref())
587 {
588 session.metadata.title = title;
589 }
590 let title = session.metadata.title.clone();
591 let message_count = session.metadata.message_count;
592
593 persist_thread_cost(&state.runtime_threads, thread_id, &mut session).await?;
594
595 manager
596 .save_session(&session)
597 .map_err(|e| ApiError::internal(format!("Failed to save session: {e}")))?;
598
599 // Link the session to the thread so that `ensure_engine_loaded` can
600 // restore the full message history from the session file.
601 state
602 .runtime_threads
603 .set_thread_session_checkpoint(&detail.thread.id, &session)
604 .await
605 .map_err(|e| {
606 ApiError::internal(format!(
607 "Session was saved but its Runtime checkpoint could not be bound: {e}"
608 ))
609 })?;
610
611 Ok((
612 if created {
613 StatusCode::CREATED
614 } else {
615 StatusCode::OK
616 },
617 Json(CreateSessionResponse {
618 session_id: session_handle,
619 thread_id: detail.thread.id,
620 message_count,
621 title,
622 }),
623 ))
624 }
625
626 pub(super) fn stamp_session_provider_from_thread(
627 config: &crate::config::Config,
628 detail: &ThreadDetail,
629 metadata: &mut crate::session_manager::SessionMetadata,
630 ) -> Result<(), String> {
631 let thread_has_route = detail
632 .thread
633 .model_provider
634 .as_deref()
635 .is_some_and(|provider| !provider.trim().is_empty())
636 || detail.thread.model_provider_id.is_some();
637 let provider_identity = if thread_has_route {
638 config.resolve_persisted_provider_identity(
639 detail.thread.model_provider.as_deref(),
640 detail.thread.model_provider_id.as_deref(),
641 )?
642 } else if let Some(turn) = detail.turns.iter().rev().find(|turn| {
643 turn.effective_provider
644 .as_deref()
645 .is_some_and(|provider| !provider.trim().is_empty())
646 || turn.effective_provider_id.is_some()
647 }) {
648 config.resolve_persisted_provider_identity(
649 turn.effective_provider.as_deref(),
650 turn.effective_provider_id.as_deref(),
651 )?
652 } else {
653 let key = config
654 .provider
655 .as_deref()
656 .unwrap_or(crate::config::ProviderKind::Deepseek.as_str());
657 config.resolve_provider_identity(key)?
658 };
659 metadata.set_model_provider_route(
660 provider_identity.provider.as_str(),
661 provider_identity.persisted_id(),
662 );
663 Ok(())
664 }
665
666 pub(super) fn thread_detail_has_live_work(detail: &ThreadDetail) -> bool {
667 detail.turns.iter().any(|turn| {
668 matches!(
669 turn.status,
670 RuntimeTurnStatus::Queued | RuntimeTurnStatus::InProgress
671 )
672 }) || detail.items.iter().any(|item| {
673 matches!(
674 item.status,
675 TurnItemLifecycleStatus::Queued | TurnItemLifecycleStatus::InProgress
676 )
677 })
678 }
679
680 pub(super) fn messages_from_thread_detail(detail: &ThreadDetail) -> anyhow::Result<Vec<Message>> {
681 let mut items_by_turn = HashMap::new();
682 for item in &detail.items {
683 items_by_turn
684 .entry(item.turn_id.clone())
685 .or_insert_with(Vec::new)
686 .push(item.clone());
687 }
688 RuntimeThreadManager::reconstruct_messages_from_turns_with(&detail.turns, &items_by_turn)
689 }
690
691 /// Merge the thread's authoritative cost into a session about to be saved.
692 ///
693 /// The engine snapshot carries messages/tokens but no cost — cost lives in
694 /// the turn records' route-audited usage — so derive it from the same
695 /// accumulation that powers `/v1/usage` (recorded-time pricing, both
696 /// published currencies). The parent/child split mirrors the TUI writer's
697 /// field semantics (`sync_cost_to_metadata`): `session_cost_*` carries
698 /// parent-turn spend and `subagent_cost_*` routed-child spend, so a session
699 /// previously saved by the TUI never gets child spend counted twice in
700 /// `total_estimate()`. Merging each side with max keeps a session resumed
701 /// across threads from losing previously persisted spend, and extends the
702 /// monotonic display guarantee (#244) to the persisted shape. Coverage
703 /// travels with the money it qualifies (#4318): the counters are
704 /// parent-turn coverage (the TUI's own session-level accounting), CNY
705 /// included, and `coverage_recorded` marks that this writer computed them
706 /// from audited turn records rather than deserializing a legacy default.
707 async fn persist_thread_cost(
708 runtime: &std::sync::Arc<RuntimeThreadManager>,
709 thread_id: &str,
710 session: &mut crate::session_manager::SavedSession,
711 ) -> Result<(), ApiError> {
712 let usage = runtime
713 .aggregate_usage_for_thread(thread_id)
714 .await
715 .map_err(|e| ApiError::internal(format!("Failed to aggregate thread usage: {e}")))?;
716 let combined = usage.combined();
717 let cost = &mut session.metadata.cost;
718 cost.session_cost_usd = cost.session_cost_usd.max(usage.parent.cost_usd);
719 cost.session_cost_cny = cost.session_cost_cny.max(usage.parent.cost_cny);
720 cost.subagent_cost_usd = cost.subagent_cost_usd.max(usage.routed_children.cost_usd);
721 cost.subagent_cost_cny = cost.subagent_cost_cny.max(usage.routed_children.cost_cny);
722 // The display total is session + subagent, so the high-water mark rides
723 // the combined figure in both currencies.
724 cost.displayed_cost_high_water_usd = cost.displayed_cost_high_water_usd.max(combined.cost_usd);
725 cost.displayed_cost_high_water_cny = cost.displayed_cost_high_water_cny.max(combined.cost_cny);
726 cost.priced_turns = cost
727 .priced_turns
728 .max(u32::try_from(usage.parent.priced_turns).unwrap_or(u32::MAX));
729 cost.unpriced_turns = cost
730 .unpriced_turns
731 .max(u32::try_from(usage.parent.unpriced_turns).unwrap_or(u32::MAX));
732 cost.cny_priced_turns = cost
733 .cny_priced_turns
734 .max(u32::try_from(usage.parent.cny_priced_turns).unwrap_or(u32::MAX));
735 cost.cny_unpriced_turns = cost
736 .cny_unpriced_turns
737 .max(u32::try_from(usage.parent.cny_unpriced_turns).unwrap_or(u32::MAX));
738 // Coverage travels with the money (#4318): reasons and classes are the
739 // qualifiers a reload needs to treat these totals as known, not a
740 // legacy-unknown complete zero. Parent-turn coverage only — the same
741 // field the TUI writer uses; routed-child spend lives in subagent_cost_*.
742 cost.unpriced_reasons
743 .extend(usage.parent.unpriced_reasons.iter().cloned());
744 cost.cny_unpriced_reasons
745 .extend(usage.parent.cny_unpriced_reasons.iter().cloned());
746 cost.unpriced_classes
747 .extend(usage.parent.unpriced_classes.iter().cloned());
748 cost.pricing_provenances
749 .extend(usage.parent.pricing_provenances.iter().cloned());
750 cost.live_pricing_defects
751 .extend(usage.parent.live_pricing_defects.iter().cloned());
752 cost.live_pricing_unusable_defects
753 .extend(usage.parent.live_pricing_unusable_defects.iter().cloned());
754 cost.route_receipts
755 .extend(usage.parent.route_receipts.iter().cloned());
756 cost.coverage_recorded = true;
757 Ok(())
758 }
759
760 /// `PUT /v1/sessions` — save a thread's current engine state as a session.
761 ///
762 /// Unlike `POST /v1/sessions` (which reconstructs messages from stored turn
763 /// items), this endpoint asks the engine for its live session snapshot so
764 /// token counts and message ordering are authoritative.
765 ///
766 /// `session_id` names the document to write. Omitted, the thread's bound
767 /// document is updated, or, for a thread bound to none, a document is created
768 /// under the thread's own conversation id — see
769 /// [`crate::core::ops::SessionSnapshot::session_id`].
770 pub(super) async fn save_current_session(
771 State(state): State<RuntimeApiState>,
772 Json(req): Json<SaveSessionRequest>,
773 ) -> Result<Json<SaveSessionResponse>, ApiError> {
774 save_session_in_runtime(&state.runtime_threads, &state.sessions_dir, req).await
775 }
776
777 pub(crate) async fn save_session_in_runtime(
778 runtime: &std::sync::Arc<RuntimeThreadManager>,
779 sessions_dir: &std::path::Path,
780 req: SaveSessionRequest,
781 ) -> Result<Json<SaveSessionResponse>, ApiError> {
782 let _checkpoint_admission = runtime.session_checkpoint_guard().await;
783 // Find the thread to save.
784 let thread_id = match req.thread_id {
785 Some(id) => id,
786 None => {
787 // Find the most recently updated thread.
788 let threads = runtime
789 .list_threads(ThreadListFilter::IncludeArchived, Some(100))
790 .await
791 .map_err(map_thread_err)?;
792 threads
793 .into_iter()
794 .max_by_key(|t| t.updated_at)
795 .map(|t| t.id)
796 .ok_or_else(|| ApiError::bad_request("No threads to save"))?
797 }
798 };
799
800 // Get the engine handle (loads the thread into an engine if needed),
801 // then request a session snapshot. This reuses the same code path as
802 // TUI's `build_session_snapshot`: the engine holds the authoritative
803 // messages and token usage, so we don't need to reconstruct from turns.
804 let engine = runtime
805 .get_engine(&thread_id)
806 .await
807 .map_err(|e| ApiError::internal(format!("Failed to get engine for thread: {e}")))?;
808
809 let snapshot = engine
810 .get_session_snapshot()
811 .await
812 .map_err(|e| ApiError::internal(format!("Failed to get session snapshot: {e}")))?;
813
814 let manager = SessionManager::new(sessions_dir.to_path_buf())
815 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
816
817 // A document another thread is bound to is that thread's conversation.
818 // Rebinding this thread onto it would leave the other thread's checkpoint
819 // describing a document it no longer owns (#6144). A thread already bound
820 // to the same document (a legacy shared link) keeps saving to it.
821 if let Some(requested) = req.session_id.as_deref() {
822 let own = runtime
823 .get_thread(&thread_id)
824 .await
825 .map_err(map_thread_err)?;
826 if own.session_id.as_deref() != Some(requested)
827 && let Some(other) = runtime.thread_bound_to_session(requested, &thread_id)
828 {
829 return Err(ApiError {
830 status: StatusCode::CONFLICT,
831 message: format!(
832 "Session '{requested}' belongs to thread {other}; save this thread without a session_id, or into its own session"
833 ),
834 code: None,
835 });
836 }
837 }
838 // Which document this save writes. A named `session_id` wins. With none,
839 // the thread's own document answers it: the one it is bound to (a
840 // `resume-thread` from document X runs bound to X, and saving it must
841 // update X, not start a second copy under another name), else a new
842 // document under the engine's conversation id, which for a Runtime thread
843 // is the thread's own id (see `ensure_engine_loaded`) and so cannot
844 // collide with another thread's document. Snapshot ownership does not
845 // depend on this binding: a thread owns the restore points recorded on
846 // its turns.
847 let document_id = match req.session_id {
848 Some(named) => named,
849 None => runtime
850 .get_thread(&thread_id)
851 .await
852 .map_err(map_thread_err)?
853 .session_id
854 .unwrap_or_else(|| snapshot.session_id.clone()),
855 };
856 let _lease = reserve_session_write(sessions_dir, &document_id, "save").await?;
857
858 // Build or update the session, mirroring TUI's `build_session_snapshot`.
859 // Only `io::ErrorKind::NotFound` falls back to creating a new session;
860 // other I/O errors (e.g. PermissionDenied) are propagated so callers
861 // don't silently overwrite a corrupt or inaccessible session file.
862 let mut session = match manager.load_session(&document_id) {
863 Ok(existing) => {
864 let mut updated = crate::session_manager::update_session(
865 existing,
866 &snapshot.messages,
867 snapshot.total_tokens,
868 snapshot.system_prompt.as_ref(),
869 );
870 updated.metadata.model = snapshot.model.clone();
871 updated.metadata.set_model_provider_route(
872 &snapshot.model_provider,
873 snapshot.model_provider_id.as_deref(),
874 );
875 updated.metadata.mode = Some(snapshot.mode.clone());
876 updated
877 }
878 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
879 let mut session = crate::session_manager::create_saved_session_with_id_and_mode(
880 document_id.clone(),
881 &snapshot.messages,
882 &snapshot.model,
883 &snapshot.workspace,
884 snapshot.total_tokens,
885 snapshot.system_prompt.as_ref(),
886 Some(snapshot.mode.as_str()),
887 );
888 session.metadata.set_model_provider_route(
889 &snapshot.model_provider,
890 snapshot.model_provider_id.as_deref(),
891 );
892 session
893 }
894 Err(e) => {
895 return Err(ApiError::internal(format!(
896 "Failed to load session {document_id}: {e}"
897 )));
898 }
899 };
900
901 persist_thread_cost(runtime, &thread_id, &mut session).await?;
902
903 if runtime.is_acp_host() {
904 session.metadata.runtime_store = Some(runtime.session_store_binding());
905 }
906
907 // Save the session.
908 manager
909 .save_session(&session)
910 .map_err(|e| ApiError::internal(format!("Failed to save session: {e}")))?;
911
912 // Link the session to the thread so that `ensure_engine_loaded` can
913 // restore the full message history (including thinking/tool blocks)
914 // from the session file instead of reconstructing from turns.
915 let session_handle = session.metadata.id.clone();
916 runtime
917 .set_thread_session_checkpoint(&thread_id, &session)
918 .await
919 .map_err(|e| {
920 ApiError::internal(format!(
921 "Session was saved but its Runtime checkpoint could not be bound: {e}"
922 ))
923 })?;
924
925 Ok(Json(SaveSessionResponse {
926 session_id: session_handle,
927 session: session_to_detail(session),
928 }))
929 }
930
931 fn total_tokens_from_thread_detail(detail: &ThreadDetail) -> u64 {
932 detail
933 .turns
934 .iter()
935 .filter_map(|turn| turn.usage.as_ref())
936 .map(|usage| u64::from(usage.input_tokens) + u64::from(usage.output_tokens))
937 .sum()
938 }
939
940 fn session_title_override(requested: Option<&str>, thread_title: Option<&str>) -> Option<String> {
941 requested
942 .and_then(nonempty_title)
943 .or_else(|| thread_title.and_then(nonempty_title))
944 }
945
946 fn nonempty_title(title: &str) -> Option<String> {
947 let trimmed = title.trim();
948 if trimmed.is_empty() {
949 None
950 } else {
951 Some(truncate_text(trimmed, 50))
952 }
953 }
954
955 pub(super) async fn delete_session(
956 State(state): State<RuntimeApiState>,
957 Path(id): Path<String>,
958 ) -> Result<StatusCode, ApiError> {
959 // Deletion validates the id (400), refuses an unknown one (404) before
960 // creating any lease file, and holds the session's live lease, refusing
961 // (409) a document an interactive session holds open: its next autosave,
962 // in whichever process holds it, would undo the delete. Taking the lease
963 // may retry with short sleeps, so all of it runs off the async worker.
964 tokio::task::spawn_blocking(move || {
965 let manager = SessionManager::new(state.sessions_dir.clone())
966 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
967 manager
968 .delete_session(&id)
969 .map_err(|e| map_session_err(&id, e, "delete"))?;
970 // Threads bound to the document keep their turns; drop the dead link
971 // so they load from those instead of failing (#6144).
972 if let Err(error) = state.runtime_threads.unbind_session_threads(&id) {
973 tracing::warn!(session_id = %id, %error, "deleted session's threads were not unbound");
974 }
975 Ok::<_, ApiError>(StatusCode::NO_CONTENT)
976 })
977 .await
978 .map_err(|_| ApiError::internal("session delete failed"))?
979 }
980
981 /// `GET /v1/sessions/repair`: what the last session-store repair did (#6144).
982 /// `null` when none has completed.
983 pub(super) async fn get_session_repair(
984 State(state): State<RuntimeApiState>,
985 ) -> Json<Option<crate::session_reconcile::ReconcileSummary>> {
986 Json(crate::session_reconcile::last_run(&state.sessions_dir))
987 }
988
989 pub(super) fn session_to_detail(session: SavedSession) -> SessionDetailResponse {
990 let messages: Vec<Value> = session
991 .messages
992 .iter()
993 .map(|msg| {
994 let content_blocks: Vec<Value> = msg
995 .content
996 .iter()
997 .map(|block| match block {
998 codewhale_models::ContentBlock::Text { text, .. } => {
999 json!({ "type": "text", "text": text })
1000 }
1001 codewhale_models::ContentBlock::Thinking { thinking, .. } => {
1002 json!({ "type": "thinking", "text": thinking })
1003 }
1004 codewhale_models::ContentBlock::ToolUse {
1005 id,
1006 name,
1007 input,
1008 caller, ..} => {
1009 let mut obj =
1010 json!({ "type": "tool_use", "id": id, "name": name, "input": input });
1011 if let Some(caller) = caller {
1012 obj["caller"] = json!(caller);
1013 }
1014 obj
1015 }
1016 codewhale_models::ContentBlock::ToolResult {
1017 tool_use_id,
1018 content,
1019 is_error,
1020 content_blocks,
1021 ..
1022 } => {
1023 let mut obj = json!({ "type": "tool_result", "tool_use_id": tool_use_id });
1024 if let Some(cbs) = content_blocks {
1025 obj["content_blocks"] = json!(cbs);
1026 if !content.is_empty() {
1027 obj["content"] = json!(content);
1028 }
1029 } else {
1030 obj["content"] = json!(content);
1031 }
1032 if let Some(e) = is_error {
1033 obj["is_error"] = json!(e);
1034 }
1035 obj
1036 }
1037 codewhale_models::ContentBlock::ServerToolUse { id, name, input } => {
1038 json!({ "type": "tool_use", "id": id, "name": name, "input": input })
1039 }
1040 codewhale_models::ContentBlock::ToolSearchToolResult {
1041 tool_use_id,
1042 content,
1043 } => {
1044 json!({ "type": "tool_result", "tool_use_id": tool_use_id, "content": content })
1045 }
1046 codewhale_models::ContentBlock::CodeExecutionToolResult {
1047 tool_use_id,
1048 content,
1049 } => {
1050 json!({ "type": "tool_result", "tool_use_id": tool_use_id, "content": content })
1051 }
1052 codewhale_models::ContentBlock::ImageUrl { .. } => Value::Null,
1053 })
1054 .collect();
1055 json!({
1056 "role": msg.role,
1057 "content": content_blocks,
1058 })
1059 })
1060 .collect();
1061 SessionDetailResponse {
1062 metadata: session.metadata,
1063 messages,
1064 system_prompt: session.system_prompt,
1065 turn_outcomes: session.turn_outcomes,
1066 }
1067 }
1068
1069 fn map_session_err(id: &str, err: std::io::Error, action: &str) -> ApiError {
1070 match err.kind() {
1071 std::io::ErrorKind::NotFound => ApiError::not_found(format!("Session '{id}' not found")),
1072 std::io::ErrorKind::InvalidData => {
1073 ApiError::bad_request(format!("Failed to parse session '{id}': {err}"))
1074 }
1075 std::io::ErrorKind::InvalidInput => {
1076 ApiError::bad_request(format!("Invalid session id '{id}'"))
1077 }
1078 // The session is open in an interactive Codewhale session, which holds
1079 // the authoritative copy in memory. Fail closed with a typed conflict
1080 // rather than write something its next autosave would revert.
1081 std::io::ErrorKind::ResourceBusy => ApiError {
1082 status: StatusCode::CONFLICT,
1083 message: err.to_string(),
1084 code: None,
1085 },
1086 _ => ApiError::internal(format!("Failed to {action} session '{id}': {err}")),
1087 }
1088 }
1089
1090 fn map_resume_thread_create_err(err: anyhow::Error) -> ApiError {
1091 let reason = err.to_string();
1092 let message = format!("Failed to create thread: {reason}");
1093 if reason.starts_with("saved session has an empty provider identity")
1094 || reason.starts_with("saved session requires custom provider")
1095 || reason.starts_with("legacy session records only the generic `custom` provider kind")
1096 || reason.starts_with("legacy `provider = \"custom\"`")
1097 {
1098 ApiError::bad_request(message)
1099 } else {
1100 // Thread-store writes, event persistence, and other runtime failures
1101 // are server-side faults; never disguise them as a client config error.
1102 ApiError::internal(message)
1103 }
1104 }
1105
1106 #[cfg(test)]
1107 mod session_query_tests {
1108 use super::*;
1109
1110 fn query(
1111 include_archived: Option<bool>,
1112 archived_only: Option<bool>,
1113 sort: Option<&str>,
1114 workspace: Option<&str>,
1115 limit: Option<usize>,
1116 ) -> SessionsQuery {
1117 SessionsQuery {
1118 limit,
1119 search: Some("whale".to_string()),
1120 include_archived,
1121 archived_only,
1122 workspace: workspace.map(PathBuf::from),
1123 sort: sort.map(str::to_string),
1124 }
1125 }
1126
1127 #[test]
1128 fn archive_params_resolve_like_the_threads_routes() {
1129 assert_eq!(
1130 projection_query(&query(None, None, None, None, None)).filter,
1131 SessionListFilter::ActiveOnly
1132 );
1133 assert_eq!(
1134 projection_query(&query(Some(true), None, None, None, None)).filter,
1135 SessionListFilter::IncludeArchived
1136 );
1137 assert_eq!(
1138 projection_query(&query(Some(true), Some(true), None, None, None)).filter,
1139 SessionListFilter::ArchivedOnly
1140 );
1141 }
1142
1143 #[test]
1144 fn sort_and_workspace_scope_flow_through_and_bad_sorts_fall_back() {
1145 let projected = projection_query(&query(None, None, Some("name"), Some("/repo"), Some(9)));
1146 assert_eq!(projected.sort, SessionSortMode::Name);
1147 // `Path` in this module is `axum::extract::Path`; spell out the std one.
1148 assert_eq!(
1149 projected.workspace_scope.as_deref(),
1150 Some(std::path::Path::new("/repo"))
1151 );
1152 assert_eq!(projected.limit, 9);
1153 assert_eq!(projected.search, "whale");
1154
1155 // An unknown sort must not fail the request — a stale client should
1156 // still get a listing, just in the default order.
1157 assert_eq!(
1158 projection_query(&query(None, None, Some("nonsense"), None, None)).sort,
1159 SessionSortMode::Recent
1160 );
1161 }
1162
1163 #[test]
1164 fn limit_is_clamped_at_both_ends() {
1165 assert_eq!(
1166 projection_query(&query(None, None, None, None, Some(0))).limit,
1167 1
1168 );
1169 assert_eq!(
1170 projection_query(&query(None, None, None, None, Some(10_000))).limit,
1171 500
1172 );
1173 // Absent limit keeps the historical page size.
1174 assert_eq!(
1175 projection_query(&query(None, None, None, None, None)).limit,
1176 50
1177 );
1178 }
1179
1180 #[test]
1181 fn absent_workspace_means_every_workspace() {
1182 assert!(
1183 projection_query(&query(None, None, None, None, None))
1184 .workspace_scope
1185 .is_none(),
1186 "the API must not silently scope to the runtime's own CWD"
1187 );
1188 }
1189 }
1190
1191 #[cfg(test)]
1192 mod resume_thread_error_tests {
1193 use super::*;
1194
1195 #[test]
1196 fn provider_config_errors_are_client_errors_but_storage_errors_stay_internal() {
1197 let provider = map_resume_thread_create_err(anyhow::anyhow!(
1198 "saved session requires custom provider 'lm-studio', but `[providers.lm-studio]` is missing"
1199 ));
1200 assert_eq!(provider.status, StatusCode::BAD_REQUEST);
1201
1202 let storage = map_resume_thread_create_err(anyhow::anyhow!(
1203 "Failed to save runtime thread: permission denied"
1204 ));
1205 assert_eq!(storage.status, StatusCode::INTERNAL_SERVER_ERROR);
1206 }
1207 }
1208
1209 // ---------------------------------------------------------------------------
1210 // Session artifacts (#6163): the oversized tool outputs a session recorded as
1211 // `ArtifactRecord`s live under `sessions/<id>/artifacts/`. These routes list
1212 // the records a saved session carries and read one artifact through the same
1213 // confined opener the workspace file routes use. Nothing is copied anywhere.
1214 // ---------------------------------------------------------------------------
1215
1216 #[derive(Debug, Serialize)]
1217 pub(super) struct SessionArtifactSummary {
1218 id: String,
1219 session_id: String,
1220 #[serde(skip_serializing_if = "Option::is_none")]
1221 content_type: Option<String>,
1222 kind: crate::artifacts::ArtifactKind,
1223 tool_call_id: String,
1224 tool_name: String,
1225 created_at: chrono::DateTime<chrono::Utc>,
1226 byte_size: u64,
1227 preview: String,
1228 /// Session-relative storage path with `/` separators.
1229 path: String,
1230 }
1231
1232 #[derive(Debug, Serialize)]
1233 pub(super) struct SessionArtifactsResponse {
1234 session_id: String,
1235 artifacts: Vec<SessionArtifactSummary>,
1236 }
1237
1238 #[derive(Deserialize)]
1239 #[serde(deny_unknown_fields)]
1240 pub(super) struct SessionArtifactReadQuery {
1241 offset: Option<usize>,
1242 limit: Option<usize>,
1243 }
1244
1245 #[derive(Debug, Serialize)]
1246 pub(super) struct SessionArtifactReadResponse {
1247 artifact: SessionArtifactSummary,
1248 size: u64,
1249 revision: String,
1250 offset: usize,
1251 bytes: usize,
1252 truncated: bool,
1253 encoding: &'static str,
1254 content: String,
1255 }
1256
1257 fn artifact_summary(record: &crate::artifacts::ArtifactRecord) -> SessionArtifactSummary {
1258 SessionArtifactSummary {
1259 id: record.id.clone(),
1260 session_id: record.session_id.clone(),
1261 content_type: None,
1262 kind: record.kind.clone(),
1263 tool_call_id: record.tool_call_id.clone(),
1264 tool_name: record.tool_name.clone(),
1265 created_at: record.created_at,
1266 byte_size: record.byte_size,
1267 preview: record.preview.clone(),
1268 path: crate::artifacts::format_artifact_relative_path(&record.storage_path),
1269 }
1270 }
1271
1272 pub(super) async fn list_session_artifacts(
1273 State(state): State<RuntimeApiState>,
1274 Path(id): Path<String>,
1275 ) -> Result<Json<SessionArtifactsResponse>, ApiError> {
1276 let manager = SessionManager::new(state.sessions_dir.clone())
1277 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
1278 let session = manager
1279 .load_session(&id)
1280 .map_err(|e| map_session_err(&id, e, "read"))?;
1281 Ok(Json(SessionArtifactsResponse {
1282 session_id: session.metadata.id.clone(),
1283 artifacts: session.artifacts.iter().map(artifact_summary).collect(),
1284 }))
1285 }
1286
1287 pub(super) async fn read_session_artifact(
1288 State(state): State<RuntimeApiState>,
1289 Path((id, artifact_id)): Path<(String, String)>,
1290 Query(query): Query<SessionArtifactReadQuery>,
1291 ) -> Result<Json<SessionArtifactReadResponse>, ApiError> {
1292 let (offset, limit) = super::workspace::parse_read_window(query.offset, query.limit)?;
1293 tokio::task::spawn_blocking(move || {
1294 read_session_artifact_window(&state.sessions_dir, &id, &artifact_id, offset, limit)
1295 })
1296 .await
1297 .map_err(|_| ApiError::internal("session artifact read failed"))?
1298 .map(Json)
1299 }
1300
1301 fn read_session_artifact_window(
1302 sessions_dir: &std::path::Path,
1303 id: &str,
1304 artifact_id: &str,
1305 offset: usize,
1306 limit: usize,
1307 ) -> Result<SessionArtifactReadResponse, ApiError> {
1308 let resolved = resolve_session_artifact(
1309 sessions_dir,
1310 id,
1311 artifact_id,
1312 ArtifactAuthority::SavedSession,
1313 )?;
1314 let summary = resolved
1315 .summary
1316 .ok_or_else(|| ApiError::internal("saved-session artifact has no summary"))?;
1317 let read = resolved.read;
1318 let (window, truncated) = super::workspace::read_window(&read.bytes, offset, limit);
1319 let (encoding, content) = super::workspace::encode_window(window);
1320 Ok(SessionArtifactReadResponse {
1321 artifact: summary,
1322 size: read.size,
1323 revision: read.revision,
1324 offset: offset.min(read.bytes.len()),
1325 bytes: window.len(),
1326 truncated,
1327 encoding,
1328 content,
1329 })
1330 }
1331
1332 /// Who vouches that `artifact_id` belongs to session `id`.
1333 pub(super) enum ArtifactAuthority<'a> {
1334 /// The SavedSession JSON's `artifacts` index.
1335 SavedSession,
1336 /// A runtime turn's recorded reference. A Runtime engine runs under its
1337 /// thread's own id (#6621), which has no SavedSession index, so the turn
1338 /// record is the ownership proof: the bytes must sit at `path` and, when `revision` is
1339 /// known, hash to it.
1340 TurnRef {
1341 path: &'a str,
1342 revision: Option<&'a str>,
1343 },
1344 }
1345
1346 /// One session artifact's bytes, read through the confined opener.
1347 pub(super) struct ResolvedSessionArtifact {
1348 /// The index or manifest record; `None` for a turn-ref read of a
1349 /// non-image artifact, whose caller already holds the reference.
1350 pub(super) summary: Option<SessionArtifactSummary>,
1351 pub(super) read: super::workspace::ConfinedFileBytes,
1352 }
1353
1354 /// The one resolver behind both session-artifact reads: the session route
1355 /// (SavedSession authority) and the turn route (turn-ref authority). Both
1356 /// get the same session-id validation, confinement, image-manifest checks
1357 /// and integrity checks.
1358 pub(super) fn resolve_session_artifact(
1359 sessions_dir: &std::path::Path,
1360 id: &str,
1361 artifact_id: &str,
1362 authority: ArtifactAuthority<'_>,
1363 ) -> Result<ResolvedSessionArtifact, ApiError> {
1364 if !crate::artifacts::is_valid_session_id(id) {
1365 return Err(ApiError::bad_request("invalid session id"));
1366 }
1367 // The reserved image namespace always requires its immutable manifest,
1368 // even if a later SavedSession index also mentions that handle.
1369 let image_handle = artifact_id.strip_prefix("art_image_").is_some_and(|hash| {
1370 hash.len() == 64
1371 && hash
1372 .bytes()
1373 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase())
1374 });
1375 let (summary, relative_path, evidence) = if image_handle {
1376 let (summary, evidence) = image_evidence_summary(sessions_dir, id, artifact_id)?;
1377 let path = summary.path.clone();
1378 (Some(summary), path, Some(evidence))
1379 } else {
1380 match &authority {
1381 ArtifactAuthority::SavedSession => {
1382 let summary = saved_session_summary(sessions_dir, id, artifact_id)?;
1383 let path = summary.path.clone();
1384 (Some(summary), path, None)
1385 }
1386 ArtifactAuthority::TurnRef { path, .. } => {
1387 let relative = PathBuf::from(path);
1388 if relative.is_absolute() || !crate::fleet::files::path_is_confined(&relative) {
1389 return Err(ApiError::forbidden(
1390 "artifact reference is not confined to its session",
1391 ));
1392 }
1393 (None, (*path).to_string(), None)
1394 }
1395 }
1396 };
1397 let relative = PathBuf::from(id).join(&relative_path);
1398 let opened = super::workspace::open_confined_file(sessions_dir, &relative, false)
1399 .and_then(|file| super::workspace::read_confined_bytes(&file));
1400 let read = match (opened, &authority) {
1401 (Ok(read), _) => read,
1402 // A turn recorded these bytes; their absence means the session
1403 // directory was pruned since.
1404 (Err(error), ArtifactAuthority::TurnRef { .. })
1405 if error.status == StatusCode::NOT_FOUND =>
1406 {
1407 return Err(ApiError::gone(
1408 "this artifact's bytes are no longer stored (its session was pruned)",
1409 ));
1410 }
1411 (Err(error), _) => return Err(error),
1412 };
1413 if let Some(evidence) = evidence
1414 && (read.size != evidence.size_bytes
1415 || read.revision != evidence.digest
1416 || crate::image_attach::sniff_media_type(&read.bytes)
1417 != Some(evidence.content_type.as_str())
1418 || crate::image_attach::decode_and_guard_image(&read.bytes).is_err())
1419 {
1420 return Err(ApiError::bad_request(
1421 "image evidence integrity check failed",
1422 ));
1423 }
1424 if let ArtifactAuthority::TurnRef {
1425 revision: Some(expected),
1426 ..
1427 } = authority
1428 && read.revision != expected
1429 {
1430 return Err(ApiError::conflict(format!(
1431 "artifact bytes changed since the turn recorded them; current revision is {}",
1432 read.revision
1433 )));
1434 }
1435 Ok(ResolvedSessionArtifact { summary, read })
1436 }
1437
1438 fn saved_session_summary(
1439 sessions_dir: &std::path::Path,
1440 id: &str,
1441 artifact_id: &str,
1442 ) -> Result<SessionArtifactSummary, ApiError> {
1443 let manager = SessionManager::new(sessions_dir.to_path_buf())
1444 .map_err(|e| ApiError::internal(format!("Failed to open sessions dir: {e}")))?;
1445 let record = match manager.load_session_snapshot(id) {
1446 Ok(session) if session.metadata.id == id => session
1447 .artifacts
1448 .into_iter()
1449 .find(|record| record.id == artifact_id),
1450 Ok(_) => return Err(ApiError::forbidden("artifact session owner does not match")),
1451 Err(error) if error.kind() == std::io::ErrorKind::NotFound => None,
1452 Err(error) => return Err(map_session_err(id, error, "read")),
1453 };
1454 let record = record.ok_or_else(|| ApiError::not_found("artifact not found"))?;
1455 if record.storage_path.is_absolute()
1456 || !crate::fleet::files::path_is_confined(&record.storage_path)
1457 || (!record.session_id.is_empty() && record.session_id != id)
1458 {
1459 return Err(ApiError::forbidden(
1460 "artifact record is not confined to its session",
1461 ));
1462 }
1463 let mut summary = artifact_summary(&record);
1464 summary.session_id = id.to_owned();
1465 Ok(summary)
1466 }
1467
1468 /// Fresh Engine sessions can publish immutable observations before a
1469 /// SavedSession JSON exists. Only the exact owned image manifest grants
1470 /// access; arbitrary relative paths and other evidence are not a fallback.
1471 fn image_evidence_summary(
1472 sessions_dir: &std::path::Path,
1473 id: &str,
1474 artifact_id: &str,
1475 ) -> Result<
1476 (
1477 SessionArtifactSummary,
1478 crate::tools::large_output_router::EvidenceArtifact,
1479 ),
1480 ApiError,
1481 > {
1482 let relative = PathBuf::from(id)
1483 .join(crate::tools::large_output_router::evidence_metadata_relative_path(artifact_id));
1484 let file = super::workspace::open_confined_file(sessions_dir, &relative, false)?;
1485 let evidence = crate::tools::large_output_router::read_evidence_metadata_file(&file)
1486 .map_err(|error| super::workspace::map_fs_error(error, "evidence metadata"))?;
1487 let expected_path =
1488 PathBuf::from(crate::artifacts::ARTIFACTS_DIR_NAME).join(format!("{artifact_id}.image"));
1489 if evidence.origin_session != id
1490 || evidence.handle != artifact_id
1491 || evidence.storage_path != expected_path
1492 || evidence.generation != 1
1493 || evidence.encoding != "binary"
1494 || evidence.call_id.is_empty()
1495 {
1496 return Err(ApiError::forbidden(
1497 "image evidence owner or path does not match",
1498 ));
1499 }
1500 if evidence.redacted
1501 || crate::tools::large_output_router::evidence_is_expired(
1502 &evidence,
1503 crate::tools::large_output_router::unix_millis_now(),
1504 )
1505 {
1506 return Err(ApiError::forbidden("image evidence is no longer available"));
1507 }
1508 if evidence.size_bytes > crate::image_attach::MAX_IMAGE_BYTES as u64
1509 || !matches!(
1510 evidence.content_type.as_str(),
1511 "image/png" | "image/jpeg" | "image/gif" | "image/webp"
1512 )
1513 {
1514 return Err(ApiError::bad_request("invalid image evidence"));
1515 }
1516 let created_at = i64::try_from(evidence.created_at_unix_ms)
1517 .ok()
1518 .and_then(chrono::DateTime::from_timestamp_millis)
1519 .ok_or_else(|| ApiError::bad_request("invalid evidence timestamp"))?;
1520 let summary = SessionArtifactSummary {
1521 id: artifact_id.to_owned(),
1522 session_id: id.to_owned(),
1523 content_type: Some(evidence.content_type.clone()),
1524 kind: crate::artifacts::ArtifactKind::ToolOutput,
1525 tool_call_id: evidence.call_id.clone(),
1526 tool_name: evidence.tool_name.clone(),
1527 created_at,
1528 byte_size: evidence.size_bytes,
1529 preview: String::new(),
1530 path: crate::artifacts::format_artifact_relative_path(&evidence.storage_path),
1531 };
1532 Ok((summary, evidence))
1533 }
1534
1535 #[cfg(test)]
1536 mod tool_media_artifact_tests {
1537 use super::*;
1538 use crate::tools::large_output_router::{
1539 EvidenceArtifact, EvidenceRetentionState, evidence_metadata_relative_path, unix_millis_now,
1540 };
1541 use base64::Engine as _;
1542
1543 fn fixture(root: &std::path::Path) -> EvidenceArtifact {
1544 let id = format!("art_image_{}", "a".repeat(64));
1545 let relative = PathBuf::from("artifacts").join(format!("{id}.image"));
1546 let mut bytes = std::io::Cursor::new(Vec::new());
1547 image::DynamicImage::new_rgba8(2, 1)
1548 .write_to(&mut bytes, image::ImageFormat::Png)
1549 .unwrap();
1550 let bytes = bytes.into_inner();
1551 let now = unix_millis_now();
1552 let evidence = EvidenceArtifact {
1553 handle: id,
1554 digest: crate::hashing::sha256_hex(&bytes),
1555 size_bytes: bytes.len() as u64,
1556 content_type: "image/png".into(),
1557 tool_name: "screenshot".into(),
1558 call_id: "image-call".into(),
1559 origin_session: "media_owner".into(),
1560 generation: 1,
1561 redacted: false,
1562 encoding: "binary".into(),
1563 retention_state: EvidenceRetentionState::Live,
1564 created_at_unix_ms: now,
1565 retain_until_unix_ms: now + 60_000,
1566 storage_path: relative,
1567 };
1568 std::fs::create_dir_all(root.join("media_owner/artifacts")).unwrap();
1569 std::fs::write(root.join("media_owner").join(&evidence.storage_path), bytes).unwrap();
1570 save_manifest(root, &evidence);
1571 evidence
1572 }
1573
1574 fn save_manifest(root: &std::path::Path, evidence: &EvidenceArtifact) {
1575 std::fs::write(
1576 root.join("media_owner")
1577 .join(evidence_metadata_relative_path(&evidence.handle)),
1578 serde_json::to_vec(evidence).unwrap(),
1579 )
1580 .unwrap();
1581 }
1582
1583 fn decode(response: &SessionArtifactReadResponse) -> Vec<u8> {
1584 if response.encoding == "base64" {
1585 base64::engine::general_purpose::STANDARD
1586 .decode(&response.content)
1587 .unwrap()
1588 } else {
1589 response.content.as_bytes().to_vec()
1590 }
1591 }
1592
1593 #[test]
1594 fn tool_media_artifact_reads_without_saved_session_and_retains_window_revision() {
1595 let temp = tempfile::tempdir().unwrap();
1596 let evidence = fixture(temp.path());
1597 assert!(!temp.path().join("media_owner.json").exists());
1598 let first =
1599 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 7)
1600 .unwrap();
1601 let rest =
1602 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 7, 1024)
1603 .unwrap();
1604 assert_eq!(first.artifact.session_id, "media_owner");
1605 assert_eq!(first.artifact.tool_call_id, "image-call");
1606 assert_eq!(first.artifact.content_type.as_deref(), Some("image/png"));
1607 assert_eq!(first.revision, rest.revision);
1608 assert!(first.truncated);
1609 assert!(!rest.truncated);
1610 let mut bytes = decode(&first);
1611 bytes.extend(decode(&rest));
1612 assert_eq!(crate::hashing::sha256_hex(&bytes), evidence.digest);
1613 assert_eq!(
1614 read_session_artifact_window(temp.path(), "other_owner", &evidence.handle, 0, 1024)
1615 .unwrap_err()
1616 .status,
1617 StatusCode::NOT_FOUND
1618 );
1619 assert_eq!(
1620 read_session_artifact_window(temp.path(), "../media_owner", &evidence.handle, 0, 1024)
1621 .unwrap_err()
1622 .status,
1623 StatusCode::BAD_REQUEST
1624 );
1625 }
1626
1627 #[test]
1628 fn tool_media_artifact_rejects_wrong_owner_expiry_spoofed_mime_and_changed_bytes() {
1629 let temp = tempfile::tempdir().unwrap();
1630 let evidence = fixture(temp.path());
1631 let mut owner = evidence.clone();
1632 owner.origin_session = "foreign".into();
1633 let mut handle = evidence.clone();
1634 handle.handle = "forged".into();
1635 let mut path = evidence.clone();
1636 path.storage_path = PathBuf::from("../outside.png");
1637 let mut generation = evidence.clone();
1638 generation.generation = 2;
1639 for invalid in [owner, handle, path, generation] {
1640 save_manifest(temp.path(), &invalid);
1641 assert_eq!(
1642 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 1024)
1643 .unwrap_err()
1644 .status,
1645 StatusCode::FORBIDDEN
1646 );
1647 }
1648 let mut expired = evidence.clone();
1649 expired.retain_until_unix_ms = 0;
1650 save_manifest(temp.path(), &expired);
1651 assert_eq!(
1652 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 1024)
1653 .unwrap_err()
1654 .status,
1655 StatusCode::FORBIDDEN
1656 );
1657 let mut redacted = evidence.clone();
1658 redacted.redacted = true;
1659 save_manifest(temp.path(), &redacted);
1660 assert_eq!(
1661 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 1024)
1662 .unwrap_err()
1663 .status,
1664 StatusCode::FORBIDDEN
1665 );
1666 let mut mime = evidence.clone();
1667 mime.content_type = "image/jpeg".into();
1668 save_manifest(temp.path(), &mime);
1669 assert_eq!(
1670 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 1024)
1671 .unwrap_err()
1672 .status,
1673 StatusCode::BAD_REQUEST
1674 );
1675 save_manifest(temp.path(), &evidence);
1676 std::fs::write(
1677 temp.path().join("media_owner").join(&evidence.storage_path),
1678 b"changed",
1679 )
1680 .unwrap();
1681 assert_eq!(
1682 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 1024)
1683 .unwrap_err()
1684 .status,
1685 StatusCode::BAD_REQUEST
1686 );
1687 }
1688
1689 #[cfg(unix)]
1690 #[test]
1691 fn tool_media_artifact_rejects_manifest_and_payload_symlinks() {
1692 use std::os::unix::fs::symlink;
1693 let temp = tempfile::tempdir().unwrap();
1694 let evidence = fixture(temp.path());
1695 let manifest = temp
1696 .path()
1697 .join("media_owner")
1698 .join(evidence_metadata_relative_path(&evidence.handle));
1699 let outside = temp.path().join("outside.json");
1700 std::fs::rename(&manifest, &outside).unwrap();
1701 symlink(&outside, &manifest).unwrap();
1702 assert!(
1703 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 1024)
1704 .is_err()
1705 );
1706 std::fs::remove_file(&manifest).unwrap();
1707 save_manifest(temp.path(), &evidence);
1708 let payload = temp.path().join("media_owner").join(&evidence.storage_path);
1709 let outside = temp.path().join("outside.png");
1710 std::fs::rename(&payload, &outside).unwrap();
1711 symlink(&outside, &payload).unwrap();
1712 assert!(
1713 read_session_artifact_window(temp.path(), "media_owner", &evidence.handle, 0, 1024)
1714 .is_err()
1715 );
1716 }
1717 }
1718
1719 /// Trusted transport creation uses the same lease/checkpoint writer before
1720 /// the first provider call. Ordinary HTTP export still rejects empty history.
1721 pub(crate) async fn initialize_empty_session(
1722 runtime: &std::sync::Arc<RuntimeThreadManager>,
1723 sessions_dir: &std::path::Path,
1724 thread_id: &str,
1725 session_id: &str,
1726 ) -> Result<(), ApiError> {
1727 let _checkpoint_admission = runtime.session_checkpoint_guard().await;
1728 let detail = runtime
1729 .get_thread_detail(thread_id)
1730 .await
1731 .map_err(map_thread_err)?;
1732 if thread_detail_has_live_work(&detail)
1733 || !detail.turns.is_empty()
1734 || detail.thread.session_id.is_some()
1735 {
1736 return Err(ApiError::conflict(
1737 "Initial ACP checkpoint requires a new idle, unbound thread",
1738 ));
1739 }
1740 let _lease = reserve_session_write(sessions_dir, session_id, "initialize").await?;
1741 let manager = SessionManager::new(sessions_dir.to_path_buf())
1742 .map_err(|error| ApiError::internal(error.to_string()))?;
1743 match manager.load_session(session_id) {
1744 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
1745 Ok(_) => {
1746 return Err(ApiError::conflict(
1747 "Initial session identity already exists",
1748 ));
1749 }
1750 Err(error) => return Err(map_session_err(session_id, error, "read")),
1751 }
1752 let mut session = create_saved_session_with_id_and_mode(
1753 session_id.to_string(),
1754 &[],
1755 &detail.thread.model,
1756 &detail.thread.workspace,
1757 0,
1758 None,
1759 Some(&detail.thread.mode),
1760 );
1761 stamp_session_provider_from_thread(&runtime.read_config(), &detail, &mut session.metadata)
1762 .map_err(ApiError::bad_request)?;
1763 session.metadata.runtime_store = Some(runtime.session_store_binding());
1764 manager
1765 .save_session(&session)
1766 .map_err(|error| ApiError::internal(error.to_string()))?;
1767 runtime
1768 .set_thread_session_checkpoint(thread_id, &session)
1769 .await
1770 .map_err(|error| {
1771 ApiError::internal(format!(
1772 "Initial session was saved but checkpoint binding failed: {error}"
1773 ))
1774 })?;
1775 Ok(())
1776 }
1777
1777 lines RUST