From 1dc5e704381cfd18f8113d4de25a0c84e9b01ac8 Mon Sep 17 00:00:00 2001 From: Lewis Date: Sun, 26 Jul 2026 21:43:39 +0300 Subject: [PATCH] knot-xrpc: byte limits & read budgets ungrouped Lewis: May this revision serve well! --- crates/knot-server/Cargo.toml | 1 + crates/knot-server/src/main.rs | 91 +++++++------ crates/knot-server/tests/invariants.rs | 51 +++++++ crates/knot-sim/Cargo.toml | 2 + crates/knot-sim/src/harness.rs | 43 +++--- crates/knot-xrpc/src/lib.rs | 179 ++++++++++--------------- crates/knot-xrpc/src/reads.rs | 105 ++++++--------- crates/knot-xrpc/src/tests.rs | 27 +--- crates/knot-xrpc/tests/common/mod.rs | 104 ++++---------- crates/knot-xrpc/tests/reads.rs | 9 +- 10 files changed, 277 insertions(+), 335 deletions(-) diff --git a/crates/knot-server/Cargo.toml b/crates/knot-server/Cargo.toml index 7909252..394b39d 100644 --- a/crates/knot-server/Cargo.toml +++ b/crates/knot-server/Cargo.toml @@ -38,6 +38,7 @@ tracing = { workspace = true } tracing-subscriber = { workspace = true } [dev-dependencies] +confique = { workspace = true } http = { workspace = true } walkdir = { workspace = true } tempfile = { workspace = true } diff --git a/crates/knot-server/src/main.rs b/crates/knot-server/src/main.rs index 2a8c3c6..3b29fa4 100644 --- a/crates/knot-server/src/main.rs +++ b/crates/knot-server/src/main.rs @@ -118,7 +118,7 @@ async fn main() -> anyhow::Result<()> { }); tracing::info!( threads = resources.threads.get(), - memory_bytes = resources.memory.map(knot_resource::MemoryBudget::bytes), + memory_bytes = resources.memory.map(knot_resource::MemoryBudget::get), memory_source = ?resources.memory_source, memory_high_bytes = resources .memory_high_bytes @@ -299,29 +299,50 @@ async fn main() -> anyhow::Result<()> { .context("load or create SSH host key")?; let xrpc_limits = knot_xrpc::LimitConfig { - burst: config.xrpc.preauth_burst, - refill_micros: config.xrpc.preauth_refill_ms.saturating_mul(1_000), - per_peer_inflight: knot_xrpc::PerPeerInflight::new(config.xrpc.per_peer_inflight as usize), - global_inflight: knot_xrpc::GlobalInflight::new(config.xrpc.global_inflight as usize), + rate: Some(knot_xrpc::RateLimit { + burst: knot_xrpc::Burst::new(config.xrpc.preauth_burst), + refill: knot_xrpc::RefillMicros::new( + config.xrpc.preauth_refill_ms.saturating_mul(1_000), + ), + }), + per_peer_inflight: Some(knot_xrpc::PerPeerInflight::new( + config.xrpc.per_peer_inflight as usize, + )), + global_inflight: Some(knot_xrpc::GlobalInflight::new( + config.xrpc.global_inflight as usize, + )), + }; + let byte_limits = knot_xrpc::ByteLimits { + body: knot_xrpc::BodyLimit::new(config.xrpc.max_body_bytes as usize), + patch: knot_xrpc::PatchLimit::new(config.xrpc.max_patch_bytes as usize), + patch_decompressed: knot_xrpc::PatchDecompressedLimit::new( + config.xrpc.max_patch_decompressed_bytes, + ), + response: knot_xrpc::ResponseLimit::new(config.xrpc.max_response_bytes as usize), + archive: knot_xrpc::ArchiveLimit::new(config.xrpc.max_archive_bytes), + fork_pack: knot_xrpc::ForkPackLimit::new(config.xrpc.fork_max_pack_bytes), + pack: knot_xrpc::MaxWireBytes::new(ssh_max_pack_bytes), + }; + let budgets = knot_xrpc::Budgets { + tree_last_commit: knot_xrpc::TreeReadBudget::new(knot_xrpc::ReadBudget::Within( + Duration::from_millis(config.xrpc.tree_last_commit_budget_ms), + )), + blob_last_commit: knot_xrpc::BlobReadBudget::new(knot_xrpc::ReadBudget::Within( + Duration::from_millis(config.xrpc.blob_last_commit_budget_ms), + )), + languages: knot_xrpc::LanguagesReadBudget::new(knot_xrpc::ReadBudget::Within( + Duration::from_millis(config.xrpc.languages_budget_ms), + )), + languages_push: knot_xrpc::LanguagesPushBudget::new(Duration::from_millis( + config.xrpc.languages_push_budget_ms, + )), }; - let xrpc_max_body_bytes = config.xrpc.max_body_bytes as usize; - let xrpc_max_patch_bytes = config.xrpc.max_patch_bytes as usize; - let xrpc_max_patch_decompressed_bytes = config.xrpc.max_patch_decompressed_bytes; - let xrpc_max_response_bytes = config.xrpc.max_response_bytes as usize; - let xrpc_max_archive_bytes = config.xrpc.max_archive_bytes; - let xrpc_tree_last_commit_budget = - Duration::from_millis(config.xrpc.tree_last_commit_budget_ms); - let xrpc_blob_last_commit_budget = - Duration::from_millis(config.xrpc.blob_last_commit_budget_ms); - let xrpc_languages_budget = Duration::from_millis(config.xrpc.languages_budget_ms); - let xrpc_languages_push_budget = Duration::from_millis(config.xrpc.languages_push_budget_ms); - let fork_max_pack_bytes = config.xrpc.fork_max_pack_bytes; let committer = knot_xrpc::Committer { name: AuthorName::new(config.git.user_name.clone()), email: Email::new(config.git.user_email.clone()), }; let reservations = Arc::new(knot_xrpc::Reservations::new( - config.xrpc.reservation_ttl_secs as i64, + knot_xrpc::ReservationTtl::new(config.xrpc.reservation_ttl_secs as i64), knot_xrpc::PerActorQuota::new(config.xrpc.per_actor_reservations as usize), knot_xrpc::GlobalQuota::new(config.xrpc.max_pending_reservations as usize), )); @@ -439,6 +460,8 @@ async fn main() -> anyhow::Result<()> { (knot_maintenance::MaintenanceHandle::disabled(), None, None) }; + let slots = knot_resource::Slots::for_machine(); + let ssh_base = knot_ssh::SshState::new( layout.clone(), Arc::clone(&index), @@ -449,11 +472,12 @@ async fn main() -> anyhow::Result<()> { appview_endpoint.clone(), admins.clone(), admission, - ssh_max_pack_bytes, - xrpc_languages_push_budget, + byte_limits.pack, + budgets.languages_push, ) .with_maintenance(maintenance_handle.clone()) .with_limits(pack_limits) + .with_slots(slots.clone()) .with_catalog(Arc::clone(&catalog)); let ssh_state = Arc::new(match &lfs_handle { Some(handle) => ssh_base.with_lfs(handle.clone(), lfs_max_ssh_transfers), @@ -477,35 +501,15 @@ async fn main() -> anyhow::Result<()> { reservations, trusted_proxy_header, committer, - max_body_bytes: knot_xrpc::BodyLimit::new(xrpc_max_body_bytes), - max_patch_bytes: knot_xrpc::PatchLimit::new(xrpc_max_patch_bytes), - max_patch_decompressed_bytes: knot_xrpc::PatchDecompressedLimit::new( - xrpc_max_patch_decompressed_bytes, - ), - max_response_bytes: knot_xrpc::ResponseLimit::new(xrpc_max_response_bytes), - max_archive_bytes: knot_xrpc::ArchiveLimit::new(xrpc_max_archive_bytes), - tree_last_commit_budget: knot_xrpc::TreeReadBudget::new(knot_xrpc::ReadBudget::Within( - xrpc_tree_last_commit_budget, - )), - blob_last_commit_budget: knot_xrpc::BlobReadBudget::new(knot_xrpc::ReadBudget::Within( - xrpc_blob_last_commit_budget, - )), - languages_budget: knot_xrpc::LanguagesReadBudget::new(knot_xrpc::ReadBudget::Within( - xrpc_languages_budget, - )), - languages_push_budget: xrpc_languages_push_budget, + byte_limits, + budgets, git_http, - fork_max_pack_bytes: knot_xrpc::ForkPackLimit::new(fork_max_pack_bytes), pack_limits, - max_pack_bytes: knot_xrpc::PackLimit::new(ssh_max_pack_bytes), service_owner, subscriber_gate, maintenance: maintenance_handle, appview: appview_endpoint, - resolve_slots: knot_xrpc::ResolveSlots::new(Arc::new(tokio::sync::Semaphore::new(16))), - receive_slots: knot_xrpc::ReceiveSlots::new(Arc::new(tokio::sync::Semaphore::new( - knot_resource::threads().get(), - ))), + slots: slots.clone(), events, lfs: lfs_handle.map(|handle| knot_xrpc::LfsWeb::new(handle, lfs_max_http_downloads)), catalog: Arc::clone(&catalog), @@ -531,6 +535,7 @@ async fn main() -> anyhow::Result<()> { resolver, Some(receive_advertiser), Some(handle_resolver), + slots.pack.clone(), pack_cache_config, Arc::clone(&catalog), xrpc_state.knot_hostname.clone(), diff --git a/crates/knot-server/tests/invariants.rs b/crates/knot-server/tests/invariants.rs index 288da8e..66e8017 100644 --- a/crates/knot-server/tests/invariants.rs +++ b/crates/knot-server/tests/invariants.rs @@ -187,3 +187,54 @@ fn nothing_depends_on_the_offline_migration_tool() { "knot-migrate is an offline one-shot tool whose rusqlite dependency must never reach the server, but it is depended on by: {dependents:?}" ); } + +#[test] +fn the_shared_limit_defaults_match_the_config_defaults() { + use confique::{Config, Layer}; + use knot_xrpc::{Budgets, ByteLimits, ReadBudget}; + + fn ms(budget: ReadBudget) -> u64 { + match budget { + ReadBudget::Within(within) => within.as_millis() as u64, + ReadBudget::Unbounded => u64::MAX, + } + } + + let xrpc = ::Layer::default_values(); + let server = ::Layer::default_values(); + let bytes = ByteLimits::default(); + let budgets = Budgets::default(); + let push_ms = budgets.languages_push.get().as_millis() as u64; + + let configured = [ + ("body", xrpc.max_body_bytes), + ("patch", xrpc.max_patch_bytes), + ("patch_decompressed", xrpc.max_patch_decompressed_bytes), + ("response", xrpc.max_response_bytes), + ("archive", xrpc.max_archive_bytes), + ("fork_pack", xrpc.fork_max_pack_bytes), + ("pack", server.ssh_max_pack_bytes), + ("tree_last_commit", xrpc.tree_last_commit_budget_ms), + ("blob_last_commit", xrpc.blob_last_commit_budget_ms), + ("languages", xrpc.languages_budget_ms), + ("languages_push", xrpc.languages_push_budget_ms), + ] + .map(|(name, value)| (name, value.expect("every limit has a config default"))); + let shared = [ + ("body", bytes.body.get() as u64), + ("patch", bytes.patch.get() as u64), + ("patch_decompressed", bytes.patch_decompressed.get()), + ("response", bytes.response.get() as u64), + ("archive", bytes.archive.get()), + ("fork_pack", bytes.fork_pack.get()), + ("pack", bytes.pack.get() as u64), + ("tree_last_commit", ms(budgets.tree_last_commit.get())), + ("blob_last_commit", ms(budgets.blob_last_commit.get())), + ("languages", ms(budgets.languages.get())), + ("languages_push", push_ms), + ]; + assert_eq!( + configured, shared, + "the config defaults and the in-code defaults must match. update ByteLimits::default and Budgets::default alongside the config defaults" + ); +} diff --git a/crates/knot-sim/Cargo.toml b/crates/knot-sim/Cargo.toml index 59302d5..65189c0 100644 --- a/crates/knot-sim/Cargo.toml +++ b/crates/knot-sim/Cargo.toml @@ -17,6 +17,7 @@ knot-runtime = { workspace = true } knot-git = { workspace = true } knot-index = { workspace = true } knot-pack = { workspace = true } +knot-resource = { workspace = true } knot-atproto = { workspace = true } knot-secrets = { workspace = true } knot-xrpc = { workspace = true } @@ -38,6 +39,7 @@ futures = { workspace = true } tempfile = { workspace = true } [dev-dependencies] +knot-fixtures = { workspace = true } knot-ssh = { workspace = true } knot-edge = { workspace = true } knot-lfs = { workspace = true } diff --git a/crates/knot-sim/src/harness.rs b/crates/knot-sim/src/harness.rs index 9cd7fbd..f153cc4 100644 --- a/crates/knot-sim/src/harness.rs +++ b/crates/knot-sim/src/harness.rs @@ -30,10 +30,10 @@ use knot_types::{ KnotServiceUrl, Oid, OwnerDid, RefName, RepoDid, RepoName, RepoRkey, UnixSeconds, }; use knot_xrpc::{ - ArchiveLimit, BlobReadBudget, BodyLimit, CobLocks, Committer, ForkPackLimit, GlobalInflight, - GlobalQuota, LanguagesReadBudget, LimitConfig, PackLimit, PatchDecompressedLimit, PatchLimit, - PerActorQuota, PerPeerInflight, PreAuthLimiter, ReadBudget, ReceiveSlots, Reservations, - ResolveSlots, ResponseLimit, TreeReadBudget, XrpcState, + BlobReadBudget, BodyLimit, Budgets, ByteLimits, CobLocks, Committer, GlobalInflight, + GlobalQuota, LanguagesPushBudget, LanguagesReadBudget, LimitConfig, MaxWireBytes, + PerActorQuota, PerPeerInflight, PreAuthLimiter, ReadBudget, ReservationTtl, Reservations, + ResponseLimit, TreeReadBudget, XrpcState, }; use crate::realdata::{ @@ -721,14 +721,13 @@ fn assemble_router(parts: StateParts) -> Router { meta_path, knot_service_url, limiter: Arc::new(PreAuthLimiter::with_config(LimitConfig { - burst: 1_000_000, - refill_micros: 1, - per_peer_inflight: PerPeerInflight::new(4096), - global_inflight: GlobalInflight::new(4096), + rate: None, + per_peer_inflight: Some(PerPeerInflight::new(4096)), + global_inflight: Some(GlobalInflight::new(4096)), })), cob_locks: Arc::new(CobLocks::default()), reservations: Arc::new(Reservations::new( - 1_000_000, + ReservationTtl::new(1_000_000), PerActorQuota::new(256), GlobalQuota::new(256), )), @@ -737,23 +736,24 @@ fn assemble_router(parts: StateParts) -> Router { name: AuthorName::new("knot"), email: Email::new("knot@nel.pet"), }, - max_body_bytes: BodyLimit::new(256 * 1024), - max_patch_bytes: PatchLimit::new(16 * 1024 * 1024), - max_patch_decompressed_bytes: PatchDecompressedLimit::new(128 * 1024 * 1024), - max_response_bytes: ResponseLimit::new(8 * 1024 * 1024), - max_archive_bytes: ArchiveLimit::new(1024 * 1024 * 1024), - tree_last_commit_budget: TreeReadBudget::new(ReadBudget::Unbounded), - blob_last_commit_budget: BlobReadBudget::new(ReadBudget::Unbounded), - languages_budget: LanguagesReadBudget::new(ReadBudget::Unbounded), - languages_push_budget: Duration::from_secs(120), + byte_limits: ByteLimits { + body: BodyLimit::new(256 * 1024), + response: ResponseLimit::new(8 * 1024 * 1024), + pack: MaxWireBytes::new(1024 * 1024 * 1024), + ..ByteLimits::default() + }, + budgets: Budgets { + tree_last_commit: TreeReadBudget::new(ReadBudget::Unbounded), + blob_last_commit: BlobReadBudget::new(ReadBudget::Unbounded), + languages: LanguagesReadBudget::new(ReadBudget::Unbounded), + languages_push: LanguagesPushBudget::new(Duration::from_secs(120)), + }, git_http: Arc::new(FakeHttp::new(|_request: &HttpRequest| { Err(NetworkError::Connect( "simulation serves no git upstream".to_string(), )) })), - fork_max_pack_bytes: ForkPackLimit::new(1024 * 1024 * 1024), pack_limits: knot_pack::PackLimits::default(), - max_pack_bytes: PackLimit::new(1024 * 1024 * 1024), service_owner, events: Arc::new(EventLog::new(Arc::clone(&clock), 4096)), subscriber_gate: Arc::new(SubscriberGate::new( @@ -762,8 +762,7 @@ fn assemble_router(parts: StateParts) -> Router { )), maintenance: knot_maintenance::MaintenanceHandle::disabled(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), - resolve_slots: ResolveSlots::new(Arc::new(tokio::sync::Semaphore::new(8))), - receive_slots: ReceiveSlots::new(Arc::new(tokio::sync::Semaphore::new(8))), + slots: knot_resource::Slots::testing(8), lfs: None, catalog: Arc::new(knot_messages::Catalog::defaults()), }); diff --git a/crates/knot-xrpc/src/lib.rs b/crates/knot-xrpc/src/lib.rs index c8cffb0..4d2f39f 100644 --- a/crates/knot-xrpc/src/lib.rs +++ b/crates/knot-xrpc/src/lib.rs @@ -7,7 +7,6 @@ mod error; mod events; mod forks; mod lfs; -mod limit; mod lists; mod locks; mod members; @@ -26,12 +25,16 @@ mod wire; mod tests; pub use error::XrpcError; +pub use knot_pack::MaxWireBytes; +pub use knot_postreceive::LanguagesPushBudget; +pub use knot_resource::{ + Burst, GlobalInflight, LimitConfig, PerPeerInflight, PreAuthLimiter, RateLimit, RefillMicros, +}; pub use lfs::LfsWeb; -pub use limit::{GlobalInflight, LimitConfig, PerPeerInflight, PreAuthLimiter}; pub use locks::CobLocks; pub use merge::Committer; pub use receive::advertiser as receive_advertiser; -pub use reservations::{GlobalQuota, PerActorQuota, Reservations}; +pub use reservations::{GlobalQuota, PerActorQuota, ReservationTtl, Reservations}; use std::collections::BTreeSet; use std::net::IpAddr; @@ -51,26 +54,26 @@ use http::{HeaderMap, HeaderValue, StatusCode, header::AUTHORIZATION}; use serde::de::DeserializeOwned; use serde_json::json; -use knot_atproto::{Atproto, AtprotoError, IdentityError, ServiceJwt}; +use knot_atproto::{Atproto, AtprotoError, ServiceJwt}; use knot_events::{EventLog, SubscriberGate}; use knot_git::Layout; use knot_index::{Index, Resolved}; use knot_maintenance::MaintenanceHandle; +use knot_resource::Slots; use knot_runtime::{Clock, Entropy, HttpTransport}; use knot_secrets::SealedStore; use knot_types::{ AccountDid, AdmissionPolicy, AppviewEndpoint, KnotHostname, KnotId, KnotServiceUrl, Nsid, OwnerDid, OwnerRef, RepoDid, RepoRkey, UnixSeconds, }; -use tokio::sync::Semaphore; use base64::Engine; use knot_pack::SocketPeer; -use limit::{Admission, AdmitGuard}; +use knot_resource::{AdmitGuard, Refusal}; pub(crate) const PUSH_NSID: &str = "sh.tangled.repo.push"; -#[derive(Debug, Clone, Copy)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum ReadBudget { Within(Duration), Unbounded, @@ -88,73 +91,62 @@ impl ReadBudget { // `XrpcState` keeps a bunch of these side by side, // some usize & some u64. // Within each group every one of them typechecked in every other one's slot. -// TODO: newtype -macro_rules! byte_limit { - ($name:ident, $prim:ty) => { - #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] - pub struct $name($prim); - - impl $name { - pub const fn new(v: $prim) -> Self { - Self(v) - } - - pub const fn get(self) -> $prim { - self.0 - } +knot_types::scalar_newtype! { + pub struct BodyLimit(usize); + pub struct PatchLimit(usize); + pub struct PatchDecompressedLimit(u64); + pub struct ResponseLimit(usize); + pub struct ArchiveLimit(u64); + pub struct ForkPackLimit(u64); + pub struct TreeReadBudget(ReadBudget); + pub struct BlobReadBudget(ReadBudget); + pub struct LanguagesReadBudget(ReadBudget); +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ByteLimits { + pub body: BodyLimit, + pub patch: PatchLimit, + pub patch_decompressed: PatchDecompressedLimit, + pub response: ResponseLimit, + pub archive: ArchiveLimit, + pub fork_pack: ForkPackLimit, + pub pack: MaxWireBytes, +} + +impl Default for ByteLimits { + fn default() -> Self { + Self { + body: BodyLimit::new(64 * 1024), + patch: PatchLimit::new(16 * 1024 * 1024), + patch_decompressed: PatchDecompressedLimit::new(128 * 1024 * 1024), + response: ResponseLimit::new(5 * 1024 * 1024), + archive: ArchiveLimit::new(1024 * 1024 * 1024), + fork_pack: ForkPackLimit::new(1024 * 1024 * 1024), + pack: MaxWireBytes::new(8 * 1024 * 1024 * 1024), } - }; + } } -byte_limit!(BodyLimit, usize); -byte_limit!(PatchLimit, usize); -byte_limit!(PatchDecompressedLimit, u64); -byte_limit!(ResponseLimit, usize); -byte_limit!(ArchiveLimit, u64); -byte_limit!(ForkPackLimit, u64); -byte_limit!(PackLimit, usize); - -macro_rules! read_budget_role { - ($name:ident) => { - #[derive(Clone, Copy)] - pub struct $name(ReadBudget); - - impl $name { - pub const fn new(b: ReadBudget) -> Self { - Self(b) - } - - pub const fn get(self) -> ReadBudget { - self.0 - } - } - }; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct Budgets { + pub tree_last_commit: TreeReadBudget, + pub blob_last_commit: BlobReadBudget, + pub languages: LanguagesReadBudget, + pub languages_push: LanguagesPushBudget, } -read_budget_role!(TreeReadBudget); -read_budget_role!(BlobReadBudget); -read_budget_role!(LanguagesReadBudget); - -macro_rules! semaphore_role { - ($name:ident) => { - #[derive(Clone)] - pub struct $name(Arc); - - impl $name { - pub fn new(s: Arc) -> Self { - Self(s) - } - - pub fn get(&self) -> &Arc { - &self.0 - } +impl Default for Budgets { + fn default() -> Self { + Self { + tree_last_commit: TreeReadBudget::new(ReadBudget::Within(Duration::from_millis(300))), + blob_last_commit: BlobReadBudget::new(ReadBudget::Within(Duration::from_millis(2_000))), + languages: LanguagesReadBudget::new(ReadBudget::Within(Duration::from_millis(1_000))), + languages_push: LanguagesPushBudget::new(Duration::from_millis(2_000)), } - }; + } } -semaphore_role!(ResolveSlots); -semaphore_role!(ReceiveSlots); - pub struct XrpcState { pub layout: Layout, pub index: Arc, @@ -172,26 +164,16 @@ pub struct XrpcState { pub reservations: Arc, pub trusted_proxy_header: Option, pub committer: Committer, - pub max_body_bytes: BodyLimit, - pub max_patch_bytes: PatchLimit, - pub max_patch_decompressed_bytes: PatchDecompressedLimit, - pub max_response_bytes: ResponseLimit, - pub max_archive_bytes: ArchiveLimit, - pub tree_last_commit_budget: TreeReadBudget, - pub blob_last_commit_budget: BlobReadBudget, - pub languages_budget: LanguagesReadBudget, - pub languages_push_budget: Duration, + pub byte_limits: ByteLimits, + pub budgets: Budgets, pub git_http: Arc, - pub fork_max_pack_bytes: ForkPackLimit, pub pack_limits: knot_pack::PackLimits, - pub max_pack_bytes: PackLimit, pub service_owner: AccountDid, pub events: Arc>, pub subscriber_gate: Arc, pub maintenance: MaintenanceHandle, pub appview: AppviewEndpoint, - pub resolve_slots: ResolveSlots, - pub receive_slots: ReceiveSlots, + pub slots: Slots, pub lfs: Option, pub catalog: Arc, } @@ -276,7 +258,7 @@ pub fn router(state: Arc>) -> Router let merge_routes = Router::new() .route(merge::MERGE_ROUTE, post(merge::merge::)) .route(merge::MERGE_CHECK_ROUTE, post(merge::merge_check::)) - .layer(DefaultBodyLimit::max(state.max_patch_bytes.get())); + .layer(DefaultBodyLimit::max(state.byte_limits.patch.get())); Router::new() .merge(merge_routes) .route(members::ADD_ROUTE, post(members::add_member::)) @@ -334,7 +316,7 @@ pub fn router(state: Arc>) -> Router ) .route(service::VERSION_ROUTE, get(service::version)) .route(service::OWNER_ROUTE, get(service::owner::)) - .layer(DefaultBodyLimit::max(state.max_body_bytes.get())) + .layer(DefaultBodyLimit::max(state.byte_limits.body.get())) .layer(from_fn_with_state( Arc::clone(&state), enforce_pre_auth_limit::, @@ -379,15 +361,17 @@ pub(crate) fn admit_pre_auth( state: &XrpcState, peer: Option, ) -> Result { - match state.limiter.admit(peer, state.atproto.now()) { - Admission::Admitted(guard) => Ok(guard), - Admission::RateLimited => Err(XrpcError::rate_limited( - "too many pre-authentication requests, slow down", - )), - Admission::Saturated => Err(XrpcError::overloaded( - "knot is shedding pre-authentication load, retry shortly", - )), - } + state + .limiter + .admit(peer, state.atproto.now()) + .map_err(|refusal| match refusal { + Refusal::RateLimited => { + XrpcError::rate_limited("too many pre-authentication requests, retry shortly") + } + Refusal::Saturated => { + XrpcError::overloaded("knot is shedding pre-authentication load, retry shortly") + } + }) } pub(crate) const BASIC_CHALLENGE: HeaderValue = HeaderValue::from_static("Basic realm=\"knot\""); @@ -504,21 +488,6 @@ where } } -pub(crate) fn map_atproto(error: AtprotoError) -> XrpcError { - if error.is_transient() { - return XrpcError::upstream_unavailable(error.to_string()); - } - match error { - AtprotoError::Resolve(_) => XrpcError::invalid_request(error.to_string()), - AtprotoError::PlcSubmit { .. } => XrpcError::bad_gateway(error.to_string()), - other => XrpcError::internal(other.to_string()), - } -} - -pub(crate) fn map_identity(error: IdentityError) -> XrpcError { - XrpcError::internal(error.to_string()) -} - #[derive(serde::Deserialize)] #[serde(transparent)] pub(crate) struct OwnerSegment(String); diff --git a/crates/knot-xrpc/src/reads.rs b/crates/knot-xrpc/src/reads.rs index 36959a7..db3bd40 100644 --- a/crates/knot-xrpc/src/reads.rs +++ b/crates/knot-xrpc/src/reads.rs @@ -140,7 +140,7 @@ fn commit_for(repo: &Repo, refspec: &Revspec) -> Result { match repo.reachable_from_public(commit) { Ok(true) => Ok(commit), Ok(false) => Err(ref_not_found()), - Err(error) => Err(XrpcError::internal(error.to_string())), + Err(error) => Err(error.into()), } } @@ -287,8 +287,8 @@ pub(crate) async fn repo_tree( ) -> Result { let did = resolve_repo(&state, ¶ms.repo)?; let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); - let tree_deadline = state.tree_last_commit_budget.get().deadline(); + let limit = state.byte_limits.response.get(); + let tree_deadline = state.budgets.tree_last_commit.get().deadline(); run_blocking(move || { let repo = open(&layout, &did)?; let commit = commit_for(&repo, ¶ms.refspec)?; @@ -302,8 +302,7 @@ pub(crate) async fn repo_tree( let dir = params.path.dir().ok_or_else(path_not_found)?; let path = params.path.as_str(); let entries = repo - .tree_entries_at(commit, dir) - .map_err(|error| XrpcError::internal(error.to_string()))? + .tree_entries_at(commit, dir)? .ok_or_else(path_not_found)?; let names: Vec = entries.iter().map(|entry| entry.name.clone()).collect(); let attributed = repo @@ -394,13 +393,12 @@ pub(crate) async fn repo_log( let offset = params.cursor.get(); let limit = params.limit.get(); let layout = state.layout.clone(); - let response_limit = state.max_response_bytes.get(); + let response_limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; let start = commit_for(&repo, ¶ms.refspec)?; - let (commits, total) = repo - .log_window(start, LogSkip::new(offset), LogLimit::new(limit)) - .map_err(|error| XrpcError::internal(error.to_string()))?; + let (commits, total) = + repo.log_window(start, LogSkip::new(offset), LogLimit::new(limit))?; json( LogOut { commits: commits.iter().map(CommitWire::of).collect(), @@ -440,12 +438,10 @@ pub(crate) async fn repo_branches( let offset = params.cursor.get(); let limit = params.limit.get(); let layout = state.layout.clone(); - let response_limit = state.max_response_bytes.get(); + let response_limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; - let mut branches = repo - .branch_list() - .map_err(|error| XrpcError::internal(error.to_string()))?; + let mut branches = repo.branch_list()?; branches.sort_by(|a, b| { b.tip .created_at() @@ -504,7 +500,7 @@ pub(crate) async fn repo_branch( return Err(XrpcError::invalid_request("missing name parameter")); }; let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); + let limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; let branch_not_found = @@ -562,12 +558,10 @@ pub(crate) async fn repo_tags( let offset = params.cursor.get(); let limit = params.limit.get(); let layout = state.layout.clone(); - let response_limit = state.max_response_bytes.get(); + let response_limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; - let mut tags = repo - .tag_list() - .map_err(|error| XrpcError::internal(error.to_string()))?; + let mut tags = repo.tag_list()?; tags.sort_by(|a, b| { b.created_at .cmp(&a.created_at) @@ -605,12 +599,11 @@ pub(crate) async fn repo_tag( return Err(XrpcError::invalid_request("missing tag parameter")); }; let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); + let limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; let info = repo - .tag_list() - .map_err(|error| XrpcError::internal(error.to_string()))? + .tag_list()? .into_iter() .find(|tag| tag.name == name) .ok_or_else(|| { @@ -736,8 +729,8 @@ pub(crate) async fn repo_blob( return Err(XrpcError::invalid_request("missing path parameter")); } let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); - let blob_deadline = state.blob_last_commit_budget.get().deadline(); + let limit = state.byte_limits.response.get(); + let blob_deadline = state.budgets.blob_last_commit.get().deadline(); run_blocking(move || { let refspec = params.refspec.as_str().to_string(); let path = params.path.as_str().to_string(); @@ -781,8 +774,7 @@ pub(crate) async fn repo_blob( }; let file_path = params.path.file().ok_or_else(file_not_found)?; let entry = repo - .entry_at(commit, file_path) - .map_err(|error| XrpcError::internal(error.to_string()))? + .entry_at(commit, file_path)? .filter(|entry| { matches!( entry.kind, @@ -869,19 +861,15 @@ pub(crate) async fn repo_diff( ) -> Result { let did = resolve_repo(&state, ¶ms.repo)?; let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); + let limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; let target = commit_for(&repo, ¶ms.refspec)?; - let commit = repo - .find_commit(target) - .map_err(|error| XrpcError::internal(error.to_string()))?; - let patches = repo - .commit_patches(knot_git::PatchRange { - base: commit.parents.first().copied(), - head: target, - }) - .map_err(|error| XrpcError::internal(error.to_string()))?; + let commit = repo.find_commit(target)?; + let patches = repo.commit_patches(knot_git::PatchRange { + base: commit.parents.first().copied(), + head: target, + })?; json( DiffOut { refspec: params.refspec.as_str().to_string(), @@ -983,7 +971,7 @@ pub(crate) async fn repo_compare( return Err(XrpcError::invalid_request("missing rev2 parameter")); } let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); + let limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; let resolve = |rev: &str| { @@ -1004,7 +992,7 @@ pub(crate) async fn repo_compare( match repo.reachable_from_public(commit) { Ok(true) => Ok(commit), Ok(false) => Err(revision_not_found()), - Err(error) => Err(XrpcError::internal(error.to_string())), + Err(error) => Err(error.into()), } }; let base = resolve(&rev1)?; @@ -1334,14 +1322,12 @@ pub(crate) async fn repo_archive( let temp = run_blocking({ let layout = state.layout.clone(); let did = did.clone(); - let archive_limit = state.max_archive_bytes.get(); + let archive_limit = state.byte_limits.archive.get(); let tree_prefix = knot_git::ArchivePrefix::new(format!("{archive_prefix}/")) .expect("validated prefix with trailing slash stays valid"); move || { let repo = open(&layout, &did)?; - let tree = repo - .peel_to_tree(resolved) - .map_err(|error| XrpcError::internal(error.to_string()))?; + let tree = repo.peel_to_tree(resolved)?; let temp = tempfile::NamedTempFile::new() .map_err(|error| XrpcError::internal(format!("cannot spool archive: {error}")))?; let file = temp @@ -1389,26 +1375,21 @@ pub(crate) async fn repo_archive( let serve_response = ServeFile::new(temp.path()) .oneshot(request) .await - .map_err(|error| match error {})?; + .unwrap_or_else(|error| match error {}); drop(temp); let mut response = serve_response.map(Body::new); let headers = response.headers_mut(); headers.insert(header::CONTENT_TYPE, HeaderValue::from_static(content_type)); - headers.insert( - header::CONTENT_DISPOSITION, - HeaderValue::from_str(&disposition) - .map_err(|error| XrpcError::internal(error.to_string()))?, - ); + let header_value = |value: &str| { + HeaderValue::from_str(value).map_err(|error| XrpcError::internal(error.to_string())) + }; + headers.insert(header::CONTENT_DISPOSITION, header_value(&disposition)?); headers.insert( header::LINK, - HeaderValue::from_str(&format!("<{immutable}>; rel=\"immutable\"")) - .map_err(|error| XrpcError::internal(error.to_string()))?, - ); - headers.insert( - header::ETAG, - HeaderValue::from_str(&etag).map_err(|error| XrpcError::internal(error.to_string()))?, + header_value(&format!("<{immutable}>; rel=\"immutable\""))?, ); + headers.insert(header::ETAG, header_value(&etag)?); headers.insert(header::CACHE_CONTROL, HeaderValue::from_static("no-cache")); Ok(response) } @@ -1444,13 +1425,12 @@ pub(crate) async fn repo_languages( ) -> Result { let did = resolve_repo(&state, ¶ms.repo)?; let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); - let languages_deadline = state.languages_budget.get().deadline(); + let limit = state.byte_limits.response.get(); + let languages_deadline = state.budgets.languages.get().deadline(); run_blocking(move || { let repo = open(&layout, &did)?; let commit = commit_for(&repo, ¶ms.refspec)?; - let sizes = knot_langs::analyze(&repo, commit, languages_deadline) - .map_err(|error| XrpcError::internal(error.to_string()))?; + let sizes = knot_langs::analyze(&repo, commit, languages_deadline)?; let total: u64 = sizes.values().map(|size| size.get()).sum(); let mut languages: Vec = sizes .iter() @@ -1494,7 +1474,7 @@ pub(crate) async fn repo_get_default_branch( ) -> Result { let did = resolve_repo(&state, ¶ms.repo)?; let layout = state.layout.clone(); - let limit = state.max_response_bytes.get(); + let limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; let name = repo @@ -1555,7 +1535,7 @@ pub(crate) async fn repo_describe_repo( owner_did: owner, rkey, }, - state.max_response_bytes.get(), + state.byte_limits.response.get(), ) } @@ -1600,12 +1580,11 @@ pub(crate) async fn git_list_refs( let offset = params.cursor; let limit = params.limit; let layout = state.layout.clone(); - let response_limit = state.max_response_bytes.get(); + let response_limit = state.byte_limits.response.get(); run_blocking(move || { let repo = open(&layout, &did)?; let mut refs: Vec<_> = repo - .references() - .map_err(|error| XrpcError::internal(error.to_string()))? + .references()? .into_iter() .filter(|record| is_public_ref(&record.name)) .collect(); @@ -1683,7 +1662,7 @@ pub(crate) async fn sync_list_repos( .collect(); let cursor = next_cursor(offset, limit, Total::new(total)); let layout = state.layout.clone(); - let response_limit = state.max_response_bytes.get(); + let response_limit = state.byte_limits.response.get(); run_blocking(move || { let repos: Vec = page.iter() diff --git a/crates/knot-xrpc/src/tests.rs b/crates/knot-xrpc/src/tests.rs index 9e4fa36..776637d 100644 --- a/crates/knot-xrpc/src/tests.rs +++ b/crates/knot-xrpc/src/tests.rs @@ -2,7 +2,6 @@ use std::collections::{BTreeSet, HashMap, HashSet}; use std::path::PathBuf; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; -use std::time::Duration; use axum::body::Bytes; use axum::extract::State; @@ -322,25 +321,10 @@ fn state_from( name: AuthorName::new("Tangled"), email: Email::new("noreply@tangled.sh"), }, - max_body_bytes: crate::BodyLimit::new(64 * 1024), - max_patch_bytes: crate::PatchLimit::new(16 * 1024 * 1024), - max_patch_decompressed_bytes: crate::PatchDecompressedLimit::new(128 * 1024 * 1024), - max_response_bytes: crate::ResponseLimit::new(5 * 1024 * 1024), - max_archive_bytes: crate::ArchiveLimit::new(1024 * 1024 * 1024), - tree_last_commit_budget: crate::TreeReadBudget::new(crate::ReadBudget::Within( - Duration::from_millis(300), - )), - blob_last_commit_budget: crate::BlobReadBudget::new(crate::ReadBudget::Within( - Duration::from_millis(2000), - )), - languages_budget: crate::LanguagesReadBudget::new(crate::ReadBudget::Within( - Duration::from_millis(1000), - )), - languages_push_budget: Duration::from_millis(2000), + byte_limits: crate::ByteLimits::default(), + budgets: crate::Budgets::default(), git_http, - fork_max_pack_bytes: crate::ForkPackLimit::new(1024 * 1024 * 1024), pack_limits: knot_pack::PackLimits::default(), - max_pack_bytes: crate::PackLimit::new(1024 * 1024 * 1024), service_owner: account(ADMIN_HOST), events: Arc::new(knot_events::EventLog::new( ManualClock::new(UnixMicros::new(1_000_000_000)), @@ -352,8 +336,7 @@ fn state_from( )), maintenance: knot_maintenance::MaintenanceHandle::disabled(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), - resolve_slots: crate::ResolveSlots::new(Arc::new(tokio::sync::Semaphore::new(8))), - receive_slots: crate::ReceiveSlots::new(Arc::new(tokio::sync::Semaphore::new(8))), + slots: knot_resource::Slots::testing(8), lfs: None, catalog: Arc::new(knot_messages::Catalog::defaults()), }) @@ -484,7 +467,7 @@ impl World { responder, admission, Arc::new(crate::Reservations::new( - 1_000_000, + crate::ReservationTtl::new(1_000_000), crate::PerActorQuota::new(per_actor), crate::GlobalQuota::new(global), )), @@ -527,7 +510,7 @@ fn build_state(responder: Responder, rebuild: bool) -> (TempDir, SharedState) { responder, AdmissionPolicy::Closed, Arc::new(crate::Reservations::new( - 1_000_000, + crate::ReservationTtl::new(1_000_000), crate::PerActorQuota::new(256), crate::GlobalQuota::new(256), )), diff --git a/crates/knot-xrpc/tests/common/mod.rs b/crates/knot-xrpc/tests/common/mod.rs index 279d756..a353967 100644 --- a/crates/knot-xrpc/tests/common/mod.rs +++ b/crates/knot-xrpc/tests/common/mod.rs @@ -3,7 +3,6 @@ use std::collections::BTreeSet; use std::path::Path; use std::sync::Arc; -use std::time::Duration; use std::sync::atomic::{AtomicU64, Ordering}; @@ -33,10 +32,8 @@ use knot_types::{ RepoName, RepoRkey, UnixSeconds, }; use knot_xrpc::{ - ArchiveLimit, BlobReadBudget, BodyLimit, CobLocks, ForkPackLimit, GlobalInflight, GlobalQuota, - LanguagesReadBudget, LimitConfig, PackLimit, PatchDecompressedLimit, PatchLimit, PerActorQuota, - PerPeerInflight, PreAuthLimiter, ReadBudget, ReceiveSlots, Reservations, ResolveSlots, - ResponseLimit, TreeReadBudget, XrpcState, + ArchiveLimit, Budgets, ByteLimits, CobLocks, GlobalQuota, LimitConfig, PerActorQuota, + PreAuthLimiter, Reservations, ResponseLimit, XrpcState, }; pub const KNOT_HOST: &str = "knot.nel.pet"; @@ -54,79 +51,55 @@ pub struct World { impl World { pub fn new() -> Self { - Self::build( - true, - 5 * 1024 * 1024, - 1024 * 1024 * 1024, - ObjectFormat::SHA1, - ) + Self::build(true, ByteLimits::default(), ObjectFormat::SHA1) } pub fn unshed() -> Self { Self::build_with_limits( true, - 5 * 1024 * 1024, - 1024 * 1024 * 1024, + ByteLimits::default(), ObjectFormat::SHA1, - LimitConfig { - burst: u32::MAX, - refill_micros: 1, - per_peer_inflight: PerPeerInflight::new(usize::MAX), - global_inflight: GlobalInflight::new(usize::MAX), - }, + LimitConfig::unmetered(), ) } pub fn sha256() -> Self { - Self::build( - true, - 5 * 1024 * 1024, - 1024 * 1024 * 1024, - ObjectFormat::SHA256, - ) + Self::build(true, ByteLimits::default(), ObjectFormat::SHA256) } pub fn warming() -> Self { + Self::build(false, ByteLimits::default(), ObjectFormat::SHA1) + } + + pub fn with_response_limit(response: ResponseLimit) -> Self { Self::build( - false, - 5 * 1024 * 1024, - 1024 * 1024 * 1024, + true, + ByteLimits { + response, + ..ByteLimits::default() + }, ObjectFormat::SHA1, ) } - pub fn with_response_limit(max_response_bytes: usize) -> Self { + pub fn with_archive_limit(archive: ArchiveLimit) -> Self { Self::build( true, - max_response_bytes, - 1024 * 1024 * 1024, + ByteLimits { + archive, + ..ByteLimits::default() + }, ObjectFormat::SHA1, ) } - pub fn with_archive_limit(max_archive_bytes: u64) -> Self { - Self::build(true, 5 * 1024 * 1024, max_archive_bytes, ObjectFormat::SHA1) - } - - fn build( - rebuilt: bool, - max_response_bytes: usize, - max_archive_bytes: u64, - object_format: ObjectFormat, - ) -> Self { - Self::build_with_limits( - rebuilt, - max_response_bytes, - max_archive_bytes, - object_format, - LimitConfig::default(), - ) + fn build(rebuilt: bool, byte_limits: ByteLimits, object_format: ObjectFormat) -> Self { + Self::build_with_limits(rebuilt, byte_limits, object_format, LimitConfig::default()) } fn build_with_limits( rebuilt: bool, - max_response_bytes: usize, - max_archive_bytes: u64, + byte_limits: ByteLimits, object_format: ObjectFormat, limits: LimitConfig, ) -> Self { @@ -206,7 +179,7 @@ impl World { limiter: Arc::new(PreAuthLimiter::with_config(limits)), cob_locks: Arc::new(CobLocks::default()), reservations: Arc::new(Reservations::new( - 1_000_000, + knot_xrpc::ReservationTtl::new(1_000_000), PerActorQuota::new(16), GlobalQuota::new(16), )), @@ -215,29 +188,14 @@ impl World { name: AuthorName::new("Tangled"), email: Email::new("noreply@tangled.sh"), }, - max_body_bytes: BodyLimit::new(64 * 1024), - max_patch_bytes: PatchLimit::new(16 * 1024 * 1024), - max_patch_decompressed_bytes: PatchDecompressedLimit::new(128 * 1024 * 1024), - max_response_bytes: ResponseLimit::new(max_response_bytes), - max_archive_bytes: ArchiveLimit::new(max_archive_bytes), - tree_last_commit_budget: TreeReadBudget::new(ReadBudget::Within( - Duration::from_millis(300), - )), - blob_last_commit_budget: BlobReadBudget::new(ReadBudget::Within( - Duration::from_millis(2000), - )), - languages_budget: LanguagesReadBudget::new(ReadBudget::Within(Duration::from_millis( - 1000, - ))), - languages_push_budget: Duration::from_millis(2000), + byte_limits, + budgets: Budgets::default(), git_http: Arc::new(FakeHttp::new(|_request: &HttpRequest| { Err(NetworkError::Connect( "no git upstream is served in this test".to_string(), )) })), - fork_max_pack_bytes: ForkPackLimit::new(1024 * 1024 * 1024), pack_limits: knot_pack::PackLimits::default(), - max_pack_bytes: PackLimit::new(1024 * 1024 * 1024), service_owner: AccountDid::new(OWNER).unwrap(), events: Arc::new(knot_events::EventLog::new( ManualClock::new(UnixMicros::new(1_000_000_000)), @@ -249,8 +207,7 @@ impl World { )), maintenance: knot_maintenance::MaintenanceHandle::disabled(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), - resolve_slots: ResolveSlots::new(Arc::new(tokio::sync::Semaphore::new(8))), - receive_slots: ReceiveSlots::new(Arc::new(tokio::sync::Semaphore::new(8))), + slots: knot_resource::Slots::testing(8), lfs: Some(knot_xrpc::LfsWeb::new(lfs_handle, 8)), catalog: Arc::new(knot_messages::Catalog::defaults()), }); @@ -452,17 +409,12 @@ pub async fn post_json( } pub fn git_run(cwd: &Path, when: &str, author: (&str, &str), args: &[&str]) -> String { - let output = std::process::Command::new("git") + let output = knot_fixtures::command_at(cwd, when) .args(args) - .current_dir(cwd) - .env("GIT_CONFIG_GLOBAL", "/dev/null") - .env("GIT_CONFIG_SYSTEM", "/dev/null") .env("GIT_AUTHOR_NAME", author.0) .env("GIT_AUTHOR_EMAIL", author.1) - .env("GIT_AUTHOR_DATE", when) .env("GIT_COMMITTER_NAME", author.0) .env("GIT_COMMITTER_EMAIL", author.1) - .env("GIT_COMMITTER_DATE", when) .output() .expect("git is available"); assert!( diff --git a/crates/knot-xrpc/tests/reads.rs b/crates/knot-xrpc/tests/reads.rs index ff69941..cb9fc6b 100644 --- a/crates/knot-xrpc/tests/reads.rs +++ b/crates/knot-xrpc/tests/reads.rs @@ -11,6 +11,7 @@ use tokio_tungstenite::tungstenite; use knot_events::{EventCursor, GitRefUpdate}; use knot_types::{AccountDid, ObjectFormat, Oid, OwnerDid, RepoDid}; +use knot_xrpc::{ArchiveLimit, ResponseLimit}; use common::{ OWNER, World, archive_full, assert_immutable_round_trip, assert_post_rejected, assert_warming, @@ -1224,7 +1225,7 @@ async fn a_submodule_path_in_the_tree_is_path_not_found() { #[tokio::test] async fn a_blob_past_the_derived_serving_limit_is_a_named_error() { - let world = World::with_response_limit(1024); + let world = World::with_response_limit(ResponseLimit::new(1024)); let (did, bare, work) = empty_repo(&world, "auger"); commit_file( work.path(), @@ -1257,7 +1258,7 @@ async fn a_blob_past_the_derived_serving_limit_is_a_named_error() { #[tokio::test] async fn an_oversized_readme_is_omitted_from_the_tree() { - let world = World::with_response_limit(2_048); + let world = World::with_response_limit(ResponseLimit::new(2_048)); let (did, bare, work) = empty_repo(&world, "cowrie"); commit_file( work.path(), @@ -1352,7 +1353,7 @@ async fn archive_rejects_traversal_prefixes_and_sanitizes_the_filename() { #[tokio::test] async fn an_archive_larger_than_the_configured_limit_is_refused() { - let world = World::with_archive_limit(64); + let world = World::with_archive_limit(ArchiveLimit::new(64)); let (did, _work) = seeded(&world, "murex"); let (status, error) = get_error( @@ -1366,7 +1367,7 @@ async fn an_archive_larger_than_the_configured_limit_is_refused() { #[tokio::test] async fn an_oversized_read_response_is_refused() { - let world = World::with_response_limit(256); + let world = World::with_response_limit(ResponseLimit::new(256)); let (did, _work) = seeded(&world, "clam"); let (status, error) = get_error( -- 2.51.2