diff --git a/Cargo.lock b/Cargo.lock index 88aa942aa..5ccce1279 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4838,6 +4838,7 @@ dependencies = [ "knot-runtime", "knot-types", "tempfile", + "tokio", ] [[package]] @@ -4939,6 +4940,21 @@ dependencies = [ "url", ] +[[package]] +name = "knot-consent" +version = "2.0.0" +dependencies = [ + "bytes", + "http", + "knot-atproto", + "knot-index", + "knot-runtime", + "knot-types", + "tokio", + "tracing", + "url", +] + [[package]] name = "knot-edge" version = "2.0.0" @@ -5350,6 +5366,7 @@ dependencies = [ "knot-cob", "knot-cobs", "knot-config", + "knot-consent", "knot-edge", "knot-events", "knot-fixtures", @@ -5389,6 +5406,7 @@ dependencies = [ "knot-atproto", "knot-cob", "knot-cobs", + "knot-consent", "knot-events", "knot-fixtures", "knot-git", @@ -5462,6 +5480,7 @@ dependencies = [ "knot-cob", "knot-cobs", "knot-config", + "knot-consent", "knot-events", "knot-fixtures", "knot-git", diff --git a/Cargo.toml b/Cargo.toml index e9796510a..83ca0fd77 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -63,6 +63,7 @@ knot-index = { path = "knot2/crates/knot-index" } knot-keyfill = { path = "knot2/crates/knot-keyfill" } knot-cache = { path = "knot2/crates/knot-cache" } knot-acl = { path = "knot2/crates/knot-acl" } +knot-consent = { path = "knot2/crates/knot-consent" } knot-atproto = { path = "knot2/crates/knot-atproto" } knot-messages = { path = "knot2/crates/knot-messages" } knot-postreceive = { path = "knot2/crates/knot-postreceive" } diff --git a/knot2/crates/knot-atproto/src/jwt.rs b/knot2/crates/knot-atproto/src/jwt.rs index b433e596c..15696e2ca 100644 --- a/knot2/crates/knot-atproto/src/jwt.rs +++ b/knot2/crates/knot-atproto/src/jwt.rs @@ -5,7 +5,9 @@ use knot_types::crypto::{KeyCodec, PublicKey as CryptoKey}; use knot_types::service_auth::{ JwtHeader, ParsedJwt, PublicKey as VerifyKey, ServiceAuthClaims, ServiceAuthError, parse_jwt, }; -use knot_types::{AccountDid, CowStr, Did, DidService, KnotId, Nsid, ServiceDid, UnixSeconds}; +use knot_types::{ + AccountDid, CowStr, Did, DidService, KnotId, Nsid, RepoDid, ServiceDid, UnixSeconds, +}; pub(crate) const CLOCK_SKEW_SECS: i64 = 60; pub(crate) const SERVICE_TOKEN_LIFETIME_SECS: i64 = 60; @@ -172,25 +174,21 @@ pub trait TokenAudience { fn canonicalizes(&self, claimed: &str) -> bool; } -impl TokenAudience for KnotId { - fn as_str(&self) -> &str { - KnotId::as_str(self) - } +macro_rules! token_audience { + ($($id:ty),+ $(,)?) => { + $(impl TokenAudience for $id { + fn as_str(&self) -> &str { + <$id>::as_str(self) + } - fn canonicalizes(&self, claimed: &str) -> bool { - KnotId::new(claimed).is_ok_and(|aud| aud.as_str() == KnotId::as_str(self)) - } + fn canonicalizes(&self, claimed: &str) -> bool { + <$id>::new(claimed).is_ok_and(|aud| <$id>::as_str(&aud) == <$id>::as_str(self)) + } + })+ + }; } -impl TokenAudience for ServiceDid { - fn as_str(&self) -> &str { - ServiceDid::as_str(self) - } - - fn canonicalizes(&self, claimed: &str) -> bool { - ServiceDid::new(claimed).is_ok_and(|aud| aud.as_str() == ServiceDid::as_str(self)) - } -} +token_audience!(KnotId, ServiceDid, RepoDid); pub fn check_claims( parsed: &ParsedJwt, @@ -202,7 +200,7 @@ pub fn check_claims( let exp = UnixSeconds::new(claims.exp); let iat = UnixSeconds::new(claims.iat); - if !audience.canonicalizes(claims.aud.as_str()) { + if !audience.canonicalizes(claims.aud.audience().as_str()) { return Err(JwtError::AudienceMismatch { expected: audience.as_str().to_string(), actual: claims.aud.as_str().to_string(), @@ -295,6 +293,26 @@ mod tests { assert_eq!(issuer(&parsed).unwrap(), AccountDid::new(SQUID).unwrap()); } + #[test] + fn a_repository_audience_takes_the_same_service_reference_as_the_knot() { + let signing = signer(1); + let repo = RepoDid::new("did:plc:3fwecdnvtcscjnrx2p4n7alz").unwrap(); + let checked = |aud: &str| { + let token = mint_claims(&signing, &claims(SQUID, aud, UnixSeconds::new(1_000), METHOD)); + check_claims(&parse(&token).unwrap(), &repo, &member_method(), UnixSeconds::new(900)) + }; + + checked(repo.as_str()).unwrap(); + checked(&format!("{}#tangled_knot", repo.as_str())).unwrap(); + assert!( + matches!( + checked("did:plc:qxx2hmw5gv6ghjvcxl63q5d3#tangled_knot"), + Err(JwtError::AudienceMismatch { .. }) + ), + "a fragment let another DID pass as this repository" + ); + } + struct ClaimsCase { name: &'static str, claims: fn() -> serde_json::Value, @@ -369,6 +387,18 @@ mod tests { now: 900, expect: |r| r.is_ok(), }, + ClaimsCase { + name: "a service reference matches on the DID alone", + claims: || serde_json::json!({ "iss": SQUID, "aud": "did:web:nel.pet#tangled_knot", "exp": 1_000, "iat": 940, "lxm": METHOD }), + now: 900, + expect: |r| r.is_ok(), + }, + ClaimsCase { + name: "a fragment on another knot's DID", + claims: || serde_json::json!({ "iss": SQUID, "aud": "did:web:oyster.cafe#tangled_knot", "exp": 1_000, "iat": 940, "lxm": METHOD }), + now: 900, + expect: |r| matches!(r, Err(JwtError::AudienceMismatch { .. })), + }, ]; #[test] @@ -432,6 +462,11 @@ mod tests { assert_eq!(parsed.claims().iat, 1_000); assert_eq!(parsed.claims().exp, 1_000 + SERVICE_TOKEN_LIFETIME_SECS); + assert_eq!( + parsed.claims().aud.as_str(), + audience.as_str(), + "the reference PDS compares its configured DID byte for byte, so a fragment would 401" + ); check_claims(&parsed, &audience, &bound, UnixSeconds::new(1_005)).unwrap(); assert_eq!(nonce(&parsed).unwrap().as_str(), "nonce-minted"); assert_eq!(issuer(&parsed).unwrap(), AccountDid::new(KNOT).unwrap()); diff --git a/knot2/crates/knot-atproto/src/lib.rs b/knot2/crates/knot-atproto/src/lib.rs index c34f46320..2fa2b0275 100644 --- a/knot2/crates/knot-atproto/src/lib.rs +++ b/knot2/crates/knot-atproto/src/lib.rs @@ -53,7 +53,7 @@ use std::time::Duration; use futures::stream::{self, TryStreamExt}; use http::StatusCode; use knot_cache::{ - Admitted, AsyncCache, EntryCount, Expiring, GroupQuota, MokaFuture, Quotas, Rejected, + Admitted, AsyncCache, EntryCount, Expiring, Filled, GroupQuota, MokaFuture, Quotas, Rejected, TotalQuota, Weight, }; use knot_runtime::{ @@ -61,8 +61,8 @@ use knot_runtime::{ UnixMicros, }; use knot_types::{ - AccountDid, Collection, Handle, HttpStatus, KnotId, Nsid, OfferedKey, RepoDid, RepoRkey, Rkey, - UnixSeconds, + AccountDid, Collection, DidRkey, Handle, HttpStatus, KnotId, Nsid, OfferedKey, RepoDid, + RepoRkey, Rkey, UnixSeconds, }; use pubkeys::Cursor; use serde::{Deserialize, Serialize}; @@ -238,10 +238,10 @@ impl Atproto { } pub async fn resolve_identity(&self, did: &AccountDid) -> Result { - self.resolve_identity_inner(did).await.map_err(Into::into) + Ok(self.identity(did).await?.value) } - async fn resolve_identity_inner(&self, did: &AccountDid) -> Result { + async fn identity(&self, did: &AccountDid) -> Result, ResolveError> { let now = self.clock.now_unix_micros(); let filled = self .identities @@ -253,7 +253,7 @@ impl Atproto { .await; let fresh = filled.fresh; match filled.value.resolution { - Resolution::Found(identity) => Ok(identity), + Resolution::Found(value) => Ok(Filled { value, fresh }), Resolution::Transient(error) => Err(error), Resolution::Failed(error) if fresh || error.is_gone() => Err(error), Resolution::Failed(_) => Err(ResolveError::RecentlyFailed { did: did.clone() }), @@ -327,7 +327,7 @@ impl Atproto { }, Err(dns_error) => return Err(dns_error), }; - let identity = self.resolve_identity_inner(&candidate).await?; + let identity = self.identity(&candidate).await?.value; if identity.claims_handle(handle) { Ok(candidate) } else { @@ -523,8 +523,48 @@ impl Atproto { owner: &AccountDid, rkey: &RepoRkey, ) -> Result { - let identity = self.resolve_identity(owner).await?; - let url = get_record_url(&identity.pds, owner, &repo_collection(), rkey)?; + self.record_present(owner, &repo_collection(), rkey.as_str()) + .await + } + + pub async fn subject_record_present( + &self, + owner: &AccountDid, + collection: &Nsid, + subject: &DidRkey, + ) -> Result { + self.record_present(owner, collection, subject.as_str()) + .await + } + + async fn record_present( + &self, + owner: &AccountDid, + collection: &Nsid, + rkey: &str, + ) -> Result { + let cached = self.identity(owner).await?; + let read = self.read_record(&cached.value, owner, collection, rkey).await?; + match read { + RecordPresence::Absent if !cached.fresh => { + let current = self.fetch_identity(owner).await?; + match current.pds.url() == cached.value.pds.url() { + true => Ok(RecordPresence::Absent), + false => self.read_record(¤t, owner, collection, rkey).await, + } + } + settled => Ok(settled), + } + } + + async fn read_record( + &self, + identity: &Identity, + owner: &AccountDid, + collection: &Nsid, + rkey: &str, + ) -> Result { + let url = get_record_url(&identity.pds, owner, collection, rkey)?; resolve::guard_fetch_url(&url)?; let response = self.http.execute(HttpRequest::get(url)).await?; match response.status { @@ -554,13 +594,34 @@ impl Atproto { token: &ServiceJwt, method: &Nsid, replay: ReplayGuard, + ) -> Result { + self.verify_addressed(token, method, &self.knot_did, replay) + .await + } + + pub async fn verify_service_jwt_for_repo( + &self, + token: &ServiceJwt, + method: &Nsid, + repo: &RepoDid, + ) -> Result { + self.verify_addressed(token, method, repo, ReplayGuard::SingleUse) + .await + } + + async fn verify_addressed( + &self, + token: &ServiceJwt, + method: &Nsid, + audience: &impl jwt::TokenAudience, + replay: ReplayGuard, ) -> Result { let parsed = jwt::parse(token)?; let issuer = jwt::issuer(&parsed)?; let now_micros = self.clock.now_unix_micros(); let now = UnixSeconds::new((now_micros.get() / 1_000_000) as i64); - jwt::check_claims(&parsed, &self.knot_did, method, now)?; + jwt::check_claims(&parsed, audience, method, now)?; let jti = jwt::nonce(&parsed)?; let identity = self.resolve_identity(&issuer).await?; @@ -757,14 +818,14 @@ fn get_record_url( pds: &PdsEndpoint, owner: &AccountDid, collection: &Nsid, - rkey: &RepoRkey, + rkey: &str, ) -> Result { let method: Nsid = Nsid::new_static("com.atproto.repo.getRecord").expect("literal nsid parses"); let mut url = xrpc_url(pds, &method)?; url.query_pairs_mut() .append_pair("repo", owner.as_str()) .append_pair("collection", collection.as_str()) - .append_pair("rkey", rkey.as_str()); + .append_pair("rkey", rkey); Ok(url) } @@ -1204,6 +1265,83 @@ mod tests { ); } + fn asking() -> (AccountDid, Nsid, DidRkey) { + let nsid = Nsid::new_static("sh.tangled.knot.memberAcceptance").unwrap(); + (did(SQUID), nsid, DidRkey::new(KNOT).unwrap()) + } + + #[derive(Clone, Copy)] + enum Hosting { + Moved, + Blind, + } + + fn serving(then: Hosting) -> (Atproto, Arc) { + let (signing, lookups) = (signer(9), Arc::new(AtomicUsize::new(0))); + let counted = Arc::clone(&lookups); + let missing = Bytes::from_static(b"{\"error\":\"RecordNotFound\"}"); + let http = FakeHttp::new(move |request| match request.url.host_str() { + Some("plc.directory") => match counted.fetch_add(1, Ordering::SeqCst) { + 0 => Ok(ok(squid_doc(&signing))), + _ => match then { + Hosting::Blind => Ok(status(StatusCode::BAD_REQUEST, Bytes::new())), + Hosting::Moved => Ok(ok(did_doc(DocSpec { + id: SQUID, + signing: &signing, + handle: "nel.pet", + pds: "https://pds.anemone.town", + method: MethodKind::Multikey, + }))), + }, + }, + Some("pds.anemone.town") => Ok(ok(Bytes::new())), + _ => Ok(status(StatusCode::BAD_REQUEST, missing.clone())), + }); + let atproto = Atproto::new(http, clock(), knot_did(KNOT), plc()); + (atproto, lookups) + } + + #[tokio::test] + async fn an_absence_is_re_read_against_the_pds_the_account_moved_to() { + let (atproto, lookups) = serving(Hosting::Moved); + let (owner, nsid, rkey) = asking(); + let sealed = atproto.subject_record_present(&owner, &nsid, &rkey).await; + assert_eq!( + (sealed.unwrap(), lookups.load(Ordering::SeqCst)), + (RecordPresence::Absent, 1), + "nothing can have moved between the document this call read and the record it asked \ + for, so a document this call resolved needs no second look" + ); + let moved = atproto.subject_record_present(&owner, &nsid, &rkey).await; + assert_eq!( + (moved.unwrap(), lookups.load(Ordering::SeqCst)), + (RecordPresence::Present, 2), + "a 404 from the pds the account left is no evidence the record was deleted, and the \ + document sat inside its ttl, so only an absence buys the second read" + ); + } + + #[tokio::test] + async fn directory_400_fails_read_and_leaves_document_intact() { + let (atproto, lookups) = serving(Hosting::Blind); + let (owner, nsid, rkey) = asking(); + atproto.resolve_identity(&owner).await.unwrap(); + let unproven = atproto.subject_record_present(&owner, &nsid, &rkey).await; + let intact = atproto.resolve_identity(&owner).await.unwrap(); + assert!( + unproven.is_err(), + "the 404 came off a pds the knot can no longer confirm is the current one, so \ + calling the record deleted would revoke on a guess" + ); + assert_eq!( + (intact.pds.url().as_str(), lookups.load(Ordering::SeqCst)), + ("https://pds.oyster.cafe/", 2), + "a confirmation that fails leaves the document every other caller reads alone: a 400 \ + from the directory mustn't cost this account its ssh keys or its service auth, and \ + the surviving document answers this third read off one fill and one confirmation" + ); + } + #[tokio::test] async fn pubkey_resolution_follows_the_cursor() { let doc_key = signer(9); diff --git a/knot2/crates/knot-consent/Cargo.toml b/knot2/crates/knot-consent/Cargo.toml new file mode 100644 index 000000000..835a2cbbb --- /dev/null +++ b/knot2/crates/knot-consent/Cargo.toml @@ -0,0 +1,19 @@ +[package] +name = "knot-consent" +version = "2.0.0" +edition.workspace = true +rust-version.workspace = true +license.workspace = true + +[dependencies] +knot-types = { workspace = true } +knot-runtime = { workspace = true } +knot-index = { workspace = true } +knot-atproto = { workspace = true } +tracing = { workspace = true } + +[dev-dependencies] +bytes = { workspace = true } +http = { workspace = true } +tokio = { workspace = true } +url = { workspace = true } diff --git a/knot2/crates/knot-consent/src/lib.rs b/knot2/crates/knot-consent/src/lib.rs new file mode 100644 index 000000000..41f188552 --- /dev/null +++ b/knot2/crates/knot-consent/src/lib.rs @@ -0,0 +1,124 @@ +use std::pin::Pin; + +use knot_atproto::{Atproto, RecordPresence}; +use knot_index::{Acceptance, Acceptances, Read, Reading}; +use knot_runtime::{Clock, HttpTransport}; +use knot_types::{AccountDid, DidRkey, KnotId, Nsid}; + +pub const OF_MEMBERSHIP: &str = "sh.tangled.knot.memberAcceptance"; +pub const OF_COLLABORATION: &str = "sh.tangled.repo.collaboratorAcceptance"; + +pub struct Pds<'a, H, C> { + atproto: &'a Atproto, + knot: &'a KnotId, +} + +impl<'a, H, C> Pds<'a, H, C> { + pub const fn reading(atproto: &'a Atproto, knot: &'a KnotId) -> Self { + Self { atproto, knot } + } +} + +impl Acceptances for Pds<'_, H, C> { + fn read<'a>( + &'a self, + scope: Acceptance<'a>, + subject: &'a AccountDid, + ) -> Pin + Send + 'a>> { + Box::pin(async move { + let (collection, named) = match scope { + Acceptance::OfMembership => (OF_MEMBERSHIP, self.knot.as_str()), + Acceptance::OfCollaboration(repo) => (OF_COLLABORATION, repo.as_str()), + }; + let collection = + Nsid::new_static(collection).expect("acceptance collections are literal nsids"); + let taken = |reading| Read::taken(reading, self.atproto.now().seconds()); + let Ok(rkey) = DidRkey::new(named) else { + tracing::warn!( + named, + "acceptance record key isn't a DID this knot can ask a PDS for, so no \ + acceptance for it can exist anywhere" + ); + return taken(Reading::Absent); + }; + match self + .atproto + .subject_record_present(subject, &collection, &rkey) + .await + { + Ok(RecordPresence::Present) => taken(Reading::Present), + Ok(RecordPresence::Absent) => taken(Reading::Absent), + Err(error) => { + let gone = error.is_gone(); + tracing::debug!(subject = subject.as_str(), %error, gone, "acceptance record didn't resolve"); + match gone { + true => taken(Reading::Absent), + false => taken(Reading::Unreachable), + } + } + } + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + use bytes::Bytes; + use http::StatusCode; + use knot_atproto::PlcDirectory; + use knot_runtime::{FakeHttp, HttpResponse, ManualClock}; + use knot_types::{RepoDid, UnixMicros}; + + const NOW: u64 = 1_000_000_000; + + fn knot() -> KnotId { + KnotId::new("did:web:nel.pet").unwrap() + } + + async fn reading_for(status: StatusCode, scope: Acceptance<'_>) -> Read { + let atproto = Atproto::new( + FakeHttp::new(move |_| { + Ok(HttpResponse { status, headers: http::HeaderMap::new(), body: Bytes::new() }) + }), + ManualClock::new(UnixMicros::new(NOW)), + knot(), + PlcDirectory::new(url::Url::parse("https://plc.directory/").unwrap()).unwrap(), + ); + let knot = knot(); + Pds::reading(&atproto, &knot) + .read(scope, &AccountDid::new("did:plc:squid").unwrap()) + .await + } + + fn taken(reading: Reading) -> Read { + Read::taken(reading, UnixMicros::new(NOW).seconds()) + } + + #[tokio::test] + async fn a_404_settles_the_read_and_a_500_opens_the_grace_window() { + assert_eq!( + reading_for(StatusCode::NOT_FOUND, Acceptance::OfMembership).await, + taken(Reading::Absent), + "a directory that answers 404 answered, so the grace window has nothing to cover \ + and a deleted account loses the standing at the next use" + ); + assert_eq!( + reading_for(StatusCode::INTERNAL_SERVER_ERROR, Acceptance::OfMembership).await, + taken(Reading::Unreachable), + "a 500 is an outage rather than an answer, and telling the two apart is the whole \ + difference between a bounded grace and a revocation" + ); + } + + #[tokio::test] + async fn ported_repository_settles_absent_from_its_did_alone() { + let ported = RepoDid::new("did:web:nel.pet%3A8080").unwrap(); + assert_eq!( + reading_for(StatusCode::OK, Acceptance::OfCollaboration(&ported)).await, + taken(Reading::Absent), + "a record key comes from the bare host, so a ported did settles Absent unasked" + ); + } +} diff --git a/knot2/crates/knot-index/Cargo.toml b/knot2/crates/knot-index/Cargo.toml index fe5502d4f..58871831d 100644 --- a/knot2/crates/knot-index/Cargo.toml +++ b/knot2/crates/knot-index/Cargo.toml @@ -5,6 +5,9 @@ edition.workspace = true rust-version.workspace = true license.workspace = true +[features] +test-support = [] + [dependencies] knot-types = { workspace = true } knot-git = { workspace = true } diff --git a/knot2/crates/knot-index/src/consent.rs b/knot2/crates/knot-index/src/consent.rs new file mode 100644 index 000000000..d1b95d368 --- /dev/null +++ b/knot2/crates/knot-index/src/consent.rs @@ -0,0 +1,650 @@ +use std::pin::Pin; +use std::sync::Arc; +use std::time::Duration; + +use tokio::sync::Mutex; + +use knot_cobs::Offer; +use knot_types::{AccountDid, RepoDid, UnixSeconds}; + +use crate::after; +use crate::intern::{AccountKey, RepoKey}; +use crate::projections::Provenance; +use crate::{Index, Resolved}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct AcceptanceTtl(Duration); + +impl AcceptanceTtl { + pub const fn from_secs(secs: u64) -> Self { + Self(Duration::from_secs(secs)) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct AcceptanceGrace(Duration); + +impl AcceptanceGrace { + pub const fn from_secs(secs: u64) -> Self { + Self(Duration::from_secs(secs)) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct Freshness { + ttl: AcceptanceTtl, + grace: AcceptanceGrace, +} + +impl Freshness { + pub const DEFAULT: Self = Self::new( + AcceptanceTtl::from_secs(60), + AcceptanceGrace::from_secs(21_600), + ); + + pub const fn new(ttl: AcceptanceTtl, grace: AcceptanceGrace) -> Self { + Self { ttl, grace } + } + + const fn leased(self, read_at: UnixSeconds, now: UnixSeconds) -> Lease { + Lease { + read_at, + expires_at: after(now, self.ttl.0), + } + } + + const fn graces(self, read_at: UnixSeconds, now: UnixSeconds) -> bool { + after(read_at, self.grace.0).get() > now.get() + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +struct Lease { + read_at: UnixSeconds, + expires_at: UnixSeconds, +} + +impl Lease { + const fn is_live(self, now: UnixSeconds) -> bool { + self.expires_at.get() > now.get() + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) struct Offered { + at: UnixSeconds, + verified_at: Option, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum Consent { + Unconditional, + Conditional(Offered), +} + +impl Consent { + pub(crate) const fn of( + offer: Offer, + created_at: UnixSeconds, + verified_at: Option, + ) -> Self { + match offer { + Offer::Granted => Self::Unconditional, + Offer::Invited => Self::Conditional(Offered { + at: created_at, + verified_at, + }), + } + } + + pub(crate) const fn is_conditional(self) -> bool { + matches!(self, Self::Conditional(_)) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Acceptance<'a> { + OfMembership, + OfCollaboration(&'a RepoDid), +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Reading { + Present, + Absent, + Unreachable, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct Read { + reading: Reading, + at: UnixSeconds, +} + +impl Read { + pub const fn taken(reading: Reading, at: UnixSeconds) -> Self { + Self { reading, at } + } +} + +pub trait Acceptances: Sync { + fn read<'a>( + &'a self, + scope: Acceptance<'a>, + subject: &'a AccountDid, + ) -> Pin + Send + 'a>>; +} + +pub struct Answer<'a, S> { + scope: S, + subject: &'a AccountDid, + granted: bool, +} + +pub type MemberConsent<'a> = Answer<'a, ()>; +pub type CollaborationConsent<'a> = Answer<'a, &'a RepoDid>; + +impl<'a, S> Answer<'a, S> { + pub const fn subject(&self) -> &'a AccountDid { + self.subject + } + + pub const fn granted(&self) -> bool { + self.granted + } + + #[cfg(feature = "test-support")] + #[doc(hidden)] + pub const fn assumed(scope: S, subject: &'a AccountDid, granted: bool) -> Self { + Self { + scope, + subject, + granted, + } + } +} + +impl<'a> MemberConsent<'a> { + pub const fn unread(subject: &'a AccountDid) -> Self { + Self { + scope: (), + subject, + granted: false, + } + } +} + +impl<'a> CollaborationConsent<'a> { + pub const fn repo(&self) -> &'a RepoDid { + self.scope + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub(crate) enum Scope { + Membership, + Collaboration(RepoKey), +} + +type Seat = (Scope, AccountKey); + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Seen { + Present, + Absent, +} + +#[derive(Debug, Clone, Copy)] +struct Held { + offered_at: UnixSeconds, + seen: Seen, + lease: Lease, +} + +impl Held { + const fn grants(self) -> bool { + matches!(self.seen, Seen::Present) + } + + const fn supersedes(self, seated: Self) -> bool { + let (mine, theirs) = (self.lease.read_at.get(), seated.lease.read_at.get()); + mine > theirs || (mine == theirs && !self.grants()) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Consented { + Granted, + Refused, + Stale, +} + +pub(crate) struct ConsentProjection { + readings: scc::HashMap, + turns: scc::HashMap>>, +} + +impl ConsentProjection { + pub(crate) fn new() -> Self { + Self { + readings: scc::HashMap::new(), + turns: scc::HashMap::new(), + } + } + + pub(crate) fn forget_repo(&self, repo: RepoKey) { + self.drop_where(|(scope, _)| *scope == Scope::Collaboration(repo)); + } + + pub(crate) fn retain_scope(&self, scope: Scope, live: impl Fn(AccountKey) -> bool) { + self.drop_where(|(seated, account)| *seated == scope && !live(*account)); + } + + fn drop_where(&self, gone: impl Fn(&Seat) -> bool) { + self.readings.retain_sync(|seat, _| !gone(seat)); + self.turns.retain_sync(|seat, _| !gone(seat)); + } + + fn turn(&self, seat: Seat) -> Arc> { + self.turns + .entry_sync(seat) + .or_insert_with(|| Arc::new(Mutex::new(()))) + .get() + .clone() + } + + fn read(&self, seat: Seat, offered_at: UnixSeconds) -> Option { + self.readings + .read_sync(&seat, |_, held| *held) + .filter(|held| held.offered_at == offered_at) + } + + fn write(&self, seat: Seat, held: Held) -> bool { + let mut occupied = self.readings.entry_sync(seat).or_insert(held); + let seated = occupied.get_mut(); + if seated.offered_at != held.offered_at || held.supersedes(*seated) { + *seated = held; + } + seated.grants() + } +} + +pub struct Consents<'a>(pub(crate) &'a Index); + +impl Consents<'_> { + pub async fn of_membership<'s, R: Acceptances>( + &self, + pds: &R, + subject: &'s AccountDid, + now: UnixSeconds, + ) -> MemberConsent<'s> { + Answer { + scope: (), + subject, + granted: self + .resolve(pds, Acceptance::OfMembership, subject, now) + .await, + } + } + + pub async fn of_collaboration<'s, R: Acceptances>( + &self, + pds: &R, + repo: &'s RepoDid, + subject: &'s AccountDid, + now: UnixSeconds, + ) -> CollaborationConsent<'s> { + Answer { + scope: repo, + subject, + granted: self + .resolve(pds, Acceptance::OfCollaboration(repo), subject, now) + .await, + } + } + + async fn resolve( + &self, + pds: &R, + scope: Acceptance<'_>, + subject: &AccountDid, + now: UnixSeconds, + ) -> bool { + let Resolved::Ready(Some(consent)) = self.offered(scope, subject) else { + return false; + }; + let decided = |consents: &Self| match consents.check(scope, subject, consent, now) { + Consented::Granted => Some(true), + Consented::Refused => Some(false), + Consented::Stale => None, + }; + if let Some(answer) = decided(self) { + return answer; + } + let turn = self.0.consents.turn(self.seat(scope, subject)); + let _first = turn.lock().await; + match decided(self) { + Some(answer) => answer, + None => { + let read = pds.read(scope, subject).await; + self.settle(scope, subject, consent, read) + } + } + } + + fn offered(&self, scope: Acceptance<'_>, subject: &AccountDid) -> Resolved> { + let (index, read) = (self.0, Provenance::consent); + match scope { + Acceptance::OfMembership => index.members.slot(&index.interner, subject, read), + Acceptance::OfCollaboration(repo) => { + index.collaborators.slot(&index.interner, repo, subject, read) + } + } + } + + fn check( + &self, + scope: Acceptance<'_>, + subject: &AccountDid, + consent: Consent, + now: UnixSeconds, + ) -> Consented { + let Consent::Conditional(offered) = consent else { + return Consented::Granted; + }; + match self.recall(scope, subject, offered) { + Some(held) if held.lease.is_live(now) => match held.grants() { + true => Consented::Granted, + false => Consented::Refused, + }, + _ => Consented::Stale, + } + } + + fn settle( + &self, + scope: Acceptance<'_>, + subject: &AccountDid, + consent: Consent, + read: Read, + ) -> bool { + let Consent::Conditional(offered) = consent else { + return true; + }; + let freshness = self.0.freshness; + let seen_at = |seen, read_at| Held { + offered_at: offered.at, + seen, + lease: freshness.leased(read_at, read.at), + }; + let held = match (read.reading, self.recall(scope, subject, offered)) { + (Reading::Present, _) => seen_at(Seen::Present, read.at), + (Reading::Absent, _) => seen_at(Seen::Absent, read.at), + (Reading::Unreachable, Some(held)) => match held.seen { + Seen::Present if freshness.graces(held.lease.read_at, read.at) => { + seen_at(Seen::Present, held.lease.read_at) + } + Seen::Present | Seen::Absent => seen_at(Seen::Absent, held.lease.read_at), + }, + (Reading::Unreachable, None) => seen_at(Seen::Absent, read.at), + }; + self.0.consents.write(self.seat(scope, subject), held) + } + + fn recall( + &self, + scope: Acceptance<'_>, + subject: &AccountDid, + offered: Offered, + ) -> Option { + let cached = self.0.consents.read(self.seat(scope, subject), offered.at); + let vouched = offered.verified_at.map(|verified_at| Held { + offered_at: offered.at, + seen: Seen::Present, + lease: self.0.freshness.leased(verified_at, verified_at), + }); + match (vouched, cached) { + (Some(vouched), Some(cached)) if vouched.supersedes(cached) => Some(vouched), + (Some(vouched), None) => Some(vouched), + (_, cached) => cached, + } + } + + fn seat(&self, scope: Acceptance<'_>, subject: &AccountDid) -> Seat { + let scope = match scope { + Acceptance::OfMembership => Scope::Membership, + Acceptance::OfCollaboration(repo) => { + Scope::Collaboration(self.0.interner.intern_repo(repo)) + } + }; + (scope, self.0.interner.intern_account(subject)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + use knot_git::Layout; + use tempfile::TempDir; + + const GRACE: i64 = 21_600; + const MEMBERSHIP: Acceptance<'static> = Acceptance::OfMembership; + + fn at(seconds: i64) -> UnixSeconds { + UnixSeconds::new(seconds) + } + + fn nel() -> AccountDid { + AccountDid::new("did:plc:nel").unwrap() + } + + fn knot() -> (TempDir, Index) { + let dir = tempfile::tempdir().unwrap(); + let index = Index::new( + dir.path().join("meta"), + Layout::new(dir.path().join("repos")), + ); + (dir, index) + } + + fn accepted(offered_at: i64, verified_at: i64) -> Consent { + Consent::of(Offer::Invited, at(offered_at), Some(at(verified_at))) + } + + fn invited(offered_at: i64) -> Consent { + Consent::of(Offer::Invited, at(offered_at), None) + } + + fn checks(index: &Index, who: &AccountDid, consent: Consent, now: i64) -> Consented { + index.consents().check(MEMBERSHIP, who, consent, at(now)) + } + + fn settles(index: &Index, consent: Consent, reading: Reading, when: i64) -> bool { + index + .consents() + .settle(MEMBERSHIP, &nel(), consent, Read::taken(reading, at(when))) + } + + #[test] + fn grandfathered_entry_skips_lease_and_grants_past_grace() { + let (_dir, index) = knot(); + let granted = Consent::of(Offer::Granted, at(1), None); + assert_eq!( + checks(&index, &nel(), granted, GRACE * 4), + Consented::Granted, + "Consent::Unconditional skips the lease, so GRACE * 4 later it still grants" + ); + assert!( + settles(&index, granted, Reading::Absent, GRACE * 4), + "settle seats no reading for a grant, so an Absent read leaves it alone" + ); + } + + #[test] + fn a_verified_acceptance_serves_its_ttl_and_then_follows_the_record() { + let (_dir, index) = knot(); + let consent = accepted(10, 100); + assert_eq!(checks(&index, &nel(), consent, 120), Consented::Granted); + assert_eq!(checks(&index, &nel(), consent, 200), Consented::Stale); + + assert!(!settles(&index, consent, Reading::Absent, 200)); + assert_eq!( + checks(&index, &nel(), consent, 250), + Consented::Refused, + "an Absent read is leased too, so no push inside the ttl re-reads it" + ); + assert_eq!(checks(&index, &nel(), consent, 300), Consented::Stale); + assert!(settles(&index, consent, Reading::Present, 300)); + } + + #[test] + fn outage_serves_last_reading_it_took_and_a_refusal_stands() { + let (_dir, index) = knot(); + let consent = accepted(10, 100); + assert!(settles(&index, consent, Reading::Unreachable, 200)); + assert_eq!( + checks(&index, &nel(), consent, 250), + Consented::Granted, + "the grace renews a Present read, so a push during the outage still passes" + ); + + assert!( + !settles(&index, consent, Reading::Unreachable, 100 + GRACE), + "the grace runs from read_at 100, so the attempt at 200 can't extend it" + ); + assert_eq!( + checks(&index, &nel(), consent, 100 + GRACE + 1), + Consented::Refused + ); + + settles(&index, consent, Reading::Absent, GRACE + 200); + assert!( + !settles(&index, consent, Reading::Unreachable, GRACE + 300), + "an outage can't undo a deletion the knot already read" + ); + } + + #[test] + fn an_outage_refuses_an_unread_invitation_and_a_later_acceptance_still_grants() { + let (_dir, index) = knot(); + let consent = invited(10); + assert_eq!( + checks(&index, &nel(), consent, 11), + Consented::Stale, + "no reading means Stale, however fresh the invitation is" + ); + assert!( + !settles(&index, consent, Reading::Unreachable, 11), + "no prior reading, so Unreachable seats Absent at 11 and the invite is refused" + ); + assert!( + settles(&index, consent, Reading::Present, 70), + "a Present read replaces the refused reading, so the invite grants at 70" + ); + } + + #[test] + fn a_second_invitation_discards_what_the_knot_read_about_the_first() { + let (_dir, index) = knot(); + settles(&index, invited(10), Reading::Absent, 20); + assert_eq!( + checks(&index, &nel(), invited(30), 40), + Consented::Stale, + "the cached reading is keyed by offered_at, so the invite at 30 starts Stale" + ); + } + + #[test] + fn a_later_read_supersedes_an_earlier_one_and_a_tie_refuses() { + let (_dir, index) = knot(); + settles(&index, invited(10), Reading::Absent, 20); + assert_eq!( + checks(&index, &nel(), accepted(10, 25), 26), + Consented::Granted, + "acceptMembership signed at 25, newer than the 404 at 20, so 26 grants" + ); + + let consent = invited(10); + settles(&index, consent, Reading::Present, 300); + settles(&index, consent, Reading::Absent, 299); + assert_eq!( + checks(&index, &nel(), consent, 301), + Consented::Granted, + "read_at 299 loses to read_at 300, whatever order the requests finished in" + ); + + settles(&index, consent, Reading::Absent, 300); + assert!( + !settles(&index, consent, Reading::Present, 300), + "a deletion and a presence read in the same second aren't ordered by the clock, \ + so the one that takes access away is the reading the map keeps" + ); + assert_eq!(checks(&index, &nel(), consent, 301), Consented::Refused); + } + + #[test] + fn each_repository_leases_its_own_readings_and_forget_repo_drops_them() { + let (_dir, index) = knot(); + let squid = RepoDid::new("did:plc:squid").unwrap(); + let scope = Acceptance::OfCollaboration(&squid); + let consent = accepted(10, 100); + settles(&index, consent, Reading::Absent, 200); + assert_eq!( + index.consents().check(scope, &nel(), consent, at(200)), + Consented::Stale, + "the membership reading at 200 is in another Seat, so collaboration is still Stale" + ); + + let consent = invited(10); + index.consents().settle( + scope, + &nel(), + consent, + Read::taken(Reading::Present, at(20)), + ); + assert_eq!( + index.consents().check(scope, &nel(), consent, at(30)), + Consented::Granted + ); + + index + .consents + .forget_repo(index.interner.intern_repo(&squid)); + assert_eq!( + index.consents().check(scope, &nel(), consent, at(30)), + Consented::Stale, + "forget_repo dropped the seat, so this check must resolve from scratch" + ); + } + + #[test] + fn retain_scope_forgets_a_dropped_subject_and_leaves_the_kept_one_alone() { + let (_dir, index) = knot(); + let consent = invited(10); + let leaving = AccountDid::new("did:plc:teq").unwrap(); + [&nel(), &leaving].iter().for_each(|subject| { + index.consents().settle( + MEMBERSHIP, + subject, + consent, + Read::taken(Reading::Present, at(20)), + ); + }); + + let kept = index.interner.intern_account(&nel()); + index + .consents + .retain_scope(Scope::Membership, |account| account == kept); + assert_eq!( + checks(&index, &nel(), consent, 30), + Consented::Granted, + "retain_scope kept nel, so the reading from 20 is still leased at 30" + ); + assert_eq!( + checks(&index, &leaving, consent, 30), + Consented::Stale, + "the readings map can't outgrow the roster, so a dropped subject resolves again" + ); + } +} diff --git a/knot2/crates/knot-index/src/coverage.rs b/knot2/crates/knot-index/src/coverage.rs index b63e8138c..d9049f260 100644 --- a/knot2/crates/knot-index/src/coverage.rs +++ b/knot2/crates/knot-index/src/coverage.rs @@ -23,6 +23,16 @@ impl Resolved { Resolved::Warming => Resolved::Warming, } } + + pub fn ready_or_default(self) -> T + where + T: Default, + { + match self { + Resolved::Ready(value) => value, + Resolved::Warming => T::default(), + } + } } const WARMING: u8 = 0; diff --git a/knot2/crates/knot-index/src/lib.rs b/knot2/crates/knot-index/src/lib.rs index 297a82100..c7c9af12b 100644 --- a/knot2/crates/knot-index/src/lib.rs +++ b/knot2/crates/knot-index/src/lib.rs @@ -1,8 +1,13 @@ +mod consent; mod coverage; mod error; mod intern; mod projections; +pub use consent::{ + Acceptance, AcceptanceGrace, AcceptanceTtl, Acceptances, Answer, CollaborationConsent, + Consents, Freshness, MemberConsent, Read, Reading, +}; pub use coverage::{Coverage, Resolved}; pub use error::IndexError; pub use knot_types::OfferedKey; @@ -13,15 +18,19 @@ use std::time::Duration; use knot_cob::{ChangePayload, CobStore}; use knot_cobs::{ - BlocklistChange, BlocklistCob, CollaboratorsChange, CollaboratorsCob, Grant, MembersChange, + BlocklistChange, BlocklistCob, CollaboratorsChange, CollaboratorsCob, MembersChange, MembersCob, RegistryChange, RepoRef, RepoRegistryCob, }; +pub use knot_cobs::{EffectiveSince, Entry, Offer, Standing}; use knot_git::{Layout, Repo}; use knot_types::{AccountDid, ClonePath, OwnerDid, RepoDid, RepoRkey, UnixSeconds}; use tokio::sync::watch; +use consent::ConsentProjection; use intern::{Interner, RepoKey}; -use projections::{CollaboratorsProjection, GrantSetProjection, KeyProjection, RegistryProjection}; +use projections::{ + CollaboratorsProjection, GrantSetProjection, KeyProjection, Provenance, RegistryProjection, +}; knot_types::scalar_newtype! { pub struct IndexGeneration(u64); @@ -91,19 +100,25 @@ impl HostedCoverage { struct Granted { subjects: Vec, + invited: Vec, hosted: HostedCoverage, } +struct Seats { + grants: Vec, + invited: Vec, +} + enum Folded { - Grants(Vec), + Seated(Seats), Pending, Unreadable, } impl Folded { - fn into_grants(self) -> Option> { + fn into_seats(self) -> Option { match self { - Folded::Grants(subjects) => Some(subjects), + Folded::Seated(seats) => Some(seats), Folded::Pending | Folded::Unreadable => None, } } @@ -233,7 +248,7 @@ pub(crate) const fn whole_secs(span: Duration) -> i64 { } } -const fn after(now: UnixSeconds, span: Duration) -> UnixSeconds { +pub(crate) const fn after(now: UnixSeconds, span: Duration) -> UnixSeconds { now.saturating_add_secs(whole_secs(span)) } @@ -308,6 +323,8 @@ pub struct Index { collaborators: CollaboratorsProjection, registry: RegistryProjection, keys: KeyProjection, + consents: ConsentProjection, + freshness: Freshness, unreadable: scc::HashSet, generation: AtomicU64, generations: watch::Sender, @@ -332,12 +349,19 @@ impl Index { collaborators: CollaboratorsProjection::new(), registry: RegistryProjection::new(), keys: KeyProjection::new(budget), + consents: ConsentProjection::new(), + freshness: Freshness::DEFAULT, unreadable: scc::HashSet::new(), generation: AtomicU64::new(0), generations: watch::Sender::new(IndexGeneration::new(0)), } } + pub fn with_freshness(mut self, freshness: Freshness) -> Self { + self.freshness = freshness; + self + } + pub fn generation(&self) -> IndexGeneration { IndexGeneration(self.generation.load(Ordering::Acquire)) } @@ -394,6 +418,10 @@ impl Index { }); } } + self.consents + .retain_scope(consent::Scope::Membership, |key| { + self.members.conditional(key) + }); self.bump_generation(); Ok(()) } @@ -431,6 +459,7 @@ impl Index { evacuated.iter().for_each(|repo| { if let Some(key) = self.interner.repo(repo) { self.collaborators.drop_repo(key); + self.consents.forget_repo(key); self.unreadable.remove_sync(&key); } }); @@ -472,37 +501,93 @@ impl Index { }); } } + self.consents + .retain_scope(consent::Scope::Collaboration(repo_key), |key| { + self.collaborators.conditional(repo_key, key) + }); self.bump_generation(); Ok(()) } - pub fn is_member(&self, did: &AccountDid) -> Resolved { - self.members.contains(&self.interner, did) + pub fn effective_member(&self, did: &AccountDid) -> Resolved { + self.member_standing(did) + .map(|standing| standing.is_some_and(Standing::is_effective)) + } + + pub fn member_standing(&self, did: &AccountDid) -> Resolved> { + self.members.slot(&self.interner, did, Provenance::standing) + } + + pub fn member_entries(&self) -> Resolved> { + self.member_entries_where(|_| true) } - pub fn member_entries(&self) -> Resolved> { - self.members.entries(&self.interner) + pub fn member_entries_where( + &self, + keep: impl Fn(Standing) -> bool, + ) -> Resolved> { + self.members.entries_where(&self.interner, keep) + } + + pub fn effective_members(&self) -> Resolved> { + self.members + .subjects_where(&self.interner, Standing::is_effective) } pub fn is_blocked(&self, did: &AccountDid) -> Resolved { - self.blocklist.contains(&self.interner, did) + self.blocklist + .slot(&self.interner, did, Provenance::standing) + .map(|standing| standing.is_some()) + } + + pub fn blocked(&self) -> Resolved> { + self.blocklist.subjects_where(&self.interner, |_| true) + } + + pub fn effective_collaborator(&self, repo: &RepoDid, did: &AccountDid) -> Resolved { + self.collaborator_standing(repo, did) + .map(|standing| standing.is_some_and(Standing::is_effective)) + } + + pub fn collaborator_standing( + &self, + repo: &RepoDid, + did: &AccountDid, + ) -> Resolved> { + self.collaborators + .slot(&self.interner, repo, did, Provenance::standing) } - pub fn blocked_entries(&self) -> Resolved> { - self.blocklist.entries(&self.interner) + pub fn collaborator_entries(&self, repo: &RepoDid) -> Resolved> { + self.collaborator_entries_where(repo, |_| true) } - pub fn is_collaborator(&self, repo: &RepoDid, did: &AccountDid) -> Resolved { - self.collaborators.contains(&self.interner, repo, did) + pub fn collaborator_entries_where( + &self, + repo: &RepoDid, + keep: impl Fn(Standing) -> bool, + ) -> Resolved> { + self.collaborators.entries_where(&self.interner, repo, keep) + } + + pub fn effective_collaborators(&self, repo: &RepoDid) -> Resolved> { + self.collaborators + .subjects_where(&self.interner, repo, Standing::is_effective) + } + + pub fn invited_collaborators(&self, repo: &RepoDid) -> Resolved> { + self.collaborators + .subjects_where(&self.interner, repo, |standing| { + matches!(standing, Standing::Invited) + }) } - pub fn collaborator_entries(&self, repo: &RepoDid) -> Resolved> { - self.collaborators.entries(&self.interner, repo) + pub fn any_invited_collaborator(&self) -> bool { + self.collaborators.any_invited() } - pub fn collaborators_of(&self, repo: &RepoDid) -> Resolved> { - self.collaborator_entries(repo) - .map(|entries| entries.into_iter().map(|grant| grant.subject).collect()) + pub fn consents(&self) -> Consents<'_> { + Consents(self) } pub fn resolve_repo(&self, owner: &OwnerDid, rkey: &RepoRkey) -> Resolved> { @@ -549,21 +634,7 @@ impl Index { let folded: Vec = self .hosted_repos() .iter() - .map( - |repo| match (self.owner_of(repo), self.collaborators_of(repo)) { - (Resolved::Ready(owner), Resolved::Ready(collaborators)) => { - Folded::Grants( - owner - .map(AccountDid::from) - .into_iter() - .chain(collaborators) - .collect(), - ) - } - _ if self.unreadable_repo(repo) => Folded::Unreadable, - _ => Folded::Pending, - }, - ) + .map(|repo| self.seats_of(repo)) .collect(); let hosted = HostedCoverage::over( folded @@ -571,23 +642,38 @@ impl Index { .filter(|repo| matches!(repo, Folded::Pending)) .count(), ); + let (grants, invited) = folded.into_iter().filter_map(Folded::into_seats).fold( + (Vec::new(), Vec::new()), + |(mut grants, mut invited), seats| { + grants.extend(seats.grants); + invited.extend(seats.invited); + (grants, invited) + }, + ); Resolved::Ready(Granted { - subjects: sorted( - folded - .into_iter() - .filter_map(Folded::into_grants) - .flatten() - .collect(), - ), + subjects: sorted(grants), + invited: sorted(invited), hosted, }) } } } - fn member_grants(&self) -> Resolved> { - self.member_entries() - .map(|entries| sorted(entries.into_iter().map(|grant| grant.subject).collect())) + fn seats_of(&self, repo: &RepoDid) -> Folded { + match ( + self.owner_of(repo), + self.effective_collaborators(repo), + self.invited_collaborators(repo), + ) { + (Resolved::Ready(owner), Resolved::Ready(granted), Resolved::Ready(invited)) => { + Folded::Seated(Seats { + grants: owner.map(AccountDid::from).into_iter().chain(granted).collect(), + invited, + }) + } + _ if self.unreadable_repo(repo) => Folded::Unreadable, + _ => Folded::Pending, + } } fn due_among( @@ -686,19 +772,33 @@ impl KeySet<'_> { true => Recheck::Everything, false => Recheck::Renewals, }; - let members = index.member_grants().map(|members| { + let vouched: Vec = granted + .invited + .into_iter() + .filter(|did| { + pushers.binary_search(did).is_err() && index.keys.on_file(&index.interner, did) + }) + .collect(); + let members = index.effective_members().map(|members| { let outside: Vec = members .into_iter() - .filter(|did| pushers.binary_search(did).is_err()) + .filter(|did| { + pushers.binary_search(did).is_err() && vouched.binary_search(did).is_err() + }) .collect(); - let (unread, due) = index + let (unread, due): (Vec, Vec) = index .due_among(&outside, now, against) .into_iter() .partition(|did| !index.keys.on_file(&index.interner, did)); + let due = due + .into_iter() + .chain(index.due_among(&vouched, now, against)) + .collect(); + let kept = pushers.iter().chain(vouched.iter()).cloned().chain(outside).collect(); MemberWork { unread: UnreadMembers(unread), due: StaleMembers(due), - kept: KeptAccounts(pushers.iter().cloned().chain(outside).collect()), + kept: KeptAccounts(kept), } }); let tracked = match &members { diff --git a/knot2/crates/knot-index/src/projections.rs b/knot2/crates/knot-index/src/projections.rs index 0d4746ecd..488fb1d1c 100644 --- a/knot2/crates/knot-index/src/projections.rs +++ b/knot2/crates/knot-index/src/projections.rs @@ -6,11 +6,12 @@ use std::sync::{Arc, Mutex}; use knot_cob::{Change, ChangeId, ChangePayload, Checkpoint, CobId, CobStore, Evaluate}; use knot_cobs::{ - CollaboratorsChange, CollaboratorsCob, Grant, GrantChange, Registration, Registry, - RegistryChange, Rename, RepoRef, RepoRegistryCob, Roster, + CollaboratorsChange, CollaboratorsCob, Effect, Entry, Grant, GrantChange, Offer, Registration, + Registry, RegistryChange, Rename, RepoRef, RepoRegistryCob, Roster, Standing, }; use knot_types::{AccountDid, ClonePath, OfferedKey, OwnerDid, RepoDid, RepoRkey, UnixSeconds}; +use crate::consent::Consent; use crate::coverage::{Coverage, CoverageCell, Resolved}; use crate::error::IndexError; use crate::intern::{AccountKey, Interner, NameKey, OwnerKey, RepoKey, RkeyKey}; @@ -19,24 +20,51 @@ use crate::{ }; #[derive(Debug, Clone, Copy)] -struct Provenance { +pub(crate) struct Provenance { + offer: Offer, added_by: AccountKey, created_at: UnixSeconds, + verified_at: Option, } impl Provenance { - fn intern(interner: &Interner, grant: &Grant) -> Self { + fn admitted(interner: &Interner, offer: Offer, grant: &Grant) -> Self { Self { + offer, added_by: interner.intern_account(&grant.added_by), created_at: grant.created_at, + verified_at: None, } } - fn grant(self, interner: &Interner, subject: AccountKey) -> Grant { - Grant { - subject: interner.resolve_account(subject), + fn seeded(interner: &Interner, entry: &Entry) -> Self { + Self { + offer: entry.offer, + added_by: interner.intern_account(&entry.added_by), + created_at: entry.created_at, + verified_at: entry.verified_at, + } + } + + fn verified(mut self, verified_at: UnixSeconds) -> Self { + self.verified_at.get_or_insert(verified_at); + self + } + + pub(crate) fn standing(self) -> Standing { + Standing::of(self.offer, self.created_at, self.verified_at) + } + + pub(crate) fn consent(self) -> Consent { + Consent::of(self.offer, self.created_at, self.verified_at) + } + + fn entry(self, interner: &Interner) -> Entry { + Entry { + offer: self.offer, added_by: interner.resolve_account(self.added_by), created_at: self.created_at, + verified_at: self.verified_at, } } } @@ -106,25 +134,43 @@ where self.coverage.set(Coverage::Ready); } - pub(crate) fn contains(&self, interner: &Interner, did: &AccountDid) -> Resolved { + pub(crate) fn slot( + &self, + interner: &Interner, + did: &AccountDid, + read: impl Fn(Provenance) -> T, + ) -> Resolved> { match self.coverage.get() { Coverage::Warming => Resolved::Warming, Coverage::Ready => Resolved::Ready( interner .account(did) - .is_some_and(|did| self.membership.contains_sync(&did)), + .and_then(|key| self.membership.read_sync(&key, |_, slot| read(*slot))), ), } } - pub(crate) fn entries(&self, interner: &Interner) -> Resolved> { + pub(crate) fn conditional(&self, subject: AccountKey) -> bool { + self.membership + .read_sync(&subject, |_, slot| slot.consent().is_conditional()) + .unwrap_or(false) + } + + fn rows_where( + &self, + interner: &Interner, + keep: impl Fn(Standing) -> bool, + row: impl Fn(AccountDid, Provenance) -> T, + ) -> Resolved> { match self.coverage.get() { Coverage::Warming => Resolved::Warming, Coverage::Ready => { let mut out = BTreeMap::new(); self.membership.iter_sync(|&subject, slot| { - let grant = slot.grant(interner, subject); - out.insert(grant.subject.clone(), grant); + if keep(slot.standing()) { + let did = interner.resolve_account(subject); + out.insert(did.clone(), row(did, *slot)); + } true }); Resolved::Ready(out.into_values().collect()) @@ -132,6 +178,22 @@ where } } + pub(crate) fn entries_where( + &self, + interner: &Interner, + keep: impl Fn(Standing) -> bool, + ) -> Resolved> { + self.rows_where(interner, keep, |did, slot| (did, slot.entry(interner))) + } + + pub(crate) fn subjects_where( + &self, + interner: &Interner, + keep: impl Fn(Standing) -> bool, + ) -> Resolved> { + self.rows_where(interner, keep, |did, _| did) + } + pub(crate) fn refresh( &self, interner: &Interner, @@ -170,13 +232,10 @@ where fn seed(&self, interner: &Interner, roster: &Roster) { self.membership.clear_sync(); roster.entries().for_each(|(subject, entry)| { - let slot = Provenance { - added_by: interner.intern_account(&entry.added_by), - created_at: entry.created_at, - }; - let _ = self - .membership - .insert_sync(interner.intern_account(subject), slot); + let _ = self.membership.insert_sync( + interner.intern_account(subject), + Provenance::seeded(interner, entry), + ); }); } @@ -188,10 +247,15 @@ where .account(&did) .and_then(|key| self.membership.read_sync(&key, |_, slot| *slot)); let net = ops - .into_iter() - .fold(current, |slot, change| match change.as_grant() { - Some(grant) => slot.or_else(|| Some(Provenance::intern(interner, grant))), - None => None, + .iter() + .fold(current, |slot, change| match change.effect() { + Effect::Admit(offer, grant) => { + slot.or_else(|| Some(Provenance::admitted(interner, offer, grant))) + } + Effect::Accept(accept) => { + slot.map(|slot| slot.verified(accept.verified_at)) + } + Effect::Revoke(_) => None, }); match net { Some(slot) => { @@ -231,6 +295,15 @@ impl CollaboratorsProjection { } } + pub(crate) fn any_invited(&self) -> bool { + !self.rosters.iter_sync(|_, roster| { + roster + .entries + .values() + .all(|slot| !matches!(slot.standing(), Standing::Invited)) + }) + } + fn repo_lock(&self, repo: RepoKey) -> &Mutex<()> { let mut hasher = std::collections::hash_map::DefaultHasher::new(); repo.hash(&mut hasher); @@ -247,38 +320,86 @@ impl CollaboratorsProjection { .is_some_and(|repo| self.rosters.contains_sync(&repo)) } - pub(crate) fn contains( + pub(crate) fn slot( &self, interner: &Interner, repo: &RepoDid, did: &AccountDid, - ) -> Resolved { + read: impl Fn(Provenance) -> T, + ) -> Resolved> { let Some(repo) = interner.repo(repo) else { return Resolved::Warming; }; match self.rosters.read_sync(&repo, |_, roster| { interner .account(did) - .is_some_and(|account| roster.entries.contains_key(&account)) + .and_then(|account| roster.entries.get(&account)) + .map(|slot| read(*slot)) }) { - Some(present) => Resolved::Ready(present), + Some(value) => Resolved::Ready(value), None => Resolved::Warming, } } - pub(crate) fn entries(&self, interner: &Interner, repo: &RepoDid) -> Resolved> { + pub(crate) fn conditional(&self, repo: RepoKey, subject: AccountKey) -> bool { + self.rosters + .read_sync(&repo, |_, roster| { + roster + .entries + .get(&subject) + .is_some_and(|slot| slot.consent().is_conditional()) + }) + .unwrap_or(false) + } + + fn rows_where( + &self, + interner: &Interner, + repo: &RepoDid, + keep: impl Fn(Standing) -> bool, + row: impl Fn(AccountDid, Provenance) -> T, + ) -> Resolved> { let Some(repo) = interner.repo(repo) else { return Resolved::Warming; }; - match self - .rosters - .read_sync(&repo, |_, roster| roster_grants(interner, roster)) - { - Some(grants) => Resolved::Ready(grants), + match self.rosters.read_sync(&repo, |_, roster| { + roster + .entries + .iter() + .filter(|(_, slot)| keep(slot.standing())) + .map(|(&subject, slot)| { + let did = interner.resolve_account(subject); + (did.clone(), row(did, *slot)) + }) + .collect::>() + .into_values() + .collect() + }) { + Some(rows) => Resolved::Ready(rows), None => Resolved::Warming, } } + pub(crate) fn entries_where( + &self, + interner: &Interner, + repo: &RepoDid, + keep: impl Fn(Standing) -> bool, + ) -> Resolved> { + self.rows_where(interner, repo, keep, |did, slot| { + (did, slot.entry(interner)) + }) + } + + pub(crate) fn subjects_where( + &self, + interner: &Interner, + repo: &RepoDid, + keep: impl Fn(Standing) -> bool, + ) -> Resolved> { + self.rows_where(interner, repo, keep, |did, _| did) + } + pub(crate) fn mark_repo_empty(&self, repo: RepoKey) { let lock = self.repo_lock(repo); let _guard = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner()); @@ -323,10 +444,7 @@ impl CollaboratorsProjection { .map(|(subject, entry)| { ( interner.intern_account(subject), - Provenance { - added_by: interner.intern_account(&entry.added_by), - created_at: entry.created_at, - }, + Provenance::seeded(interner, entry), ) }) .collect(); @@ -379,32 +497,27 @@ impl CollaboratorsProjection { } } -fn roster_grants(interner: &Interner, roster: &RepoRoster) -> Vec { - roster - .entries - .iter() - .map(|(&subject, slot)| { - let grant = slot.grant(interner, subject); - (grant.subject.clone(), grant) - }) - .collect::>() - .into_values() - .collect() -} - fn apply( interner: &Interner, mut entries: BTreeMap, change: CollaboratorsChange, ) -> BTreeMap { - match change { - CollaboratorsChange::Add(grant) => { + match change.effect() { + Effect::Admit(offer, grant) => { entries .entry(interner.intern_account(&grant.subject)) - .or_insert_with(|| Provenance::intern(interner, &grant)); + .or_insert_with(|| Provenance::admitted(interner, offer, grant)); + } + Effect::Accept(accept) => { + if let Some(slot) = interner + .account(&accept.subject) + .and_then(|key| entries.get_mut(&key)) + { + *slot = slot.verified(accept.verified_at); + } } - CollaboratorsChange::Remove(removal) => { - if let Some(key) = interner.account(&removal.subject) { + Effect::Revoke(subject) => { + if let Some(key) = interner.account(subject) { entries.remove(&key); } } diff --git a/knot2/crates/knot-index/tests/common/mod.rs b/knot2/crates/knot-index/tests/common/mod.rs index 05195059b..bebf25d98 100644 --- a/knot2/crates/knot-index/tests/common/mod.rs +++ b/knot2/crates/knot-index/tests/common/mod.rs @@ -2,7 +2,7 @@ use std::path::PathBuf; -use knot_cob::{CobHome, CobId, CobStore}; +use knot_cob::{ChangePayload, CobHome, CobId, CobStore}; use knot_cobs::{CollaboratorsChange, Grant, MembersChange, Registration, RegistryChange, Removal}; use knot_git::{Layout, Repo}; use knot_index::{Index, KeyBudget}; @@ -102,6 +102,20 @@ impl World { Index::with_key_budget(&self.meta_path, self.layout.clone(), budget) } + pub fn on_meta(&self, object: CobId, change: &P, t: i64) { + let meta = Repo::open(&self.meta_path).unwrap(); + CobStore::new(&meta) + .update(&meta_home(), object, change, &self.signer, at(t)) + .unwrap(); + } + + pub fn on_repo(&self, repo: &RepoDid, object: CobId, change: &P, t: i64) { + let git = self.layout.open(repo).unwrap(); + CobStore::new(&git) + .update(&CobHome::from(repo), object, change, &self.signer, at(t)) + .unwrap(); + } + pub fn seed_members(&self) -> CobId { let meta = Repo::open(&self.meta_path).unwrap(); let store = CobStore::new(&meta); @@ -185,6 +199,20 @@ impl World { .object } + pub fn invite_collaborator(&self, repo: &RepoDid, subject: &str) -> CobId { + let git = self.layout.create(repo).unwrap(); + let store = CobStore::new(&git); + store + .create( + &CobHome::from(repo), + &CollaboratorsChange::Invite(knot_cobs::Invite(grant(subject, "nel", 1))), + &self.signer, + at(1), + ) + .unwrap() + .object + } + pub fn remove_collaborator(&self, repo: &RepoDid, object: CobId, subject: &str, seconds: i64) { let git = self.layout.open(repo).unwrap(); let store = CobStore::new(&git); diff --git a/knot2/crates/knot-index/tests/consent.rs b/knot2/crates/knot-index/tests/consent.rs new file mode 100644 index 000000000..6dce02afe --- /dev/null +++ b/knot2/crates/knot-index/tests/consent.rs @@ -0,0 +1,71 @@ +mod common; + +use std::pin::Pin; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::Duration; + +use common::*; +use knot_cobs::{Invite, MembersChange}; +use knot_index::{Acceptance, Acceptances, Read, Reading}; +use knot_types::{AccountDid, UnixSeconds}; +use tokio::task::JoinSet; + +const WAVE: usize = 8; +const INFLIGHT: Duration = Duration::from_millis(100); + +struct SlowPds { + reads: AtomicUsize, +} + +impl Acceptances for SlowPds { + fn read<'a>( + &'a self, + _scope: Acceptance<'a>, + _subject: &'a AccountDid, + ) -> Pin + Send + 'a>> { + Box::pin(async move { + self.reads.fetch_add(1, Ordering::SeqCst); + tokio::time::sleep(INFLIGHT).await; + Read::taken(Reading::Present, UnixSeconds::new(20)) + }) + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 8)] +async fn a_wave_of_gated_calls_shares_one_turn_and_one_request() { + let world = World::new(); + let members = world.seed_members(); + let invite = MembersChange::Invite(Invite(grant("teq", "nel", 10))); + world.on_meta(members, &invite, 10); + let index = Arc::new(world.index()); + index.rebuild().unwrap(); + let pds = Arc::new(SlowPds { + reads: AtomicUsize::new(0), + }); + + let mut wave = JoinSet::new(); + (0..WAVE).for_each(|_| { + let (index, pds) = (Arc::clone(&index), Arc::clone(&pds)); + wave.spawn(async move { + index + .consents() + .of_membership(pds.as_ref(), &acc("teq"), at(30)) + .await + .granted() + }); + }); + wave.join_all().await.into_iter().for_each(|granted| { + assert!( + granted, + "the waiters must see the read the first caller settled" + ); + }); + assert_eq!( + pds.reads.load(Ordering::SeqCst), + 1, + "the lease bounds the request rate only where the callers that miss it wait for the one \ + that went, and {WAVE} readings of a subject who published once is the amplifier a knot \ + aims at somebody else's pds" + ); +} diff --git a/knot2/crates/knot-index/tests/lifecycle.rs b/knot2/crates/knot-index/tests/lifecycle.rs index f493eb11b..eed036013 100644 --- a/knot2/crates/knot-index/tests/lifecycle.rs +++ b/knot2/crates/knot-index/tests/lifecycle.rs @@ -2,9 +2,10 @@ use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use knot_cob::{CobHome, CobStore}; -use knot_cobs::{CollaboratorsChange, CollaboratorsCob, Grant, MembersChange}; +use knot_cobs::{Accept, CollaboratorsChange, CollaboratorsCob, Invite, MembersChange, Offer}; use knot_git::Repo; -use knot_index::{Coverage, IndexCoverage, IndexError, OfferedKey, Resolved}; +use knot_index::{Coverage, Entry, IndexCoverage, IndexError, OfferedKey, Resolved, Standing}; +use knot_types::AccountDid; mod common; use common::{World, acc, at, grant, meta_home, own, repo_did, rkey}; @@ -20,9 +21,12 @@ fn rebuild_folds_members_and_registry_and_collaborators_fold_on_access() { let index = world.index(); index.rebuild().unwrap(); - assert_eq!(index.is_member(&acc("nel")), Resolved::Ready(true)); - assert_eq!(index.is_member(&acc("olaren")), Resolved::Ready(true)); - assert_eq!(index.is_member(&acc("teq")), Resolved::Ready(false)); + assert_eq!(index.effective_member(&acc("nel")), Resolved::Ready(true)); + assert_eq!( + index.effective_member(&acc("olaren")), + Resolved::Ready(true) + ); + assert_eq!(index.effective_member(&acc("teq")), Resolved::Ready(false)); assert_eq!( index.resolve_repo(&own("nel"), &rkey("anemone")), Resolved::Ready(Some(repo.clone())) @@ -32,18 +36,18 @@ fn rebuild_folds_members_and_registry_and_collaborators_fold_on_access() { Resolved::Ready(None) ); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Warming, "rebuild doesn't fold collaborators, so roster reads warming until first access" ); index.ensure_collaborators(&repo).unwrap(); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true) ); assert_eq!( - index.is_collaborator(&repo, &acc("bailey")), + index.effective_collaborator(&repo, &acc("bailey")), Resolved::Ready(false) ); assert_eq!( @@ -65,10 +69,10 @@ fn every_accessor_fails_closed_while_warming() { world.seed_collaborator(&repo_did("squid"), "lyna"); let index = world.index(); - assert_eq!(index.is_member(&acc("nel")), Resolved::Warming); + assert_eq!(index.effective_member(&acc("nel")), Resolved::Warming); assert_eq!(index.member_entries(), Resolved::Warming); assert_eq!( - index.is_collaborator(&repo_did("squid"), &acc("lyna")), + index.effective_collaborator(&repo_did("squid"), &acc("lyna")), Resolved::Warming ); assert_eq!( @@ -94,6 +98,15 @@ fn every_accessor_fails_closed_while_warming() { ); } +fn asserted(added_by: &str, seconds: i64) -> Entry { + Entry { + offer: Offer::Granted, + added_by: acc(added_by), + created_at: at(seconds), + verified_at: None, + } +} + #[test] fn member_entries_record_provenance() { let world = World::new(); @@ -102,7 +115,10 @@ fn member_entries_record_provenance() { index.rebuild().unwrap(); assert_eq!( index.member_entries(), - Resolved::Ready(vec![grant("nel", "nel", 1), grant("olaren", "nel", 2)]) + Resolved::Ready(vec![ + (acc("nel"), asserted("nel", 1)), + (acc("olaren"), asserted("nel", 2)), + ]) ); } @@ -128,7 +144,10 @@ fn a_re_added_member_keeps_the_first_provenance() { assert_eq!( index.member_entries(), - Resolved::Ready(vec![grant("nel", "nel", 1), grant("olaren", "nel", 2)]), + Resolved::Ready(vec![ + (acc("nel"), asserted("nel", 1)), + (acc("olaren"), asserted("nel", 2)), + ]), "duplicate add never rewrites original provenance, matching canonical roster" ); } @@ -166,14 +185,10 @@ fn collaborator_entries_match_the_canonical_roster() { index.ensure_collaborators(&repo).unwrap(); let canonical = store.get::(object).unwrap(); - let expected: Vec = canonical + let expected: Vec<(AccountDid, Entry)> = canonical .state() .entries() - .map(|(subject, entry)| Grant { - subject: subject.clone(), - added_by: entry.added_by.clone(), - created_at: entry.created_at, - }) + .map(|(subject, entry)| (subject.clone(), entry.clone())) .collect(); assert_eq!( index.collaborator_entries(&repo), @@ -182,6 +197,80 @@ fn collaborator_entries_match_the_canonical_roster() { ); } +#[test] +fn invited_subject_is_on_roster_and_gains_access_when_accept_lands() { + let world = World::new(); + let members = world.seed_members(); + let repo = repo_did("squid"); + world.seed_registry(&repo); + let collaborators = world.seed_collaborator(&repo, "lyna"); + let meta = Repo::open(&world.meta_path).unwrap(); + let meta_store = CobStore::new(&meta); + let git = world.layout.open(&repo).unwrap(); + let repo_store = CobStore::new(&git); + let home = CobHome::from(&repo); + let on_knot = |change: MembersChange, seconds: i64| { + meta_store + .update(&meta_home(), members, &change, &world.signer, at(seconds)) + .unwrap(); + }; + let on_repo = |change: CollaboratorsChange, seconds: i64| { + repo_store + .update(&home, collaborators, &change, &world.signer, at(seconds)) + .unwrap(); + }; + let invite = |added_by: &str| Invite(grant("teq", added_by, 4)); + let accepted = Accept { + subject: acc("teq"), + verified_at: at(6), + }; + + on_knot(MembersChange::Invite(invite("nel")), 4); + on_repo(CollaboratorsChange::Invite(invite("olaren")), 4); + let index = world.index(); + index.rebuild().unwrap(); + index.ensure_collaborators(&repo).unwrap(); + let acl = || { + ( + index.effective_member(&acc("teq")), + index.effective_collaborator(&repo, &acc("teq")), + index.effective_members(), + index.effective_collaborators(&repo), + ) + }; + + assert_eq!( + index.member_standing(&acc("teq")), + Resolved::Ready(Some(Standing::Invited)), + "an invite is an entry on the roster with Standing::Invited" + ); + assert_eq!( + acl(), + ( + Resolved::Ready(false), + Resolved::Ready(false), + Resolved::Ready(vec![acc("nel"), acc("olaren")]), + Resolved::Ready(vec![acc("lyna")]), + ), + "can_push and the listings read these, so teq mustn't appear until the accept" + ); + + on_knot(MembersChange::Accept(accepted.clone()), 6); + on_repo(CollaboratorsChange::Accept(accepted), 6); + index.refresh_members().unwrap(); + index.refresh_collaborators(&repo).unwrap(); + assert_eq!( + acl(), + ( + Resolved::Ready(true), + Resolved::Ready(true), + Resolved::Ready(vec![acc("nel"), acc("olaren"), acc("teq")]), + Resolved::Ready(vec![acc("lyna"), acc("teq")]), + ), + "the Accept at 6 makes teq a member and a pusher on both folds" + ); +} + #[test] fn a_folded_repo_serves_while_unaccessed_repos_stay_warming() { let world = World::new(); @@ -201,13 +290,13 @@ fn a_folded_repo_serves_while_unaccessed_repos_stay_warming() { index.ensure_collaborators(&present).unwrap(); assert_eq!( - index.is_collaborator(&present, &acc("lyna")), + index.effective_collaborator(&present, &acc("lyna")), Resolved::Ready(true) ); assert_eq!( index .collaborator_entries(&present) - .map(|grants| grants.len()), + .map(|entries| entries.len()), Resolved::Ready(1), "folded repo serves its roster" ); @@ -217,7 +306,7 @@ fn a_folded_repo_serves_while_unaccessed_repos_stay_warming() { "registered repo that was never accessed stays fail-closed until folded" ); assert_eq!( - index.is_collaborator(&repo_did("conch"), &acc("lyna")), + index.effective_collaborator(&repo_did("conch"), &acc("lyna")), Resolved::Warming, "repo the index never folded cannot answer, so it fails closed" ); @@ -273,7 +362,7 @@ fn a_repo_missing_on_disk_does_not_fail_the_boot_and_isolates_its_fold() { index.ensure_collaborators(&present).unwrap(); assert_eq!( - index.is_collaborator(&present, &acc("lyna")), + index.effective_collaborator(&present, &acc("lyna")), Resolved::Ready(true) ); assert!( @@ -281,7 +370,7 @@ fn a_repo_missing_on_disk_does_not_fail_the_boot_and_isolates_its_fold() { "folding repo with no dir on disk fails for that repo alone" ); assert_eq!( - index.is_collaborator(&absent, &acc("lyna")), + index.effective_collaborator(&absent, &acc("lyna")), Resolved::Warming, "repo the index couldn't fold stays fail-closed" ); @@ -306,7 +395,7 @@ fn concurrent_refreshes_of_distinct_repos_all_land() { repos.iter().for_each(|repo| { assert_eq!( - index.is_collaborator(&repo_did(repo), &acc("lyna")), + index.effective_collaborator(&repo_did(repo), &acc("lyna")), Resolved::Ready(true) ); }); @@ -335,12 +424,12 @@ fn a_refresh_is_eventually_consistent_not_an_atomic_snapshot() { scope.spawn(move || { while !reader_done.load(Ordering::Acquire) { assert_eq!( - reader.is_member(&acc("nel")), + reader.effective_member(&acc("nel")), Resolved::Ready(true), "stable member stays visible and read never blocks on writer" ); assert!( - !reader.is_member(&acc("m0")).is_warming(), + !reader.effective_member(&acc("m0")).is_warming(), "already-ready projection serves reads mid-refresh, it never re-warms" ); } @@ -349,7 +438,7 @@ fn a_refresh_is_eventually_consistent_not_an_atomic_snapshot() { (0..32).for_each(|i| { assert_eq!( - index.is_member(&acc(&format!("m{i}"))), + index.effective_member(&acc(&format!("m{i}"))), Resolved::Ready(true), "once writer returns, whole delta has converged" ); diff --git a/knot2/crates/knot-index/tests/meta_repo.rs b/knot2/crates/knot-index/tests/meta_repo.rs index 71c09e8ef..74debd173 100644 --- a/knot2/crates/knot-index/tests/meta_repo.rs +++ b/knot2/crates/knot-index/tests/meta_repo.rs @@ -70,14 +70,17 @@ fn the_meta_repo_round_trips_every_projection_from_git() { }; let first = boot(); - assert_eq!(first.is_member(&acc("nel")), Resolved::Ready(true)); - assert_eq!(first.is_member(&acc("olaren")), Resolved::Ready(true)); + assert_eq!(first.effective_member(&acc("nel")), Resolved::Ready(true)); + assert_eq!( + first.effective_member(&acc("olaren")), + Resolved::Ready(true) + ); assert_eq!( first.resolve_repo(&own("nel"), &rkey("anemone")), Resolved::Ready(Some(repo.clone())) ); assert_eq!( - first.is_collaborator(&repo, &acc("lyna")), + first.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true) ); @@ -85,10 +88,13 @@ fn the_meta_repo_round_trips_every_projection_from_git() { ["nel", "olaren", "lyna", "stranger"] .into_iter() .for_each(|who| { - assert_eq!(first.is_member(&acc(who)), second.is_member(&acc(who))); assert_eq!( - first.is_collaborator(&repo, &acc(who)), - second.is_collaborator(&repo, &acc(who)), + first.effective_member(&acc(who)), + second.effective_member(&acc(who)) + ); + assert_eq!( + first.effective_collaborator(&repo, &acc(who)), + second.effective_collaborator(&repo, &acc(who)), ); }); assert_eq!( diff --git a/knot2/crates/knot-index/tests/projections.rs b/knot2/crates/knot-index/tests/projections.rs index f0b583f4d..33e022f82 100644 --- a/knot2/crates/knot-index/tests/projections.rs +++ b/knot2/crates/knot-index/tests/projections.rs @@ -31,7 +31,7 @@ fn a_concurrent_reader_never_sees_a_net_absent_subject() { let object = world.seed_members(); let index = Arc::new(world.index()); index.rebuild().unwrap(); - assert_eq!(index.is_member(&acc("teq")), Resolved::Ready(false)); + assert_eq!(index.effective_member(&acc("teq")), Resolved::Ready(false)); let meta = Repo::open(&world.meta_path).unwrap(); let store = CobStore::new(&meta); @@ -75,7 +75,7 @@ fn a_concurrent_reader_never_sees_a_net_absent_subject() { let reader_leaked = Arc::clone(&leaked); scope.spawn(move || { while !reader_done.load(Ordering::Acquire) { - if reader.is_member(&acc("teq")) == Resolved::Ready(true) { + if reader.effective_member(&acc("teq")) == Resolved::Ready(true) { reader_leaked.store(true, Ordering::Release); } } @@ -88,9 +88,9 @@ fn a_concurrent_reader_never_sees_a_net_absent_subject() { !leaked.load(Ordering::Acquire), "net-absent subject is never written, so no reader can observe it mid-delta" ); - assert_eq!(index.is_member(&acc("teq")), Resolved::Ready(false)); - assert_eq!(index.is_member(&acc("f0")), Resolved::Ready(true)); - assert_eq!(index.is_member(&acc("f999")), Resolved::Ready(true)); + assert_eq!(index.effective_member(&acc("teq")), Resolved::Ready(false)); + assert_eq!(index.effective_member(&acc("f0")), Resolved::Ready(true)); + assert_eq!(index.effective_member(&acc("f999")), Resolved::Ready(true)); } #[test] @@ -124,7 +124,7 @@ fn a_concurrent_reader_never_sees_a_collaborator_roster_emptied_mid_refresh() { index.rebuild().unwrap(); index.ensure_collaborators(&repo).unwrap(); assert_eq!( - index.is_collaborator(&repo, &acc("anchor")), + index.effective_collaborator(&repo, &acc("anchor")), Resolved::Ready(true) ); @@ -137,7 +137,7 @@ fn a_concurrent_reader_never_sees_a_collaborator_roster_emptied_mid_refresh() { let target = repo.clone(); scope.spawn(move || { while !reader_done.load(Ordering::Acquire) { - if reader.is_collaborator(&target, &acc("anchor")) != Resolved::Ready(true) { + if reader.effective_collaborator(&target, &acc("anchor")) != Resolved::Ready(true) { reader_leaked.store(true, Ordering::Release); } } @@ -152,7 +152,7 @@ fn a_concurrent_reader_never_sees_a_collaborator_roster_emptied_mid_refresh() { pre- and post-refresh roster is never observed absent or warming mid-refresh" ); assert_eq!( - index.is_collaborator(&repo, &acc("anchor")), + index.effective_collaborator(&repo, &acc("anchor")), Resolved::Ready(true) ); } @@ -189,7 +189,7 @@ fn a_diverged_collaborators_tip_purges_the_roster_instead_of_serving_it_stale() index.rebuild().unwrap(); index.ensure_collaborators(&repo).unwrap(); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true) ); @@ -205,9 +205,9 @@ fn a_diverged_collaborators_tip_purges_the_roster_instead_of_serving_it_stale() "tip that no longer descends from folded tip is structural error" ); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Warming, - "diverged COB tip purges roster and fails closed, it does not serve \ + "diverged COB tip purges roster and fails closed, it doesn't serve \ pre-divergence collaborators" ); } @@ -230,7 +230,7 @@ fn a_diverged_members_tip_fails_closed_to_warming() { let index = world.index(); index.rebuild().unwrap(); - assert_eq!(index.is_member(&acc("nel")), Resolved::Ready(true)); + assert_eq!(index.effective_member(&acc("nel")), Resolved::Ready(true)); meta.update_ref(&RefUpdate::Update { name: cob_ref(MembersChange::TYPE, object), @@ -245,7 +245,7 @@ fn a_diverged_members_tip_fails_closed_to_warming() { ); assert_eq!(index.coverage().members, Coverage::Warming); assert_eq!( - index.is_member(&acc("nel")), + index.effective_member(&acc("nel")), Resolved::Warming, "diverged members COB fails projection closed instead of serving stale members" ); @@ -297,12 +297,12 @@ fn an_undecodable_change_fails_closed_with_no_partial_apply() { assert_eq!(index.coverage().members, Coverage::Warming); assert_eq!( - index.is_member(&acc("teq")), + index.effective_member(&acc("teq")), Resolved::Warming, "no partial apply: teq from pre-error change was never committed" ); assert_eq!( - index.is_member(&acc("nel")), + index.effective_member(&acc("nel")), Resolved::Warming, "structurally broken COB fails whole projection closed" ); @@ -346,7 +346,7 @@ fn deregister_purges_collaborators_fail_closed() { index.rebuild().unwrap(); index.warm_collaborators(); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true) ); @@ -369,7 +369,7 @@ fn deregister_purges_collaborators_fail_closed() { Resolved::Ready(None) ); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Warming, "deregistered repo's collaborators are purged and fail closed, not served stale" ); @@ -438,7 +438,7 @@ fn a_renamed_repo_keeps_both_rkeys_and_its_collaborators() { "new rkey is canonical" ); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true), "rename never evacuates repo, so its collaborators survive" ); @@ -462,7 +462,7 @@ fn a_renamed_repo_keeps_both_rkeys_and_its_collaborators() { "deregistering through retained alias removes repo and every alias" ); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Warming, "deregistered repo's collaborators are evacuated" ); @@ -499,7 +499,7 @@ fn a_repo_moved_within_one_delta_is_not_evacuated() { index.rebuild().unwrap(); index.warm_collaborators(); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true) ); @@ -535,7 +535,7 @@ fn a_repo_moved_within_one_delta_is_not_evacuated() { Resolved::Ready(Some(repo.clone())) ); assert_eq!( - index.is_collaborator(&repo, &acc("lyna")), + index.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true), "deregister and re-register within single delta leaves repo hosted, so its collaborators survive" ); @@ -932,6 +932,42 @@ fn a_warming_member_roll_still_yields_the_accounts_that_may_push() { ); } +#[test] +fn the_sweep_retains_keys_a_push_recorded_for_an_invited_collaborator() { + let world = World::new(); + let repo = repo_did("squid"); + world.invite_collaborator(&repo, "lyna"); + world.seed_registry(&repo); + let index = world.index(); + index.rebuild().unwrap(); + index.warm_collaborators(); + + let kept = || match index.keys().work(at(0), sweep()).map(|work| work.members) { + Resolved::Ready(Resolved::Ready(members)) => members.kept, + _ => panic!("the registry or the member roll is still warming"), + }; + assert!( + !kept().as_slice().contains(&acc("lyna")), + "the sweep is for accounts that may push, and lyna is only invited" + ); + + index + .keys() + .record(&acc("lyna"), vec![OfferedKey::from_bytes(vec![7])], hour(0)); + let vouching = kept(); + assert!( + vouching.as_slice().contains(&acc("lyna")), + "keeping lyna here saves the next push a second read of the same keys" + ); + + index.keys().retain(&vouching); + assert_eq!( + index.owner_of_key(&OfferedKey::from_bytes(vec![7]), at(0)), + Resolved::Ready(Some(acc("lyna"))), + "retain kept the key, so owner_of_key still finds lyna between sweeps" + ); +} + #[test] fn a_reprieve_coasts_on_the_last_read_until_its_budget_runs_out() { let (_world, index) = folded(); diff --git a/knot2/crates/knot-index/tests/properties.rs b/knot2/crates/knot-index/tests/properties.rs index 0cf5f0c05..b772fb32d 100644 --- a/knot2/crates/knot-index/tests/properties.rs +++ b/knot2/crates/knot-index/tests/properties.rs @@ -1,10 +1,10 @@ use knot_cob::{CobHome, CobStore}; use knot_cobs::{ - CollaboratorsChange, CollaboratorsCob, Grant, MembersChange, MembersCob, Registration, - RegistryChange, Removal, Rename, RepoRef, RepoRegistryCob, + Accept, CollaboratorsChange, CollaboratorsCob, Grant, Invite, MembersChange, MembersCob, + Registration, RegistryChange, Removal, Rename, RepoRef, RepoRegistryCob, }; use knot_git::Repo; -use knot_index::{Index, Resolved}; +use knot_index::{Index, Resolved, Standing}; use knot_types::{AccountDid, OwnerDid, RepoDid, RepoName, RepoRkey, UnixSeconds}; use proptest::prelude::*; @@ -38,9 +38,18 @@ fn grant(subject: u8, t: i64) -> Grant { } } +fn accept(subject: u8, t: i64) -> Accept { + Accept { + subject: acc(subject), + verified_at: UnixSeconds::new(t), + } +} + fn member_change(op: u8, subject: u8, t: i64) -> MembersChange { match op { 0 => MembersChange::Add(grant(subject, t)), + 1 => MembersChange::Accept(accept(subject, t)), + 2 => MembersChange::Invite(Invite(grant(subject, t))), _ => MembersChange::Remove(Removal { subject: acc(subject), }), @@ -50,6 +59,8 @@ fn member_change(op: u8, subject: u8, t: i64) -> MembersChange { fn collaborator_change(op: u8, subject: u8, t: i64) -> CollaboratorsChange { match op { 0 => CollaboratorsChange::Add(grant(subject, t)), + 1 => CollaboratorsChange::Accept(accept(subject, t)), + 2 => CollaboratorsChange::Invite(Invite(grant(subject, t))), _ => CollaboratorsChange::Remove(Removal { subject: acc(subject), }), @@ -83,7 +94,7 @@ proptest! { #[test] fn members_fold_equals_canonical_evaluate( - ops in prop::collection::vec((0u8..2, 0u8..4), 1..14) + ops in prop::collection::vec((0u8..4, 0u8..4), 1..14) ) { let world = World::seeded(7); let meta = Repo::open(&world.meta_path).unwrap(); @@ -111,22 +122,22 @@ proptest! { let canonical = store.get::(object).unwrap(); let roster = canonical.state(); - let expected: Vec> = (0u8..4) - .map(|subject| Resolved::Ready(roster.contains(&acc(subject)))) + let expected: Vec>> = (0u8..4) + .map(|subject| Resolved::Ready(roster.standing(&acc(subject)))) .collect(); prop_assert_eq!( - (0u8..4).map(|s| incremental.is_member(&acc(s))).collect::>(), + (0u8..4).map(|s| incremental.member_standing(&acc(s))).collect::>(), expected.clone() ); prop_assert_eq!( - (0u8..4).map(|s| full.is_member(&acc(s))).collect::>(), + (0u8..4).map(|s| full.member_standing(&acc(s))).collect::>(), expected ); } #[test] fn collaborators_fold_equals_canonical_evaluate( - ops in prop::collection::vec((0u8..2, 0u8..4), 1..14) + ops in prop::collection::vec((0u8..4, 0u8..4), 1..14) ) { let world = World::seeded(8); let repo = repo_did(0); @@ -157,15 +168,15 @@ proptest! { let canonical = store.get::(object).unwrap(); let roster = canonical.state(); - let expected: Vec> = (0u8..4) - .map(|subject| Resolved::Ready(roster.contains(&acc(subject)))) + let expected: Vec>> = (0u8..4) + .map(|subject| Resolved::Ready(roster.standing(&acc(subject)))) .collect(); prop_assert_eq!( - (0u8..4).map(|s| incremental.is_collaborator(&repo, &acc(s))).collect::>(), + (0u8..4).map(|s| incremental.collaborator_standing(&repo, &acc(s))).collect::>(), expected.clone() ); prop_assert_eq!( - (0u8..4).map(|s| full.is_collaborator(&repo, &acc(s))).collect::>(), + (0u8..4).map(|s| full.collaborator_standing(&repo, &acc(s))).collect::>(), expected ); } @@ -248,7 +259,7 @@ proptest! { #[test] fn two_rebuilds_are_observably_identical( - members in prop::collection::vec((0u8..2, 0u8..6), 0..16), + members in prop::collection::vec((0u8..3, 0u8..6), 0..16), registry in prop::collection::vec((0u8..3, 0u8..2, 0u8..4, 0u8..4), 0..10), ) { let world = World::seeded(10); @@ -287,8 +298,8 @@ proptest! { second.rebuild().unwrap(); prop_assert_eq!( - (0u8..6).map(|s| first.is_member(&acc(s))).collect::>(), - (0u8..6).map(|s| second.is_member(&acc(s))).collect::>() + (0u8..6).map(|s| first.effective_member(&acc(s))).collect::>(), + (0u8..6).map(|s| second.effective_member(&acc(s))).collect::>() ); let lookups: Vec<(u8, u8)> = (0u8..2) .flat_map(|who| (0u8..4).map(move |rkey| (who, rkey))) diff --git a/knot2/crates/knot-xrpc/src/consent.rs b/knot2/crates/knot-xrpc/src/consent.rs new file mode 100644 index 000000000..ed3729c81 --- /dev/null +++ b/knot2/crates/knot-xrpc/src/consent.rs @@ -0,0 +1,166 @@ +use knot_atproto::RecordPresence; +use knot_cobs::Standing; +use knot_index::Resolved; +use knot_runtime::{Clock, HttpTransport}; +use knot_types::{AccountDid, AtUri, DidRkey, KnotId, Nsid}; +use serde::Deserialize; + +use crate::XrpcState; +use crate::error::XrpcError; + +#[derive(Clone, Copy)] +pub(crate) enum Acceptances { + OfMembership, + OfCollaboration, +} + +impl Acceptances { + fn as_str(self) -> &'static str { + match self { + Self::OfMembership => knot_consent::OF_MEMBERSHIP, + Self::OfCollaboration => knot_consent::OF_COLLABORATION, + } + } + + fn nsid(self) -> Nsid { + Nsid::new_static(self.as_str()).expect("acceptance collection names must parse as nsids") + } +} + +#[derive(Deserialize)] +pub(crate) struct AcceptanceInput { + pub(crate) acceptance: AtUri, +} + +pub(crate) struct AcceptanceRef { + author: AccountDid, + subject: DidRkey, + collection: Acceptances, +} + +impl AcceptanceRef { + pub(crate) fn parse(uri: &AtUri, collection: Acceptances) -> Result { + let author = AccountDid::new(uri.authority().as_str()) + .map_err(|_| XrpcError::invalid_request("acceptance authority must be a DID"))?; + let named = uri + .collection() + .ok_or_else(|| XrpcError::invalid_request("no collection in the acceptance at-uri"))?; + if named.as_str() != collection.as_str() { + let (wanted, named) = (collection.as_str(), named.as_str()); + return Err(XrpcError::invalid_request(format!( + "acceptance belongs in {wanted}, and this one is in {named}" + ))); + } + let rkey = uri + .rkey() + .ok_or_else(|| XrpcError::invalid_request("no record key in the acceptance at-uri"))?; + let subject = DidRkey::new(rkey.as_str()).map_err(|_| { + XrpcError::invalid_request("acceptance record key must be a bare DID") + })?; + Ok(Self { + author, + subject, + collection, + }) + } + + pub(crate) fn subject(&self) -> &DidRkey { + &self.subject + } + + pub(crate) fn written_by(self, actor: &AccountDid) -> Result { + match self.author == *actor { + true => Ok(Acceptance(self)), + false => Err(XrpcError::forbidden( + "acceptance at-uri must point at your own repository", + )), + } + } +} + +pub(crate) struct Acceptance(AcceptanceRef); + +impl Acceptance { + pub(crate) fn names_knot(&self, knot: &KnotId) -> Result<(), XrpcError> { + match self.0.subject.as_str() == knot.as_str() { + true => Ok(()), + false => Err(XrpcError::forbidden("this acceptance is for another knot")), + } + } + + pub(crate) async fn require_published( + &self, + state: &XrpcState, + ) -> Result<(), XrpcError> { + match state + .atproto + .subject_record_present(&self.0.author, &self.0.collection.nsid(), &self.0.subject) + .await + { + Ok(RecordPresence::Present) => Ok(()), + Ok(RecordPresence::Absent) => Err(XrpcError::forbidden( + "no acceptance record in your repository, so write it first", + )), + Err(error) => Err(XrpcError::upstream_unavailable(format!( + "acceptance lookup failed: {error}" + ))), + } + } +} + +pub(crate) enum Consent { + Awaited, + Recorded, +} + +pub(crate) fn consent_gate( + standing: Resolved>, + no_offer: &'static str, +) -> Result { + match standing { + Resolved::Ready(Some(Standing::Accepted { .. })) => Ok(Consent::Recorded), + Resolved::Ready(Some(Standing::Invited | Standing::Grandfathered { .. })) => { + Ok(Consent::Awaited) + } + Resolved::Ready(None) => Err(XrpcError::forbidden(no_offer)), + Resolved::Warming => Err(XrpcError::warming( + "roster projection is still warming, retry shortly", + )), + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn uri(authority: &str, collection: &str, rkey: &str) -> String { + format!("at://{authority}/{collection}/{rkey}") + } + + #[test] + fn an_acceptance_at_uri_parses_with_its_own_collection_and_a_did_at_both_ends() { + let nel = "did:web:nel.pet"; + let knot = "did:web:knot.nel.pet"; + let mine = Acceptances::OfMembership.as_str(); + let theirs = Acceptances::OfCollaboration.as_str(); + [ + (uri(nel, mine, knot), true), + (uri(nel, theirs, knot), false), + (uri("nel.pet", mine, knot), false), + (uri(nel, mine, "did:web:localhost:3000"), false), + (format!("at://{nel}/{mine}"), false), + ] + .into_iter() + .for_each(|(raw, parses)| { + let input: AcceptanceInput = serde_json::from_value(serde_json::json!({ + "acceptance": &raw, + })) + .expect("every row is a syntactic at-uri; parse judges the rest"); + assert_eq!( + AcceptanceRef::parse(&input.acceptance, Acceptances::OfMembership).is_ok(), + parses, + "{raw}" + ); + }); + } +}