diff --git a/Cargo.lock b/Cargo.lock index 733175a5e..8d44f81cb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -801,6 +801,7 @@ dependencies = [ "itertools 0.14.0", "jacquard-common", "lasso", + "parking_lot", "scc", "serde", "smallvec", diff --git a/api/org_tangled/notificationlistNotifications.go b/api/org_tangled/notificationlistNotifications.go index f15fd5682..4649fd404 100644 --- a/api/org_tangled/notificationlistNotifications.go +++ b/api/org_tangled/notificationlistNotifications.go @@ -23,12 +23,14 @@ type TempNotificationListNotifications_Notification struct { CreatedAt string `json:"createdAt" cborgen:"createdAt"` // issueAt: AT-URI of the related org.tangled.issue.issue record, if applicable. IssueAt *string `json:"issueAt,omitempty" cborgen:"issueAt,omitempty"` + // knotDid: did:web of knot that offered membership. Only knot_invited has it for now! Though in future more notifs may use it. + KnotDid *string `json:"knotDid,omitempty" cborgen:"knotDid,omitempty"` // pullAt: AT-URI of the related org.tangled.pulls.pull record, if applicable. PullAt *string `json:"pullAt,omitempty" cborgen:"pullAt,omitempty"` Read bool `json:"read" cborgen:"read"` // repoDid: DID of the related repository, if applicable. RepoDid *string `json:"repoDid,omitempty" cborgen:"repoDid,omitempty"` - // type: Notification type: repo_starred, issue_created, issue_commented, issue_closed, issue_reopen, issue_assigned, issue_unassigned, pull_created, pull_commented, pull_merged, pull_closed, pull_reopen, pull_assigned, pull_unassigned, followed, user_mentioned. + // type: Notification type: repo_starred, issue_created, issue_commented, issue_closed, issue_reopen, issue_assigned, issue_unassigned, pull_created, pull_commented, pull_merged, pull_closed, pull_reopen, pull_assigned, pull_unassigned, followed, user_mentioned, knot_invited, collaborator_invited. Type string `json:"type" cborgen:"type"` // uri: notification key; pass it back to updateSeen and deleteNotification. Uri string `json:"uri" cborgen:"uri"` diff --git a/api/tangled/knotlistMemberInvitesBy.go b/api/tangled/knotlistMemberInvitesBy.go new file mode 100644 index 000000000..ce862a66c --- /dev/null +++ b/api/tangled/knotlistMemberInvitesBy.go @@ -0,0 +1,53 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.knot.listMemberInvitesBy + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + KnotListMemberInvitesByNSID = "sh.tangled.knot.listMemberInvitesBy" +) + +// KnotListMemberInvitesBy_InviteItem is a "inviteItem" in the sh.tangled.knot.listMemberInvitesBy schema. +// +// Offer awaiting the subject's acceptance. +type KnotListMemberInvitesBy_InviteItem struct { + // addedBy: DID that made the offer. + AddedBy string `json:"addedBy" cborgen:"addedBy"` + // createdAt: When the knot made the offer. This listing sorts on it by newest first. + CreatedAt string `json:"createdAt" cborgen:"createdAt"` + // knot: did:web of knot that made this offer. Check against pending. + Knot string `json:"knot" cborgen:"knot"` + // uri: Invite record on knot's repo. + Uri string `json:"uri" cborgen:"uri"` +} + +// KnotListMemberInvitesBy_Output is the output of a sh.tangled.knot.listMemberInvitesBy call. +type KnotListMemberInvitesBy_Output struct { + Items []*KnotListMemberInvitesBy_InviteItem `json:"items" cborgen:"items"` + // pending: The knots that this indexation hasn't yet read. + Pending []string `json:"pending" cborgen:"pending"` + // truncated: Per-subject cut dropped offers from items. If truncated is true, missing offers won't advertise that they were withdrawn, same as with pending knots. + Truncated bool `json:"truncated" cborgen:"truncated"` +} + +// KnotListMemberInvitesBy calls the XRPC method "sh.tangled.knot.listMemberInvitesBy". +// +// subject: Actor DID offered membership. +func KnotListMemberInvitesBy(ctx context.Context, c util.LexClient, subject string) (*KnotListMemberInvitesBy_Output, error) { + var out KnotListMemberInvitesBy_Output + + params := map[string]interface{}{} + params["subject"] = subject + if err := c.LexDo(ctx, util.Query, "", "sh.tangled.knot.listMemberInvitesBy", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/api/tangled/repolistCollaboratorInvitesBy.go b/api/tangled/repolistCollaboratorInvitesBy.go new file mode 100644 index 000000000..6e3963352 --- /dev/null +++ b/api/tangled/repolistCollaboratorInvitesBy.go @@ -0,0 +1,55 @@ +// Code generated by cmd/lexgen (see Makefile's lexgen); DO NOT EDIT. + +package tangled + +// schema: sh.tangled.repo.listCollaboratorInvitesBy + +import ( + "context" + + "github.com/bluesky-social/indigo/lex/util" +) + +const ( + RepoListCollaboratorInvitesByNSID = "sh.tangled.repo.listCollaboratorInvitesBy" +) + +// RepoListCollaboratorInvitesBy_InviteItem is a "inviteItem" in the sh.tangled.repo.listCollaboratorInvitesBy schema. +// +// An offer awaiting the subject's acceptance. +type RepoListCollaboratorInvitesBy_InviteItem struct { + // addedBy: DID that made the offer. + AddedBy string `json:"addedBy" cborgen:"addedBy"` + // createdAt: When the repository made the offer. This listing sorts it by newest first. + CreatedAt string `json:"createdAt" cborgen:"createdAt"` + // knot: did:web of knot hosting this repo. + Knot string `json:"knot" cborgen:"knot"` + // repo: DID of repo that made the offer. + Repo string `json:"repo" cborgen:"repo"` + // uri: Invite record on the repository. + Uri string `json:"uri" cborgen:"uri"` +} + +// RepoListCollaboratorInvitesBy_Output is the output of a sh.tangled.repo.listCollaboratorInvitesBy call. +type RepoListCollaboratorInvitesBy_Output struct { + Items []*RepoListCollaboratorInvitesBy_InviteItem `json:"items" cborgen:"items"` + // pending: Knots that this index hasn't read yet + Pending []string `json:"pending" cborgen:"pending"` + // truncated: Per-subject dropped offers from items. If truncated is true, missing offers won't say if they were withdrawn, same with pending knots. + Truncated bool `json:"truncated" cborgen:"truncated"` +} + +// RepoListCollaboratorInvitesBy calls the XRPC method "sh.tangled.repo.listCollaboratorInvitesBy". +// +// subject: Actor DID offered collaboration. +func RepoListCollaboratorInvitesBy(ctx context.Context, c util.LexClient, subject string) (*RepoListCollaboratorInvitesBy_Output, error) { + var out RepoListCollaboratorInvitesBy_Output + + params := map[string]interface{}{} + params["subject"] = subject + if err := c.LexDo(ctx, util.Query, "", "sh.tangled.repo.listCollaboratorInvitesBy", params, nil, &out); err != nil { + return nil, err + } + + return &out, nil +} diff --git a/bobbin/crates/edge-index/Cargo.toml b/bobbin/crates/edge-index/Cargo.toml index 00777dfd5..212e6f458 100644 --- a/bobbin/crates/edge-index/Cargo.toml +++ b/bobbin/crates/edge-index/Cargo.toml @@ -13,6 +13,7 @@ ftree = "1.3.0" itertools = { workspace = true } jacquard-common = { workspace = true } lasso = { workspace = true } +parking_lot = { workspace = true } scc = { workspace = true } smallvec = "1" serde = { workspace = true } diff --git a/bobbin/crates/edge-index/src/invites.rs b/bobbin/crates/edge-index/src/invites.rs new file mode 100644 index 000000000..8a5ac6c7e --- /dev/null +++ b/bobbin/crates/edge-index/src/invites.rs @@ -0,0 +1,226 @@ +use std::collections::{HashMap, HashSet}; + +use bobbin_runtime::UnixMicros; +use itertools::Itertools; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use parking_lot::RwLock; + +#[derive(Clone, Debug, Eq, Hash, PartialEq)] +pub enum InviteScope { + Knot(Did), + Repo { + knot: Did, + repo: Did, + }, +} + +impl InviteScope { + pub fn knot(&self) -> &Did { + match self { + Self::Knot(knot) | Self::Repo { knot, .. } => knot, + } + } + + pub fn subject(&self) -> &Did { + match self { + Self::Knot(knot) => knot, + Self::Repo { repo, .. } => repo, + } + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct Offer { + pub invited_by: Did, + pub offered_at: UnixMicros, +} + +#[derive(Default)] +struct Held { + offers: HashMap, HashMap>, + discovered: bool, + awaiting: HashSet>, +} + +#[derive(Default)] +pub struct InviteIndex(RwLock); + +impl InviteIndex { + pub fn new() -> Self { + Self::default() + } + + pub fn awaiting(&self, knot: Did) { + self.0.write().awaiting.insert(knot); + } + + pub fn backfilled(&self, knot: &Did) { + self.0.write().awaiting.remove(knot); + } + + pub fn all_knots_discovered(&self) { + self.0.write().discovered = true; + } + + pub fn discovered(&self) -> bool { + self.0.read().discovered + } + + pub fn pending(&self) -> Vec> { + self.0 + .read() + .awaiting + .iter() + .sorted_by(|left, right| left.as_ref().cmp(right.as_ref())) + .cloned() + .collect() + } + + pub fn offer(&self, invitee: Did, scope: InviteScope, offer: Offer) { + self.0 + .write() + .offers + .entry(invitee) + .or_default() + .insert(scope, offer); + } + + pub fn withdraw(&self, invitee: &Did, scope: &InviteScope) { + let mut guard = self.0.write(); + let emptied = guard.offers.get_mut(invitee).is_some_and(|scoped| { + scoped.remove(scope); + scoped.is_empty() + }); + if emptied { + guard.offers.remove(invitee); + } + } + + pub fn outstanding(&self, invitee: &Did) -> Vec<(InviteScope, Offer)> { + self.0 + .read() + .offers + .get(invitee) + .into_iter() + .flatten() + .map(|(scope, offer)| (scope.clone(), offer.clone())) + .sorted_by(|(left_scope, left), (right_scope, right)| { + right.offered_at.cmp(&left.offered_at).then_with(|| { + left_scope + .subject() + .as_ref() + .cmp(right_scope.subject().as_ref()) + }) + }) + .collect() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const INVITEE: &str = "did:plc:boltless"; + + fn did(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + + fn knot_scope() -> InviteScope { + InviteScope::Knot(did("did:web:knot.oyster.cafe")) + } + + fn repo_scope(repo: &str) -> InviteScope { + InviteScope::Repo { + knot: did("did:web:knot.oyster.cafe"), + repo: did(repo), + } + } + + fn offer(index: &InviteIndex, scope: InviteScope, invited_by: &str, micros: u64) { + index.offer( + did(INVITEE), + scope, + Offer { + invited_by: did(invited_by), + offered_at: UnixMicros::new(micros), + }, + ); + } + + #[test] + fn an_invitee_reads_their_outstanding_offers_newest_first() { + let index = InviteIndex::new(); + offer(&index, knot_scope(), "did:plc:akshay", 10); + offer(&index, repo_scope("did:plc:scallop"), "did:plc:akshay", 30); + offer(&index, repo_scope("did:plc:limpet"), "did:plc:olaren", 20); + + let outstanding = index.outstanding(&did(INVITEE)); + assert_eq!( + outstanding + .iter() + .map(|(_, offer)| offer.offered_at.raw()) + .collect::>(), + vec![30, 20, 10] + ); + assert_eq!(outstanding[0].0, repo_scope("did:plc:scallop")); + assert_eq!(outstanding[0].1.invited_by, did("did:plc:akshay")); + assert_eq!( + outstanding[0].0.knot(), + &did("did:web:knot.oyster.cafe"), + "repo-scoped offer still reports knot, since pending is keyed on knot" + ); + assert_eq!(outstanding[0].0.subject(), &did("did:plc:scallop")); + } + + #[test] + fn re_offer_replaces_original_and_withdrawal_clears_it() { + let index = InviteIndex::new(); + offer(&index, knot_scope(), "did:plc:akshay", 10); + offer(&index, knot_scope(), "did:plc:olaren", 40); + + let outstanding = index.outstanding(&did(INVITEE)); + assert_eq!(outstanding.len(), 1); + assert_eq!(outstanding[0].1.invited_by, did("did:plc:olaren")); + assert_eq!(outstanding[0].1.offered_at, UnixMicros::new(40)); + + index.withdraw(&did(INVITEE), &knot_scope()); + assert!(index.outstanding(&did(INVITEE)).is_empty()); + } + + #[test] + fn knot_stays_pending_from_discovery_until_backfill_finishes() { + let index = InviteIndex::new(); + assert!( + !index.discovered(), + "index still to read its first knot can't say which knots it is missing" + ); + + index.awaiting(did("did:web:knot.oyster.cafe")); + index.awaiting(did("did:web:knot.nel.pet")); + index.all_knots_discovered(); + assert!(index.discovered()); + assert_eq!( + index.pending(), + vec![did("did:web:knot.nel.pet"), did("did:web:knot.oyster.cafe")] + ); + + index.backfilled(&did("did:web:knot.nel.pet")); + assert_eq!( + index.pending(), + vec![did("did:web:knot.oyster.cafe")], + "one knot mid-backfill affects its own offers only" + ); + + index.backfilled(&did("did:web:knot.oyster.cafe")); + assert!(index.pending().is_empty()); + + index.awaiting(did("did:web:knot.late.pet")); + assert_eq!( + index.pending(), + vec![did("did:web:knot.late.pet")], + "knot registered after replay is pending on its own" + ); + } +} diff --git a/bobbin/crates/edge-index/src/lib.rs b/bobbin/crates/edge-index/src/lib.rs index bb16d98cf..fb1c45a0f 100644 --- a/bobbin/crates/edge-index/src/lib.rs +++ b/bobbin/crates/edge-index/src/lib.rs @@ -50,11 +50,13 @@ impl BucketKey { } pub mod coverage; +pub mod invites; pub mod settle; mod sorted_blocks; pub mod state_index; mod state_view; pub use coverage::{Coverage, CoverageWatch, HydrantCursor, PromotionSignal}; +pub use invites::{InviteIndex, InviteScope, Offer}; pub use settle::{ParsedCid, RecordOutcome, Rejection, Settlement, Settlements, Watch}; use sorted_blocks::{PageKey, SortedBlocks}; pub use state_index::{ diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index 77e038133..50b585c9f 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -24,16 +24,17 @@ use axum::{ }; pub use bobbin_codesearch::{CodeSearch, CodeSearchError}; use bobbin_edge_index::{ - Coverage, CoverageWatch, CursorParseError, EdgeItem, EdgePage, EdgeStore, IssueStateKind, - PageCursor, PageLimit, PageOffset, PageStart, PageToken, ParsedCid, PullStatusKind, Rejection, - Settlement, Settlements, SortDir, StateIndex, StateKind, TotalCount, Watch, + Coverage, CoverageWatch, CursorParseError, EdgeItem, EdgePage, EdgeStore, InviteIndex, + InviteScope, IssueStateKind, Offer, PageCursor, PageLimit, PageOffset, PageStart, PageToken, + ParsedCid, PullStatusKind, Rejection, Settlement, Settlements, SortDir, StateIndex, StateKind, + TotalCount, Watch, }; use bobbin_knot_proxy::{ KnotHost, KnotProxy, KnotProxyError, MirrorNsid, MirrorProxy, ProxyResponse, RepoSlug, }; use bobbin_record_lru::RecordStore; use bobbin_resolver::{IdentityResolveError, IdentityResolver, RepoIdResolver}; -use bobbin_runtime::ReqwestHttp; +use bobbin_runtime::{ReqwestHttp, UnixMicros}; use bobbin_search::{ ActorIndex, ActorIndexError, SearchCursor, SearchError, SearchFilters, SearchHit, SearchOffset, SearchReader, @@ -89,6 +90,7 @@ use bobbin_types::sh_tangled::string::{ use futures::Stream; use futures::stream::{self, StreamExt, TryStreamExt}; use jacquard_axum::service_auth::{self, ExtractOptionalServiceAuth, ServiceAuthConfig}; +use jacquard_common::types::datetime::Datetime; use jacquard_common::types::did::Did; use jacquard_common::types::ident::AtIdentifier; use jacquard_common::types::nsid::Nsid; @@ -136,6 +138,7 @@ const ACTOR_TYPEAHEAD_MAX_LIMIT: u32 = 100; const ACTOR_TYPEAHEAD_MAX_QUERY_BYTES: usize = 253; const FETCH_CONCURRENCY: usize = 8; const TANGLED_NSID_PREFIX: &str = "sh.tangled."; +const MAX_OUTSTANDING_OFFERS: usize = 128; pub type Directory = JacquardResolver; @@ -165,6 +168,7 @@ pub struct AppState { pub awaiting: Arc, pub client_address: Arc, pub settlements: Arc, + pub invites: Arc, service_auth_config: service_auth::ServiceAuthConfig>, enrich_router: Arc>, trending_cache: Arc>>, @@ -209,6 +213,7 @@ impl AppState { awaiting: Arc::new(AwaitLimiter::new(MaxAwaiting::default())), client_address: Arc::new(ClientAddress::default()), settlements: Arc::new(Settlements::new()), + invites: Arc::new(InviteIndex::new()), service_auth_config: ServiceAuthConfig::new( Did::new_static("did:web:localhost").unwrap(), directory, @@ -262,6 +267,11 @@ impl AppState { self } + pub fn with_invites(mut self, invites: Arc) -> Self { + self.invites = invites; + self + } + pub fn with_service_did(mut self, did: Did) -> Self { self.service_auth_config = ServiceAuthConfig::new(did, self.directory.clone()); self @@ -412,6 +422,14 @@ pub fn router(state: AppState) -> Router { "/xrpc/sh.tangled.knot.listMembersBy", get(list_knot_members_by), ) + .route( + "/xrpc/sh.tangled.knot.listMemberInvitesBy", + get(list_knot_member_invites_by), + ) + .route( + "/xrpc/sh.tangled.repo.listCollaboratorInvitesBy", + get(list_collaborator_invites_by), + ) .route( "/xrpc/sh.tangled.knot.countMembersBy", get(count_knot_members_by), @@ -907,6 +925,8 @@ pub enum XrpcError { Internal(String), #[error("overloaded, shedding under memory pressure")] Overloaded, + #[error("invite index is still warming, so it can't list offers yet")] + Warming, #[error("record has not been ingested yet")] NotSettled, #[error("not implemented: {0}")] @@ -936,6 +956,7 @@ impl IntoResponse for XrpcError { Self::InvalidRecord(_) => (StatusCode::BAD_GATEWAY, "InvalidRecord"), Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "InternalError"), Self::Overloaded => (StatusCode::SERVICE_UNAVAILABLE, "Overloaded"), + Self::Warming => (StatusCode::SERVICE_UNAVAILABLE, "Warming"), Self::NotSettled => (StatusCode::GATEWAY_TIMEOUT, "NotSettled"), Self::NotImplemented(_) => (StatusCode::NOT_IMPLEMENTED, "MethodNotImplemented"), }; @@ -1727,7 +1748,7 @@ where } fn synth_knot_owned_value(source: KnotOwnedSource, sort_micros: u64) -> Option { - let created_at = micros_to_rfc3339(sort_micros)?; + let created_at = datetime_of(UnixMicros::new(sort_micros))?; match source { KnotOwnedSource::Member { knot, subject } => Some(serde_json::json!({ "domain": knot_did_host(&knot)?, @@ -1742,10 +1763,9 @@ fn synth_knot_owned_value(source: KnotOwnedSource, sort_micros: u64) -> Option Option { - let micros = i64::try_from(micros).ok()?; - chrono::DateTime::from_timestamp_micros(micros) - .map(|dt| dt.to_rfc3339_opts(chrono::SecondsFormat::Micros, true)) +fn datetime_of(at: UnixMicros) -> Option { + let micros = i64::try_from(at.raw()).ok()?; + chrono::DateTime::from_timestamp_micros(micros).map(|dt| Datetime::new(dt.fixed_offset())) } #[derive(Clone, Copy)] @@ -1814,6 +1834,7 @@ fn drop_unhydratable( Err( err @ (XrpcError::Internal(_) | XrpcError::Overloaded + | XrpcError::Warming | XrpcError::NotSettled | XrpcError::AuthRequired(_) | XrpcError::NotImplemented(_)), @@ -2591,6 +2612,126 @@ async fn count_knot_members_by( count_mirror::(&state, q).map(Json) } +#[derive(serde::Deserialize)] +struct InvitesByQuery { + subject: Did, +} + +#[derive(serde::Serialize)] +#[serde(rename_all = "camelCase")] +struct InviteItem { + uri: AtUri, + knot: Did, + #[serde(skip_serializing_if = "Option::is_none")] + repo: Option>, + added_by: Did, + created_at: Datetime, +} + +#[derive(serde::Serialize)] +struct InvitesByOutput { + items: Vec, + pending: Vec>, + truncated: bool, +} + +#[derive(Clone, Copy)] +enum InviteFacet { + Knot, + Repo, +} + +impl InviteFacet { + fn collection(self) -> &'static str { + match self { + Self::Knot => "sh.tangled.knot.memberInvite", + Self::Repo => "sh.tangled.repo.collaboratorInvite", + } + } + + fn accepts(self, scope: &InviteScope) -> bool { + matches!( + (self, scope), + (Self::Knot, InviteScope::Knot(_)) | (Self::Repo, InviteScope::Repo { .. }) + ) + } +} + +fn list_invites_by( + state: &AppState, + q: InvitesByQuery, + facet: InviteFacet, +) -> Result, XrpcError> { + let subject = q.subject; + if !state.invites.discovered() { + return Err(XrpcError::Warming); + } + + let offered: Vec<(InviteScope, Offer)> = state + .invites + .outstanding(&subject) + .into_iter() + .filter(|(scope, _)| facet.accepts(scope)) + .collect(); + let truncated = offered.len() > MAX_OUTSTANDING_OFFERS; + + let items = offered + .into_iter() + .take(MAX_OUTSTANDING_OFFERS) + .map(|(scope, offer)| { + let created_at = datetime_of(offer.offered_at).ok_or_else(|| { + XrpcError::Internal(format!( + "offer to {} has timestamp we can't render, {}", + subject.as_ref(), + offer.offered_at.raw() + )) + })?; + let uri = AtUri::::new_owned(format!( + "at://{}/{}/{}", + scope.subject().as_ref(), + facet.collection(), + subject.as_ref() + )) + .map_err(|e| { + XrpcError::Internal(format!( + "offer to {} won't form record uri, {e}", + subject.as_ref() + )) + })?; + Ok(InviteItem { + uri, + knot: scope.knot().clone(), + repo: match &scope { + InviteScope::Knot(_) => None, + InviteScope::Repo { repo, .. } => Some(repo.clone()), + }, + added_by: offer.invited_by, + created_at, + }) + }) + .collect::, XrpcError>>()?; + + Ok(Json(InvitesByOutput { + items, + pending: state.invites.pending(), + truncated, + })) +} + +async fn list_knot_member_invites_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + list_invites_by(&state, q, InviteFacet::Knot) +} + +async fn list_collaborator_invites_by( + State(state): State, + XrpcQuery(q): XrpcQuery, +) -> Result, XrpcError> { + list_invites_by(&state, q, InviteFacet::Repo) +} + async fn list_label_ops_by( State(state): State, XrpcQuery(q): XrpcQuery>, diff --git a/bobbin/crates/xrpc/tests/aggregation.rs b/bobbin/crates/xrpc/tests/aggregation.rs index ca58fdc9f..98bdb9115 100644 --- a/bobbin/crates/xrpc/tests/aggregation.rs +++ b/bobbin/crates/xrpc/tests/aggregation.rs @@ -2,13 +2,13 @@ use std::sync::Arc; use axum::body::{Body, to_bytes}; use bobbin_edge_index::{ - Coverage, CoverageWatch, EdgeStore, HydrantCursor, IssueStateKind, PageToken, PullStatusKind, - StateIndex, + Coverage, CoverageWatch, EdgeStore, HydrantCursor, InviteIndex, InviteScope, IssueStateKind, + Offer, PageToken, PullStatusKind, StateIndex, }; use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; use bobbin_record_lru::{CacheCapacity, LruRecordStore}; use bobbin_resolver::RepoIdResolver; -use bobbin_runtime::{RuntimeHasher, SystemClock}; +use bobbin_runtime::{RuntimeHasher, SystemClock, UnixMicros}; use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; use bobbin_types::edges::Edge; @@ -56,6 +56,7 @@ struct Harness { server: MockServer, edges: Arc, coverage: Arc, + invites: Arc, state: AppState, } @@ -72,6 +73,7 @@ impl Harness { let issue_states = Arc::new(StateIndex::new(RuntimeHasher::default())); let pull_statuses = Arc::new(StateIndex::new(RuntimeHasher::default())); let coverage = Arc::new(CoverageWatch::new()); + let invites = Arc::new(InviteIndex::new()); let state = AppState::new( Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), SlingshotClient::with_default_http(Url::parse(&server.uri()).unwrap()).unwrap(), @@ -93,11 +95,13 @@ impl Harness { ) as Arc, Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), Arc::new(bobbin_xrpc::default_directory()), - ); + ) + .with_invites(invites.clone()); Self { server, edges, coverage, + invites, state, } } @@ -3089,3 +3093,149 @@ async fn list_recipients_empty_for_unknown_subject() { serde_json::from_slice(&to_bytes(resp.into_body(), 1 << 20).await.unwrap()).unwrap(); assert_eq!(body["dids"], json!([])); } + +const KNOT: &str = "did:web:knot.oyster.cafe"; +const MEMBER_INVITES_BY: &str = "sh.tangled.knot.listMemberInvitesBy"; +const COLLABORATOR_INVITES_BY: &str = "sh.tangled.repo.listCollaboratorInvitesBy"; +const JUNE_FIRST: u64 = 1_780_272_000_000_000; +const JUNE_FOURTH: u64 = 1_780_531_200_000_000; + +fn offer(h: &Harness, repo: Option<&str>, invitee: &str, added_by: &str, micros: u64) { + h.invites.offer( + did(invitee), + match repo { + None => InviteScope::Knot(did(KNOT)), + Some(repo) => InviteScope::Repo { + knot: did(KNOT), + repo: did(repo), + }, + }, + Offer { + invited_by: did(added_by), + offered_at: UnixMicros::new(micros), + }, + ); +} + +async fn invites_body(h: &Harness, endpoint: &str, subject: &str) -> (StatusCode, Value) { + let resp = router(h.state.clone()) + .oneshot( + Request::builder() + .uri(format!("/xrpc/{endpoint}?subject={subject}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + let status = resp.status(); + let body: Value = + serde_json::from_slice(&to_bytes(resp.into_body(), 1 << 20).await.unwrap()).unwrap(); + (status, body) +} + +#[tokio::test] +async fn each_listing_serves_invitee_one_facet_of_their_offers() { + let h = Harness::new().await; + h.invites.all_knots_discovered(); + offer(&h, None, "did:plc:limpet", "did:plc:akshay", JUNE_FIRST); + offer( + &h, + Some("did:plc:scallop"), + "did:plc:limpet", + "did:plc:boltless", + JUNE_FOURTH, + ); + + let (status, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + body, + json!({ + "items": [{ + "uri": "at://did:web:knot.oyster.cafe/sh.tangled.knot.memberInvite/did:plc:limpet", + "knot": KNOT, + "addedBy": "did:plc:akshay", + "createdAt": "2026-06-01T00:00:00.000000Z", + }], + "pending": [], + "truncated": false, + }) + ); + + let (status, body) = invites_body(&h, COLLABORATOR_INVITES_BY, "did:plc:limpet").await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + body["items"], + json!([{ + "uri": "at://did:plc:scallop/sh.tangled.repo.collaboratorInvite/did:plc:limpet", + "knot": KNOT, + "repo": "did:plc:scallop", + "addedBy": "did:plc:boltless", + "createdAt": "2026-06-04T00:00:00.000000Z", + }]), + "deliberi reads knot to tell withdrawal from unread knot, and repo to title row" + ); +} + +#[tokio::test] +async fn index_still_finding_its_knots_answers_warming() { + let h = Harness::new().await; + + let (status, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; + + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); + assert_eq!( + body["error"], "Warming", + "before replay finishes bobbin can't even list which knots it is missing, so partial answers would be guesses" + ); +} + +#[tokio::test] +async fn unread_knot_goes_into_pending_and_answer_still_serves() { + let h = Harness::new().await; + h.invites.all_knots_discovered(); + h.invites.awaiting(did(KNOT)); + offer(&h, None, "did:plc:limpet", "did:plc:akshay", JUNE_FIRST); + + let (status, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + body["pending"], + json!([KNOT]), + "one knot mid-backfill skips withdrawal check for its own offers only" + ); + assert_eq!(body["items"].as_array().unwrap().len(), 1); + + h.invites.backfilled(&did(KNOT)); + let (_, body) = invites_body(&h, MEMBER_INVITES_BY, "did:plc:limpet").await; + assert_eq!(body["pending"], json!([])); +} + +#[tokio::test] +async fn flood_of_offers_is_limited_and_says_so() { + let h = Harness::new().await; + h.invites.all_knots_discovered(); + (0..200u64).for_each(|n| { + offer( + &h, + Some(&format!("did:plc:repo{n:0>3}")), + "did:plc:limpet", + "did:plc:boltless", + JUNE_FOURTH + n, + ) + }); + + let (status, body) = invites_body(&h, COLLABORATOR_INVITES_BY, "did:plc:limpet").await; + + assert_eq!(status, StatusCode::OK); + assert_eq!( + body["items"].as_array().unwrap().len(), + 128, + "repositories are free to create, so one operator could otherwise mint one row per repo in strangers' inboxes" + ); + assert_eq!(body["truncated"], json!(true)); + assert_eq!( + body["items"][0]["createdAt"], "2026-06-04T00:00:00.000199Z", + "newest offers are kept, since those are the ones somebody might still answer" + ); +} diff --git a/lexicons/knot/listMemberInvitesBy.json b/lexicons/knot/listMemberInvitesBy.json new file mode 100644 index 000000000..8ee24b529 --- /dev/null +++ b/lexicons/knot/listMemberInvitesBy.json @@ -0,0 +1,70 @@ +{ + "lexicon": 1, + "id": "sh.tangled.knot.listMemberInvitesBy", + "defs": { + "main": { + "type": "query", + "description": "Grabbing the membership offers that a given DID has lined up. The inverse of sh.tangled.knot.listMemberInvites.", + "parameters": { + "type": "params", + "required": ["subject"], + "properties": { + "subject": { + "type": "string", + "format": "did", + "description": "Actor DID offered membership." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["items", "pending", "truncated"], + "properties": { + "items": { + "type": "array", + "items": { "type": "ref", "ref": "#inviteItem" } + }, + "pending": { + "type": "array", + "description": "The knots that this indexation hasn't yet read.", + "items": { "type": "string", "format": "did" } + }, + "truncated": { + "type": "boolean", + "description": "Per-subject cut dropped offers from items. If truncated is true, missing offers won't advertise that they were withdrawn, same as with pending knots." + } + } + } + } + }, + "inviteItem": { + "type": "object", + "description": "Offer awaiting the subject's acceptance.", + "required": ["uri", "knot", "addedBy", "createdAt"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri", + "description": "Invite record on knot's repo." + }, + "knot": { + "type": "string", + "format": "did", + "description": "did:web of knot that made this offer. Check against pending." + }, + "addedBy": { + "type": "string", + "format": "did", + "description": "DID that made the offer." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When the knot made the offer. This listing sorts on it by newest first." + } + } + } + } +} diff --git a/lexicons/repo/listCollaboratorInvitesBy.json b/lexicons/repo/listCollaboratorInvitesBy.json new file mode 100644 index 000000000..96b1cbbca --- /dev/null +++ b/lexicons/repo/listCollaboratorInvitesBy.json @@ -0,0 +1,75 @@ +{ + "lexicon": 1, + "id": "sh.tangled.repo.listCollaboratorInvitesBy", + "defs": { + "main": { + "type": "query", + "description": "Collaboration offers an actor DID has waiting for them. The inverse of sh.tangled.repo.listCollaboratorInvites, which knots serve only for their own repo.", + "parameters": { + "type": "params", + "required": ["subject"], + "properties": { + "subject": { + "type": "string", + "format": "did", + "description": "Actor DID offered collaboration." + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["items", "pending", "truncated"], + "properties": { + "items": { + "type": "array", + "items": { "type": "ref", "ref": "#inviteItem" } + }, + "pending": { + "type": "array", + "description": "Knots that this index hasn't read yet", + "items": { "type": "string", "format": "did" } + }, + "truncated": { + "type": "boolean", + "description": "Per-subject dropped offers from items. If truncated is true, missing offers won't say if they were withdrawn, same with pending knots." + } + } + } + } + }, + "inviteItem": { + "type": "object", + "description": "An offer awaiting the subject's acceptance.", + "required": ["uri", "knot", "repo", "addedBy", "createdAt"], + "properties": { + "uri": { + "type": "string", + "format": "at-uri", + "description": "Invite record on the repository." + }, + "knot": { + "type": "string", + "format": "did", + "description": "did:web of knot hosting this repo." + }, + "repo": { + "type": "string", + "format": "did", + "description": "DID of repo that made the offer." + }, + "addedBy": { + "type": "string", + "format": "did", + "description": "DID that made the offer." + }, + "createdAt": { + "type": "string", + "format": "datetime", + "description": "When the repository made the offer. This listing sorts it by newest first." + } + } + } + } +} diff --git a/lexicons/temp/notification/listNotifications.json b/lexicons/temp/notification/listNotifications.json index d6fcbee61..dc2f427b7 100644 --- a/lexicons/temp/notification/listNotifications.json +++ b/lexicons/temp/notification/listNotifications.json @@ -63,7 +63,7 @@ }, "type": { "type": "string", - "description": "Notification type: repo_starred, issue_created, issue_commented, issue_closed, issue_reopen, issue_assigned, issue_unassigned, pull_created, pull_commented, pull_merged, pull_closed, pull_reopen, pull_assigned, pull_unassigned, followed, user_mentioned." + "description": "Notification type: repo_starred, issue_created, issue_commented, issue_closed, issue_reopen, issue_assigned, issue_unassigned, pull_created, pull_commented, pull_merged, pull_closed, pull_reopen, pull_assigned, pull_unassigned, followed, user_mentioned, knot_invited, collaborator_invited." }, "category": { "type": "string", @@ -86,6 +86,11 @@ "format": "did", "description": "DID of the related repository, if applicable." }, + "knotDid": { + "type": "string", + "format": "did", + "description": "did:web of knot that offered membership. Only knot_invited has it for now! Though in future more notifs may use it." + }, "issueAt": { "type": "string", "format": "at-uri", diff --git a/web/lex.config.ts b/web/lex.config.ts index e3ac88f6a..ceacce749 100644 --- a/web/lex.config.ts +++ b/web/lex.config.ts @@ -30,7 +30,6 @@ export default defineLexiconConfig({ "../lexicons/string/**/*.json", "../lexicons/temp/**/*.json", "../lexicons/sync/**/*.json", - "../bobbin/crates/types/lexicons/**/*.json" ] } }); diff --git a/web/src/lib/api/lexicons/index.ts b/web/src/lib/api/lexicons/index.ts index 737ca2580..c00b6b3a2 100644 --- a/web/src/lib/api/lexicons/index.ts +++ b/web/src/lib/api/lexicons/index.ts @@ -129,6 +129,7 @@ export * as ShTangledKnotCountMembersBy from "./types/sh/tangled/knot/countMembe export * as ShTangledKnotListKeys from "./types/sh/tangled/knot/listKeys.js"; export * as ShTangledKnotListKnots from "./types/sh/tangled/knot/listKnots.js"; export * as ShTangledKnotListMemberInvites from "./types/sh/tangled/knot/listMemberInvites.js"; +export * as ShTangledKnotListMemberInvitesBy from "./types/sh/tangled/knot/listMemberInvitesBy.js"; export * as ShTangledKnotListMembers from "./types/sh/tangled/knot/listMembers.js"; export * as ShTangledKnotListMembersBy from "./types/sh/tangled/knot/listMembersBy.js"; export * as ShTangledKnotMember from "./types/sh/tangled/knot/member.js"; @@ -220,6 +221,7 @@ export * as ShTangledRepoLanguages from "./types/sh/tangled/repo/languages.js"; export * as ShTangledRepoListArtifacts from "./types/sh/tangled/repo/listArtifacts.js"; export * as ShTangledRepoListArtifactsBy from "./types/sh/tangled/repo/listArtifactsBy.js"; export * as ShTangledRepoListCollaboratorInvites from "./types/sh/tangled/repo/listCollaboratorInvites.js"; +export * as ShTangledRepoListCollaboratorInvitesBy from "./types/sh/tangled/repo/listCollaboratorInvitesBy.js"; export * as ShTangledRepoListCollaborators from "./types/sh/tangled/repo/listCollaborators.js"; export * as ShTangledRepoListCollaboratorsBy from "./types/sh/tangled/repo/listCollaboratorsBy.js"; export * as ShTangledRepoListIssues from "./types/sh/tangled/repo/listIssues.js"; diff --git a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts index c64dd1a38..dee13f586 100644 --- a/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts +++ b/web/src/lib/api/lexicons/types/org/tangled/temp/notification/listNotifications.ts @@ -58,6 +58,10 @@ const _notificationSchema = /*#__PURE__*/ v.object({ * AT-URI of the related org.tangled.issue.issue record, if applicable. */ issueAt: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.resourceUriString()), + /** + * did:web of knot that offered membership. Only knot_invited has it for now! Though in future more notifs may use it. + */ + knotDid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.didString()), /** * AT-URI of the related org.tangled.pulls.pull record, if applicable. */ @@ -68,7 +72,7 @@ const _notificationSchema = /*#__PURE__*/ v.object({ */ repoDid: /*#__PURE__*/ v.optional(/*#__PURE__*/ v.didString()), /** - * Notification type: repo_starred, issue_created, issue_commented, issue_closed, issue_reopen, issue_assigned, issue_unassigned, pull_created, pull_commented, pull_merged, pull_closed, pull_reopen, pull_assigned, pull_unassigned, followed, user_mentioned. + * Notification type: repo_starred, issue_created, issue_commented, issue_closed, issue_reopen, issue_assigned, issue_unassigned, pull_created, pull_commented, pull_merged, pull_closed, pull_reopen, pull_assigned, pull_unassigned, followed, user_mentioned, knot_invited, collaborator_invited. */ type: /*#__PURE__*/ v.string(), /** diff --git a/web/src/lib/api/lexicons/types/sh/tangled/knot/listMemberInvitesBy.ts b/web/src/lib/api/lexicons/types/sh/tangled/knot/listMemberInvitesBy.ts new file mode 100644 index 000000000..fa3a5a657 --- /dev/null +++ b/web/src/lib/api/lexicons/types/sh/tangled/knot/listMemberInvitesBy.ts @@ -0,0 +1,72 @@ +import type {} from "@atcute/lexicons"; +import * as v from "@atcute/lexicons/validations"; +import type {} from "@atcute/lexicons/ambient"; + +const _inviteItemSchema = /*#__PURE__*/ v.object({ + $type: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.literal("sh.tangled.knot.listMemberInvitesBy#inviteItem"), + ), + /** + * DID that made the offer. + */ + addedBy: /*#__PURE__*/ v.didString(), + /** + * When the knot made the offer. This listing sorts on it by newest first. + */ + createdAt: /*#__PURE__*/ v.datetimeString(), + /** + * did:web of knot that made this offer. Check against pending. + */ + knot: /*#__PURE__*/ v.didString(), + /** + * Invite record on knot's repo. + */ + uri: /*#__PURE__*/ v.resourceUriString(), +}); +const _mainSchema = /*#__PURE__*/ v.query( + "sh.tangled.knot.listMemberInvitesBy", + { + params: /*#__PURE__*/ v.object({ + /** + * Actor DID offered membership. + */ + subject: /*#__PURE__*/ v.didString(), + }), + output: { + type: "lex", + schema: /*#__PURE__*/ v.object({ + get items() { + return /*#__PURE__*/ v.array(inviteItemSchema); + }, + /** + * The knots that this indexation hasn't yet read. + */ + pending: /*#__PURE__*/ v.array(/*#__PURE__*/ v.didString()), + /** + * Per-subject cut dropped offers from items. If truncated is true, missing offers won't advertise that they were withdrawn, same as with pending knots. + */ + truncated: /*#__PURE__*/ v.boolean(), + }), + }, + }, +); + +type inviteItem$schematype = typeof _inviteItemSchema; +type main$schematype = typeof _mainSchema; + +export interface inviteItemSchema extends inviteItem$schematype {} +export interface mainSchema extends main$schematype {} + +export const inviteItemSchema = _inviteItemSchema as inviteItemSchema; +export const mainSchema = _mainSchema as mainSchema; + +export interface InviteItem extends v.InferInput {} + +export interface $params extends v.InferInput {} +export interface $output extends v.InferXRPCBodyInput {} + +declare module "@atcute/lexicons/ambient" { + interface XRPCQueries { + "sh.tangled.knot.listMemberInvitesBy": mainSchema; + } +} diff --git a/web/src/lib/api/lexicons/types/sh/tangled/repo/listCollaboratorInvitesBy.ts b/web/src/lib/api/lexicons/types/sh/tangled/repo/listCollaboratorInvitesBy.ts new file mode 100644 index 000000000..106ad52f4 --- /dev/null +++ b/web/src/lib/api/lexicons/types/sh/tangled/repo/listCollaboratorInvitesBy.ts @@ -0,0 +1,78 @@ +import type {} from "@atcute/lexicons"; +import * as v from "@atcute/lexicons/validations"; +import type {} from "@atcute/lexicons/ambient"; + +const _inviteItemSchema = /*#__PURE__*/ v.object({ + $type: /*#__PURE__*/ v.optional( + /*#__PURE__*/ v.literal( + "sh.tangled.repo.listCollaboratorInvitesBy#inviteItem", + ), + ), + /** + * DID that made the offer. + */ + addedBy: /*#__PURE__*/ v.didString(), + /** + * When the repository made the offer. This listing sorts it by newest first. + */ + createdAt: /*#__PURE__*/ v.datetimeString(), + /** + * did:web of knot hosting this repo. + */ + knot: /*#__PURE__*/ v.didString(), + /** + * DID of repo that made the offer. + */ + repo: /*#__PURE__*/ v.didString(), + /** + * Invite record on the repository. + */ + uri: /*#__PURE__*/ v.resourceUriString(), +}); +const _mainSchema = /*#__PURE__*/ v.query( + "sh.tangled.repo.listCollaboratorInvitesBy", + { + params: /*#__PURE__*/ v.object({ + /** + * Actor DID offered collaboration. + */ + subject: /*#__PURE__*/ v.didString(), + }), + output: { + type: "lex", + schema: /*#__PURE__*/ v.object({ + get items() { + return /*#__PURE__*/ v.array(inviteItemSchema); + }, + /** + * Knots that this index hasn't read yet + */ + pending: /*#__PURE__*/ v.array(/*#__PURE__*/ v.didString()), + /** + * Per-subject dropped offers from items. If truncated is true, missing offers won't say if they were withdrawn, same with pending knots. + */ + truncated: /*#__PURE__*/ v.boolean(), + }), + }, + }, +); + +type inviteItem$schematype = typeof _inviteItemSchema; +type main$schematype = typeof _mainSchema; + +export interface inviteItemSchema extends inviteItem$schematype {} +export interface mainSchema extends main$schematype {} + +export const inviteItemSchema = _inviteItemSchema as inviteItemSchema; +export const mainSchema = _mainSchema as mainSchema; + +export interface InviteItem extends v.InferInput {} + +export interface $params extends v.InferInput {} +export interface $output extends v.InferXRPCBodyInput {} + +declare module "@atcute/lexicons/ambient" { + interface XRPCQueries { + "sh.tangled.repo.listCollaboratorInvitesBy": mainSchema; + } +}