diff --git a/.dockerignore b/.dockerignore index 6887b922a..d5505e76d 100644 --- a/.dockerignore +++ b/.dockerignore @@ -10,7 +10,8 @@ genjwks.out blog/build build/ .wrangler/ -localinfra/certs/root.key +localinfra/certs/*.key +localinfra/ncps-secret-key appview/pages/static/tw.css appview/pages/static/x diff --git a/Cargo.lock b/Cargo.lock index 1911f4e8d..c421d9e6d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -929,6 +929,8 @@ dependencies = [ "http", "jacquard-common", "reqwest 0.13.1", + "rustls", + "rustls-native-certs", "thiserror 2.0.18", "tokio", "tokio-tungstenite 0.29.0", @@ -7474,7 +7476,6 @@ dependencies = [ "wasm-bindgen-futures", "wasm-streams", "web-sys", - "webpki-roots 1.0.8", ] [[package]] @@ -9051,11 +9052,11 @@ dependencies = [ "futures-util", "log", "rustls", + "rustls-native-certs", "rustls-pki-types", "tokio", "tokio-rustls", "tungstenite 0.29.0", - "webpki-roots 0.26.11", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 85e072461..57ca3edf3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -118,7 +118,7 @@ walkdir = "2" tokio = { version = "1.52", features = ["macros", "rt-multi-thread", "time", "signal", "io-util", "net", "sync"] } tokio-util = { version = "0.7", features = ["rt"] } tokio-stream = "0.1" -tokio-tungstenite = { version = "0.29", features = ["rustls-tls-webpki-roots"] } +tokio-tungstenite = { version = "0.29", features = ["rustls-tls-native-roots"] } futures = "0.3" either = "1" async-trait = "0.1" @@ -147,7 +147,7 @@ getrandom = "0.4" ahash = { version = "0.8", default-features = false, features = ["std"] } arc-swap = "1" -reqwest = { version = "0.13", default-features = false, features = ["rustls", "webpki-roots", "http2", "json", "gzip", "stream"] } +reqwest = { version = "0.13", default-features = false, features = ["rustls", "http2", "json", "gzip", "stream"] } axum = { version = "0.8", features = ["macros"] } hyper = { version = "1", features = ["server", "http1", "http2"] } hyper-util = { version = "0.1", features = ["server", "server-auto", "tokio", "service"] } @@ -163,6 +163,7 @@ ipnet = "2.10" rustls = { version = "0.23", features = ["aws_lc_rs", "prefer-post-quantum"] } tokio-rustls = { version = "0.26", default-features = false, features = ["aws_lc_rs", "tls12", "logging"] } rustls-pemfile = "2" +rustls-native-certs = "0.8" rustls-acme = { version = "0.15.3", default-features = false, features = ["aws-lc-rs", "tls12", "webpki-roots"] } x509-parser = "0.18" quinn = { version = "0.11.9", default-features = false, features = ["runtime-tokio", "rustls-aws-lc-rs", "log"] } diff --git a/bobbin/containerfiles/bobbin.Containerfile b/bobbin/containerfiles/bobbin.Containerfile index 4b800850a..f38a0b7b8 100644 --- a/bobbin/containerfiles/bobbin.Containerfile +++ b/bobbin/containerfiles/bobbin.Containerfile @@ -15,9 +15,7 @@ ARG BOBBIN_PROFILE=release WORKDIR /src COPY Cargo.toml Cargo.lock rust-toolchain.toml ./ COPY lexicons ./lexicons -COPY crates/knot-capability ./crates/knot-capability -COPY crates/lexicons ./crates/lexicons -COPY crates/trusted-proxies ./crates/trusted-proxies +COPY crates ./crates COPY bobbin ./bobbin # keep image lighter by not copying in shuttle, gitmirror, knot2 # does however need some cargo toml patchings diff --git a/bobbin/crates/bobbin-sim/src/runtime.rs b/bobbin/crates/bobbin-sim/src/runtime.rs index 2fb41c52a..ecaf6da98 100644 --- a/bobbin/crates/bobbin-sim/src/runtime.rs +++ b/bobbin/crates/bobbin-sim/src/runtime.rs @@ -5,8 +5,8 @@ use std::time::Duration; use bobbin_edge_index::{CoverageWatch, EdgeStore}; use bobbin_ingest::{ - DEFAULT_INGEST_PARALLELISM, DisconnectSink, IngestConfig, IngestRuntime, RepoIdResolver, - WarmingBuffer, WarmingShadowBuffer, run as run_ingest, + DEFAULT_IDLE_PROMOTE_MIN_EVENTS, DEFAULT_INGEST_PARALLELISM, DisconnectSink, IngestConfig, + IngestRuntime, RepoIdResolver, WarmingBuffer, WarmingShadowBuffer, run as run_ingest, }; use bobbin_record_lru::NoopRecordStore; use bobbin_resolver::IdentityResolver; @@ -132,6 +132,7 @@ impl Sim { hydrant_base, start_cursor: bobbin_edge_index::HydrantCursor::new(0), parallelism, + idle_promote_min_events: DEFAULT_IDLE_PROMOTE_MIN_EVENTS, }; let mut ingest_handle = tokio::spawn(async move { let _ = run_ingest(ingest_config, ingest_runtime).await; diff --git a/bobbin/crates/bobbin/src/config.rs b/bobbin/crates/bobbin/src/config.rs index 2c1da7b5f..8b0ed6143 100644 --- a/bobbin/crates/bobbin/src/config.rs +++ b/bobbin/crates/bobbin/src/config.rs @@ -20,6 +20,7 @@ const KNOWN_KEYS: &[&str] = &[ "hydrant.url", "hydrant.start_cursor", "ingest.parallelism", + "ingest.idle_promote_min_events", "backpressure.per_request_anon_bytes", "backpressure.adjust_interval_ms", "backpressure.relieve_below_ratio", @@ -47,6 +48,7 @@ const KNOWN_ENVS: &[&str] = &[ "BOBBIN_HYDRANT_URL", "BOBBIN_START_CURSOR", "BOBBIN_INGEST_PARALLELISM", + "BOBBIN_INGEST_IDLE_PROMOTE_MIN_EVENTS", "BOBBIN_BACKPRESSURE_PER_REQUEST_ANON_BYTES", "BOBBIN_BACKPRESSURE_ADJUST_INTERVAL_MS", "BOBBIN_BACKPRESSURE_RELIEVE_BELOW_RATIO", @@ -197,6 +199,12 @@ pub struct IngestConfig { /// serial so that like, cursor and Spur allocation order are preserved. #[config(env = "BOBBIN_INGEST_PARALLELISM", default = 16)] pub parallelism: usize, + + /// Events processed before a quiet stream may promote coverage to ready. + /// The invite listings will be `503 Warming` while a quiet stream is under + /// said count. + #[config(env = "BOBBIN_INGEST_IDLE_PROMOTE_MIN_EVENTS", default = 256)] + pub idle_promote_min_events: u64, } #[derive(Debug, Config)] diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index ebe99c940..ef0243333 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -1,5 +1,5 @@ use std::net::SocketAddr; -use std::num::NonZeroUsize; +use std::num::{NonZeroU64, NonZeroUsize}; use std::path::PathBuf; use std::process::ExitCode; use std::sync::Arc; @@ -18,7 +18,7 @@ use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; use bobbin_resolver::{HydrantClient, IdentityResolver}; use bobbin_runtime::{ Clock, GuardedWs, MemoryBudget, NetworkError, OsEntropy, ReqwestHttp, RuntimeHasher, - SystemClock, TungsteniteWs, WsTransport, + SystemClock, TungsteniteWs, WsTls, WsTransport, }; use bobbin_search::{ActorIndex, SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; @@ -184,7 +184,8 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { let entropy = Arc::new(OsEntropy); let hasher = RuntimeHasher::from_entropy(&*entropy); - let ws = TungsteniteWs::shared(); + let ws_tls = WsTls::from_native_roots()?; + let ws = TungsteniteWs::shared(ws_tls.clone()); let records: Arc = Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(lru_cap))); @@ -292,10 +293,13 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { effective = parallelism.get(), "ingest parallelism" ); + let idle_promote_min_events = NonZeroU64::new(cfg.ingest.idle_promote_min_events) + .ok_or_else(|| anyhow!("ingest.idle_promote_min_events must be at least 1"))?; let ingest_cfg = IngestConfig { hydrant_base: cfg.hydrant.url.clone(), start_cursor: HydrantCursor::new(cfg.hydrant.start_cursor), parallelism, + idle_promote_min_events, }; let cancel = CancellationToken::new(); let actor_rebuilder = actors.clone(); @@ -332,14 +336,17 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { let knot_ws: Arc = if knot_allow_private { ws.clone() } else { - GuardedWs::shared(Arc::new(|addrs: &[SocketAddr]| { - match addrs.iter().find_map(|sa| classify_ip(&sa.ip())) { - Some(reason) => Err(NetworkError::Connect(format!( - "knot firehose resolves to {reason} address space" - ))), - None => Ok(()), - } - })) + GuardedWs::shared( + Arc::new(|addrs: &[SocketAddr]| { + match addrs.iter().find_map(|sa| classify_ip(&sa.ip())) { + Some(reason) => Err(NetworkError::Connect(format!( + "knot firehose resolves to {reason} address space" + ))), + None => Ok(()), + } + }), + ws_tls, + ) }; let ingest_runtime = IngestRuntime { diff --git a/bobbin/crates/ingest/examples/smoke.rs b/bobbin/crates/ingest/examples/smoke.rs index 28b02d2bc..c8afc6395 100644 --- a/bobbin/crates/ingest/examples/smoke.rs +++ b/bobbin/crates/ingest/examples/smoke.rs @@ -5,7 +5,7 @@ use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex}; use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run}; use bobbin_record_lru::{NoopRecordStore, RecordStore}; use bobbin_resolver::IdentityResolver; -use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; +use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs, WsTls}; use bobbin_types::search::NoopSearchSink; use futures::stream::{self, StreamExt}; use tokio_util::sync::CancellationToken; @@ -43,7 +43,9 @@ async fn main() { identity: Arc::new(IdentityResolver::detached(hasher)), clock: Arc::new(SystemClock::new()), entropy: Arc::new(OsEntropy), - ws: TungsteniteWs::shared(), + ws: TungsteniteWs::shared( + WsTls::from_native_roots().expect("a system trust store for wss certificates"), + ), cancel: cancel.clone(), disconnects: None, warming_shadow: None, diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index b6e381705..0458596bf 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -1,6 +1,7 @@ use std::collections::{HashSet, VecDeque}; -use std::num::NonZeroUsize; +use std::num::{NonZeroU64, NonZeroUsize}; use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; use bobbin_edge_index::{ @@ -66,12 +67,17 @@ pub const DEFAULT_INGEST_PARALLELISM: NonZeroUsize = match NonZeroUsize::new(16) Some(n) => n, None => unreachable!(), }; +pub const DEFAULT_IDLE_PROMOTE_MIN_EVENTS: NonZeroU64 = match NonZeroU64::new(256) { + Some(n) => n, + None => unreachable!(), +}; #[derive(Clone, Debug)] pub struct IngestConfig { pub hydrant_base: Url, pub start_cursor: HydrantCursor, pub parallelism: NonZeroUsize, + pub idle_promote_min_events: NonZeroU64, } impl IngestConfig { @@ -80,6 +86,7 @@ impl IngestConfig { hydrant_base, start_cursor: HydrantCursor::new(0), parallelism: DEFAULT_INGEST_PARALLELISM, + idle_promote_min_events: DEFAULT_IDLE_PROMOTE_MIN_EVENTS, } } @@ -293,16 +300,13 @@ pub async fn run( config: IngestConfig, runtime: IngestRuntime, ) -> Result<(), IngestError> { + let frames = FrameCounter::default(); let metrics_dumper = spawn_metrics_dumper(&runtime); - let idle_promoter = spawn_idle_promoter(&runtime); let warming_flusher = spawn_warming_flusher(&runtime); - let result = run_inner(config, &runtime).await; + let result = run_inner(config, &runtime, &frames).await; if let Err(join) = metrics_dumper.await { warn!(?join, "metrics dumper task panicked"); } - if let Err(join) = idle_promoter.await { - warn!(?join, "idle promoter task panicked"); - } if let Some(handle) = warming_flusher && let Err(join) = handle.await { @@ -314,11 +318,14 @@ pub async fn run( async fn run_inner( config: IngestConfig, runtime: &IngestRuntime, + frames: &FrameCounter, ) -> Result<(), IngestError> { let mut backoff = RECONNECT_INITIAL_DELAY; loop { let cursor = next_connect_cursor(runtime.coverage.snapshot(), config.start_cursor); - let SessionEnd { outcome, error } = run_session(&config, cursor, runtime).await; + let opened = runtime.clock.now_instant(); + let SessionEnd { outcome, error } = run_session(&config, cursor, runtime, frames).await; + let open_for = runtime.clock.now_instant().duration_since(opened); if runtime.cancel.is_cancelled() { info!( last_cursor = runtime.coverage.snapshot().last_cursor().raw(), @@ -327,13 +334,19 @@ async fn run_inner( return Ok(()); } match (outcome, &error) { - (SessionOutcome::Progressed, None) => { - info!("hydrant stream closed after delivering frames, reconnecting") - } - (SessionOutcome::Empty, None) => { - warn!("hydrant stream closed without delivering frames") - } - (_, Some(err)) => warn!(?err, "hydrant stream errored"), + (SessionOutcome::Progressed, None) => info!( + open_secs = open_for.as_secs(), + "hydrant stream closed after delivering frames, reconnecting" + ), + (SessionOutcome::Empty, None) => warn!( + open_secs = open_for.as_secs(), + "hydrant stream closed without delivering frames" + ), + (_, Some(err)) => warn!( + ?err, + open_secs = open_for.as_secs(), + "hydrant stream errored" + ), } if let (Some(sink), Some(err)) = (runtime.disconnects.as_ref(), error.as_ref()) { sink.record(DisconnectSnapshot { @@ -453,39 +466,39 @@ async fn flush_warming_buffer( ); } -const IDLE_PROMOTE_WINDOW: Duration = Duration::from_secs(15); -const IDLE_PROMOTE_MIN_EVENTS: u64 = 256; +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +struct FramesDelivered(u64); + +#[derive(Clone, Debug, Default)] +struct FrameCounter(Arc); + +impl FrameCounter { + fn delivered(&self) -> FramesDelivered { + FramesDelivered(self.0.load(Ordering::Acquire)) + } + + fn deliver(&self) { + self.0.fetch_add(1, Ordering::Release); + } +} -fn spawn_idle_promoter( +fn promote_if_caught_up( runtime: &IngestRuntime, -) -> tokio::task::JoinHandle<()> { - let rt = runtime.clone(); - tokio::spawn(async move { - let mut prev = rt.coverage.snapshot().events_processed(); - loop { - tokio::select! { - biased; - _ = rt.cancel.cancelled() => return, - _ = rt.clock.sleep(IDLE_PROMOTE_WINDOW) => {} - } - let snap = rt.coverage.snapshot(); - if snap.is_ready() { - return; - } - let processed = snap.events_processed(); - if processed >= IDLE_PROMOTE_MIN_EVENTS && processed == prev { - rt.coverage.update(|c| c.force_ready()); - info!( - target: "bobbin_ingest::coverage", - events_processed = processed, - last_cursor = snap.last_cursor().raw(), - "stream idle, promoting coverage to ready", - ); - return; - } - prev = processed; - } - }) + min_events: NonZeroU64, + delivered: FramesDelivered, +) { + let snap = runtime.coverage.snapshot(); + if snap.is_ready() || snap.events_processed() < min_events.get() { + return; + } + runtime.coverage.update(|c| c.force_ready()); + info!( + target: "bobbin_ingest::coverage", + events_processed = snap.events_processed(), + frames_delivered = delivered.0, + last_cursor = snap.last_cursor().raw(), + "a live socket didn't deliver a frame across a keepalive interval, promoting coverage to ready", + ); } fn spawn_metrics_dumper( @@ -536,6 +549,7 @@ async fn run_session( config: &IngestConfig, cursor: HydrantCursor, runtime: &IngestRuntime, + frames: &FrameCounter, ) -> SessionEnd { let url = match config.stream_url(cursor) { Ok(u) => u, @@ -596,10 +610,17 @@ async fn run_session( let (control_tx, mut control_rx) = tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); let session_cancel = runtime.cancel.child_token(); let reader_cancel = session_cancel.clone(); - let reader = tokio::spawn(reader_loop(ws_stream, frame_tx, control_tx, reader_cancel)); + let reader = tokio::spawn(reader_loop( + ws_stream, + frame_tx, + control_tx, + frames.clone(), + reader_cancel, + )); let mut next_ping = runtime.clock.now_instant() + PING_INTERVAL; let mut pong_deadline: Option = None; + let mut quiet: Option = None; let writer_error: Option = loop { tokio::select! { @@ -614,6 +635,10 @@ async fn run_session( _ = runtime.clock.sleep_until(next_ping) => { next_ping = runtime.clock.now_instant() + PING_INTERVAL; if pong_deadline.is_none() { + let delivered = frames.delivered(); + if quiet.replace(delivered) == Some(delivered) { + promote_if_caught_up(runtime, config.idle_promote_min_events, delivered); + } if let Err(e) = timed_send(&mut ws_sink, WsMessage::Ping(Bytes::new())).await { break Some(e); } @@ -671,6 +696,7 @@ async fn reader_loop( mut ws_stream: Box, frame_tx: tokio::sync::mpsc::Sender, control_tx: tokio::sync::mpsc::Sender, + frames: FrameCounter, cancel: CancellationToken, ) -> SessionEnd { let mut outcome = SessionOutcome::Empty; @@ -736,6 +762,7 @@ async fn reader_loop( Ok(f) => f, Err(e) => break Some(e), }; + frames.deliver(); held.push_back(frame); } WsMessage::Binary(_) => { @@ -1743,7 +1770,7 @@ mod tests { use super::*; use bobbin_edge_index::Coverage; use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; - use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; + use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, WsConnectFuture, WsTransport}; use bobbin_types::search::NoopSearchSink; use jacquard_common::types::handle::Handle; use jacquard_common::types::nsid::Nsid; @@ -1752,6 +1779,64 @@ mod tests { const VALID_CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; + struct ScriptedWs(std::sync::Mutex>); + + impl ScriptedWs { + fn undialed() -> Arc { + Arc::new(Self(std::sync::Mutex::new(None))) + } + + fn dialing(conn: WsConn) -> Arc { + Arc::new(Self(std::sync::Mutex::new(Some(conn)))) + } + } + + impl WsTransport for ScriptedWs { + fn connect(&self, url: Url) -> WsConnectFuture { + let conn = self.0.lock().expect("scripted ws lock").take(); + Box::pin(async move { + conn.ok_or_else(|| { + NetworkError::Connect(format!("no scripted stream left for {url}")) + }) + }) + } + } + + enum Backlog { + Drained, + Delivering(u64), + } + + struct PongingSink { + ws: tokio::sync::mpsc::Sender>, + backlog: Backlog, + } + + impl bobbin_runtime::WsSink for PongingSink { + fn send<'a>(&'a mut self, msg: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { + Box::pin(async move { + let WsMessage::Ping(payload) = msg else { + return Ok(()); + }; + let next = match &mut self.backlog { + Backlog::Drained => None, + Backlog::Delivering(id) => { + *id += 1; + Some(*id) + } + }; + let closed = || NetworkError::Connect("scripted stream closed".to_owned()); + if let Some(id) = next { + self.ws.send(other_frame(id)).await.map_err(|_| closed())?; + } + self.ws + .send(Ok(WsMessage::Pong(payload))) + .await + .map_err(|_| closed()) + }) + } + } + fn did_subj(s: &str) -> SubjectRef { SubjectRef::Did(Did::new_owned(s).unwrap()) } @@ -2411,7 +2496,14 @@ mod tests { tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); let cancel = CancellationToken::new(); - let end = reader_loop(stream, frame_tx, control_tx, cancel).await; + let end = reader_loop( + stream, + frame_tx, + control_tx, + FrameCounter::default(), + cancel, + ) + .await; match end.error { Some(IngestError::ConsumerTooSlow { message }) => { @@ -2434,7 +2526,14 @@ mod tests { tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); let cancel = CancellationToken::new(); - let end = reader_loop(stream, frame_tx, control_tx, cancel).await; + let end = reader_loop( + stream, + frame_tx, + control_tx, + FrameCounter::default(), + cancel, + ) + .await; match end.error { Some(IngestError::HydrantStream { code, message }) => { @@ -2510,6 +2609,7 @@ mod tests { stream, frame_tx.clone(), control_tx.clone(), + FrameCounter::default(), cancel.clone(), )); @@ -2580,6 +2680,7 @@ mod tests { stream, frame_tx.clone(), control_tx, + FrameCounter::default(), cancel.clone(), )); @@ -2947,7 +3048,7 @@ mod tests { identity: Arc::new(IdentityResolver::detached(RuntimeHasher::default())), clock: Arc::new(SystemClock::new()), entropy: Arc::new(OsEntropy), - ws: TungsteniteWs::shared(), + ws: ScriptedWs::undialed(), cancel, disconnects: None, warming_shadow: None, @@ -2972,6 +3073,97 @@ mod tests { assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); } + struct Keepalives { + coverage: Arc, + cancel: CancellationToken, + session: tokio::task::JoinHandle, + } + + fn keepalives(backlog: Backlog) -> Keepalives { + let (ws, ws_rx) = tokio::sync::mpsc::channel(8); + let cancel = CancellationToken::new(); + let mut runtime = fresh_runtime(cancel.clone()); + runtime.ws = ScriptedWs::dialing(WsConn { + sink: Box::new(PongingSink { ws, backlog }), + stream: Box::new(ChannelWsStream { rx: ws_rx }), + }); + runtime + .coverage + .update(|c| c.advance(HydrantCursor::new(1))); + let coverage = runtime.coverage.clone(); + let mut cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").expect("hydrant url")); + cfg.idle_promote_min_events = NonZeroU64::new(1).expect("one event is nonzero"); + let session = tokio::spawn(async move { + run_session( + &cfg, + HydrantCursor::new(0), + &runtime, + &FrameCounter::default(), + ) + .await + }); + Keepalives { + coverage, + cancel, + session, + } + } + + async fn promoted_within(k: &Keepalives, intervals: u32) -> bool { + let mut coverage = k.coverage.subscribe(); + tokio::time::timeout( + PING_INTERVAL * intervals, + coverage.wait_for(|c| c.is_ready()), + ) + .await + .is_ok() + } + + fn other_frame(id: u64) -> Result { + Ok(WsMessage::Text( + json!({"id": id, "type": "future_event"}).to_string(), + )) + } + + #[tokio::test(start_paused = true)] + async fn a_socket_answering_keepalives_without_a_frame_promotes_coverage() { + let idle = keepalives(Backlog::Drained); + assert!( + promoted_within(&idle, 4).await, + "the socket answers keepalives and doesn't deliver a record, so the ingest is at the tip", + ); + idle.cancel.cancel(); + idle.session.await.expect("session panicked"); + } + + #[tokio::test(start_paused = true)] + async fn a_socket_delivering_a_record_every_keepalive_stays_warming() { + let busy = keepalives(Backlog::Delivering(1)); + assert!( + !promoted_within(&busy, 8).await, + "a record arrives in every interval, so the ingest is behind the tip", + ); + busy.cancel.cancel(); + busy.session.await.expect("session panicked"); + } + + #[tokio::test] + async fn a_quiet_socket_leaves_a_cold_ingest_warming() { + let runtime = fresh_runtime(CancellationToken::new()); + runtime + .coverage + .update(|c| c.advance(HydrantCursor::new(1))); + promote_if_caught_up( + &runtime, + DEFAULT_IDLE_PROMOTE_MIN_EVENTS, + FramesDelivered(0), + ); + assert!( + !runtime.coverage.snapshot().is_ready(), + "a stream one event in just connected, so the default 256-event minimum leaves it warming", + ); + } + #[test] fn jittered_stays_within_one_quarter_of_base() { let base = Duration::from_secs(1); @@ -3732,49 +3924,29 @@ mod tests { assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); } - struct CloseOnConnectTransport { - used: std::sync::Mutex, - } - impl bobbin_runtime::WsTransport for CloseOnConnectTransport { - fn connect(&self, _url: Url) -> bobbin_runtime::WsConnectFuture { - let mut used = self.used.lock().unwrap(); - if *used { - return Box::pin(async move { - Err(NetworkError::Connect("only one connect allowed".to_owned())) - }); - } - *used = true; - Box::pin(async move { - let mut q = std::collections::VecDeque::new(); - q.push_back(Ok(WsMessage::Close { - code: 1000, - reason: "bye".to_owned(), - })); - let stream: Box = Box::new(ScriptedWsStream { messages: q }); - struct NoopSink; - impl bobbin_runtime::WsSink for NoopSink { - fn send<'a>(&'a mut self, _m: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { - Box::pin(async move { Ok(()) }) - } - } - let sink: Box = Box::new(NoopSink); - Ok(bobbin_runtime::WsConn { sink, stream }) - }) - } - } - #[tokio::test] async fn run_session_returns_after_remote_close_when_outer_cancel_unfired() { let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); let cancel = CancellationToken::new(); let mut runtime = fresh_runtime(cancel.clone()); - runtime.ws = Arc::new(CloseOnConnectTransport { - used: std::sync::Mutex::new(false), + runtime.ws = ScriptedWs::dialing(WsConn { + sink: Box::new(OkSink), + stream: Box::new(ScriptedWsStream { + messages: VecDeque::from([Ok(WsMessage::Close { + code: 1000, + reason: "bye".to_owned(), + })]), + }), }); let res = tokio::time::timeout( Duration::from_secs(3), - run_session(&cfg, HydrantCursor::new(0), &runtime), + run_session( + &cfg, + HydrantCursor::new(0), + &runtime, + &FrameCounter::default(), + ), ) .await; assert!( @@ -3892,7 +4064,7 @@ mod tests { identity: Arc::new(IdentityResolver::detached(RuntimeHasher::default())), clock, entropy: Arc::new(OsEntropy), - ws: TungsteniteWs::shared(), + ws: ScriptedWs::undialed(), cancel: CancellationToken::new(), disconnects: None, warming_shadow: None, diff --git a/bobbin/crates/runtime/Cargo.toml b/bobbin/crates/runtime/Cargo.toml index 2c03198f1..00cc71337 100644 --- a/bobbin/crates/runtime/Cargo.toml +++ b/bobbin/crates/runtime/Cargo.toml @@ -13,6 +13,8 @@ getrandom = { workspace = true } http = { workspace = true } jacquard-common = { workspace = true } reqwest = { workspace = true } +rustls = { workspace = true } +rustls-native-certs = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true } tokio-tungstenite = { workspace = true } diff --git a/bobbin/crates/runtime/src/lib.rs b/bobbin/crates/runtime/src/lib.rs index c8cb3809a..817ecbd0e 100644 --- a/bobbin/crates/runtime/src/lib.rs +++ b/bobbin/crates/runtime/src/lib.rs @@ -16,5 +16,5 @@ pub use mem_network::{ pub use network::{ AddrGuard, BodyStream, GuardedWs, HttpRequest, HttpResponseFuture, HttpResponseHead, HttpResult, HttpTransport, NetworkError, ReqwestHttp, TungsteniteWs, WsConn, WsConnectFuture, - WsMessage, WsMessageFuture, WsSendFuture, WsSink, WsStream, WsTransport, + WsMessage, WsMessageFuture, WsSendFuture, WsSink, WsStream, WsTls, WsTlsError, WsTransport, }; diff --git a/bobbin/crates/runtime/src/network.rs b/bobbin/crates/runtime/src/network.rs index ede0d91cf..69f77dd0d 100644 --- a/bobbin/crates/runtime/src/network.rs +++ b/bobbin/crates/runtime/src/network.rs @@ -7,6 +7,7 @@ use bytes::Bytes; use futures::stream::{Stream, StreamExt}; use http::{HeaderMap, StatusCode}; use thiserror::Error; +use tokio_tungstenite::Connector; use tokio_tungstenite::tungstenite::{ Bytes as WsBytes, Message as TungsteniteMessage, protocol::CloseFrame as TungsteniteClose, protocol::frame::coding::CloseCode as TungsteniteCloseCode, @@ -170,45 +171,94 @@ pub trait WsTransport: Send + Sync + 'static { pub type AddrGuard = Arc Result<(), NetworkError> + Send + Sync>; -#[derive(Clone, Copy, Debug, Default)] -pub struct TungsteniteWs; +#[derive(Debug, Error)] +pub enum WsTlsError { + #[error("no root certificates in the system trust store: {0:?}")] + NoRoots(Vec), + #[error("all {0} certificates in the system trust store failed to parse")] + Unparsable(usize), + #[error("the aws-lc-rs provider doesn't support any of the default protocol versions: {0}")] + Versions(rustls::Error), +} + +#[derive(Clone)] +pub struct WsTls(Arc); + +impl WsTls { + pub fn from_native_roots() -> Result { + let loaded = rustls_native_certs::load_native_certs(); + if loaded.certs.is_empty() { + return Err(WsTlsError::NoRoots(loaded.errors)); + } + let mut roots = rustls::RootCertStore::empty(); + let (accepted, rejected) = roots.add_parsable_certificates(loaded.certs); + if accepted == 0 { + return Err(WsTlsError::Unparsable(rejected)); + } + let provider = Arc::new(rustls::crypto::aws_lc_rs::default_provider()); + let config = rustls::ClientConfig::builder_with_provider(provider) + .with_safe_default_protocol_versions() + .map_err(WsTlsError::Versions)? + .with_root_certificates(roots) + .with_no_client_auth(); + Ok(Self(Arc::new(config))) + } + + fn connector(&self) -> Connector { + Connector::Rustls(Arc::clone(&self.0)) + } +} + +pub struct TungsteniteWs { + tls: WsTls, +} impl TungsteniteWs { - pub fn shared() -> Arc { - Arc::new(Self) + pub fn shared(tls: WsTls) -> Arc { + Arc::new(Self { tls }) + } +} + +fn wired(ws: TungsteniteWsStream) -> WsConn { + let (sink, stream) = futures::StreamExt::split(ws); + WsConn { + sink: Box::new(TungsteniteSink { inner: sink }), + stream: Box::new(TungsteniteStream { inner: stream }), } } impl WsTransport for TungsteniteWs { fn connect(&self, url: Url) -> WsConnectFuture { + let connector = self.tls.connector(); Box::pin(async move { - let url_str = url.as_str().to_owned(); - let (ws, _resp) = tokio_tungstenite::connect_async(&url_str) - .await - .map_err(|e| NetworkError::Connect(e.to_string()))?; - let (sink_inner, stream_inner) = futures::StreamExt::split(ws); - let sink: Box = Box::new(TungsteniteSink { inner: sink_inner }); - let stream: Box = Box::new(TungsteniteStream { - inner: stream_inner, - }); - Ok(WsConn { sink, stream }) + let (ws, _resp) = tokio_tungstenite::connect_async_tls_with_config( + url.as_str(), + None, + false, + Some(connector), + ) + .await + .map_err(|e| NetworkError::Connect(e.to_string()))?; + Ok(wired(ws)) }) } } pub struct GuardedWs { guard: AddrGuard, + tls: WsTls, } impl GuardedWs { - pub fn shared(guard: AddrGuard) -> Arc { - Arc::new(Self { guard }) + pub fn shared(guard: AddrGuard, tls: WsTls) -> Arc { + Arc::new(Self { guard, tls }) } } impl WsTransport for GuardedWs { fn connect(&self, url: Url) -> WsConnectFuture { let guard = self.guard.clone(); + let connector = self.tls.connector(); Box::pin(async move { let host = url .host_str() @@ -229,15 +279,15 @@ impl WsTransport for GuardedWs { let tcp = tokio::net::TcpStream::connect(addr) .await .map_err(|e| NetworkError::Connect(e.to_string()))?; - let (ws, _resp) = tokio_tungstenite::client_async_tls(url.as_str(), tcp) - .await - .map_err(|e| NetworkError::Connect(e.to_string()))?; - let (sink_inner, stream_inner) = futures::StreamExt::split(ws); - let sink: Box = Box::new(TungsteniteSink { inner: sink_inner }); - let stream: Box = Box::new(TungsteniteStream { - inner: stream_inner, - }); - Ok(WsConn { sink, stream }) + let (ws, _resp) = tokio_tungstenite::client_async_tls_with_config( + url.as_str(), + tcp, + None, + Some(connector), + ) + .await + .map_err(|e| NetworkError::Connect(e.to_string()))?; + Ok(wired(ws)) }) } } diff --git a/bobbin/crates/xrpc/tests/await_record.rs b/bobbin/crates/xrpc/tests/await_record.rs index b55e9909e..499c80713 100644 --- a/bobbin/crates/xrpc/tests/await_record.rs +++ b/bobbin/crates/xrpc/tests/await_record.rs @@ -14,7 +14,7 @@ use bobbin_slingshot_client::SlingshotClient; use bobbin_xrpc::{AppState, MaxAwaiting, router}; use http::{Request, StatusCode}; use jacquard_common::DefaultStr; -use jacquard_common::types::string::{AtUri, Cid}; +use jacquard_common::types::string::AtUri; use serde_json::{Value, json}; use tower::ServiceExt; use url::Url; diff --git a/bobbin/crates/xrpc/tests/await_record_e2e.rs b/bobbin/crates/xrpc/tests/await_record_e2e.rs index 60033b704..29e2423da 100644 --- a/bobbin/crates/xrpc/tests/await_record_e2e.rs +++ b/bobbin/crates/xrpc/tests/await_record_e2e.rs @@ -144,6 +144,7 @@ async fn a_record_off_the_stream_answers_an_await() { hydrant_base: Url::parse("ws://hydrant.invalid").unwrap(), start_cursor: HydrantCursor::new(0), parallelism: bobbin_ingest::DEFAULT_INGEST_PARALLELISM, + idle_promote_min_events: bobbin_ingest::DEFAULT_IDLE_PROMOTE_MIN_EVENTS, }; let ingesting = tokio::spawn(run_ingest(config, ingest)); let answer = tokio::time::timeout(Duration::from_secs(10), router(state).oneshot(request)) diff --git a/bobbin/example.toml b/bobbin/example.toml index 751a23af6..e203ff367 100644 --- a/bobbin/example.toml +++ b/bobbin/example.toml @@ -70,6 +70,15 @@ # Default value: 16 #parallelism = 16 +# Events processed before a quiet stream may promote coverage to ready. +# The invite listings will be `503 Warming` while a quiet stream is under +# said count. +# +# Can also be specified via environment variable `BOBBIN_INGEST_IDLE_PROMOTE_MIN_EVENTS`. +# +# Default value: 256 +#idle_promote_min_events = 256 + [backpressure] # Estimated peak transient anonymous bytes a single hydrating request holds: # roughly the page of decoded records plus its serialized copy. Combined with diff --git a/docker-compose.yml b/docker-compose.yml index ce8166019..8f2bc68fe 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -43,7 +43,9 @@ services: networks: [tngl] init-accounts: - image: alpine:3.22 + build: + context: localinfra + dockerfile: seed.Dockerfile restart: "no" env_file: localinfra/pds.env environment: @@ -54,14 +56,16 @@ services: volumes: - ./localinfra/scripts:/scripts:ro - init-state:/shared - command: sh -c "apk add --no-cache curl jq >/dev/null && sh /scripts/init-accounts.sh" + command: sh /scripts/init-accounts.sh depends_on: pds: condition: service_healthy networks: [tngl] init-data: - image: alpine:3.22 + build: + context: localinfra + dockerfile: seed.Dockerfile restart: "no" env_file: localinfra/pds.env environment: @@ -71,8 +75,10 @@ services: volumes: - ./localinfra/scripts:/scripts:ro - init-state:/shared:ro - command: sh -c "apk add --no-cache curl jq >/dev/null && sh /scripts/init-data.sh" + command: sh /scripts/init-data.sh depends_on: + init-accounts: + condition: service_completed_successfully knot2: condition: service_healthy networks: [tngl] @@ -104,7 +110,7 @@ services: JETSTREAM_WS_URL: wss://pds.tngl.boltless.dev/xrpc/com.atproto.sync.subscribeRepos volumes: - jetstream-data:/data - - ./localinfra/certs/root.crt:/etc/ssl/certs/ca-certificates.crt:ro + - ./localinfra/certs/root.crt:/etc/ssl/certs/tngl-dev-root.crt:ro depends_on: pds: condition: service_healthy @@ -151,7 +157,7 @@ services: condition: service_started init-accounts: condition: service_completed_successfully - networks: [tngl] + networks: [tngl, upstream-cache] ncps-migrate: image: &ncps-image ghcr.io/kalbasit/ncps:v0.9.4 @@ -181,7 +187,7 @@ services: depends_on: ncps-migrate: condition: service_completed_successfully - networks: [tngl] + networks: [tngl, upstream-cache] spindle: profiles: ["linux"] @@ -264,7 +270,7 @@ services: condition: service_completed_successfully ncps: condition: service_started - networks: [tngl] + networks: [tngl, upstream-cache] knotmirror-tap: image: ghcr.io/bluesky-social/indigo/tap:sha-4f47add43060c27e8a37d9d76482ecddf001fcd8 # 0.1.10 @@ -279,7 +285,7 @@ services: TAP_RESYNC_PARALLELISM: "10" TAP_RETRY_TIMEOUT: 60s volumes: - - ./localinfra/certs/root.crt:/etc/ssl/certs/ca-certificates.crt:ro + - ./localinfra/certs/root.crt:/etc/ssl/certs/tngl-dev-root.crt:ro depends_on: postgres: condition: service_started @@ -361,7 +367,7 @@ services: environment: TANGLED_ZOEKT_INDEX_DIR: /data/index TANGLED_ZOEKT_PLC_URL: https://plc.tngl.boltless.dev - TANGLED_ZOEKT_APPVIEW_URL: http://127.0.0.1:3000 + TANGLED_ZOEKT_APPVIEW_URL: http://appview:3000 TANGLED_ZOEKT_ALLOW_HTTP: "true" volumes: - zoekt-index:/data/index @@ -495,7 +501,7 @@ services: RUST_LOG: info volumes: - hydrant-data:/data - - ./localinfra/certs/root.crt:/etc/ssl/certs/ca-certificates.crt:ro + - ./localinfra/certs/root.crt:/etc/ssl/certs/tngl-dev-root.crt:ro healthcheck: test: ["CMD", "bash", "-c", "echo > /dev/tcp/127.0.0.1/3000"] interval: 5s @@ -523,15 +529,16 @@ services: BOBBIN_SLINGSHOT_URL: http://hydrant:3000 BOBBIN_SERVICE_DID: did:web:bobbin.tngl.boltless.dev BOBBIN_KNOT_ALLOW_PRIVATE: "true" - BOBBIN_KNOT_REQUIRE_HTTPS: "false" BOBBIN_MIRROR_V2_URL: http://knotmirror:7000 BOBBIN_CODESEARCH_ZOEKT_URL: https://zoekt.tngl.boltless.dev BOBBIN_LOG: info + BOBBIN_KNOT_REQUIRE_HTTPS: "true" + BOBBIN_INGEST_IDLE_PROMOTE_MIN_EVENTS: "1" volumes: - .:/src:cached - bobbin-cargo:/cargo - bobbin-target:/target - - ./localinfra/certs/root.crt:/etc/ssl/certs/ca-certificates.crt:ro + - ./localinfra/certs/root.crt:/etc/ssl/certs/tngl-dev-root.crt:ro healthcheck: test: ["CMD", "wget", "-qO-", "http://localhost:8090/xrpc/sh.tangled.bobbin.getCoverage"] interval: 5s @@ -671,6 +678,7 @@ services: - bobbin.tngl.boltless.dev - camo.tngl.boltless.dev - avatar.tngl.boltless.dev + - deliberi.tngl.boltless.dev prometheus: image: prom/prometheus:v2.54.1 @@ -762,9 +770,13 @@ volumes: networks: tngl: driver: bridge + internal: true # Public-looking subnet so SSRF checks see container IPs as "public". # RFC1918 + doc/benchmark ranges are blocklisted; 11.x is unrouted on # the public internet, so it passes the check and won't collide. ipam: config: - subnet: 11.0.0.0/24 + + upstream-cache: + driver: bridge diff --git a/knot2/crates/knot-edge/Cargo.toml b/knot2/crates/knot-edge/Cargo.toml index 866e09c77..bbc839dd6 100644 --- a/knot2/crates/knot-edge/Cargo.toml +++ b/knot2/crates/knot-edge/Cargo.toml @@ -11,7 +11,7 @@ axum = { workspace = true } hyper = { workspace = true } hyper-util = { workspace = true } tower = { workspace = true, features = ["limit", "load-shed"] } -tower-http = { workspace = true, features = ["compression-zstd", "compression-br", "compression-gzip", "timeout", "map-request-body"] } +tower-http = { workspace = true, features = ["compression-zstd", "compression-br", "compression-gzip", "timeout", "map-request-body", "cors", "set-header"] } tower_governor = { workspace = true } governor = { workspace = true } http-body = { workspace = true } diff --git a/knot2/crates/knot-edge/src/cors.rs b/knot2/crates/knot-edge/src/cors.rs new file mode 100644 index 000000000..489ef5c64 --- /dev/null +++ b/knot2/crates/knot-edge/src/cors.rs @@ -0,0 +1,30 @@ +//! Preflight is answered under the guards, so a burst of them is rate +//! limited like any other request. The allow-origin header alone is set +//! outside the governor, which keeps a 429 or a shed 503 readable to a browser. + +use std::time::Duration; + +use axum::Router; +use http::header::{ACCESS_CONTROL_ALLOW_ORIGIN, AUTHORIZATION, CONTENT_TYPE}; +use http::{HeaderValue, Method}; +use tower_http::cors::{Any, CorsLayer}; +use tower_http::set_header::SetResponseHeaderLayer; + +const PREFLIGHT_CACHE: Duration = Duration::from_secs(86_400); + +pub(crate) fn preflight(router: Router) -> Router { + router.layer( + CorsLayer::new() + .allow_origin(Any) + .allow_methods([Method::GET, Method::POST, Method::OPTIONS]) + .allow_headers([AUTHORIZATION, CONTENT_TYPE]) + .max_age(PREFLIGHT_CACHE), + ) +} + +pub(crate) fn origin_header(router: Router) -> Router { + router.layer(SetResponseHeaderLayer::if_not_present( + ACCESS_CONTROL_ALLOW_ORIGIN, + HeaderValue::from_static("*"), + )) +} diff --git a/knot2/crates/knot-edge/src/lib.rs b/knot2/crates/knot-edge/src/lib.rs index 5248b7845..37c274a15 100644 --- a/knot2/crates/knot-edge/src/lib.rs +++ b/knot2/crates/knot-edge/src/lib.rs @@ -1,6 +1,7 @@ mod acme; mod altsvc; mod compression; +mod cors; mod limits; mod peer; mod protocol; @@ -184,9 +185,9 @@ pub async fn serve( .as_ref() .and_then(|setup| setup.internal.as_ref()) .is_some(); - let base = base_router(app, early_data_safe); + let base = cors::preflight(base_router(app, early_data_safe)); let internal_router = wants_internal.then(|| finish(base.clone())); - let router = finish(robustness::apply(base, layers)); + let router = finish(cors::origin_header(robustness::apply(base, layers))); let listener = TcpListener::bind(http_addr.get()).await?; let Some(setup) = tls else { @@ -301,7 +302,8 @@ mod tests { let full = RequiresFullHandshake::new( Router::new().route("/git-upload-pack", post(|| async { "pack" })), ); - finish(robustness::apply(base_router(full, safe), test_layers())) + let base = cors::preflight(base_router(full, safe)); + finish(cors::origin_header(robustness::apply(base, test_layers()))) } fn tight_layers() -> robustness::GuardLayers { @@ -350,6 +352,21 @@ mod tests { assert_eq!(status_of(request).await, StatusCode::OK); } + #[tokio::test] + async fn a_preflight_in_early_data_is_answered_and_not_deferred_to_425() { + let request = Request::options("/git-upload-pack") + .header("early-data", "1") + .header(http::header::ORIGIN, "https://tangled.org") + .header(http::header::ACCESS_CONTROL_REQUEST_METHOD, "POST") + .body(Body::empty()) + .unwrap(); + let answer = status_of(request).await; + assert!( + answer.is_success(), + "a preflight answer is a constant, so replaying one doesn't get an attacker a write: got {answer}" + ); + } + #[tokio::test] async fn the_internal_admin_router_shares_no_rate_limit_budget_with_the_data_plane() { let safe = ZeroRttRoutes::new().get("/info/refs", ZeroRttSafe::new(|| async { "ok" })); diff --git a/knot2/crates/knot-edge/src/robustness.rs b/knot2/crates/knot-edge/src/robustness.rs index 97fd8fea8..be8e3c7c0 100644 --- a/knot2/crates/knot-edge/src/robustness.rs +++ b/knot2/crates/knot-edge/src/robustness.rs @@ -545,6 +545,76 @@ mod tests { ); } + fn browser_edge(router: Router, guards: EdgeGuards) -> Router { + crate::cors::origin_header(guarded_router(crate::cors::preflight(router), guards)) + } + + fn preflight(method: &str, headers: &str) -> Request { + Request::options("/") + .header(http::header::ORIGIN, "https://tangled.org") + .header(http::header::ACCESS_CONTROL_REQUEST_METHOD, method) + .header(http::header::ACCESS_CONTROL_REQUEST_HEADERS, headers) + .body(Body::empty()) + .unwrap() + } + + fn from_tangled_org() -> Request { + Request::get("/") + .header(http::header::ORIGIN, "https://tangled.org") + .body(Body::empty()) + .unwrap() + } + + #[tokio::test] + async fn the_cors_layer_answers_a_full_preflight_from_under_the_governor() { + let app = browser_edge( + Router::new().route("/", get(|| async { "ok" })), + guards(1, 1, 1_024, 60_000, 30_000, ProxyTrust::default()), + ); + let ask = || from_peer(preflight("POST", "authorization,content-type"), 12); + let answer = app.clone().oneshot(ask()).await.unwrap(); + assert!(answer.status().is_success(), "got {:?}", answer.status()); + let allowed = answer.headers()[http::header::ACCESS_CONTROL_ALLOW_HEADERS] + .to_str() + .map(str::to_ascii_lowercase) + .unwrap(); + assert!( + allowed.contains("authorization") && allowed.contains("content-type"), + "every authed xrpc call sends both headers: allowed {allowed}" + ); + assert_eq!( + answer.headers()[http::header::ACCESS_CONTROL_MAX_AGE], + "86400", + "without a max-age a browser re-preflights and every call costs two tokens" + ); + assert_eq!( + app.oneshot(ask()).await.unwrap().status(), + StatusCode::TOO_MANY_REQUESTS, + "an unmetered OPTIONS would bypass every guard" + ); + } + + #[tokio::test] + async fn the_allow_origin_header_outlives_a_handler_403_and_a_governor_429() { + let app = browser_edge( + Router::new().route("/", get(|| async { StatusCode::FORBIDDEN })), + guards(1, 1, 1_024, 60_000, 30_000, ProxyTrust::default()), + ); + let ask = || from_peer(from_tangled_org(), 13); + let denied = app.clone().oneshot(ask()).await.unwrap(); + let refused = app.oneshot(ask()).await.unwrap(); + assert_eq!(denied.status(), StatusCode::FORBIDDEN); + assert_eq!(refused.status(), StatusCode::TOO_MANY_REQUESTS); + [denied, refused].into_iter().for_each(|answer| { + assert_eq!( + answer.headers()[http::header::ACCESS_CONTROL_ALLOW_ORIGIN], + "*", + "the frontend prints the {} itself, since this header is set outside the governor", + answer.status() + ); + }); + } + #[tokio::test] async fn requests_beyond_the_inflight_limit_are_shed_with_503() { let app = guarded_router( diff --git a/knot2/crates/knot-events/src/lib.rs b/knot2/crates/knot-events/src/lib.rs index 24d551d70..0584536d5 100644 --- a/knot2/crates/knot-events/src/lib.rs +++ b/knot2/crates/knot-events/src/lib.rs @@ -173,8 +173,14 @@ impl EventLog { } } - pub fn with_floor(self, floor: UnixMicros) -> Self { - self.raise_floor(floor); + pub fn with_boot_floor(self, floor: UnixMicros) -> Self { + { + let mut ring = self.inner.lock(); + ring.last_micros = ring.last_micros.max(floor); + let booted = (ring.last_micros > UnixMicros::new(0)) + .then(|| EventCursor::from_unix_micros(ring.last_micros)); + ring.evicted_through = ring.evicted_through.max(booted); + } self } @@ -440,6 +446,38 @@ mod tests { assert!(!log.lost_before(fourth)); } + #[test] + fn a_boot_floor_loses_every_cursor_below_it_and_survives_a_lower_floor() { + let floor = UnixMicros::new(1_700_000_000_000_000); + let low = UnixMicros::new(1_600_000_000_000_000); + let at_floor = EventCursor::from_unix_micros(floor); + let booted = log(8).with_boot_floor(floor); + assert!( + booted.lost_before(EventCursor::START) + && booted.lost_before(EventCursor::new(at_floor.get() - 1)), + "a just-booted ring can't replay an event from below its floor" + ); + assert!( + !booted.lost_before(at_floor), + "the floor is the highest seq issued, so a cursor there didn't miss an event" + ); + let unbooted = log(8).with_boot_floor(UnixMicros::new(0)); + assert!( + !unbooted.lost_before(EventCursor::START), + "a knot that hasn't issued an event hasn't evicted one either" + ); + let lowered = booted.with_boot_floor(low); + assert!( + lowered.lost_before(EventCursor::from_unix_micros(low)), + "lowering the floor after boot doesn't bring back events from before it" + ); + assert_eq!( + lowered.last_issued(), + floor, + "the higher floor is still the highest seq issued" + ); + } + #[test] fn a_wide_event_evicts_by_bytes_long_before_the_ring_fills() { let log = EventLog::new( diff --git a/knot2/crates/knot-server/src/main.rs b/knot2/crates/knot-server/src/main.rs index af0150367..d4dc74bab 100644 --- a/knot2/crates/knot-server/src/main.rs +++ b/knot2/crates/knot-server/src/main.rs @@ -602,7 +602,7 @@ async fn main() -> anyhow::Result<()> { appview: appview_endpoint, slots: slots.clone(), firehose: Arc::new( - knot_events::EventLog::new(SystemClock, replay_bounds).with_floor(firehose_floor), + knot_events::EventLog::new(SystemClock, replay_bounds).with_boot_floor(firehose_floor), ), lfs: lfs_handle .clone() @@ -673,7 +673,6 @@ async fn main() -> anyhow::Result<()> { HomepageSource::Default => base_router.route("/", get(|| async { Html(DEFAULT_HOMEPAGE) })), HomepageSource::File(path) => base_router.route_service("/", ServeFile::new(path)), }; - let base_router = base_router.layer(knot_xrpc::browser_cors()); let app = knot_edge::RequiresFullHandshake::new(base_router); let scheme = if tls_setup.is_some() { "https" } else { "http" }; let edge_config = knot_edge::EdgeConfig { diff --git a/knot2/crates/knot-xrpc/Cargo.toml b/knot2/crates/knot-xrpc/Cargo.toml index 13bc01539..965fd9fe5 100644 --- a/knot2/crates/knot-xrpc/Cargo.toml +++ b/knot2/crates/knot-xrpc/Cargo.toml @@ -38,7 +38,7 @@ jacquard-axum = { workspace = true } tracing = { workspace = true } axum = { workspace = true, features = ["ws"] } tower = { workspace = true } -tower-http = { workspace = true, features = ["fs", "cors"] } +tower-http = { workspace = true, features = ["fs"] } http-body = { workspace = true } tokio-util = { workspace = true, features = ["io-util"] } tokio = { workspace = true } diff --git a/knot2/crates/knot-xrpc/src/firehose.rs b/knot2/crates/knot-xrpc/src/firehose.rs index 4c880fbad..75afcdf6e 100644 --- a/knot2/crates/knot-xrpc/src/firehose.rs +++ b/knot2/crates/knot-xrpc/src/firehose.rs @@ -1161,7 +1161,7 @@ mod tests { ManualClock::new(UnixMicros::new(1_700_000_000_000_000 - 10_000_000)), drain_bounds(), ) - .with_floor(floor); + .with_boot_floor(floor); let persisted = stepped_back.last_issued().get(); assert!( !future_cursor(&stepped_back, persisted), diff --git a/knot2/crates/knot-xrpc/src/lib.rs b/knot2/crates/knot-xrpc/src/lib.rs index 0269145f3..d87a8b09d 100644 --- a/knot2/crates/knot-xrpc/src/lib.rs +++ b/knot2/crates/knot-xrpc/src/lib.rs @@ -58,13 +58,9 @@ use axum::middleware::{Next, from_fn_with_state}; use axum::response::{IntoResponse, Response}; use axum::routing::{get, post}; use http::request::Parts; -use http::{ - HeaderMap, HeaderValue, StatusCode, - header::{AUTHORIZATION, CONTENT_TYPE}, -}; +use http::{HeaderMap, HeaderValue, StatusCode, header::AUTHORIZATION}; use serde::de::DeserializeOwned; use serde_json::json; -use tower_http::cors::{Any, CorsLayer}; use knot_atproto::{Atproto, AtprotoError, ServiceJwt}; use knot_events::{EventLog, SubscriberGate}; @@ -289,13 +285,6 @@ impl FromRequestParts for Method { } } -pub fn browser_cors() -> CorsLayer { - CorsLayer::new() - .allow_origin(Any) - .allow_methods([http::Method::GET, http::Method::POST, http::Method::OPTIONS]) - .allow_headers([AUTHORIZATION, CONTENT_TYPE]) -} - pub fn router(state: Arc>) -> Router { let merge_routes = Router::new() // .route(merge::MERGE_ROUTE, post(merge::merge::)) diff --git a/knot2/crates/knot-xrpc/src/plc.rs b/knot2/crates/knot-xrpc/src/plc.rs index d5951f778..a1df0c775 100644 --- a/knot2/crates/knot-xrpc/src/plc.rs +++ b/knot2/crates/knot-xrpc/src/plc.rs @@ -488,4 +488,21 @@ mod backoff_tests { "the sweep parks the repository on the twelfth failure" ); } + + #[test] + fn a_repository_parks_after_five_hours_of_failing_at_every_due_time() { + let mut memory = SweepMemory::default(); + let did = "did:plc:nelpet"; + let due_at_last_retry = (1..MAX_ATTEMPTS).fold(0u64, |now, _| { + memory.schedule(did, now); + memory.waiting[did].due_at_micros + }); + assert_eq!( + (memory.is_parked(did), due_at_last_retry / 1_000_000), + (false, 18_210), + "the eleven delays sum to five hours, so the test asserts the arithmetic instead of sleeping" + ); + memory.schedule(did, due_at_last_retry); + assert!(memory.is_parked(did), "the twelfth failure parks the repo"); + } } diff --git a/knot2/crates/knot-xrpc/src/tests.rs b/knot2/crates/knot-xrpc/src/tests.rs index 4ebb9078b..d7b3736dd 100644 --- a/knot2/crates/knot-xrpc/src/tests.rs +++ b/knot2/crates/knot-xrpc/src/tests.rs @@ -3838,7 +3838,7 @@ mod rosters { } #[tokio::test] - async fn a_browser_can_preflight_the_acceptance_route_and_read_what_it_answers() { + async fn the_acceptance_route_takes_no_token_minted_for_another_subject() { use axum::body::Body; use axum::extract::ConnectInfo; use std::net::SocketAddr; @@ -3846,53 +3846,12 @@ mod rosters { let offer = Offer::opened(Roster::Knot, Seen::Projected).await; offer.answered(StatusCode::OK); - let app = crate::router(Arc::clone(&offer.world.state)).layer(crate::browser_cors()); + let app = crate::router(Arc::clone(&offer.world.state)); let peer = SocketAddr::from(([203, 0, 113, 42], 5555)); - let origin = "https://tangled.org"; - - let mut preflight = http::Request::builder() - .method("OPTIONS") - .uri(crate::members::ACCEPT_ROUTE) - .header(http::header::ORIGIN, origin) - .header(http::header::ACCESS_CONTROL_REQUEST_METHOD, "POST") - .header( - http::header::ACCESS_CONTROL_REQUEST_HEADERS, - "authorization,content-type", - ) - .body(Body::empty()) - .unwrap(); - preflight.extensions_mut().insert(ConnectInfo(peer)); - - let answer = app.clone().oneshot(preflight).await.unwrap(); - assert!( - answer.status().is_success(), - "the frontend sends a preflight before every accept, and it must pass: {:?}", - answer.status() - ); - assert_eq!( - answer - .headers() - .get(http::header::ACCESS_CONTROL_ALLOW_ORIGIN), - Some(&HeaderValue::from_static("*")), - "frontend and knot sit on different origins" - ); - let allowed = answer - .headers() - .get(http::header::ACCESS_CONTROL_ALLOW_HEADERS) - .and_then(|value| value.to_str().ok()) - .unwrap_or_default() - .to_ascii_lowercase(); - for wanted in ["authorization", "content-type"] { - assert!( - allowed.contains(wanted), - "the accept call sends {wanted}, and the preflight allowed only {allowed}" - ); - } let mut refused = http::Request::builder() .method("POST") .uri(crate::members::ACCEPT_ROUTE) - .header(http::header::ORIGIN, origin) .header( http::header::AUTHORIZATION, format!("Bearer {}", mint(&offer.world.admin, ACCEPT_MEMBERSHIP)), @@ -3909,14 +3868,10 @@ mod rosters { .unwrap(); refused.extensions_mut().insert(ConnectInfo(peer)); - let answer = app.oneshot(refused).await.unwrap(); - assert_eq!(answer.status(), StatusCode::FORBIDDEN); assert_eq!( - answer - .headers() - .get(http::header::ACCESS_CONTROL_ALLOW_ORIGIN), - Some(&HeaderValue::from_static("*")), - "the row prints the knot's own sentence, so the browser needs to read it" + app.oneshot(refused).await.unwrap().status(), + StatusCode::FORBIDDEN, + "the admin minted a well-formed token for an invite addressed to somebody else" ); } diff --git a/localinfra/.dockerignore b/localinfra/.dockerignore new file mode 100644 index 000000000..ed34186cf --- /dev/null +++ b/localinfra/.dockerignore @@ -0,0 +1,2 @@ +certs/*.key +ncps-secret-key diff --git a/localinfra/Caddyfile b/localinfra/Caddyfile index e8ce9dcfb..56ec092c2 100644 --- a/localinfra/Caddyfile +++ b/localinfra/Caddyfile @@ -52,10 +52,13 @@ jetstream.tngl.boltless.dev { reverse_proxy jetstream:6008 } -# knot2 serves no cors headers of its own. +# knot2 answers preflight itself (knot2/crates/knot-edge/src/cors.rs), so this block +# doesn't import cors: a browser rejects a response with two Access-Control-Allow-Origin +# values, and web/ posts sh.tangled.knot.acceptMembership and sh.tangled.git.keepCommit +# straight from the browser to the knot serving the repo. knot2.tngl.boltless.dev { tls internal - import cors knot2:5555 + reverse_proxy knot2:5555 } # spindle diff --git a/localinfra/appview.Dockerfile b/localinfra/appview.Dockerfile index d199ecc31..28794e46c 100644 --- a/localinfra/appview.Dockerfile +++ b/localinfra/appview.Dockerfile @@ -1,6 +1,6 @@ # Development only. Not for production use. -FROM golang:1.25-alpine +FROM golang:1.26-alpine RUN apk add --no-cache git build-base sqlite-dev tini sqlite-libs ca-certificates diff --git a/localinfra/bobbin.Dockerfile b/localinfra/bobbin.Dockerfile index 3963ed14f..fde109904 100644 --- a/localinfra/bobbin.Dockerfile +++ b/localinfra/bobbin.Dockerfile @@ -1,5 +1,5 @@ # Development only. Not for production use. -FROM rust:1.96-slim-trixie +FROM rust:1.96.0-slim-trixie RUN apt-get update && apt-get install -y --no-install-recommends \ ca-certificates pkg-config perl make cmake clang mold git curl wget \ diff --git a/localinfra/readme.md b/localinfra/readme.md index 97086b6a8..3d19024b9 100644 --- a/localinfra/readme.md +++ b/localinfra/readme.md @@ -1,4 +1,5 @@ Heavily inspired by [frontpage dev environment](https://github.com/frontpagefyi/frontpage/blob/10678df9c3f72cbd82f0856a9f99c74dd22326d8/apps/frontpage/local-infra/README.md). + Tangled's setup is slightly more involved because services inside the network need to reach the PDS over its **public** hostname with **valid TLS** — federation paths (DID resolution, OAuth, etc.) round-trip through the same URLs an external client would use. For example, resolving `alice.pds.tngl.boltless.dev` yields an `#atproto_pds` service pointing at `https://pds.tngl.boltless.dev`. Knot and spindle running inside docker must hit that exact URL and trust its cert. @@ -15,7 +16,7 @@ To make that work: - jetstream () - knot () - spindle () -- knotmirror () +- knotmirror () - appview () (live reloading) - [ncps](https://github.com/kalbasit/ncps) nix binary cache (internal, `http://ncps:8501`) - pdsls () @@ -56,14 +57,65 @@ To make that work: ./localinfra/certs/root.crt ``` - Depending on your browser you may have to import the certificate into your browser profiles too as some have their own certs do not use your system ones -3. run `./localinfra/scripts/appview-static-files.sh` -4. Prepare the spindle microVM images: +3. Fetch the appview's vendored static assets. The tailwind service writes + `tw.css` alone, so skipping this leaves the appview without htmx, mermaid, + mathjax, fonts or icons: + ```bash + ./localinfra/scripts/appview-static-files.sh + ``` +4. For the `linux` profile's spindle, prepare its microVM images, written under + `out/localinfra-spindle-images`: ```bash ./localinfra/scripts/prepare-spindle-images.sh ``` - This writes the image directory under `out/localinfra-spindle-images`. -5. `docker compose up` -6. AppView will be running on `127.0.0.1:3000` with four test users: `alice`, `bob`, `charlie`, and `david`, all under `pds.tngl.boltless.dev` and using the password `password`. Deliberi also bootstraps a verified `${user}@pds.tngl.boltless.dev` address for each account. Use that address as `user.email` when making local Git commits that should appear in profile activity. + + spindle's nix builds pull store paths through `ncps`, which proxies + `cache.nixos.org`. They also fetch `github:` flakerefs straight from + `github.com` as a tarball, so `ncps` never sees those. Both services join a + second bridge, `upstream-cache`, and so does knot2. `repo.create` for `core` + passes `https://knot1.tangled.sh/...` as its `source`, and the clone runs + inside knot2. The seed only posts the URL, so knot2 is the container that + needs the route out. Everything else stays on `tngl` alone, which is + `internal: true` and answers NXDOMAIN for any name outside the project. + Without that bridge the profile fails on the first flake input or store + path it tries to fetch. + +5. `podman compose build` (or `docker compose build`), then warm the dev-loop + cache volumes while you still have a route out. air rebuilds the appview, + cargo rebuilds bobbin, and pnpm installs web/camo/avatar, all into volumes + that a container can't fill from inside `internal: true`: + ```bash + ./localinfra/scripts/warm-caches.sh + ``` + The script calls `podman run` with the volume and image names podman-compose + uses, so on docker you run its five commands by hand. Run it again after a + `Cargo.lock` or `pnpm-lock.yaml` change. Otherwise the container-start + `cargo build` and `pnpm install` try crates.io and the npm registry from + inside `tngl`, where a stale cache reads as a network error. +6. `podman compose up` (or `docker compose up`) +7. AppView will be running on `127.0.0.1:3000` with four test users: `alice`, `bob`, `charlie`, and `david`, all under `pds.tngl.boltless.dev` and using the password `password`. Deliberi also bootstraps a verified `${user}@pds.tngl.boltless.dev` address for each account. Use that address as `user.email` when making local Git commits that should appear in profile activity. + + On rootless podman, caddy's published `80:80`/`443:443` need + `net.ipv4.ip_unprivileged_port_start` at or below 80 (`sudo sysctl -w + net.ipv4.ip_unprivileged_port_start=80`, persisted in the host config). + + Rootless podman with `internal: true` has one more trap: netavark doesn't + write firewall rules for an internal network, and that includes the accept + rule for aardvark-dns on the bridge gateway. Where the host input policy + drops, containers lose name resolution rather than egress, and `postgres` + and `redis` failing to resolve reads as a broken DNS setup. See + [netavark#1055](https://github.com/containers/netavark/issues/1055) and + [podman#26917](https://github.com/containers/podman/issues/26917). Docker + answers on 127.0.0.11 without a gateway hop and is unaffected. + + Signup is the one flow `internal: true` shuts off: the appview and deliberi + both verify a Turnstile token against `challenges.cloudflare.com` from + `tngl`. With the secret key unset, which is how the dev stack runs, each + bails out before the call, so signup answers a captcha error. Set + `TANGLED_CLOUDFLARE_TURNSTILE_SECRET_KEY` or + `DELIBERI_TURNSTILE_SECRET_KEY` and it answers 500 until you put both + services on `upstream-cache` too. + `TANGLED_APPVIEW_HOST` must be a loopback IP with the mapped port (`127.0.0.1:3000`), not `localhost`: atproto's dev OAuth client requires a loopback IP for the redirect URI. If you remap the published appview port, update `TANGLED_APPVIEW_HOST` in `docker-compose.yml` to match. ## Observability @@ -130,6 +182,6 @@ The executors differ on the two axes placement cares about, so you can watch can | executor-b | `linux, slow` | full | `image/alpine, image/nixos` | | executor-c | `linux, gpu` | alpine only | `image/alpine` | -So `image: nixos` has two candidates, `image: alpine` has three, and `runs_on: [gpu]` pins to executor-c. The alpine-only image set is staged by `prepare-spindle-images.sh` (step 4) alongside the full one, no extra step. +So `image: nixos` has two candidates, `image: alpine` has three, and `runs_on: [gpu]` pins to executor-c. The alpine-only image set is staged by `prepare-spindle-images.sh` (step 5) alongside the full one, no extra step. Each executor needs its own identity (one live session per token), so the `mill-tokens` service registers a token per executor in the mill db and drops it into the shared volume for the executor to read. diff --git a/localinfra/scripts/init-data.sh b/localinfra/scripts/init-data.sh index f9424a368..1c8c67809 100644 --- a/localinfra/scripts/init-data.sh +++ b/localinfra/scripts/init-data.sh @@ -10,14 +10,16 @@ SHARED_DIR="${SHARED_DIR:-/shared}" . /scripts/lib.sh +[ -f "${SHARED_DIR}/owner-did" ] || fail "no owner-did in $SHARED_DIR." \ + ' init-accounts.sh writes it into the init-state volume, so it either never ran or failed.' OWNER_DID=$(cat "${SHARED_DIR}/owner-did") OWNER_JWT=$(login "$OWNER_DID") -BOB_DID=$(resolve_handle "bob.${PDS_HOSTNAME}") +BOB_DID=$(require_handle "bob.${PDS_HOSTNAME}") BOB_JWT=$(login "$BOB_DID") -CHARLIE_DID=$(resolve_handle "charlie.${PDS_HOSTNAME}") -DAVID_DID=$(resolve_handle "david.${PDS_HOSTNAME}") +CHARLIE_DID=$(require_handle "charlie.${PDS_HOSTNAME}") +DAVID_DID=$(require_handle "david.${PDS_HOSTNAME}") CORE_SOURCE=https://knot1.tangled.sh/did:plc:ioighzh4jnrhybayw65gfbia @@ -67,9 +69,10 @@ VOUCH_BOB_DAVID=at://$BOB_DID/sh.tangled.graph.vouch/$DAVID_DID ( # idempotent: the knot no-ops a subject that is already a member or an admin for did in "$BOB_DID" "$CHARLIE_DID" "$DAVID_DID"; do + token=$(service_token "$OWNER_JWT" sh.tangled.knot.addMember "$KNOT_DID") curl -fsS -o /dev/null -X POST \ -H "Content-Type: application/json" \ - -H "Authorization: Bearer $(service_token "$OWNER_JWT" sh.tangled.knot.addMember "$KNOT_DID")" \ + -H "Authorization: Bearer $token" \ -d "$(jq -nc --arg subject "$did" '{subject:$subject}')" \ "${KNOT_URL}/xrpc/sh.tangled.knot.addMember" printf '[knot2] member %s\n' "$did" >&2 @@ -80,12 +83,42 @@ VOUCH_BOB_DAVID=at://$BOB_DID/sh.tangled.graph.vouch/$DAVID_DID # a fork adopts the source's object format, so this lands as sha1 despite knot2 # defaulting to sha256. -REPO_DID=$(curl -fsS \ - -H "Content-Type: application/json" \ - -H "Authorization: Bearer $(service_token "$OWNER_JWT" sh.tangled.repo.create "$KNOT_DID")" \ - -d "$(jq -nc --arg source "$CORE_SOURCE" \ - '{rkey:"core", name:"core", defaultBranch:"master", source:$source}')" \ - "${KNOT_URL}/xrpc/sh.tangled.repo.create" | jq -er '.repoDid') +attempt=1 +while :; do + token=$(service_token "$OWNER_JWT" sh.tangled.repo.create "$KNOT_DID") + resp=$(curl -sS -w '\n%{http_code}' \ + -H "Content-Type: application/json" \ + -H "Authorization: Bearer $token" \ + -d "$(jq -nc --arg source "$CORE_SOURCE" \ + '{rkey:"core", name:"core", defaultBranch:"master", source:$source}')" \ + "${KNOT_URL}/xrpc/sh.tangled.repo.create") + body=$(printf '%s\n' "$resp" | sed '$d') + status=$(printf '%s\n' "$resp" | tail -n1) + [ "$status" = 503 ] && [ "$attempt" -lt 30 ] || break + printf '[knot] registry projection warming, retrying repo.create (%s/30)\n' "$attempt" >&2 + attempt=$((attempt + 1)) + sleep 2 +done + +case "$status" in + 200) REPO_DID=$(printf '%s\n' "$body" | jq -er '.repoDid') ;; + 409) + REPO_DID=$(curl -fsS \ + "${PDS_URL}/xrpc/com.atproto.repo.getRecord?repo=${OWNER_DID}&collection=sh.tangled.repo&rkey=core" \ + | jq -er '.value.repoDid') || fail \ + "core exists on the knot, and the pds doesn't have a core record with a repoDid in it." \ + " The knot volume outlived the pds one, so the seed can't recover the repo DID." \ + ' Wipe both with "compose down -v", or delete core on the knot and rerun.' + refs_status=$(curl -sS -o /dev/null -w '%{http_code}' \ + "${KNOT_URL}/xrpc/sh.tangled.git.listRefs?repo=${REPO_DID}&limit=1") + [ "$refs_status" = 200 ] || fail \ + "the pds has $REPO_DID as core, and the knot answered HTTP $refs_status for its refs." \ + " The pds volume outlived the knot one, so that repo DID doesn't point at a repository." \ + ' Wipe both with "compose down -v".' + ;; + 503) fail "the knot stayed warming through a minute of repo.create retries: $body" ;; + *) fail "repo.create core: HTTP $status: $body" ;; +esac printf '[knot] core = %s\n' "$REPO_DID" >&2 @@ -96,9 +129,10 @@ put_record "$REPO_CORE" "$(jq -nc \ --arg createdAt "$CREATED_AT" \ '{knot:$knot, name:"core", repoDid:$repoDid, source:$source, createdAt:$createdAt}')" >/dev/null +token=$(service_token "$BOB_JWT" sh.tangled.repo.create "$KNOT_DID") EMPTY_REPO_DID=$(curl -fsS \ -H "Content-Type: application/json" \ - -H "Authorization: Bearer $(service_token "$BOB_JWT" sh.tangled.repo.create "$KNOT_DID")" \ + -H "Authorization: Bearer $token" \ -d '{"rkey":"empty-repo", "name":"empty-repo", "defaultBranch":"master"}' \ "${KNOT_URL}/xrpc/sh.tangled.repo.create" | jq -er '.repoDid') @@ -180,22 +214,29 @@ PULL_HEAD_1=f3c356fc8da46818a8ab0c24f22ee61854fe4b5f PULL_HEAD_2=9a925efef6a0e6dfcd2d4317b4a1eee8752928b8 patch_blob() { + pb_patch=$(mktemp) curl -fsS "${KNOT_URL}/xrpc/sh.tangled.repo.compare?repo=${REPO_DID}&rev1=${PULL_BASE}&rev2=$1" \ - | jq -j '.patch' \ - | gzip \ - | curl -fsS --data-binary @- \ - -H "Content-Type: application/gzip" \ - -H "Authorization: Bearer ${OWNER_JWT}" \ - "${PDS_URL}/xrpc/com.atproto.repo.uploadBlob" \ - | jq -ec '.blob' + | jq -je '.patch | select(length > 0)' >"$pb_patch" || fail \ + "repo.compare $PULL_BASE..$1 didn't give back any patch text." \ + ' An empty patch uploads as a valid blob, so the seeded rounds would show an empty diff.' \ + " Both revisions have to be in the clone the knot took from $CORE_SOURCE." + gzip -c "$pb_patch" | curl -fsS --data-binary @- \ + -H "Content-Type: application/gzip" \ + -H "Authorization: Bearer ${OWNER_JWT}" \ + "${PDS_URL}/xrpc/com.atproto.repo.uploadBlob" \ + | jq -ec '.blob' || fail "uploadBlob for the $1 patch failed." + rm -f "$pb_patch" } +ROUND_1=$(patch_blob "$PULL_HEAD_1") +ROUND_2=$(patch_blob "$PULL_HEAD_2") + put_record "$PULL_1" "$(jq -nc \ --arg repo "$REPO_DID" \ --arg title 'legacy PR' \ --arg body 'two rounds, each a gzipped format-patch blob' \ - --argjson round1 "$(patch_blob "$PULL_HEAD_1")" \ - --argjson round2 "$(patch_blob "$PULL_HEAD_2")" \ + --argjson round1 "$ROUND_1" \ + --argjson round2 "$ROUND_2" \ --arg createdAt1 '2025-09-22T10:40:35Z' \ --arg createdAt2 '2025-09-22T10:41:35Z' \ '{title:$title, body:$body, createdAt:$createdAt1, @@ -205,18 +246,22 @@ put_record "$PULL_1" "$(jq -nc \ {patchBlob:$round2, createdAt:$createdAt2}]}')" >/dev/null keep_commit() { - did=${3#at://} - did=${did%%/*} + kc_did=${3#at://} + kc_did=${kc_did%%/*} + kc_token=$(service_token "$(login "$kc_did")" sh.tangled.git.keepCommit "$KNOT_DID") - curl -fsS \ + kc_kept=$(curl -sS -w '\n%{http_code}' \ -H "Content-Type: application/json" \ - -H "Authorization: Bearer $(service_token "$(login "$did")" sh.tangled.git.keepCommit "$KNOT_DID")" \ + -H "Authorization: Bearer $kc_token" \ -d "$(jq -nc --arg repo "$1" --arg oid "$2" --arg record "$3" \ '{repo:$repo, record:$record, source:{"$type":"sh.tangled.git.keepCommit#commit", repo:$repo, oid:$oid}}')" \ - "${KNOT_URL}/xrpc/sh.tangled.git.keepCommit" >/dev/null \ - && printf '[keep] %s %s\n' "$2" "$3" >&2 \ - || printf '[keep] %s failed, knot has no keepCommit\n' "$2" >&2 + "${KNOT_URL}/xrpc/sh.tangled.git.keepCommit") + kc_status=$(printf '%s\n' "$kc_kept" | tail -n1) + [ "$kc_status" = 200 ] || fail \ + "keepCommit $2 for $3: HTTP $kc_status: $(printf '%s\n' "$kc_kept" | sed '$d')" \ + ' The record cites this commit, and a commit the knot never kept is one it may gc.' + printf '[keep] %s %s\n' "$2" "$3" >&2 } keep_commit "$REPO_DID" "$PULL_HEAD_1" "$PULL_2" diff --git a/localinfra/scripts/lib.sh b/localinfra/scripts/lib.sh index ff71059c2..69122c6f5 100644 --- a/localinfra/scripts/lib.sh +++ b/localinfra/scripts/lib.sh @@ -5,19 +5,32 @@ PASSWORD="${PASSWORD:-password}" CREATED_AT="${CREATED_AT:-2025-09-22T11:14:35+01:00}" +# fail LINE... -> every line on stderr, then abort +fail() { + printf '%s\n' "$@" >&2 + exit 1 +} + # resolve_handle HANDLE → DID on stdout resolve_handle() { - resp=$(curl -sS -w '\n%{http_code}' \ + rh_resp=$(curl -sS -w '\n%{http_code}' \ "${PDS_URL}/xrpc/com.atproto.identity.resolveHandle?handle=$1") - body=$(printf '%s\n' "$resp" | sed '$d') - status=$(printf '%s\n' "$resp" | tail -n1) - case "$status" in - 200) printf '%s\n' "$body" | jq -er '.did' ;; + rh_body=$(printf '%s\n' "$rh_resp" | sed '$d') + rh_status=$(printf '%s\n' "$rh_resp" | tail -n1) + case "$rh_status" in + 200) printf '%s\n' "$rh_body" | jq -er '.did' ;; 400) : ;; # not found — expected - *) printf 'resolveHandle %s: HTTP %s: %s\n' "$1" "$status" "$body" >&2; return 1 ;; + *) fail "resolveHandle $1: HTTP $rh_status: $rh_body" ;; esac } +require_handle() { + rq_did=$(resolve_handle "$1") + [ -n "$rq_did" ] || fail "resolveHandle $1: the pds doesn't know that handle." \ + ' init-accounts.sh creates the seed accounts, so it either never ran or failed.' + printf '%s\n' "$rq_did" +} + # login DID/Handle → access JWT on stdout login() { curl -fsS -H "Content-Type: application/json" \ @@ -29,10 +42,10 @@ login() { # service_token JWT LXM [AUD] → service auth JWT for the knot, on stdout. # AUD defaults to the did:web of $KNOT_HOSTNAME. service_token() { - aud=${3:-"did:web:$(printf '%s' "${KNOT_HOSTNAME:?KNOT_HOSTNAME must be set}" | sed 's/:/%3A/g')"} + st_aud=${3:-"did:web:$(printf '%s' "${KNOT_HOSTNAME:?KNOT_HOSTNAME must be set}" | sed 's/:/%3A/g')"} curl -fsS -H "Authorization: Bearer $1" \ - "${PDS_URL}/xrpc/com.atproto.server.getServiceAuth?aud=${aud}&lxm=$2&exp=$(( $(date +%s) + 240 ))" \ + "${PDS_URL}/xrpc/com.atproto.server.getServiceAuth?aud=${st_aud}&lxm=$2&exp=$(( $(date +%s) + 240 ))" \ | jq -er '.token' } @@ -41,33 +54,33 @@ service_token() { # callers that don't need it can ignore stdout. put_record() { case "$1" in - at://did:*/*/*) ;; - *) printf 'put_record: expected at://DID/COLLECTION/RKEY, got %s\n' "$1" >&2; return 1 ;; + at://did:*/*/?*) ;; + *) fail "put_record: expected at://DID/COLLECTION/RKEY, got $1" ;; esac - rest=${1#at://} - did=${rest%%/*} - rest=${rest#*/} - collection=${rest%%/*} - rkey=${rest#*/} - record="$2" + pr_rest=${1#at://} + pr_did=${pr_rest%%/*} + pr_rest=${pr_rest#*/} + pr_collection=${pr_rest%%/*} + pr_rkey=${pr_rest#*/} + pr_record="$2" # one login per record — the dev PDS runs with rate limits disabled - jwt=$(login "$did") + pr_jwt=$(login "$pr_did") - payload=$(jq -nc \ - --arg repo "$did" \ - --arg collection "$collection" \ - --arg rkey "$rkey" \ - --argjson record "$record" \ + pr_payload=$(jq -nc \ + --arg repo "$pr_did" \ + --arg collection "$pr_collection" \ + --arg rkey "$pr_rkey" \ + --argjson record "$pr_record" \ '{repo:$repo, collection:$collection, rkey:$rkey, record:$record}') - resp=$(curl -fsS \ + pr_resp=$(curl -fsS \ -H "Content-Type: application/json" \ - -H "Authorization: Bearer ${jwt}" \ - -d "$payload" \ + -H "Authorization: Bearer ${pr_jwt}" \ + -d "$pr_payload" \ "${PDS_URL}/xrpc/com.atproto.repo.putRecord") printf '[record] %s\n' "$1" >&2 - printf '%s\n' "$resp" | jq -er '.cid' # -e: fail loudly if the cid is absent + printf '%s\n' "$pr_resp" | jq -er '.cid' # -e: fail loudly if the cid is absent } diff --git a/localinfra/scripts/prepare-spindle-images.sh b/localinfra/scripts/prepare-spindle-images.sh index 90fedfdf9..9cd29eb65 100755 --- a/localinfra/scripts/prepare-spindle-images.sh +++ b/localinfra/scripts/prepare-spindle-images.sh @@ -4,6 +4,13 @@ set -euo pipefail repo=$(cd "$(dirname "$0")/../.." && pwd) image_root="${1:-$repo/out/localinfra-spindle-images}" +case "$image_root" in + /) + printf 'refusing / as the image root: this script deletes paths under it.\n' >&2 + exit 1 + ;; +esac + mkdir -p "$image_root" extract_image() { @@ -15,13 +22,13 @@ extract_image() { tarball=$(nix build "$repo#$package" --no-link --print-out-paths) [ -d "$image_root/$name" ] && chmod -R +w "$image_root/$name" || true - rm -rf "$image_root/$name" + rm -rf "${image_root:?}/$name" mkdir -p "$image_root/$name" tar -C "$image_root/$name" -xzf "$tarball" local alias for alias in "$@"; do - rm -rf "$image_root/$alias" + rm -rf "${image_root:?}/$alias" ln -s "$name" "$image_root/$alias" done } diff --git a/localinfra/scripts/warm-caches.sh b/localinfra/scripts/warm-caches.sh new file mode 100755 index 000000000..5406d019a --- /dev/null +++ b/localinfra/scripts/warm-caches.sh @@ -0,0 +1,54 @@ +#!/usr/bin/env bash +set -euo pipefail + +REPO_ROOT="$(cd "$(dirname "$0")/../.." && pwd)" +RAW_PROJECT="${COMPOSE_PROJECT_NAME:-$(basename "$REPO_ROOT")}" +PROJECT="$(printf '%s' "$RAW_PROJECT" | tr '[:upper:]' '[:lower:]' | tr -cd 'a-z0-9_-')" + +[ -n "$PROJECT" ] || { + printf "compose can't use %s as a project name.\n" "$RAW_PROJECT" >&2 + printf ' Set COMPOSE_PROJECT_NAME to something matching [a-z0-9_-] and run this again.\n' >&2 + exit 1 +} +[ "$PROJECT" = "$RAW_PROJECT" ] || + printf 'addressing volumes and images as %s, which is what compose normalizes %s to.\n' \ + "$PROJECT" "$RAW_PROJECT" >&2 + +failed=() + +warm() { + local name="$1" + shift + "$@" || { printf 'warming %s failed.\n' "$name" >&2; failed+=("$name"); } +} + +warm appview podman run --rm -v "$REPO_ROOT:/src" -w /src \ + -v "${PROJECT}_go-cache:/go/cache" -v "${PROJECT}_go-mod-cache:/go/mod" \ + -e CGO_ENABLED=1 -e GOCACHE=/go/cache -e GOMODCACHE=/go/mod \ + --entrypoint sh "localhost/${PROJECT}_appview:latest" \ + -c 'go mod download && go build ./cmd/appview' + +while read -r svc dir vol; do + warm "$svc" podman run --rm -v "$REPO_ROOT/$svc:/src" -w /src \ + -v "${PROJECT}_${svc}-node-modules:/src/node_modules" \ + -v "${PROJECT}_${vol}:/src/${dir}" \ + --entrypoint sh "localhost/${PROJECT}_${svc}:latest" \ + -c 'corepack pnpm install --frozen-lockfile --store-dir /src/node_modules/.pnpm-store' &2 + printf ' Their containers fetch at start-up, and the tngl network is internal, so\n' >&2 + printf ' "compose up" will loop on a registry or crates.io fetch until this succeeds.\n' >&2 + exit 1 +} diff --git a/localinfra/seed.Dockerfile b/localinfra/seed.Dockerfile new file mode 100644 index 000000000..78ae190c1 --- /dev/null +++ b/localinfra/seed.Dockerfile @@ -0,0 +1,2 @@ +FROM alpine:3.22 +RUN apk add --no-cache curl jq diff --git a/localinfra/spindle.Dockerfile b/localinfra/spindle.Dockerfile index eec929835..f5b341f81 100644 --- a/localinfra/spindle.Dockerfile +++ b/localinfra/spindle.Dockerfile @@ -1,6 +1,6 @@ # Development only. Not for production use. -FROM golang:1.25-alpine AS builder +FROM golang:1.26-alpine AS builder RUN apk add --no-cache git build-base sqlite-dev diff --git a/localinfra/web.Dockerfile b/localinfra/web.Dockerfile index 6eafb6392..6579ca37c 100644 --- a/localinfra/web.Dockerfile +++ b/localinfra/web.Dockerfile @@ -1,7 +1,7 @@ # Development only. Not for production use. FROM node:24-slim -RUN corepack enable && corepack install -g pnpm@11.10.0 +RUN corepack enable && corepack install -g pnpm@11.17.0 COPY <<'EOF' /usr/local/bin/web-entrypoint.sh #!/bin/sh diff --git a/localinfra/zoekt-tngl-indexserver.Dockerfile b/localinfra/zoekt-tngl-indexserver.Dockerfile index 4f7ec0cf1..b0ad07a07 100644 --- a/localinfra/zoekt-tngl-indexserver.Dockerfile +++ b/localinfra/zoekt-tngl-indexserver.Dockerfile @@ -1,6 +1,6 @@ # Development only. Not for production use. -FROM golang:1.25-alpine AS builder +FROM golang:1.26-alpine AS builder RUN apk add --no-cache git