返回 CodeWhale
turn_artifacts.rs
根目录 / crates / tui / src / runtime_api / turn_artifacts.rs
1 //! Turn artifact routes (#6653): what a runtime turn produced.
2 //!
3 //! The authority is the runtime store's turn record and its items. Nothing
4 //! is scanned: every reference was recorded where its bytes were written,
5 //! and the workspace delta comes from the engine's own snapshot pair. Reads
6 //! go through the same confined openers as the workspace file and session
7 //! artifact routes; there is no second store.
8
9 use std::path::Path as FsPath;
10
11 use axum::Json;
12 use axum::extract::{Path, Query, State};
13 use serde::{Deserialize, Serialize};
14
15 use super::sessions::{ArtifactAuthority, resolve_session_artifact};
16 use super::workspace::{
17 ConfinedFileBytes, FILE_SERVE_MAX_BYTES, canonical_workspace, encode_window,
18 open_confined_file, parse_read_window, precheck_file_target, read_confined_bytes, read_window,
19 relative_request_path,
20 };
21 use super::{ApiError, RuntimeApiState, map_thread_err};
22 use crate::runtime_threads::{
23 FileChangeKind, TurnArtifactKind, TurnArtifactRef, TurnArtifactsView, TurnWorkspaceArtifacts,
24 };
25
26 pub(super) async fn list_turn_artifacts(
27 State(state): State<RuntimeApiState>,
28 Path((thread_id, turn_id)): Path<(String, String)>,
29 ) -> Result<Json<TurnArtifactsView>, ApiError> {
30 load_view(&state, &thread_id, &turn_id).await.map(Json)
31 }
32
33 async fn load_view(
34 state: &RuntimeApiState,
35 thread_id: &str,
36 turn_id: &str,
37 ) -> Result<TurnArtifactsView, ApiError> {
38 state
39 .runtime_threads
40 .turn_artifacts(thread_id, turn_id)
41 .await
42 .map_err(map_thread_err)?
43 .ok_or_else(|| ApiError::not_found(format!("turn '{turn_id}' not found in this thread")))
44 }
45
46 #[derive(Deserialize)]
47 #[serde(deny_unknown_fields)]
48 pub(super) struct TurnArtifactReadQuery {
49 offset: Option<usize>,
50 limit: Option<usize>,
51 /// Read an intermediate revision one of the turn's items recorded,
52 /// instead of the reference's final one.
53 revision: Option<String>,
54 }
55
56 #[derive(Debug, Serialize)]
57 pub(super) struct TurnArtifactReadResponse {
58 artifact: TurnArtifactRef,
59 /// `workspace` (the live file), `snapshot` (the turn's post-turn
60 /// snapshot), or `session_artifact` (the session artifact directory).
61 source: &'static str,
62 /// Whether the workspace still holds exactly these bytes. `null` for a
63 /// session artifact, or a file reference with no recorded revision.
64 current: Option<bool>,
65 size: u64,
66 /// SHA-256 of the whole content served.
67 revision: String,
68 offset: usize,
69 bytes: usize,
70 truncated: bool,
71 encoding: &'static str,
72 content: String,
73 }
74
75 pub(super) async fn read_turn_artifact(
76 State(state): State<RuntimeApiState>,
77 Path((thread_id, turn_id, artifact_id)): Path<(String, String, String)>,
78 Query(query): Query<TurnArtifactReadQuery>,
79 ) -> Result<Json<TurnArtifactReadResponse>, ApiError> {
80 let (offset, limit) = parse_read_window(query.offset, query.limit)?;
81 let view = load_view(&state, &thread_id, &turn_id).await?;
82 let artifact = select_reference(&view, &artifact_id, query.revision.as_deref())?;
83 let workspace = view.workspace.clone();
84 let thread_workspace = view.thread_workspace.clone();
85 #[cfg(test)]
86 let env_ticket = crate::test_support::env_scope_ticket();
87 tokio::task::spawn_blocking(move || {
88 #[cfg(test)]
89 let _membership = crate::test_support::join_env_scope(env_ticket);
90 let (source, current, read) = match artifact.kind {
91 TurnArtifactKind::File => {
92 read_file_artifact(&thread_workspace, workspace.as_ref(), &artifact)?
93 }
94 TurnArtifactKind::ToolOutput | TurnArtifactKind::Media => {
95 let session_id = artifact
96 .session_id
97 .as_deref()
98 .ok_or_else(|| ApiError::internal("artifact reference has no session"))?;
99 let root = crate::artifacts::artifact_sessions_root()
100 .ok_or_else(|| ApiError::internal("session artifact root is unavailable"))?;
101 let resolved = resolve_session_artifact(
102 &root,
103 session_id,
104 &artifact.id,
105 ArtifactAuthority::TurnRef {
106 path: &artifact.path,
107 revision: artifact.revision.as_deref(),
108 },
109 )?;
110 ("session_artifact", None, resolved.read)
111 }
112 };
113 let (window, truncated) = read_window(&read.bytes, offset, limit);
114 let (encoding, content) = encode_window(window);
115 Ok(TurnArtifactReadResponse {
116 artifact,
117 source,
118 current,
119 size: read.size,
120 revision: read.revision,
121 offset: offset.min(read.bytes.len()),
122 bytes: window.len(),
123 truncated,
124 encoding,
125 content,
126 })
127 })
128 .await
129 .map_err(|_| ApiError::internal("turn artifact read failed"))?
130 .map(Json)
131 }
132
133 /// The reference to serve: the turn aggregate's, or an item's when the turn
134 /// aggregate no longer lists it (a file written and then deleted still has
135 /// item-level history). `?revision=` selects the item ref that recorded
136 /// exactly that revision.
137 fn select_reference(
138 view: &TurnArtifactsView,
139 artifact_id: &str,
140 revision: Option<&str>,
141 ) -> Result<TurnArtifactRef, ApiError> {
142 let aggregate = view.artifacts.iter().find(|r| r.id == artifact_id);
143 let items = view
144 .item_artifacts
145 .iter()
146 .rev()
147 .filter(|r| r.id == artifact_id);
148 let Some(wanted) = revision.map(|revision| revision.trim().to_ascii_lowercase()) else {
149 return aggregate
150 .or_else(|| {
151 view.item_artifacts
152 .iter()
153 .rev()
154 .find(|r| r.id == artifact_id)
155 })
156 .cloned()
157 .ok_or_else(|| ApiError::not_found("artifact not found in this turn"));
158 };
159 aggregate
160 .into_iter()
161 .chain(items)
162 .find(|r| r.revision.as_deref() == Some(wanted.as_str()))
163 .cloned()
164 .ok_or_else(|| ApiError::not_found("this turn recorded no such revision of the artifact"))
165 }
166
167 /// Serve a file reference: from the workspace when it still holds the
168 /// recorded revision, otherwise from the turn's post-turn snapshot,
169 /// otherwise a conflict naming what the workspace holds now.
170 fn read_file_artifact(
171 thread_workspace: &FsPath,
172 workspace: Option<&TurnWorkspaceArtifacts>,
173 artifact: &TurnArtifactRef,
174 ) -> Result<(&'static str, Option<bool>, ConfinedFileBytes), ApiError> {
175 if artifact.change == Some(FileChangeKind::Deleted) {
176 return Err(ApiError::gone(
177 "this turn deleted the file; restore its prior content with file-revert and the reference's restore_snapshot_id",
178 ));
179 }
180 if artifact
181 .size
182 .is_some_and(|size| size > FILE_SERVE_MAX_BYTES)
183 {
184 return Err(ApiError::payload_too_large(format!(
185 "file is larger than the {FILE_SERVE_MAX_BYTES}-byte serving limit"
186 )));
187 }
188 let relative = relative_request_path(&artifact.path, false)?;
189 let root = canonical_workspace(thread_workspace)?;
190 let live = match precheck_file_target(&root, &relative)? {
191 Some(_) => Some(read_confined_bytes(&open_confined_file(
192 &root, &relative, false,
193 )?)?),
194 None => None,
195 };
196 let Some(wanted) = artifact.revision.as_deref() else {
197 // No recorded revision (an older receipt): all that can be served
198 // is the live file, and whether it is the turn's version is unknown.
199 return live
200 .map(|read| ("workspace", None, read))
201 .ok_or_else(|| ApiError::gone("the file is no longer in the workspace"));
202 };
203 if let Some(read) = live.as_ref()
204 && read.revision == wanted
205 {
206 return Ok(("workspace", Some(true), live.expect("checked above")));
207 }
208 if let Some(post) = workspace.and_then(|w| w.post_turn_snapshot_id.as_deref())
209 && let Some(bytes) = read_snapshot_blob(&root, post, &artifact.path)?
210 {
211 let read = ConfinedFileBytes {
212 size: bytes.len() as u64,
213 revision: super::workspace::content_revision(&bytes),
214 modified: None,
215 bytes,
216 };
217 if read.revision == wanted {
218 return Ok(("snapshot", Some(false), read));
219 }
220 }
221 Err(ApiError::conflict(match live {
222 Some(read) => format!(
223 "this turn's revision is no longer in the workspace or the snapshot store; current revision is {}",
224 read.revision
225 ),
226 None => "this turn's revision is no longer in the workspace or the snapshot store; the file is gone".to_string(),
227 }))
228 }
229
230 fn read_snapshot_blob(
231 workspace: &FsPath,
232 snapshot_id: &str,
233 path: &str,
234 ) -> Result<Option<Vec<u8>>, ApiError> {
235 let Ok(id) = crate::snapshot::SnapshotId::parse(snapshot_id) else {
236 return Ok(None);
237 };
238 let Some(repo) = crate::snapshot::SnapshotRepo::open_existing(workspace)
239 .map_err(|error| ApiError::internal(format!("snapshot repo unavailable: {error}")))?
240 else {
241 return Ok(None);
242 };
243 match repo.read_blob(&id, path, FILE_SERVE_MAX_BYTES) {
244 Ok(bytes) => Ok(bytes),
245 Err(error) if error.kind() == std::io::ErrorKind::FileTooLarge => {
246 Err(ApiError::payload_too_large(format!(
247 "file is larger than the {FILE_SERVE_MAX_BYTES}-byte serving limit"
248 )))
249 }
250 // A pruned snapshot is not an error: the revision is just gone.
251 Err(error) => {
252 tracing::debug!(%error, "snapshot blob read failed");
253 Ok(None)
254 }
255 }
256 }
257
257 lines RUST