Compare commits

...
9 Commits
Author SHA1 Message Date
LewisandTangled 8ccdd30cb3 fix(repo): use mst diff instead of full tree walk for obsolete blocks
Lewis: May this revision serve well! <lu5a@proton.me>
2026-04-13 17:43:07 +00:00
Lewis 7a67361993 fix(tranquil-store): checkpoint-hint race & missing dedup hints
Lewis: May this revision serve well! <lu5a@proton.me>
2026-04-13 11:10:17 +03:00
isabelandTangled cdbbaaccdf fix(nix/frontend): add nodejs 2026-04-12 23:08:39 +00:00
isabelandTangled 55d3b7f83d fix(nix/module): don't import using self 2026-04-12 22:50:00 +00:00
Gavin MoganandTangled f00b0231fb fix(Dockerfile): duplicate named stage causes failure to build 2026-04-12 18:53:21 +00:00
LewisandTangled 1dfbd27cce fix(postgres): semaphore on car endpoint & more efficient query
Lewis: May this revision serve well! <lu5a@proton.me>
2026-04-12 17:18:45 +00:00
isabelandTangled bc5e0e0446 build(frontend): use pnpm
deno is evil
2026-04-12 16:01:21 +00:00
Lewis 1c87ef5536 fix(tranquil-store): blockstore tweaks
Lewis: May this revision serve well! <lu5a@proton.me>
2026-04-12 17:35:59 +03:00
Lewis 255c7135f9 fix(auth): no bsky chat access when not specifically privileged to have it
Lewis: May this revision serve well! <lu5a@proton.me>
2026-04-12 15:56:17 +03:00
39 changed files with 3994 additions and 2376 deletions
Generated
+23 -22
View File
@@ -7405,7 +7405,7 @@ dependencies = [
[[package]]
name = "tranquil-api"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"anyhow",
"axum",
@@ -7456,7 +7456,7 @@ dependencies = [
[[package]]
name = "tranquil-auth"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"anyhow",
"base32",
@@ -7479,7 +7479,7 @@ dependencies = [
[[package]]
name = "tranquil-cache"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"base64 0.22.1",
@@ -7493,7 +7493,7 @@ dependencies = [
[[package]]
name = "tranquil-comms"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"base64 0.22.1",
@@ -7511,7 +7511,7 @@ dependencies = [
[[package]]
name = "tranquil-config"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"confique",
"serde",
@@ -7519,7 +7519,7 @@ dependencies = [
[[package]]
name = "tranquil-crypto"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"aes-gcm",
"base64 0.22.1",
@@ -7535,7 +7535,7 @@ dependencies = [
[[package]]
name = "tranquil-db"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"chrono",
@@ -7552,7 +7552,7 @@ dependencies = [
[[package]]
name = "tranquil-db-traits"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"base64 0.22.1",
@@ -7568,7 +7568,7 @@ dependencies = [
[[package]]
name = "tranquil-infra"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"bytes",
@@ -7579,7 +7579,7 @@ dependencies = [
[[package]]
name = "tranquil-lexicon"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"chrono",
"hickory-resolver",
@@ -7597,7 +7597,7 @@ dependencies = [
[[package]]
name = "tranquil-oauth"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"anyhow",
"axum",
@@ -7620,7 +7620,7 @@ dependencies = [
[[package]]
name = "tranquil-oauth-server"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"axum",
"base64 0.22.1",
@@ -7653,7 +7653,7 @@ dependencies = [
[[package]]
name = "tranquil-pds"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"aes-gcm",
"anyhow",
@@ -7745,7 +7745,7 @@ dependencies = [
[[package]]
name = "tranquil-repo"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"bytes",
"cid",
@@ -7757,7 +7757,7 @@ dependencies = [
[[package]]
name = "tranquil-ripple"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"backon",
@@ -7782,7 +7782,7 @@ dependencies = [
[[package]]
name = "tranquil-scopes"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"axum",
"futures",
@@ -7798,7 +7798,7 @@ dependencies = [
[[package]]
name = "tranquil-server"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"axum",
"clap",
@@ -7819,7 +7819,7 @@ dependencies = [
[[package]]
name = "tranquil-signal"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"chrono",
@@ -7842,7 +7842,7 @@ dependencies = [
[[package]]
name = "tranquil-storage"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"aws-config",
@@ -7859,7 +7859,7 @@ dependencies = [
[[package]]
name = "tranquil-store"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"async-trait",
"bytes",
@@ -7872,6 +7872,7 @@ dependencies = [
"jacquard-common",
"jacquard-repo",
"k256",
"libc",
"lsm-tree",
"memmap2",
"multihash",
@@ -7905,7 +7906,7 @@ dependencies = [
[[package]]
name = "tranquil-sync"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"anyhow",
"axum",
@@ -7927,7 +7928,7 @@ dependencies = [
[[package]]
name = "tranquil-types"
version = "0.5.1"
version = "0.5.3"
dependencies = [
"chrono",
"cid",
+1 -1
View File
@@ -26,7 +26,7 @@ members = [
]
[workspace.package]
version = "0.5.1"
version = "0.5.3"
edition = "2024"
license = "AGPL-3.0-or-later"
+5 -2
View File
@@ -1,7 +1,10 @@
FROM denoland/deno:alpine AS frontend
FROM node:24-alpine AS frontend
RUN corepack enable && corepack prepare pnpm@latest --activate
WORKDIR /app
COPY frontend/package.json frontend/pnpm-lock.yaml ./
RUN pnpm install --frozen-lockfile
COPY frontend/ ./
RUN deno task build
RUN pnpm build
FROM rust:1.92-alpine AS builder
RUN apk add --no-cache ca-certificates musl-dev pkgconfig openssl-dev openssl-libs-static mold clang protoc
+13 -4
View File
@@ -231,10 +231,19 @@ pub async fn verify_credential(
app_passwords
.into_iter()
.find(|app| bcrypt::verify(password, &app.password_hash).unwrap_or(false))
.map(|app| CredentialMatch::AppPassword {
name: app.name,
scopes: app.scopes,
controller_did: app.created_by_controller_did,
.map(|app| {
let scopes = app.scopes.unwrap_or_else(|| {
if app.privilege.is_privileged() {
"transition:generic transition:chat.bsky".to_string()
} else {
"transition:generic".to_string()
}
});
CredentialMatch::AppPassword {
name: app.name,
scopes: Some(scopes),
controller_did: app.created_by_controller_did,
}
})
}
@@ -132,7 +132,14 @@ pub async fn create_app_password(
};
(scope_result, Some(controller.clone()))
} else {
(input.scopes.clone(), None)
let scopes = match input.scopes {
Some(ref s) => s.clone(),
None => match input.privileged {
Some(false) => "transition:generic".to_string(),
_ => "transition:generic transition:chat.bsky".to_string(),
},
};
(Some(scopes), None)
};
let password = generate_app_password();
@@ -401,7 +401,7 @@ pub async fn create_passkey_account(
refresh_expires_at: refresh_expires,
login_type: tranquil_db_traits::LoginType::Modern,
mfa_verified: false,
scope: Some("transition:generic".to_string()),
scope: Some("transition:generic transition:chat.bsky".to_string()),
controller_did: None,
app_password_name: None,
};
+4
View File
@@ -720,6 +720,10 @@ pub struct FirehoseConfig {
#[config(env = "FIREHOSE_MAX_LAG", default = 5000)]
pub max_lag: u64,
/// Maximum concurrent full-repo exports, eg. getRepo without `since`.
#[config(env = "MAX_CONCURRENT_REPO_EXPORTS", default = 4)]
pub max_concurrent_repo_exports: usize,
/// List of relay / crawler notification URLs.
#[config(env = "CRAWLERS", parse_env = split_comma_list)]
pub crawlers: Option<Vec<String>>,
@@ -1339,7 +1339,7 @@ pub async fn complete_registration(
refresh_expires_at: refresh_meta.expires_at,
login_type: tranquil_db_traits::LoginType::Modern,
mfa_verified: false,
scope: Some("transition:generic".to_string()),
scope: Some("transition:generic transition:chat.bsky".to_string()),
controller_did: None,
app_password_name: None,
};
+20 -9
View File
@@ -3,7 +3,7 @@ use std::sync::Arc;
use std::time::Duration;
use chrono::Utc;
use tokio::time::interval;
use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
use tranquil_comms::{
@@ -75,17 +75,28 @@ impl CommsService {
);
}
info!(
poll_interval_secs = self.poll_interval.as_secs(),
poll_interval_ms = self.poll_interval.as_millis() as u64,
batch_size = self.batch_size,
channels = ?self.senders.keys().collect::<Vec<_>>(),
"Starting comms service"
);
let mut ticker = interval(self.poll_interval);
let base = self.poll_interval;
let max_backoff = Duration::from_secs(30);
let mut current_delay = base;
loop {
tokio::select! {
_ = ticker.tick() => {
if let Err(e) = self.process_batch().await {
error!(error = %e, "Failed to process comms batch");
_ = tokio::time::sleep(current_delay) => {
match self.process_batch().await {
Ok(had_work) => {
current_delay = match had_work {
true => base,
false => max_backoff.min(current_delay.saturating_mul(2)),
};
}
Err(e) => {
error!(error = %e, "Failed to process comms batch");
current_delay = max_backoff.min(current_delay.saturating_mul(2));
}
}
}
_ = shutdown.cancelled() => {
@@ -96,14 +107,14 @@ impl CommsService {
}
}
async fn process_batch(&self) -> Result<(), tranquil_db_traits::DbError> {
async fn process_batch(&self) -> Result<bool, tranquil_db_traits::DbError> {
let items = self.fetch_pending().await?;
if items.is_empty() {
return Ok(());
return Ok(false);
}
debug!(count = items.len(), "Processing comms batch");
futures::future::join_all(items.into_iter().map(|item| self.process_item(item))).await;
Ok(())
Ok(true)
}
async fn fetch_pending(&self) -> Result<Vec<QueuedComms>, tranquil_db_traits::DbError> {
+22 -54
View File
@@ -13,7 +13,6 @@ use jacquard_repo::mst::util::compute_cid;
use jacquard_repo::storage::BlockStore;
use k256::ecdsa::SigningKey;
use serde_json::{Value, json};
use std::collections::BTreeSet;
use std::str::FromStr;
use std::sync::Arc;
use tokio::sync::OwnedMutexGuard;
@@ -226,30 +225,6 @@ pub async fn begin_repo_write(
Ok((ctx, mst))
}
pub async fn compute_obsolete_cids(
original_mst: &Mst<TrackingBlockStore>,
new_mst: &Mst<TrackingBlockStore>,
original_root_cid: CommitCid,
) -> Result<Vec<Cid>, jacquard_repo::error::RepoError> {
let (old_nodes, new_nodes, old_leaves, new_leaves) = tokio::try_join!(
original_mst.collect_node_cids(),
new_mst.collect_node_cids(),
original_mst.leaves(),
new_mst.leaves(),
)?;
let old_nodes_set: BTreeSet<Cid> = old_nodes.into_iter().collect();
let new_nodes_set: BTreeSet<Cid> = new_nodes.into_iter().collect();
let old_leaf_set: BTreeSet<Cid> = old_leaves.iter().map(|(_, cid)| *cid).collect();
let new_leaf_set: BTreeSet<Cid> = new_leaves.iter().map(|(_, cid)| *cid).collect();
let removed_nodes = old_nodes_set.difference(&new_nodes_set).copied();
let removed_leaves = old_leaf_set.difference(&new_leaf_set).copied();
let obsolete: BTreeSet<Cid> = std::iter::once(original_root_cid.into_cid())
.chain(removed_nodes)
.chain(removed_leaves)
.collect();
Ok(obsolete.into_iter().collect())
}
pub async fn finalize_repo_write(
state: &AppState,
ctx: RepoWriteContext,
@@ -266,35 +241,28 @@ pub async fn finalize_repo_write(
let storage_for_diff = Arc::new(ctx.tracking_store.clone());
let original_settled = Mst::load(storage_for_diff.clone(), ctx.prev_data_cid, None);
let new_settled = Mst::load(storage_for_diff, new_mst_root, None);
let (obsolete_cids, new_tree_cids) = tokio::try_join!(
async {
compute_obsolete_cids(
&original_settled,
&new_settled,
CommitCid::from(ctx.current_root_cid),
)
.await
.map_err(|e| {
error!("MST diff failed during finalize_repo_write: {}", e);
ApiError::InternalError(Some("MST diff failed".into()))
})
},
async {
let (nodes, leaves) =
tokio::try_join!(new_settled.collect_node_cids(), new_settled.leaves(),).map_err(
|e| {
error!("new tree walk failed: {}", e);
ApiError::InternalError(None)
},
)?;
Ok::<Vec<Cid>, ApiError>(
nodes
.into_iter()
.chain(leaves.iter().map(|(_, cid)| *cid))
.collect(),
)
},
)?;
let new_tree_cids: Vec<Cid> = block_bytes.keys().copied().collect();
let obsolete_cids = match original_settled.diff(&new_settled).await {
Ok(diff) => {
let mut obsolete: Vec<Cid> = Vec::with_capacity(
1 + diff.removed_mst_blocks.len() + diff.removed_cids.len(),
);
obsolete.push(ctx.current_root_cid);
obsolete.extend(diff.removed_mst_blocks);
obsolete.extend(diff.removed_cids);
obsolete
}
Err(e) => {
error!(
"MST diff failed during finalize_repo_write: {e}. \
Proceeding with commit CID only; leaked blocks \
will be reclaimed by reachability GC."
);
vec![ctx.current_root_cid]
}
};
let result = commit_and_log(
state,
+14 -13
View File
@@ -667,6 +667,8 @@ async fn delete_account_data(
Ok(())
}
const CAR_BLOCK_BATCH_SIZE: usize = 500;
pub async fn generate_repo_car(
block_store: &AnyBlockStore,
head_cid: &Cid,
@@ -683,21 +685,20 @@ pub async fn generate_repo_car(
})
.collect();
let car_bytes = encode_car_header(head_cid).context("Failed to encode CAR header")?;
let mut car_bytes = encode_car_header(head_cid).context("Failed to encode CAR header")?;
let blocks = block_store
.get_many(&block_cids)
.await
.context("Failed to fetch blocks")?;
for chunk in block_cids.chunks(CAR_BLOCK_BATCH_SIZE) {
let blocks = block_store
.get_many(chunk)
.await
.context("Failed to fetch blocks")?;
let car_bytes = block_cids
.iter()
.zip(blocks.iter())
.filter_map(|(cid, block_opt)| block_opt.as_ref().map(|block| (cid, block)))
.fold(car_bytes, |mut acc, (cid, block)| {
acc.extend(encode_car_block(cid, block));
acc
});
chunk
.iter()
.zip(blocks.iter())
.filter_map(|(cid, block_opt)| block_opt.as_ref().map(|block| (cid, block)))
.for_each(|(cid, block)| car_bytes.extend(encode_car_block(cid, block)));
}
Ok(car_bytes)
}
+4
View File
@@ -50,6 +50,7 @@ pub struct AppState {
pub signal_sender: Option<Arc<tranquil_signal::SignalSlot>>,
pub signal_store_provider: Option<Arc<dyn tranquil_signal::SignalStoreProvider>>,
pub eventlog_segments_dir: Option<PathBuf>,
pub repo_export_semaphore: Arc<tokio::sync::Semaphore>,
}
#[derive(Debug, Clone, Copy)]
@@ -394,6 +395,9 @@ impl AppState {
signal_sender: None,
signal_store_provider,
eventlog_segments_dir,
repo_export_semaphore: Arc::new(tokio::sync::Semaphore::new(
cfg.firehose.max_concurrent_repo_exports,
)),
}
}
@@ -99,9 +99,8 @@ async fn test_check_account_status_returns_correct_block_count() {
after_delete_blocks
);
assert!(
after_delete_blocks >= initial_blocks,
"Block count after delete should be at least initial count (initial {}, now {})",
initial_blocks,
after_delete_blocks >= 2,
"Block count after delete should have at least commit + MST root (got {})",
after_delete_blocks
);
}
@@ -597,3 +597,155 @@ async fn test_request_account_delete() {
"Token should not be expired"
);
}
async fn create_app_password_session(
client: &reqwest::Client,
did: &str,
main_jwt: &str,
name: &str,
body: Value,
) -> (String, Value) {
let base = base_url().await;
let create_res = client
.post(format!(
"{}/xrpc/com.atproto.server.createAppPassword",
base
))
.bearer_auth(main_jwt)
.json(&body)
.send()
.await
.expect("Failed to create app password");
assert_eq!(create_res.status(), StatusCode::OK);
let app_pass: Value = create_res.json().await.unwrap();
let password = app_pass["password"].as_str().unwrap().to_string();
let scopes_response = app_pass.clone();
let login_res = client
.post(format!("{}/xrpc/com.atproto.server.createSession", base))
.json(&json!({ "identifier": did, "password": password }))
.send()
.await
.expect("Failed to login with app password");
assert_eq!(login_res.status(), StatusCode::OK, "App password login for '{}' failed", name);
let session: Value = login_res.json().await.unwrap();
let jwt = session["accessJwt"].as_str().unwrap().to_string();
(jwt, scopes_response)
}
async fn try_chat_service_auth(client: &reqwest::Client, jwt: &str) -> StatusCode {
let base = base_url().await;
let res = client
.get(format!(
"{}/xrpc/com.atproto.server.getServiceAuth",
base
))
.bearer_auth(jwt)
.query(&[
("aud", "did:web:api.bsky.app"),
("lxm", "chat.bsky.convo.listConvos"),
])
.send()
.await
.expect("Failed to call getServiceAuth");
res.status()
}
#[tokio::test]
async fn test_app_password_non_privileged_blocks_chat() {
let client = client();
let (did, jwt) = setup_new_user("appscope-nonchat").await;
let (app_jwt, create_body) = create_app_password_session(
&client,
&did,
&jwt,
"non-privileged",
json!({ "name": "NoChatApp", "privileged": false }),
)
.await;
assert_eq!(
create_body["scopes"].as_str().unwrap(),
"transition:generic",
"Non-privileged app password should not have chat scope"
);
let status = try_chat_service_auth(&client, &app_jwt).await;
assert_eq!(
status,
StatusCode::FORBIDDEN,
"Non-privileged app password must not access chat methods"
);
}
#[tokio::test]
async fn test_app_password_privileged_allows_chat() {
let client = client();
let (did, jwt) = setup_new_user("appscope-chat").await;
let (app_jwt, create_body) = create_app_password_session(
&client,
&did,
&jwt,
"privileged",
json!({ "name": "ChatApp", "privileged": true }),
)
.await;
assert_eq!(
create_body["scopes"].as_str().unwrap(),
"transition:generic transition:chat.bsky",
"Privileged app password should have chat scope"
);
let status = try_chat_service_auth(&client, &app_jwt).await;
assert_eq!(
status,
StatusCode::OK,
"Privileged app password should access chat methods"
);
}
#[tokio::test]
async fn test_app_password_no_privileged_field_allows_chat() {
let client = client();
let (did, jwt) = setup_new_user("appscope-full").await;
let (app_jwt, create_body) = create_app_password_session(
&client,
&did,
&jwt,
"full-access",
json!({ "name": "FullApp" }),
)
.await;
assert_eq!(
create_body["scopes"].as_str().unwrap(),
"transition:generic transition:chat.bsky",
"App password without privileged field should default to full access"
);
let status = try_chat_service_auth(&client, &app_jwt).await;
assert_eq!(
status,
StatusCode::OK,
"Full-access app password should access chat methods"
);
}
#[tokio::test]
async fn test_app_password_explicit_scopes_respected() {
let client = client();
let (did, jwt) = setup_new_user("appscope-explicit").await;
let (app_jwt, create_body) = create_app_password_session(
&client,
&did,
&jwt,
"explicit-scopes",
json!({ "name": "ScopedApp", "scopes": "transition:generic" }),
)
.await;
assert_eq!(
create_body["scopes"].as_str().unwrap(),
"transition:generic",
"Explicit scopes should be stored as-is"
);
let status = try_chat_service_auth(&client, &app_jwt).await;
assert_eq!(
status,
StatusCode::FORBIDDEN,
"App password with only transition:generic should not access chat"
);
}
@@ -0,0 +1,543 @@
use std::collections::BTreeSet;
use std::sync::Arc;
use cid::Cid;
use jacquard_repo::mst::Mst;
use jacquard_repo::storage::MemoryBlockStore;
fn test_cid(n: u32) -> Cid {
let data = n.to_be_bytes();
let mut buf = [0u8; 32];
buf[..4].copy_from_slice(&data);
buf[4] = (n >> 8) as u8 ^ 0xAB;
buf[5] = (n & 0xFF) as u8 ^ 0xCD;
let mh = multihash::Multihash::wrap(0x12, &buf).unwrap();
Cid::new_v1(0x71, mh)
}
async fn compute_obsolete_full_walk<S: jacquard_repo::storage::BlockStore + Sync + Send + 'static>(
old: &Mst<S>,
new: &Mst<S>,
) -> BTreeSet<Cid> {
let old_nodes = old.collect_node_cids().await.unwrap();
let new_nodes = new.collect_node_cids().await.unwrap();
let old_leaves = old.leaves().await.unwrap();
let new_leaves = new.leaves().await.unwrap();
let old_nodes_set: BTreeSet<Cid> = old_nodes.into_iter().collect();
let new_nodes_set: BTreeSet<Cid> = new_nodes.into_iter().collect();
let old_leaf_set: BTreeSet<Cid> = old_leaves.iter().map(|(_, cid)| *cid).collect();
let new_leaf_set: BTreeSet<Cid> = new_leaves.iter().map(|(_, cid)| *cid).collect();
old_nodes_set
.difference(&new_nodes_set)
.copied()
.chain(old_leaf_set.difference(&new_leaf_set).copied())
.collect()
}
fn compute_obsolete_from_diff(
diff: &jacquard_repo::mst::diff::MstDiff,
) -> BTreeSet<Cid> {
diff.removed_mst_blocks
.iter()
.copied()
.chain(diff.removed_cids.iter().copied())
.collect()
}
async fn assert_equivalence(
old_records: &[(String, u32)],
new_records: &[(String, u32)],
scenario: &str,
) {
let storage = Arc::new(MemoryBlockStore::new());
let mut old_tree = Mst::new(storage.clone());
for (key, val) in old_records {
old_tree = old_tree.add(key, test_cid(*val)).await.unwrap();
}
let old_root = old_tree.persist().await.unwrap();
let mut new_tree = Mst::new(storage.clone());
for (key, val) in new_records {
new_tree = new_tree.add(key, test_cid(*val)).await.unwrap();
}
let new_root = new_tree.persist().await.unwrap();
let old_settled = Mst::load(storage.clone(), old_root, None);
let new_settled = Mst::load(storage.clone(), new_root, None);
let full_walk_obsolete = compute_obsolete_full_walk(&old_settled, &new_settled).await;
let old_for_diff = Mst::load(storage.clone(), old_root, None);
let new_for_diff = Mst::load(storage, new_root, None);
let diff = old_for_diff.diff(&new_for_diff).await.unwrap();
let diff_obsolete = compute_obsolete_from_diff(&diff);
assert_eq!(
full_walk_obsolete, diff_obsolete,
"MISMATCH in scenario: {scenario}\n full_walk count: {}\n diff count: {}\n in full_walk but not diff: {:?}\n in diff but not full_walk: {:?}",
full_walk_obsolete.len(),
diff_obsolete.len(),
full_walk_obsolete.difference(&diff_obsolete).collect::<Vec<_>>(),
diff_obsolete.difference(&full_walk_obsolete).collect::<Vec<_>>(),
);
}
fn make_key(collection: &str, i: u32) -> String {
format!("{collection}/{i:06}")
}
fn generate_records(collection: &str, range: std::ops::Range<u32>) -> Vec<(String, u32)> {
range.map(|i| (make_key(collection, i), i)).collect()
}
fn generate_multi_collection_records(
collections: &[&str],
per_collection: u32,
) -> Vec<(String, u32)> {
collections
.iter()
.enumerate()
.flat_map(|(ci, coll)| {
let base = ci as u32 * per_collection;
(0..per_collection).map(move |i| (make_key(coll, i), base + i))
})
.collect()
}
fn apply_scattered_updates(
records: &[(String, u32)],
stride: usize,
cid_offset: u32,
) -> Vec<(String, u32)> {
records
.iter()
.enumerate()
.map(|(idx, (key, val))| {
if idx % stride == 0 {
(key.clone(), val + cid_offset)
} else {
(key.clone(), *val)
}
})
.collect()
}
fn remove_every_nth(records: &[(String, u32)], n: usize) -> Vec<(String, u32)> {
records
.iter()
.enumerate()
.filter(|(idx, _)| idx % n != 0)
.map(|(_, r)| r.clone())
.collect()
}
fn remove_range(records: &[(String, u32)], start: usize, count: usize) -> Vec<(String, u32)> {
records
.iter()
.enumerate()
.filter(|(idx, _)| *idx < start || *idx >= start + count)
.map(|(_, r)| r.clone())
.collect()
}
fn keep_only_collection(records: &[(String, u32)], collection: &str) -> Vec<(String, u32)> {
records
.iter()
.filter(|(key, _)| key.starts_with(collection))
.cloned()
.collect()
}
fn append_records(
base: &[(String, u32)],
collection: &str,
range: std::ops::Range<u32>,
cid_base: u32,
) -> Vec<(String, u32)> {
let mut result = base.to_vec();
result.extend(range.map(|i| (make_key(collection, i), cid_base + i)));
result.sort_by(|(a, _), (b, _)| a.cmp(b));
result
}
#[tokio::test]
async fn massive_tree_single_create() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec = append_records(&old, "app.bsky.feed.post", 2000..2001, 2000);
assert_equivalence(&old, &new_rec, "2000 records + 1 create").await;
}
#[tokio::test]
async fn massive_tree_single_delete() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec = remove_range(&old, 1000, 1);
assert_equivalence(&old, &new_rec, "2000 records - 1 delete from middle").await;
}
#[tokio::test]
async fn massive_tree_single_update() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec: Vec<_> = old
.iter()
.map(|(k, v)| {
if k == "app.bsky.feed.post/001000" {
(k.clone(), v + 50000)
} else {
(k.clone(), *v)
}
})
.collect();
assert_equivalence(&old, &new_rec, "2000 records - 1 update in middle").await;
}
#[tokio::test]
async fn massive_tree_scattered_updates_every_3rd() {
let old = generate_records("app.bsky.feed.post", 0..1500);
let new_rec = apply_scattered_updates(&old, 3, 10000);
assert_equivalence(&old, &new_rec, "1500 records - update every 3rd").await;
}
#[tokio::test]
async fn massive_tree_scattered_updates_every_7th() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec = apply_scattered_updates(&old, 7, 20000);
assert_equivalence(&old, &new_rec, "2000 records - update every 7th").await;
}
#[tokio::test]
async fn massive_tree_delete_every_2nd() {
let old = generate_records("app.bsky.feed.post", 0..1000);
let new_rec = remove_every_nth(&old, 2);
assert_equivalence(&old, &new_rec, "1000 records - delete every 2nd").await;
}
#[tokio::test]
async fn massive_tree_delete_every_5th() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec = remove_every_nth(&old, 5);
assert_equivalence(&old, &new_rec, "2000 records - delete every 5th").await;
}
#[tokio::test]
async fn massive_tree_delete_first_half() {
let old = generate_records("app.bsky.feed.post", 0..1500);
let new_rec = remove_range(&old, 0, 750);
assert_equivalence(&old, &new_rec, "1500 records - delete first 750").await;
}
#[tokio::test]
async fn massive_tree_delete_last_half() {
let old = generate_records("app.bsky.feed.post", 0..1500);
let new_rec = remove_range(&old, 750, 750);
assert_equivalence(&old, &new_rec, "1500 records - delete last 750").await;
}
#[tokio::test]
async fn massive_tree_delete_middle_chunk() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec = remove_range(&old, 800, 400);
assert_equivalence(&old, &new_rec, "2000 records - delete 400 from middle").await;
}
#[tokio::test]
async fn empty_to_massive() {
let new_rec = generate_records("app.bsky.feed.post", 0..1500);
assert_equivalence(&[], &new_rec, "empty to 1500 records").await;
}
#[tokio::test]
async fn massive_to_empty() {
let old = generate_records("app.bsky.feed.post", 0..1500);
assert_equivalence(&old, &[], "1500 records to empty").await;
}
#[tokio::test]
async fn massive_complete_replacement() {
let old = generate_records("app.bsky.feed.post", 0..1000);
let new_rec = generate_records("app.bsky.feed.post", 1000..2000);
assert_equivalence(&old, &new_rec, "1000 records fully replaced with 1000 different").await;
}
#[tokio::test]
async fn massive_no_change() {
let records = generate_records("app.bsky.feed.post", 0..1500);
assert_equivalence(&records, &records, "1500 records unchanged").await;
}
#[tokio::test]
async fn multi_collection_5_collections_500_each() {
let collections = [
"app.bsky.feed.like",
"app.bsky.feed.post",
"app.bsky.feed.repost",
"app.bsky.graph.follow",
"app.bsky.graph.block",
];
let old = generate_multi_collection_records(&collections, 500);
let new_rec = apply_scattered_updates(&old, 4, 30000);
assert_equivalence(&old, &new_rec, "5 collections x 500 records - update every 4th").await;
}
#[tokio::test]
async fn multi_collection_wipe_one_collection() {
let collections = [
"app.bsky.feed.like",
"app.bsky.feed.post",
"app.bsky.feed.repost",
"app.bsky.graph.follow",
];
let old = generate_multi_collection_records(&collections, 400);
let new_rec: Vec<_> = old
.iter()
.filter(|(key, _)| !key.starts_with("app.bsky.feed.repost"))
.cloned()
.collect();
assert_equivalence(&old, &new_rec, "4 collections x 400 - wipe repost collection").await;
}
#[tokio::test]
async fn multi_collection_keep_only_one() {
let collections = [
"app.bsky.feed.like",
"app.bsky.feed.post",
"app.bsky.feed.repost",
"app.bsky.graph.follow",
"app.bsky.graph.block",
];
let old = generate_multi_collection_records(&collections, 300);
let new_rec = keep_only_collection(&old, "app.bsky.feed.post");
assert_equivalence(&old, &new_rec, "5 collections x 300 - keep only posts").await;
}
#[tokio::test]
async fn multi_collection_add_new_collection() {
let old_collections = [
"app.bsky.feed.like",
"app.bsky.feed.post",
];
let old = generate_multi_collection_records(&old_collections, 500);
let new_rec = append_records(&old, "app.bsky.graph.follow", 0..500, 40000);
assert_equivalence(&old, &new_rec, "2 collections x 500 + add 500 follows").await;
}
#[tokio::test]
async fn mixed_ops_massive_tree() {
let collections = [
"app.bsky.feed.like",
"app.bsky.feed.post",
"app.bsky.feed.repost",
"app.bsky.graph.follow",
];
let old = generate_multi_collection_records(&collections, 400);
let mut new_rec: Vec<_> = old
.iter()
.filter(|(key, _)| !key.starts_with("app.bsky.feed.repost"))
.enumerate()
.map(|(idx, (key, val))| {
if key.starts_with("app.bsky.feed.like") && idx % 3 == 0 {
(key.clone(), val + 50000)
} else {
(key.clone(), *val)
}
})
.collect();
new_rec.extend((0..200u32).map(|i| (make_key("app.bsky.graph.block", i), 60000 + i)));
new_rec.sort_by(|(a, _), (b, _)| a.cmp(b));
assert_equivalence(
&old,
&new_rec,
"4 collections x 400: wipe reposts, update every 3rd like, add 200 blocks",
)
.await;
}
#[tokio::test]
async fn grow_tree_by_double() {
let old = generate_records("app.bsky.feed.post", 0..1000);
let new_rec = generate_records("app.bsky.feed.post", 0..2000);
assert_equivalence(&old, &new_rec, "grow from 1000 to 2000").await;
}
#[tokio::test]
async fn shrink_tree_by_half() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec = generate_records("app.bsky.feed.post", 0..1000);
assert_equivalence(&old, &new_rec, "shrink from 2000 to 1000").await;
}
#[tokio::test]
async fn interleaved_keys_disjoint_ranges() {
let old: Vec<_> = (0..1000u32)
.map(|i| (make_key("app.bsky.feed.post", i * 2), i))
.collect();
let new_rec: Vec<_> = (0..1000u32)
.map(|i| (make_key("app.bsky.feed.post", i * 2 + 1), i + 10000))
.collect();
assert_equivalence(&old, &new_rec, "1000 even-keyed records replaced by 1000 odd-keyed").await;
}
#[tokio::test]
async fn sparse_keys_wide_gaps() {
let old: Vec<_> = (0..500u32)
.map(|i| (make_key("app.bsky.feed.post", i * 100), i))
.collect();
let new_rec: Vec<_> = (0..500u32)
.map(|i| {
if i % 10 == 0 {
(make_key("app.bsky.feed.post", i * 100), i + 70000)
} else {
(make_key("app.bsky.feed.post", i * 100), i)
}
})
.collect();
assert_equivalence(&old, &new_rec, "500 sparse keys - update every 10th").await;
}
#[tokio::test]
async fn many_collections_few_records_each() {
let collections: Vec<String> = (0..50u32)
.map(|i| format!("com.example.lexicon{i:02}.record"))
.collect();
let old: Vec<_> = collections
.iter()
.enumerate()
.flat_map(|(ci, coll)| {
let base = ci as u32 * 20;
(0..20u32).map(move |i| (make_key(coll, i), base + i))
})
.collect();
let new_rec: Vec<_> = old
.iter()
.enumerate()
.filter_map(|(idx, (key, val))| {
if idx % 15 == 0 {
None
} else if idx % 7 == 0 {
Some((key.clone(), val + 80000))
} else {
Some((key.clone(), *val))
}
})
.collect();
assert_equivalence(&old, &new_rec, "50 collections x 20 records - delete every 15th, update every 7th").await;
}
#[tokio::test]
async fn update_all_records() {
let old = generate_records("app.bsky.feed.post", 0..1000);
let new_rec: Vec<_> = old
.iter()
.map(|(key, val)| (key.clone(), val + 90000))
.collect();
assert_equivalence(&old, &new_rec, "1000 records - update every single one").await;
}
#[tokio::test]
async fn delete_all_but_one() {
let old = generate_records("app.bsky.feed.post", 0..1500);
let new_rec = vec![old[750].clone()];
assert_equivalence(&old, &new_rec, "1500 records - delete all but middle one").await;
}
#[tokio::test]
async fn one_to_massive() {
let old = vec![(make_key("app.bsky.feed.post", 500), 500u32)];
let new_rec = generate_records("app.bsky.feed.post", 0..1500);
assert_equivalence(&old, &new_rec, "1 record to 1500 records").await;
}
#[tokio::test]
async fn delete_head_and_tail() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let new_rec: Vec<_> = old[200..1800].to_vec();
assert_equivalence(&old, &new_rec, "2000 records - delete first 200 and last 200").await;
}
#[tokio::test]
async fn keep_head_and_tail_only() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let mut new_rec: Vec<_> = old[..100].to_vec();
new_rec.extend_from_slice(&old[1900..]);
assert_equivalence(&old, &new_rec, "2000 records - keep only first 100 and last 100").await;
}
#[tokio::test]
async fn massive_tree_update_first_and_last() {
let old = generate_records("app.bsky.feed.post", 0..2000);
let mut new_rec = old.clone();
new_rec[0].1 += 99000;
new_rec[1999].1 += 99000;
assert_equivalence(&old, &new_rec, "2000 records - update only first and last").await;
}
#[tokio::test]
async fn overlapping_collection_swap() {
let old_collections = [
"app.bsky.feed.like",
"app.bsky.feed.post",
"app.bsky.feed.repost",
];
let old = generate_multi_collection_records(&old_collections, 500);
let mut new_rec: Vec<_> = old
.iter()
.filter(|(key, _)| key.starts_with("app.bsky.feed.post"))
.cloned()
.collect();
new_rec.extend((0..500u32).map(|i| (make_key("app.bsky.graph.follow", i), 70000 + i)));
new_rec.extend((0..500u32).map(|i| (make_key("app.bsky.graph.block", i), 71000 + i)));
new_rec.sort_by(|(a, _), (b, _)| a.cmp(b));
assert_equivalence(
&old,
&new_rec,
"swap 2 of 3 collections, keep 1 (posts), 500 each",
)
.await;
}
#[tokio::test]
async fn swiss_cheese_deletions() {
let old = generate_records("app.bsky.feed.post", 0..1500);
let new_rec: Vec<_> = old
.iter()
.enumerate()
.filter(|(idx, _)| {
let bucket = idx / 50;
bucket % 3 != 0
})
.map(|(_, r)| r.clone())
.collect();
assert_equivalence(&old, &new_rec, "1500 records - delete every 3rd chunk of 50").await;
}
#[tokio::test]
async fn mixed_ops_with_key_density_change() {
let old: Vec<_> = (0..1000u32)
.map(|i| (make_key("app.bsky.feed.post", i * 3), i))
.collect();
let mut new_rec: Vec<_> = old
.iter()
.filter(|(_, val)| val % 4 != 0)
.cloned()
.collect();
new_rec.extend((0..500u32).map(|i| {
(make_key("app.bsky.feed.post", i * 3 + 1), i + 100000)
}));
new_rec.sort_by(|(a, _), (b, _)| a.cmp(b));
assert_equivalence(
&old,
&new_rec,
"1000 sparse records: delete every 4th, insert 500 in gaps",
)
.await;
}
@@ -0,0 +1,147 @@
mod common;
mod helpers;
use common::*;
use helpers::*;
use reqwest::StatusCode;
use std::sync::Once;
static SET_SEMAPHORE: Once = Once::new();
fn ensure_low_semaphore() {
SET_SEMAPHORE.call_once(|| unsafe {
std::env::set_var("MAX_CONCURRENT_REPO_EXPORTS", "1");
});
}
#[tokio::test]
async fn test_get_repo_succeeds_with_many_records() {
ensure_low_semaphore();
let client = client();
let (did, jwt) = setup_new_user("sync-batched-car").await;
let create_futures = (0..20).map(|i| {
let client = &client;
let did = &did;
let jwt = &jwt;
async move {
create_post(client, did, jwt, &format!("Batch test post {}", i)).await;
}
});
futures::future::join_all(create_futures).await;
let res = client
.get(format!(
"{}/xrpc/com.atproto.sync.getRepo",
base_url().await
))
.query(&[("did", did.as_str())])
.send()
.await
.expect("Failed to send getRepo request");
assert_eq!(res.status(), StatusCode::OK);
assert_eq!(
res.headers()
.get("content-type")
.and_then(|h| h.to_str().ok()),
Some("application/vnd.ipld.car")
);
let car_bytes = res.bytes().await.expect("Failed to read response body");
assert!(
car_bytes.len() > 200,
"CAR with 20 records should have substantial data, got {} bytes",
car_bytes.len()
);
}
#[tokio::test]
async fn test_get_repo_semaphore_rejects_excess_concurrency() {
ensure_low_semaphore();
let client = client();
let (did, jwt) = setup_new_user("sync-semaphore").await;
for i in 0..50 {
create_post(&client, &did, &jwt, &format!("Padding post {}", i)).await;
}
let base = base_url().await;
let concurrent_requests = 10;
let request_futures = (0..concurrent_requests).map(|_| {
let client = client.clone();
let did = did.clone();
async move {
client
.get(format!("{}/xrpc/com.atproto.sync.getRepo", base))
.query(&[("did", did.as_str())])
.send()
.await
.expect("Failed to send request")
.status()
}
});
let statuses: Vec<StatusCode> = futures::future::join_all(request_futures).await;
let ok_count = statuses.iter().filter(|s| **s == StatusCode::OK).count();
let rejected_count = statuses
.iter()
.filter(|s| **s == StatusCode::SERVICE_UNAVAILABLE)
.count();
assert!(ok_count >= 1, "at least one request should succeed");
assert!(
rejected_count > 0,
"semaphore=1 with {} concurrent requests, expected some 503 rejections",
concurrent_requests
);
assert!(
ok_count + rejected_count == statuses.len(),
"expected only 200 or 503 responses: {:?}",
statuses
);
}
#[tokio::test]
async fn test_get_repo_since_not_affected_by_semaphore() {
ensure_low_semaphore();
let client = client();
let (did, jwt) = setup_new_user("sync-since-no-sem").await;
create_post(&client, &did, &jwt, "First post").await;
let latest_res = client
.get(format!(
"{}/xrpc/com.atproto.sync.getLatestCommit",
base_url().await
))
.query(&[("did", did.as_str())])
.send()
.await
.expect("Failed to get latest commit");
let body: serde_json::Value = latest_res.json().await.unwrap();
let rev = body["rev"].as_str().unwrap();
create_post(&client, &did, &jwt, "Second post").await;
let base = base_url().await;
let request_futures = (0..10).map(|_| {
let client = client.clone();
let did = did.clone();
let rev = rev.to_string();
async move {
client
.get(format!("{}/xrpc/com.atproto.sync.getRepo", base))
.query(&[("did", did.as_str()), ("since", rev.as_str())])
.send()
.await
.expect("Failed to send request")
.status()
}
});
let statuses: Vec<StatusCode> = futures::future::join_all(request_futures).await;
assert!(
statuses.iter().all(|s| *s == StatusCode::OK),
"getRepo with since should bypass semaphore, got: {:?}",
statuses
);
}
+28 -3
View File
@@ -157,11 +157,19 @@ impl ScopePermissions {
}
pub fn assert_rpc(&self, aud: &str, lxm: &str) -> Result<(), ScopeError> {
if self.has_transition_generic {
return Ok(());
if lxm.starts_with("chat.bsky.") {
if self.has_transition_chat {
return Ok(());
}
if self.has_transition_generic && !self.has_transition_chat {
return Err(ScopeError::InsufficientScope {
required: "transition:chat.bsky".to_string(),
message: format!("Chat access requires transition:chat.bsky scope to call {}", lxm),
});
}
}
if lxm.starts_with("chat.bsky.") && self.has_transition_chat {
if self.has_transition_generic {
return Ok(());
}
@@ -347,6 +355,23 @@ mod tests {
assert!(perms.allows_blob("image/png"));
}
#[test]
fn test_transition_generic_without_chat_blocks_chat() {
let perms = ScopePermissions::from_scope_string(Some("transition:generic"));
assert!(perms.allows_rpc("did:web:api.bsky.app", "app.bsky.feed.getTimeline"));
assert!(!perms.allows_rpc("did:web:api.bsky.app", "chat.bsky.convo.listConvos"));
assert!(!perms.allows_rpc("did:web:api.bsky.app", "chat.bsky.convo.getMessages"));
}
#[test]
fn test_transition_generic_with_chat_allows_chat() {
let perms =
ScopePermissions::from_scope_string(Some("transition:generic transition:chat.bsky"));
assert!(perms.allows_rpc("did:web:api.bsky.app", "app.bsky.feed.getTimeline"));
assert!(perms.allows_rpc("did:web:api.bsky.app", "chat.bsky.convo.listConvos"));
assert!(perms.allows_rpc("did:web:api.bsky.app", "chat.bsky.convo.getMessages"));
}
#[test]
fn test_transition_chat_only_allows_chat() {
let perms = ScopePermissions::from_scope_string(Some("transition:chat.bsky"));
+1
View File
@@ -52,6 +52,7 @@ k256 = { workspace = true }
rand = { workspace = true }
tikv-jemallocator = "0.6"
tracing-subscriber = { workspace = true, features = ["env-filter"] }
libc = "0.2"
[[bench]]
name = "blockstore"
+1 -3
View File
@@ -100,9 +100,7 @@ impl<'a> DirectSeeder<'a> {
}
let loc = self.data_writer.append_block(cid, data).unwrap();
self.hint_writer
.append_hint(cid, loc.file_id, loc.offset, loc.length)
.unwrap();
self.hint_writer.append_hint(cid, &loc).unwrap();
self.blocks_in_file += 1;
if self.blocks_in_file.is_multiple_of(10_000) {
@@ -93,8 +93,7 @@ pub(super) fn compact_on_writer_thread<S: StorageIO>(
Err(e)
}
Ok((new_size, live_count, dead_count)) => {
let positions = hint_positions.snapshot();
if let Err(e) = index.write_checkpoint(epoch.current(), &positions) {
if let Err(e) = index.write_checkpoint(epoch.current(), hint_positions) {
tracing::warn!(error = %e, "pre-delete checkpoint failed during compaction");
}
@@ -162,13 +161,7 @@ fn stream_compact<S: StorageIO>(
} => match index.get(&cid_bytes) {
Some(e) if e.location.file_id == source_file_id && !e.refcount.is_zero() => {
let loc = writer.append_block(&cid_bytes, &data)?;
hint_writer.append_relocate(
&cid_bytes,
loc.file_id,
loc.offset,
loc.length,
e.refcount.raw(),
)?;
hint_writer.append_relocate(&cid_bytes, &loc, e.refcount.raw())?;
relocations.push((cid_bytes, loc));
live_count = live_count.saturating_add(1);
}
@@ -188,13 +181,7 @@ fn stream_compact<S: StorageIO>(
}
false => {
let loc = writer.append_block(&cid_bytes, &data)?;
hint_writer.append_relocate(
&cid_bytes,
loc.file_id,
loc.offset,
loc.length,
e.refcount.raw(),
)?;
hint_writer.append_relocate(&cid_bytes, &loc, e.refcount.raw())?;
relocations.push((cid_bytes, loc));
live_count = live_count.saturating_add(1);
}
@@ -813,8 +813,7 @@ fn maybe_checkpoint(
if !elapsed && !threshold {
return;
}
let positions = hint_positions.snapshot();
match index.write_checkpoint(epoch.current(), &positions) {
match index.write_checkpoint(epoch.current(), hint_positions) {
Ok(()) => {
*last_checkpoint = std::time::Instant::now();
*writes_since_checkpoint = 0;
@@ -831,8 +830,7 @@ fn shutdown_checkpoint(
epoch: &EpochCounter,
hint_positions: &ShardHintPositions,
) {
let positions = hint_positions.snapshot();
match index.write_checkpoint(epoch.current(), &positions) {
match index.write_checkpoint(epoch.current(), hint_positions) {
Ok(()) => tracing::debug!("shutdown checkpoint written"),
Err(e) => tracing::warn!(error = %e, "shutdown checkpoint failed"),
}
@@ -924,8 +922,6 @@ fn commit_loop<S: StorageIO>(
if let Ok((ref dedup, _)) = result {
writes_since_checkpoint =
writes_since_checkpoint.saturating_add(dedup.len() as u64);
ctx.hint_positions
.update(ctx.shard_id, state.file_id, state.hint_position);
}
dispatch_responses(drain.entries, result.map(|(dedup, _proof)| dedup));
@@ -1024,8 +1020,6 @@ fn drain_and_process_remaining<S: StorageIO>(
if let Ok((ref _dedup, ref proof)) = result {
run_post_sync_hook(post_sync_hook, proof);
ctx.hint_positions
.update(ctx.shard_id, state.file_id, state.hint_position);
}
dispatch_responses(entries, result.map(|(dedup, _proof)| dedup));
@@ -1102,6 +1096,7 @@ fn process_batch<S: StorageIO>(
let location = match dedup.get(cid_bytes) {
Some(&loc) => {
dedup_hits = dedup_hits.saturating_add(1);
hint_writer.append_hint(cid_bytes, &loc)?;
loc
}
None => {
@@ -1136,7 +1131,7 @@ fn process_batch<S: StorageIO>(
}
let loc = data_writer.append_block(cid_bytes, data)?;
hint_writer.append_hint(cid_bytes, loc.file_id, loc.offset, loc.length)?;
hint_writer.append_hint(cid_bytes, &loc)?;
block_bytes = block_bytes.saturating_add(data.len() as u64);
block_count = block_count.saturating_add(1);
@@ -1194,7 +1189,19 @@ fn process_batch<S: StorageIO>(
};
let t = std::time::Instant::now();
index
.batch_put(&index_entries, &all_decrements, cursor, current_epoch, now)
.batch_put_and_advance_position(
&index_entries,
&all_decrements,
cursor,
current_epoch,
now,
super::hash_index::PositionUpdate {
hint_positions: &ctx.hint_positions,
shard_id: ctx.shard_id,
file_id: state.file_id,
offset: state.hint_position,
},
)
.map_err(CommitError::from)?;
let index_nanos = t.elapsed().as_nanos() as u64;
@@ -5,11 +5,19 @@ use std::path::{Path, PathBuf};
use parking_lot::RwLock;
use super::data_file::CID_SIZE;
use super::group_commit::ShardHintPositions;
use super::types::{
BlockLength, BlockLocation, BlockOffset, CidBytes, CollectionResult, CommitEpoch, DataFileId,
HintOffset, IndexEntry, LivenessInfo, RefCount, WallClockMs, WriteCursor,
HintOffset, IndexEntry, LivenessInfo, RefCount, ShardId, WallClockMs, WriteCursor,
};
pub struct PositionUpdate<'a> {
pub hint_positions: &'a ShardHintPositions,
pub shard_id: ShardId,
pub file_id: DataFileId,
pub offset: HintOffset,
}
const EMPTY_CID: [u8; CID_SIZE] = [0u8; CID_SIZE];
fn is_empty(cid: &[u8; CID_SIZE]) -> bool {
@@ -1197,6 +1205,30 @@ impl BlockIndex {
cursor: WriteCursor,
epoch: CommitEpoch,
now: WallClockMs,
) -> Result<(), BlockIndexError> {
self.batch_put_inner(entries, decrements, cursor, epoch, now, None)
}
pub fn batch_put_and_advance_position(
&self,
entries: &[([u8; CID_SIZE], BlockLocation)],
decrements: &[[u8; CID_SIZE]],
cursor: WriteCursor,
epoch: CommitEpoch,
now: WallClockMs,
position_update: PositionUpdate<'_>,
) -> Result<(), BlockIndexError> {
self.batch_put_inner(entries, decrements, cursor, epoch, now, Some(position_update))
}
fn batch_put_inner(
&self,
entries: &[([u8; CID_SIZE], BlockLocation)],
decrements: &[[u8; CID_SIZE]],
cursor: WriteCursor,
epoch: CommitEpoch,
now: WallClockMs,
position_update: Option<PositionUpdate<'_>>,
) -> Result<(), BlockIndexError> {
let mut table = self.table.write();
@@ -1217,6 +1249,12 @@ impl BlockIndex {
});
table.set_write_cursor(cursor);
if let Some(pos) = position_update {
pos.hint_positions
.update(pos.shard_id, pos.file_id, pos.offset);
}
Ok(())
}
@@ -1418,6 +1456,17 @@ impl BlockIndex {
}
pub fn write_checkpoint(
&self,
epoch: CommitEpoch,
hint_positions: &ShardHintPositions,
) -> io::Result<()> {
let _guard = self.checkpoint_lock.lock();
let table = self.table.read();
let positions = hint_positions.snapshot();
write_checkpoint_ab(&table, &self.index_dir, epoch, &positions)
}
pub fn write_checkpoint_with_positions(
&self,
epoch: CommitEpoch,
positions: &CheckpointPositions,
+66 -108
View File
@@ -55,22 +55,26 @@ fn write_hint_record<S: StorageIO>(
io.write_all_at(fd, write_offset.raw(), record)
}
fn encode_location_fields(record: &mut [u8; HINT_RECORD_SIZE], loc: &BlockLocation) {
record[FIELD_A_OFFSET..FIELD_A_OFFSET + 4].copy_from_slice(&loc.file_id.raw().to_le_bytes());
record[FIELD_A_OFFSET + 4..FIELD_A_OFFSET + 8]
.copy_from_slice(&loc.length.raw().to_le_bytes());
record[FIELD_B_OFFSET..FIELD_B_OFFSET + 8]
.copy_from_slice(&loc.offset.raw().to_le_bytes());
}
pub(crate) fn encode_hint_record<S: StorageIO>(
io: &S,
fd: FileId,
write_offset: HintOffset,
cid_bytes: &[u8; CID_SIZE],
file_id: DataFileId,
block_offset: BlockOffset,
length: BlockLength,
loc: &BlockLocation,
) -> io::Result<()> {
let mut record = [0u8; HINT_RECORD_SIZE];
record[TYPE_OFFSET] = RECORD_TYPE_PUT;
record[VERSION_OFFSET] = HINT_FORMAT_VERSION;
record[CID_OFFSET..CID_OFFSET + CID_SIZE].copy_from_slice(cid_bytes);
record[FIELD_A_OFFSET..FIELD_A_OFFSET + 4].copy_from_slice(&file_id.raw().to_le_bytes());
record[FIELD_A_OFFSET + 4..FIELD_A_OFFSET + 8].copy_from_slice(&length.raw().to_le_bytes());
record[FIELD_B_OFFSET..FIELD_B_OFFSET + 8].copy_from_slice(&block_offset.raw().to_le_bytes());
encode_location_fields(&mut record, loc);
let checksum = hint_checksum(&record[..HINT_PAYLOAD_SIZE]);
record[CHECKSUM_OFFSET..].copy_from_slice(&checksum.to_le_bytes());
@@ -85,9 +89,7 @@ pub(crate) fn encode_relocate_record<S: StorageIO>(
fd: FileId,
write_offset: HintOffset,
cid_bytes: &[u8; CID_SIZE],
file_id: DataFileId,
block_offset: BlockOffset,
length: BlockLength,
loc: &BlockLocation,
refcount: u32,
) -> io::Result<()> {
let mut record = [0u8; HINT_RECORD_SIZE];
@@ -96,9 +98,7 @@ pub(crate) fn encode_relocate_record<S: StorageIO>(
let rc16 = u16::try_from(refcount).unwrap_or(u16::MAX);
record[REFCOUNT_OFFSET..REFCOUNT_OFFSET + 2].copy_from_slice(&rc16.to_le_bytes());
record[CID_OFFSET..CID_OFFSET + CID_SIZE].copy_from_slice(cid_bytes);
record[FIELD_A_OFFSET..FIELD_A_OFFSET + 4].copy_from_slice(&file_id.raw().to_le_bytes());
record[FIELD_A_OFFSET + 4..FIELD_A_OFFSET + 8].copy_from_slice(&length.raw().to_le_bytes());
record[FIELD_B_OFFSET..FIELD_B_OFFSET + 8].copy_from_slice(&block_offset.raw().to_le_bytes());
encode_location_fields(&mut record, loc);
let checksum = hint_checksum(&record[..HINT_PAYLOAD_SIZE]);
record[CHECKSUM_OFFSET..].copy_from_slice(&checksum.to_le_bytes());
@@ -323,19 +323,9 @@ impl<'a, S: StorageIO> HintFileWriter<'a, S> {
pub fn append_hint(
&mut self,
cid_bytes: &[u8; CID_SIZE],
file_id: DataFileId,
offset: BlockOffset,
length: BlockLength,
loc: &BlockLocation,
) -> io::Result<()> {
encode_hint_record(
self.io,
self.fd,
self.position,
cid_bytes,
file_id,
offset,
length,
)?;
encode_hint_record(self.io, self.fd, self.position, cid_bytes, loc)?;
self.position = self.position.advance(HINT_RECORD_SIZE as u64);
Ok(())
}
@@ -354,21 +344,10 @@ impl<'a, S: StorageIO> HintFileWriter<'a, S> {
pub fn append_relocate(
&mut self,
cid_bytes: &[u8; CID_SIZE],
file_id: DataFileId,
offset: BlockOffset,
length: BlockLength,
loc: &BlockLocation,
refcount: u32,
) -> io::Result<()> {
encode_relocate_record(
self.io,
self.fd,
self.position,
cid_bytes,
file_id,
offset,
length,
refcount,
)?;
encode_relocate_record(self.io, self.fd, self.position, cid_bytes, loc, refcount)?;
self.position = self.position.advance(HINT_RECORD_SIZE as u64);
Ok(())
}
@@ -862,7 +841,8 @@ mod tests {
let offset = BlockOffset::new(1024);
let length = BlockLength::new(256);
encode_hint_record(&sim, fd, HintOffset::new(0), &cid, file_id, offset, length).unwrap();
let loc = BlockLocation { file_id, offset, length };
encode_hint_record(&sim, fd, HintOffset::new(0), &cid, &loc).unwrap();
let file_size = sim.file_size(fd).unwrap();
let record = decode_hint_record(&sim, fd, HintOffset::new(0), file_size)
@@ -920,16 +900,12 @@ mod tests {
(0u8..5).for_each(|i| {
let cid = test_cid(i);
let write_offset = HintOffset::new(i as u64 * HINT_RECORD_SIZE as u64);
encode_hint_record(
&sim,
fd,
write_offset,
&cid,
DataFileId::new(i as u32),
BlockOffset::new(i as u64 * 100),
BlockLength::new(50 + i as u32),
)
.unwrap();
let loc = BlockLocation {
file_id: DataFileId::new(i as u32),
offset: BlockOffset::new(i as u64 * 100),
length: BlockLength::new(50 + i as u32),
};
encode_hint_record(&sim, fd, write_offset, &cid, &loc).unwrap();
});
let file_size = sim.file_size(fd).unwrap();
@@ -971,16 +947,12 @@ mod tests {
fn detects_corrupted_hint() {
let (sim, fd) = setup();
let cid = test_cid(1);
encode_hint_record(
&sim,
fd,
HintOffset::new(0),
&cid,
DataFileId::new(0),
BlockOffset::new(0),
BlockLength::new(100),
)
.unwrap();
let loc = BlockLocation {
file_id: DataFileId::new(0),
offset: BlockOffset::new(0),
length: BlockLength::new(100),
};
encode_hint_record(&sim, fd, HintOffset::new(0), &cid, &loc).unwrap();
sim.write_all_at(fd, 10, &[0xFF]).unwrap();
@@ -1006,16 +978,12 @@ mod tests {
fn oversized_length_treated_as_corrupted() {
let (sim, fd) = setup();
let cid = test_cid(1);
encode_hint_record(
&sim,
fd,
HintOffset::new(0),
&cid,
DataFileId::new(0),
BlockOffset::new(0),
BlockLength::new(100),
)
.unwrap();
let loc = BlockLocation {
file_id: DataFileId::new(0),
offset: BlockOffset::new(0),
length: BlockLength::new(100),
};
encode_hint_record(&sim, fd, HintOffset::new(0), &cid, &loc).unwrap();
let length_offset = FIELD_A_OFFSET as u64 + 4;
let oversized = (MAX_BLOCK_SIZE + 1).to_le_bytes();
@@ -1040,14 +1008,12 @@ mod tests {
let mut writer = HintFileWriter::new(&sim, fd);
(0u8..5).for_each(|i| {
writer
.append_hint(
&test_cid(i),
DataFileId::new(0),
BlockOffset::new(i as u64 * 100),
BlockLength::new(50 + i as u32),
)
.unwrap();
let loc = BlockLocation {
file_id: DataFileId::new(0),
offset: BlockOffset::new(i as u64 * 100),
length: BlockLength::new(50 + i as u32),
};
writer.append_hint(&test_cid(i), &loc).unwrap();
});
assert_eq!(
@@ -1074,25 +1040,21 @@ mod tests {
fn hint_writer_resume_continues_at_position() {
let (sim, fd) = setup();
let mut writer = HintFileWriter::new(&sim, fd);
writer
.append_hint(
&test_cid(0),
DataFileId::new(0),
BlockOffset::new(0),
BlockLength::new(100),
)
.unwrap();
let loc0 = BlockLocation {
file_id: DataFileId::new(0),
offset: BlockOffset::new(0),
length: BlockLength::new(100),
};
writer.append_hint(&test_cid(0), &loc0).unwrap();
let pos = writer.position();
let mut writer2 = HintFileWriter::resume(&sim, fd, pos);
writer2
.append_hint(
&test_cid(1),
DataFileId::new(0),
BlockOffset::new(100),
BlockLength::new(200),
)
.unwrap();
let loc1 = BlockLocation {
file_id: DataFileId::new(0),
offset: BlockOffset::new(100),
length: BlockLength::new(200),
};
writer2.append_hint(&test_cid(1), &loc1).unwrap();
let reader = HintFileReader::open(&sim, fd).unwrap();
let valid_count = reader
@@ -1115,14 +1077,12 @@ mod tests {
fn hint_reader_stops_on_truncated() {
let (sim, fd) = setup();
let mut writer = HintFileWriter::new(&sim, fd);
writer
.append_hint(
&test_cid(0),
DataFileId::new(0),
BlockOffset::new(0),
BlockLength::new(100),
)
.unwrap();
let loc = BlockLocation {
file_id: DataFileId::new(0),
offset: BlockOffset::new(0),
length: BlockLength::new(100),
};
writer.append_hint(&test_cid(0), &loc).unwrap();
sim.write_all_at(fd, writer.position().raw(), &[0u8; HINT_RECORD_SIZE - 1])
.unwrap();
@@ -1140,14 +1100,12 @@ mod tests {
let mut writer = HintFileWriter::new(&sim, fd);
(0u8..3).for_each(|i| {
writer
.append_hint(
&test_cid(i),
DataFileId::new(0),
BlockOffset::new(i as u64 * 100),
BlockLength::new(50),
)
.unwrap();
let loc = BlockLocation {
file_id: DataFileId::new(0),
offset: BlockOffset::new(i as u64 * 100),
length: BlockLength::new(50),
};
writer.append_hint(&test_cid(i), &loc).unwrap();
});
sim.write_all_at(fd, HINT_RECORD_SIZE as u64 + 5, &[0xFF])
@@ -0,0 +1,223 @@
mod common;
use std::io;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use tranquil_store::blockstore::{
BlockStoreConfig, BlocksSynced, CidBytes, GroupCommitConfig, TranquilBlockStore,
};
use tranquil_store::PostBlockstoreHook;
struct SlowHook;
impl PostBlockstoreHook for SlowHook {
fn on_blocks_synced(&self, _proof: &BlocksSynced) -> io::Result<()> {
std::thread::sleep(std::time::Duration::from_millis(1));
Ok(())
}
}
fn refcount(store: &TranquilBlockStore, cid: &CidBytes) -> Option<u32> {
store.block_index().get(cid).map(|e| e.refcount.raw())
}
fn race_config(dir: &std::path::Path) -> BlockStoreConfig {
BlockStoreConfig {
data_dir: dir.join("data"),
index_dir: dir.join("index"),
max_file_size: 256 * 1024,
group_commit: GroupCommitConfig {
checkpoint_interval_ms: 10,
checkpoint_write_threshold: 20,
..GroupCommitConfig::default()
},
shard_count: 4,
}
}
fn cid_for(shard: u8, seq: u32) -> CidBytes {
let mut cid = [0u8; 36];
cid[0] = 0x01;
cid[1] = 0x71;
cid[2] = 0x12;
cid[3] = 0x20;
cid[4] = shard;
cid[8..12].copy_from_slice(&seq.to_le_bytes());
(12..36).for_each(|i| cid[i] = (seq as u8).wrapping_add(i as u8));
cid
}
fn write_phase(base: &std::path::Path, use_hook: bool) -> Vec<CidBytes> {
let config = race_config(base);
let hook: Option<Arc<dyn PostBlockstoreHook>> = use_hook.then(|| Arc::new(SlowHook) as _);
let store = Arc::new(TranquilBlockStore::open_with_hook(config, hook).unwrap());
let running = Arc::new(AtomicBool::new(true));
let total_cycles = Arc::new(AtomicU64::new(0));
let writers: Vec<_> = (0..4u8)
.map(|shard| {
let store = Arc::clone(&store);
let running = Arc::clone(&running);
let total_cycles = Arc::clone(&total_cycles);
std::thread::spawn(move || {
let mut targets = Vec::new();
let mut seq = 0u32;
while running.load(Ordering::Relaxed) {
let cid = cid_for(shard, seq);
store
.put_blocks_blocking(vec![(cid, vec![shard; 60])])
.unwrap();
store
.put_blocks_blocking(vec![(cid, vec![shard; 60])])
.unwrap();
store
.apply_commit_blocking(vec![], vec![cid])
.unwrap();
targets.push(cid);
seq += 1;
total_cycles.fetch_add(1, Ordering::Relaxed);
}
targets
})
})
.collect();
while total_cycles.load(Ordering::Relaxed) < 500 {
std::thread::yield_now();
}
running.store(false, Ordering::Relaxed);
let all_targets: Vec<CidBytes> = writers
.into_iter()
.flat_map(|w| w.join().unwrap())
.collect();
all_targets.iter().for_each(|cid| {
assert_eq!(refcount(&store, cid), Some(1), "pre-crash sanity");
});
let store = Arc::try_unwrap(store).ok().unwrap();
std::mem::forget(store);
all_targets
}
fn verify_phase(base: &std::path::Path, targets: &[CidBytes]) -> usize {
let config = race_config(base);
let store = TranquilBlockStore::open(config).unwrap();
let bad = targets
.iter()
.filter(|cid| refcount(&store, cid) != Some(1))
.count();
drop(store);
bad
}
#[test]
fn crash_recovery_preserves_refcounts() {
common::with_runtime(|| {
let mut corrupted = 0u32;
let total = 20u32;
(0..total).for_each(|_| {
let dir = tempfile::TempDir::new().unwrap();
let exe = std::env::current_exe().unwrap();
let dir_str = dir.path().to_str().unwrap();
let output = std::process::Command::new(&exe)
.arg("--exact")
.arg("__crash_write_phase")
.env("CRASH_TEST_DIR", dir_str)
.env("CRASH_TEST_HOOK", "0")
.output()
.unwrap();
assert!(output.status.success() || output.status.code() == Some(0));
let target_bytes = std::fs::read(dir.path().join("targets.bin")).unwrap();
let targets: Vec<CidBytes> = target_bytes
.chunks_exact(36)
.map(|chunk| {
let mut cid = [0u8; 36];
cid.copy_from_slice(chunk);
cid
})
.collect();
if verify_phase(dir.path(), &targets) > 0 {
corrupted += 1;
}
});
assert_eq!(
corrupted, 0,
"{corrupted}/{total} iterations had refcount corruption after crash recovery"
);
});
}
#[test]
fn crash_with_slow_hook_preserves_refcounts() {
common::with_runtime(|| {
let mut corrupted = 0u32;
let total = 20u32;
(0..total).for_each(|_| {
let dir = tempfile::TempDir::new().unwrap();
let exe = std::env::current_exe().unwrap();
let dir_str = dir.path().to_str().unwrap();
let output = std::process::Command::new(&exe)
.arg("--exact")
.arg("__crash_write_phase")
.env("CRASH_TEST_DIR", dir_str)
.env("CRASH_TEST_HOOK", "1")
.output()
.unwrap();
assert!(output.status.success() || output.status.code() == Some(0));
let target_bytes = std::fs::read(dir.path().join("targets.bin")).unwrap();
let targets: Vec<CidBytes> = target_bytes
.chunks_exact(36)
.map(|chunk| {
let mut cid = [0u8; 36];
cid.copy_from_slice(chunk);
cid
})
.collect();
if verify_phase(dir.path(), &targets) > 0 {
corrupted += 1;
}
});
assert_eq!(
corrupted, 0,
"{corrupted}/{total} iterations had refcount corruption after crash with slow hook"
);
});
}
#[test]
fn __crash_write_phase() {
let dir = match std::env::var("CRASH_TEST_DIR") {
Ok(d) => d,
Err(_) => return,
};
let use_hook = std::env::var("CRASH_TEST_HOOK").map(|v| v == "1").unwrap_or(false);
let base = std::path::Path::new(&dir);
let rt = tokio::runtime::Runtime::new().unwrap();
let _guard = rt.enter();
let targets = write_phase(base, use_hook);
let target_bytes: Vec<u8> = targets.iter().flat_map(|cid| cid.iter().copied()).collect();
std::fs::write(base.join("targets.bin"), &target_bytes).unwrap();
unsafe { libc::_exit(0) }
}
@@ -75,7 +75,7 @@ impl SimHarness {
let data = vec![seed as u8; data_size];
let loc = writer.append_block(&cid, &data).unwrap();
hint_writer
.append_hint(&cid, loc.file_id, loc.offset, loc.length)
.append_hint(&cid, &loc)
.unwrap();
(cid, loc)
})
@@ -112,7 +112,7 @@ impl SimHarness {
HintOffset::new(entries.len() as u64 * HINT_RECORD_SIZE as u64),
);
index
.write_checkpoint(CommitEpoch::zero(), &positions)
.write_checkpoint_with_positions(CommitEpoch::zero(), &positions)
.unwrap();
}
@@ -539,7 +539,7 @@ fn sim_aggressive_faults_data_integrity() {
let data = vec![i as u8; 64];
let loc = writer.append_block(&cid, &data).ok()?;
hint_writer
.append_hint(&cid, loc.file_id, loc.offset, loc.length)
.append_hint(&cid, &loc)
.ok()?;
Some(())
})?;
+28 -12
View File
@@ -142,6 +142,17 @@ pub async fn get_repo(
return get_repo_since(&state, &did, &head_cid, since).await;
}
let _permit = match state.repo_export_semaphore.try_acquire() {
Ok(permit) => permit,
Err(_) => {
return (
StatusCode::SERVICE_UNAVAILABLE,
"Too many concurrent repo exports",
)
.into_response();
}
};
let car_bytes = match generate_repo_car_from_user_blocks(
state.repos.repo.as_ref(),
&state.block_store,
@@ -213,19 +224,24 @@ async fn get_repo_since(state: &AppState, did: &Did, head_cid: &Cid, since: &str
.into_response();
}
let blocks = match state.block_store.get_many(&block_cids).await {
Ok(b) => b,
Err(e) => {
error!("Block store error in get_repo_since: {:?}", e);
return ApiError::InternalError(Some("Failed to get blocks".into())).into_response();
}
};
for chunk_start in (0..block_cids.len()).step_by(500) {
let chunk_end = (chunk_start + 500).min(block_cids.len());
let chunk = &block_cids[chunk_start..chunk_end];
let blocks = match state.block_store.get_many(chunk).await {
Ok(b) => b,
Err(e) => {
error!("Block store error in get_repo_since: {:?}", e);
return ApiError::InternalError(Some("Failed to get blocks".into()))
.into_response();
}
};
blocks
.into_iter()
.enumerate()
.filter_map(|(i, block_opt)| block_opt.map(|block| (block_cids[i], block)))
.for_each(|(cid, block)| car_bytes.extend_from_slice(&encode_car_block(&cid, &block)));
chunk
.iter()
.zip(blocks.into_iter())
.filter_map(|(cid, block_opt)| block_opt.map(|block| (*cid, block)))
.for_each(|(cid, block)| car_bytes.extend_from_slice(&encode_car_block(&cid, &block)));
}
(
StatusCode::OK,
+7 -6
View File
@@ -49,12 +49,12 @@ mkdir -p /var/lib/tranquil/blobs
We'll set ownership after creating the service user.
## Install deno (for frontend build)
## Install Node.js and pnpm (for frontend build)
```bash
curl -fsSL https://deno.land/install.sh | sh
export PATH="$HOME/.deno/bin:$PATH"
echo 'export PATH="$HOME/.deno/bin:$PATH"' >> ~/.bashrc
curl -fsSL https://deb.nodesource.com/setup_24.x | bash -
apt install -y nodejs
npm install -g pnpm
```
## Clone and build Tranquil PDS
@@ -64,7 +64,8 @@ cd /opt
git clone https://tangled.org/tranquil.farm/tranquil-pds tranquil-pds
cd tranquil-pds
cd frontend
deno task build
pnpm install --frozen-lockfile
pnpm build
cd ..
cargo build --release
```
@@ -330,7 +331,7 @@ Update Tranquil PDS:
```bash
cd /opt/tranquil-pds
git pull
cd frontend && deno task build && cd ..
cd frontend && pnpm install --frozen-lockfile && pnpm build && cd ..
cargo build --release
systemctl stop tranquil-pds
cp target/release/tranquil-pds /usr/local/bin/
Generated
+4 -21
View File
@@ -2,11 +2,11 @@
"nodes": {
"nixpkgs": {
"locked": {
"lastModified": 1766314097,
"narHash": "sha256-laJftWbghBehazn/zxVJ8NdENVgjccsWAdAqKXhErrM=",
"lastModified": 1775888245,
"narHash": "sha256-nwASzrRDD1JBEu/o8ekKYEXm/oJW6EMCzCRdrwcLe90=",
"owner": "nixos",
"repo": "nixpkgs",
"rev": "306ea70f9eb0fb4e040f8540e2deab32ed7e2055",
"rev": "13043924aaa7375ce482ebe2494338e058282925",
"type": "github"
},
"original": {
@@ -16,26 +16,9 @@
"type": "github"
}
},
"nixpkgs-fetch-deno": {
"locked": {
"lastModified": 1766410835,
"narHash": "sha256-dRhVt0aFDyTqppyzRLxiO1JZEAoIA2fUnaeyJTe+UwU=",
"owner": "aMOPel",
"repo": "nixpkgs",
"rev": "c9801acc8c4fac6377d076bc1c102b15bd9cfa6f",
"type": "github"
},
"original": {
"owner": "aMOPel",
"ref": "feat/fetchDenoDeps",
"repo": "nixpkgs",
"type": "github"
}
},
"root": {
"inputs": {
"nixpkgs": "nixpkgs",
"nixpkgs-fetch-deno": "nixpkgs-fetch-deno"
"nixpkgs": "nixpkgs"
}
}
},
+42 -40
View File
@@ -1,51 +1,53 @@
{
inputs = {
nixpkgs.url = "github:nixos/nixpkgs/nixpkgs-unstable";
# tranquil frontend uses deno as its package manager and build time runtime.
# nixpkgs does not have deno support yet but its being worked on in https://github.com/NixOS/nixpkgs/pull/419255
# for now we important that PR as well purely for its fetchDenoDeps
nixpkgs-fetch-deno.url = "github:aMOPel/nixpkgs/feat/fetchDenoDeps";
};
outputs = {
self,
nixpkgs,
...
} @ inputs: let
forAllSystems = function:
nixpkgs.lib.genAttrs nixpkgs.lib.systems.flakeExposed (
system: (function system nixpkgs.legacyPackages.${system})
);
in {
packages = forAllSystems (system: pkgs: {
tranquil-pds = pkgs.callPackage ./default.nix {};
tranquil-frontend = pkgs.callPackage ./frontend.nix {
inherit (inputs.nixpkgs-fetch-deno.legacyPackages.${system}) fetchDenoDeps;
outputs =
{
self,
nixpkgs,
}:
let
forAllSystems =
function:
nixpkgs.lib.genAttrs nixpkgs.lib.systems.flakeExposed (
system: function nixpkgs.legacyPackages.${system}
);
in
{
packages = forAllSystems (pkgs: {
tranquil-pds = pkgs.callPackage ./default.nix { };
tranquil-frontend = pkgs.callPackage ./frontend.nix { };
default = self.packages.${pkgs.stdenv.hostPlatform.system}.tranquil-pds;
});
devShells = forAllSystems (pkgs: {
default = pkgs.callPackage ./shell.nix { };
});
nixosModules = {
default = self.nixosModules.tranquil-pds;
tranquil-pds =
{ lib, pkgs, ... }:
{
_file = "${self.outPath}/flake.nix#nixosModules.tranquil-pds";
imports = [ ./module.nix ];
config.services.tranquil-pds = {
package = self.packages.${pkgs.stdenv.hostPlatform.system}.tranquil-pds;
settings.frontend.package = self.packages.${pkgs.stdenv.hostPlatform.system}.tranquil-frontend;
};
};
};
default = self.packages.${pkgs.stdenv.hostPlatform.system}.tranquil-pds;
});
devShells = forAllSystems (system: pkgs: {
default = pkgs.callPackage ./shell.nix {};
});
checks.x86_64-linux.integration = import ./test.nix {
pkgs = nixpkgs.legacyPackages.x86_64-linux;
inherit self;
};
nixosModules = {
default = self.nixosModules.tranquil-pds;
tranquil-pds = {
_file = "${self.outPath}/flake.nix#nixosModules.tranquil-pds";
imports = [(import ./module.nix self)];
checks.aarch64-linux.integration = import ./test.nix {
pkgs = nixpkgs.legacyPackages.aarch64-linux;
inherit self;
};
};
checks.x86_64-linux.integration = import ./test.nix {
pkgs = nixpkgs.legacyPackages.x86_64-linux;
inherit self;
};
checks.aarch64-linux.integration = import ./test.nix {
pkgs = nixpkgs.legacyPackages.aarch64-linux;
inherit self;
};
};
}
+43 -60
View File
@@ -1,66 +1,49 @@
{
lib,
stdenvNoCC,
fetchDenoDeps,
fetchFromGitHub,
buildGoModule,
deno,
esbuild,
}: let
nodejs,
pnpm,
pnpmConfigHook,
fetchPnpmDeps,
nix-update-script,
}:
let
toml = (lib.importTOML ./Cargo.toml).workspace.package;
deno-deps = fetchDenoDeps {
pname = "tranquil-frontend-deno-deps";
denoLock = ./frontend/deno.lock;
hash = "sha256-UB+E00TjWX0fTUZ7XwcwRJ/OUOSSJpz6Ss04U5i8dGI=";
};
# the esbuild in upstream nixpkgs is too old.
esbuild' = esbuild.override {
buildGoModule = args: buildGoModule (
args // (
let
version = "0.27.2";
in {
inherit version;
src = fetchFromGitHub {
owner = "evanw";
repo = "esbuild";
tag = "v${version}";
hash = "sha256-JbJB3F1NQlmA5d0rdsLm4RVD24OPdV4QXpxW8VWbESA";
};
vendorHash = "sha256-+BfxCyg0KkDQpHt/wycy/8CTG6YBA/VJvJFhhzUnSiQ";
}
)
);
};
in stdenvNoCC.mkDerivation {
pname = "tranquil-frontend";
inherit (toml) version;
src = ./frontend;
in
stdenvNoCC.mkDerivation (finalAttrs: {
pname = "tranquil-frontend";
inherit (toml) version;
nativeBuildInputs = [
deno
src = ./frontend;
pnpmDeps = fetchPnpmDeps {
inherit (finalAttrs) pname version src;
fetcherVersion = 3;
hash = "sha256-E0S8dOaTOpY9m7Ft59tUQ6CLlLriWPE4WE1+S45vomY=";
};
nativeBuildInputs = [
pnpm
nodejs
pnpmConfigHook
];
buildPhase = ''
runHook preBuild
pnpm build
runHook postBuild
'';
installPhase = ''
runHook preInstall
cp -r ./dist $out
runHook postInstall
'';
passthru.updateScript = nix-update-script {
extraArgs = [
"--version"
"SKIP"
];
# tell vite (through the esbuild api) where the nix provided esbuild binary is
env.ESBUILD_BINARY_PATH = lib.getExe esbuild';
buildPhase = ''
# copy the deps to the required location
cp -r --no-preserve=mode ${deno-deps.denoDeps}/.deno ./
cp -r --no-preserve=mode ${deno-deps.denoDeps}/vendor ./
pwd
ls /build/frontend/vendor
# Now you can run the project using deps
# you need to activate [deno's vendor feature](https://docs.deno.com/runtime/fundamentals/modules/#vendoring-remote-modules)
# you need to use the `$DENO_DIR` env var, to point deno to the correct local cache
DENO_DIR=./.deno deno run --frozen --cached-only build
'';
installPhase = ''
cp -r ./dist $out
'';
}
};
})
+5 -2
View File
@@ -1,7 +1,10 @@
FROM denoland/deno:alpine AS builder
FROM node:24-alpine AS builder
RUN corepack enable && corepack prepare pnpm@latest --activate
WORKDIR /app
COPY package.json pnpm-lock.yaml ./
RUN pnpm install --frozen-lockfile
COPY . ./
RUN deno task build
RUN pnpm build
FROM nginx:1.29-alpine
COPY --from=builder /app/dist /usr/share/nginx/html
-22
View File
@@ -1,22 +0,0 @@
{
"tasks": {
"dev": "deno run -A npm:vite",
"build": "deno run -A npm:vite build",
"preview": "deno run -A npm:vite preview",
"check": "deno run -A npm:svelte-check --tsconfig ./tsconfig.json",
"test": "deno run -A npm:vitest",
"test:run": "deno run -A npm:vitest run",
"test:watch": "deno run -A npm:vitest watch",
"test:ui": "deno run -A npm:vitest --ui",
"test:coverage": "deno run -A npm:vitest run --coverage"
},
"nodeModulesDir": "auto",
"lint": {
"rules": {
"exclude": [
"require-await",
"prefer-const"
]
}
}
}
-1930
View File
File diff suppressed because it is too large Load Diff
+2475
View File
File diff suppressed because it is too large Load Diff
-1
View File
@@ -56,7 +56,6 @@ test-misc:
./scripts/run-tests.sh --test actor --test commit_signing --test image_processing --test lifecycle_social --test notifications --test server --test signing_key --test verify_live_commit
test *args:
@just test-store
@just test-unit
./scripts/run-tests.sh {{args}}
+16 -11
View File
@@ -1,15 +1,17 @@
self: {
{
lib,
pkgs,
config,
...
}: let
}:
let
cfg = config.services.tranquil-pds;
inherit (lib) types mkOption;
settingsFormat = pkgs.formats.toml { };
in {
in
{
_class = "nixos";
options.services.tranquil-pds = {
@@ -17,8 +19,8 @@ in {
package = mkOption {
type = types.package;
default = self.packages.${pkgs.stdenv.hostPlatform.system}.tranquil-pds;
defaultText = lib.literalExpression "self.packages.\${pkgs.stdenv.hostPlatform.system}.tranquil-pds";
default = pkgs.callPackage ./default.nix { };
defaultText = lib.literalExpression "pkgs.tranquil-pds";
description = "The tranquil-pds package to use";
};
@@ -97,13 +99,16 @@ in {
};
frontend = {
enabled = lib.mkEnableOption "serving the frontend from the backend. Disable to serve the frontend manually"
// { default = true; };
enabled =
lib.mkEnableOption "serving the frontend from the backend. Disable to serve the frontend manually"
// {
default = true;
};
dir = mkOption {
type = types.nullOr types.package;
default = self.packages.${pkgs.stdenv.hostPlatform.system}.tranquil-frontend;
defaultText = lib.literalExpression "self.packages.\${pkgs.stdenv.hostPlatform.system}.tranquil-frontend";
default = pkgs.callPackage ./frontend.nix { };
defaultText = lib.literalExpression "pkgs.tranquil-frontend";
description = "Frontend package to be served by the backend";
};
};
@@ -137,7 +142,7 @@ in {
};
config = lib.mkIf cfg.enable (
lib.mkMerge [
lib.mkMerge [
(lib.mkIf cfg.database.createLocally {
services.postgresql = {
enable = true;
@@ -159,7 +164,7 @@ in {
};
})
{
{
users.users.${cfg.user} = {
isSystemUser = true;
inherit (cfg) group;
+10 -6
View File
@@ -207,11 +207,15 @@ if ! command -v rustc &>/dev/null; then
source "$HOME/.cargo/env"
fi
log_info "Installing deno..."
export PATH="$HOME/.deno/bin:$PATH"
if ! command -v deno &>/dev/null && [[ ! -f "$HOME/.deno/bin/deno" ]]; then
curl -fsSL https://deno.land/install.sh | sh
grep -q 'deno/bin' ~/.bashrc 2>/dev/null || echo 'export PATH="$HOME/.deno/bin:$PATH"' >> ~/.bashrc
log_info "Installing Node.js..."
if ! command -v node &>/dev/null; then
curl -fsSL https://deb.nodesource.com/setup_24.x | bash -
apt install -y nodejs
fi
log_info "Installing pnpm..."
if ! command -v pnpm &>/dev/null; then
npm install -g pnpm
fi
log_info "Cloning Tranquil PDS..."
@@ -223,7 +227,7 @@ fi
cd /opt/tranquil-pds
log_info "Building frontend..."
"$HOME/.deno/bin/deno" task build --filter=frontend
cd frontend && pnpm install --frozen-lockfile && pnpm build && cd ..
log_success "Frontend built"
log_info "Building Tranquil PDS (this takes a while)..."
+1 -1
View File
@@ -23,7 +23,7 @@ cargo test --no-run 2>&1 | tail -1
echo "Running tests..."
echo ""
cargo nextest run "$@"
cargo nextest run -E 'not package(tranquil-store)' "$@"
echo ""
echo "All tests passed."
+11 -9
View File
@@ -1,5 +1,4 @@
{
lib,
mkShell,
callPackage,
rustPlatform,
@@ -17,13 +16,18 @@
cargo-nextest,
# frontend tooling
deno,
svelte-language-server,
typescript-language-server,
}: let
defaultPackage = callPackage ./default.nix { };
in mkShell {
inputsFrom = [ defaultPackage ];
}:
let
pds = callPackage ./default.nix { };
frontend = callPackage ./frontend.nix { };
in
mkShell {
inputsFrom = [
pds
frontend
];
env = {
RUST_SRC_PATH = rustPlatform.rustLibSrc;
@@ -39,10 +43,8 @@ in mkShell {
rust-analyzer
sqlx-cli
cargo-nextest
deno
svelte-language-server
typescript-language-server
];
}