diff --git a/knot2/crates/knot-xrpc/src/firehose.rs b/knot2/crates/knot-xrpc/src/firehose.rs index b4b07cd23..b2ec57f2d 100644 --- a/knot2/crates/knot-xrpc/src/firehose.rs +++ b/knot2/crates/knot-xrpc/src/firehose.rs @@ -18,7 +18,9 @@ use knot_events::{ use knot_git::{Layout, RefUpdate, Repo}; use knot_record::chain::{self, CommitChunk, Emit, MAX_COMMIT_BLOCKS_BYTES, Op, blocks_bytes}; use knot_runtime::{Clock, HttpTransport, UnixMicros}; -use knot_types::{Datetime, Oid, RefName, RepoDid, UnixSeconds}; +use knot_types::{ + AccountDid, Datetime, Oid, PushOption, PushOptions, RefName, RepoDid, UnixSeconds, +}; use serde::{Deserialize, Serialize}; use crate::XrpcState; @@ -33,6 +35,8 @@ const DRAIN_BYTES: usize = 4 << 20; const KEEPALIVE: Duration = Duration::from_secs(30); const WRITE_DEADLINE: Duration = Duration::from_secs(10); const TRY_AGAIN_LATER: u16 = 1013; +const MAX_PUSH_EDITS: usize = 200; +const MAX_PUSH_OPTIONS: usize = 10; pub struct FirehoseEmit<'a, C: Clock> { log: &'a EventLog, @@ -202,6 +206,44 @@ pub fn publish_sync( ); } +pub fn publish_push( + log: &EventLog, + did: &RepoDid, + pusher: &AccountDid, + applied: &[RefUpdate], + options: &PushOptions, + now: UnixSeconds, +) { + let options = options.as_slice(); + if applied.len() > MAX_PUSH_EDITS || options.len() > MAX_PUSH_OPTIONS { + tracing::warn!( + repo = did.as_str(), + edits = applied.len(), + options = options.len(), + "a push overran what its event can carry, so the frame names only the refs and options that fit" + ); + } + let reservation = log.reserve(); + let seq = reservation.micros(); + fulfill( + reservation, + FrameKind::Push, + &PushEventMessage { + seq, + repo: did, + pusher, + edits: applied + .iter() + .take(MAX_PUSH_EDITS) + .map(RefEdit::of) + .collect(), + time: datetime_of(now), + options: (!options.is_empty()).then(|| options.iter().take(MAX_PUSH_OPTIONS).collect()), + sig: b"", + }, + ); +} + fn fulfill(reservation: knot_events::Reservation, t: FrameKind, message: &M) { reservation.fulfill_frame(frame(t, message)); } @@ -218,6 +260,8 @@ enum FrameKind { Account, #[serde(rename = "#info")] Info, + #[serde(rename = "org.tangled.git.pushEvent")] + Push, } #[derive(Debug, Clone, Copy)] @@ -334,6 +378,49 @@ struct SyncMessage { time: Datetime, } +#[derive(Serialize)] +struct RefEdit<'a> { + name: &'a RefName, + new: Option<&'a Oid>, + old: Option<&'a Oid>, +} + +impl<'a> RefEdit<'a> { + fn of(update: &'a RefUpdate) -> Self { + match update { + RefUpdate::Create { name, new } => Self { + name, + new: Some(new), + old: None, + }, + RefUpdate::Update { name, old, new } => Self { + name, + new: Some(new), + old: Some(old), + }, + RefUpdate::Delete { name, old } => Self { + name, + new: None, + old: Some(old), + }, + } + } +} + +#[derive(Serialize)] +struct PushEventMessage<'a> { + seq: UnixMicros, + repo: &'a RepoDid, + pusher: &'a AccountDid, + edits: Vec>, + time: Datetime, + #[serde(skip_serializing_if = "Option::is_none")] + options: Option>, + // the signing scheme this field names isn't defined yet, so it goes out empty + #[serde(with = "serde_bytes")] + sig: &'static [u8], +} + fn encode(value: &T) -> Bytes { Bytes::from(serde_ipld_dagcbor::to_vec(value).expect("frame values encode as dag-cbor")) } diff --git a/knot2/crates/knot-xrpc/src/refrecords.rs b/knot2/crates/knot-xrpc/src/refrecords.rs index 76a93651e..2d39572d2 100644 --- a/knot2/crates/knot-xrpc/src/refrecords.rs +++ b/knot2/crates/knot-xrpc/src/refrecords.rs @@ -214,5 +214,13 @@ impl knot_receive::RefSurface for RefProjection