diff --git a/Cargo.lock b/Cargo.lock index 1a24f04..36b8cc9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4001,7 +4001,6 @@ dependencies = [ "gix-hash", "gix-pack", "knot-config", - "knot-events", "knot-fixtures", "knot-git", "knot-lfs", @@ -4320,7 +4319,6 @@ name = "knot-workflow" version = "0.1.0" dependencies = [ "globset", - "knot-events", "knot-types", "serde", "serde_json", diff --git a/crates/knot-config/src/lib.rs b/crates/knot-config/src/lib.rs index 6a6d893..6c90106 100644 --- a/crates/knot-config/src/lib.rs +++ b/crates/knot-config/src/lib.rs @@ -42,6 +42,9 @@ pub struct KnotConfig { pub resources: ResourcesConfig, #[config(nested)] pub homepage: HomepageConfig, + + #[config(nested)] + pub ci: CiConfig, #[config(nested)] pub messages: knot_messages::MessagesConfig, } @@ -148,6 +151,12 @@ pub struct RepoConfig { pub default_branch: String, } +#[derive(Debug, Config)] +pub struct CiConfig { + #[config(env = "KNOT_CI_LOGS_ADDR")] + pub logs_addr: Option, +} + #[derive(Debug, Config)] pub struct HomepageConfig { #[config(env = "KNOT_HOMEPAGE_ENABLED", default = true)] @@ -294,6 +303,9 @@ pub struct XrpcConfig { #[config(env = "KNOT_XRPC_EVENTS_REPLAY_BUFFER", default = 4096)] pub events_replay_buffer: u32, + #[config(env = "KNOT_XRPC_EVENTS_REPLAY_BYTES", default = 67_108_864)] + pub events_replay_bytes: u64, + #[config(env = "KNOT_XRPC_EVENTS_MAX_SUBSCRIBERS", default = 256)] pub events_max_subscribers: u32, @@ -691,6 +703,10 @@ impl KnotConfig { self.xrpc.events_replay_buffer > 0, "xrpc.events_replay_buffer must be greater than zero", ), + check( + self.xrpc.events_replay_bytes > 0, + "xrpc.events_replay_bytes must be greater than zero", + ), check( self.xrpc.events_max_subscribers > 0, "xrpc.events_max_subscribers must be greater than zero", @@ -1161,6 +1177,7 @@ mod tests { fork_fetch_timeout_ms: 600_000, trusted_proxy_header: None, events_replay_buffer: 4_096, + events_replay_bytes: 67_108_864, events_max_subscribers: 256, events_max_per_peer: 8, }, @@ -1206,6 +1223,9 @@ mod tests { enabled: true, path: None, }, + ci: CiConfig { + logs_addr: Some("logs.oyster.cafe:3333".to_string()), + }, } } @@ -1448,6 +1468,11 @@ mod tests { |config| config.xrpc.events_replay_buffer = 0, "events_replay_buffer", ), + ( + "zero_events_replay_bytes", + |config| config.xrpc.events_replay_bytes = 0, + "events_replay_bytes", + ), ( "zero_events_max_subscribers", |config| config.xrpc.events_max_subscribers = 0, diff --git a/crates/knot-events/src/lib.rs b/crates/knot-events/src/lib.rs index e9328dc..edf5ccf 100644 --- a/crates/knot-events/src/lib.rs +++ b/crates/knot-events/src/lib.rs @@ -1,20 +1,18 @@ use std::collections::{BTreeSet, HashMap, VecDeque}; -use std::fmt::Display; use std::net::IpAddr; use std::sync::{Arc, Mutex}; +use serde::ser::SerializeMap; use serde::{Serialize, Serializer}; use serde_json::Value; use tokio::sync::{OwnedSemaphorePermit, Semaphore, watch}; use knot_runtime::{Clock, UnixMicros}; use knot_types::{ - AccountDid, BranchName, Email, KnotHostname, LanguageBytes, LanguageName, ObjectFormat, Oid, - OwnerDid, ParseError, RefName, RefTransition, RepoDid, RepoRkey, Tid, UnixSeconds, + AccountDid, ChangedFiles, Email, LanguageBytes, LanguageName, ObjectFormat, Oid, OwnerDid, + PushOptions, RefName, RefTransition, RepoDid, RepoPath, Tid, }; -pub const RECONSTRUCT_WINDOW_SECS: i64 = 30 * 24 * 60 * 60; - #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize)] #[serde(transparent)] pub struct EventCursor(i64); @@ -48,38 +46,56 @@ pub struct Event { pub created: EventCursor, } -// `sh.tangled.git.refUpdate` requires ref + newSha & oldSha, -// and `sh.tangled.pipeline` requires defaultBranch, -// so the ones we mightn't know use "" rather than `null`. -fn or_empty( - value: &Option, - serializer: S, -) -> Result { - match value { - Some(inner) => serializer.collect_str(inner), - None => serializer.serialize_str(""), - } +// `sh.tangled.git.refUpdate` requires ref, oldSha and newSha, +// so a record about no ref at all sends "" for the three rather than `null`. +#[derive(Debug, Clone)] +enum RefChange { + Absent, + Applied { + ref_name: RefName, + old_sha: Oid, + new_sha: Oid, + }, } -fn via_display(value: &T, serializer: S) -> Result { - serializer.collect_str(value) +impl Serialize for RefChange { + fn serialize(&self, serializer: S) -> Result { + let mut map = serializer.serialize_map(Some(3))?; + match self { + Self::Absent => { + map.serialize_entry("ref", "")?; + map.serialize_entry("oldSha", "")?; + map.serialize_entry("newSha", "")?; + } + Self::Applied { + ref_name, + old_sha, + new_sha, + } => { + map.serialize_entry("ref", ref_name.as_str())?; + map.serialize_entry("oldSha", &old_sha.to_hex())?; + map.serialize_entry("newSha", &new_sha.to_hex())?; + } + } + map.end() + } } #[derive(Debug, Clone, Serialize)] pub struct GitRefUpdate { #[serde(rename = "$type")] record_type: &'static str, + #[serde(rename = "changedFiles", skip_serializing_if = "Vec::is_empty")] + changed_files: Vec, #[serde(rename = "committerDid")] committer_did: AccountDid, meta: Option, - #[serde(rename = "newSha", serialize_with = "or_empty")] - new_sha: Option, - #[serde(rename = "oldSha", serialize_with = "or_empty")] - old_sha: Option, #[serde(rename = "ownerDid", skip_serializing_if = "Option::is_none")] owner_did: Option, - #[serde(rename = "ref", serialize_with = "or_empty")] - ref_name: Option, + #[serde(rename = "pushOptions", skip_serializing_if = "PushOptions::is_empty")] + push_options: PushOptions, + #[serde(flatten)] + change: RefChange, repo: RepoDid, } @@ -87,20 +103,37 @@ impl GitRefUpdate { pub fn new(repo: RepoDid, owner: Option, committer: AccountDid) -> Self { Self { record_type: Self::NSID, + changed_files: Vec::new(), committer_did: committer, meta: None, - new_sha: None, - old_sha: None, owner_did: owner, - ref_name: None, + push_options: PushOptions::default(), + change: RefChange::Absent, repo, } } - pub fn on_ref(mut self, ref_name: RefName, transition: RefTransition) -> Self { - self.ref_name = Some(ref_name); - self.old_sha = transition.old_oid(); - self.new_sha = transition.new_oid(); + pub fn on_ref( + mut self, + ref_name: RefName, + transition: RefTransition, + format: ObjectFormat, + ) -> Self { + self.change = RefChange::Applied { + ref_name, + old_sha: transition.old_oid().unwrap_or_else(|| format.null_oid()), + new_sha: transition.new_oid().unwrap_or_else(|| format.null_oid()), + }; + self + } + + pub fn with_changed_files(mut self, changed: ChangedFiles) -> Self { + self.changed_files = changed.into_paths(); + self + } + + pub fn with_push_options(mut self, options: &PushOptions) -> Self { + self.push_options = options.clone(); self } @@ -196,207 +229,6 @@ impl LanguageSize { } } -#[derive(Debug, Clone, Serialize)] -pub struct Pipeline { - #[serde(rename = "$type")] - record_type: &'static str, - #[serde(rename = "triggerMetadata")] - trigger_metadata: TriggerMetadata, - workflows: Vec, -} - -impl Pipeline { - pub fn new(trigger_metadata: TriggerMetadata, workflows: Vec) -> Self { - Self { - record_type: Self::NSID, - trigger_metadata, - workflows, - } - } -} - -impl Publish for Pipeline { - const NSID: &'static str = "sh.tangled.pipeline"; -} - -#[derive(Debug, Clone, Copy, Serialize)] -#[serde(rename_all = "lowercase")] -enum TriggerKind { - Push, -} - -#[derive(Debug, Clone, Serialize)] -pub struct TriggerMetadata { - kind: TriggerKind, - #[serde(skip_serializing_if = "Option::is_none")] - push: Option, - repo: TriggerRepo, -} - -impl TriggerMetadata { - pub fn push(push: PushTriggerData, repo: TriggerRepo) -> Self { - Self { - kind: TriggerKind::Push, - push: Some(push), - repo, - } - } -} - -#[derive(Debug, Clone, Serialize)] -pub struct PushTriggerData { - #[serde(rename = "ref", serialize_with = "via_display")] - ref_name: RefName, - #[serde(rename = "newSha", serialize_with = "via_display")] - new_sha: Oid, - #[serde(rename = "oldSha", serialize_with = "via_display")] - old_sha: Oid, -} - -impl PushTriggerData { - // There's nothing for a pipeline to run against on deletion, - // so it doesn't trigger. - // On creation there's a trigger with the 0-oid - // where the old sha would be, the same way a create comes - // in the first place. - pub fn new(ref_name: RefName, transition: RefTransition, format: ObjectFormat) -> Option { - transition.new_oid().map(|new_sha| Self { - ref_name, - new_sha, - old_sha: transition.old_oid().unwrap_or_else(|| format.null_oid()), - }) - } -} - -#[derive(Debug, Clone, Serialize)] -pub struct TriggerRepo { - knot: KnotHostname, - did: OwnerDid, - #[serde(rename = "repoDid", skip_serializing_if = "Option::is_none")] - repo_did: Option, - #[serde(skip_serializing_if = "Option::is_none")] - repo: Option, - #[serde(rename = "defaultBranch", serialize_with = "or_empty")] - default_branch: Option, -} - -impl TriggerRepo { - pub fn new( - knot: KnotHostname, - did: OwnerDid, - repo_did: Option, - repo: Option, - default_branch: Option, - ) -> Self { - Self { - knot, - did, - repo_did, - repo, - default_branch, - } - } -} - -#[derive(Debug, Clone, PartialEq, Eq, Serialize)] -#[serde(transparent)] -pub struct WorkflowName(String); - -impl WorkflowName { - pub fn new(value: impl Into) -> Result { - let value = value.into(); - let valid = - !value.is_empty() && !value.contains('/') && !value.chars().any(char::is_control); - match valid { - true => Ok(Self(value)), - false => Err(ParseError::Invalid { - kind: "workflow name", - value, - }), - } - } - - pub fn as_str(&self) -> &str { - &self.0 - } -} - -#[derive(Debug, Clone, Serialize)] -#[serde(transparent)] -pub struct EngineRef(String); - -impl EngineRef { - pub fn new(value: impl Into) -> Result { - let value = value.into(); - let valid = - !value.is_empty() && !value.chars().any(|c| c.is_whitespace() || c.is_control()); - match valid { - true => Ok(Self(value)), - false => Err(ParseError::Invalid { - kind: "engine reference", - value, - }), - } - } -} - -#[derive(Debug, Clone, Serialize)] -pub struct WorkflowSpec { - name: WorkflowName, - engine: EngineRef, - clone: CloneSpec, - raw: String, -} - -impl WorkflowSpec { - pub fn new(name: WorkflowName, engine: EngineRef, clone: CloneSpec, raw: String) -> Self { - Self { - name, - engine, - clone, - raw, - } - } -} - -#[derive(Debug, Clone, Serialize)] -pub struct CloneSpec { - pub skip: bool, - pub depth: CloneDepth, - pub submodules: bool, -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum CloneDepth { - Full, - Limited(std::num::NonZeroU32), -} - -impl CloneDepth { - pub fn new(depth: i64) -> Result { - match u32::try_from(depth) { - Ok(raw) => Ok(std::num::NonZeroU32::new(raw).map_or(Self::Full, Self::Limited)), - Err(_) => Err(ParseError::Invalid { - kind: "clone depth", - value: depth.to_string(), - }), - } - } - - pub fn is_limited(self) -> bool { - matches!(self, Self::Limited(_)) - } -} - -impl Serialize for CloneDepth { - fn serialize(&self, serializer: S) -> Result { - serializer.serialize_i64(match self { - Self::Full => 0, - Self::Limited(depth) => depth.get() as i64, - }) - } -} - #[derive(Debug, Clone, Copy, Serialize)] #[serde(rename_all = "lowercase")] enum AclOp { @@ -459,14 +291,69 @@ impl Publish for RepoCollaboratorUpdate { const NSID: &'static str = "sh.tangled.repo.collaboratorUpdate"; } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ReplayEvents(std::num::NonZeroUsize); + +impl ReplayEvents { + pub fn new(value: usize) -> Option { + std::num::NonZeroUsize::new(value).map(Self) + } + + pub fn get(self) -> usize { + self.0.get() + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ReplayBytes(std::num::NonZeroUsize); + +impl ReplayBytes { + pub fn new(value: usize) -> Option { + std::num::NonZeroUsize::new(value).map(Self) + } + + pub fn get(self) -> usize { + self.0.get() + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ReplayBounds { + events: ReplayEvents, + bytes: ReplayBytes, +} + +impl ReplayBounds { + pub fn new(events: ReplayEvents, bytes: ReplayBytes) -> Self { + Self { events, bytes } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum BatchEnd { + CaughtUp, + Bounded, +} + +pub struct Replayed { + pub events: Vec>, + pub end: BatchEnd, +} + +struct Entry { + event: Arc, + bytes: usize, +} + struct Ring { - events: VecDeque, + entries: VecDeque, + bytes: usize, last_micros: UnixMicros, pending: BTreeSet, } struct Inner { - capacity: usize, + bounds: ReplayBounds, ring: Mutex, head: watch::Sender, } @@ -492,69 +379,72 @@ impl Inner { fn stable_head(ring: &Ring) -> EventCursor { let stable = match ring.pending.iter().next().copied() { - Some(horizon) => ring.events.partition_point(|event| event.created < horizon), - None => ring.events.len(), + Some(horizon) => ring + .entries + .partition_point(|entry| entry.event.created < horizon), + None => ring.entries.len(), }; stable .checked_sub(1) - .and_then(|index| ring.events.get(index)) - .map(|event| event.created) + .and_then(|index| ring.entries.get(index)) + .map(|entry| entry.event.created) .unwrap_or(EventCursor::START) } -fn insert_sorted(ring: &mut Ring, event: Event, capacity: usize) { - let position = ring - .events - .partition_point(|existing| existing.created < event.created); - ring.events.insert(position, event); - while ring.events.len() > capacity { - ring.events.pop_front(); +fn value_bytes(value: &Value) -> usize { + let node = std::mem::size_of::(); + match value { + Value::Null | Value::Bool(_) | Value::Number(_) => node, + Value::String(text) => node + text.len(), + Value::Array(items) => node + items.iter().map(value_bytes).sum::(), + Value::Object(fields) => { + node + fields + .iter() + .map(|(key, field)| key.len() + value_bytes(field)) + .sum::() + } } } -knot_types::scalar_newtype! { - pub struct SeedOrder(u64) => ordered; -} - -pub struct Seed { - seconds: UnixSeconds, - order: SeedOrder, - nsid: &'static str, - payload: Value, +fn insert_sorted(ring: &mut Ring, event: Event, bounds: ReplayBounds) { + let bytes = std::mem::size_of::() + value_bytes(&event.payload); + let position = ring + .entries + .partition_point(|existing| existing.event.created < event.created); + ring.entries.insert( + position, + Entry { + event: Arc::new(event), + bytes, + }, + ); + ring.bytes += bytes; + evict_oldest(ring, bounds); } -impl Seed { - pub fn from_event(seconds: UnixSeconds, order: SeedOrder, payload: &P) -> Self { - Self { - seconds, - order, - nsid: P::NSID, - payload: serde_json::to_value(payload).expect("event payload serializes to JSON"), - } +fn evict_oldest(ring: &mut Ring, bounds: ReplayBounds) { + let over = ring.entries.len() > bounds.events.get() + || (ring.bytes > bounds.bytes.get() && ring.entries.len() > 1); + if let Some(evicted) = over.then(|| ring.entries.pop_front()).flatten() { + ring.bytes -= evicted.bytes; + evict_oldest(ring, bounds); } } -fn end_of_second_micros(seconds: UnixSeconds) -> UnixMicros { - UnixMicros::new( - (seconds.get().max(0) as u64) - .saturating_mul(1_000_000) - .saturating_add(999_999), - ) -} - pub struct EventLog { clock: C, inner: Arc, } impl EventLog { - pub fn new(clock: C, capacity: usize) -> Self { + pub fn new(clock: C, bounds: ReplayBounds) -> Self { Self { clock, inner: Arc::new(Inner { - capacity: capacity.max(1), + bounds, ring: Mutex::new(Ring { - events: VecDeque::new(), + entries: VecDeque::new(), + bytes: 0, last_micros: UnixMicros::new(0), pending: BTreeSet::new(), }), @@ -581,43 +471,12 @@ impl EventLog { payload, created, }, - self.inner.capacity, + self.inner.bounds, ); self.inner.settle(ring); created } - pub fn capacity(&self) -> usize { - self.inner.capacity - } - - pub fn seed(&self, seeds: Vec) { - let mut seeds = seeds; - seeds.sort_by(|left, right| { - left.seconds - .cmp(&right.seconds) - .then_with(|| left.order.cmp(&right.order)) - }); - let mut ring = self.inner.lock(); - let mut floor = UnixMicros::new(0); - seeds.into_iter().for_each(|seed| { - let micros = end_of_second_micros(seed.seconds).max(floor.next()); - floor = micros; - insert_sorted( - &mut ring, - Event { - rkey: Tid::from_time(micros.get(), 0), - nsid: seed.nsid, - payload: seed.payload, - created: EventCursor::from_unix_micros(micros), - }, - self.inner.capacity, - ); - }); - ring.last_micros = ring.last_micros.max(floor); - self.inner.settle(ring); - } - pub fn reserve(&self) -> Reservation { let mut ring = self.inner.lock(); let (micros, cursor) = self.next_cursor(&mut ring); @@ -631,18 +490,38 @@ impl EventLog { } } - pub fn replay(&self, after: EventCursor, limit: usize) -> Vec { + pub fn replay(&self, after: EventCursor, bounds: ReplayBounds) -> Replayed { let ring = self.inner.lock(); // The corresponding read side guarantee of `reserve`. let horizon = ring.pending.iter().next().copied(); - ring.events + let visible = |entry: &&Entry| { + entry.event.created > after + && horizon.is_none_or(|horizon| entry.event.created < horizon) + }; + let events: Vec> = ring + .entries .iter() - .filter(|event| { - event.created > after && horizon.is_none_or(|horizon| event.created < horizon) + .filter(visible) + .take(bounds.events.get()) + // Why the first event gets to ignore the byte bound? + // Imagine a consumer whose next event is by itself wider + // than the entire bound, right - + // every batch it requests would come back empty, + // its cursor would never advance, + // it would ask again, repeat. + // Sending that one event alone over the bound + // is the only way. + .scan(0usize, |spent, entry| { + let first = *spent == 0; + *spent += entry.bytes; + (first || *spent <= bounds.bytes.get()).then(|| Arc::clone(&entry.event)) }) - .take(limit) - .cloned() - .collect() + .collect(); + let end = match ring.entries.iter().filter(visible).nth(events.len()) { + Some(_) => BatchEnd::Bounded, + None => BatchEnd::CaughtUp, + }; + Replayed { events, end } } pub fn subscribe(&self) -> watch::Receiver { @@ -673,7 +552,7 @@ impl Reservation { payload, created: self.cursor, }, - self.inner.capacity, + self.inner.bounds, ); ring.pending.remove(&self.cursor); self.inner.settle(ring); @@ -759,13 +638,24 @@ mod tests { use knot_runtime::{ManualClock, UnixMicros}; + fn bounds(events: usize, bytes: usize) -> ReplayBounds { + ReplayBounds::new( + ReplayEvents::new(events).unwrap(), + ReplayBytes::new(bytes).unwrap(), + ) + } + fn log(capacity: usize) -> EventLog { EventLog::new( ManualClock::new(UnixMicros::new(1_700_000_000_000_000)), - capacity, + bounds(capacity, 1 << 20), ) } + fn replay(log: &EventLog, after: EventCursor, limit: usize) -> Vec> { + log.replay(after, bounds(limit, 1 << 30)).events + } + fn update() -> GitRefUpdate { GitRefUpdate::new( RepoDid::new("did:plc:limpet").unwrap(), @@ -776,14 +666,13 @@ mod tests { fn wire(log: &EventLog, payload: &P) -> serde_json::Value { log.publish(payload); - serde_json::to_value(log.replay(EventCursor::START, 1).remove(0)).unwrap() + serde_json::to_value(&*replay(log, EventCursor::START, 1).remove(0)).unwrap() } #[test] fn publish_nsids_are_valid_type_names() { [ GitRefUpdate::NSID, - Pipeline::NSID, KnotMemberUpdate::NSID, RepoCollaboratorUpdate::NSID, ] @@ -798,7 +687,7 @@ mod tests { let log = log(8); let cursors: Vec<_> = (0..3).map(|_| log.publish(&update())).collect(); assert!(cursors.windows(2).all(|pair| pair[0] < pair[1])); - let events = log.replay(EventCursor::START, 8); + let events = replay(&log, EventCursor::START, 8); let rkeys: std::collections::BTreeSet<_> = events .iter() .map(|event| event.rkey.as_str().to_string()) @@ -812,16 +701,60 @@ mod tests { let first = log.publish(&update()); log.publish(&update()); log.publish(&update()); - let replayed = log.replay(EventCursor::START, 8); + let replayed = replay(&log, EventCursor::START, 8); assert_eq!(replayed.len(), 2); assert!(replayed.iter().all(|event| event.created > first)); } + #[test] + fn a_wide_event_evicts_by_bytes_long_before_the_ring_fills() { + let wide = |count: usize| { + update().with_changed_files(fill_changed( + (0..count).map(|index| format!("crates/knot-events/src/f{index}.rs")), + )) + }; + let log = EventLog::new( + ManualClock::new(UnixMicros::new(1_700_000_000_000_000)), + bounds(1_024, 64 * 1_024), + ); + (0..16).for_each(|_| { + log.publish(&wide(512)); + }); + let replayed = replay(&log, EventCursor::START, 1_024); + assert!( + (1..16).contains(&replayed.len()), + "the byte maximum evicts before the event maximum does: {}", + replayed.len() + ); + + let one = EventLog::new( + ManualClock::new(UnixMicros::new(1_700_000_000_000_000)), + bounds(1_024, 1), + ); + let only = one.publish(&wide(512)); + assert_eq!( + replay(&one, EventCursor::START, 8) + .iter() + .map(|event| event.created) + .collect::>(), + vec![only], + "the ring keeps the one event wider than the whole byte maximum" + ); + } + + fn fill_changed(paths: impl Iterator) -> ChangedFiles { + let mut budget = knot_types::ChangedFilesBudget::new(); + let _ = paths.into_iter().try_for_each(|path| { + budget.admit(knot_types::RepoPath::new(path).expect("test path is well-formed")) + }); + budget.finish() + } + #[test] fn replay_honors_the_cursor_and_the_limit() { let log = log(8); let cursors: Vec<_> = (0..4).map(|_| log.publish(&update())).collect(); - let after_second = log.replay(cursors[1], 8); + let after_second = replay(&log, cursors[1], 8); assert_eq!( after_second .iter() @@ -829,111 +762,54 @@ mod tests { .collect::>(), cursors[2..].to_vec() ); - assert_eq!(log.replay(EventCursor::START, 2).len(), 2); - assert!(log.replay(cursors[3], 8).is_empty()); + assert_eq!(replay(&log, EventCursor::START, 2).len(), 2); + assert!(replay(&log, cursors[3], 8).is_empty()); } #[test] - fn seeded_events_reach_a_stale_cursor_consumer_after_restart() { - let log = log(8); - log.seed(vec![ - Seed::from_event(UnixSeconds::new(1_000), SeedOrder::new(0), &update()), - Seed::from_event(UnixSeconds::new(2_000), SeedOrder::new(1), &update()), - ]); - let live = log.publish(&update()); - let all = log.replay(EventCursor::START, 8); - assert_eq!( - all.len(), - 3, - "consumer at start of stream sees both reconstructed events and live one" - ); - assert!( - all[0].created < all[1].created && all[1].created < all[2].created, - "seeded events keep their derived order" + fn a_replay_batch_stops_at_the_byte_maximum_and_reports_whether_more_remains() { + let log = EventLog::new( + ManualClock::new(UnixMicros::new(1_700_000_000_000_000)), + bounds(1_024, 1 << 20), ); assert_eq!( - all[2].created, live, - "post-restart live event sorts after every reconstructed one" + log.replay(EventCursor::START, bounds(8, 1 << 20)).end, + BatchEnd::CaughtUp, + "an empty ring has nothing left to send" ); + let wide = update().with_changed_files(fill_changed( + (0..512).map(|index| format!("crates/knot-events/src/f{index}.rs")), + )); + let cursors: Vec = (0..8).map(|_| log.publish(&wide)).collect(); + + let batch = log.replay(EventCursor::START, bounds(1_024, 16 * 1_024)); assert!( - all[0].created > EventCursor::START, - "reconstructed event keeps the wall-second of the push it replays" + (1..8).contains(&batch.events.len()), + "the byte maximum stops the batch before the event maximum does: {}", + batch.events.len() ); - assert_eq!(all[0].nsid, "sh.tangled.git.refUpdate"); - let after_first = log.replay(all[0].created, 8); - assert_eq!( - after_first.len(), - 2, - "consumer whose durable cursor sits at first reconstructed event receives the rest" - ); - assert_eq!(after_first[0].created, all[1].created); - } + assert_eq!(batch.end, BatchEnd::Bounded); - #[test] - fn a_seed_for_an_old_second_is_not_redelivered_to_a_caught_up_consumer() { - let log = log(8); - let live = log.publish(&update()); - let drained = log.replay(EventCursor::START, 8); - assert_eq!(drained.len(), 1); - let cursor = drained[0].created; - assert_eq!(cursor, live); - log.seed(vec![Seed::from_event( - UnixSeconds::new(1_000), - SeedOrder::new(0), - &update(), - )]); - assert!( - log.replay(cursor, 8).is_empty(), - "seed whose push second predates caught-up consumer's cursor isn't redelivered, \ - so restart never replays whole window to consumer that already drained it" + let rest = log.replay( + batch.events.last().expect("the batch is nonempty").created, + bounds(1_024, 1 << 30), ); + assert_eq!(rest.end, BatchEnd::CaughtUp); assert_eq!( - log.replay(EventCursor::START, 8).len(), - 2, - "consumer still at start of stream recovers it" + batch.events.len() + rest.events.len(), + cursors.len(), + "the two batches together are every event, with none repeated or skipped" ); - } + let head = log.replay(cursors[7], bounds(8, 1 << 20)); + assert!(head.events.is_empty() && head.end == BatchEnd::CaughtUp); - #[test] - fn a_boundary_second_seed_reaches_a_drained_consumer() { - let log = log(8); - let live = log.publish(&update()); - let cursor = log.replay(EventCursor::START, 8)[0].created; - assert_eq!(cursor, live); - log.seed(vec![Seed::from_event( - UnixSeconds::new(1_700_000_000), - SeedOrder::new(0), - &update(), - )]); - let after = log.replay(cursor, 8); + let single = log.replay(EventCursor::START, bounds(1_024, 1)); assert_eq!( - after.len(), + single.events.len(), 1, - "seed in same wall-second as consumer's cursor lands at end of that \ - second, above cursor, so genuinely missed pipeline near restart still arrives" - ); - assert!(after[0].created > cursor); - } - - #[test] - fn a_live_event_sorts_after_a_seed_in_the_same_wall_second() { - let log = log(8); - log.seed(vec![Seed::from_event( - UnixSeconds::new(1_700_000_000), - SeedOrder::new(0), - &update(), - )]); - let live = log.publish(&update()); - let all = log.replay(EventCursor::START, 8); - assert_eq!(all.len(), 2); - assert!( - all[0].created < live, - "seeding bumps cursor floor so live event never sorts below same-second seed" - ); - assert_eq!( - all[1].created, live, - "post-seed live event sorts after seed even when clock reads earlier micro" + "an event wider than the whole batch maximum is sent alone" ); + assert_eq!(single.end, BatchEnd::Bounded); } #[test] @@ -958,9 +834,54 @@ mod tests { assert_eq!(payload["ownerDid"], "did:web:olaren.dev"); assert_eq!(payload["repo"], "did:plc:limpet"); assert_eq!(payload["meta"], serde_json::Value::Null); - assert_eq!(payload["newSha"], ""); - assert_eq!(payload["oldSha"], ""); assert_eq!(payload["ref"], ""); + assert_eq!( + payload["oldSha"], "", + "a record about no ref sends the empty sha" + ); + assert_eq!(payload["newSha"], ""); + } + + #[test] + fn an_absent_sha_of_a_transition_is_the_null_oid_of_the_repo_object_format() { + let new = Oid::from_hex(&"cd".repeat(32)).unwrap(); + let created = &wire( + &log(8), + &GitRefUpdate::new( + RepoDid::new("did:plc:limpet").unwrap(), + None, + AccountDid::new("did:plc:nel").unwrap(), + ) + .on_ref( + RefName::new("refs/heads/fresh").unwrap(), + RefTransition::Create { new }, + ObjectFormat::SHA256, + ), + )["event"]; + assert_eq!(created["oldSha"], "0".repeat(64)); + assert_eq!(created["newSha"], new.to_hex()); + + let old = Oid::from_hex(&"ab".repeat(20)).unwrap(); + let rebuilt = update() + .on_ref( + RefName::new("refs/heads/fresh").unwrap(), + RefTransition::Create { + new: Oid::from_hex(&"cd".repeat(20)).unwrap(), + }, + ObjectFormat::SHA1, + ) + .on_ref( + RefName::new("refs/heads/gone").unwrap(), + RefTransition::Delete { old }, + ObjectFormat::SHA1, + ); + let payload = &wire(&log(8), &rebuilt)["event"]; + assert_eq!( + payload["ref"], "refs/heads/gone", + "a later transition replaces the earlier one whole" + ); + assert_eq!(payload["oldSha"], old.to_hex()); + assert_eq!(payload["newSha"], "0".repeat(40)); } fn peer(last: u8) -> IpAddr { @@ -1068,58 +989,6 @@ mod tests { assert!(meta.get("langBreakdown").is_none()); } - #[test] - fn a_pipeline_event_matches_the_lexicon_shape() { - let log = log(8); - let trigger = TriggerMetadata::push( - PushTriggerData::new( - RefName::new("refs/heads/main").unwrap(), - RefTransition::Create { - new: Oid::from_hex("3333333333333333333333333333333333333333").unwrap(), - }, - ObjectFormat::SHA1, - ) - .unwrap(), - TriggerRepo::new( - KnotHostname::new("knot.test").unwrap(), - OwnerDid::new("did:web:olaren.dev").unwrap(), - Some(RepoDid::new("did:plc:limpet").unwrap()), - Some(RepoRkey::new("anemone").unwrap()), - Some(BranchName::new("main").unwrap()), - ), - ); - let workflows = vec![WorkflowSpec::new( - WorkflowName::new("ci.yml").unwrap(), - EngineRef::new("nixery.dev/x").unwrap(), - CloneSpec { - skip: false, - depth: CloneDepth::new(1).unwrap(), - submodules: true, - }, - "engine: nixery.dev/x\n".to_string(), - )]; - let wire = wire(&log, &Pipeline::new(trigger, workflows)); - assert_eq!(wire["nsid"], "sh.tangled.pipeline"); - let payload = &wire["event"]; - assert_eq!(payload["$type"], "sh.tangled.pipeline"); - let meta = &payload["triggerMetadata"]; - assert_eq!(meta["kind"], "push"); - assert_eq!(meta["push"]["ref"], "refs/heads/main"); - assert_eq!(meta["push"]["oldSha"], "0".repeat(40)); - assert_eq!(meta["push"]["newSha"], "3".repeat(40)); - assert_eq!(meta["repo"]["knot"], "knot.test"); - assert_eq!(meta["repo"]["did"], "did:web:olaren.dev"); - assert_eq!(meta["repo"]["repoDid"], "did:plc:limpet"); - assert_eq!(meta["repo"]["repo"], "anemone"); - assert_eq!(meta["repo"]["defaultBranch"], "main"); - let workflow = &payload["workflows"][0]; - assert_eq!(workflow["name"], "ci.yml"); - assert_eq!(workflow["engine"], "nixery.dev/x"); - assert_eq!(workflow["clone"]["depth"], 1); - assert_eq!(workflow["clone"]["submodules"], true); - assert_eq!(workflow["clone"]["skip"], false); - } - #[test] fn acl_updates_match_eventstream_shape() { let log = log(8); @@ -1137,30 +1006,30 @@ mod tests { AccountDid::new("did:plc:nel").unwrap(), RepoDid::new("did:plc:limpet").unwrap(), )); - let events = log.replay(EventCursor::START, 8); + let events = replay(&log, EventCursor::START, 8); - let member_added = serde_json::to_value(&events[0]).unwrap(); + let member_added = serde_json::to_value(&*events[0]).unwrap(); assert_eq!(member_added["nsid"], "sh.tangled.knot.memberUpdate"); assert_eq!(member_added["event"]["op"], "add"); assert_eq!(member_added["event"]["subject"], "did:plc:nel"); assert!(member_added["event"].get("$type").is_none()); - let member_removed = serde_json::to_value(&events[1]).unwrap(); + let member_removed = serde_json::to_value(&*events[1]).unwrap(); assert_eq!(member_removed["event"]["op"], "remove"); assert_eq!(member_removed["event"]["subject"], "did:plc:olaren"); - let collab_added = serde_json::to_value(&events[2]).unwrap(); + let collab_added = serde_json::to_value(&*events[2]).unwrap(); assert_eq!(collab_added["nsid"], "sh.tangled.repo.collaboratorUpdate"); assert_eq!(collab_added["event"]["op"], "add"); assert_eq!(collab_added["event"]["subject"], "did:plc:nel"); assert_eq!(collab_added["event"]["repo"], "did:plc:limpet"); assert!(collab_added["event"].get("$type").is_none()); - let collab_removed = serde_json::to_value(&events[3]).unwrap(); + let collab_removed = serde_json::to_value(&*events[3]).unwrap(); assert_eq!(collab_removed["event"]["op"], "remove"); assert_eq!(collab_removed["event"]["repo"], "did:plc:limpet"); } fn cursors(log: &EventLog) -> Vec { - log.replay(EventCursor::START, 64) + replay(log, EventCursor::START, 64) .iter() .map(|event| event.created) .collect() diff --git a/crates/knot-maintenance/Cargo.toml b/crates/knot-maintenance/Cargo.toml index 2175c3c..0261646 100644 --- a/crates/knot-maintenance/Cargo.toml +++ b/crates/knot-maintenance/Cargo.toml @@ -12,7 +12,6 @@ knot-config = { workspace = true } knot-git = { workspace = true } knot-lfs = { workspace = true } knot-pack = { workspace = true } -knot-events = { workspace = true } knot-runtime = { workspace = true } tracing = { workspace = true } gix = { workspace = true } diff --git a/crates/knot-maintenance/src/lib.rs b/crates/knot-maintenance/src/lib.rs index ae4867e..f6b36e5 100644 --- a/crates/knot-maintenance/src/lib.rs +++ b/crates/knot-maintenance/src/lib.rs @@ -2,7 +2,6 @@ use std::collections::HashSet; use std::path::PathBuf; use std::time::Duration; -use knot_events::RECONSTRUCT_WINDOW_SECS; use knot_git::{PackRefsReport, ReflogReport, Repo}; use knot_types::{Oid, UnixSeconds}; @@ -20,6 +19,8 @@ mod test_support; pub use midx::MidxStatus; pub use scheduler::{MaintenanceHandle, PushBytes, RepoSource, Scheduler}; +pub const MIN_REFLOG_RETENTION_SECS: i64 = 30 * 24 * 60 * 60; + #[derive(Debug, thiserror::Error)] pub enum MaintError { #[error("git: {0}")] @@ -151,7 +152,7 @@ impl SweepInterval { pub fn lfs_grace(gc_grace: GcGrace, reflog_retention: ReflogRetention) -> LfsGrace { let ceiling = reflog_retention .0 - .max(Duration::from_secs(RECONSTRUCT_WINDOW_SECS as u64)) + .max(Duration::from_secs(MIN_REFLOG_RETENTION_SECS as u64)) .max(LFS_GRACE_MIN); LfsGrace(gc_grace.0.clamp(LFS_GRACE_MIN, ceiling)) } @@ -256,7 +257,7 @@ pub fn run_repo( } else { PackRefsReport { packed: 0 } }; - let floor_secs = (opts.reflog_floor.get().as_secs() as i64).max(RECONSTRUCT_WINDOW_SECS); + let floor_secs = (opts.reflog_floor.get().as_secs() as i64).max(MIN_REFLOG_RETENTION_SECS); let reflog = repo.expire_reflogs(now_seconds.saturating_sub_secs(floor_secs))?; let commit_graph = if opts.commit_graph { @@ -357,8 +358,7 @@ fn collect_roots(repo: &Repo, retention_floor: UnixSeconds) -> Result = ["pipeline compiled with no diagnostics"], + pipeline_none: Lines = ["no pipelines to compile"], + ci_logs: Lines = [ + "-> Browse CI logs in your terminal:", + " ssh -t -p {port} {host} {repo} {sha}" + ], } } diff --git a/crates/knot-server/src/main.rs b/crates/knot-server/src/main.rs index 3b29fa4..4c8370d 100644 --- a/crates/knot-server/src/main.rs +++ b/crates/knot-server/src/main.rs @@ -1,5 +1,4 @@ mod allocator; -mod reconstruct; #[global_allocator] static GLOBAL: tikv_jemallocator::Jemalloc = tikv_jemallocator::Jemalloc; @@ -24,10 +23,10 @@ use axum::routing::get; use base64::Engine; use knot_atproto::Atproto; use knot_config::HomepageSource; -use knot_index::{Index, Resolved}; +use knot_index::Index; use knot_runtime::{Clock, HttpTransport, OsEntropy, ReqwestHttp, SystemClock}; use knot_secrets::{MasterKey, SealedStore}; -use knot_types::{ActorId, AuthorName, BranchName, Email, KnotHostname, ObjectCount, UnixSeconds}; +use knot_types::{ActorId, AuthorName, BranchName, CiLogsAddr, Email, KnotHostname, ObjectCount}; use knot_xrpc::XrpcState; use tower_http::services::ServeFile; @@ -346,10 +345,13 @@ async fn main() -> anyhow::Result<()> { knot_xrpc::PerActorQuota::new(config.xrpc.per_actor_reservations as usize), knot_xrpc::GlobalQuota::new(config.xrpc.max_pending_reservations as usize), )); - let events = Arc::new(knot_events::EventLog::new( - SystemClock, - config.xrpc.events_replay_buffer as usize, - )); + let replay_bounds = knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(config.xrpc.events_replay_buffer as usize) + .context("xrpc.events_replay_buffer must be greater than zero")?, + knot_events::ReplayBytes::new(config.xrpc.events_replay_bytes as usize) + .context("xrpc.events_replay_bytes must be greater than zero")?, + ); + let events = Arc::new(knot_events::EventLog::new(SystemClock, replay_bounds)); let subscriber_gate = Arc::new(knot_events::SubscriberGate::new( knot_events::GlobalSubscriberLimit::new(config.xrpc.events_max_subscribers as usize), knot_events::PerPeerSubscriberLimit::new(config.xrpc.events_max_per_peer as usize), @@ -400,32 +402,13 @@ async fn main() -> anyhow::Result<()> { ) as usize), }; - let recon_index = Arc::clone(&index); - let recon_layout = layout.clone(); - let recon_knot = hostname.clone(); - let recon_events = Arc::clone(&events); - tokio::task::spawn_blocking(move || { - let repos: Vec = recon_index - .hosted_repos() - .into_iter() - .filter_map( - |repo| match (recon_index.owner_of(&repo), recon_index.rkey_of(&repo)) { - (Resolved::Ready(Some(owner)), Resolved::Ready(Some(rkey))) => { - Some(reconstruct::RepoContext { repo, owner, rkey }) - } - _ => None, - }, - ) - .collect(); - let now_seconds = - UnixSeconds::new((SystemClock.now_unix_micros().get() / 1_000_000) as i64); - let capacity = recon_events.capacity(); - let seeds = - reconstruct::pipeline_seeds(&recon_layout, &recon_knot, &repos, now_seconds, capacity); - if !seeds.is_empty() { - recon_events.seed(seeds); - } - }); + let ci_logs = config + .ci + .logs_addr + .as_deref() + .map(CiLogsAddr::new) + .transpose() + .context("ci.logs_addr must be host:port")?; let homepage = config.homepage.source(); let catalog = Arc::new( @@ -474,6 +457,7 @@ async fn main() -> anyhow::Result<()> { admission, byte_limits.pack, budgets.languages_push, + ci_logs.clone(), ) .with_maintenance(maintenance_handle.clone()) .with_limits(pack_limits) @@ -490,6 +474,7 @@ async fn main() -> anyhow::Result<()> { atproto: Arc::clone(&atproto), secrets, entropy: Arc::new(OsEntropy), + ci_logs, admins, admission, knot_did, diff --git a/crates/knot-server/src/reconstruct.rs b/crates/knot-server/src/reconstruct.rs deleted file mode 100644 index 5ebbb59..0000000 --- a/crates/knot-server/src/reconstruct.rs +++ /dev/null @@ -1,232 +0,0 @@ -use knot_events::{RECONSTRUCT_WINDOW_SECS, Seed, SeedOrder}; -use knot_git::{Layout, ReflogUpdate}; -use knot_types::{KnotHostname, OwnerDid, RepoDid, RepoRkey, UnixSeconds}; - -pub struct RepoContext { - pub repo: RepoDid, - pub owner: OwnerDid, - pub rkey: RepoRkey, -} - -struct Candidate { - repo_index: usize, - update: ReflogUpdate, -} - -pub fn pipeline_seeds( - layout: &Layout, - knot: &KnotHostname, - repos: &[RepoContext], - now_seconds: UnixSeconds, - capacity: usize, -) -> Vec { - let since = now_seconds.saturating_sub_secs(RECONSTRUCT_WINDOW_SECS); - let mut candidates: Vec = repos - .iter() - .enumerate() - .flat_map(|(repo_index, context)| match layout.open(&context.repo) { - Ok(opened) => opened - .reflog_updates_since(since) - .into_iter() - .map(|update| Candidate { repo_index, update }) - .collect::>(), - Err(_) => Vec::new(), - }) - .collect(); - candidates.sort_by(|left, right| { - let left_repo = repos[left.repo_index].repo.as_str(); - let right_repo = repos[right.repo_index].repo.as_str(); - left.update - .seconds - .cmp(&right.update.seconds) - .then_with(|| left_repo.cmp(right_repo)) - .then_with(|| left.update.name.as_str().cmp(right.update.name.as_str())) - .then_with(|| left.update.new.to_hex().cmp(&right.update.new.to_hex())) - }); - let start = candidates.len().saturating_sub(capacity.max(1)); - candidates - .split_off(start) - .into_iter() - .enumerate() - .filter_map(|(order, candidate)| { - let context = &repos[candidate.repo_index]; - let opened = layout.open(&context.repo).ok()?; - let pipeline = knot_postreceive::build_pipeline( - &opened, - &candidate.update.name, - candidate.update.transition(), - &context.repo, - &context.owner, - Some(&context.rkey), - knot, - )?; - Some(Seed::from_event( - candidate.update.seconds, - SeedOrder::new(order as u64), - &pipeline, - )) - }) - .collect() -} - -#[cfg(test)] -mod tests { - use knot_events::{EventCursor, EventLog}; - use knot_git::{ - EntryKind, Identity, Layout, NewCommit, RefUpdate, Repo, StagedAction, StagedChange, - }; - use knot_runtime::{Clock, SystemClock}; - use knot_types::{ - AuthorName, BranchName, Email, KnotHostname, Oid, OwnerDid, RefName, RepoDid, RepoRkey, - UnixSeconds, - }; - - use super::{RepoContext, pipeline_seeds}; - - const WORKFLOW: &str = "engine: nixery.dev/x\nwhen:\n - event: push\n branch: [main]\n"; - const EMPTY_TREE: &str = "4b825dc642cb6eb9a060e54bf8d69288fbee4904"; - - fn identity() -> Identity { - Identity { - name: AuthorName::new("nel"), - email: Email::new("nel@oyster.cafe"), - time: UnixSeconds::new(1_700_000_000), - offset_seconds: 0, - } - } - - fn push_commit( - bare: &Repo, - name: &RefName, - parent: Option, - extra: &[StagedChange], - ) -> Oid { - let workflow = StagedChange { - path: knot_types::RepoPath::new(".tangled/workflows/ci.yml").unwrap(), - action: StagedAction::Put { - content: WORKFLOW.as_bytes().to_vec(), - kind: EntryKind::Blob, - }, - }; - let changes: Vec = std::iter::once(workflow) - .chain(extra.iter().cloned()) - .collect(); - let tree = bare - .write_staged_tree(Oid::from_hex(EMPTY_TREE).unwrap(), &changes) - .unwrap(); - let tip = bare - .write_commit(&NewCommit { - tree, - parents: parent.into_iter().collect(), - author: identity(), - committer: identity(), - message: "add ci".to_string(), - extra_headers: Vec::new(), - }) - .unwrap(); - let update = match parent { - None => RefUpdate::Create { - name: name.clone(), - new: tip, - }, - Some(old) => RefUpdate::Update { - name: name.clone(), - old, - new: tip, - }, - }; - bare.update_ref(&update).unwrap(); - tip - } - - #[test] - fn a_missed_pipeline_is_reconstructed_from_the_reflog_after_restart() { - let scan = tempfile::tempdir().unwrap(); - let layout = Layout::new(scan.path()).with_default_branch(BranchName::new("main").unwrap()); - let did = RepoDid::new("did:plc:squid").unwrap(); - let bare = layout.create(&did).unwrap(); - let tip = push_commit(&bare, &RefName::new("refs/heads/main").unwrap(), None, &[]); - - let context = RepoContext { - repo: did.clone(), - owner: OwnerDid::new("did:web:olaren.dev").unwrap(), - rkey: RepoRkey::new("anemone").unwrap(), - }; - let now_seconds = - UnixSeconds::new((SystemClock.now_unix_micros().get() / 1_000_000) as i64); - let seeds = pipeline_seeds( - &layout, - &KnotHostname::new("knot.test").unwrap(), - &[context], - now_seconds, - 32, - ); - assert_eq!( - seeds.len(), - 1, - "push that compiled workflow yields exactly one reconstructed pipeline" - ); - - let log = EventLog::new(SystemClock, 32); - log.seed(seeds); - let replayed = log.replay(EventCursor::START, 32); - assert_eq!(replayed.len(), 1); - let event = &replayed[0]; - assert_eq!(event.nsid, "sh.tangled.pipeline"); - let payload = &event.payload; - assert_eq!(payload["$type"], "sh.tangled.pipeline"); - assert_eq!(payload["workflows"][0]["name"], "ci.yml"); - assert_eq!(payload["triggerMetadata"]["kind"], "push"); - assert_eq!(payload["triggerMetadata"]["repo"]["knot"], "knot.test"); - assert_eq!(payload["triggerMetadata"]["repo"]["repoDid"], did.as_str()); - assert_eq!(payload["triggerMetadata"]["push"]["ref"], "refs/heads/main"); - assert_eq!( - payload["triggerMetadata"]["push"]["newSha"], - tip.to_hex().to_string() - ); - } - - #[test] - fn pipeline_seeds_are_bounded_by_the_ring_capacity() { - let scan = tempfile::tempdir().unwrap(); - let layout = Layout::new(scan.path()).with_default_branch(BranchName::new("main").unwrap()); - let did = RepoDid::new("did:plc:limpet").unwrap(); - let bare = layout.create(&did).unwrap(); - let main = RefName::new("refs/heads/main").unwrap(); - - (0u8..4).fold(None, |parent, index| { - Some(push_commit( - &bare, - &main, - parent, - &[StagedChange { - path: knot_types::RepoPath::new(format!("file{index}.txt")).unwrap(), - action: StagedAction::Put { - content: vec![index], - kind: EntryKind::Blob, - }, - }], - )) - }); - - let context = RepoContext { - repo: did.clone(), - owner: OwnerDid::new("did:web:olaren.dev").unwrap(), - rkey: RepoRkey::new("anemone").unwrap(), - }; - let now_seconds = - UnixSeconds::new((SystemClock.now_unix_micros().get() / 1_000_000) as i64); - let seeds = pipeline_seeds( - &layout, - &KnotHostname::new("knot.test").unwrap(), - &[context], - now_seconds, - 2, - ); - assert_eq!( - seeds.len(), - 2, - "four reflogged workflow pushes are truncated to ring capacity" - ); - } -} diff --git a/crates/knot-xrpc/src/events.rs b/crates/knot-xrpc/src/events.rs index af1611e..ea8f6e0 100644 --- a/crates/knot-xrpc/src/events.rs +++ b/crates/knot-xrpc/src/events.rs @@ -10,7 +10,9 @@ use futures::{SinkExt, StreamExt}; use http::HeaderMap; use serde::Deserialize; -use knot_events::{EventCursor, EventLog}; +use knot_events::{ + BatchEnd, EventCursor, EventLog, ReplayBounds, ReplayBytes, ReplayEvents, Replayed, +}; use knot_runtime::{Clock, HttpTransport}; use crate::XrpcState; @@ -19,6 +21,7 @@ use crate::error::XrpcError; pub(crate) const EVENTS_ROUTE: &str = "/events"; const DRAIN_BATCH: usize = 100; +const DRAIN_BYTES: usize = 4 << 20; const MAX_BATCHES_PER_DRAIN: usize = 1_000; const KEEPALIVE: Duration = Duration::from_secs(30); const WRITE_DEADLINE: Duration = Duration::from_secs(10); @@ -106,6 +109,13 @@ async fn stream_events(socket: WebSocket, log: Arc>, start } } +fn drain_bounds() -> ReplayBounds { + ReplayBounds::new( + ReplayEvents::new(DRAIN_BATCH).expect("drain event maximum is nonzero"), + ReplayBytes::new(DRAIN_BYTES).expect("drain byte maximum is nonzero"), + ) +} + // who up draining they clock async fn drain( sink: &mut SplitSink, @@ -114,8 +124,7 @@ async fn drain( ) -> Result { let mut batches = 0; loop { - let events = log.replay(*cursor, DRAIN_BATCH); - let caught_up = events.len() < DRAIN_BATCH; + let Replayed { events, end } = log.replay(*cursor, drain_bounds()); if let Some(last) = events.last() { *cursor = last.created; } @@ -123,12 +132,13 @@ async fn drain( .iter() .map(|event| { Ok(Message::Text( - serde_json::to_string(event) + serde_json::to_string(event.as_ref()) .expect("wire event serializes to JSON") .into(), )) }) .collect(); + drop(events); let sent = tokio::time::timeout( WRITE_DEADLINE, sink.send_all(&mut futures::stream::iter(messages)), @@ -137,7 +147,7 @@ async fn drain( if !matches!(sent, Ok(Ok(()))) { return Err(()); } - if caught_up { + if end == BatchEnd::CaughtUp { return Ok(Drained::CaughtUp); } batches += 1; diff --git a/crates/knot-xrpc/src/tests.rs b/crates/knot-xrpc/src/tests.rs index 776637d..69d3b73 100644 --- a/crates/knot-xrpc/src/tests.rs +++ b/crates/knot-xrpc/src/tests.rs @@ -223,18 +223,25 @@ fn resolve(world: &World, rkey: &str) -> Resolved> { .resolve_repo(&member_owner(), &RepoRkey::new(rkey).unwrap()) } -fn replay(world: &World) -> Vec { +fn replay(world: &World) -> Vec> { world .state .events - .replay(knot_events::EventCursor::START, 64) + .replay( + knot_events::EventCursor::START, + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(64).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), + ) + .events } fn event_count(world: &World) -> usize { replay(world).len() } -fn last_event(world: &World, nsid: &str) -> knot_events::Event { +fn last_event(world: &World, nsid: &str) -> std::sync::Arc { replay(world) .into_iter() .rev() @@ -242,7 +249,7 @@ fn last_event(world: &World, nsid: &str) -> knot_events::Event { .unwrap_or_else(|| panic!("{nsid} event is emitted")) } -fn git_events(world: &World) -> Vec { +fn git_events(world: &World) -> Vec> { replay(world) .into_iter() .filter(|event| { @@ -254,7 +261,7 @@ fn git_events(world: &World) -> Vec { .collect() } -fn only_git_event(world: &World) -> knot_events::Event { +fn only_git_event(world: &World) -> std::sync::Arc { let mut events = git_events(world); assert_eq!(events.len(), 1, "expected exactly one non-acl event"); events.remove(0) @@ -307,6 +314,7 @@ fn state_from( atproto, secrets, entropy: Arc::new(OsEntropy), + ci_logs: None, admins: BTreeSet::from([account(ADMIN_HOST)]), admission, knot_did: knot, @@ -328,7 +336,10 @@ fn state_from( service_owner: account(ADMIN_HOST), events: Arc::new(knot_events::EventLog::new( ManualClock::new(UnixMicros::new(1_000_000_000)), - 1024, + knot_events::ReplayBounds::new( + knot_events::ReplayEvents::new(1024).unwrap(), + knot_events::ReplayBytes::new(16 << 20).unwrap(), + ), )), subscriber_gate: Arc::new(knot_events::SubscriberGate::new( knot_events::GlobalSubscriberLimit::new(16), @@ -1846,7 +1857,11 @@ async fn delete_branch_removes_a_branch_and_refuses_the_default() { oid.to_string(), "the deletion event includes the branch's old tip" ); - assert_eq!(event.payload["newSha"], "", "deletion has no new sha"); + assert_eq!( + event.payload["newSha"], + git.object_format().null_oid().to_string(), + "deletion reports the null oid as the new sha" + ); assert_eq!( event.payload["committerDid"], account(MEMBER_HOST).to_string() @@ -2367,7 +2382,7 @@ mod merge_endpoints { } #[tokio::test] - async fn a_native_merge_triggers_the_pipeline() { + async fn a_native_merge_advances_the_branch_without_a_pipeline_event() { let world = World::new(); add_member_helper(&world).await; let repo_did = create_repo_helper(&world, "kelp").await; @@ -2403,13 +2418,14 @@ mod merge_endpoints { StatusCode::OK ); - let pipeline = last_event(&world, "sh.tangled.pipeline"); - assert_eq!(pipeline.payload["triggerMetadata"]["kind"], "push"); - assert_eq!( - pipeline.payload["triggerMetadata"]["push"]["ref"], - "refs/heads/main" + let update = last_event(&world, "sh.tangled.git.refUpdate"); + assert_eq!(update.payload["ref"], "refs/heads/main"); + assert!( + replay(&world) + .iter() + .all(|event| event.nsid != "sh.tangled.pipeline"), + "the knot emits no pipeline record" ); - assert_eq!(pipeline.payload["workflows"][0]["name"], "ci.yml"); } #[tokio::test] diff --git a/crates/knot-xrpc/tests/reads.rs b/crates/knot-xrpc/tests/reads.rs index cb9fc6b..594d09f 100644 --- a/crates/knot-xrpc/tests/reads.rs +++ b/crates/knot-xrpc/tests/reads.rs @@ -1143,9 +1143,13 @@ async fn branch_tips_render_edge_shapes() { &["merge", "-q", "--no-ff", "-m", "merge side", "side"], ); sh_git(work.path(), &["push", "-q", &bare_str, "main"]); + let first_parent = sh_git(work.path(), &["rev-parse", "HEAD^1"]); + let second_parent = sh_git(work.path(), &["rev-parse", "HEAD^2"]); let tag_object = sh_git(work.path(), &["rev-parse", "v1.0.0"]); std::fs::write(bare.join("refs/heads/tagtip"), format!("{tag_object}\n")).unwrap(); + let root_commit = sh_git(work.path(), &["rev-list", "--max-parents=0", "HEAD"]); + std::fs::write(bare.join("refs/heads/roottip"), format!("{root_commit}\n")).unwrap(); let value = get_json( &world, @@ -1153,28 +1157,47 @@ async fn branch_tips_render_edge_shapes() { ) .await; let branches = value["branches"].as_array().unwrap(); - assert_eq!(branches.len(), 2); + assert_eq!(branches.len(), 3); - let main = branches - .iter() - .find(|branch| branch["reference"]["name"] == "main") - .unwrap(); - let parents = main["commit"]["ParentHashes"].as_array().unwrap(); - assert_eq!(parents.len(), 1); - assert!( - parents[0] + let branch = |name: &str| { + branches + .iter() + .find(|branch| branch["reference"]["name"] == name) + .unwrap() + }; + let parents = |branch: &serde_json::Value| -> Vec { + branch["commit"]["ParentHashes"] .as_array() .unwrap() .iter() - .all(|byte| byte.as_u64() == Some(0)), - "merge tip must have the zero parent hash" + .map(|parent| { + parent + .as_array() + .unwrap() + .iter() + .map(|byte| format!("{:02x}", byte.as_u64().unwrap())) + .collect() + }) + .collect() + }; + + let main = branch("main"); + assert_eq!( + parents(main), + vec![first_parent, second_parent], + "merge tip must report both parents in order" ); assert_eq!(main["commit"]["Author"]["Name"], "nel"); + assert!( + parents(branch("roottip")).is_empty(), + "root tip must report no parents" + ); - let tagtip = branches - .iter() - .find(|branch| branch["reference"]["name"] == "tagtip") - .unwrap(); + let tagtip = branch("tagtip"); + assert!( + parents(tagtip).is_empty(), + "a non-commit tip must report no parents" + ); assert_eq!(tagtip["reference"]["hash"], tag_object.as_str()); assert_eq!(tagtip["commit"]["Author"]["Name"], ""); assert_eq!(tagtip["commit"]["Author"]["When"], "0001-01-01T00:00:00Z"); diff --git a/example.toml b/example.toml index 213c429..1a57ef4 100644 --- a/example.toml +++ b/example.toml @@ -267,6 +267,10 @@ # Default value: 4096 #events_replay_buffer = 4096 +# Can also be specified via environment variable `KNOT_XRPC_EVENTS_REPLAY_BYTES`. +# Default value: 67108864 +#events_replay_bytes = 67108864 + # Can also be specified via environment variable `KNOT_XRPC_EVENTS_MAX_SUBSCRIBERS`. # Default value: 256 #events_max_subscribers = 256 @@ -395,6 +399,10 @@ # Can also be specified via environment variable `KNOT_HOMEPAGE_PATH`. #path = +[ci] +# Can also be specified via environment variable `KNOT_CI_LOGS_ADDR`. +#logs_addr = + [messages] [messages.push] # Default value: ["{knot} received {refs}."] @@ -406,6 +414,12 @@ # Default value: ["pipeline compiled with no diagnostics"] #pipeline_clean = ["pipeline compiled with no diagnostics"] +# Default value: ["no pipelines to compile"] +#pipeline_none = ["no pipelines to compile"] + +# Default value: ["-> Browse CI logs in your terminal:", " ssh -t -p {port} {host} {repo} {sha}"] +#ci_logs = ["-> Browse CI logs in your terminal:", " ssh -t -p {port} {host} {repo} {sha}"] + [messages.fetch] # Default value: ["Thanks for using {knot}!"] #motd = ["Thanks for using {knot}!"] diff --git a/lexicons/git/refUpdate.json b/lexicons/git/refUpdate.json index 6b8c04d..c37f611 100644 --- a/lexicons/git/refUpdate.json +++ b/lexicons/git/refUpdate.json @@ -4,7 +4,7 @@ "defs": { "main": { "type": "record", - "description": "An update to a git repository, emitted by knots.", + "description": "An event record representing git-push operation to git repository, emitted by knots.", "key": "tid", "record": { "type": "object", @@ -50,6 +50,26 @@ "minLength": 40, "maxLength": 40 }, + "changedFiles": { + "type": "array", + "description": "files changed between commits", + "items": { "type": "string" } + }, + "pushOptions": { + "type": "array", + "description": "push options passed on git-push", + "maxLength": 50, + "items": { + "type": "string", + "maxLength": 1024, + "knownValues": [ + "ci-skip", + "ci-verbose", + "skip-ci", + "verbose-ci" + ] + } + }, "meta": { "type": "ref", "ref": "#meta" diff --git a/lexicons/issue/state.json b/lexicons/issue/state.json index ba9d0ec..ebd9107 100644 --- a/lexicons/issue/state.json +++ b/lexicons/issue/state.json @@ -11,13 +11,18 @@ "type": "object", "required": [ "issue", - "state" + "state", + "createdAt" ], "properties": { "issue": { "type": "string", "format": "at-uri" }, + "createdAt": { + "type": "string", + "format": "datetime" + }, "state": { "type": "string", "description": "state of the issue", diff --git a/lexicons/pipeline/cancelPipeline.json b/lexicons/pipeline/cancelPipeline.json index f1ecf48..12b9b95 100644 --- a/lexicons/pipeline/cancelPipeline.json +++ b/lexicons/pipeline/cancelPipeline.json @@ -4,7 +4,7 @@ "defs": { "main": { "type": "procedure", - "description": "Cancel a running pipeline", + "description": "DEPRECATED: use sh.tangled.ci.cancelPipeline instead - Cancel a running pipeline", "input": { "encoding": "application/json", "schema": { diff --git a/lexicons/pipeline/pipeline.json b/lexicons/pipeline/pipeline.json index 3e0f932..9deb76f 100644 --- a/lexicons/pipeline/pipeline.json +++ b/lexicons/pipeline/pipeline.json @@ -9,6 +9,7 @@ "key": "tid", "record": { "type": "object", + "description": "DEPRECATED: use sh.tangled.ci.pipeline instead", "required": [ "triggerMetadata", "workflows" @@ -58,6 +59,11 @@ "manual": { "type": "ref", "ref": "#manualTriggerData" + }, + "sourceRepo": { + "type": "string", + "format": "did", + "description": "Repository DID that code and workflow definitions are checked out from, when different from repo (e.g. a fork's commit for a fork-based manual trigger). If absent, source uses repo itself." } } }, @@ -117,8 +123,7 @@ "required": [ "sourceBranch", "targetBranch", - "sourceSha", - "action" + "sourceSha" ], "properties": { "sourceBranch": { @@ -132,14 +137,29 @@ "minLength": 40, "maxLength": 40 }, - "action": { - "type": "string" + "pull": { + "type": "string", + "format": "at-uri", + "description": "AT-URI of the sh.tangled.repo.pull record this run belongs to" } } }, "manualTriggerData": { "type": "object", + "required": [ + "sha" + ], "properties": { + "sha": { + "type": "string", + "description": "commit SHA the manual run targets", + "minLength": 40, + "maxLength": 40 + }, + "ref": { + "type": "string", + "description": "optional ref the SHA was resolved from, for display and TANGLED_REF" + }, "inputs": { "type": "array", "items": { @@ -178,7 +198,8 @@ "required": [ "skip", "depth", - "submodules" + "submodules", + "tags" ], "properties": { "skip": { @@ -189,6 +210,9 @@ }, "submodules": { "type": "boolean" + }, + "tags": { + "type": "boolean" } } }, diff --git a/lexicons/pipeline/status.json b/lexicons/pipeline/status.json index 4f93318..a9689f7 100644 --- a/lexicons/pipeline/status.json +++ b/lexicons/pipeline/status.json @@ -9,6 +9,7 @@ "key": "tid", "record": { "type": "object", + "description": "DEPRECATED: use sh.tangled.ci.pipeline instead", "required": ["pipeline", "workflow", "status", "createdAt"], "properties": { "pipeline": { diff --git a/lexicons/pulls/state.json b/lexicons/pulls/state.json index d33422f..f8b7e80 100644 --- a/lexicons/pulls/state.json +++ b/lexicons/pulls/state.json @@ -11,13 +11,18 @@ "type": "object", "required": [ "pull", - "status" + "status", + "createdAt" ], "properties": { "pull": { "type": "string", "format": "at-uri" }, + "createdAt": { + "type": "string", + "format": "datetime" + }, "status": { "type": "string", "description": "status of the pull request", diff --git a/lexicons/repo/forkSync.json b/lexicons/repo/forkSync.json index 3af3fca..f3c22ca 100644 --- a/lexicons/repo/forkSync.json +++ b/lexicons/repo/forkSync.json @@ -30,6 +30,11 @@ "type": "string", "description": "Name of the forked repository" }, + "repo": { + "type": "string", + "format": "did", + "description": "DID of the repository" + }, "branch": { "type": "string", "description": "Branch to sync" diff --git a/lexicons/repo/merge.json b/lexicons/repo/merge.json index 7cecb55..ffd2104 100644 --- a/lexicons/repo/merge.json +++ b/lexicons/repo/merge.json @@ -20,6 +20,11 @@ "type": "string", "description": "Name of the repository" }, + "repo": { + "type": "string", + "format": "did", + "description": "DID of the repository" + }, "patch": { "type": "string", "description": "Patch content to merge" diff --git a/lexicons/repo/mergeCheck.json b/lexicons/repo/mergeCheck.json index 054b7f1..7a0ee6e 100644 --- a/lexicons/repo/mergeCheck.json +++ b/lexicons/repo/mergeCheck.json @@ -20,6 +20,11 @@ "type": "string", "description": "Name of the repository" }, + "repo": { + "type": "string", + "format": "did", + "description": "DID of the repository" + }, "patch": { "type": "string", "description": "Patch or pull request to check for merge conflicts"