store: typed revs thru metastore keys & requests

Lewis: May this revision serve well! <lu5a@proton.me>
This commit is contained in:
Lewis
2026-07-25 08:27:39 +03:00
committed by Tangled
parent bbe9f6f3b3
commit 4f37ac26cd
30 changed files with 751 additions and 510 deletions
@@ -0,0 +1,23 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT block_cid AS \"block_cid!\" FROM user_blocks\n WHERE user_id = $1 AND repo_rev > $2\n ORDER BY repo_rev ASC\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "block_cid!",
"type_info": "Bytea"
}
],
"parameters": {
"Left": [
"Uuid",
"Text"
]
},
"nullable": [
false
]
},
"hash": "30570ed3866840d1258c8768a5c8a23ade40700c05ddbbf7fc4f64bfa95b1ed4"
}
@@ -0,0 +1,22 @@
{
"db_name": "PostgreSQL",
"query": "\n SELECT block_cid AS \"block_cid!\" FROM user_blocks\n WHERE user_id = $1\n ORDER BY repo_rev ASC\n ",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "block_cid!",
"type_info": "Bytea"
}
],
"parameters": {
"Left": [
"Uuid"
]
},
"nullable": [
false
]
},
"hash": "85cc0cd1e62a30fa67d415b7a01164f962a422513e8f0737553321fd9987a56c"
}
+1 -1
View File
@@ -432,7 +432,7 @@ pub trait RepoRepository: Send + Sync {
async fn get_user_block_cids_since_rev(
&self,
user_id: Uuid,
since_rev: &Tid,
since_rev: Option<&Tid>,
) -> Result<Vec<Vec<u8>>, DbError>;
async fn count_user_blocks(&self, user_id: Uuid) -> Result<i64, DbError>;
+29 -14
View File
@@ -793,22 +793,37 @@ impl RepoRepository for PostgresRepoRepository {
async fn get_user_block_cids_since_rev(
&self,
user_id: Uuid,
since_rev: &Tid,
since_rev: Option<&Tid>,
) -> Result<Vec<Vec<u8>>, DbError> {
let rows: Vec<(Vec<u8>,)> = sqlx::query_as(
r#"
SELECT block_cid FROM user_blocks
WHERE user_id = $1 AND repo_rev > $2
ORDER BY repo_rev ASC
"#,
)
.bind(user_id)
.bind(since_rev.as_str())
.fetch_all(&self.pool)
.await
.map_err(map_sqlx_error)?;
let rows = match since_rev {
None => {
sqlx::query_scalar!(
r#"
SELECT block_cid AS "block_cid!" FROM user_blocks
WHERE user_id = $1
ORDER BY repo_rev ASC
"#,
user_id
)
.fetch_all(&self.pool)
.await
}
Some(rev) => {
sqlx::query_scalar!(
r#"
SELECT block_cid AS "block_cid!" FROM user_blocks
WHERE user_id = $1 AND repo_rev > $2
ORDER BY repo_rev ASC
"#,
user_id,
rev.as_str()
)
.fetch_all(&self.pool)
.await
}
};
Ok(rows.into_iter().map(|(cid,)| cid).collect())
rows.map_err(map_sqlx_error)
}
async fn insert_commit_event(&self, data: &CommitEventData) -> Result<(), DbError> {
@@ -144,7 +144,7 @@ async fn repair_fails_loud_on_missing_leaf_block() {
let recorded = repos
.repo
.get_user_block_cids_since_rev(user_id, &tranquil_types::Tid::from(String::new()))
.get_user_block_cids_since_rev(user_id, None)
.await
.expect("read user_blocks");
assert!(
+56 -52
View File
@@ -86,7 +86,22 @@ fn test_cid_bytes(seed: u8) -> Vec<u8> {
}
fn make_rev(n: u64) -> Tid {
Tid::from(format!("rev{n:010}"))
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let value = n << 40;
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((value >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
fn post_collection() -> Nsid {
Nsid::new("app.bsky.feed.post").expect("app.bsky.feed.post is a valid NSID")
}
fn bench_handle(user_id: Uuid) -> Handle {
Handle::new(format!("bench{}.oyster.cafe", user_id.as_simple()))
.expect("generated handle is valid")
}
struct BenchHarness {
@@ -141,9 +156,9 @@ async fn create_user(pool: &HandlerPool, user_id: Uuid, did: &Did, cid: &CidLink
pool.send(MetastoreRequest::Repo(RepoRequest::CreateRepoFull {
user_id,
did: did.clone(),
handle: Handle::from(format!("bench.{}.invalid", user_id.as_simple())),
handle: bench_handle(user_id),
repo_root_cid: cid.clone(),
repo_rev: Tid::from("rev0000000000".to_string()),
repo_rev: make_rev(0),
tx,
}))
.unwrap();
@@ -167,7 +182,7 @@ fn make_commit_input(
obsolete_block_cids: vec![],
record_upserts: vec![RecordUpsert {
collection: collection.clone(),
rkey: Rkey::from(format!("r{rev_n:010}")),
rkey: Rkey::new(format!("r{rev_n:010}")).expect("generated rkey is valid"),
cid: test_cid(cid_seed),
}],
record_deletes: vec![],
@@ -213,7 +228,7 @@ async fn seed_records(
let record_upserts: Vec<RecordUpsert> = (batch_start..batch_end)
.map(|i| RecordUpsert {
collection: collection.clone(),
rkey: Rkey::from(format!("rec{i:08}")),
rkey: Rkey::new(format!("rec{i:08}")).expect("generated rkey is valid"),
cid: test_cid(((i * 7 + 3) & 0xFF) as u8),
})
.collect();
@@ -265,7 +280,9 @@ async fn bench_apply_commit(pool: &Arc<HandlerPool>, concurrency: usize, ops_per
let user_ids: Vec<Uuid> = (0..concurrency).map(|_| Uuid::new_v4()).collect();
let dids: Vec<Did> = user_ids
.iter()
.map(|u| Did::from(format!("did:plc:bench{}", u.as_simple())))
.map(|u| {
Did::new(format!("did:plc:bench{}", u.as_simple())).expect("generated DID is valid")
})
.collect();
futures::stream::iter(user_ids.iter().zip(dids.iter()))
@@ -274,7 +291,7 @@ async fn bench_apply_commit(pool: &Arc<HandlerPool>, concurrency: usize, ops_per
})
.await;
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let start = Instant::now();
let handles: Vec<_> = (0..concurrency)
.map(|task_id| {
@@ -321,8 +338,9 @@ async fn bench_apply_commit(pool: &Arc<HandlerPool>, concurrency: usize, ops_per
async fn bench_get_record_cid(pool: &Arc<HandlerPool>, concurrency: usize, ops_per_task: usize) {
let user_id = Uuid::new_v4();
let did = Did::from(format!("did:plc:getrecord{}", user_id.as_simple()));
let collection = Nsid::from("app.bsky.feed.post".to_string());
let did = Did::new(format!("did:plc:getrecord{}", user_id.as_simple()))
.expect("generated DID is valid");
let collection = post_collection();
create_user(pool, user_id, &did, &test_cid(1)).await;
seed_records(pool, user_id, &did, &collection, 1000).await;
@@ -339,7 +357,8 @@ async fn bench_get_record_cid(pool: &Arc<HandlerPool>, concurrency: usize, ops_p
let collection = &collection;
async move {
let rec_idx = (task_id * 7 + i * 13) % total_records;
let rkey = Rkey::from(format!("rec{rec_idx:08}"));
let rkey = Rkey::new(format!("rec{rec_idx:08}"))
.expect("generated rkey is valid");
let t = Instant::now();
let (tx, rx) = oneshot::channel();
pool.send(MetastoreRequest::Record(RecordRequest::GetRecordCid {
@@ -369,8 +388,9 @@ async fn bench_get_record_cid(pool: &Arc<HandlerPool>, concurrency: usize, ops_p
async fn bench_list_records(pool: &Arc<HandlerPool>, concurrency: usize, ops_per_task: usize) {
let user_id = Uuid::new_v4();
let did = Did::from(format!("did:plc:listrecords{}", user_id.as_simple()));
let collection = Nsid::from("app.bsky.feed.post".to_string());
let did = Did::new(format!("did:plc:listrecords{}", user_id.as_simple()))
.expect("generated DID is valid");
let collection = post_collection();
create_user(pool, user_id, &did, &test_cid(1)).await;
seed_records(pool, user_id, &did, &collection, 1000).await;
@@ -485,7 +505,7 @@ async fn setup_pg_bench_schema(pool: &sqlx::PgPool) {
async fn pg_create_user(pg: &sqlx::PgPool, user_id: Uuid, did: &str) {
sqlx::query("INSERT INTO users (id, handle, did) VALUES ($1, $2, $3) ON CONFLICT DO NOTHING")
.bind(user_id)
.bind(format!("bench.{}.invalid", Uuid::new_v4().as_simple()))
.bind(format!("bench{}.oyster.cafe", Uuid::new_v4().as_simple()))
.bind(did)
.execute(pg)
.await
@@ -507,21 +527,15 @@ async fn bench_pg_upsert_records(
futures::stream::iter(user_ids.iter().zip(dids.iter()))
.fold((), |(), (uid, did)| async {
pg_create_user(pg, *uid, did).await;
let did = Did::from(did.clone());
let handle = Handle::from(format!("bench.{}.invalid", uid.as_simple()));
repo.create_repo(
*uid,
&did,
&handle,
&test_cid(1),
&Tid::from("rev0000000000".to_string()),
)
.await
.unwrap();
let did = Did::new(did.clone()).expect("generated DID is valid");
let handle = bench_handle(*uid);
repo.create_repo(*uid, &did, &handle, &test_cid(1), &make_rev(0))
.await
.unwrap();
})
.await;
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let start = Instant::now();
let handles: Vec<_> = (0..concurrency)
.map(|task_id| {
@@ -537,7 +551,8 @@ async fn bench_pg_upsert_records(
async move {
let rev_n = (task_id * ops_per_task + i + 1) as u64;
let cid_seed = ((task_id * 31 + i * 7) & 0xFF) as u8;
let rkey = Rkey::from(format!("r{rev_n:010}"));
let rkey = Rkey::new(format!("r{rev_n:010}"))
.expect("generated rkey is valid");
let t = Instant::now();
repo.upsert_records(
user_id,
@@ -587,7 +602,7 @@ async fn pg_seed_records(
.map(|_| collection.clone())
.collect();
let rkeys: Vec<Rkey> = (batch_start..batch_end)
.map(|i| Rkey::from(format!("rec{i:08}")))
.map(|i| Rkey::new(format!("rec{i:08}")).expect("generated rkey is valid"))
.collect();
let cids: Vec<CidLink> = (batch_start..batch_end)
.map(|i| test_cid(((i * 7 + 3) & 0xFF) as u8))
@@ -609,20 +624,14 @@ async fn pg_seed_records(
async fn bench_pg_get_record_cid(pg: &sqlx::PgPool, concurrency: usize, ops_per_task: usize) {
let user_id = Uuid::new_v4();
let did = format!("did:plc:pgget{}", user_id.as_simple());
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
pg_create_user(pg, user_id, &did).await;
let repo = tranquil_db::postgres::PostgresRepoRepository::new(pg.clone());
let did = Did::from(did);
let handle = Handle::from(format!("bench.{}.invalid", user_id.as_simple()));
repo.create_repo(
user_id,
&did,
&handle,
&test_cid(1),
&Tid::from("rev0000000000".to_string()),
)
.await
.unwrap();
let did = Did::new(did).expect("generated DID is valid");
let handle = bench_handle(user_id);
repo.create_repo(user_id, &did, &handle, &test_cid(1), &make_rev(0))
.await
.unwrap();
pg_seed_records(&repo, user_id, &collection, 1000).await;
let total_records = 1000usize;
@@ -639,7 +648,8 @@ async fn bench_pg_get_record_cid(pg: &sqlx::PgPool, concurrency: usize, ops_per_
let collection = &collection;
async move {
let rec_idx = (task_id * 7 + i * 13) % total_records;
let rkey = Rkey::from(format!("rec{rec_idx:08}"));
let rkey = Rkey::new(format!("rec{rec_idx:08}"))
.expect("generated rkey is valid");
let t = Instant::now();
let result = repo
.get_record_cid(user_id, collection, &rkey)
@@ -665,20 +675,14 @@ async fn bench_pg_get_record_cid(pg: &sqlx::PgPool, concurrency: usize, ops_per_
async fn bench_pg_list_records(pg: &sqlx::PgPool, concurrency: usize, ops_per_task: usize) {
let user_id = Uuid::new_v4();
let did = format!("did:plc:pglist{}", user_id.as_simple());
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
pg_create_user(pg, user_id, &did).await;
let repo = tranquil_db::postgres::PostgresRepoRepository::new(pg.clone());
let did = Did::from(did);
let handle = Handle::from(format!("bench.{}.invalid", user_id.as_simple()));
repo.create_repo(
user_id,
&did,
&handle,
&test_cid(1),
&Tid::from("rev0000000000".to_string()),
)
.await
.unwrap();
let did = Did::new(did).expect("generated DID is valid");
let handle = bench_handle(user_id);
repo.create_repo(user_id, &did, &handle, &test_cid(1), &make_rev(0))
.await
.unwrap();
pg_seed_records(&repo, user_id, &collection, 1000).await;
let start = Instant::now();
@@ -86,7 +86,17 @@ fn test_cid_bytes(seed: u8) -> Vec<u8> {
}
fn make_rev(n: u64) -> Tid {
Tid::from(format!("rev{n:010}"))
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let value = n << 40;
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((value >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
fn post_collection() -> Nsid {
Nsid::new("app.bsky.feed.post").expect("app.bsky.feed.post is a valid NSID")
}
struct BenchHarness {
@@ -173,7 +183,8 @@ async fn seed_users(pool: &HandlerPool, count: usize) -> Vec<UserInfo> {
.map(|i| {
let user_id = Uuid::new_v4();
UserInfo {
did: Did::from(format!("did:plc:scale{i:06x}{}", user_id.as_simple())),
did: Did::new(format!("did:plc:scale{i:06x}{}", user_id.as_simple()))
.expect("generated DID is valid"),
user_id,
}
})
@@ -192,12 +203,10 @@ async fn seed_users(pool: &HandlerPool, count: usize) -> Vec<UserInfo> {
pool.send(MetastoreRequest::Repo(RepoRequest::CreateRepoFull {
user_id: user.user_id,
did: user.did.clone(),
handle: Handle::from(format!(
"u{}.scale.invalid",
user.user_id.as_simple()
)),
handle: Handle::new(format!("u{}.oyster.cafe", user.user_id.as_simple()))
.expect("generated handle is valid"),
repo_root_cid: test_cid(1),
repo_rev: Tid::from("rev0000000000".to_string()),
repo_rev: make_rev(0),
tx,
}))
.unwrap();
@@ -225,11 +234,11 @@ async fn seed_users(pool: &HandlerPool, count: usize) -> Vec<UserInfo> {
}
async fn seed_records_for_user(pool: &HandlerPool, user: &UserInfo, record_count: usize) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let record_upserts: Vec<RecordUpsert> = (0..record_count)
.map(|i| RecordUpsert {
collection: collection.clone(),
rkey: Rkey::from(format!("rec{i:08}")),
rkey: Rkey::new(format!("rec{i:08}")).expect("generated rkey is valid"),
cid: test_cid(((i * 7 + 3) & 0xFF) as u8),
})
.collect();
@@ -313,7 +322,7 @@ async fn seed_all_records(pool: &Arc<HandlerPool>, users: &[UserInfo], records_p
}
async fn bench_single_user_commit(pool: &Arc<HandlerPool>, user: &UserInfo, ops: usize) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let start = Instant::now();
let mut latencies: Vec<Duration> = Vec::with_capacity(ops);
@@ -335,7 +344,8 @@ async fn bench_single_user_commit(pool: &Arc<HandlerPool>, user: &UserInfo, ops:
obsolete_block_cids: vec![],
record_upserts: vec![RecordUpsert {
collection: collection.clone(),
rkey: Rkey::from(format!("new{rev_n:010}")),
rkey: Rkey::new(format!("new{rev_n:010}"))
.expect("generated rkey is valid"),
cid: test_cid(cid_seed),
}],
record_deletes: vec![],
@@ -380,7 +390,7 @@ async fn bench_multi_user_commit(
concurrency: usize,
ops_per_task: usize,
) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let active_users: Vec<&UserInfo> = users.iter().take(concurrency).collect();
let start = Instant::now();
@@ -411,7 +421,8 @@ async fn bench_multi_user_commit(
obsolete_block_cids: vec![],
record_upserts: vec![RecordUpsert {
collection: collection.clone(),
rkey: Rkey::from(format!("mu{rev_n:010}")),
rkey: Rkey::new(format!("mu{rev_n:010}"))
.expect("generated rkey is valid"),
cid: test_cid(cid_seed),
}],
record_deletes: vec![],
@@ -466,7 +477,7 @@ async fn bench_list_records_at_scale(
concurrency: usize,
ops_per_task: usize,
) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let user_count = users.len();
let start = Instant::now();
@@ -526,7 +537,7 @@ async fn bench_get_record_at_scale(
concurrency: usize,
ops_per_task: usize,
) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let user_count = users.len();
let records_per_user = 10usize;
@@ -545,7 +556,8 @@ async fn bench_get_record_at_scale(
async move {
let user_idx = (task_id * 997 + i * 31) % user_count;
let rec_idx = (task_id * 13 + i * 7) % records_per_user;
let rkey = Rkey::from(format!("rec{rec_idx:08}"));
let rkey = Rkey::new(format!("rec{rec_idx:08}"))
.expect("generated rkey is valid");
let t = Instant::now();
let (tx, rx) = oneshot::channel();
pool.send(MetastoreRequest::Record(RecordRequest::GetRecordCid {
@@ -582,7 +594,8 @@ async fn pg_seed_users(pg: &sqlx::PgPool, count: usize) -> Vec<UserInfo> {
.map(|i| {
let user_id = Uuid::new_v4();
UserInfo {
did: Did::from(format!("did:plc:pgscale{i:06x}{}", user_id.as_simple())),
did: Did::new(format!("did:plc:pgscale{i:06x}{}", user_id.as_simple()))
.expect("generated DID is valid"),
user_id,
}
})
@@ -601,23 +614,22 @@ async fn pg_seed_users(pg: &sqlx::PgPool, count: usize) -> Vec<UserInfo> {
"INSERT INTO users (id, handle, did) VALUES ($1, $2, $3) ON CONFLICT DO NOTHING",
)
.bind(user.user_id)
.bind(format!("u{}.pgscale.invalid", user.user_id.as_simple()))
.bind(format!("pg{}.oyster.cafe", user.user_id.as_simple()))
.bind(user.did.as_str())
.execute(pg)
.await
.unwrap();
let repo = tranquil_db::postgres::PostgresRepoRepository::new(pg.clone());
let handle = Handle::from(format!(
"u{}.pgscale.invalid",
user.user_id.as_simple()
));
let handle =
Handle::new(format!("pg{}.oyster.cafe", user.user_id.as_simple()))
.expect("generated handle is valid");
repo.create_repo(
user.user_id,
&user.did,
&handle,
&test_cid(1),
&Tid::from("rev0000000000".to_string()),
&make_rev(0),
)
.await
.unwrap();
@@ -646,7 +658,7 @@ async fn pg_seed_users(pg: &sqlx::PgPool, count: usize) -> Vec<UserInfo> {
async fn pg_seed_all_records(pg: &sqlx::PgPool, users: &[UserInfo], records_per_user: usize) {
let start = Instant::now();
let total = users.len();
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let chunk_size = 500;
let chunks: Vec<&[UserInfo]> = users.chunks(chunk_size).collect();
let total_chunks = chunks.len();
@@ -665,7 +677,10 @@ async fn pg_seed_all_records(pg: &sqlx::PgPool, users: &[UserInfo], records_per_
let collections: Vec<Nsid> =
(0..records_per_user).map(|_| collection.clone()).collect();
let rkeys: Vec<Rkey> = (0..records_per_user)
.map(|i| Rkey::from(format!("rec{i:08}")))
.map(|i| {
Rkey::new(format!("rec{i:08}"))
.expect("generated rkey is valid")
})
.collect();
let cids: Vec<CidLink> = (0..records_per_user)
.map(|i| test_cid(((i * 7 + 3) & 0xFF) as u8))
@@ -675,7 +690,7 @@ async fn pg_seed_all_records(pg: &sqlx::PgPool, users: &[UserInfo], records_per_
&collections,
&rkeys,
&cids,
&Tid::from("rev0000000001".to_string()),
&make_rev(1),
)
.await
.unwrap();
@@ -702,7 +717,7 @@ async fn pg_seed_all_records(pg: &sqlx::PgPool, users: &[UserInfo], records_per_
async fn bench_pg_single_user_commit(pg: &sqlx::PgPool, user: &UserInfo, ops: usize) {
let repo = tranquil_db::postgres::PostgresRepoRepository::new(pg.clone());
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let start = Instant::now();
let mut latencies: Vec<Duration> = Vec::with_capacity(ops);
@@ -714,7 +729,7 @@ async fn bench_pg_single_user_commit(pg: &sqlx::PgPool, user: &UserInfo, ops: us
async move {
let rev_n = (i + 100) as u64;
let cid_seed = ((i * 7 + 42) & 0xFF) as u8;
let rkey = Rkey::from(format!("new{rev_n:010}"));
let rkey = Rkey::new(format!("new{rev_n:010}")).expect("generated rkey is valid");
let t = Instant::now();
repo.upsert_records(
user.user_id,
@@ -742,7 +757,7 @@ async fn bench_pg_multi_user_commit(
concurrency: usize,
ops_per_task: usize,
) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let active_users: Vec<&UserInfo> = users.iter().take(concurrency).collect();
let start = Instant::now();
@@ -762,7 +777,8 @@ async fn bench_pg_multi_user_commit(
async move {
let rev_n = (task_id * ops_per_task + i + 200) as u64;
let cid_seed = ((task_id * 31 + i * 7) & 0xFF) as u8;
let rkey = Rkey::from(format!("mu{rev_n:010}"));
let rkey = Rkey::new(format!("mu{rev_n:010}"))
.expect("generated rkey is valid");
let t = Instant::now();
repo.upsert_records(
user_id,
@@ -800,7 +816,7 @@ async fn bench_pg_list_records(
concurrency: usize,
ops_per_task: usize,
) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let user_count = users.len();
let start = Instant::now();
@@ -858,7 +874,7 @@ async fn bench_pg_get_record(
concurrency: usize,
ops_per_task: usize,
) {
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
let user_count = users.len();
let records_per_user = 10usize;
@@ -878,7 +894,8 @@ async fn bench_pg_get_record(
async move {
let user_idx = (task_id * 997 + i * 31) % user_count;
let rec_idx = (task_id * 13 + i * 7) % records_per_user;
let rkey = Rkey::from(format!("rec{rec_idx:08}"));
let rkey = Rkey::new(format!("rec{rec_idx:08}"))
.expect("generated rkey is valid");
let t = Instant::now();
let _result = repo
.get_record_cid(user_ids[user_idx], collection, &rkey)
+31 -10
View File
@@ -27,6 +27,20 @@ fn test_cid_bytes(seed: u8) -> Vec<u8> {
cid::Cid::new_v1(0x71, mh).to_bytes()
}
fn test_rev(sequence: u64, discriminator: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let value = (sequence << 40) | (discriminator & 0xFF_FFFF_FFFF);
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((value >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
fn post_collection() -> Nsid {
Nsid::new("app.bsky.feed.post").expect("app.bsky.feed.post is a valid NSID")
}
struct UserInfo {
user_id: Uuid,
did: Did,
@@ -37,7 +51,8 @@ async fn seed_users(pool: &HandlerPool, count: usize) -> Vec<UserInfo> {
.map(|i| {
let user_id = Uuid::new_v4();
UserInfo {
did: Did::from(format!("did:plc:prof{i:06x}{}", user_id.as_simple())),
did: Did::new(format!("did:plc:squid{i:06x}{}", user_id.as_simple()))
.expect("generated DID is valid"),
user_id,
}
})
@@ -54,9 +69,13 @@ async fn seed_users(pool: &HandlerPool, count: usize) -> Vec<UserInfo> {
pool.send(MetastoreRequest::Repo(RepoRequest::CreateRepoFull {
user_id: user.user_id,
did: user.did.clone(),
handle: Handle::from(format!("u{}.prof.invalid", user.user_id.as_simple())),
handle: Handle::new(format!(
"squid{}.oyster.cafe",
user.user_id.as_simple()
))
.expect("generated handle is valid"),
repo_root_cid: test_cid(1),
repo_rev: Tid::from("rev0000000000".to_string()),
repo_rev: test_rev(0, 0),
tx,
}))
.unwrap();
@@ -87,7 +106,7 @@ async fn seed_records(pool: &Arc<HandlerPool>, users: &[UserInfo], records_per_u
let total = users.len();
let batch_size = 500;
let total_batches = total.div_ceil(batch_size);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = post_collection();
futures::stream::iter(users.chunks(batch_size).enumerate())
.fold((), |(), (chunk_idx, chunk)| {
@@ -102,7 +121,8 @@ async fn seed_records(pool: &Arc<HandlerPool>, users: &[UserInfo], records_per_u
let record_upserts: Vec<RecordUpsert> = (0..records_per_user)
.map(|i| RecordUpsert {
collection: collection.clone(),
rkey: Rkey::from(format!("rec{i:08}")),
rkey: Rkey::new(format!("rec{i:08}"))
.expect("generated rkey is valid"),
cid: test_cid(((i * 7 + 3) & 0xFF) as u8),
})
.collect();
@@ -114,7 +134,7 @@ async fn seed_records(pool: &Arc<HandlerPool>, users: &[UserInfo], records_per_u
did: user.did.clone(),
expected_root_cid: None,
new_root_cid: test_cid(2),
new_rev: Tid::from("rev0000000001".to_string()),
new_rev: test_rev(1, 0),
new_block_cids,
obsolete_block_cids: vec![],
record_upserts,
@@ -130,7 +150,7 @@ async fn seed_records(pool: &Arc<HandlerPool>, users: &[UserInfo], records_per_u
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(Tid::from("rev0000000001".to_string())),
rev: Some(test_rev(1, 0)),
},
};
let (tx, rx) = oneshot::channel();
@@ -190,7 +210,7 @@ async fn profile_list_records(
let (tx, rx) = oneshot::channel();
pool.send(MetastoreRequest::Record(RecordRequest::ListRecords {
repo_id: user_ids[idx],
collection: Nsid::from("app.bsky.feed.post".to_string()),
collection: post_collection(),
cursor: None,
limit: 50,
reverse: false,
@@ -238,11 +258,12 @@ async fn profile_get_record_cid(
async move {
let user_idx = (task_id * 997 + i * 31) % user_count;
let rec_idx = (task_id * 13 + i * 7) % records_per_user;
let rkey = Rkey::from(format!("rec{rec_idx:08}"));
let rkey =
Rkey::new(format!("rec{rec_idx:08}")).expect("generated rkey is valid");
let (tx, rx) = oneshot::channel();
pool.send(MetastoreRequest::Record(RecordRequest::GetRecordCid {
repo_id: user_ids[user_idx],
collection: Nsid::from("app.bsky.feed.post".to_string()),
collection: post_collection(),
rkey,
tx,
}))
@@ -243,6 +243,15 @@ mod tests {
tranquil_types::CidLink::from_cid(&c)
}
fn test_rev(seq: u64) -> tranquil_types::Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
tranquil_types::Tid::new(s).expect("generated TID is valid")
}
fn create_repo(h: &TestHarness, name: &str, seed: u8) -> (Uuid, UserHash) {
let user_id = Uuid::new_v4();
let did = Did::from(format!("did:plc:{name}"));
@@ -250,7 +259,14 @@ mod tests {
let cid = test_cid_link(seed);
h.metastore
.repo_ops()
.create_repo(h.metastore.database(), user_id, &did, &handle, &cid, "rev0")
.create_repo(
h.metastore.database(),
user_id,
&did,
&handle,
&cid,
&test_rev(0),
)
.unwrap();
let user_hash = h.metastore.user_hashes().get(&user_id).unwrap();
(user_id, user_hash)
@@ -185,7 +185,7 @@ impl<S: StorageIO + 'static> tranquil_db_traits::RepoRepository for MetastoreCli
.send(MetastoreRequest::Repo(RepoRequest::UpdateRepoRoot {
user_id,
repo_root_cid: repo_root_cid.clone(),
repo_rev: repo_rev.to_string(),
repo_rev: repo_rev.clone(),
tx,
}))?;
recv(rx).await
@@ -196,7 +196,7 @@ impl<S: StorageIO + 'static> tranquil_db_traits::RepoRepository for MetastoreCli
self.pool
.send(MetastoreRequest::Repo(RepoRequest::UpdateRepoRev {
user_id,
repo_rev: repo_rev.to_string(),
repo_rev: repo_rev.clone(),
tx,
}))?;
recv(rx).await
@@ -270,7 +270,7 @@ impl<S: StorageIO + 'static> tranquil_db_traits::RepoRepository for MetastoreCli
collections: collections.to_vec(),
rkeys: rkeys.to_vec(),
record_cids: record_cids.to_vec(),
repo_rev: repo_rev.to_string(),
repo_rev: repo_rev.clone(),
tx,
}))?;
recv(rx).await
@@ -427,7 +427,7 @@ impl<S: StorageIO + 'static> tranquil_db_traits::RepoRepository for MetastoreCli
UserBlockRequest::InsertUserBlocks {
user_id,
block_cids: block_cids.to_vec(),
repo_rev: repo_rev.to_string(),
repo_rev: repo_rev.clone(),
tx,
},
))?;
@@ -453,13 +453,13 @@ impl<S: StorageIO + 'static> tranquil_db_traits::RepoRepository for MetastoreCli
async fn get_user_block_cids_since_rev(
&self,
user_id: Uuid,
since_rev: &Tid,
since_rev: Option<&Tid>,
) -> Result<Vec<Vec<u8>>, DbError> {
let (tx, rx) = oneshot::channel();
self.pool.send(MetastoreRequest::UserBlock(
UserBlockRequest::GetUserBlockCidsSinceRev {
user_id,
since_rev: since_rev.to_string(),
since_rev: since_rev.cloned(),
tx,
},
))?;
@@ -522,7 +522,7 @@ impl<S: StorageIO + 'static> tranquil_db_traits::RepoRepository for MetastoreCli
.send(MetastoreRequest::Event(EventRequest::InsertSyncEvent {
did: did.clone(),
commit_cid: commit_cid.clone(),
rev: rev.map(|r| r.to_string()),
rev: rev.cloned(),
commit_bytes: commit_bytes.to_vec(),
tx,
}))?;
@@ -544,7 +544,7 @@ impl<S: StorageIO + 'static> tranquil_db_traits::RepoRepository for MetastoreCli
did: did.clone(),
commit_cid: commit_cid.clone(),
mst_root_cid: mst_root_cid.clone(),
rev: rev.to_string(),
rev: rev.clone(),
commit_bytes: commit_bytes.to_vec(),
mst_root_bytes: mst_root_bytes.to_vec(),
tx,
@@ -916,7 +916,7 @@ impl<S: StorageIO + 'static> tranquil_db_traits::BlobRepository for MetastoreCli
self.pool
.send(MetastoreRequest::Blob(BlobRequest::ListBlobsSinceRev {
did: did.clone(),
since: since.to_string(),
since: since.clone(),
tx,
}))?;
recv(rx).await
@@ -203,7 +203,7 @@ impl<S: StorageIO + 'static> CommitOps<S> {
let mutation_set = CommitMutationSet {
new_root_cid: new_cid_bytes.clone(),
new_rev: input.new_rev.to_string(),
new_rev: input.new_rev.clone(),
record_upserts: input
.record_upserts
.iter()
@@ -534,6 +534,15 @@ mod tests {
Handle::from(format!("{name}.test.invalid"))
}
fn test_rev(seq: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
fn make_commit_ops(h: &TestHarness) -> CommitOps<RealIO> {
use crate::metastore::partitions::Partition;
CommitOps::new(
@@ -552,7 +561,14 @@ mod tests {
let cid = test_cid_link(seed);
h.metastore
.repo_ops()
.create_repo(h.metastore.database(), user_id, &did, &handle, &cid, "rev0")
.create_repo(
h.metastore.database(),
user_id,
&did,
&handle,
&cid,
&test_rev(0),
)
.unwrap();
(user_id, did, cid)
}
@@ -1110,7 +1126,7 @@ mod tests {
&did,
&handle,
&initial_root,
"rev0",
&test_rev(0),
)
.unwrap();
metastore.persist().unwrap();
@@ -1235,7 +1251,7 @@ mod tests {
&did,
&handle,
&initial_root,
"rev0",
&test_rev(0),
)
.unwrap();
metastore.persist().unwrap();
@@ -1243,7 +1259,7 @@ mod tests {
let event_ops = metastore.event_ops(Arc::clone(&bridge));
let mutation_set = super::CommitMutationSet {
new_root_cid: super::cid_link_to_bytes(&new_root).unwrap(),
new_rev: "rev1".to_string(),
new_rev: test_rev(1),
record_upserts: vec![super::RecordMutationUpsert {
collection: collection.as_str().to_owned(),
rkey: rkey.as_str().to_owned(),
@@ -2,8 +2,9 @@ use smallvec::SmallVec;
use super::encoding::KeyBuilder;
use super::keys::{KeyTag, UserHash};
use tranquil_types::Tid;
pub fn rev_to_seq_key(user_hash: UserHash, rev: &str) -> SmallVec<[u8; 128]> {
pub fn rev_to_seq_key(user_hash: UserHash, rev: &Tid) -> SmallVec<[u8; 128]> {
KeyBuilder::new()
.tag(KeyTag::REV_TO_SEQ)
.u64(user_hash.raw())
@@ -52,14 +53,24 @@ mod tests {
use super::*;
use crate::metastore::encoding::KeyReader;
fn test_rev(seq: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
#[test]
fn rev_to_seq_key_roundtrip() {
let hash = UserHash::from_raw(0xDEAD_BEEF_CAFE_BABE);
let key = rev_to_seq_key(hash, "3k2abcde");
let rev = test_rev(10);
let key = rev_to_seq_key(hash, &rev);
let mut reader = KeyReader::new(&key);
assert_eq!(reader.tag(), Some(KeyTag::REV_TO_SEQ.raw()));
assert_eq!(reader.u64(), Some(0xDEAD_BEEF_CAFE_BABE));
assert_eq!(reader.string(), Some("3k2abcde".to_owned()));
assert_eq!(reader.string(), Some(rev.as_str().to_owned()));
assert!(reader.is_empty());
}
@@ -67,9 +78,9 @@ mod tests {
fn rev_to_seq_keys_sort_by_user_then_rev() {
let h1 = UserHash::from_raw(1);
let h2 = UserHash::from_raw(2);
let k1 = rev_to_seq_key(h1, "abc");
let k2 = rev_to_seq_key(h1, "def");
let k3 = rev_to_seq_key(h2, "abc");
let k1 = rev_to_seq_key(h1, &test_rev(10));
let k2 = rev_to_seq_key(h1, &test_rev(20));
let k3 = rev_to_seq_key(h2, &test_rev(10));
assert!(k1.as_slice() < k2.as_slice());
assert!(k2.as_slice() < k3.as_slice());
}
@@ -78,7 +89,7 @@ mod tests {
fn rev_to_seq_user_prefix_is_prefix_of_full_key() {
let hash = UserHash::from_raw(42);
let prefix = rev_to_seq_user_prefix(hash);
let full = rev_to_seq_key(hash, "some_rev");
let full = rev_to_seq_key(hash, &test_rev(10));
assert!(full.as_slice().starts_with(prefix.as_slice()));
}
@@ -43,7 +43,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
pub fn insert_commit_event(&self, data: &CommitEventData) -> Result<SequenceNumber, DbError> {
let event = Self::build_commit_event(data);
self.append_and_index(&event, &data.did, data.rev.as_deref())
self.append_and_index(&event, &data.did, data.rev.as_ref())
}
pub fn append_commit_event_into_batch(
@@ -151,7 +151,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
&self,
did: &Did,
commit_cid: &CidLink,
rev: Option<&str>,
rev: Option<&Tid>,
commit_bytes: &[u8],
) -> Result<SequenceNumber, DbError> {
let inline = tranquil_db_traits::EventBlockInline {
@@ -175,7 +175,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
handle: None,
active: None,
status: None,
rev: rev.map(|r| Tid::from(r.to_owned())),
rev: rev.cloned(),
};
self.append_and_index(&event, did, rev)
@@ -186,7 +186,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
did: &Did,
commit_cid: &CidLink,
mst_root_cid: &CidLink,
rev: &str,
rev: &Tid,
commit_bytes: &[u8],
mst_root_bytes: &[u8],
) -> Result<SequenceNumber, DbError> {
@@ -221,7 +221,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
handle: None,
active: None,
status: None,
rev: Some(Tid::from(rev.to_owned())),
rev: Some(rev.clone()),
};
self.append_and_index(&event, did, Some(rev))
@@ -281,7 +281,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
pub fn get_blob_cids_since_rev(
&self,
did: &Did,
since_rev: &str,
since_rev: &Tid,
) -> Result<Vec<CidLink>, DbError> {
let user_hash = UserHash::from_did(did.as_str());
@@ -489,11 +489,24 @@ impl<S: StorageIO + 'static> EventOps<S> {
self.stage_rev_to_seq(&mut batch, user_hash, rev, seq_u64);
}
if let Some(ms_bytes) = &ewm.mutation_set {
let ms = CommitMutationSet::deserialize(ms_bytes).ok_or_else(|| {
DbError::Query(format!("corrupt CommitMutationSet at seq {seq_u64}"))
})?;
let mutation_set = match &ewm.mutation_set {
None => None,
Some(ms_bytes) => match CommitMutationSet::deserialize(ms_bytes) {
Some(ms) => Some(ms),
None => {
tracing::error!(
seq = seq_u64,
did = %ewm.event.did,
"skipping a CommitMutationSet that does not decode; the \
metastore stays behind the eventlog for this commit until a \
structural repair rewrites it"
);
None
}
},
};
if let Some(ms) = mutation_set {
let meta_key = super::repo_meta::repo_meta_key(user_hash);
let current_meta = self
.repo_data
@@ -586,7 +599,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
batch: &mut fjall::OwnedWriteBatch,
event: &SequencedEvent,
did: &Did,
rev: Option<&str>,
rev: Option<&Tid>,
) -> Result<SequenceNumber, DbError> {
let seq = self
.bridge
@@ -608,7 +621,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
&self,
event: &SequencedEvent,
did: &Did,
rev: Option<&str>,
rev: Option<&Tid>,
) -> Result<SequenceNumber, DbError> {
let mut batch = self.db.batch();
let seq = self.append_and_stage_indexes(&mut batch, event, did, rev)?;
@@ -625,7 +638,7 @@ impl<S: StorageIO + 'static> EventOps<S> {
&self,
batch: &mut fjall::OwnedWriteBatch,
user_hash: UserHash,
rev: &str,
rev: &Tid,
seq: u64,
) {
let key = rev_to_seq_key(user_hash, rev);
@@ -727,8 +740,13 @@ mod tests {
use sha2::Digest;
use tranquil_db_traits::RepoEventType;
fn tid(s: &str) -> Tid {
Tid::from(s.to_owned())
fn test_rev(seq: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
struct TestHarness {
@@ -793,7 +811,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("3k2abcde")),
rev: Some(test_rev(1)),
};
let seq = h.event_ops.insert_commit_event(&data).unwrap();
@@ -803,7 +821,7 @@ mod tests {
assert_eq!(event.did.as_str(), test_did().as_str());
assert_eq!(event.event_type, RepoEventType::Commit);
assert_eq!(event.commit_cid, Some(cid));
assert_eq!(event.rev, Some(tid("3k2abcde")));
assert_eq!(event.rev, Some(test_rev(1)));
}
#[test]
@@ -847,14 +865,14 @@ mod tests {
let seq = h
.event_ops
.insert_sync_event(&test_did(), &cid, Some("rev1"), b"sync_commit_bytes")
.insert_sync_event(&test_did(), &cid, Some(&test_rev(2)), b"sync_commit_bytes")
.unwrap();
assert!(seq.as_i64() > 0);
let event = h.event_ops.get_event_by_seq(seq).unwrap().unwrap();
assert_eq!(event.event_type, RepoEventType::Sync);
assert_eq!(event.commit_cid, Some(cid));
assert_eq!(event.rev, Some(tid("rev1")));
assert_eq!(event.rev, Some(test_rev(2)));
}
#[test]
@@ -869,7 +887,7 @@ mod tests {
&test_did(),
&commit_cid,
&mst_cid,
"genesis_rev",
&test_rev(3),
b"genesis_commit_bytes",
b"genesis_mst_bytes",
)
@@ -880,7 +898,7 @@ mod tests {
assert_eq!(event.event_type, RepoEventType::Commit);
assert_eq!(event.commit_cid, Some(commit_cid));
assert_eq!(event.prev_data_cid, Some(mst_cid));
assert_eq!(event.rev, Some(tid("genesis_rev")));
assert_eq!(event.rev, Some(test_rev(3)));
}
#[test]
@@ -1002,7 +1020,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_a")),
rev: Some(test_rev(4)),
})
.unwrap();
@@ -1017,7 +1035,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_b")),
rev: Some(test_rev(5)),
})
.unwrap();
@@ -1044,7 +1062,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_x")),
rev: Some(test_rev(6)),
})
.unwrap();
@@ -1059,7 +1077,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_y")),
rev: Some(test_rev(7)),
})
.unwrap();
@@ -1145,7 +1163,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("r1")),
rev: Some(test_rev(8)),
})
.unwrap();
let s2 = h.event_ops.insert_identity_event(&did, None).unwrap();
@@ -1155,7 +1173,7 @@ mod tests {
.unwrap();
let s4 = h
.event_ops
.insert_sync_event(&did, &cid, Some("r2"), b"sync_commit_bytes")
.insert_sync_event(&did, &cid, Some(&test_rev(9)), b"sync_commit_bytes")
.unwrap();
let events = h
@@ -1190,7 +1208,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_keep")),
rev: Some(test_rev(10)),
})
.unwrap();
@@ -1203,7 +1221,7 @@ mod tests {
let keep_seq = h
.event_ops
.insert_sync_event(&did, &cid, Some("rev_sync"), b"sync_commit_bytes")
.insert_sync_event(&did, &cid, Some(&test_rev(11)), b"sync_commit_bytes")
.unwrap();
h.event_ops.delete_sequences_except(&did, keep_seq).unwrap();
@@ -1245,7 +1263,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_old")),
rev: Some(test_rev(12)),
})
.unwrap();
@@ -1260,11 +1278,11 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_keep")),
rev: Some(test_rev(10)),
})
.unwrap();
let old_key = super::super::event_keys::rev_to_seq_key(user_hash, "rev_old");
let old_key = super::super::event_keys::rev_to_seq_key(user_hash, &test_rev(12));
assert!(
h.event_ops
.repo_data
@@ -1283,7 +1301,7 @@ mod tests {
.is_none()
);
let keep_key = super::super::event_keys::rev_to_seq_key(user_hash, "rev_keep");
let keep_key = super::super::event_keys::rev_to_seq_key(user_hash, &test_rev(10));
assert!(
h.event_ops
.repo_data
@@ -1311,7 +1329,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_a")),
rev: Some(test_rev(4)),
})
.unwrap();
@@ -1326,7 +1344,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_b")),
rev: Some(test_rev(5)),
})
.unwrap();
@@ -1368,7 +1386,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_1")),
rev: Some(test_rev(13)),
})
.unwrap();
@@ -1397,7 +1415,7 @@ mod tests {
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tid("rev_1")),
rev: Some(test_rev(13)),
})
.unwrap();
@@ -1419,7 +1437,7 @@ mod tests {
handle: None,
active: None,
status: None,
rev: Some(tid("rev_2")),
rev: Some(test_rev(14)),
};
h.event_ops.bridge.insert_event(&crash_event).unwrap();
@@ -1455,12 +1473,12 @@ mod tests {
handle: None,
active: None,
status: None,
rev: Some(tid("rev_3")),
rev: Some(test_rev(15)),
};
h.event_ops.bridge.insert_event(&crash_event_3).unwrap();
let user_hash = super::UserHash::from_did(did.as_str());
let rev2_key = super::super::event_keys::rev_to_seq_key(user_hash, "rev_2");
let rev2_key = super::super::event_keys::rev_to_seq_key(user_hash, &test_rev(14));
assert!(
h.event_ops
.repo_data
+53 -29
View File
@@ -198,12 +198,12 @@ pub enum RepoRequest {
UpdateRepoRoot {
user_id: Uuid,
repo_root_cid: CidLink,
repo_rev: String,
repo_rev: Tid,
tx: Tx<()>,
},
UpdateRepoRev {
user_id: Uuid,
repo_rev: String,
repo_rev: Tid,
tx: Tx<()>,
},
DeleteRepo {
@@ -276,7 +276,7 @@ pub enum RecordRequest {
collections: Vec<Nsid>,
rkeys: Vec<Rkey>,
record_cids: Vec<CidLink>,
repo_rev: String,
repo_rev: Tid,
tx: Tx<()>,
},
DeleteRecords {
@@ -357,7 +357,7 @@ pub enum UserBlockRequest {
InsertUserBlocks {
user_id: Uuid,
block_cids: Vec<Vec<u8>>,
repo_rev: String,
repo_rev: Tid,
tx: Tx<()>,
},
DeleteUserBlocks {
@@ -367,7 +367,7 @@ pub enum UserBlockRequest {
},
GetUserBlockCidsSinceRev {
user_id: Uuid,
since_rev: String,
since_rev: Option<Tid>,
tx: Tx<Vec<Vec<u8>>>,
},
CountUserBlocks {
@@ -405,7 +405,7 @@ pub enum EventRequest {
InsertSyncEvent {
did: Did,
commit_cid: CidLink,
rev: Option<String>,
rev: Option<Tid>,
commit_bytes: Vec<u8>,
tx: Tx<SequenceNumber>,
},
@@ -413,7 +413,7 @@ pub enum EventRequest {
did: Did,
commit_cid: CidLink,
mst_root_cid: CidLink,
rev: String,
rev: Tid,
commit_bytes: Vec<u8>,
mst_root_bytes: Vec<u8>,
tx: Tx<SequenceNumber>,
@@ -579,7 +579,7 @@ pub enum BlobRequest {
},
ListBlobsSinceRev {
did: Did,
since: String,
since: Tid,
tx: Tx<Vec<CidLink>>,
},
CountBlobsByUser {
@@ -2446,12 +2446,19 @@ impl OAuthRequest {
}
}
fn convert_repo_info(r: super::repo_ops::RepoInfo) -> tranquil_db_traits::RepoInfo {
tranquil_db_traits::RepoInfo {
fn convert_repo_info(
r: super::repo_ops::RepoInfo,
) -> Result<tranquil_db_traits::RepoInfo, DbError> {
Ok(tranquil_db_traits::RepoInfo {
user_id: r.user_id,
repo_root_cid: r.repo_root_cid,
repo_rev: r.repo_rev.map(Tid::from),
}
repo_rev: r
.repo_rev
.map(|rev| {
Tid::new(rev).map_err(|_| DbError::CorruptData("corrupt repo_meta repo_rev"))
})
.transpose()?,
})
}
fn convert_repo_account(
@@ -2473,11 +2480,16 @@ fn convert_repo_list_entry(
.did
.ok_or(DbError::CorruptData("repo_meta missing DID field"))?;
Ok(tranquil_db_traits::RepoListItem {
did: Did::from(did),
did: Did::new(did).map_err(|_| DbError::CorruptData("corrupt repo_meta did"))?,
deactivated_at: r.deactivated_at,
takedown_ref: r.takedown_ref,
repo_root_cid: r.repo_root_cid,
repo_rev: r.repo_rev.map(Tid::from),
repo_rev: r
.repo_rev
.map(|rev| {
Tid::new(rev).map_err(|_| DbError::CorruptData("corrupt repo_meta repo_rev"))
})
.transpose()?,
})
}
@@ -2597,8 +2609,8 @@ fn dispatch_repo<S: StorageIO>(state: &HandlerState<S>, req: RepoRequest) {
.metastore
.repo_ops()
.get_repo(user_id)
.map(|opt| opt.map(convert_repo_info))
.map_err(metastore_to_db);
.map_err(metastore_to_db)
.and_then(|opt| opt.map(convert_repo_info).transpose());
let _ = tx.send(result);
}
RepoRequest::GetRepoRootByDid { did, tx } => {
@@ -2712,7 +2724,7 @@ fn dispatch_record<S: StorageIO>(state: &HandlerState<S>, req: RecordRequest) {
.record_ops()
.upsert_records(&mut batch, user_hash, &writes)
.map_err(metastore_to_db)?;
meta.repo_rev = repo_rev;
meta.repo_rev = repo_rev.as_str().to_owned();
state
.metastore
.repo_ops()
@@ -2923,7 +2935,7 @@ fn dispatch_user_block<S: StorageIO>(state: &HandlerState<S>, req: UserBlockRequ
let result = state
.metastore
.user_block_ops()
.get_user_block_cids_since_rev(user_id, &since_rev)
.get_user_block_cids_since_rev(user_id, since_rev.as_ref())
.map_err(metastore_to_db);
let _ = tx.send(result);
}
@@ -2962,7 +2974,7 @@ fn dispatch_event<S: StorageIO + 'static>(state: &HandlerState<S>, req: EventReq
let result =
state
.event_ops
.insert_sync_event(&did, &commit_cid, rev.as_deref(), &commit_bytes);
.insert_sync_event(&did, &commit_cid, rev.as_ref(), &commit_bytes);
let _ = tx.send(result);
}
EventRequest::InsertGenesisCommitEvent {
@@ -6139,13 +6151,23 @@ mod tests {
CidLink::from_cid(&c)
}
fn test_rev(seq: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
#[tokio::test]
async fn create_and_get_roundtrip() {
let h = setup();
let user_id = Uuid::new_v4();
let did = Did::from("did:plc:handler_test".to_string());
let handle = Handle::from("handler.test.invalid".to_string());
let did = Did::new("did:plc:whelk").expect("test DID is well-formed");
let handle = Handle::new("whelk.oyster.cafe").expect("test handle is valid");
let cid = test_cid_link(1);
let rev = test_rev(1);
let (tx, rx) = oneshot::channel();
h.pool
@@ -6154,7 +6176,7 @@ mod tests {
did,
handle,
repo_root_cid: cid.clone(),
repo_rev: Tid::from("rev1".to_string()),
repo_rev: rev.clone(),
tx,
}))
.unwrap();
@@ -6166,7 +6188,7 @@ mod tests {
.unwrap();
let repo = rx.await.unwrap().unwrap().unwrap();
assert_eq!(repo.repo_root_cid, cid);
assert_eq!(repo.repo_rev.as_deref(), Some("rev1"));
assert_eq!(repo.repo_rev.as_deref(), Some(rev.as_str()));
}
#[test]
@@ -6181,9 +6203,10 @@ mod tests {
.unwrap();
let infra = metastore.infra_ops();
let squid = InviteCode::new("squid-invite");
let owner = Did::new("did:plc:whelk").expect("valid DID");
let whelk = InviteCode::new("whelk");
assert!(infra.create_invite_code(&squid, 1, None).unwrap());
assert!(infra.create_invite_code(&squid, 1, Some(&owner)).unwrap());
infra.reserve_invite_code(&squid).unwrap();
assert_eq!(
@@ -6219,9 +6242,10 @@ mod tests {
.unwrap();
let infra = metastore.infra_ops();
let squid = InviteCode::new("squid-invite");
let owner = Did::new("did:plc:whelk").expect("valid DID");
let whelk = InviteCode::new("whelk");
assert!(infra.create_invite_code(&squid, 1, None).unwrap());
assert!(infra.create_invite_code(&squid, 1, Some(&owner)).unwrap());
infra.reserve_invite_code(&squid).unwrap();
infra.refund_invite_code(&squid).unwrap();
@@ -6275,7 +6299,7 @@ mod tests {
)
.unwrap();
let user_hashes = ms.user_hashes().as_ref();
let did = Did::from("did:plc:limpet".to_string());
let did = Did::new("did:plc:limpet").expect("test DID is well-formed");
let expected = did_to_routing(&did);
let sid = SessionId::new(7);
@@ -6320,8 +6344,8 @@ mod tests {
async fn shutdown_completes_inflight() {
let h = setup();
let user_id = Uuid::new_v4();
let did = Did::from("did:plc:shutdown_test".to_string());
let handle = Handle::from("shutdown.test.invalid".to_string());
let did = Did::new("did:plc:scallop").expect("test DID is well-formed");
let handle = Handle::new("scallop.oyster.cafe").expect("test handle is valid");
let cid = test_cid_link(2);
let (tx, rx) = oneshot::channel();
@@ -6331,7 +6355,7 @@ mod tests {
did,
handle,
repo_root_cid: cid,
repo_rev: Tid::from("rev1".to_string()),
repo_rev: test_rev(1),
tx,
}))
.unwrap();
@@ -506,13 +506,22 @@ mod tests {
tranquil_types::Handle::from(format!("{name}.test.invalid"))
}
fn test_rev(seq: u64) -> tranquil_types::Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
tranquil_types::Tid::new(s).expect("generated TID is valid")
}
fn setup_user(ms: &Metastore) -> (Uuid, super::super::keys::UserHash) {
let user_id = Uuid::new_v4();
let did = test_did("testuser");
let handle = test_handle("testuser");
let cid = test_cid_link(0);
ms.repo_ops()
.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev0")
.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(0))
.unwrap();
let user_hash = ms.user_hashes().get(&user_id).unwrap();
(user_id, user_hash)
@@ -976,7 +985,7 @@ mod tests {
&did1,
&handle1,
&test_cid_link(0),
"r",
&test_rev(0),
)
.unwrap();
let hash1 = ms.user_hashes().get(&user1).unwrap();
@@ -991,7 +1000,7 @@ mod tests {
&did2,
&handle2,
&test_cid_link(0),
"r",
&test_rev(0),
)
.unwrap();
let hash2 = ms.user_hashes().get(&user2).unwrap();
@@ -1165,7 +1174,7 @@ mod tests {
&did,
&handle,
&test_cid_link(0),
"r",
&test_rev(0),
)
.unwrap();
let user_hash = ms.user_hashes().get(&user_id).unwrap();
@@ -1,7 +1,7 @@
use std::collections::HashSet;
use serde::{Deserialize, Serialize};
use tranquil_types::{Nsid, Rkey};
use tranquil_types::{Nsid, Rkey, Tid};
use super::backlink_ops::remove_backlinks_for_record;
use super::backlinks::{BacklinkValue, backlink_by_user_key, backlink_key, discriminant_to_path};
@@ -17,7 +17,7 @@ const MUTATION_SET_VERSION: u8 = 1;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CommitMutationSet {
pub new_root_cid: Vec<u8>,
pub new_rev: String,
pub new_rev: Tid,
pub record_upserts: Vec<RecordMutationUpsert>,
pub record_deletes: Vec<RecordMutationDelete>,
pub block_inserts: Vec<Vec<u8>>,
@@ -111,7 +111,7 @@ pub fn replay_mutation_set(
let updated_meta = RepoMetaValue {
repo_root_cid: mutation_set.new_root_cid.clone(),
repo_rev: mutation_set.new_rev.clone(),
repo_rev: mutation_set.new_rev.as_str().to_owned(),
..current_meta.clone()
};
let meta_key = repo_meta_key(user_hash);
@@ -254,7 +254,7 @@ mod tests {
fn mutation_set_roundtrip() {
let ms = CommitMutationSet {
new_root_cid: vec![0x01, 0x71, 0x12, 0x20],
new_rev: "rev1".to_owned(),
new_rev: Tid::new("3k2abcdefghij").unwrap(),
record_upserts: vec![RecordMutationUpsert {
collection: "app.bsky.feed.post".to_owned(),
rkey: "3k2abc".to_owned(),
@@ -284,7 +284,7 @@ mod tests {
fn mutation_set_empty_roundtrip() {
let ms = CommitMutationSet {
new_root_cid: vec![],
new_rev: String::new(),
new_rev: Tid::new("3k2abcdefghij").unwrap(),
record_upserts: vec![],
record_deletes: vec![],
block_inserts: vec![],
@@ -301,7 +301,7 @@ mod tests {
fn unknown_version_returns_none() {
let ms = CommitMutationSet {
new_root_cid: vec![],
new_rev: String::new(),
new_rev: Tid::new("3k2abcdefghij").unwrap(),
record_upserts: vec![],
record_deletes: vec![],
block_inserts: vec![],
+106 -80
View File
@@ -14,7 +14,7 @@ use super::scan::{count_prefix, delete_all_by_prefix, point_lookup};
use super::user_blocks::user_block_user_prefix;
use super::user_hash::UserHashMap;
use tranquil_types::{CidLink, Did, Handle};
use tranquil_types::{CidLink, Did, Handle, Tid};
pub struct RepoOps {
repo_data: Keyspace,
@@ -36,7 +36,7 @@ impl RepoOps {
did: &Did,
handle: &Handle,
repo_root_cid: &CidLink,
repo_rev: &str,
repo_rev: &Tid,
) -> Result<(), MetastoreError> {
let user_hash = UserHash::from_did(did.as_str());
let mut batch = db.batch();
@@ -49,7 +49,7 @@ impl RepoOps {
let value = RepoMetaValue {
repo_root_cid: cid_bytes,
repo_rev: repo_rev.to_string(),
repo_rev: repo_rev.as_str().to_owned(),
handle: handle_lower.clone(),
status: RepoStatus::Active,
deactivated_at_ms: None,
@@ -114,7 +114,7 @@ impl RepoOps {
db: &fjall::Database,
user_id: Uuid,
repo_root_cid: &CidLink,
repo_rev: &str,
repo_rev: &Tid,
) -> Result<(), MetastoreError> {
let user_hash = self.resolve_user_hash(user_id)?;
let key = repo_meta_key(user_hash);
@@ -122,7 +122,7 @@ impl RepoOps {
let mut value = self.get_meta_value(key.as_slice())?;
let cid_bytes = cid_link_to_bytes(repo_root_cid)?;
value.repo_root_cid = cid_bytes;
value.repo_rev = repo_rev.to_string();
value.repo_rev = repo_rev.as_str().to_owned();
let mut batch = db.batch();
batch.insert(&self.repo_data, key.as_slice(), value.serialize());
@@ -133,13 +133,13 @@ impl RepoOps {
&self,
db: &fjall::Database,
user_id: Uuid,
repo_rev: &str,
repo_rev: &Tid,
) -> Result<(), MetastoreError> {
let user_hash = self.resolve_user_hash(user_id)?;
let key = repo_meta_key(user_hash);
let mut value = self.get_meta_value(key.as_slice())?;
value.repo_rev = repo_rev.to_string();
value.repo_rev = repo_rev.as_str().to_owned();
let mut batch = db.batch();
batch.insert(&self.repo_data, key.as_slice(), value.serialize());
@@ -316,7 +316,10 @@ impl RepoOps {
Ok(RepoInfo {
user_id,
repo_root_cid: cid,
repo_rev: Some(value.repo_rev),
repo_rev: match value.repo_rev.is_empty() {
true => None,
false => Some(value.repo_rev),
},
})
})
.transpose()
@@ -651,11 +654,20 @@ mod tests {
}
fn test_did(name: &str) -> Did {
Did::from(format!("did:plc:{name}"))
Did::new(format!("did:plc:{name}")).expect("test DID is well-formed")
}
fn test_handle(name: &str) -> Handle {
Handle::from(format!("{name}.test.invalid"))
Handle::new(format!("{name}.oyster.cafe")).expect("test handle is valid")
}
fn test_rev(seq: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
#[test]
@@ -667,13 +679,13 @@ mod tests {
let handle = test_handle("olaren");
let cid = test_cid_link(1);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
let repo = ops.get_repo(user_id).unwrap().unwrap();
assert_eq!(repo.user_id, user_id);
assert_eq!(repo.repo_root_cid, cid);
assert_eq!(repo.repo_rev.as_deref(), Some("rev1"));
assert_eq!(repo.repo_rev.as_deref(), Some(test_rev(1).as_str()));
}
#[test]
@@ -693,14 +705,14 @@ mod tests {
let cid1 = test_cid_link(1);
let cid2 = test_cid_link(2);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid1, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid1, &test_rev(1))
.unwrap();
ops.update_repo_root(ms.database(), user_id, &cid2, "rev2")
ops.update_repo_root(ms.database(), user_id, &cid2, &test_rev(2))
.unwrap();
let repo = ops.get_repo(user_id).unwrap().unwrap();
assert_eq!(repo.repo_root_cid, cid2);
assert_eq!(repo.repo_rev.as_deref(), Some("rev2"));
assert_eq!(repo.repo_rev.as_deref(), Some(test_rev(2).as_str()));
}
#[test]
@@ -712,14 +724,14 @@ mod tests {
let handle = test_handle("nel");
let cid = test_cid_link(3);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
ops.update_repo_rev(ms.database(), user_id, "rev_updated")
ops.update_repo_rev(ms.database(), user_id, &test_rev(3))
.unwrap();
let repo = ops.get_repo(user_id).unwrap().unwrap();
assert_eq!(repo.repo_root_cid, cid);
assert_eq!(repo.repo_rev.as_deref(), Some("rev_updated"));
assert_eq!(repo.repo_rev.as_deref(), Some(test_rev(3).as_str()));
}
#[test]
@@ -731,7 +743,7 @@ mod tests {
let handle = test_handle("lyna");
let cid = test_cid_link(4);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
assert!(ops.get_repo(user_id).unwrap().is_some());
assert!(ops.lookup_handle(&handle).unwrap().is_some());
@@ -750,7 +762,7 @@ mod tests {
let handle = test_handle("mapped");
let cid = test_cid_link(4);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
assert!(ms.user_hashes().get(&user_id).is_some());
@@ -763,22 +775,22 @@ mod tests {
let (_dir, ms) = open_fresh();
let ops = ms.repo_ops();
let did = test_did("recreate");
let handle_a = test_handle("recreate_a");
let handle_b = test_handle("recreate_b");
let handle_a = test_handle("recreate-a");
let handle_b = test_handle("recreate-b");
let uid_a = uuid::Uuid::new_v4();
let uid_b = uuid::Uuid::new_v4();
let cid = test_cid_link(70);
ops.create_repo(ms.database(), uid_a, &did, &handle_a, &cid, "r1")
ops.create_repo(ms.database(), uid_a, &did, &handle_a, &cid, &test_rev(1))
.unwrap();
ops.delete_repo(ms.database(), uid_a).unwrap();
ops.create_repo(ms.database(), uid_b, &did, &handle_b, &cid, "r2")
ops.create_repo(ms.database(), uid_b, &did, &handle_b, &cid, &test_rev(2))
.unwrap();
let repo = ops.get_repo(uid_b).unwrap().unwrap();
assert_eq!(repo.user_id, uid_b);
assert_eq!(repo.repo_rev.as_deref(), Some("r2"));
assert_eq!(repo.repo_rev.as_deref(), Some(test_rev(2).as_str()));
assert!(ops.get_repo(uid_a).unwrap().is_none());
}
@@ -800,7 +812,7 @@ mod tests {
&orphan_did,
&orphan_handle,
&cid,
"rev1",
&test_rev(1),
)
.unwrap();
ops.create_repo(
@@ -809,7 +821,7 @@ mod tests {
&live_did,
&live_handle,
&cid,
"rev1",
&test_rev(1),
)
.unwrap();
@@ -853,7 +865,7 @@ mod tests {
&orphan_did,
&orphan_handle,
&cid,
"rev1",
&test_rev(1),
)
.unwrap();
ops.create_repo(
@@ -862,7 +874,7 @@ mod tests {
&live_did,
&live_handle,
&cid,
"rev1",
&test_rev(1),
)
.unwrap();
@@ -914,10 +926,10 @@ mod tests {
let handle = test_handle("bailey");
let cid = test_cid_link(5);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
let upper_handle = Handle::from("BAILEY.TEST.INVALID".to_string());
let upper_handle = Handle::new("BAILEY.OYSTER.CAFE").expect("test handle is valid");
let found = ops.lookup_handle(&upper_handle).unwrap();
assert_eq!(found, Some(user_id));
}
@@ -931,7 +943,7 @@ mod tests {
let handle = test_handle("olaren");
let cid = test_cid_link(6);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
let root = ops.get_repo_root_by_did(&did).unwrap().unwrap();
@@ -950,7 +962,7 @@ mod tests {
let handle = test_handle("teq");
let cid = test_cid_link(7);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
let root = ops.get_repo_root_for_update(user_id).unwrap().unwrap();
@@ -969,7 +981,7 @@ mod tests {
let did = test_did(&format!("user{i}"));
let handle = test_handle(&format!("user{i}"));
let cid = test_cid_link(i);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
});
@@ -985,7 +997,7 @@ mod tests {
let handle = test_handle("nel");
let cid = test_cid_link(8);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
let account = ops.get_account_with_repo(&did).unwrap().unwrap();
@@ -1012,7 +1024,7 @@ mod tests {
&did,
&handle,
&cid,
&format!("rev{i}"),
&test_rev(i as u64),
)
.unwrap();
});
@@ -1041,7 +1053,7 @@ mod tests {
&did,
&handle,
&cid,
&format!("rev{i}"),
&test_rev(i as u64),
)
.unwrap();
});
@@ -1074,7 +1086,7 @@ mod tests {
{
let ms = Metastore::open(dir.path(), test_config()).unwrap();
let ops = ms.repo_ops();
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev_persist")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(4))
.unwrap();
ms.persist().unwrap();
}
@@ -1084,7 +1096,7 @@ mod tests {
let ops = ms.repo_ops();
let repo = ops.get_repo(user_id).unwrap().unwrap();
assert_eq!(repo.repo_root_cid, cid);
assert_eq!(repo.repo_rev.as_deref(), Some("rev_persist"));
assert_eq!(repo.repo_rev.as_deref(), Some(test_rev(4).as_str()));
let found = ops.lookup_handle(&handle).unwrap();
assert_eq!(found, Some(user_id));
@@ -1097,31 +1109,38 @@ mod tests {
let ops = ms.repo_ops();
let uid_with = uuid::Uuid::new_v4();
let did_with = test_did("with_rev");
let handle_with = test_handle("with_rev");
let did_with = test_did("with-rev");
let handle_with = test_handle("with-rev");
ops.create_repo(
ms.database(),
uid_with,
&did_with,
&handle_with,
&test_cid_link(40),
"some_rev",
&test_rev(6),
)
.unwrap();
let uid_without = uuid::Uuid::new_v4();
let did_without = test_did("without_rev");
let handle_without = test_handle("without_rev");
let did_without = test_did("without-rev");
let handle_without = test_handle("without-rev");
ops.create_repo(
ms.database(),
uid_without,
&did_without,
&handle_without,
&test_cid_link(41),
"",
&test_rev(7),
)
.unwrap();
let (user_hash_without, mut meta_without) =
ops.get_repo_meta(uid_without).unwrap().unwrap();
meta_without.repo_rev = String::new();
let mut batch = ms.database().batch();
ops.write_repo_meta(&mut batch, user_hash_without, &meta_without);
batch.commit().unwrap();
let result = ops.get_repos_without_rev(100).unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0].user_id, uid_without);
@@ -1132,11 +1151,11 @@ mod tests {
let (_dir, ms) = open_fresh();
let ops = ms.repo_ops();
let user_id = uuid::Uuid::new_v4();
let did = test_did("root_cid");
let handle = test_handle("root_cid");
let did = test_did("root-cid");
let handle = test_handle("root-cid");
let cid = test_cid_link(50);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
let root = ops.get_repo_root_cid_by_user_id(user_id).unwrap().unwrap();
@@ -1155,28 +1174,28 @@ mod tests {
let ops = ms.repo_ops();
let uid_a = uuid::Uuid::new_v4();
let did_a = test_did("keep_a");
let handle_a = test_handle("keep_a");
let did_a = test_did("keep-a");
let handle_a = test_handle("keep-a");
ops.create_repo(
ms.database(),
uid_a,
&did_a,
&handle_a,
&test_cid_link(60),
"r",
&test_rev(1),
)
.unwrap();
let uid_b = uuid::Uuid::new_v4();
let did_b = test_did("delete_b");
let handle_b = test_handle("delete_b");
let did_b = test_did("delete-b");
let handle_b = test_handle("delete-b");
ops.create_repo(
ms.database(),
uid_b,
&did_b,
&handle_b,
&test_cid_link(61),
"r",
&test_rev(1),
)
.unwrap();
@@ -1193,17 +1212,17 @@ mod tests {
let ops = ms.repo_ops();
let uid_a = uuid::Uuid::new_v4();
let did_a = test_did("collision_a");
let handle_a = test_handle("collision_a");
let did_a = test_did("collision-a");
let handle_a = test_handle("collision-a");
let cid = test_cid_link(80);
ops.create_repo(ms.database(), uid_a, &did_a, &handle_a, &cid, "r1")
ops.create_repo(ms.database(), uid_a, &did_a, &handle_a, &cid, &test_rev(1))
.unwrap();
let uid_b = uuid::Uuid::new_v4();
let handle_b = test_handle("collision_b");
let handle_b = test_handle("collision-b");
let result = ops.create_repo(ms.database(), uid_b, &did_a, &handle_b, &cid, "r2");
let result = ops.create_repo(ms.database(), uid_b, &did_a, &handle_b, &cid, &test_rev(2));
match result {
Ok(()) => {
@@ -1213,7 +1232,7 @@ mod tests {
Err(MetastoreError::UserHashCollision { .. }) => {
let repo = ops.get_repo(uid_a).unwrap().unwrap();
assert_eq!(repo.repo_root_cid, cid);
assert_eq!(repo.repo_rev.as_deref(), Some("r1"));
assert_eq!(repo.repo_rev.as_deref(), Some(test_rev(1).as_str()));
}
Err(e) => panic!("unexpected error: {e}"),
}
@@ -1224,17 +1243,17 @@ mod tests {
let (_dir, ms) = open_fresh();
let ops = ms.repo_ops();
let user_id = uuid::Uuid::new_v4();
let did = test_did("meta_raw");
let handle = test_handle("meta_raw");
let did = test_did("meta-raw");
let handle = test_handle("meta-raw");
let cid = test_cid_link(81);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "rev_meta")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(5))
.unwrap();
let (user_hash, value) = ops.get_repo_meta(user_id).unwrap().unwrap();
assert_eq!(user_hash, UserHash::from_did(did.as_str()));
assert_eq!(value.repo_rev, "rev_meta");
assert_eq!(value.handle, "meta_raw.test.invalid");
assert_eq!(value.repo_rev, test_rev(5).as_str());
assert_eq!(value.handle, "meta-raw.oyster.cafe");
assert_eq!(value.status, RepoStatus::Active);
}
@@ -1250,17 +1269,17 @@ mod tests {
let (_dir, ms) = open_fresh();
let ops = ms.repo_ops();
let user_id = uuid::Uuid::new_v4();
let did = test_did("batch_write");
let handle = test_handle("batch_write");
let did = test_did("batch-write");
let handle = test_handle("batch-write");
let cid1 = test_cid_link(82);
let cid2 = test_cid_link(83);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid1, "rev1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid1, &test_rev(1))
.unwrap();
let (user_hash, mut value) = ops.get_repo_meta(user_id).unwrap().unwrap();
value.repo_root_cid = cid_link_to_bytes(&cid2).unwrap();
value.repo_rev = "rev2".to_string();
value.repo_rev = test_rev(2).to_string();
let mut batch = ms.database().batch();
ops.write_repo_meta(&mut batch, user_hash, &value);
@@ -1268,7 +1287,7 @@ mod tests {
let repo = ops.get_repo(user_id).unwrap().unwrap();
assert_eq!(repo.repo_root_cid, cid2);
assert_eq!(repo.repo_rev.as_deref(), Some("rev2"));
assert_eq!(repo.repo_rev.as_deref(), Some(test_rev(2).as_str()));
}
#[test]
@@ -1276,13 +1295,20 @@ mod tests {
let (_dir, ms) = open_fresh();
let ops = ms.repo_ops();
let user_id = uuid::Uuid::new_v4();
let did = test_did("handle_swap");
let old_handle = test_handle("old_name");
let new_handle = test_handle("new_name");
let did = test_did("handle-swap");
let old_handle = test_handle("old-name");
let new_handle = test_handle("new-name");
let cid = test_cid_link(84);
ops.create_repo(ms.database(), user_id, &did, &old_handle, &cid, "r1")
.unwrap();
ops.create_repo(
ms.database(),
user_id,
&did,
&old_handle,
&cid,
&test_rev(1),
)
.unwrap();
assert!(ops.lookup_handle(&old_handle).unwrap().is_some());
ops.update_handle(ms.database(), user_id, &new_handle)
@@ -1292,7 +1318,7 @@ mod tests {
assert_eq!(ops.lookup_handle(&new_handle).unwrap(), Some(user_id));
let (_, meta) = ops.get_repo_meta(user_id).unwrap().unwrap();
assert_eq!(meta.handle, "new_name.test.invalid");
assert_eq!(meta.handle, "new-name.oyster.cafe");
}
#[test]
@@ -1300,18 +1326,18 @@ mod tests {
let (_dir, ms) = open_fresh();
let ops = ms.repo_ops();
let user_id = uuid::Uuid::new_v4();
let did = test_did("handle_case");
let did = test_did("handle-case");
let handle = test_handle("original");
let cid = test_cid_link(85);
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, "r1")
ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1))
.unwrap();
let mixed_case = Handle::from("UPPER.TEST.INVALID".to_string());
let mixed_case = Handle::new("UPPER.OYSTER.CAFE").expect("test handle is valid");
ops.update_handle(ms.database(), user_id, &mixed_case)
.unwrap();
let lower_lookup = Handle::from("upper.test.invalid".to_string());
let lower_lookup = Handle::new("upper.oyster.cafe").expect("test handle is valid");
assert_eq!(ops.lookup_handle(&lower_lookup).unwrap(), Some(user_id));
assert!(ops.lookup_handle(&handle).unwrap().is_none());
}
@@ -2,6 +2,7 @@ use std::collections::HashSet;
use std::sync::Arc;
use fjall::Keyspace;
use tranquil_types::Tid;
use uuid::Uuid;
use super::MetastoreError;
@@ -29,7 +30,7 @@ impl UserBlockOps {
batch: &mut fjall::OwnedWriteBatch,
user_hash: UserHash,
block_cids: &[C],
repo_rev: &str,
repo_rev: &Tid,
) -> Result<(), MetastoreError> {
let existing: HashSet<Vec<u8>> = match block_cids.is_empty() {
true => HashSet::new(),
@@ -37,11 +38,12 @@ impl UserBlockOps {
let prefix = user_block_user_prefix(user_hash);
self.repo_data
.prefix(prefix.as_slice())
.filter_map(|guard| {
let (key_bytes, _) = guard.into_inner().ok()?;
extract_cid_from_key(&key_bytes).map(|c| c.to_vec())
.map(|guard| {
let (key_bytes, _) = guard.into_inner().map_err(MetastoreError::Fjall)?;
Ok(extract_cid_from_key(&key_bytes).map(|c| c.to_vec()))
})
.collect()
.filter_map(Result::transpose)
.collect::<Result<_, MetastoreError>>()?
}
};
@@ -60,26 +62,6 @@ impl UserBlockOps {
})
}
pub fn delete_user_blocks<C: AsRef<[u8]>>(
&self,
batch: &mut fjall::OwnedWriteBatch,
user_hash: UserHash,
block_cids: &[C],
rev: &str,
) -> Result<(), MetastoreError> {
block_cids.iter().try_for_each(|cid| {
let cid = cid.as_ref();
match cid.is_empty() {
true => Err(MetastoreError::InvalidInput("block CID must not be empty")),
false => {
let key = user_block_key(user_hash, rev, cid);
batch.remove(&self.repo_data, key.as_slice());
Ok(())
}
}
})
}
pub fn delete_user_blocks_by_cid<C: AsRef<[u8]>>(
&self,
batch: &mut fjall::OwnedWriteBatch,
@@ -119,23 +101,25 @@ impl UserBlockOps {
pub fn get_user_block_cids_since_rev(
&self,
user_id: Uuid,
since_rev: &str,
since_rev: Option<&Tid>,
) -> Result<Vec<Vec<u8>>, MetastoreError> {
let user_hash = match self.user_hashes.get(&user_id) {
Some(h) => h,
None => return Ok(Vec::new()),
};
let since_prefix = user_block_rev_prefix(user_hash, since_rev);
let since_upper = exclusive_upper_bound(since_prefix.as_slice())
.expect("user block rev prefix always contains non-0xFF bytes");
let user_prefix = user_block_user_prefix(user_hash);
let user_upper = exclusive_upper_bound(user_prefix.as_slice())
.expect("user block user prefix always contains non-0xFF bytes");
let lower = match since_rev {
None => user_prefix,
Some(rev) => exclusive_upper_bound(user_block_rev_prefix(user_hash, rev).as_slice())
.expect("user block rev prefix always contains non-0xFF bytes"),
};
self.repo_data
.range(since_upper.as_slice()..user_upper.as_slice())
.range(lower.as_slice()..user_upper.as_slice())
.map(|guard| {
let (key_bytes, _) = guard.into_inner().map_err(MetastoreError::Fjall)?;
extract_cid_from_key(&key_bytes)
@@ -182,9 +166,18 @@ mod tests {
(dir, ms)
}
fn test_rev(seq: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
fn setup_user(ms: &Metastore) -> (Uuid, UserHash) {
let uuid = Uuid::new_v4();
let hash = UserHash::from_did("did:plc:testuser1");
let hash = UserHash::from_did("did:plc:limpet");
let mut batch = ms.database().batch();
ms.user_hashes()
.stage_insert(&mut batch, uuid, hash)
@@ -202,7 +195,7 @@ mod tests {
let cids = vec![vec![0x01, 0x71], vec![0x02, 0x72], vec![0x03, 0x73]];
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &cids, "rev1")
ops.insert_user_blocks(&mut batch, hash, &cids, &test_rev(10))
.unwrap();
batch.commit().unwrap();
@@ -219,17 +212,21 @@ mod tests {
let cids_def = vec![vec![0x03]];
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &cids_abc, "abc")
ops.insert_user_blocks(&mut batch, hash, &cids_abc, &test_rev(10))
.unwrap();
ops.insert_user_blocks(&mut batch, hash, &cids_def, "def")
ops.insert_user_blocks(&mut batch, hash, &cids_def, &test_rev(20))
.unwrap();
batch.commit().unwrap();
let since_abc = ops.get_user_block_cids_since_rev(uuid, "abc").unwrap();
let since_abc = ops
.get_user_block_cids_since_rev(uuid, Some(&test_rev(10)))
.unwrap();
assert_eq!(since_abc.len(), 1);
assert_eq!(since_abc[0], vec![0x03]);
let since_def = ops.get_user_block_cids_since_rev(uuid, "def").unwrap();
let since_def = ops
.get_user_block_cids_since_rev(uuid, Some(&test_rev(20)))
.unwrap();
assert!(since_def.is_empty());
}
@@ -243,62 +240,46 @@ mod tests {
let cids_r2 = vec![vec![0x02]];
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &cids_r1, "aaa")
ops.insert_user_blocks(&mut batch, hash, &cids_r1, &test_rev(10))
.unwrap();
ops.insert_user_blocks(&mut batch, hash, &cids_r2, "bbb")
ops.insert_user_blocks(&mut batch, hash, &cids_r2, &test_rev(20))
.unwrap();
batch.commit().unwrap();
let since_before = ops.get_user_block_cids_since_rev(uuid, "aaa").unwrap();
let since_before = ops
.get_user_block_cids_since_rev(uuid, Some(&test_rev(10)))
.unwrap();
assert_eq!(since_before.len(), 1);
assert_eq!(since_before[0], vec![0x02]);
let all = ops.get_user_block_cids_since_rev(uuid, "").unwrap();
let all = ops.get_user_block_cids_since_rev(uuid, None).unwrap();
assert_eq!(all.len(), 2);
}
#[test]
fn delete_blocks_at_rev() {
fn delete_by_cid_ignores_the_rev_the_block_was_recorded_at() {
let (_dir, ms) = open_fresh();
let (uuid, hash) = setup_user(&ms);
let ops = ms.user_block_ops();
let cids = vec![vec![0x01], vec![0x02], vec![0x03]];
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &cids, "rev1")
ops.insert_user_blocks(&mut batch, hash, &[vec![0x01], vec![0x02]], &test_rev(10))
.unwrap();
ops.insert_user_blocks(&mut batch, hash, &[vec![0x03]], &test_rev(20))
.unwrap();
batch.commit().unwrap();
assert_eq!(ops.count_user_blocks(uuid).unwrap(), 3);
let to_delete = vec![vec![0x01], vec![0x03]];
let mut batch = ms.database().batch();
ops.delete_user_blocks(&mut batch, hash, &to_delete, "rev1")
ops.delete_user_blocks_by_cid(&mut batch, hash, &[vec![0x01], vec![0x03]])
.unwrap();
batch.commit().unwrap();
assert_eq!(ops.count_user_blocks(uuid).unwrap(), 1);
}
#[test]
fn delete_wrong_rev_is_noop() {
let (_dir, ms) = open_fresh();
let (uuid, hash) = setup_user(&ms);
let ops = ms.user_block_ops();
let cids = vec![vec![0x01], vec![0x02]];
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &cids, "rev1")
.unwrap();
batch.commit().unwrap();
let mut batch = ms.database().batch();
ops.delete_user_blocks(&mut batch, hash, &cids, "wrong_rev")
.unwrap();
batch.commit().unwrap();
assert_eq!(ops.count_user_blocks(uuid).unwrap(), 2);
assert_eq!(
ops.get_user_block_cids_since_rev(uuid, None).unwrap(),
vec![vec![0x02]]
);
}
#[test]
@@ -308,9 +289,9 @@ mod tests {
let ops = ms.user_block_ops();
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &[vec![0x01], vec![0x02]], "rev1")
ops.insert_user_blocks(&mut batch, hash, &[vec![0x01], vec![0x02]], &test_rev(10))
.unwrap();
ops.insert_user_blocks(&mut batch, hash, &[vec![0x03]], "rev2")
ops.insert_user_blocks(&mut batch, hash, &[vec![0x03]], &test_rev(20))
.unwrap();
batch.commit().unwrap();
assert_eq!(ops.count_user_blocks(uuid).unwrap(), 3);
@@ -328,7 +309,9 @@ mod tests {
let (uuid, _hash) = setup_user(&ms);
let ops = ms.user_block_ops();
let result = ops.get_user_block_cids_since_rev(uuid, "anything").unwrap();
let result = ops
.get_user_block_cids_since_rev(uuid, Some(&test_rev(10)))
.unwrap();
assert!(result.is_empty());
}
@@ -363,7 +346,7 @@ mod tests {
let ops = ms.user_block_ops();
let cids = vec![vec![0xAA, 0xBB], vec![0xCC, 0xDD]];
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &cids, "rev1")
ops.insert_user_blocks(&mut batch, hash, &cids, &test_rev(10))
.unwrap();
batch.commit().unwrap();
ms.persist().unwrap();
@@ -387,9 +370,9 @@ mod tests {
let (_dir, ms) = open_fresh();
let uuid1 = Uuid::new_v4();
let hash1 = UserHash::from_did("did:plc:user1");
let hash1 = UserHash::from_did("did:plc:whelk");
let uuid2 = Uuid::new_v4();
let hash2 = UserHash::from_did("did:plc:user2");
let hash2 = UserHash::from_did("did:plc:scallop");
let mut batch = ms.database().batch();
ms.user_hashes()
@@ -403,9 +386,9 @@ mod tests {
let ops = ms.user_block_ops();
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash1, &[vec![0x01], vec![0x02]], "rev1")
ops.insert_user_blocks(&mut batch, hash1, &[vec![0x01], vec![0x02]], &test_rev(10))
.unwrap();
ops.insert_user_blocks(&mut batch, hash2, &[vec![0x03]], "rev1")
ops.insert_user_blocks(&mut batch, hash2, &[vec![0x03]], &test_rev(10))
.unwrap();
batch.commit().unwrap();
@@ -413,7 +396,7 @@ mod tests {
assert_eq!(ops.count_user_blocks(uuid2).unwrap(), 1);
let mut batch = ms.database().batch();
ops.delete_user_blocks(&mut batch, hash1, &[vec![0x01]], "rev1")
ops.delete_user_blocks_by_cid(&mut batch, hash1, &[vec![0x01]])
.unwrap();
batch.commit().unwrap();
@@ -426,9 +409,9 @@ mod tests {
let (_dir, ms) = open_fresh();
let uuid1 = Uuid::new_v4();
let hash1 = UserHash::from_did("did:plc:user1");
let hash1 = UserHash::from_did("did:plc:whelk");
let uuid2 = Uuid::new_v4();
let hash2 = UserHash::from_did("did:plc:user2");
let hash2 = UserHash::from_did("did:plc:scallop");
let mut batch = ms.database().batch();
ms.user_hashes()
@@ -442,9 +425,9 @@ mod tests {
let ops = ms.user_block_ops();
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash1, &[vec![0x01]], "rev1")
ops.insert_user_blocks(&mut batch, hash1, &[vec![0x01]], &test_rev(10))
.unwrap();
ops.insert_user_blocks(&mut batch, hash2, &[vec![0x02]], "rev1")
ops.insert_user_blocks(&mut batch, hash2, &[vec![0x02]], &test_rev(10))
.unwrap();
batch.commit().unwrap();
@@ -469,13 +452,13 @@ mod tests {
];
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &cids, "rev1")
ops.insert_user_blocks(&mut batch, hash, &cids, &test_rev(10))
.unwrap();
batch.commit().unwrap();
assert_eq!(ops.count_user_blocks(uuid).unwrap(), 3);
let retrieved = ops.get_user_block_cids_since_rev(uuid, "").unwrap();
let retrieved = ops.get_user_block_cids_since_rev(uuid, None).unwrap();
assert_eq!(retrieved.len(), 3);
let mut expected = cids.clone();
expected.sort();
@@ -490,19 +473,7 @@ mod tests {
let cids = vec![vec![]];
let mut batch = ms.database().batch();
let result = ops.insert_user_blocks(&mut batch, hash, &cids, "rev1");
assert!(matches!(result, Err(MetastoreError::InvalidInput(_))));
}
#[test]
fn delete_empty_cid_is_rejected() {
let (_dir, ms) = open_fresh();
let (_uuid, hash) = setup_user(&ms);
let ops = ms.user_block_ops();
let cids = vec![vec![]];
let mut batch = ms.database().batch();
let result = ops.delete_user_blocks(&mut batch, hash, &cids, "rev1");
let result = ops.insert_user_blocks(&mut batch, hash, &cids, &test_rev(10));
assert!(matches!(result, Err(MetastoreError::InvalidInput(_))));
}
@@ -513,20 +484,24 @@ mod tests {
let ops = ms.user_block_ops();
let mut batch = ms.database().batch();
ops.insert_user_blocks(&mut batch, hash, &[vec![0x01]], "aaa")
ops.insert_user_blocks(&mut batch, hash, &[vec![0x01]], &test_rev(10))
.unwrap();
ops.insert_user_blocks(&mut batch, hash, &[vec![0x02]], "bbb")
ops.insert_user_blocks(&mut batch, hash, &[vec![0x02]], &test_rev(20))
.unwrap();
ops.insert_user_blocks(&mut batch, hash, &[vec![0x03]], "ddd")
ops.insert_user_blocks(&mut batch, hash, &[vec![0x03]], &test_rev(30))
.unwrap();
batch.commit().unwrap();
let result = ops.get_user_block_cids_since_rev(uuid, "aab").unwrap();
let result = ops
.get_user_block_cids_since_rev(uuid, Some(&test_rev(15)))
.unwrap();
assert_eq!(result.len(), 2);
assert_eq!(result[0], vec![0x02]);
assert_eq!(result[1], vec![0x03]);
let result = ops.get_user_block_cids_since_rev(uuid, "ccc").unwrap();
let result = ops
.get_user_block_cids_since_rev(uuid, Some(&test_rev(25)))
.unwrap();
assert_eq!(result.len(), 1);
assert_eq!(result[0], vec![0x03]);
}
@@ -1,9 +1,10 @@
use smallvec::SmallVec;
use tranquil_types::Tid;
use super::encoding::KeyBuilder;
use super::keys::{KeyTag, UserHash};
pub fn user_block_key(user_hash: UserHash, rev: &str, cid_bytes: &[u8]) -> SmallVec<[u8; 128]> {
pub fn user_block_key(user_hash: UserHash, rev: &Tid, cid_bytes: &[u8]) -> SmallVec<[u8; 128]> {
KeyBuilder::new()
.tag(KeyTag::USER_BLOCKS)
.u64(user_hash.raw())
@@ -19,7 +20,7 @@ pub fn user_block_user_prefix(user_hash: UserHash) -> SmallVec<[u8; 128]> {
.build()
}
pub fn user_block_rev_prefix(user_hash: UserHash, rev: &str) -> SmallVec<[u8; 128]> {
pub fn user_block_rev_prefix(user_hash: UserHash, rev: &Tid) -> SmallVec<[u8; 128]> {
KeyBuilder::new()
.tag(KeyTag::USER_BLOCKS)
.u64(user_hash.raw())
@@ -32,16 +33,26 @@ mod tests {
use super::*;
use crate::metastore::encoding::KeyReader;
fn test_rev(seq: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((seq >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
#[test]
fn user_block_key_roundtrip() {
let hash = UserHash::from_raw(0xDEAD_BEEF_CAFE_BABE);
let cid = [0x01, 0x71, 0x12, 0x20, 0xAB];
let key = user_block_key(hash, "3k2abcde", &cid);
let rev = test_rev(10);
let key = user_block_key(hash, &rev, &cid);
let mut reader = KeyReader::new(&key);
assert_eq!(reader.tag(), Some(KeyTag::USER_BLOCKS.raw()));
assert_eq!(reader.u64(), Some(0xDEAD_BEEF_CAFE_BABE));
assert_eq!(reader.string(), Some("3k2abcde".to_string()));
assert_eq!(reader.string(), Some(rev.as_str().to_owned()));
assert_eq!(reader.remaining(), &cid);
}
@@ -49,11 +60,13 @@ mod tests {
fn keys_sort_by_user_then_rev_then_cid() {
let h1 = UserHash::from_raw(1);
let h2 = UserHash::from_raw(2);
let rev_low = test_rev(10);
let rev_high = test_rev(20);
let k1 = user_block_key(h1, "abc", &[0x01]);
let k2 = user_block_key(h1, "abc", &[0x02]);
let k3 = user_block_key(h1, "def", &[0x01]);
let k4 = user_block_key(h2, "abc", &[0x01]);
let k1 = user_block_key(h1, &rev_low, &[0x01]);
let k2 = user_block_key(h1, &rev_low, &[0x02]);
let k3 = user_block_key(h1, &rev_high, &[0x01]);
let k4 = user_block_key(h2, &rev_low, &[0x01]);
assert!(k1.as_slice() < k2.as_slice());
assert!(k2.as_slice() < k3.as_slice());
@@ -64,15 +77,16 @@ mod tests {
fn user_prefix_is_prefix_of_rev_prefix() {
let hash = UserHash::from_raw(42);
let user_pfx = user_block_user_prefix(hash);
let rev_pfx = user_block_rev_prefix(hash, "some_rev");
let rev_pfx = user_block_rev_prefix(hash, &test_rev(10));
assert!(rev_pfx.as_slice().starts_with(user_pfx.as_slice()));
}
#[test]
fn rev_prefix_is_prefix_of_full_key() {
let hash = UserHash::from_raw(42);
let rev_pfx = user_block_rev_prefix(hash, "some_rev");
let full = user_block_key(hash, "some_rev", &[0x01, 0x02]);
let rev = test_rev(10);
let rev_pfx = user_block_rev_prefix(hash, &rev);
let full = user_block_key(hash, &rev, &[0x01, 0x02]);
assert!(full.as_slice().starts_with(rev_pfx.as_slice()));
}
@@ -80,12 +94,13 @@ mod tests {
fn cid_with_null_bytes_roundtrips() {
let hash = UserHash::from_raw(99);
let cid = [0x00, 0x01, 0x00, 0xFF];
let key = user_block_key(hash, "rev1", &cid);
let rev = test_rev(10);
let key = user_block_key(hash, &rev, &cid);
let mut reader = KeyReader::new(&key);
assert_eq!(reader.tag(), Some(KeyTag::USER_BLOCKS.raw()));
assert_eq!(reader.u64(), Some(99));
assert_eq!(reader.string(), Some("rev1".to_string()));
assert_eq!(reader.string(), Some(rev.as_str().to_owned()));
assert_eq!(reader.remaining(), &cid);
}
@@ -93,20 +108,22 @@ mod tests {
fn cid_with_double_null_bytes_roundtrips() {
let hash = UserHash::from_raw(99);
let cid = [0x00, 0x00, 0x01, 0x00, 0x00];
let key = user_block_key(hash, "rev1", &cid);
let rev = test_rev(10);
let key = user_block_key(hash, &rev, &cid);
let mut reader = KeyReader::new(&key);
assert_eq!(reader.tag(), Some(KeyTag::USER_BLOCKS.raw()));
assert_eq!(reader.u64(), Some(99));
assert_eq!(reader.string(), Some("rev1".to_string()));
assert_eq!(reader.string(), Some(rev.as_str().to_owned()));
assert_eq!(reader.remaining(), &cid);
}
#[test]
fn empty_cid_produces_key_equal_to_rev_prefix() {
let hash = UserHash::from_raw(42);
let rev_pfx = user_block_rev_prefix(hash, "rev1");
let full = user_block_key(hash, "rev1", &[]);
let rev = test_rev(10);
let rev_pfx = user_block_rev_prefix(hash, &rev);
let full = user_block_key(hash, &rev, &[]);
assert_eq!(full.as_slice(), rev_pfx.as_slice());
}
}
@@ -11,8 +11,8 @@ use tranquil_store::archival::{
use tranquil_store::eventlog::{EventLog, EventLogConfig, SegmentId};
use tranquil_store::sim_seed_range;
use common::{test_did, test_rev};
use tranquil_db_traits::{RepoEventType, SequenceNumber, SequencedEvent};
use tranquil_types::Did;
struct FlakyDestination {
inner: LocalArchivalDestination,
@@ -42,7 +42,7 @@ fn build_segments(segments_dir: &std::path::Path, n_events: u32) {
)
.unwrap();
(0..n_events).for_each(|i| {
let did = Did::from(format!("did:plc:archive{}", i % 8));
let did = test_did((i % 8) as u64);
let event = SequencedEvent {
seq: SequenceNumber::from_raw(0),
did: did.clone(),
@@ -57,7 +57,7 @@ fn build_segments(segments_dir: &std::path::Path, n_events: u32) {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("rev{i}"))),
rev: Some(test_rev(i as u64, 0)),
};
el.append_event(&did, RepoEventType::Commit, &event)
.unwrap();
+20 -14
View File
@@ -1,3 +1,7 @@
mod common;
use common::test_rev;
use std::path::Path;
use std::sync::Arc;
@@ -95,8 +99,9 @@ fn seed_blocks(store: &TestStore, range: std::ops::Range<u16>) {
}
fn seed_repo(store: &TestStore, name: &str, seed: u8) {
let did = tranquil_types::Did::from(format!("did:plc:{name}"));
let handle = tranquil_types::Handle::from(format!("{name}.test.invalid"));
let did = tranquil_types::Did::new(format!("did:plc:{name}")).expect("seed DID is valid");
let handle =
tranquil_types::Handle::new(format!("{name}.oyster.cafe")).expect("seed handle is valid");
let digest: [u8; 32] = std::array::from_fn(|i| seed.wrapping_add(i as u8));
let mh = multihash::Multihash::<64>::wrap(0x12, &digest).unwrap();
let cid = cid::Cid::new_v1(0x71, mh);
@@ -111,13 +116,13 @@ fn seed_repo(store: &TestStore, name: &str, seed: u8) {
&did,
&handle,
&cid_link,
&format!("rev_{seed}"),
&test_rev(0, seed as u64),
)
.unwrap();
}
fn seed_events(store: &TestStore, count: u16) {
let did = tranquil_types::Did::from("did:plc:evtest".to_string());
let did = tranquil_types::Did::new("did:plc:anemone").expect("did:plc:anemone is a valid DID");
(0..count).for_each(|i| {
let event = tranquil_db_traits::SequencedEvent {
seq: tranquil_db_traits::SequenceNumber::from_raw(i as i64 + 1),
@@ -133,7 +138,7 @@ fn seed_events(store: &TestStore, count: u16) {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("rev{i}"))),
rev: Some(common::test_rev(i as u64, 0)),
};
store
.eventlog
@@ -221,7 +226,7 @@ fn verify_restored_metastore(restored_dir: &Path, repo_names: &[&str]) {
let ops = ms.repo_ops();
repo_names.iter().for_each(|name| {
let did = tranquil_types::Did::from(format!("did:plc:{name}"));
let did = tranquil_types::Did::new(format!("did:plc:{name}")).expect("seed DID is valid");
let root = ops.get_repo_root_by_did(&did).unwrap();
assert!(
root.is_some(),
@@ -395,7 +400,7 @@ fn full_backup_and_restore_cycle() {
},
)
.unwrap();
let nel_did = tranquil_types::Did::from("did:plc:nel".to_string());
let nel_did = tranquil_types::Did::new("did:plc:nel").expect("did:plc:nel is a valid DID");
assert!(
ms.repo_ops()
.get_repo_root_by_did(&nel_did)
@@ -536,7 +541,7 @@ fn point_in_time_recovery_to_exact_sequence() {
let store = open_test_store();
seed_blocks(&store, 0..20);
seed_repo(&store, "pitr_user", 1);
seed_repo(&store, "limpet", 1);
seed_events(&store, 30);
store.metastore.persist().unwrap();
@@ -572,7 +577,7 @@ fn point_in_time_recovery_to_exact_sequence() {
assert!(pitr_result.events_replayed > 0);
verify_restored_blocks(restore_dir.path(), 0..20);
verify_restored_metastore(restore_dir.path(), &["pitr_user"]);
verify_restored_metastore(restore_dir.path(), &["limpet"]);
let restored_el = EventLog::open(
EventLogConfig {
@@ -621,7 +626,7 @@ fn crash_during_backup_does_not_corrupt_live_store() {
with_runtime(|| {
let store = open_test_store();
seed_blocks(&store, 0..50);
seed_repo(&store, "crash_test", 1);
seed_repo(&store, "whelk", 1);
seed_events(&store, 20);
store.metastore.persist().unwrap();
@@ -641,9 +646,10 @@ fn crash_during_backup_does_not_corrupt_live_store() {
seed_events(&store, 10);
assert!(store.eventlog.max_seq().raw() >= 30);
seed_repo(&store, "post_backup", 2);
seed_repo(&store, "scallop", 2);
store.metastore.persist().unwrap();
let did = tranquil_types::Did::from("did:plc:post_backup".to_string());
let did =
tranquil_types::Did::new("did:plc:scallop").expect("did:plc:scallop is a valid DID");
assert!(
store
.metastore
@@ -685,7 +691,7 @@ fn backup_during_concurrent_writes_produces_consistent_snapshot() {
let el_clone = Arc::clone(&store.eventlog);
let flag_clone2 = Arc::clone(&writer_flag);
let event_handle = std::thread::spawn(move || {
let did = tranquil_types::Did::from("did:plc:concurrent".to_string());
let did = tranquil_types::Did::new("did:plc:conch").expect("did:plc:conch is a valid DID");
let mut i = 0u32;
while flag_clone2.load(std::sync::atomic::Ordering::Relaxed) {
let event = tranquil_db_traits::SequencedEvent {
@@ -702,7 +708,7 @@ fn backup_during_concurrent_writes_produces_consistent_snapshot() {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("concurrent-{i}"))),
rev: Some(common::test_rev(i as u64, 0)),
};
let _ = el_clone.append_event(&did, tranquil_db_traits::RepoEventType::Commit, &event);
let _ = el_clone.sync();
+12 -2
View File
@@ -9,7 +9,7 @@ use tranquil_store::blockstore::{
use tranquil_store::eventlog::{EventLog, EventLogConfig};
use tranquil_store::metastore::{Metastore, MetastoreConfig};
use tranquil_store::{RealIO, SystemClock};
use tranquil_types::{CidLink, Did, Handle};
use tranquil_types::{CidLink, Did, Handle, Tid};
use uuid::Uuid;
pub const NAMES: &[&str] = &["olaren", "teq", "nel", "lyna", "bailey"];
@@ -33,7 +33,17 @@ pub fn block_data(seed: u32) -> Vec<u8> {
pub fn test_did(seed: u64) -> Did {
let name = NAMES[(seed as usize) % NAMES.len()];
Did::from(format!("did:plc:{name}{seed}"))
Did::new(format!("did:plc:{name}{seed}")).expect("generated DID is valid")
}
pub fn test_rev(sequence: u64, discriminator: u64) -> Tid {
const ALPHABET: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz";
let value = (sequence << 40) | (discriminator & 0xFF_FFFF_FFFF);
let s: String = (0..13)
.rev()
.map(|i| ALPHABET[((value >> (i * 5)) & 0x1F) as usize] as char)
.collect();
Tid::new(s).expect("generated TID is valid")
}
pub fn test_handle(seed: u64) -> Handle {
@@ -1,3 +1,5 @@
mod common;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
@@ -222,14 +224,12 @@ fn payload_round_trip() {
ops: Some(
serde_json::json!([{"action": "create", "path": "app.bsky.feed.post/abc"}]),
),
blobs: Some(vec![tranquil_types::CidLink::from(
"bafkreibtest".to_owned(),
)]),
blobs: Some(vec![common::test_cid_link(7)]),
blocks: None,
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from("rev1".to_owned())),
rev: Some(common::test_rev(1, 0)),
},
),
(
+16 -10
View File
@@ -1,5 +1,8 @@
mod common;
use std::path::Path;
use common::test_rev;
use proptest::prelude::*;
use rayon::prelude::*;
use tranquil_store::metastore::recovery::{
@@ -7,7 +10,7 @@ use tranquil_store::metastore::recovery::{
};
use tranquil_store::metastore::{Metastore, MetastoreConfig};
use tranquil_store::{sim_proptest_cases, sim_seed_range};
use tranquil_types::{CidLink, Did, Handle};
use tranquil_types::{CidLink, Did, Handle, Tid};
use uuid::Uuid;
const NAMES: &[&str] = &["olaren", "teq", "nel", "lyna", "bailey"];
@@ -43,6 +46,8 @@ fn test_uuid(seed: u64) -> Uuid {
Uuid::from_u128(seed as u128 | 0x4000_0000_0000_0000_8000_0000_0000_0000)
}
const TID_PATTERN: &str = "[234567abcdefghij][234567abcdefghijklmnopqrstuvwxyz]{12}";
fn arb_mutation_set() -> impl Strategy<Value = CommitMutationSet> {
let arb_upsert = (
"[a-z\\.]{5,20}",
@@ -67,7 +72,7 @@ fn arb_mutation_set() -> impl Strategy<Value = CommitMutationSet> {
(
prop::collection::vec(any::<u8>(), 0..64),
"[a-z0-9]{1,16}",
TID_PATTERN.prop_map(|rev| Tid::new(rev).expect("generated TID is valid")),
prop::collection::vec(arb_upsert, 0..20),
prop::collection::vec(arb_delete, 0..20),
prop::collection::vec(prop::collection::vec(any::<u8>(), 4..36), 0..20),
@@ -155,7 +160,7 @@ fn metastore_survives_abrupt_drop() {
let did = test_did(idx);
let handle = test_handle(idx);
let cid = test_cid_link((idx & 0xFF) as u8);
let rev = format!("rev{idx}");
let rev = test_rev(idx, 0);
repo_ops
.create_repo(db, uid, &did, &handle, &cid, &rev)
.unwrap();
@@ -173,9 +178,10 @@ fn metastore_survives_abrupt_drop() {
"seed={seed} user {i} repo_meta missing after abrupt drop"
);
let (_, meta) = result.unwrap();
let expected_rev = format!("rev{}", seed * 100 + i as u64);
let expected_rev = test_rev(seed * 100 + i as u64, 0);
assert_eq!(
meta.repo_rev, expected_rev,
meta.repo_rev,
expected_rev.as_str(),
"seed={seed} user {i} rev mismatch"
);
});
@@ -213,7 +219,7 @@ fn metastore_multi_crash_cycle() {
let handle = test_handle(idx);
let cid = test_cid_link((idx & 0xFF) as u8);
repo_ops
.create_repo(db, uid, &did, &handle, &cid, &format!("rev{idx}"))
.create_repo(db, uid, &did, &handle, &cid, &test_rev(idx, 0))
.unwrap();
expected_repos.push((uid, idx));
});
@@ -253,7 +259,7 @@ fn metastore_persisted_survives_unpersisted_lost() {
&test_did(idx),
&test_handle(idx),
&test_cid_link((idx & 0xFF) as u8),
&format!("rev{idx}"),
&test_rev(idx, 0),
)
.unwrap();
});
@@ -268,7 +274,7 @@ fn metastore_persisted_survives_unpersisted_lost() {
&test_did(extra_idx),
&test_handle(extra_idx),
&test_cid_link((extra_idx & 0xFF) as u8),
&format!("rev{extra_idx}"),
&test_rev(extra_idx, 0),
)
.unwrap();
}
@@ -316,7 +322,7 @@ fn metastore_user_hashes_reload_after_crash() {
&test_did(idx),
&test_handle(idx),
&test_cid_link((idx & 0xFF) as u8),
&format!("rev{idx}"),
&test_rev(idx, 0),
)
.unwrap();
});
@@ -354,7 +360,7 @@ fn metastore_handle_lookup_survives_crash() {
&did,
&handle,
&test_cid_link((idx & 0xFF) as u8),
&format!("rev{idx}"),
&test_rev(idx, 0),
)
.unwrap();
}
+18 -11
View File
@@ -16,7 +16,7 @@ use tranquil_store::{RealIO, SimClock, SimulatedIO, SystemClock, sim_seed_range,
use common::{
TestStores, assert_store_consistent, block_data, open_test_stores, test_cid, test_cid_link,
test_did, test_handle, test_uuid, with_runtime,
test_did, test_handle, test_rev, test_uuid, with_runtime,
};
use tranquil_db_traits::{
ImportBlock, ImportRecord, RepoEventType, SequenceNumber, SequencedEvent,
@@ -45,7 +45,7 @@ fn seed_repo(stores: &TestStores, idx: u64) -> Uuid {
&did,
&handle,
&cid,
&format!("rev{idx}"),
&test_rev(idx, 0),
)
.unwrap();
uid
@@ -66,7 +66,7 @@ fn make_commit_event(did: &Did, idx: u64) -> SequencedEvent {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("rev{idx}"))),
rev: Some(test_rev(idx, 0)),
}
}
@@ -477,7 +477,7 @@ fn sim_backup_during_concurrent_block_and_event_writes() {
});
let event_handle = s.spawn(|| {
let did = Did::from("did:plc:concurrent_writer".to_string());
let did = Did::new("did:plc:conch").expect("did:plc:conch is a valid DID");
std::iter::from_fn(|| writer_flag.load(Ordering::Relaxed).then_some(())).fold(
100u32,
|i, ()| {
@@ -495,7 +495,7 @@ fn sim_backup_during_concurrent_block_and_event_writes() {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("concurrent-{i}"))),
rev: Some(test_rev(i as u64, 0)),
};
let _ = stores
.eventlog
@@ -639,7 +639,7 @@ fn create_import_dirs(dir: &std::path::Path) {
fn import_fixture(seed: u64) -> (Vec<ImportBlock>, Vec<ImportRecord>) {
let block_count = ((seed % 12) + 4) as u32;
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = Nsid::new("app.bsky.feed.post").expect("app.bsky.feed.post is a valid NSID");
let blocks: Vec<ImportBlock> = (0..block_count)
.map(|i| ImportBlock {
cid_bytes: test_cid(i).to_vec(),
@@ -651,7 +651,7 @@ fn import_fixture(seed: u64) -> (Vec<ImportBlock>, Vec<ImportRecord>) {
let cid = cid::Cid::try_from(&test_cid(i)[..]).unwrap();
ImportRecord {
collection: collection.clone(),
rkey: Rkey::from(format!("3kimport{i:04}")),
rkey: Rkey::new(format!("3kimport{i:04}")).expect("generated rkey is valid"),
record_cid: CidLink::from_cid(&cid),
}
})
@@ -714,7 +714,14 @@ fn run_import_crash_scenario(seed: u64) {
{
let ms = import_metastore(dir.path());
ms.repo_ops()
.create_repo(ms.database(), user_id, &did, &handle, &root_cid, "rev0")
.create_repo(
ms.database(),
user_id,
&did,
&handle,
&root_cid,
&test_rev(0, 0),
)
.unwrap();
ms.persist().unwrap();
}
@@ -855,7 +862,7 @@ fn commit_block_records(stores: &SimStores, uid: Uuid, collection: &Nsid, block_
};
let rkeys: Vec<Rkey> = block_ids
.iter()
.map(|i| Rkey::from(format!("3k{i:08}")))
.map(|i| Rkey::new(format!("3k{i:08}")).expect("generated rkey is valid"))
.collect();
let links: Vec<CidLink> = block_ids.iter().map(|&i| block_cid_link(i)).collect();
let writes: Vec<RecordWrite<'_>> = rkeys
@@ -887,7 +894,7 @@ fn run_cross_store_fault_scenario(seed: u64) {
let repo_count = ((seed % 4) + 1) as u32;
let blocks_per_round = ((seed % 6) + 3) as u32;
let rounds = ((seed % 4) + 2) as u32;
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection = Nsid::new("app.bsky.feed.post").expect("app.bsky.feed.post is a valid NSID");
let root_base = (seed as u32).wrapping_mul(100_000);
let mut acked_blocks: Vec<u32> = Vec::new();
@@ -921,7 +928,7 @@ fn run_cross_store_fault_scenario(seed: u64) {
&did,
&handle,
&root,
&format!("rev{idx}"),
&test_rev(idx, 0),
)
.unwrap();
committed_repos.push((idx, uid));
+3 -4
View File
@@ -7,9 +7,8 @@ use rayon::prelude::*;
use tranquil_store::eventlog::{EventLog, EventLogConfig, EventSequence};
use tranquil_store::{FaultConfig, SimulatedIO, sim_seed_range};
use common::with_runtime;
use common::{test_did, test_rev, with_runtime};
use tranquil_db_traits::{RepoEventType, SequenceNumber, SequencedEvent};
use tranquil_types::Did;
fn open_sim_eventlog(
dir: &std::path::Path,
@@ -29,7 +28,7 @@ fn open_sim_eventlog(
}
fn append_seq(el: &EventLog<Arc<SimulatedIO>>, idx: u32) {
let did = Did::from(format!("did:plc:firehose{}", idx % 16));
let did = test_did((idx % 16) as u64);
let event = SequencedEvent {
seq: SequenceNumber::from_raw(0),
did: did.clone(),
@@ -44,7 +43,7 @@ fn append_seq(el: &EventLog<Arc<SimulatedIO>>, idx: u32) {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("rev{idx}"))),
rev: Some(test_rev(idx as u64, 0)),
};
el.append_event(&did, RepoEventType::Commit, &event)
.unwrap();
+31 -35
View File
@@ -15,7 +15,7 @@ use tranquil_store::{sim_seed_range, sim_single_seed};
use tranquil_types::{AtUri, CidLink, Nsid, Rkey};
use common::{
NAMES, open_test_stores, test_cid_link, test_did, test_handle, test_uuid, with_runtime,
open_test_stores, test_cid_link, test_did, test_handle, test_rev, test_uuid, with_runtime,
};
const CACHE_SIZE: u64 = 16 * 1024 * 1024;
@@ -55,7 +55,7 @@ fn build_commit_event(
did: &tranquil_types::Did,
prev_cid: &CidLink,
new_cid: &CidLink,
rev: &str,
rev: &tranquil_types::Tid,
) -> CommitEventData {
CommitEventData {
did: did.clone(),
@@ -66,7 +66,7 @@ fn build_commit_event(
blobs: None,
blocks: None,
prev_data_cid: None,
rev: Some(tranquil_types::Tid::from(rev.to_owned())),
rev: Some(rev.clone()),
}
}
@@ -110,7 +110,7 @@ impl MetastoreTestHarness {
&did,
&handle,
&cid,
&format!("rev0_{idx}"),
&test_rev(0, idx),
)
.unwrap_or_else(|e| panic!("create_repo idx={idx}: {e:?}"));
(uid, did, cid)
@@ -135,7 +135,6 @@ fn sim_apply_commit_crash_before_batch_commit_is_invisible() {
with_runtime(|| {
sim_seed_range().into_par_iter().for_each(|seed| {
let dir = tempfile::TempDir::new().unwrap();
let name = NAMES[(seed as usize) % NAMES.len()];
let record_count = (seed % 5) + 1;
{
@@ -146,7 +145,7 @@ fn sim_apply_commit_crash_before_batch_commit_is_invisible() {
.unwrap_or_else(|e| panic!("seed={seed} persist after create: {e:?}"));
let new_cid = test_cid_link(((seed + 50) & 0xFF) as u8);
let rev = format!("rev1_{name}{seed}");
let rev = test_rev(1, seed);
let upserts: Vec<RecordUpsert> = (0..record_count)
.map(|i| {
@@ -164,7 +163,7 @@ fn sim_apply_commit_crash_before_batch_commit_is_invisible() {
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: new_cid.clone(),
new_rev: tranquil_types::Tid::from(rev.to_owned()),
new_rev: rev.clone(),
record_upserts: upserts,
record_deletes: vec![],
backlinks_to_add: vec![],
@@ -202,7 +201,7 @@ fn sim_apply_commit_crash_before_batch_commit_is_invisible() {
let (_, repo_meta) = meta.unwrap();
assert_eq!(
repo_meta.repo_rev,
format!("rev0_{seed}"),
test_rev(0, seed).as_str(),
"seed={seed} repo rev must be initial if eventlog empty"
);
}
@@ -244,8 +243,6 @@ fn sim_apply_commit_atomicity_all_or_nothing() {
sim_seed_range().into_par_iter().for_each(|seed| {
let dir = tempfile::TempDir::new().unwrap();
let record_count = (seed % 8) + 2;
let name = NAMES[(seed as usize) % NAMES.len()];
{
let h = MetastoreTestHarness::open(dir.path());
let (uid, did, root_cid) = h.create_repo(seed);
@@ -254,7 +251,7 @@ fn sim_apply_commit_atomicity_all_or_nothing() {
.unwrap_or_else(|e| panic!("seed={seed} persist after create: {e:?}"));
let new_cid = test_cid_link(((seed + 77) & 0xFF) as u8);
let rev = format!("rev1_{name}{seed}");
let rev = test_rev(1, seed);
let collection = collection_nsid(seed);
let upserts: Vec<RecordUpsert> = (0..record_count)
@@ -294,7 +291,7 @@ fn sim_apply_commit_atomicity_all_or_nothing() {
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: new_cid.clone(),
new_rev: tranquil_types::Tid::from(rev.to_owned()),
new_rev: rev.clone(),
record_upserts: upserts,
record_deletes: vec![],
backlinks_to_add: backlinks,
@@ -326,7 +323,7 @@ fn sim_apply_commit_atomicity_all_or_nothing() {
.get_repo_meta(uid)
.unwrap_or_else(|e| panic!("seed={seed} get_repo_meta: {e:?}"))
.unwrap_or_else(|| panic!("seed={seed} repo meta missing"));
let rev = format!("rev1_{name}{seed}");
let rev = test_rev(1, seed);
assert_eq!(
repo_meta.repo_rev, rev,
"seed={seed} repo rev must match committed value"
@@ -378,7 +375,7 @@ fn sim_crash_recovery_cursor_tracks_last_durable_commit() {
.fold((root_cid, 0i64), |(prev_cid, persisted), commit_idx| {
let new_cid =
test_cid_link(((seed + commit_idx + 50) & 0xFF) as u8);
let rev = format!("rev{commit_idx}_{seed}");
let rev = test_rev(commit_idx, seed);
let collection = collection_nsid(seed + commit_idx);
let rkey = test_rkey(seed * 1000 + commit_idx);
@@ -387,7 +384,7 @@ fn sim_crash_recovery_cursor_tracks_last_durable_commit() {
did: did.clone(),
expected_root_cid: Some(prev_cid.clone()),
new_root_cid: new_cid.clone(),
new_rev: tranquil_types::Tid::from(rev.to_owned()),
new_rev: rev.clone(),
record_upserts: vec![RecordUpsert {
collection,
rkey,
@@ -480,7 +477,7 @@ fn sim_multi_commit_crash_cycle_consistency() {
};
let (total_records, last_rev) = (0..cycles).fold(
(0i64, String::new()),
(0i64, test_rev(0, seed)),
|(prev_total, _), cycle| {
let records_this_cycle = (seed.wrapping_add(cycle as u64) % 4) + 1;
@@ -526,14 +523,14 @@ fn sim_multi_commit_crash_cycle_consistency() {
let new_cid =
test_cid_link(((seed + cycle as u64 + 100) & 0xFF) as u8);
let rev = format!("rev{cycle}_{seed}");
let rev = test_rev(cycle as u64, seed);
let input = ApplyCommitInput {
user_id: uid,
did: did.clone(),
expected_root_cid: Some(current_cid.clone()),
new_root_cid: new_cid.clone(),
new_rev: tranquil_types::Tid::from(rev.to_owned()),
new_rev: rev.clone(),
record_upserts: upserts,
record_deletes: vec![],
backlinks_to_add: vec![],
@@ -571,7 +568,8 @@ fn sim_multi_commit_crash_cycle_consistency() {
.unwrap_or_else(|e| panic!("seed={seed} final get_repo_meta: {e:?}"))
.unwrap_or_else(|| panic!("seed={seed} final repo meta missing"));
assert_eq!(
repo_meta.repo_rev, last_rev,
repo_meta.repo_rev,
last_rev.as_str(),
"seed={seed} final rev must match last committed"
);
@@ -612,8 +610,6 @@ fn sim_handler_pool_shutdown_with_inflight_commits() {
let dir = tempfile::TempDir::new().unwrap();
let stores = open_test_stores(dir.path(), MAX_FILE_SIZE, CACHE_SIZE);
let bridge = Arc::new(EventLogBridge::new(Arc::clone(&stores.eventlog)));
let name = NAMES[(seed as usize) % NAMES.len()];
let repo_count = (seed % 5) + 2;
let repos: Vec<(uuid::Uuid, tranquil_types::Did, CidLink)> = (0..repo_count)
.map(|i| {
@@ -631,7 +627,7 @@ fn sim_handler_pool_shutdown_with_inflight_commits() {
&did,
&handle,
&cid,
&format!("rev0_{name}{idx}"),
&test_rev(0, idx),
)
.unwrap_or_else(|e| panic!("seed={seed} create_repo idx={idx}: {e:?}"));
(uid, did, cid)
@@ -651,13 +647,13 @@ fn sim_handler_pool_shutdown_with_inflight_commits() {
.map(|(i, (uid, did, root_cid))| {
let idx = seed * 100 + i as u64;
let new_cid = test_cid_link(((seed + i as u64 + 50) & 0xFF) as u8);
let rev = format!("rev1_{name}{idx}");
let rev = test_rev(1, idx);
let input = ApplyCommitInput {
user_id: *uid,
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: new_cid.clone(),
new_rev: tranquil_types::Tid::from(rev.to_owned()),
new_rev: rev.clone(),
record_upserts: vec![RecordUpsert {
collection: collection_nsid(seed + i as u64),
rkey: test_rkey(idx),
@@ -716,7 +712,7 @@ fn sim_handler_pool_shutdown_with_inflight_commits() {
(0..repo_count).for_each(|i| {
let idx = seed * 100 + i;
let uid = test_uuid(idx);
let expected_rev = format!("rev1_{name}{idx}");
let expected_rev = test_rev(1, idx);
let meta = h
.metastore
@@ -775,13 +771,13 @@ fn sim_record_deletes_through_crash_recovery() {
.collect();
let mid_cid = test_cid_link(((seed + 60) & 0xFF) as u8);
let rev1 = format!("rev1_{seed}");
let rev1 = test_rev(1, seed);
let insert_input = ApplyCommitInput {
user_id: uid,
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: mid_cid.clone(),
new_rev: tranquil_types::Tid::from(rev1.clone()),
new_rev: rev1.clone(),
record_upserts: upserts,
record_deletes: vec![],
backlinks_to_add: vec![],
@@ -810,13 +806,13 @@ fn sim_record_deletes_through_crash_recovery() {
.collect();
let final_cid = test_cid_link(((seed + 70) & 0xFF) as u8);
let rev2 = format!("rev2_{seed}");
let rev2 = test_rev(2, seed);
let delete_input = ApplyCommitInput {
user_id: uid,
did: did.clone(),
expected_root_cid: Some(mid_cid.clone()),
new_root_cid: final_cid.clone(),
new_rev: tranquil_types::Tid::from(rev2.clone()),
new_rev: rev2.clone(),
record_upserts: vec![],
record_deletes: deletes,
backlinks_to_add: vec![],
@@ -849,7 +845,7 @@ fn sim_record_deletes_through_crash_recovery() {
true => {
assert_eq!(
repo_meta.repo_rev,
format!("rev2_{seed}"),
test_rev(2, seed).as_str(),
"seed={seed} rev must reflect delete commit after recovery"
);
@@ -935,13 +931,13 @@ fn sim_obsolete_block_cids_through_crash_recovery() {
.collect();
let mid_cid = test_cid_link(((seed + 60) & 0xFF) as u8);
let rev1 = format!("rev1_{seed}");
let rev1 = test_rev(1, seed);
let insert_input = ApplyCommitInput {
user_id: uid,
did: did.clone(),
expected_root_cid: Some(root_cid.clone()),
new_root_cid: mid_cid.clone(),
new_rev: tranquil_types::Tid::from(rev1.clone()),
new_rev: rev1.clone(),
record_upserts: vec![RecordUpsert {
collection: collection_nsid(seed),
rkey: test_rkey(seed),
@@ -968,13 +964,13 @@ fn sim_obsolete_block_cids_through_crash_recovery() {
.collect();
let final_cid = test_cid_link(((seed + 70) & 0xFF) as u8);
let rev2 = format!("rev2_{seed}");
let rev2 = test_rev(2, seed);
let obsolete_input = ApplyCommitInput {
user_id: uid,
did: did.clone(),
expected_root_cid: Some(mid_cid.clone()),
new_root_cid: final_cid.clone(),
new_rev: tranquil_types::Tid::from(rev2.clone()),
new_rev: rev2.clone(),
record_upserts: vec![],
record_deletes: vec![],
backlinks_to_add: vec![],
@@ -1012,7 +1008,7 @@ fn sim_obsolete_block_cids_through_crash_recovery() {
true => {
assert_eq!(
repo_meta.repo_rev,
format!("rev2_{seed}"),
test_rev(2, seed).as_str(),
"seed={seed} rev must reflect obsolete commit after recovery"
);
@@ -10,7 +10,7 @@ use tranquil_store::metastore::repo_meta::RepoStatus;
use tranquil_store::metastore::{Metastore, MetastoreConfig};
use tranquil_store::sim_seed_range;
use common::{test_cid_link, test_did, test_handle, test_uuid};
use common::{test_cid_link, test_did, test_handle, test_rev, test_uuid};
use tranquil_db_traits::{RepoEventType, SequenceNumber, SequencedEvent};
use tranquil_types::{Did, Handle, Nsid, Rkey};
use uuid::Uuid;
@@ -35,7 +35,7 @@ fn seed_one_repo(ms: &Metastore, user_id: Uuid, did: &Did, handle: &Handle) {
did,
handle,
&test_cid_link(7),
"rev0",
&test_rev(0, 0),
)
.unwrap();
}
@@ -204,11 +204,12 @@ fn sim_list_records_range_scan_survives_restart() {
let user_id = test_uuid(seed);
let did = test_did(seed);
let handle = test_handle(seed);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection =
Nsid::new("app.bsky.feed.post").expect("app.bsky.feed.post is a valid NSID");
let count = ((seed % 30) + 10) as usize;
let rkeys: Vec<Rkey> = (0..count)
.map(|i| Rkey::from(format!("3k{i:04}")))
.map(|i| Rkey::new(format!("3k{i:04}")).expect("generated rkey is valid"))
.collect();
{
@@ -337,10 +338,11 @@ fn sim_major_compact_preserves_data() {
let user_id = test_uuid(seed);
let did = test_did(seed);
let handle = test_handle(seed);
let collection = Nsid::from("app.bsky.feed.post".to_string());
let collection =
Nsid::new("app.bsky.feed.post").expect("app.bsky.feed.post is a valid NSID");
let count = ((seed % 40) + 20) as usize;
let rkeys: Vec<Rkey> = (0..count)
.map(|i| Rkey::from(format!("3k{i:04}")))
.map(|i| Rkey::new(format!("3k{i:04}")).expect("generated rkey is valid"))
.collect();
let cids: Vec<_> = (0..count).map(|i| test_cid_link((i % 200) as u8)).collect();
let blob_cid = test_cid_link(123);
@@ -444,7 +446,7 @@ fn open_eventlog(dir: &std::path::Path, max_segment_size: u64) -> Arc<EventLog<R
fn append_n_events(el: &EventLog<RealIO>, n: u32) {
(0..n).for_each(|i| {
let did = Did::from(format!("did:plc:contiguity{i}"));
let did = test_did(i as u64);
let event = SequencedEvent {
seq: SequenceNumber::from_raw(0),
did: did.clone(),
@@ -459,7 +461,7 @@ fn append_n_events(el: &EventLog<RealIO>, n: u32) {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("rev{i}"))),
rev: Some(test_rev(i as u64, 0)),
};
el.append_event(&did, RepoEventType::Commit, &event)
.unwrap();
+3 -3
View File
@@ -9,7 +9,7 @@ use tranquil_store::sim_single_seed;
use common::{
Rng, TestStores, block_data, compact_all_sealed, open_test_stores, test_cid, test_cid_link,
test_did, test_handle, test_uuid,
test_did, test_handle, test_rev, test_uuid,
};
use tranquil_db_traits::{RepoEventType, SequenceNumber, SequencedEvent};
@@ -235,7 +235,7 @@ fn sim_soak_continuous_operations_with_crash_recovery() {
&did,
&handle,
&cid_link,
&format!("rev{idx}"),
&test_rev(idx, 0),
)
.unwrap();
oracle.add_repo(idx);
@@ -257,7 +257,7 @@ fn sim_soak_continuous_operations_with_crash_recovery() {
handle: None,
active: None,
status: None,
rev: Some(tranquil_types::Tid::from(format!("soak-rev-{event_idx}"))),
rev: Some(test_rev(event_idx, 0)),
};
s.eventlog
.append_event(&did, RepoEventType::Commit, &event)
+7 -6
View File
@@ -1,5 +1,5 @@
use axum::{
extract::{Query, RawQuery, State},
extract::{RawQuery, State},
http::StatusCode,
response::{IntoResponse, Response},
};
@@ -9,6 +9,7 @@ use serde::Deserialize;
use std::str::FromStr;
use tracing::error;
use tranquil_pds::api::error::ApiError;
use tranquil_pds::api::query::XrpcQuery;
use tranquil_pds::scheduled::generate_repo_car_from_user_blocks;
use tranquil_pds::state::AppState;
use tranquil_pds::sync::car::{encode_car_block, encode_car_header};
@@ -114,12 +115,12 @@ pub async fn get_blocks(State(state): State<AppState>, RawQuery(query): RawQuery
#[derive(Deserialize)]
pub struct GetRepoQuery {
pub did: Did,
pub since: Option<String>,
pub since: Option<Tid>,
}
pub async fn get_repo(
State(state): State<AppState>,
Query(query): Query<GetRepoQuery>,
XrpcQuery(query): XrpcQuery<GetRepoQuery>,
) -> Response {
let did = query.did;
let account =
@@ -139,7 +140,7 @@ pub async fn get_repo(
};
if let Some(since) = &query.since {
return get_repo_since(&state, &did, &head_cid, &Tid::from(since.clone())).await;
return get_repo_since(&state, &did, &head_cid, since).await;
}
let _permit = match state.repo_export_semaphore.try_acquire() {
@@ -195,7 +196,7 @@ async fn get_repo_since(state: &AppState, did: &Did, head_cid: &Cid, since: &Tid
let block_cid_bytes = match state
.repos
.repo
.get_user_block_cids_since_rev(user_id, since)
.get_user_block_cids_since_rev(user_id, Some(since))
.await
{
Ok(cids) => cids,
@@ -266,7 +267,7 @@ pub struct GetRecordQuery {
pub async fn get_record(
State(state): State<AppState>,
Query(query): Query<GetRecordQuery>,
XrpcQuery(query): XrpcQuery<GetRecordQuery>,
) -> Response {
use jacquard_repo::commit::Commit;
use jacquard_repo::mst::Mst;