feat(server): improve doc gc (#15329)
#### PR Dependency Tree * **PR #15329** 👈 This tree was auto-generated by [Charcoal](https://github.com/danerwilliams/charcoal) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Improvements** * Enhanced reliability of document reconciliation/projection by improving failure tracking and resumability after partial runs. * Upgraded checkpoint persistence to include additional state (including parser version and failure counts) so interrupted processing can resume accurately. * Improved recovery behavior so repeated rebuilds and parser-upgrade scenarios correctly complete and reset failure indicators when appropriate. * **Tests** * Added integration coverage for checkpoint failure/resume semantics, including partial limits and upgrade-style recovery behavior. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
@@ -14,6 +14,12 @@ struct ExtractedRef {
|
|||||||
flavour: String,
|
flavour: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Default)]
|
||||||
|
struct ProjectionState {
|
||||||
|
cursor: Option<String>,
|
||||||
|
failed_docs: i64,
|
||||||
|
}
|
||||||
|
|
||||||
async fn load_workspace_doc_ids(pool: &PgPool, workspace_id: &str) -> RuntimeResult<Vec<String>> {
|
async fn load_workspace_doc_ids(pool: &PgPool, workspace_id: &str) -> RuntimeResult<Vec<String>> {
|
||||||
let mut ids = load_workspace_live_doc_ids(pool, workspace_id).await?;
|
let mut ids = load_workspace_live_doc_ids(pool, workspace_id).await?;
|
||||||
let retained = sqlx::query_scalar::<_, String>(
|
let retained = sqlx::query_scalar::<_, String>(
|
||||||
@@ -34,15 +40,16 @@ async fn upsert_projection_checkpoint(
|
|||||||
pool: &PgPool,
|
pool: &PgPool,
|
||||||
workspace_id: &str,
|
workspace_id: &str,
|
||||||
result: &RuntimeDocBlobRefsResult,
|
result: &RuntimeDocBlobRefsResult,
|
||||||
|
failed_docs: i64,
|
||||||
) -> RuntimeResult<()> {
|
) -> RuntimeResult<()> {
|
||||||
let completed = result.next_cursor.is_none();
|
let status = if result.next_cursor.is_some() {
|
||||||
let status = if completed && result.failed_docs == 0 {
|
"running"
|
||||||
"completed"
|
} else if failed_docs > 0 {
|
||||||
} else if result.failed_docs > 0 {
|
|
||||||
"failed"
|
"failed"
|
||||||
} else {
|
} else {
|
||||||
"running"
|
"completed"
|
||||||
};
|
};
|
||||||
|
let completed = status == "completed";
|
||||||
sqlx::query(
|
sqlx::query(
|
||||||
r#"
|
r#"
|
||||||
INSERT INTO storage_reconciliation_checkpoints
|
INSERT INTO storage_reconciliation_checkpoints
|
||||||
@@ -59,9 +66,10 @@ async fn upsert_projection_checkpoint(
|
|||||||
.bind(workspace_id)
|
.bind(workspace_id)
|
||||||
.bind(status)
|
.bind(status)
|
||||||
.bind(serde_json::json!({ "lastDocId": result.next_cursor }))
|
.bind(serde_json::json!({ "lastDocId": result.next_cursor }))
|
||||||
.bind(completed && result.failed_docs == 0)
|
.bind(completed)
|
||||||
.bind(serde_json::json!({
|
.bind(serde_json::json!({
|
||||||
"parserVersion": PARSER_VERSION,
|
"parserVersion": PARSER_VERSION,
|
||||||
|
"failedDocs": failed_docs,
|
||||||
}))
|
}))
|
||||||
.execute(pool)
|
.execute(pool)
|
||||||
.await
|
.await
|
||||||
@@ -94,23 +102,39 @@ async fn upsert_projection_failure_checkpoint(pool: &PgPool, workspace_id: &str,
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn load_projection_cursor(pool: &PgPool, workspace_id: &str) -> RuntimeResult<Option<String>> {
|
async fn load_projection_state(pool: &PgPool, workspace_id: &str) -> RuntimeResult<ProjectionState> {
|
||||||
let checkpoint = sqlx::query_as::<_, (String, serde_json::Value)>(
|
let checkpoint = sqlx::query_as::<_, (String, serde_json::Value, serde_json::Value)>(
|
||||||
"SELECT status, cursor FROM storage_reconciliation_checkpoints WHERE kind = 'doc_blob_refs' AND scope = $1",
|
"SELECT status, cursor, metadata FROM storage_reconciliation_checkpoints WHERE kind = 'doc_blob_refs' AND scope = \
|
||||||
|
$1",
|
||||||
)
|
)
|
||||||
.bind(workspace_id)
|
.bind(workspace_id)
|
||||||
.fetch_optional(pool)
|
.fetch_optional(pool)
|
||||||
.await
|
.await
|
||||||
.map_err(|err| RuntimeError::database("Doc blob refs checkpoint load failed", err))?;
|
.map_err(|err| RuntimeError::database("Doc blob refs checkpoint load failed", err))?;
|
||||||
Ok(checkpoint.and_then(|(status, cursor)| {
|
let Some((status, cursor, metadata)) = checkpoint else {
|
||||||
if status != "running" {
|
return Ok(ProjectionState::default());
|
||||||
return None;
|
};
|
||||||
|
if status != "running" && status != "failed" {
|
||||||
|
return Ok(ProjectionState::default());
|
||||||
}
|
}
|
||||||
cursor
|
if metadata.get("parserVersion").and_then(serde_json::Value::as_i64) != Some(i64::from(PARSER_VERSION)) {
|
||||||
|
return Ok(ProjectionState::default());
|
||||||
|
}
|
||||||
|
let cursor = cursor
|
||||||
.get("lastDocId")
|
.get("lastDocId")
|
||||||
.and_then(|value| value.as_str())
|
.and_then(|value| value.as_str())
|
||||||
.map(ToString::to_string)
|
.map(ToString::to_string);
|
||||||
}))
|
let Some(cursor) = cursor else {
|
||||||
|
return Ok(ProjectionState::default());
|
||||||
|
};
|
||||||
|
let failed_docs = metadata
|
||||||
|
.get("failedDocs")
|
||||||
|
.and_then(serde_json::Value::as_i64)
|
||||||
|
.unwrap_or(i64::from(status == "failed"));
|
||||||
|
Ok(ProjectionState {
|
||||||
|
cursor: Some(cursor),
|
||||||
|
failed_docs,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn purge_removed_doc_refs(pool: &PgPool, workspace_id: &str, current_doc_ids: &[String]) -> RuntimeResult<i64> {
|
async fn purge_removed_doc_refs(pool: &PgPool, workspace_id: &str, current_doc_ids: &[String]) -> RuntimeResult<i64> {
|
||||||
@@ -351,11 +375,11 @@ impl StorageRuntime {
|
|||||||
return Err(err.into());
|
return Err(err.into());
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
let cursor = load_projection_cursor(&pool, &workspace_id).await?;
|
let state = load_projection_state(&pool, &workspace_id).await?;
|
||||||
let current_doc_ids = doc_ids.clone();
|
let current_doc_ids = doc_ids.clone();
|
||||||
let doc_ids = doc_ids
|
let doc_ids = doc_ids
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.filter(|doc_id| cursor.as_ref().is_none_or(|cursor| doc_id > cursor))
|
.filter(|doc_id| state.cursor.as_ref().is_none_or(|cursor| doc_id > cursor))
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
let has_more = doc_ids.len() > limit as usize;
|
let has_more = doc_ids.len() > limit as usize;
|
||||||
let mut total = RuntimeDocBlobRefsResult {
|
let mut total = RuntimeDocBlobRefsResult {
|
||||||
@@ -377,13 +401,14 @@ impl StorageRuntime {
|
|||||||
total.refs_deleted += result.refs_deleted;
|
total.refs_deleted += result.refs_deleted;
|
||||||
total.failed_docs += result.failed_docs;
|
total.failed_docs += result.failed_docs;
|
||||||
}
|
}
|
||||||
|
let failed_docs = state.failed_docs + total.failed_docs;
|
||||||
if has_more {
|
if has_more {
|
||||||
total.next_cursor = last_doc_id;
|
total.next_cursor = last_doc_id;
|
||||||
} else if total.failed_docs == 0 {
|
} else if failed_docs == 0 {
|
||||||
total.refs_deleted += purge_removed_doc_refs(&pool, &workspace_id, ¤t_doc_ids).await?;
|
total.refs_deleted += purge_removed_doc_refs(&pool, &workspace_id, ¤t_doc_ids).await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
upsert_projection_checkpoint(&pool, &workspace_id, &total).await?;
|
upsert_projection_checkpoint(&pool, &workspace_id, &total, failed_docs).await?;
|
||||||
|
|
||||||
Ok(total)
|
Ok(total)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1026,6 +1026,91 @@ mod tests {
|
|||||||
1
|
1
|
||||||
);
|
);
|
||||||
|
|
||||||
|
sqlx::query("UPDATE snapshots SET blob = $2 WHERE workspace_id = $1 AND guid = 'live-doc'")
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.bind(vec![0xff_u8])
|
||||||
|
.execute(&pool)
|
||||||
|
.await?;
|
||||||
|
let partial = runtime
|
||||||
|
.rebuild_workspace_doc_blob_refs(workspace_id.clone(), 1)
|
||||||
|
.await
|
||||||
|
.map_err(|err| anyhow::anyhow!(err.to_string()))?;
|
||||||
|
assert_eq!(
|
||||||
|
(partial.failed_docs, partial.next_cursor.as_deref()),
|
||||||
|
(1, Some("live-doc"))
|
||||||
|
);
|
||||||
|
let partial_checkpoint = sqlx::query(
|
||||||
|
"SELECT status, metadata FROM storage_reconciliation_checkpoints WHERE kind = 'doc_blob_refs' AND scope = $1",
|
||||||
|
)
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.fetch_one(&pool)
|
||||||
|
.await?;
|
||||||
|
assert_eq!(partial_checkpoint.get::<String, _>("status"), "running");
|
||||||
|
assert_eq!(partial_checkpoint.get::<Value, _>("metadata")["failedDocs"], 1);
|
||||||
|
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE storage_reconciliation_checkpoints SET status = 'failed', metadata = '{\"parserVersion\":1}' WHERE kind \
|
||||||
|
= 'doc_blob_refs' AND scope = $1",
|
||||||
|
)
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.execute(&pool)
|
||||||
|
.await?;
|
||||||
|
let resumed = runtime
|
||||||
|
.rebuild_workspace_doc_blob_refs(workspace_id.clone(), 1)
|
||||||
|
.await
|
||||||
|
.map_err(|err| anyhow::anyhow!(err.to_string()))?;
|
||||||
|
assert_eq!((resumed.failed_docs, resumed.next_cursor), (0, None));
|
||||||
|
let resumed_checkpoint = sqlx::query(
|
||||||
|
"SELECT status, metadata FROM storage_reconciliation_checkpoints WHERE kind = 'doc_blob_refs' AND scope = $1",
|
||||||
|
)
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.fetch_one(&pool)
|
||||||
|
.await?;
|
||||||
|
assert_eq!(resumed_checkpoint.get::<String, _>("status"), "failed");
|
||||||
|
assert_eq!(resumed_checkpoint.get::<Value, _>("metadata")["failedDocs"], 1);
|
||||||
|
|
||||||
|
sqlx::query("UPDATE snapshots SET blob = $2 WHERE workspace_id = $1 AND guid = 'live-doc'")
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.bind(affine_doc_loader::build_full_doc("Live", "", "live-doc")?)
|
||||||
|
.execute(&pool)
|
||||||
|
.await?;
|
||||||
|
let recovered_projection = runtime
|
||||||
|
.rebuild_workspace_doc_blob_refs(workspace_id.clone(), 100)
|
||||||
|
.await
|
||||||
|
.map_err(|err| anyhow::anyhow!(err.to_string()))?;
|
||||||
|
assert_eq!(recovered_projection.failed_docs, 0);
|
||||||
|
assert_eq!(
|
||||||
|
sqlx::query_scalar::<_, String>(
|
||||||
|
"SELECT status FROM storage_reconciliation_checkpoints WHERE kind = 'doc_blob_refs' AND scope = $1",
|
||||||
|
)
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.fetch_one(&pool)
|
||||||
|
.await?,
|
||||||
|
"completed"
|
||||||
|
);
|
||||||
|
|
||||||
|
sqlx::query(
|
||||||
|
"UPDATE storage_reconciliation_checkpoints SET status = 'failed', cursor = '{\"lastDocId\":\"live-doc\"}', \
|
||||||
|
metadata = '{\"parserVersion\":0,\"failedDocs\":99}' WHERE kind = 'doc_blob_refs' AND scope = $1",
|
||||||
|
)
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.execute(&pool)
|
||||||
|
.await?;
|
||||||
|
let parser_upgrade = runtime
|
||||||
|
.rebuild_workspace_doc_blob_refs(workspace_id.clone(), 100)
|
||||||
|
.await
|
||||||
|
.map_err(|err| anyhow::anyhow!(err.to_string()))?;
|
||||||
|
assert_eq!((parser_upgrade.scanned_docs, parser_upgrade.failed_docs), (2, 0));
|
||||||
|
assert_eq!(
|
||||||
|
sqlx::query_scalar::<_, String>(
|
||||||
|
"SELECT status FROM storage_reconciliation_checkpoints WHERE kind = 'doc_blob_refs' AND scope = $1",
|
||||||
|
)
|
||||||
|
.bind(&workspace_id)
|
||||||
|
.fetch_one(&pool)
|
||||||
|
.await?,
|
||||||
|
"completed"
|
||||||
|
);
|
||||||
|
|
||||||
let second = reconcile_workspace(&runtime, &workspace_id).await?;
|
let second = reconcile_workspace(&runtime, &workspace_id).await?;
|
||||||
assert_eq!((second.marked, second.reset, second.recovered), (0, 0, 0));
|
assert_eq!((second.marked, second.reset, second.recovered), (0, 0, 0));
|
||||||
let unchanged_missing_since = sqlx::query_scalar::<_, DateTime<Utc>>(
|
let unchanged_missing_since = sqlx::query_scalar::<_, DateTime<Utc>>(
|
||||||
|
|||||||
Reference in New Issue
Block a user