mirror of
https://tangled.org/tranquil.farm/tranquil-pds
synced 2026-08-25 10:46:11 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1dfbd27cce | ||
|
|
bc5e0e0446 | ||
|
|
1c87ef5536 | ||
|
|
255c7135f9 |
Generated
+22
-22
@@ -7405,7 +7405,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-api"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum",
|
||||
@@ -7456,7 +7456,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-auth"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"base32",
|
||||
@@ -7479,7 +7479,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-cache"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"base64 0.22.1",
|
||||
@@ -7493,7 +7493,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-comms"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"base64 0.22.1",
|
||||
@@ -7511,7 +7511,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-config"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"confique",
|
||||
"serde",
|
||||
@@ -7519,7 +7519,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-crypto"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"base64 0.22.1",
|
||||
@@ -7535,7 +7535,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-db"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"chrono",
|
||||
@@ -7552,7 +7552,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-db-traits"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"base64 0.22.1",
|
||||
@@ -7568,7 +7568,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-infra"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"bytes",
|
||||
@@ -7579,7 +7579,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-lexicon"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"hickory-resolver",
|
||||
@@ -7597,7 +7597,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-oauth"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum",
|
||||
@@ -7620,7 +7620,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-oauth-server"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"axum",
|
||||
"base64 0.22.1",
|
||||
@@ -7653,7 +7653,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-pds"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"anyhow",
|
||||
@@ -7745,7 +7745,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-repo"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"cid",
|
||||
@@ -7757,7 +7757,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-ripple"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"backon",
|
||||
@@ -7782,7 +7782,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-scopes"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"axum",
|
||||
"futures",
|
||||
@@ -7798,7 +7798,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-server"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"axum",
|
||||
"clap",
|
||||
@@ -7819,7 +7819,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-signal"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"chrono",
|
||||
@@ -7842,7 +7842,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-storage"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"aws-config",
|
||||
@@ -7859,7 +7859,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-store"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"bytes",
|
||||
@@ -7905,7 +7905,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-sync"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum",
|
||||
@@ -7927,7 +7927,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "tranquil-types"
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"cid",
|
||||
|
||||
+1
-1
@@ -26,7 +26,7 @@ members = [
|
||||
]
|
||||
|
||||
[workspace.package]
|
||||
version = "0.5.1"
|
||||
version = "0.5.2"
|
||||
edition = "2024"
|
||||
license = "AGPL-3.0-or-later"
|
||||
|
||||
|
||||
+5
-2
@@ -1,7 +1,10 @@
|
||||
FROM denoland/deno:alpine AS frontend
|
||||
FROM node:24-alpine AS builder
|
||||
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
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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> {
|
||||
|
||||
@@ -266,35 +266,48 @@ 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()))
|
||||
})
|
||||
},
|
||||
|
||||
let (obsolete_result, new_tree_result) = tokio::join!(
|
||||
compute_obsolete_cids(
|
||||
&original_settled,
|
||||
&new_settled,
|
||||
CommitCid::from(ctx.current_root_cid),
|
||||
),
|
||||
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>(
|
||||
tokio::try_join!(new_settled.collect_node_cids(), new_settled.leaves(),)?;
|
||||
Ok::<Vec<Cid>, jacquard_repo::error::RepoError>(
|
||||
nodes
|
||||
.into_iter()
|
||||
.chain(leaves.iter().map(|(_, cid)| *cid))
|
||||
.collect(),
|
||||
)
|
||||
},
|
||||
)?;
|
||||
);
|
||||
|
||||
let new_tree_cids = match new_tree_result {
|
||||
Ok(cids) => cids,
|
||||
Err(e) => {
|
||||
error!(
|
||||
"new tree walk failed: {e}. \
|
||||
Falling back to written-block CIDs only; \
|
||||
shared subtree ownership already tracked by prior commits."
|
||||
);
|
||||
block_bytes.keys().copied().collect()
|
||||
}
|
||||
};
|
||||
|
||||
let obsolete_cids = match obsolete_result {
|
||||
Ok(cids) => cids,
|
||||
Err(e) => {
|
||||
error!(
|
||||
"MST diff failed during finalize_repo_write: {e}. \
|
||||
Proceeding with empty obsolete set; leaked blocks \
|
||||
will be reclaimed by reachability GC."
|
||||
);
|
||||
vec![ctx.current_root_cid]
|
||||
}
|
||||
};
|
||||
|
||||
let result = commit_and_log(
|
||||
state,
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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,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
|
||||
);
|
||||
}
|
||||
@@ -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"));
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -162,13 +162,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 +182,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);
|
||||
}
|
||||
|
||||
@@ -1136,7 +1136,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);
|
||||
|
||||
@@ -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])
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
@@ -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(())
|
||||
})?;
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
@@ -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"
|
||||
}
|
||||
}
|
||||
},
|
||||
|
||||
@@ -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 = [ (import ./module.nix self) ];
|
||||
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;
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
+41
-60
@@ -1,66 +1,47 @@
|
||||
{
|
||||
lib,
|
||||
stdenvNoCC,
|
||||
|
||||
fetchDenoDeps,
|
||||
fetchFromGitHub,
|
||||
|
||||
buildGoModule,
|
||||
|
||||
deno,
|
||||
esbuild,
|
||||
}: let
|
||||
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
|
||||
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
@@ -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
|
||||
|
||||
@@ -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"
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
Generated
-1930
File diff suppressed because it is too large
Load Diff
Generated
+2475
File diff suppressed because it is too large
Load Diff
+16
-11
@@ -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;
|
||||
|
||||
@@ -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,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
|
||||
];
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user