feat(server): blob reconciliation (#15165)

#### PR Dependency Tree


* **PR #15165** 👈

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

* **New Features**
* Added automated backend maintenance for missing blob metadata
backfill, document-to-blob reference rebuilding, and unreferenced blob
cleanup planning/execution.
* Introduced scheduled batch processing (workspace-paged) and paginated
object-storage listing.
* **Bug Fixes**
* Improved reliability of object-storage reads by treating expected “not
found” results as non-errors.
* Strengthened blob/expired cleanup flows with runtime-driven batching
and reduced coupling to metadata synchronization.
* **Tests**
* Expanded unit and e2e coverage for partial blob metadata and updated
runtime/job cleanup test assertions.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
DarkSky
2026-06-29 00:02:38 +08:00
committed by GitHub
parent 4a7c931eca
commit 0a422aa158
42 changed files with 2494 additions and 264 deletions
@@ -0,0 +1,647 @@
use chrono::{DateTime, Duration, Utc};
use napi::Result;
use sqlx::{FromRow, PgPool};
use super::{
BackendRuntime,
error::napi_error,
types::{RuntimeBlobCleanupExecuteResult, RuntimeBlobCleanupPlanResult},
};
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn blob_cleanup_plan_result_keeps_run_id_for_execute() {
let result = RuntimeBlobCleanupPlanResult {
run_id: Some("00000000-0000-0000-0000-000000000000".to_string()),
scanned_blobs: 1,
candidates_marked: 1,
protected_by_doc_refs: 0,
protected_by_metadata: 0,
protected_by_other_refs: 0,
next_cursor: None,
};
assert!(result.run_id.is_some());
assert_eq!(result.candidates_marked, 1);
}
#[test]
fn blob_cleanup_execute_result_tracks_skipped_and_failed_counts() {
let result = RuntimeBlobCleanupExecuteResult {
scanned_candidates: 3,
deleted_objects: 1,
deleted_metadata: 1,
skipped_still_referenced: 1,
failed: 1,
workspace_ids: vec!["workspace".to_string()],
};
assert_eq!(result.scanned_candidates, 3);
assert_eq!(
result.skipped_still_referenced + result.failed + result.deleted_objects,
3
);
}
}
#[derive(FromRow)]
struct BlobCandidateRow {
workspace_id: String,
key: String,
size: i32,
}
#[derive(FromRow)]
struct MarkedCandidateRow {
workspace_id: String,
blob_key: String,
}
fn push_workspace_once(workspace_ids: &mut Vec<String>, workspace_id: &str) {
if !workspace_ids.iter().any(|id| id == workspace_id) {
workspace_ids.push(workspace_id.to_string());
}
}
async fn checkpoint_completed(pool: &PgPool, kind: &str, scope: &str) -> Result<bool> {
sqlx::query_scalar::<_, bool>(
"SELECT EXISTS(SELECT 1 FROM blob_reconciliation_checkpoints WHERE kind = $1 AND scope = $2 AND status = \
'completed')",
)
.bind(kind)
.bind(scope)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup checkpoint check failed: {err}")))
}
async fn projection_is_stale(pool: &PgPool, workspace_id: &str) -> Result<bool> {
let checkpoint_fresh = checkpoint_completed(pool, "doc_blob_refs", workspace_id).await?;
let has_stale_rows = sqlx::query_scalar::<_, bool>(
"SELECT EXISTS(SELECT 1 FROM doc_blob_refs WHERE workspace_id = $1 AND status <> 'fresh')",
)
.bind(workspace_id)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup projection freshness check failed: {err}")))?;
Ok(!checkpoint_fresh || has_stale_rows)
}
async fn stale_projection_workspaces(pool: &PgPool, workspace_id: &str) -> Result<Vec<String>> {
if projection_is_stale(pool, workspace_id).await? {
Ok(vec![workspace_id.to_string()])
} else {
Ok(Vec::new())
}
}
async fn metadata_backfill_is_complete(pool: &PgPool, workspace_id: &str) -> Result<bool> {
checkpoint_completed(pool, "blob_metadata_backfill", workspace_id).await
}
async fn has_doc_ref(pool: &PgPool, workspace_id: &str, key: &str) -> Result<bool> {
sqlx::query_scalar::<_, bool>(
"SELECT EXISTS(SELECT 1 FROM doc_blob_refs WHERE workspace_id = $1 AND blob_key = $2 AND status = 'fresh')",
)
.bind(workspace_id)
.bind(key)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup doc ref check failed: {err}")))
}
async fn has_other_ref(pool: &PgPool, workspace_id: &str, key: &str) -> Result<bool> {
let required_ref = sqlx::query_scalar::<_, bool>(
r#"
SELECT EXISTS(SELECT 1 FROM workspaces WHERE id = $1 AND avatar_key = $2)
OR EXISTS(SELECT 1 FROM ai_transcript_tasks WHERE workspace_id = $1 AND blob_id = $2)
OR EXISTS(SELECT 1 FROM ai_jobs WHERE workspace_id = $1 AND blob_id = $2)
OR EXISTS(
SELECT 1
FROM ai_contexts c
JOIN ai_sessions_metadata s ON s.id = c.session_id
WHERE s.workspace_id = $1
AND jsonb_path_exists(
c.config::jsonb,
'$.** ? (@ == $blobKey)',
jsonb_build_object('blobKey', to_jsonb($2::text))
)
)
"#,
)
.bind(workspace_id)
.bind(key)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup protected ref check failed: {err}")))?;
if required_ref {
return Ok(true);
}
if table_exists(pool, "ai_workspace_files").await?
&& sqlx::query_scalar::<_, bool>(
"SELECT EXISTS(SELECT 1 FROM ai_workspace_files WHERE workspace_id = $1 AND blob_id = $2)",
)
.bind(workspace_id)
.bind(key)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup workspace file ref check failed: {err}")))?
{
return Ok(true);
}
if table_exists(pool, "ai_workspace_blob_embeddings").await?
&& sqlx::query_scalar::<_, bool>(
"SELECT EXISTS(SELECT 1 FROM ai_workspace_blob_embeddings WHERE workspace_id = $1 AND blob_id = $2)",
)
.bind(workspace_id)
.bind(key)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup workspace blob embedding ref check failed: {err}")))?
{
return Ok(true);
}
Ok(false)
}
async fn table_exists(pool: &PgPool, table: &str) -> Result<bool> {
sqlx::query_scalar::<_, bool>("SELECT to_regclass($1) IS NOT NULL")
.bind(format!("public.{table}"))
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup table existence check failed: {err}")))
}
async fn load_completed_blobs(
pool: &PgPool,
workspace_id: &str,
after_key: Option<&str>,
limit: i64,
) -> Result<Vec<BlobCandidateRow>> {
sqlx::query_as::<_, BlobCandidateRow>(
r#"
SELECT workspace_id, key, size
FROM blobs
WHERE workspace_id = $1
AND status = 'completed'
AND deleted_at IS NULL
AND ($2::text IS NULL OR key > $2)
ORDER BY key ASC
LIMIT $3
"#,
)
.bind(workspace_id)
.bind(after_key)
.bind(limit)
.fetch_all(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup load completed blobs failed: {err}")))
}
async fn load_plan_cursor(pool: &PgPool, workspace_id: &str) -> Result<Option<String>> {
let row = sqlx::query_as::<_, (String, serde_json::Value)>(
"SELECT status, cursor FROM blob_reconciliation_checkpoints WHERE kind = 'blob_cleanup_plan' AND scope = $1",
)
.bind(workspace_id)
.fetch_optional(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup plan checkpoint load failed: {err}")))?;
let Some((status, cursor)) = row else {
return Ok(None);
};
if status == "completed" {
return Ok(None);
}
Ok({
cursor
.get("lastBlobKey")
.and_then(|value| value.as_str())
.map(ToString::to_string)
})
}
async fn upsert_plan_checkpoint(
pool: &PgPool,
workspace_id: &str,
last_blob_key: Option<&str>,
completed: bool,
) -> Result<()> {
let status = if completed { "completed" } else { "running" };
sqlx::query(
r#"
INSERT INTO blob_reconciliation_checkpoints
(kind, scope, status, cursor, last_key, completed_at)
VALUES ('blob_cleanup_plan', $1, $2, $3, $4, CASE WHEN $5 THEN CURRENT_TIMESTAMP ELSE NULL END)
ON CONFLICT (kind, scope) DO UPDATE
SET status = EXCLUDED.status,
cursor = EXCLUDED.cursor,
last_key = COALESCE(EXCLUDED.last_key, blob_reconciliation_checkpoints.last_key),
completed_at = CASE WHEN $5 THEN CURRENT_TIMESTAMP ELSE NULL END,
updated_at = CURRENT_TIMESTAMP
"#,
)
.bind(workspace_id)
.bind(status)
.bind(serde_json::json!({ "lastBlobKey": last_blob_key }))
.bind(last_blob_key)
.bind(completed)
.execute(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup plan checkpoint write failed: {err}")))?;
Ok(())
}
async fn create_run(pool: &PgPool, workspace_id: &str) -> Result<String> {
sqlx::query_scalar::<_, String>(
r#"
INSERT INTO blob_reconciliation_runs (kind, mode, status, workspace_id)
VALUES ('blob_cleanup_plan', 'mark_only', 'running', $1)
RETURNING id::text
"#,
)
.bind(workspace_id)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup create run failed: {err}")))
}
async fn finish_run(
pool: &PgPool,
run_id: &str,
workspace_id: &str,
result: &RuntimeBlobCleanupPlanResult,
stale_projection_workspaces: Vec<String>,
) -> Result<()> {
let candidate_bytes = sqlx::query_scalar::<_, Option<i64>>(
"SELECT SUM(object_size)::bigint FROM blob_cleanup_candidates WHERE run_id = $1::uuid AND status = 'marked'",
)
.bind(run_id)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup candidate bytes audit failed: {err}")))?
.unwrap_or(0);
sqlx::query(
r#"
UPDATE blob_reconciliation_runs
SET status = 'finished',
finished_at = CURRENT_TIMESTAMP,
scanned = $2,
changed = $3,
metadata = $4
WHERE id = $1::uuid
"#,
)
.bind(run_id)
.bind(result.scanned_blobs as i32)
.bind(result.candidates_marked as i32)
.bind(serde_json::json!({
"protectedByDocRefs": result.protected_by_doc_refs,
"protectedByMetadata": result.protected_by_metadata,
"protectedByOtherRefs": result.protected_by_other_refs,
"topWorkspaceCandidateBytes": [{
"workspaceId": workspace_id,
"candidateBytes": candidate_bytes,
}],
"staleOrFailedProjectionWorkspaces": stale_projection_workspaces,
}))
.execute(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup finish run failed: {err}")))?;
Ok(())
}
async fn mark_candidate_status(
pool: &PgPool,
run_id: &str,
workspace_id: &str,
blob_key: &str,
status: &str,
evidence: serde_json::Value,
error: Option<&str>,
) -> Result<()> {
sqlx::query(
r#"
UPDATE blob_cleanup_candidates
SET status = $3,
executed_at = CURRENT_TIMESTAMP,
evidence = evidence || $4,
error = $5
WHERE workspace_id = $1 AND blob_key = $2 AND run_id = $6::uuid
"#,
)
.bind(workspace_id)
.bind(blob_key)
.bind(status)
.bind(evidence)
.bind(error)
.bind(run_id)
.execute(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup mark candidate status failed: {err}")))?;
Ok(())
}
async fn finish_execute_run(pool: &PgPool, run_id: &str, result: &RuntimeBlobCleanupExecuteResult) -> Result<()> {
sqlx::query(
r#"
UPDATE blob_reconciliation_runs
SET status = 'finished',
finished_at = CURRENT_TIMESTAMP,
scanned = $2,
changed = $3,
failed = $4,
metadata = metadata || $5
WHERE id = $1::uuid
"#,
)
.bind(run_id)
.bind(result.scanned_candidates as i32)
.bind(result.deleted_metadata as i32)
.bind(result.failed as i32)
.bind(serde_json::json!({
"deletedObjects": result.deleted_objects,
"deletedMetadata": result.deleted_metadata,
"skippedStillReferenced": result.skipped_still_referenced,
"failed": result.failed,
}))
.execute(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup execute run finish failed: {err}")))?;
Ok(())
}
async fn mark_candidate(
pool: &PgPool,
run_id: &str,
row: &BlobCandidateRow,
object_size: i64,
object_last_modified: DateTime<Utc>,
) -> Result<i64> {
let result = sqlx::query(
r#"
INSERT INTO blob_cleanup_candidates
(workspace_id, blob_key, reason, status, object_size, object_last_modified, run_id, evidence)
VALUES ($1, $2, 'unreferenced_completed_blob', 'marked', $3, $4, $5::uuid, $6)
ON CONFLICT (workspace_id, blob_key) DO UPDATE
SET reason = EXCLUDED.reason,
status = 'marked',
object_size = EXCLUDED.object_size,
object_last_modified = EXCLUDED.object_last_modified,
planned_at = CURRENT_TIMESTAMP,
run_id = EXCLUDED.run_id,
evidence = EXCLUDED.evidence,
error = NULL
"#,
)
.bind(&row.workspace_id)
.bind(&row.key)
.bind(object_size)
.bind(object_last_modified)
.bind(run_id)
.bind(serde_json::json!({ "metadataSize": row.size }))
.execute(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup mark candidate failed: {err}")))?;
Ok(result.rows_affected() as i64)
}
async fn load_marked_candidates(pool: &PgPool, run_id: &str, limit: i64) -> Result<Vec<MarkedCandidateRow>> {
sqlx::query_as::<_, MarkedCandidateRow>(
r#"
SELECT workspace_id, blob_key
FROM blob_cleanup_candidates
WHERE run_id = $1::uuid AND status IN ('marked', 'failed')
ORDER BY CASE WHEN status = 'marked' THEN 0 ELSE 1 END, planned_at ASC
LIMIT $2
"#,
)
.bind(run_id)
.bind(limit)
.fetch_all(pool)
.await
.map_err(|err| napi_error(format!("Blob cleanup load marked candidates failed: {err}")))
}
#[napi_derive::napi]
impl BackendRuntime {
#[napi]
pub async fn plan_unreferenced_workspace_blobs(
&self,
workspace_id: String,
grace_period_days: i64,
limit: i64,
) -> Result<RuntimeBlobCleanupPlanResult> {
if limit <= 0 {
return Err(napi_error("blob cleanup plan limit must be positive"));
}
if grace_period_days < 0 {
return Err(napi_error("blob cleanup grace period must be non-negative"));
}
let pool = self.pool().await?;
let run_id = create_run(&pool, &workspace_id).await?;
let mut result = RuntimeBlobCleanupPlanResult {
run_id: Some(run_id.clone()),
scanned_blobs: 0,
candidates_marked: 0,
protected_by_doc_refs: 0,
protected_by_metadata: 0,
protected_by_other_refs: 0,
next_cursor: None,
};
let cursor = load_plan_cursor(&pool, &workspace_id).await?;
let stale_projection_workspaces = stale_projection_workspaces(&pool, &workspace_id).await?;
if !metadata_backfill_is_complete(&pool, &workspace_id).await? || !stale_projection_workspaces.is_empty() {
result.protected_by_metadata = load_completed_blobs(&pool, &workspace_id, cursor.as_deref(), limit)
.await?
.len() as i64;
finish_run(&pool, &run_id, &workspace_id, &result, stale_projection_workspaces).await?;
return Ok(result);
}
let min_last_modified = Utc::now() - Duration::days(grace_period_days);
let rows = load_completed_blobs(&pool, &workspace_id, cursor.as_deref(), limit).await?;
let has_more = rows.len() == limit as usize;
let mut last_blob_key = None;
for row in rows {
result.scanned_blobs += 1;
last_blob_key = Some(row.key.clone());
if has_doc_ref(&pool, &row.workspace_id, &row.key).await? {
result.protected_by_doc_refs += 1;
continue;
}
if has_other_ref(&pool, &row.workspace_id, &row.key).await? {
result.protected_by_other_refs += 1;
continue;
}
let object_key = format!("{}/{}", row.workspace_id, row.key);
let Some(metadata) = self.object_storage_head(object_key).await? else {
result.protected_by_metadata += 1;
continue;
};
let last_modified = DateTime::<Utc>::from_timestamp_millis(metadata.last_modified_ms)
.ok_or_else(|| napi_error("blob cleanup object last modified is invalid"))?;
if metadata.content_length != row.size as i64 || last_modified > min_last_modified {
result.protected_by_metadata += 1;
continue;
}
result.candidates_marked += mark_candidate(&pool, &run_id, &row, metadata.content_length, last_modified).await?;
}
if has_more {
result.next_cursor = last_blob_key.clone();
}
upsert_plan_checkpoint(&pool, &workspace_id, last_blob_key.as_deref(), !has_more).await?;
finish_run(&pool, &run_id, &workspace_id, &result, Vec::new()).await?;
Ok(result)
}
#[napi]
pub async fn execute_blob_cleanup_candidates(
&self,
run_id: String,
grace_period_days: i64,
limit: i64,
) -> Result<RuntimeBlobCleanupExecuteResult> {
if limit <= 0 {
return Err(napi_error("blob cleanup execute limit must be positive"));
}
if grace_period_days < 0 {
return Err(napi_error("blob cleanup grace period must be non-negative"));
}
let pool = self.pool().await?;
let min_last_modified = Utc::now() - Duration::days(grace_period_days);
let rows = load_marked_candidates(&pool, &run_id, limit).await?;
let mut result = RuntimeBlobCleanupExecuteResult {
scanned_candidates: rows.len() as i64,
deleted_objects: 0,
deleted_metadata: 0,
skipped_still_referenced: 0,
failed: 0,
workspace_ids: Vec::new(),
};
for row in rows {
if projection_is_stale(&pool, &row.workspace_id).await?
|| has_doc_ref(&pool, &row.workspace_id, &row.blob_key).await?
|| has_other_ref(&pool, &row.workspace_id, &row.blob_key).await?
{
result.skipped_still_referenced += 1;
mark_candidate_status(
&pool,
&run_id,
&row.workspace_id,
&row.blob_key,
"skipped",
serde_json::json!({ "skipReason": "referenced_or_projection_stale" }),
None,
)
.await?;
continue;
}
let object_key = format!("{}/{}", row.workspace_id, row.blob_key);
let mut object_was_missing = false;
let metadata = match self.object_storage_head(object_key.clone()).await {
Ok(metadata) => metadata,
Err(err) => {
result.failed += 1;
mark_candidate_status(
&pool,
&run_id,
&row.workspace_id,
&row.blob_key,
"failed",
serde_json::json!({ "failure": "object_head_failed" }),
Some(&err.to_string()),
)
.await?;
continue;
}
};
if let Some(metadata) = metadata {
let last_modified = DateTime::<Utc>::from_timestamp_millis(metadata.last_modified_ms)
.ok_or_else(|| napi_error("blob cleanup execute object last modified is invalid"))?;
if last_modified > min_last_modified {
result.skipped_still_referenced += 1;
mark_candidate_status(
&pool,
&run_id,
&row.workspace_id,
&row.blob_key,
"skipped",
serde_json::json!({ "skipReason": "object_inside_grace_period" }),
None,
)
.await?;
continue;
}
if let Err(err) = self.object_storage_delete(object_key).await {
result.failed += 1;
mark_candidate_status(
&pool,
&run_id,
&row.workspace_id,
&row.blob_key,
"failed",
serde_json::json!({ "failure": "object_delete_failed" }),
Some(&err.to_string()),
)
.await?;
continue;
}
result.deleted_objects += 1;
} else {
object_was_missing = true;
}
let deleted_metadata =
match sqlx::query("DELETE FROM blobs WHERE workspace_id = $1 AND key = $2 AND deleted_at IS NULL")
.bind(&row.workspace_id)
.bind(&row.blob_key)
.execute(&pool)
.await
{
Ok(result) => result.rows_affected() as i64,
Err(err) => {
result.failed += 1;
mark_candidate_status(
&pool,
&run_id,
&row.workspace_id,
&row.blob_key,
"failed",
serde_json::json!({ "failure": "metadata_delete_failed" }),
Some(&err.to_string()),
)
.await?;
continue;
}
};
result.deleted_metadata += deleted_metadata;
push_workspace_once(&mut result.workspace_ids, &row.workspace_id);
mark_candidate_status(
&pool,
&run_id,
&row.workspace_id,
&row.blob_key,
"executed",
serde_json::json!({
"deletedMetadata": deleted_metadata,
"objectMissingBeforeDelete": object_was_missing,
}),
None,
)
.await?;
}
finish_execute_run(&pool, &run_id, &result).await?;
Ok(result)
}
}
@@ -0,0 +1,276 @@
use chrono::{DateTime, Utc};
use napi::Result;
use sqlx::{FromRow, PgPool};
use super::{
BackendRuntime,
error::napi_error,
types::{RuntimeBlobMetadataBackfillResult, RuntimeObjectMetadata},
};
async fn workspace_exists(pool: &PgPool, workspace_id: &str) -> Result<bool> {
sqlx::query_scalar::<_, bool>("SELECT EXISTS(SELECT 1 FROM workspaces WHERE id = $1)")
.bind(workspace_id)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob metadata backfill workspace check failed: {err}")))
}
async fn blob_exists(pool: &PgPool, workspace_id: &str, key: &str) -> Result<bool> {
sqlx::query_scalar::<_, bool>("SELECT EXISTS(SELECT 1 FROM blobs WHERE workspace_id = $1 AND key = $2)")
.bind(workspace_id)
.bind(key)
.fetch_one(pool)
.await
.map_err(|err| napi_error(format!("Blob metadata backfill blob check failed: {err}")))
}
async fn upsert_blob_metadata(
pool: &PgPool,
workspace_id: &str,
key: &str,
metadata: RuntimeObjectMetadata,
) -> Result<i64> {
let last_modified = DateTime::<Utc>::from_timestamp_millis(metadata.last_modified_ms)
.ok_or_else(|| napi_error("Blob metadata backfill object last modified is invalid"))?;
let result = sqlx::query(
r#"
INSERT INTO blobs (workspace_id, key, size, mime, status, upload_id, created_at, deleted_at)
VALUES ($1, $2, $3, $4, 'completed', NULL, $5, NULL)
ON CONFLICT (workspace_id, key) DO UPDATE
SET size = EXCLUDED.size,
mime = EXCLUDED.mime,
status = 'completed',
upload_id = NULL,
deleted_at = NULL
WHERE blobs.deleted_at IS NULL
"#,
)
.bind(workspace_id)
.bind(key)
.bind(metadata.content_length as i32)
.bind(metadata.content_type)
.bind(last_modified)
.execute(pool)
.await
.map_err(|err| napi_error(format!("Blob metadata backfill upsert failed: {err}")))?;
Ok(result.rows_affected() as i64)
}
fn split_workspace_blob_key(full_key: &str) -> Option<(&str, &str)> {
let (workspace_id, key) = full_key.split_once('/')?;
if workspace_id.is_empty() || key.is_empty() || key.contains('/') {
return None;
}
Some((workspace_id, key))
}
fn checkpoint_scope(workspace_id: Option<&str>) -> String {
workspace_id.unwrap_or("__all__").to_string()
}
#[derive(FromRow)]
struct BackfillCheckpoint {
last_key: Option<String>,
cursor: serde_json::Value,
}
impl BackfillCheckpoint {
fn continuation_token(&self) -> Option<String> {
self
.cursor
.get("continuationToken")
.and_then(|value| value.as_str())
.map(ToString::to_string)
}
}
async fn load_checkpoint(pool: &PgPool, scope: &str) -> Result<Option<BackfillCheckpoint>> {
sqlx::query_as::<_, BackfillCheckpoint>(
"SELECT last_key, cursor FROM blob_reconciliation_checkpoints WHERE kind = 'blob_metadata_backfill' AND scope = $1",
)
.bind(scope)
.fetch_optional(pool)
.await
.map_err(|err| napi_error(format!("Blob metadata backfill checkpoint load failed: {err}")))
}
async fn upsert_checkpoint(
pool: &PgPool,
scope: &str,
last_key: Option<&str>,
continuation_token: Option<&str>,
completed: bool,
) -> Result<()> {
let status = if completed { "completed" } else { "running" };
sqlx::query(
r#"
INSERT INTO blob_reconciliation_checkpoints
(kind, scope, status, cursor, last_key, completed_at, metadata)
VALUES ('blob_metadata_backfill', $1, $2, $3, $4, CASE WHEN $5 THEN CURRENT_TIMESTAMP ELSE NULL END, $6)
ON CONFLICT (kind, scope) DO UPDATE
SET status = EXCLUDED.status,
cursor = EXCLUDED.cursor,
last_key = COALESCE(EXCLUDED.last_key, blob_reconciliation_checkpoints.last_key),
completed_at = CASE WHEN $5 THEN CURRENT_TIMESTAMP ELSE NULL END,
updated_at = CURRENT_TIMESTAMP,
metadata = EXCLUDED.metadata
"#,
)
.bind(scope)
.bind(status)
.bind(serde_json::json!({
"lastKey": last_key,
"continuationToken": continuation_token,
}))
.bind(last_key)
.bind(completed)
.bind(serde_json::json!({
"quotaReportingReconciliationRequired": true,
}))
.execute(pool)
.await
.map_err(|err| napi_error(format!("Blob metadata backfill checkpoint write failed: {err}")))?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn blob_metadata_backfill_splits_workspace_blob_keys() {
assert_eq!(
split_workspace_blob_key("workspace/blob-key"),
Some(("workspace", "blob-key"))
);
assert_eq!(split_workspace_blob_key("workspace/nested/blob-key"), None);
assert_eq!(split_workspace_blob_key("workspace/"), None);
assert_eq!(split_workspace_blob_key("blob-key"), None);
}
#[test]
fn blob_metadata_backfill_checkpoint_scope_is_explicit() {
assert_eq!(checkpoint_scope(Some("workspace")), "workspace");
assert_eq!(checkpoint_scope(None), "__all__");
}
}
fn push_workspace_once(workspace_ids: &mut Vec<String>, workspace_id: &str) {
if !workspace_ids.iter().any(|id| id == workspace_id) {
workspace_ids.push(workspace_id.to_string());
}
}
fn checked_list_page_limit(limit: i64) -> Result<i32> {
i32::try_from(limit).map_err(|_| napi_error("blob metadata backfill limit exceeds i32::MAX"))
}
#[napi_derive::napi]
impl BackendRuntime {
#[napi]
pub async fn backfill_missing_blob_metadata(
&self,
workspace_id: Option<String>,
limit: i64,
) -> Result<RuntimeBlobMetadataBackfillResult> {
if limit <= 0 {
return Err(napi_error("blob metadata backfill limit must be positive"));
}
let page_limit = checked_list_page_limit(limit)?;
let pool = self.pool().await?;
let prefix = workspace_id.as_ref().map(|id| format!("{id}/"));
let scope = checkpoint_scope(workspace_id.as_deref());
let checkpoint = load_checkpoint(&pool, &scope).await?;
let page = self
.object_storage_list_page(
prefix,
checkpoint.as_ref().and_then(BackfillCheckpoint::continuation_token),
checkpoint.as_ref().and_then(|checkpoint| checkpoint.last_key.clone()),
page_limit,
)
.await?;
let has_more = page.next_continuation_token.is_some();
let mut result = RuntimeBlobMetadataBackfillResult {
scanned_objects: 0,
headed_objects: 0,
upserted_metadata: 0,
skipped_existing: 0,
skipped_workspace_missing: 0,
failed: 0,
next_cursor: None,
workspace_ids: Vec::new(),
};
let mut last_scanned_key = None;
for object in &page.entries {
result.scanned_objects += 1;
last_scanned_key = Some(object.key.clone());
let Some((object_workspace_id, key)) = split_workspace_blob_key(&object.key) else {
result.failed += 1;
continue;
};
if workspace_id.as_deref().is_some_and(|id| id != object_workspace_id) {
result.failed += 1;
continue;
}
if !workspace_exists(&pool, object_workspace_id).await? {
result.skipped_workspace_missing += 1;
continue;
}
if blob_exists(&pool, object_workspace_id, key).await? {
result.skipped_existing += 1;
continue;
}
result.headed_objects += 1;
let Some(metadata) = self.object_storage_head(object.key.clone()).await? else {
result.failed += 1;
continue;
};
let affected = upsert_blob_metadata(&pool, object_workspace_id, key, metadata).await?;
if affected > 0 {
result.upserted_metadata += affected;
push_workspace_once(&mut result.workspace_ids, object_workspace_id);
}
}
if has_more {
result.next_cursor = last_scanned_key.clone();
}
upsert_checkpoint(
&pool,
&scope,
last_scanned_key.as_deref(),
page.next_continuation_token.as_deref(),
!has_more,
)
.await?;
sqlx::query(
r#"
INSERT INTO blob_reconciliation_runs
(kind, mode, status, workspace_id, finished_at, scanned, changed, failed, metadata)
VALUES ('blob_metadata_backfill', 'execute', 'finished', $1, CURRENT_TIMESTAMP, $2, $3, $4, $5)
"#,
)
.bind(workspace_id)
.bind(result.scanned_objects as i32)
.bind(result.upserted_metadata as i32)
.bind(result.failed as i32)
.bind(serde_json::json!({
"headedObjects": result.headed_objects,
"skippedExisting": result.skipped_existing,
"skippedWorkspaceMissing": result.skipped_workspace_missing,
"checkpointScope": scope,
"nextCursor": result.next_cursor,
"quotaReportingReconciliationRequired": true,
}))
.execute(&pool)
.await
.map_err(|err| napi_error(format!("Blob metadata backfill run record failed: {err}")))?;
Ok(result)
}
}
@@ -20,8 +20,9 @@ pub(super) struct RuntimeConfig {
impl RuntimeConfig {
pub(super) fn from_config_files() -> Result<Self> {
let database_url =
database_url_from_config_files()?.unwrap_or_else(|| "postgresql://localhost:5432/affine".to_string());
let database_url = database_url_from_env()
.or(database_url_from_config_files()?)
.unwrap_or_else(|| "postgresql://localhost:5432/affine".to_string());
let storage = ObjectStorageConfig::from_config_files()?;
Ok(Self { database_url, storage })
}
@@ -39,6 +40,14 @@ struct DbConfigFile {
datasource_url: Option<String>,
}
fn database_url_from_env() -> Option<String> {
env::var("DATABASE_URL").ok().and_then(non_empty_string)
}
fn non_empty_string(value: String) -> Option<String> {
if value.trim().is_empty() { None } else { Some(value) }
}
fn database_url_from_config_files() -> Result<Option<String>> {
let mut database_url = None;
for path in config_json_paths() {
@@ -49,9 +58,7 @@ fn database_url_from_config_files() -> Result<Option<String>> {
.map_err(|err| napi_error(format!("failed to read config file {}: {err}", path.display())))?;
let config: AppConfigFile = serde_json::from_str(&raw)
.map_err(|err| napi_error(format!("failed to parse config file {}: {err}", path.display())))?;
if let Some(next) = config.db.and_then(|db| db.datasource_url)
&& !next.trim().is_empty()
{
if let Some(next) = config.db.and_then(|db| db.datasource_url).and_then(non_empty_string) {
database_url = Some(next);
}
}
@@ -125,4 +132,14 @@ mod tests {
.all(|path| !path.to_string_lossy().contains("packages/backend/server"))
);
}
#[test]
fn blank_database_urls_are_ignored() {
assert_eq!(non_empty_string("".to_string()), None);
assert_eq!(non_empty_string(" ".to_string()), None);
assert_eq!(
non_empty_string("postgresql://affine:affine@localhost:5432/affine".to_string()),
Some("postgresql://affine:affine@localhost:5432/affine".to_string())
);
}
}
@@ -0,0 +1,416 @@
use affine_common::doc_parser;
use chrono::{DateTime, Utc};
use napi::Result;
use sqlx::{FromRow, PgPool};
use y_octo::Doc;
use super::{BackendRuntime, error::napi_error, types::RuntimeDocBlobRefsResult};
const PARSER_VERSION: i32 = 1;
#[derive(FromRow)]
struct SnapshotRow {
workspace_id: String,
doc_id: String,
blob: Vec<u8>,
updated_at: DateTime<Utc>,
}
#[derive(FromRow)]
struct UpdateRow {
blob: Vec<u8>,
created_at: DateTime<Utc>,
}
struct ExtractedRef {
blob_key: String,
block_id: String,
flavour: String,
}
async fn load_snapshot(pool: &PgPool, workspace_id: &str, doc_id: &str) -> Result<Option<SnapshotRow>> {
sqlx::query_as::<_, SnapshotRow>(
r#"
SELECT workspace_id, guid AS doc_id, blob, updated_at
FROM snapshots
WHERE workspace_id = $1 AND guid = $2
"#,
)
.bind(workspace_id)
.bind(doc_id)
.fetch_optional(pool)
.await
.map_err(|err| napi_error(format!("Doc blob refs load snapshot failed: {err}")))
}
async fn load_updates(pool: &PgPool, workspace_id: &str, doc_id: &str) -> Result<Vec<UpdateRow>> {
sqlx::query_as::<_, UpdateRow>(
r#"
SELECT blob, created_at
FROM updates
WHERE workspace_id = $1 AND guid = $2
ORDER BY created_at ASC
"#,
)
.bind(workspace_id)
.bind(doc_id)
.fetch_all(pool)
.await
.map_err(|err| napi_error(format!("Doc blob refs load updates failed: {err}")))
}
fn apply_doc_updates(updates: impl IntoIterator<Item = Vec<u8>>) -> Result<Vec<u8>> {
let mut doc = Doc::default();
for update in updates {
doc
.apply_update_from_binary_v1(&update)
.map_err(|err| napi_error(format!("Doc blob refs merge failed: {err}")))?;
}
doc
.encode_update_v1()
.map_err(|err| napi_error(format!("Doc blob refs encode failed: {err}")))
}
async fn load_current_doc(pool: &PgPool, workspace_id: &str, doc_id: &str) -> Result<Option<SnapshotRow>> {
let snapshot = load_snapshot(pool, workspace_id, doc_id).await?;
let updates = load_updates(pool, workspace_id, doc_id).await?;
if snapshot.is_none() && updates.is_empty() {
return Ok(None);
}
let mut merge_inputs = Vec::with_capacity(updates.len() + usize::from(snapshot.is_some()));
let mut updated_at = snapshot
.as_ref()
.map(|snapshot| snapshot.updated_at)
.unwrap_or_else(Utc::now);
if let Some(snapshot) = snapshot {
merge_inputs.push(snapshot.blob);
}
for update in updates {
updated_at = update.created_at;
merge_inputs.push(update.blob);
}
Ok(Some(SnapshotRow {
workspace_id: workspace_id.to_string(),
doc_id: doc_id.to_string(),
blob: apply_doc_updates(merge_inputs)?,
updated_at,
}))
}
async fn load_workspace_doc_ids(pool: &PgPool, workspace_id: &str) -> Result<Vec<String>> {
let Some(root) = load_current_doc(pool, workspace_id, workspace_id).await? else {
return Ok(Vec::new());
};
let ids = doc_parser::get_doc_ids_from_binary(root.blob, false)
.map_err(|err| napi_error(format!("Doc blob refs root doc parse failed: {err}")))?;
let mut ids = ids;
ids.sort();
Ok(ids)
}
async fn upsert_projection_checkpoint(
pool: &PgPool,
workspace_id: &str,
result: &RuntimeDocBlobRefsResult,
) -> Result<()> {
let completed = result.next_cursor.is_none();
let status = if completed && result.failed_docs == 0 {
"completed"
} else if result.failed_docs > 0 {
"failed"
} else {
"running"
};
sqlx::query(
r#"
INSERT INTO blob_reconciliation_checkpoints
(kind, scope, status, cursor, completed_at, metadata)
VALUES ('doc_blob_refs', $1, $2, $3, CASE WHEN $4 THEN CURRENT_TIMESTAMP ELSE NULL END, $5)
ON CONFLICT (kind, scope) DO UPDATE
SET status = EXCLUDED.status,
cursor = EXCLUDED.cursor,
completed_at = CASE WHEN $4 THEN CURRENT_TIMESTAMP ELSE NULL END,
updated_at = CURRENT_TIMESTAMP,
metadata = EXCLUDED.metadata
"#,
)
.bind(workspace_id)
.bind(status)
.bind(serde_json::json!({ "lastDocId": result.next_cursor }))
.bind(completed && result.failed_docs == 0)
.bind(serde_json::json!({
"parserVersion": PARSER_VERSION,
}))
.execute(pool)
.await
.map_err(|err| napi_error(format!("Doc blob refs checkpoint write failed: {err}")))?;
Ok(())
}
async fn upsert_projection_failure_checkpoint(pool: &PgPool, workspace_id: &str, error: &str) -> Result<()> {
sqlx::query(
r#"
INSERT INTO blob_reconciliation_checkpoints
(kind, scope, status, cursor, completed_at, metadata)
VALUES ('doc_blob_refs', $1, 'failed', '{}', NULL, $2)
ON CONFLICT (kind, scope) DO UPDATE
SET status = 'failed',
cursor = '{}',
completed_at = NULL,
updated_at = CURRENT_TIMESTAMP,
metadata = EXCLUDED.metadata
"#,
)
.bind(workspace_id)
.bind(serde_json::json!({
"parserVersion": PARSER_VERSION,
"error": error,
}))
.execute(pool)
.await
.map_err(|err| napi_error(format!("Doc blob refs failure checkpoint write failed: {err}")))?;
Ok(())
}
async fn load_projection_cursor(pool: &PgPool, workspace_id: &str) -> Result<Option<String>> {
let cursor = sqlx::query_scalar::<_, serde_json::Value>(
"SELECT cursor FROM blob_reconciliation_checkpoints WHERE kind = 'doc_blob_refs' AND scope = $1",
)
.bind(workspace_id)
.fetch_optional(pool)
.await
.map_err(|err| napi_error(format!("Doc blob refs checkpoint load failed: {err}")))?;
Ok(cursor.and_then(|cursor| {
cursor
.get("lastDocId")
.and_then(|value| value.as_str())
.map(ToString::to_string)
}))
}
async fn purge_removed_doc_refs(pool: &PgPool, workspace_id: &str, current_doc_ids: &[String]) -> Result<i64> {
let result = sqlx::query(
r#"
DELETE FROM doc_blob_refs
WHERE workspace_id = $1
AND NOT (doc_id = ANY($2))
"#,
)
.bind(workspace_id)
.bind(current_doc_ids)
.execute(pool)
.await
.map_err(|err| napi_error(format!("Doc blob refs purge removed docs failed: {err}")))?;
Ok(result.rows_affected() as i64)
}
fn extract_refs(snapshot: &SnapshotRow) -> Result<Vec<ExtractedRef>> {
let parsed = doc_parser::parse_doc_from_binary(snapshot.blob.clone(), snapshot.doc_id.clone())
.map_err(|err| napi_error(format!("Doc blob refs parse failed: {err}")))?;
let mut refs = Vec::new();
for block in parsed.blocks {
let Some(blob_keys) = block.blob else {
continue;
};
for blob_key in blob_keys {
refs.push(ExtractedRef {
blob_key,
block_id: block.block_id.clone(),
flavour: block.flavour.clone(),
});
}
}
Ok(refs)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn doc_blob_refs_extracts_image_refs() {
let doc_id = "doc-blob-ref-test".to_string();
let blob =
doc_parser::build_full_doc("Doc", "![Alt](blob://image-blob-key)", &doc_id).expect("doc fixture should build");
let snapshot = SnapshotRow {
workspace_id: "workspace".to_string(),
doc_id,
blob,
updated_at: Utc::now(),
};
let refs = extract_refs(&snapshot).expect("refs should parse");
assert!(
refs
.iter()
.any(|reference| { reference.blob_key == "image-blob-key" && reference.flavour == "affine:image" })
);
}
}
async fn replace_doc_refs(pool: &PgPool, snapshot: &SnapshotRow, refs: Vec<ExtractedRef>) -> Result<(i64, i64)> {
let mut tx = pool
.begin()
.await
.map_err(|err| napi_error(format!("Doc blob refs transaction failed: {err}")))?;
let deleted = sqlx::query("DELETE FROM doc_blob_refs WHERE workspace_id = $1 AND doc_id = $2")
.bind(&snapshot.workspace_id)
.bind(&snapshot.doc_id)
.execute(&mut *tx)
.await
.map_err(|err| napi_error(format!("Doc blob refs delete failed: {err}")))?
.rows_affected() as i64;
let mut written = 0;
for reference in refs {
let affected = sqlx::query(
r#"
INSERT INTO doc_blob_refs
(workspace_id, doc_id, blob_key, block_id, flavour, snapshot_updated_at, parser_version, status)
VALUES ($1, $2, $3, $4, $5, $6, $7, 'fresh')
ON CONFLICT (workspace_id, doc_id, blob_key, block_id) DO UPDATE
SET flavour = EXCLUDED.flavour,
snapshot_updated_at = EXCLUDED.snapshot_updated_at,
indexed_at = CURRENT_TIMESTAMP,
parser_version = EXCLUDED.parser_version,
status = 'fresh',
error = NULL
"#,
)
.bind(&snapshot.workspace_id)
.bind(&snapshot.doc_id)
.bind(reference.blob_key)
.bind(reference.block_id)
.bind(reference.flavour)
.bind(snapshot.updated_at)
.bind(PARSER_VERSION)
.execute(&mut *tx)
.await
.map_err(|err| napi_error(format!("Doc blob refs insert failed: {err}")))?
.rows_affected() as i64;
written += affected;
}
tx.commit()
.await
.map_err(|err| napi_error(format!("Doc blob refs transaction commit failed: {err}")))?;
Ok((written, deleted))
}
async fn mark_doc_failed(pool: &PgPool, workspace_id: &str, doc_id: &str, error: &str) -> Result<()> {
sqlx::query(
r#"
INSERT INTO doc_blob_refs
(workspace_id, doc_id, blob_key, block_id, flavour, snapshot_updated_at, parser_version, status, error)
VALUES ($1, $2, '__parse_failed__', '__parse_failed__', '__parse_failed__', CURRENT_TIMESTAMP, $3, 'failed', $4)
ON CONFLICT (workspace_id, doc_id, blob_key, block_id) DO UPDATE
SET indexed_at = CURRENT_TIMESTAMP,
status = 'failed',
error = EXCLUDED.error
"#,
)
.bind(workspace_id)
.bind(doc_id)
.bind(PARSER_VERSION)
.bind(error)
.execute(pool)
.await
.map_err(|err| napi_error(format!("Doc blob refs mark failure failed: {err}")))?;
Ok(())
}
#[napi_derive::napi]
impl BackendRuntime {
#[napi]
pub async fn rebuild_doc_blob_refs(&self, workspace_id: String, doc_id: String) -> Result<RuntimeDocBlobRefsResult> {
let pool = self.pool().await?;
let mut result = RuntimeDocBlobRefsResult {
scanned_docs: 1,
parsed_docs: 0,
refs_written: 0,
refs_deleted: 0,
failed_docs: 0,
next_cursor: None,
};
let Some(snapshot) = load_current_doc(&pool, &workspace_id, &doc_id).await? else {
result.failed_docs = 1;
mark_doc_failed(&pool, &workspace_id, &doc_id, "snapshot_missing").await?;
return Ok(result);
};
match extract_refs(&snapshot) {
Ok(refs) => {
let (written, deleted) = replace_doc_refs(&pool, &snapshot, refs).await?;
result.parsed_docs = 1;
result.refs_written = written;
result.refs_deleted = deleted;
}
Err(err) => {
result.failed_docs = 1;
mark_doc_failed(&pool, &workspace_id, &doc_id, &err.to_string()).await?;
}
}
Ok(result)
}
#[napi]
pub async fn rebuild_workspace_doc_blob_refs(
&self,
workspace_id: String,
limit: i64,
) -> Result<RuntimeDocBlobRefsResult> {
if limit <= 0 {
return Err(napi_error("doc blob refs rebuild limit must be positive"));
}
let pool = self.pool().await?;
let doc_ids = match load_workspace_doc_ids(&pool, &workspace_id).await {
Ok(doc_ids) => doc_ids,
Err(err) => {
upsert_projection_failure_checkpoint(&pool, &workspace_id, &err.to_string()).await?;
return Err(err);
}
};
let cursor = load_projection_cursor(&pool, &workspace_id).await?;
let current_doc_ids = doc_ids.clone();
let doc_ids = doc_ids
.into_iter()
.filter(|doc_id| cursor.as_ref().is_none_or(|cursor| doc_id > cursor))
.collect::<Vec<_>>();
let has_more = doc_ids.len() > limit as usize;
let mut total = RuntimeDocBlobRefsResult {
scanned_docs: 0,
parsed_docs: 0,
refs_written: 0,
refs_deleted: 0,
failed_docs: 0,
next_cursor: None,
};
let mut last_doc_id = None;
for doc_id in doc_ids.into_iter().take(limit as usize) {
last_doc_id = Some(doc_id.clone());
let result = self.rebuild_doc_blob_refs(workspace_id.clone(), doc_id).await?;
total.scanned_docs += result.scanned_docs;
total.parsed_docs += result.parsed_docs;
total.refs_written += result.refs_written;
total.refs_deleted += result.refs_deleted;
total.failed_docs += result.failed_docs;
}
if has_more {
total.next_cursor = last_doc_id;
} else if total.failed_docs == 0 {
total.refs_deleted += purge_removed_doc_refs(&pool, &workspace_id, &current_doc_ids).await?;
}
upsert_projection_checkpoint(&pool, &workspace_id, &total).await?;
Ok(total)
}
}
@@ -1,8 +1,11 @@
mod blob_cleanup;
mod blob_complete;
mod blob_reclaimer;
mod blob_reconciliation;
mod config;
mod constants;
mod coordination_lease;
mod doc_blob_refs;
mod doc_compactor;
mod doc_storage;
mod error;
@@ -9,8 +9,8 @@ use aws_sdk_s3::{
use napi::Result;
use super::types::{
MultipartUploadInitResult, MultipartUploadPart, ObjectGetResult, ObjectListEntry, ObjectMetadata, ObjectPutMetadata,
PresignedObjectRequest, completed_multipart_parts, trim_etag,
MultipartUploadInitResult, MultipartUploadPart, ObjectGetResult, ObjectListEntry, ObjectListPage, ObjectMetadata,
ObjectPutMetadata, PresignedObjectRequest, completed_multipart_parts, trim_etag,
};
use crate::backend_runtime::error::napi_error;
@@ -233,14 +233,11 @@ impl ObjectStorageClient {
}
pub(super) async fn head(&self, key: &str) -> Result<Option<ObjectMetadata>> {
let result = self
.client
.head_object()
.bucket(&self.bucket)
.key(key)
.send()
.await
.map_err(|err| napi_error(format!("ObjectStorage head failed for {key}: {err:?}")))?;
let result = match self.client.head_object().bucket(&self.bucket).key(key).send().await {
Ok(result) => result,
Err(err) if is_not_found_error(&err) => return Ok(None),
Err(err) => return Err(napi_error(format!("ObjectStorage head failed for {key}: {err:?}"))),
};
Ok(Some(ObjectMetadata {
content_type: result
@@ -253,14 +250,11 @@ impl ObjectStorageClient {
}
pub(super) async fn get(&self, key: &str) -> Result<Option<ObjectGetResult>> {
let result = self
.client
.get_object()
.bucket(&self.bucket)
.key(key)
.send()
.await
.map_err(|err| napi_error(format!("ObjectStorage get failed for {key}: {err:?}")))?;
let result = match self.client.get_object().bucket(&self.bucket).key(key).send().await {
Ok(result) => result,
Err(err) if is_not_found_error(&err) => return Ok(None),
Err(err) => return Err(napi_error(format!("ObjectStorage get failed for {key}: {err:?}"))),
};
let metadata = ObjectMetadata {
content_type: result
.content_type
@@ -314,6 +308,43 @@ impl ObjectStorageClient {
Ok(entries)
}
pub(super) async fn list_page(
&self,
prefix: Option<String>,
continuation_token: Option<String>,
start_after: Option<String>,
max_keys: i32,
) -> Result<ObjectListPage> {
let mut request = self.client.list_objects_v2().bucket(&self.bucket).max_keys(max_keys);
if let Some(prefix) = prefix {
request = request.prefix(prefix);
}
if let Some(continuation_token) = continuation_token {
request = request.continuation_token(continuation_token);
} else if let Some(start_after) = start_after {
request = request.start_after(start_after);
}
let result = request
.send()
.await
.map_err(|err| napi_error(format!("ObjectStorage list page failed: {err:?}")))?;
Ok(ObjectListPage {
entries: result
.contents()
.iter()
.filter_map(|object| {
Some(ObjectListEntry {
key: object.key.as_ref()?.clone(),
content_length: object.size.unwrap_or(0),
last_modified_ms: optional_datetime_ms(object.last_modified),
})
})
.collect(),
next_continuation_token: result.next_continuation_token,
})
}
pub(super) async fn delete(&self, key: &str) -> Result<()> {
self
.client
@@ -327,6 +358,11 @@ impl ObjectStorageClient {
}
}
fn is_not_found_error(error: &impl std::fmt::Debug) -> bool {
let message = format!("{error:?}");
message.contains("NoSuchKey") || (message.contains("NotFound") && !message.contains("NoSuchBucket"))
}
fn expires_at_ms(expires_in_seconds: u64) -> Result<i64> {
let expires_at = SystemTime::now()
.checked_add(Duration::from_secs(expires_in_seconds))
@@ -7,6 +7,7 @@ mod types;
use client::ObjectStorageClient;
pub(super) use config::ObjectStorageConfig;
use napi::{Result, bindgen_prelude::Buffer};
use types::ObjectListPage;
pub(super) use types::StorageProviderConfig;
use super::{
@@ -32,6 +33,19 @@ impl BackendRuntime {
self.object_storage_client()?.delete(key).await
}
pub(super) async fn object_storage_list_page(
&self,
prefix: Option<String>,
continuation_token: Option<String>,
start_after: Option<String>,
max_keys: i32,
) -> Result<ObjectListPage> {
self
.object_storage_client()?
.list_page(prefix, continuation_token, start_after, max_keys)
.await
}
pub(super) async fn object_storage_abort_upload(&self, key: &str, upload_id: &str) -> Result<()> {
self
.object_storage_client()?
@@ -28,10 +28,16 @@ pub(super) struct ObjectMetadata {
}
#[derive(Clone, Debug, PartialEq)]
pub(super) struct ObjectListEntry {
pub(super) key: String,
pub(super) content_length: i64,
pub(super) last_modified_ms: i64,
pub(in crate::backend_runtime) struct ObjectListEntry {
pub(in crate::backend_runtime) key: String,
pub(in crate::backend_runtime) content_length: i64,
pub(in crate::backend_runtime) last_modified_ms: i64,
}
#[derive(Clone, Debug, PartialEq)]
pub(in crate::backend_runtime) struct ObjectListPage {
pub(in crate::backend_runtime) entries: Vec<ObjectListEntry>,
pub(in crate::backend_runtime) next_continuation_token: Option<String>,
}
#[derive(Clone, Debug, PartialEq)]
@@ -38,3 +38,75 @@ CREATE TABLE IF NOT EXISTS runtime_leases (
CREATE INDEX IF NOT EXISTS runtime_leases_expires_at_idx
ON runtime_leases (expires_at);
CREATE TABLE IF NOT EXISTS blob_reconciliation_runs (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
kind TEXT NOT NULL,
mode TEXT NOT NULL,
status TEXT NOT NULL,
workspace_id TEXT,
started_at TIMESTAMPTZ(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
finished_at TIMESTAMPTZ(3),
cursor JSONB NOT NULL DEFAULT '{}',
scanned INTEGER NOT NULL DEFAULT 0,
changed INTEGER NOT NULL DEFAULT 0,
failed INTEGER NOT NULL DEFAULT 0,
metadata JSONB NOT NULL DEFAULT '{}'
);
CREATE INDEX IF NOT EXISTS blob_reconciliation_runs_workspace_idx
ON blob_reconciliation_runs (workspace_id, started_at DESC);
CREATE TABLE IF NOT EXISTS blob_reconciliation_checkpoints (
kind TEXT NOT NULL,
scope TEXT NOT NULL,
status TEXT NOT NULL,
cursor JSONB NOT NULL DEFAULT '{}',
last_key TEXT,
last_sid INTEGER,
completed_at TIMESTAMPTZ(3),
updated_at TIMESTAMPTZ(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
metadata JSONB NOT NULL DEFAULT '{}',
PRIMARY KEY (kind, scope)
);
CREATE INDEX IF NOT EXISTS blob_reconciliation_checkpoints_status_idx
ON blob_reconciliation_checkpoints (kind, status, updated_at DESC);
CREATE TABLE IF NOT EXISTS doc_blob_refs (
workspace_id TEXT NOT NULL,
doc_id TEXT NOT NULL,
blob_key TEXT NOT NULL,
block_id TEXT NOT NULL,
flavour TEXT NOT NULL,
snapshot_updated_at TIMESTAMPTZ(3) NOT NULL,
indexed_at TIMESTAMPTZ(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
parser_version INTEGER NOT NULL,
status TEXT NOT NULL DEFAULT 'fresh',
error TEXT,
PRIMARY KEY (workspace_id, doc_id, blob_key, block_id)
);
CREATE INDEX IF NOT EXISTS doc_blob_refs_workspace_blob_idx
ON doc_blob_refs (workspace_id, blob_key);
CREATE INDEX IF NOT EXISTS doc_blob_refs_workspace_status_idx
ON doc_blob_refs (workspace_id, status);
CREATE TABLE IF NOT EXISTS blob_cleanup_candidates (
workspace_id TEXT NOT NULL,
blob_key TEXT NOT NULL,
reason TEXT NOT NULL,
status TEXT NOT NULL,
object_size BIGINT NOT NULL,
object_last_modified TIMESTAMPTZ(3),
planned_at TIMESTAMPTZ(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
executed_at TIMESTAMPTZ(3),
run_id UUID NOT NULL,
evidence JSONB NOT NULL DEFAULT '{}',
error TEXT,
PRIMARY KEY (workspace_id, blob_key)
);
CREATE INDEX IF NOT EXISTS blob_cleanup_candidates_run_idx
ON blob_cleanup_candidates (run_id, status);
@@ -14,6 +14,10 @@ fn migrations_include_runtime_tables_without_worker_heartbeats() {
assert!(RUNTIME_MIGRATIONS.contains("runtime_states"));
assert!(RUNTIME_MIGRATIONS.contains("runtime_gates"));
assert!(RUNTIME_MIGRATIONS.contains("runtime_leases"));
assert!(RUNTIME_MIGRATIONS.contains("blob_reconciliation_runs"));
assert!(RUNTIME_MIGRATIONS.contains("blob_reconciliation_checkpoints"));
assert!(RUNTIME_MIGRATIONS.contains("doc_blob_refs"));
assert!(RUNTIME_MIGRATIONS.contains("blob_cleanup_candidates"));
assert!(!RUNTIME_MIGRATIONS.contains("runtime_worker_heartbeats"));
}
@@ -138,6 +138,49 @@ pub struct RuntimeBlobCompleteResult {
pub last_modified_ms: Option<i64>,
}
#[napi_derive::napi(object)]
pub struct RuntimeBlobMetadataBackfillResult {
pub scanned_objects: i64,
pub headed_objects: i64,
pub upserted_metadata: i64,
pub skipped_existing: i64,
pub skipped_workspace_missing: i64,
pub failed: i64,
pub next_cursor: Option<String>,
pub workspace_ids: Vec<String>,
}
#[napi_derive::napi(object)]
pub struct RuntimeDocBlobRefsResult {
pub scanned_docs: i64,
pub parsed_docs: i64,
pub refs_written: i64,
pub refs_deleted: i64,
pub failed_docs: i64,
pub next_cursor: Option<String>,
}
#[napi_derive::napi(object)]
pub struct RuntimeBlobCleanupPlanResult {
pub run_id: Option<String>,
pub scanned_blobs: i64,
pub candidates_marked: i64,
pub protected_by_doc_refs: i64,
pub protected_by_metadata: i64,
pub protected_by_other_refs: i64,
pub next_cursor: Option<String>,
}
#[napi_derive::napi(object)]
pub struct RuntimeBlobCleanupExecuteResult {
pub scanned_candidates: i64,
pub deleted_objects: i64,
pub deleted_metadata: i64,
pub skipped_still_referenced: i64,
pub failed: i64,
pub workspace_ids: Vec<String>,
}
#[napi_derive::napi(object)]
pub struct RuntimeDocCompactionResult {
pub lease_acquired: bool,