Skip to main content

codex_rollout/
metadata.rs

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                    // Paginated metadata updates are SQLite-only. Use the rollout mode to seed a
271                    // missing row, then keep the value from SQLite.
272                    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;