diff --git a/crates/knot-pack/src/lib.rs b/crates/knot-pack/src/lib.rs index 7099411..1abd02f 100644 --- a/crates/knot-pack/src/lib.rs +++ b/crates/knot-pack/src/lib.rs @@ -27,12 +27,13 @@ use axum::response::Response; use axum::routing::post; use knot_git::{Layout, Repo}; use knot_messages::{Catalog, ErrorKey, FetchMessages}; +use knot_resource::{PackSlots, SlotPermit}; use knot_runtime::Clock; use knot_types::{ AccountDid, Handle, KnotHostname, OwnerDid, OwnerRef, ParseError, RepoDid, RepoRkey, }; use std::sync::Arc; -use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc}; +use tokio::sync::mpsc; use tokio_stream::wrappers::ReceiverStream; pub use cache::{CacheConfig, MaxCacheBytes, MaxEntryBytes}; @@ -42,7 +43,7 @@ pub use frame::{ ReceiveFramer, UploadFramer, archive_request_complete, receive_request_complete, upload_v0_nak, }; pub use guard::PushGuard; -pub use ids::{DeltaDepth, MaxObjectBytes, MaxTotalBytes}; +pub use ids::{DeltaDepth, MaxObjectBytes, MaxTotalBytes, MaxWireBytes}; pub use meter::PackLimits; pub use objects::{ExpandedPack, count_expanded, write_expanded, write_pack}; pub use oids::{HaveOids, WantOids}; @@ -353,20 +354,25 @@ struct PackState { resolver: Arc, receive: Option>, handle_resolver: Option>, - pack_slots: Arc, + pack_slots: PackSlots, cache: Arc, catalog: Arc, hostname: KnotHostname, } pub fn router(layout: Layout, resolver: Arc, clock: Arc) -> Router { - router_with_pack_slots(layout, resolver, knot_resource::threads().get(), clock) + router_with_pack_slots( + layout, + resolver, + PackSlots::new(knot_resource::threads().get()), + clock, + ) } pub fn router_with_pack_slots( layout: Layout, resolver: Arc, - pack_slots: usize, + pack_slots: PackSlots, clock: Arc, ) -> Router { let state = pack_state( @@ -389,6 +395,7 @@ pub fn edge_routes( resolver: Arc, receive: Option>, handle_resolver: Option>, + pack_slots: PackSlots, cache: CacheConfig, catalog: Arc, hostname: KnotHostname, @@ -399,7 +406,7 @@ pub fn edge_routes( resolver, receive, handle_resolver, - knot_resource::threads().get(), + pack_slots, cache, catalog, hostname, @@ -414,7 +421,7 @@ fn pack_state( resolver: Arc, receive: Option>, handle_resolver: Option>, - pack_slots: usize, + pack_slots: PackSlots, cache: CacheConfig, catalog: Arc, hostname: KnotHostname, @@ -425,7 +432,7 @@ fn pack_state( resolver, receive, handle_resolver, - pack_slots: Arc::new(Semaphore::new(pack_slots.max(1))), + pack_slots, cache: cache::PackCache::new(cache, clock), catalog, hostname, @@ -676,10 +683,7 @@ async fn dispatch_plan( body: &[u8], lease: Option, ) -> Result { - let permit = Arc::clone(&state.pack_slots) - .acquire_owned() - .await - .expect("pack concurrency semaphore is never closed"); + let permit = state.pack_slots.acquire().await; let owned_body = body.to_vec(); let planned = tokio::task::spawn_blocking(move || { let outcome = upload::plan(&repo, &owned_body); @@ -743,7 +747,7 @@ fn stream_response( wants: WantOids, haves: HaveOids, opts: upload::StreamOpts, - permit: OwnedSemaphorePermit, + permit: SlotPermit, lease: Option, catalog: Arc, hostname: KnotHostname, @@ -845,14 +849,11 @@ async fn archive_dispatch( repo: Repo, body: Vec, ) -> Result { - let permit = Arc::clone(&state.pack_slots) - .acquire_owned() - .await - .expect("pack concurrency semaphore is never closed"); + let permit = state.pack_slots.acquire().await; Ok(archive_response(repo, body, permit)) } -fn archive_response(repo: Repo, body: Vec, permit: OwnedSemaphorePermit) -> Response { +fn archive_response(repo: Repo, body: Vec, permit: SlotPermit) -> Response { let (tx, rx) = mpsc::channel::>(16); tokio::task::spawn_blocking(move || { let _permit = permit; diff --git a/crates/knot-pack/src/quarantine.rs b/crates/knot-pack/src/quarantine.rs index 8c490a7..cf6d315 100644 --- a/crates/knot-pack/src/quarantine.rs +++ b/crates/knot-pack/src/quarantine.rs @@ -1,3 +1,16 @@ +//! The same way one would do this for uploading images onto a server, +//! we stage pushes such that objects unpack into a separate bare repo +//! under the incoming prefix, because we treat a push as real only once +//! every object it referenced actually made it over. +//! +//! Hence aborting is as simple as removal of the directory that +//! represents the push in flight. +//! +//! "Why not use `GIT_QUARANTINE_PATH`?" - because we never forked `receive-pack`. +//! +//! If you were wondering, this `sweep_incoming` is for a crash between a stage +//! and migration that would otherwise leak a staging dir. + use std::path::Path; use knot_git::{INCOMING_PREFIX, Repo, Staging}; diff --git a/crates/knot-receive/Cargo.toml b/crates/knot-receive/Cargo.toml index 9dce2a5..abcb05b 100644 --- a/crates/knot-receive/Cargo.toml +++ b/crates/knot-receive/Cargo.toml @@ -9,6 +9,7 @@ license.workspace = true knot-types = { workspace = true } knot-git = { workspace = true } knot-pack = { workspace = true } +knot-resource = { workspace = true } knot-cob = { workspace = true } knot-index = { workspace = true } knot-atproto = { workspace = true } diff --git a/crates/knot-receive/src/lib.rs b/crates/knot-receive/src/lib.rs index 9c6a5be..c3b8092 100644 --- a/crates/knot-receive/src/lib.rs +++ b/crates/knot-receive/src/lib.rs @@ -68,7 +68,6 @@ use std::cell::RefCell; use std::sync::Arc; -use std::time::Duration; use knot_atproto::Atproto; use knot_cob::CobHome; @@ -81,12 +80,12 @@ use knot_pack::{ PackError, PackLimits, PushGuard, ReceiveOutcome, ReceivedPack, frame_report, receive_pack_guarded_streamed, receive_preflight, }; -use knot_postreceive::{Actor, Ci, OwnerLabel, PullLink, post_receive}; +use knot_postreceive::{Actor, Ci, LanguagesPushBudget, OwnerLabel, PullLink, post_receive}; +use knot_resource::ResolveSlots; use knot_runtime::{Clock, HttpTransport}; use knot_types::{ AccountDid, ActorId, AppviewEndpoint, Handle, KnotHostname, OwnerDid, RepoDid, RepoRkey, }; -use tokio::sync::Semaphore; type Applied = Vec<(RefUpdate, Reservation)>; @@ -100,11 +99,11 @@ pub struct Push<'a, H: HttpTransport, C: Clock> { pub events: Arc>, pub index: &'a Index, pub atproto: &'a Atproto, - pub resolve_slots: &'a Semaphore, + pub resolve_slots: &'a ResolveSlots, pub appview: &'a AppviewEndpoint, pub maintenance: &'a MaintenanceHandle, pub hostname: &'a KnotHostname, - pub languages_push_budget: Duration, + pub languages_push_budget: LanguagesPushBudget, pub catalog: Arc, } @@ -292,7 +291,7 @@ pub fn ci_from_push_options(options: &[String], knot: &KnotHostname, rkey: Optio async fn resolve_pull_link( index: &Index, atproto: &Atproto, - resolve_slots: &Semaphore, + resolve_slots: &ResolveSlots, appview: &AppviewEndpoint, repo_did: &RepoDid, ) -> Option { @@ -313,7 +312,7 @@ async fn resolve_pull_link( async fn resolve_owner_label( atproto: &Atproto, - resolve_slots: &Semaphore, + resolve_slots: &ResolveSlots, owner: &OwnerDid, ) -> OwnerLabel { let did = AccountDid::from(owner.clone()); @@ -325,10 +324,10 @@ async fn resolve_owner_label( pub async fn resolve_handle( atproto: &Atproto, - resolve_slots: &Semaphore, + resolve_slots: &ResolveSlots, did: &AccountDid, ) -> Option { - let _permit = resolve_slots.acquire().await.ok()?; + let _permit = resolve_slots.try_acquire()?; let identity = atproto.resolve_identity(did).await.ok()?; identity.primary_handle().cloned() } diff --git a/crates/knot-resource/src/lib.rs b/crates/knot-resource/src/lib.rs index 300fe50..51ed432 100644 --- a/crates/knot-resource/src/lib.rs +++ b/crates/knot-resource/src/lib.rs @@ -1,12 +1,23 @@ +mod admission; mod cpu; mod disk; +mod fsio; mod mem; +mod slots; +pub use admission::{ + AdmitGuard, Burst, GlobalInflight, LimitConfig, PerPeerInflight, PreAuthLimiter, RateLimit, + RefillMicros, Refusal, +}; pub use cpu::{Saturate, ThreadCount, gix_thread_limit, map_chunks, map_spans, saturate, threads}; pub use disk::{ DiskFloorBytes, DiskGovernor, DiskReservation, FreeBytes, ReserveBytes, ReserveError, free_bytes as disk_free_bytes, }; +pub use fsio::{ + FileMode, FsError, atomic_write, atomic_write_bytes, clear_staging, clear_stale, clear_temps, + fsync_path, staging_nonce, +}; pub use mem::{ AvailableBytes, BudgetSource, ChurnBytes, ConnectivityObjects, DecayMs, MemoryBudget, MemoryHighBytes, PayloadBytes, WorkingSetBytes, advert_cache_bytes, available_bytes, @@ -14,6 +25,7 @@ pub use mem::{ ingest_admits_churn, ingest_base_budget, ingest_thread_limit, object_cache_bytes, pack_cache_bytes, target_decay, }; +pub use slots::{PackSlots, ReceiveSlots, ResolveSlots, SlotPermit, Slots}; #[derive(Clone, Copy, Debug, Default)] pub struct Ceilings { diff --git a/crates/knot-resource/src/slots.rs b/crates/knot-resource/src/slots.rs new file mode 100644 index 0000000..410078c --- /dev/null +++ b/crates/knot-resource/src/slots.rs @@ -0,0 +1,135 @@ +use std::sync::{Arc, OnceLock}; + +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; + +use crate::cpu::threads; + +const RESOLVE_SLOTS_PER_TRANSPORT: usize = 16; + +const SHARING_TRANSPORTS: usize = 2; + +const RESOLVE_SLOTS: usize = RESOLVE_SLOTS_PER_TRANSPORT * SHARING_TRANSPORTS; + +static PROCESS: OnceLock = OnceLock::new(); + +pub struct SlotPermit(#[allow(dead_code)] OwnedSemaphorePermit); + +macro_rules! slot_kind { + ($($name:ident),+ $(,)?) => {$( + #[derive(Clone)] + pub struct $name(Arc); + + impl $name { + pub fn new(permits: usize) -> Self { + Self(Arc::new(Semaphore::new(permits.max(1)))) + } + + pub async fn acquire(&self) -> SlotPermit { + SlotPermit( + Arc::clone(&self.0) + .acquire_owned() + .await + .expect("a slot budget closes only with the process that owns it"), + ) + } + + pub fn available(&self) -> usize { + self.0.available_permits() + } + } + )+}; +} + +slot_kind!(ResolveSlots, ReceiveSlots, PackSlots); + +impl ResolveSlots { + pub fn try_acquire(&self) -> Option { + Arc::clone(&self.0).try_acquire_owned().ok().map(SlotPermit) + } +} + +#[derive(Clone)] +pub struct Slots { + pub resolve: ResolveSlots, + pub receive: ReceiveSlots, + pub pack: PackSlots, +} + +impl Slots { + pub fn for_machine() -> Self { + PROCESS + .get_or_init(|| Self { + resolve: ResolveSlots::new(RESOLVE_SLOTS), + receive: ReceiveSlots::new(threads().get()), + pack: PackSlots::new(threads().get()), + }) + .clone() + } + + pub fn testing(permits: usize) -> Self { + Self { + resolve: ResolveSlots::new(permits), + receive: ReceiveSlots::new(permits), + pack: PackSlots::new(permits), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn a_clone_shares_one_pool_per_kind_of_work() { + let slots = Slots::testing(1); + let other = slots.clone(); + let held = slots.receive.acquire().await; + assert_eq!( + other.receive.available(), + 0, + "a cloned budget mustn't grant a second permit for the one slot" + ); + assert_eq!( + other.pack.available(), + 1, + "spending a receive slot mustn't spend the pack budget" + ); + let _resolving = slots.resolve.acquire().await; + assert!( + other.resolve.try_acquire().is_none(), + "a cosmetic lookup mustn't wait, \ + or a push queues on it while its receive and pack slots stay spent" + ); + drop(held); + assert_eq!(other.receive.available(), 1); + } + + #[tokio::test] + async fn every_caller_of_for_machine_shares_one_budget_that_no_test_budget_touches() { + let ssh = Slots::for_machine(); + let http = Slots::for_machine(); + let isolated = Slots::testing(1); + let before = http.receive.available(); + let _held = ssh.receive.acquire().await; + assert_eq!( + http.receive.available(), + before - 1, + "two transports asking the machine for a budget must get the same one, \ + or the process grants twice the concurrency it was configured for" + ); + let resolving: Vec = (0..RESOLVE_SLOTS_PER_TRANSPORT) + .filter_map(|_| ssh.resolve.try_acquire()) + .collect(); + assert_eq!(resolving.len(), RESOLVE_SLOTS_PER_TRANSPORT); + assert!( + http.resolve.try_acquire().is_some(), + "collapsing a per-transport pool into a process-wide one mustn't give a deployment \ + running both transports less outbound resolution than either had on its own" + ); + assert_eq!( + isolated.receive.available(), + 1, + "whatever the process budget is doing mustn't spend a test budget" + ); + } +} diff --git a/crates/knot-ssh/src/exec.rs b/crates/knot-ssh/src/exec.rs index 121e3a6..ac89bda 100644 --- a/crates/knot-ssh/src/exec.rs +++ b/crates/knot-ssh/src/exec.rs @@ -138,13 +138,16 @@ pub(crate) async fn run_exec( .map_or(&state.peer_slots, |lfs| &lfs.peer_slots), _ => &state.peer_slots, }; - let Some(_peer_guard) = peer_limiter.enter(peer) else { - tracing::warn!( - ?peer, - reason = "peer concurrency limit reached", - "ssh exec rejected" - ); - return fail(channel, &state.catalog.ssh.too_many_operations.text()).await; + let _peer_guard = match peer_limiter.admit(peer, state.atproto.now()) { + Ok(guard) => guard, + Err(refusal) => { + let reason = match refusal { + knot_resource::Refusal::RateLimited => "peer request rate exceeded", + knot_resource::Refusal::Saturated => "peer concurrency limit reached", + }; + tracing::warn!(?peer, reason, "ssh exec rejected"); + return fail(channel, &state.catalog.ssh.too_many_operations.text()).await; + } }; let resolved_ref = match repo_ref { RepoRef::Did(did) => ResolvedRef::Did(did), @@ -417,10 +420,7 @@ async fn serve_upload_archive( Err(_) => return fail(channel, &state.catalog.ssh.archive_timeout.text()).await, }; - let permit = match Arc::clone(&state.pack_slots).acquire_owned().await { - Ok(permit) => permit, - Err(_) => return fail(channel, &state.catalog.ssh.shutting_down.text()).await, - }; + let permit = state.slots.pack.acquire().await; let (tx, mut rx) = mpsc::channel::>(16); let layout = state.layout.clone(); let did = repo_did.clone(); @@ -593,10 +593,7 @@ where C: Clock, W: AsyncWriteExt + Unpin, { - let permit = Arc::clone(&state.pack_slots) - .acquire_owned() - .await - .map_err(|_| ())?; + let permit = state.slots.pack.acquire().await; let (tx, mut rx) = mpsc::channel::>(16); let layout = state.layout.clone(); let did = repo_did.clone(); @@ -677,10 +674,7 @@ async fn serve_receive( } }; - let _receive_permit = match Arc::clone(&state.receive_slots).acquire_owned().await { - Ok(permit) => permit, - Err(_) => return fail(channel, &state.catalog.ssh.shutting_down.text()).await, - }; + let _receive_permit = state.slots.receive.acquire().await; let limits = state.limits; let body = { @@ -726,10 +720,7 @@ async fn serve_receive( return finish(channel, 0).await; } - let _pack_permit = match Arc::clone(&state.pack_slots).acquire_owned().await { - Ok(permit) => permit, - Err(_) => return fail(channel, &state.catalog.ssh.shutting_down.text()).await, - }; + let _pack_permit = state.slots.pack.acquire().await; let landed = knot_receive::land(knot_receive::Push { layout: &state.layout, repo_did: &repo_did, @@ -740,7 +731,7 @@ async fn serve_receive( events: Arc::clone(&state.events), index: &state.index, atproto: &state.atproto, - resolve_slots: &state.resolve_slots, + resolve_slots: &state.slots.resolve, appview: &state.appview, maintenance: &state.maintenance, hostname: &state.hostname, @@ -788,7 +779,7 @@ async fn greeting_identity( let Some(did) = key.and_then(|key| state.roster.did_for(key)) else { return "there".to_string(); }; - match knot_receive::resolve_handle(&state.atproto, &state.resolve_slots, &did).await { + match knot_receive::resolve_handle(&state.atproto, &state.slots.resolve, &did).await { Some(handle) => format!("@{}", handle.as_str()), None => did.as_str().to_string(), } @@ -819,7 +810,7 @@ async fn resolve_pusher( { return Some(cached); } - let _permit = state.resolve_slots.acquire().await.ok()?; + let _permit = state.slots.resolve.acquire().await; let matches = futures::stream::iter(candidates).filter_map(|did| async move { let keys = state.atproto.resolve_pubkeys(&did).await.ok()?; keys.iter() @@ -850,7 +841,7 @@ async fn read_chunk( async fn read_receive( reader: &mut R, dir: PathBuf, - limit: usize, + limit: knot_pack::MaxWireBytes, limits: PackLimits, format: ObjectFormat, ) -> Result { @@ -907,7 +898,7 @@ fn read_error(error: knot_pack::ReceiveReadError) -> ReadError { fn frame_receive( mut rx: mpsc::Receiver>, dir: PathBuf, - limit: usize, + limit: knot_pack::MaxWireBytes, limits: PackLimits, format: ObjectFormat, ) -> Result { diff --git a/crates/knot-ssh/src/lib.rs b/crates/knot-ssh/src/lib.rs index d75cfd7..643ed9a 100644 --- a/crates/knot-ssh/src/lib.rs +++ b/crates/knot-ssh/src/lib.rs @@ -1,5 +1,4 @@ mod exec; -mod limit; mod roster; mod server; @@ -14,7 +13,8 @@ use knot_events::EventLog; use knot_git::Layout; use knot_index::Index; use knot_maintenance::MaintenanceHandle; -use knot_pack::PackLimits; +use knot_pack::{MaxWireBytes, PackLimits}; +use knot_postreceive::LanguagesPushBudget; use knot_runtime::{Clock, Entropy, HttpTransport, OsEntropy}; use knot_types::{AccountDid, ActorId, AdmissionPolicy, AppviewEndpoint, KnotHostname}; use russh::keys::ssh_key::rand_core; @@ -25,11 +25,10 @@ use tokio::sync::Semaphore; use tokio_util::sync::CancellationToken; use tokio_util::task::TaskTracker; -use limit::PeerLimiter; +use knot_resource::{LimitConfig, PerPeerInflight, PreAuthLimiter, Slots}; use roster::KeyRoster; use server::KnotSshServer; -const MAX_CONCURRENT_RESOLUTIONS: usize = 16; const MAX_INFLIGHT_PER_PEER: usize = 4; const INACTIVITY_TIMEOUT: Duration = Duration::from_secs(120); const KEEPALIVE_INTERVAL: Duration = Duration::from_secs(30); @@ -56,12 +55,10 @@ pub struct SshState { admins: BTreeSet, admission: AdmissionPolicy, limits: PackLimits, - max_pack_bytes: usize, - languages_push_budget: Duration, - resolve_slots: Arc, - pack_slots: Arc, - receive_slots: Arc, - peer_slots: Arc, + max_pack_bytes: MaxWireBytes, + languages_push_budget: LanguagesPushBudget, + slots: Slots, + peer_slots: Arc, roster: Arc, maintenance: MaintenanceHandle, lfs: Option, @@ -72,7 +69,7 @@ pub struct SshState { pub(crate) struct LfsRuntime { pub(crate) handle: knot_lfs::LfsHandle, pub(crate) slots: Arc, - pub(crate) peer_slots: Arc, + pub(crate) peer_slots: Arc, } impl SshState { @@ -87,8 +84,8 @@ impl SshState { appview: AppviewEndpoint, admins: BTreeSet, admission: AdmissionPolicy, - max_pack_bytes: usize, - languages_push_budget: Duration, + max_pack_bytes: MaxWireBytes, + languages_push_budget: LanguagesPushBudget, ) -> Self { Self { layout, @@ -103,10 +100,10 @@ impl SshState { limits: PackLimits::default(), max_pack_bytes, languages_push_budget, - resolve_slots: Arc::new(Semaphore::new(MAX_CONCURRENT_RESOLUTIONS)), - pack_slots: Arc::new(Semaphore::new(knot_resource::threads().get())), - receive_slots: Arc::new(Semaphore::new(knot_resource::threads().get())), - peer_slots: Arc::new(PeerLimiter::new(MAX_INFLIGHT_PER_PEER)), + slots: Slots::for_machine(), + peer_slots: Arc::new(PreAuthLimiter::with_config(LimitConfig::per_peer_only( + PerPeerInflight::new(MAX_INFLIGHT_PER_PEER), + ))), roster: Arc::new(KeyRoster::new()), maintenance: MaintenanceHandle::disabled(), lfs: None, @@ -119,6 +116,11 @@ impl SshState { self } + pub fn with_slots(mut self, slots: Slots) -> Self { + self.slots = slots; + self + } + pub fn with_maintenance(mut self, maintenance: MaintenanceHandle) -> Self { self.maintenance = maintenance; self @@ -128,7 +130,9 @@ impl SshState { self.lfs = Some(LfsRuntime { handle, slots: Arc::new(Semaphore::new(max_transfers)), - peer_slots: Arc::new(PeerLimiter::new(max_transfers)), + peer_slots: Arc::new(PreAuthLimiter::with_config(LimitConfig::per_peer_only( + PerPeerInflight::new(max_transfers), + ))), }); self }