diff --git a/src/admin/backfill.rs b/src/admin/backfill.rs index 529300e..7d1768c 100644 --- a/src/admin/backfill.rs +++ b/src/admin/backfill.rs @@ -121,6 +121,31 @@ pub(super) async fn create_backfill( // Determine target collections. let collections: Vec = if let Some(ref col) = body.collection { + // Validate that a record-type lexicon exists for the explicit collection. + let lexicon_exists: bool = state + .lexicons + .get(col) + .await + .is_some_and(|lex| lex.lexicon_type == crate::lexicon::LexiconType::Record); + if !lexicon_exists { + let error = format!("no record-type lexicon registered for collection '{col}'"); + let _ = sqlx::query( + "UPDATE backfill_jobs SET status = 'failed', completed_at = NOW(), error = $2 WHERE id::text = $1", + ) + .bind(&job_id) + .bind(&error) + .execute(&state.db) + .await; + + return Ok(( + StatusCode::CREATED, + Json(serde_json::json!({ + "id": job_id, + "status": "failed", + "error": error, + })), + )); + } vec![col.clone()] } else { let rows: Vec<(String,)> = sqlx::query_as( diff --git a/src/tap.rs b/src/tap.rs index 3602094..35a9bb6 100644 --- a/src/tap.rs +++ b/src/tap.rs @@ -315,6 +315,19 @@ async fn run( ) = tokio_tungstenite::connect_async(request).await?; tracing::info!("connected to tap"); + // Re-sync collection filters on every (re)connect so Tap knows which + // collections to track, even if Tap was restarted since the initial sync. + { + let collections = collections_rx.borrow().clone(); + let mut wanted = collections; + if !wanted.contains(&LEXICON_SCHEMA_COLLECTION.to_string()) { + wanted.push(LEXICON_SCHEMA_COLLECTION.to_string()); + } + if let Err(e) = sync_collections(http, tap_url, tap_admin_password, &wanted).await { + tracing::warn!("failed to sync collections to tap on reconnect: {e}"); + } + } + log_event( db, EventLog {