From bbe9f6f3b34d6daec203bebab921eb2530a5dedd Mon Sep 17 00:00:00 2001 From: Lewis Date: Sun, 19 Jul 2026 17:19:22 +0300 Subject: [PATCH] db: decode sequenced event rows leniently Lewis: May this revision serve well! --- crates/tranquil-db/src/postgres/repo.rs | 205 ++++++++++++++---------- 1 file changed, 122 insertions(+), 83 deletions(-) diff --git a/crates/tranquil-db/src/postgres/repo.rs b/crates/tranquil-db/src/postgres/repo.rs index c63cc29..b491435 100644 --- a/crates/tranquil-db/src/postgres/repo.rs +++ b/crates/tranquil-db/src/postgres/repo.rs @@ -10,7 +10,9 @@ use tranquil_db_traits::{ use tranquil_types::{AtUri, CidLink, Did, Handle, Nsid, Rkey, Tid}; use uuid::Uuid; +use super::col; use super::user::map_sqlx_error; +use super::{column, column_vec, legacy_column, opt_column}; struct RecordRow { rkey: String, @@ -43,7 +45,7 @@ fn row_to_event_blocks( ) -> Result, DbError> { match (block_cids, block_data) { (Some(cids), Some(data)) if cids.len() == data.len() => match cids.is_empty() { - true => Ok(legacy_fallback(legacy_blocks_cids)), + true => legacy_fallback(legacy_blocks_cids), false => Ok(Some(EventBlocks::Inline( cids.into_iter() .zip(data) @@ -57,16 +59,19 @@ fn row_to_event_blocks( (Some(_), None) | (None, Some(_)) => Err(DbError::CorruptData( "repo_seq.block_cids/block_data partially populated", )), - (None, None) => Ok(legacy_fallback(legacy_blocks_cids)), + (None, None) => legacy_fallback(legacy_blocks_cids), } } -fn legacy_fallback(legacy_blocks_cids: Option>) -> Option { +fn legacy_fallback( + legacy_blocks_cids: Option>, +) -> Result, DbError> { match legacy_blocks_cids { - Some(cids) if !cids.is_empty() => Some(EventBlocks::LegacyCids( - cids.into_iter().map(CidLink::from).collect(), - )), - _ => None, + Some(cids) if !cids.is_empty() => Ok(Some(EventBlocks::LegacyCids(column_vec( + cids, + col::REPO_SEQ_BLOCKS_CIDS, + )?))), + _ => Ok(None), } } @@ -97,24 +102,40 @@ fn map_sequenced_row(r: SequencedEventRow) -> Result { let blocks = row_to_event_blocks(r.block_cids, r.block_data, r.blocks_cids)?; Ok(SequencedEvent { seq: r.seq.into(), - did: Did::from(r.did), + did: column(r.did, col::REPO_SEQ_DID)?, created_at: r.created_at, event_type: r.event_type, - commit_cid: r.commit_cid.map(CidLink::from), - prev_cid: r.prev_cid.map(CidLink::from), - prev_data_cid: r.prev_data_cid.map(CidLink::from), + commit_cid: opt_column(r.commit_cid, col::REPO_SEQ_COMMIT_CID)?, + prev_cid: opt_column(r.prev_cid, col::REPO_SEQ_PREV_CID)?, + prev_data_cid: opt_column(r.prev_data_cid, col::REPO_SEQ_PREV_DATA_CID)?, ops: r.ops, blobs: r .blobs - .map(|blobs| blobs.into_iter().map(CidLink::from).collect()), + .map(|blobs| column_vec(blobs, col::REPO_SEQ_BLOBS)) + .transpose()?, blocks, - handle: r.handle.map(Handle::from), + handle: r + .handle + .and_then(|h| legacy_column(h, col::REPO_SEQ_HANDLE)), active: r.active, status, - rev: r.rev.map(Tid::from), + rev: opt_column(r.rev, col::REPO_SEQ_REV)?, }) } +fn collect_sequenced_rows(rows: Vec) -> Vec { + rows.into_iter() + .filter_map(|r| { + let seq = r.seq; + map_sequenced_row(r) + .inspect_err(|e| { + tracing::error!(seq, error = %e, "skipping a repo_seq row that doesn't decode"); + }) + .ok() + }) + .collect() +} + const SEQUENCER_LOCK_KEY: i64 = 0x0074_7261_6e73_6571; const SEQUENCER_BATCH_SIZE: i64 = 1000; @@ -266,7 +287,7 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(result.map(CidLink::from)) + opt_column(result, col::REPOS_REPO_ROOT_CID) } async fn get_repo(&self, user_id: Uuid) -> Result, DbError> { @@ -278,11 +299,14 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(row.map(|r| RepoInfo { - user_id: r.user_id, - repo_root_cid: CidLink::from(r.repo_root_cid), - repo_rev: r.repo_rev.map(Tid::from), - })) + row.map(|r| { + Ok(RepoInfo { + user_id: r.user_id, + repo_root_cid: column(r.repo_root_cid, col::REPOS_REPO_ROOT_CID)?, + repo_rev: opt_column(r.repo_rev, col::REPOS_REPO_REV)?, + }) + }) + .transpose() } async fn get_repo_root_by_did(&self, did: &Did) -> Result, DbError> { @@ -294,7 +318,7 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(result.map(CidLink::from)) + opt_column(result, col::REPOS_REPO_ROOT_CID) } async fn count_repos(&self) -> Result { @@ -312,13 +336,14 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(rows - .into_iter() - .map(|r| RepoWithoutRev { - user_id: r.user_id, - repo_root_cid: CidLink::from(r.repo_root_cid), + rows.into_iter() + .map(|r| { + Ok(RepoWithoutRev { + user_id: r.user_id, + repo_root_cid: column(r.repo_root_cid, col::REPOS_REPO_ROOT_CID)?, + }) }) - .collect()) + .collect() } async fn upsert_records( @@ -405,7 +430,7 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(result.map(CidLink::from)) + opt_column(result, col::RECORDS_RECORD_CID) } async fn list_records( @@ -418,13 +443,16 @@ impl RepoRepository for PostgresRepoRepository { rkey_start: Option<&Rkey>, rkey_end: Option<&Rkey>, ) -> Result, DbError> { - let to_record_info = |rows: Vec| { - rows.into_iter() - .map(|r| RecordInfo { - rkey: Rkey::from(r.rkey), - record_cid: CidLink::from(r.record_cid), + let to_record_info = |rows: Vec| -> Result, DbError> { + Ok(rows + .into_iter() + .filter_map(|r| { + Some(RecordInfo { + rkey: legacy_column(r.rkey, col::RECORDS_RKEY)?, + record_cid: legacy_column(r.record_cid, col::RECORDS_RECORD_CID)?, + }) }) - .collect() + .collect()) }; let collection_str = collection.as_str(); @@ -446,7 +474,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } true => { let rows = sqlx::query_as!( @@ -462,7 +490,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } }; } @@ -486,7 +514,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } true => { let rows = sqlx::query_as!( @@ -503,7 +531,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } }; } @@ -525,7 +553,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } true => { let rows = sqlx::query_as!( @@ -541,7 +569,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } }; } @@ -563,7 +591,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } true => { let rows = sqlx::query_as!( @@ -579,7 +607,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } }; } @@ -598,7 +626,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } true => { let rows = sqlx::query_as!( @@ -613,7 +641,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(to_record_info(rows)) + to_record_info(rows) } } } @@ -629,10 +657,12 @@ impl RepoRepository for PostgresRepoRepository { Ok(rows .into_iter() - .map(|r| FullRecordInfo { - collection: Nsid::from(r.collection), - rkey: Rkey::from(r.rkey), - record_cid: CidLink::from(r.record_cid), + .filter_map(|r| { + Some(FullRecordInfo { + collection: legacy_column(r.collection, col::RECORDS_COLLECTION)?, + rkey: legacy_column(r.rkey, col::RECORDS_RKEY)?, + record_cid: legacy_column(r.record_cid, col::RECORDS_RECORD_CID)?, + }) }) .collect()) } @@ -646,7 +676,10 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(rows.into_iter().map(Nsid::from).collect()) + Ok(rows + .into_iter() + .filter_map(|c| legacy_column(c, col::RECORDS_COLLECTION)) + .collect()) } async fn count_records(&self, repo_id: Uuid) -> Result { @@ -1024,13 +1057,16 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(row.map(|r| RepoAccountInfo { - user_id: r.id, - did: Did::from(r.did), - deactivated_at: r.deactivated_at, - takedown_ref: r.takedown_ref, - repo_root_cid: r.repo_root_cid.map(CidLink::from), - })) + row.map(|r| { + Ok(RepoAccountInfo { + user_id: r.id, + did: column(r.did, col::USERS_DID)?, + deactivated_at: r.deactivated_at, + takedown_ref: r.takedown_ref, + repo_root_cid: opt_column(r.repo_root_cid, col::REPOS_REPO_ROOT_CID)?, + }) + }) + .transpose() } async fn get_events_since_seq( @@ -1054,7 +1090,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - rows.into_iter().map(map_sequenced_row).collect() + Ok(collect_sequenced_rows(rows)) } None => { let rows = sqlx::query_as!( @@ -1069,7 +1105,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - rows.into_iter().map(map_sequenced_row).collect() + Ok(collect_sequenced_rows(rows)) } } } @@ -1092,7 +1128,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - rows.into_iter().map(map_sequenced_row).collect() + Ok(collect_sequenced_rows(rows)) } async fn get_event_by_seq( @@ -1132,7 +1168,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; - rows.into_iter().map(map_sequenced_row).collect() + Ok(collect_sequenced_rows(rows)) } async fn list_repos_paginated( @@ -1155,16 +1191,17 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(rows - .into_iter() - .map(|r| RepoListItem { - did: Did::from(r.did), - deactivated_at: r.deactivated_at, - takedown_ref: r.takedown_ref, - repo_root_cid: CidLink::from(r.repo_root_cid), - repo_rev: r.repo_rev.map(Tid::from), + rows.into_iter() + .map(|r| { + Ok(RepoListItem { + did: column(r.did, col::USERS_DID)?, + deactivated_at: r.deactivated_at, + takedown_ref: r.takedown_ref, + repo_root_cid: column(r.repo_root_cid, col::REPOS_REPO_ROOT_CID)?, + repo_rev: opt_column(r.repo_rev, col::REPOS_REPO_REV)?, + }) }) - .collect()) + .collect() } async fn get_repo_root_cid_by_user_id( @@ -1178,7 +1215,7 @@ impl RepoRepository for PostgresRepoRepository { .fetch_optional(&self.pool) .await .map_err(map_sqlx_error)?; - Ok(cid.map(CidLink::from)) + opt_column(cid, col::REPOS_REPO_ROOT_CID) } async fn import_repo_data( @@ -1520,14 +1557,15 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(rows - .into_iter() - .map(|(user_id, repo_root_cid, repo_rev)| UserWithoutBlocks { - user_id, - repo_root_cid: CidLink::from(repo_root_cid), - repo_rev: repo_rev.map(Tid::from), + rows.into_iter() + .map(|(user_id, repo_root_cid, repo_rev)| { + Ok(UserWithoutBlocks { + user_id, + repo_root_cid: column(repo_root_cid, col::REPOS_REPO_ROOT_CID)?, + repo_rev: opt_column(repo_rev, col::REPOS_REPO_REV)?, + }) }) - .collect()) + .collect() } async fn get_users_needing_record_blobs_backfill( @@ -1548,13 +1586,14 @@ impl RepoRepository for PostgresRepoRepository { .await .map_err(map_sqlx_error)?; - Ok(rows - .into_iter() - .map(|r| UserNeedingRecordBlobsBackfill { - user_id: r.user_id, - did: Did::from(r.did), + rows.into_iter() + .map(|r| { + Ok(UserNeedingRecordBlobsBackfill { + user_id: r.user_id, + did: column(r.did, col::USERS_DID)?, + }) }) - .collect()) + .collect() } async fn insert_record_blobs(