From c20af332c1299bcf43f24237275c3d9acef2298d Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Mon, 3 Aug 2026 02:04:08 -0400 Subject: [PATCH] Discard stale legacy materialization journals --- src/app/materialize.rs | 7 +++++++ src/storage.rs | 1 + src/storage/materializations.rs | 21 ++++++++++++++------- src/storage/tests.rs | 2 ++ 4 files changed, 24 insertions(+), 7 deletions(-) diff --git a/src/app/materialize.rs b/src/app/materialize.rs index aca9664..ccb0708 100644 --- a/src/app/materialize.rs +++ b/src/app/materialize.rs @@ -134,6 +134,13 @@ pub(super) async fn recover_pending_materialization( else { return Ok(()); }; + if pending.is_legacy_single_entry { + tracing::warn!("Discarding a stale single-entry materialization journal"); + context + .state_store + .clear_pending_materialization(context.folder.id)?; + return Ok(()); + } if pending.resulting_manifest.folder_id != context.folder.id { anyhow::bail!("pending materialization belongs to another folder"); } diff --git a/src/storage.rs b/src/storage.rs index ac0da30..dc03a29 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -62,6 +62,7 @@ pub struct FolderConfig { pub struct PendingMaterialization { pub entries: Vec, pub resulting_manifest: Manifest, + pub is_legacy_single_entry: bool, } impl StateStore { diff --git a/src/storage/materializations.rs b/src/storage/materializations.rs index 2315f3b..fcb78bc 100644 --- a/src/storage/materializations.rs +++ b/src/storage/materializations.rs @@ -13,10 +13,10 @@ enum StoredMaterializationEntries { Batch(Vec), } -fn parse_pending_entries(entries: &str) -> anyhow::Result> { +fn parse_pending_entries(entries: &str) -> anyhow::Result<(Vec, bool)> { match serde_json::from_str(entries)? { - StoredMaterializationEntries::Single(entry) => Ok(vec![entry]), - StoredMaterializationEntries::Batch(entries) => Ok(entries), + StoredMaterializationEntries::Single(entry) => Ok((vec![entry], true)), + StoredMaterializationEntries::Batch(entries) => Ok((entries, false)), } } @@ -42,10 +42,17 @@ impl StateStore { "SELECT entry, resulting_manifest FROM pending_materializations WHERE folder_id = ?1", params![folder_id.to_string()], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)), - ).optional()?.map(|(entries, manifest)| Ok(PendingMaterialization { - entries: parse_pending_entries(&entries)?, - resulting_manifest: serde_json::from_str(&manifest)?, - })).transpose() + ) + .optional()? + .map(|(entries, manifest)| { + let (entries, is_legacy_single_entry) = parse_pending_entries(&entries)?; + Ok(PendingMaterialization { + entries, + resulting_manifest: serde_json::from_str(&manifest)?, + is_legacy_single_entry, + }) + }) + .transpose() } pub fn clear_pending_materialization(&self, folder_id: FolderId) -> anyhow::Result<()> { diff --git a/src/storage/tests.rs b/src/storage/tests.rs index 2ba090d..8cdb642 100644 --- a/src/storage/tests.rs +++ b/src/storage/tests.rs @@ -212,6 +212,7 @@ fn saves_and_clears_pending_materialization() -> anyhow::Result<()> { .expect("journal exists"); assert_eq!(pending.entries, entries); assert_eq!(pending.resulting_manifest, resulting_manifest); + assert!(!pending.is_legacy_single_entry); store.connection.execute( "UPDATE pending_materializations SET entry = ?1 WHERE folder_id = ?2", @@ -221,6 +222,7 @@ fn saves_and_clears_pending_materialization() -> anyhow::Result<()> { .pending_materialization(folder.id)? .expect("legacy journal exists"); assert_eq!(pending.entries, vec![entries[0].clone()]); + assert!(pending.is_legacy_single_entry); store.clear_pending_materialization(folder.id)?; assert!(store.pending_materialization(folder.id)?.is_none()); -- 2.51.2