1use crate::ARCHIVED_SESSIONS_SUBDIR;
2use crate::SESSIONS_SUBDIR;
3use crate::compression;
4use crate::list::parse_timestamp_uuid_from_filename;
5use crate::recorder::RolloutRecorder;
6use crate::state_db::normalize_cwd_for_state_db;
7use chrono::DateTime;
8use chrono::NaiveDateTime;
9use chrono::Timelike;
10use chrono::Utc;
11use codex_protocol::ThreadId;
12use codex_protocol::protocol::AskForApproval;
13use codex_protocol::protocol::RolloutItem;
14use codex_protocol::protocol::SandboxPolicy;
15use codex_protocol::protocol::SessionMetaLine;
16use codex_protocol::protocol::SessionSource;
17use codex_protocol::protocol::ThreadHistoryMode;
18use codex_state::BackfillState;
19use codex_state::BackfillStats;
20use codex_state::BackfillStatus;
21use codex_state::DB_ERROR_METRIC;
22use codex_state::DB_METRIC_BACKFILL;
23use codex_state::DB_METRIC_BACKFILL_DURATION_MS;
24use codex_state::ExtractionOutcome;
25use codex_state::ThreadMetadataBuilder;
26use codex_state::apply_rollout_item;
27use std::path::Path;
28use std::path::PathBuf;
29use tracing::info;
30use tracing::warn;
31
32const BACKFILL_BATCH_SIZE: usize = 200;
33#[cfg(not(test))]
34const BACKFILL_LEASE_SECONDS: i64 = 900;
35#[cfg(test)]
36const BACKFILL_LEASE_SECONDS: i64 = 1;
37
38pub(crate) fn builder_from_session_meta(
39 session_meta: &SessionMetaLine,
40 rollout_path: &Path,
41) -> Option<ThreadMetadataBuilder> {
42 let created_at = parse_timestamp_to_utc(session_meta.meta.timestamp.as_str())?;
43 let mut builder = ThreadMetadataBuilder::new(
44 session_meta.meta.id,
45 rollout_path.to_path_buf(),
46 created_at,
47 session_meta.meta.source.clone(),
48 );
49 builder.history_mode = session_meta.meta.history_mode;
50 builder.model_provider = session_meta.meta.model_provider.clone();
51 builder.agent_nickname = session_meta.meta.agent_nickname.clone();
52 builder.agent_role = session_meta.meta.agent_role.clone();
53 builder.agent_path = session_meta.meta.agent_path.clone();
54 builder.cwd = session_meta.meta.cwd.clone();
55 builder.cli_version = Some(session_meta.meta.cli_version.clone());
56 builder.sandbox_policy = SandboxPolicy::new_read_only_policy();
57 builder.approval_mode = AskForApproval::OnRequest;
58 if let Some(git) = session_meta.git.as_ref() {
59 builder.git_sha = git.commit_hash.as_ref().map(|sha| sha.0.clone());
60 builder.git_branch = git.branch.clone();
61 builder.git_origin_url = git.repository_url.clone();
62 }
63 Some(builder)
64}
65
66pub fn builder_from_items(
67 items: &[RolloutItem],
68 rollout_path: &Path,
69) -> Option<ThreadMetadataBuilder> {
70 if let Some(session_meta) = items.iter().find_map(|item| match item {
71 RolloutItem::SessionMeta(meta_line) => Some(meta_line),
72 RolloutItem::ResponseItem(_)
73 | RolloutItem::InterAgentCommunication(_)
74 | RolloutItem::InterAgentCommunicationMetadata { .. }
75 | RolloutItem::Compacted(_)
76 | RolloutItem::TurnContext(_)
77 | RolloutItem::WorldState(_)
78 | RolloutItem::EventMsg(_) => None,
79 }) && let Some(builder) = builder_from_session_meta(session_meta, rollout_path)
80 {
81 return Some(builder);
82 }
83
84 let file_name = rollout_path.file_name()?.to_str()?;
85 let file_name = compression::parse_rollout_file_name(file_name)?;
86 let (created_ts, uuid) = parse_timestamp_uuid_from_filename(file_name)?;
87 let created_at =
88 DateTime::<Utc>::from_timestamp(created_ts.unix_timestamp(), 0)?.with_nanosecond(0)?;
89 let id = ThreadId::from_string(&uuid.to_string()).ok()?;
90 Some(ThreadMetadataBuilder::new(
91 id,
92 rollout_path.to_path_buf(),
93 created_at,
94 SessionSource::default(),
95 ))
96}
97
98pub async fn extract_metadata_from_rollout(
99 rollout_path: &Path,
100 default_provider: &str,
101) -> anyhow::Result<ExtractionOutcome> {
102 let (items, _thread_id, parse_errors) =
103 RolloutRecorder::load_rollout_items(rollout_path).await?;
104 if items.is_empty() {
105 return Err(anyhow::anyhow!(
106 "empty session file: {}",
107 rollout_path.display()
108 ));
109 }
110 let builder = builder_from_items(items.as_slice(), rollout_path).ok_or_else(|| {
111 anyhow::anyhow!(
112 "rollout missing metadata builder: {}",
113 rollout_path.display()
114 )
115 })?;
116 let mut metadata = builder.build(default_provider);
117 for item in &items {
118 apply_rollout_item(&mut metadata, item, default_provider);
119 }
120 if let Some(updated_at) = file_modified_time_utc(rollout_path).await {
121 metadata.updated_at = updated_at;
122 metadata.recency_at = updated_at;
123 }
124 Ok(ExtractionOutcome {
125 metadata,
126 memory_mode: items.iter().rev().find_map(|item| match item {
127 RolloutItem::SessionMeta(meta_line) => meta_line.meta.memory_mode.clone(),
128 RolloutItem::ResponseItem(_)
129 | RolloutItem::InterAgentCommunication(_)
130 | RolloutItem::InterAgentCommunicationMetadata { .. }
131 | RolloutItem::Compacted(_)
132 | RolloutItem::TurnContext(_)
133 | RolloutItem::WorldState(_)
134 | RolloutItem::EventMsg(_) => None,
135 }),
136 parse_errors,
137 })
138}
139
140pub(crate) async fn backfill_sessions(
141 runtime: &codex_state::StateRuntime,
142 codex_home: &Path,
143 default_provider: &str,
144) {
145 backfill_sessions_with_lease(
146 runtime,
147 codex_home,
148 default_provider,
149 BACKFILL_LEASE_SECONDS,
150 )
151 .await;
152}
153
154pub(crate) async fn backfill_sessions_with_lease(
155 runtime: &codex_state::StateRuntime,
156 codex_home: &Path,
157 default_provider: &str,
158 backfill_lease_seconds: i64,
159) {
160 let metric_client = codex_otel::global();
161 let timer = metric_client
162 .as_ref()
163 .and_then(|otel| otel.start_timer(DB_METRIC_BACKFILL_DURATION_MS, &[]).ok());
164 let backfill_state = match runtime.get_backfill_state().await {
165 Ok(state) => state,
166 Err(err) => {
167 warn!(
168 "failed to read backfill state at {}: {err}",
169 codex_home.display()
170 );
171 BackfillState::default()
172 }
173 };
174 if backfill_state.status == BackfillStatus::Complete {
175 return;
176 }
177 let claimed = match runtime.try_claim_backfill(backfill_lease_seconds).await {
178 Ok(claimed) => claimed,
179 Err(err) => {
180 warn!(
181 "failed to claim backfill worker at {}: {err}",
182 codex_home.display()
183 );
184 return;
185 }
186 };
187 if !claimed {
188 info!(
189 "state db backfill already running at {}; skipping duplicate worker",
190 codex_home.display()
191 );
192 return;
193 }
194 let mut backfill_state = match runtime.get_backfill_state().await {
195 Ok(state) => state,
196 Err(err) => {
197 warn!(
198 "failed to read claimed backfill state at {}: {err}",
199 codex_home.display()
200 );
201 BackfillState {
202 status: BackfillStatus::Running,
203 ..Default::default()
204 }
205 }
206 };
207 if backfill_state.status != BackfillStatus::Running {
208 if let Err(err) = runtime.mark_backfill_running().await {
209 warn!(
210 "failed to mark backfill running at {}: {err}",
211 codex_home.display()
212 );
213 } else {
214 backfill_state.status = BackfillStatus::Running;
215 }
216 }
217
218 let sessions_root = codex_home.join(SESSIONS_SUBDIR);
219 let archived_root = codex_home.join(ARCHIVED_SESSIONS_SUBDIR);
220 let mut rollout_paths: Vec<BackfillRolloutPath> = Vec::new();
221 for (root, archived) in [(sessions_root, false), (archived_root, true)] {
222 if !tokio::fs::try_exists(&root).await.unwrap_or(false) {
223 continue;
224 }
225 match collect_rollout_paths(&root).await {
226 Ok(paths) => {
227 rollout_paths.extend(paths.into_iter().map(|path| BackfillRolloutPath {
228 watermark: backfill_watermark_for_path(codex_home, &path),
229 path,
230 archived,
231 }));
232 }
233 Err(err) => {
234 warn!(
235 "failed to collect rollout paths under {}: {err}",
236 root.display()
237 );
238 }
239 }
240 }
241 rollout_paths.sort_by(|a, b| a.watermark.cmp(&b.watermark));
242 if let Some(last_watermark) = backfill_state.last_watermark.as_deref() {
243 rollout_paths.retain(|entry| entry.watermark.as_str() > last_watermark);
244 }
245
246 let mut stats = BackfillStats {
247 scanned: 0,
248 upserted: 0,
249 failed: 0,
250 };
251 let mut last_watermark = backfill_state.last_watermark.clone();
252 for batch in rollout_paths.chunks(BACKFILL_BATCH_SIZE) {
253 for rollout in batch {
254 stats.scanned = stats.scanned.saturating_add(1);
255 match extract_metadata_from_rollout(&rollout.path, default_provider).await {
256 Ok(outcome) => {
257 if outcome.parse_errors > 0
258 && let Some(ref metric_client) = metric_client
259 {
260 let _ = metric_client.counter(
261 DB_ERROR_METRIC,
262 outcome.parse_errors as i64,
263 &[("stage", "backfill_sessions")],
264 );
265 }
266 let mut metadata = outcome.metadata;
267 metadata.cwd = normalize_cwd_for_state_db(&metadata.cwd);
268 let memory_mode = outcome.memory_mode.unwrap_or_else(|| "enabled".to_string());
269 let existing_metadata = runtime.get_thread(metadata.id).await.ok().flatten();
270 let restore_memory_mode_from_rollout = existing_metadata.is_none()
273 || matches!(metadata.history_mode, ThreadHistoryMode::Legacy);
274 if let Some(existing_metadata) = existing_metadata.as_ref() {
275 metadata.prefer_existing_git_info(existing_metadata);
276 metadata.prefer_existing_explicit_title(existing_metadata);
277 }
278 if rollout.archived && metadata.archived_at.is_none() {
279 let fallback_archived_at = metadata.updated_at;
280 metadata.archived_at = file_modified_time_utc(&rollout.path)
281 .await
282 .or(Some(fallback_archived_at));
283 }
284 if let Err(err) = runtime.upsert_thread(&metadata).await {
285 stats.failed = stats.failed.saturating_add(1);
286 warn!("failed to upsert rollout {}: {err}", rollout.path.display());
287 } else {
288 if restore_memory_mode_from_rollout
289 && let Err(err) = runtime
290 .set_thread_memory_mode(metadata.id, memory_mode.as_str())
291 .await
292 {
293 stats.failed = stats.failed.saturating_add(1);
294 warn!(
295 "failed to restore memory mode for {}: {err}",
296 rollout.path.display()
297 );
298 continue;
299 }
300 stats.upserted = stats.upserted.saturating_add(1);
301 }
302 }
303 Err(err) => {
304 stats.failed = stats.failed.saturating_add(1);
305 warn!(
306 "failed to extract rollout {}: {err}",
307 rollout.path.display()
308 );
309 }
310 }
311 }
312
313 if let Some(last_entry) = batch.last() {
314 if let Err(err) = runtime
315 .checkpoint_backfill(last_entry.watermark.as_str())
316 .await
317 {
318 warn!(
319 "failed to checkpoint backfill at {}: {err}",
320 codex_home.display()
321 );
322 } else {
323 last_watermark = Some(last_entry.watermark.clone());
324 }
325 }
326 }
327 if let Err(err) = runtime
328 .mark_backfill_complete(last_watermark.as_deref())
329 .await
330 {
331 warn!(
332 "failed to mark backfill complete at {}: {err}",
333 codex_home.display()
334 );
335 }
336
337 info!(
338 "state db backfill scanned={}, upserted={}, failed={}",
339 stats.scanned, stats.upserted, stats.failed
340 );
341 if let Some(metric_client) = metric_client {
342 let _ = metric_client.counter(
343 DB_METRIC_BACKFILL,
344 stats.upserted as i64,
345 &[("status", "upserted")],
346 );
347 let _ = metric_client.counter(
348 DB_METRIC_BACKFILL,
349 stats.failed as i64,
350 &[("status", "failed")],
351 );
352 }
353 if let Some(timer) = timer.as_ref() {
354 let status = if stats.failed == 0 {
355 "success"
356 } else if stats.upserted == 0 {
357 "failed"
358 } else {
359 "partial_failure"
360 };
361 let _ = timer.record(&[("status", status)]);
362 }
363}
364
365#[derive(Debug, Clone)]
366struct BackfillRolloutPath {
367 watermark: String,
368 path: PathBuf,
369 archived: bool,
370}
371
372fn backfill_watermark_for_path(codex_home: &Path, path: &Path) -> String {
373 path.strip_prefix(codex_home)
374 .unwrap_or(path)
375 .to_string_lossy()
376 .replace('\\', "/")
377}
378
379async fn file_modified_time_utc(path: &Path) -> Option<DateTime<Utc>> {
380 let modified = compression::file_modified_time(path).await.ok()??;
381 DateTime::<Utc>::from_timestamp(modified.unix_timestamp(), modified.nanosecond())
382}
383
384fn parse_timestamp_to_utc(ts: &str) -> Option<DateTime<Utc>> {
385 const FILENAME_TS_FORMAT: &str = "%Y-%m-%dT%H-%M-%S";
386 if let Ok(naive) = NaiveDateTime::parse_from_str(ts, FILENAME_TS_FORMAT) {
387 let dt = DateTime::<Utc>::from_naive_utc_and_offset(naive, Utc);
388 return dt.with_nanosecond(0);
389 }
390 if let Ok(dt) = DateTime::parse_from_rfc3339(ts) {
391 return Some(dt.with_timezone(&Utc));
392 }
393 None
394}
395
396async fn collect_rollout_paths(root: &Path) -> std::io::Result<Vec<PathBuf>> {
397 let mut stack = vec![root.to_path_buf()];
398 let mut paths = Vec::new();
399 while let Some(dir) = stack.pop() {
400 let mut read_dir = match tokio::fs::read_dir(&dir).await {
401 Ok(read_dir) => read_dir,
402 Err(err) => {
403 warn!("failed to read directory {}: {err}", dir.display());
404 continue;
405 }
406 };
407 loop {
408 let next_entry = match read_dir.next_entry().await {
409 Ok(next_entry) => next_entry,
410 Err(err) => {
411 warn!(
412 "failed to read directory entry under {}: {err}",
413 dir.display()
414 );
415 continue;
416 }
417 };
418 let Some(entry) = next_entry else {
419 break;
420 };
421 let path = entry.path();
422 let file_type = match entry.file_type().await {
423 Ok(file_type) => file_type,
424 Err(err) => {
425 warn!("failed to read file type for {}: {err}", path.display());
426 continue;
427 }
428 };
429 if file_type.is_dir() {
430 stack.push(path);
431 continue;
432 }
433 if !file_type.is_file() {
434 continue;
435 }
436 if let Some(rollout_file) = compression::RolloutFile::from_path(path) {
437 paths.push(rollout_file.into_path());
438 }
439 }
440 }
441 Ok(paths)
442}
443
444#[cfg(test)]
445#[path = "metadata_tests.rs"]
446mod tests;