diff --git a/CHANGELOG.md b/CHANGELOG.md index 11074ba..08fd452 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -80,7 +80,15 @@ Expect a bit of downtime as a git data migraiton is performed on startup. The la repositories accessible to the service. With GRASP-06 enabled, crafted PR submissions could also write Git objects and PR refs into another hosted repository without its maintainer authorization. Operators must upgrade - to v3.0.0; GRASP-06 operators should review hosted repository integrity. + to v3.0.0. On every v3 startup, a non-blocking authorization-integrity pass + compares every served branch, tag, `HEAD`, and `refs/nostr/*` ref with the + accepted State, PR, and PR Update events (including precisely scoped + in-flight events). It repairs unambiguous differences and emits + `manual_inspection=true` errors without deleting unexplained PR refs that + may be evidence. GRASP-06 operators must check those logs and the terminal + `Git authorization-integrity startup pass completed` summary; any non-zero + `manual_inspection` or `failed` count means the named repository still + requires review. ### Added diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 5eb6f98..91ab6f1 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -90,6 +90,11 @@ runtime: identity stays local so the relay's existence is never advertised - Atomically checkpoint purgatory and rejected-event recovery state every 60 seconds without consuming the checkpoint during restore +- Start non-blocking storage-integrity and authorization-integrity passes after + database initialization. The former checks family objects and thin-view + wiring; the latter reconciles each served ref against accepted State, PR, + and PR Update events, auto-repairing only unambiguous differences and + logging preserved evidence for manual inspection - Serve HTTP + WebSocket until a caller-supplied shutdown future resolves, then stop background mutation, persist a final state snapshot (purgatory and rejected-events cache), and clean up placeholder refs diff --git a/docs/explanation/git-family-object-storage.md b/docs/explanation/git-family-object-storage.md index b3788b6..23ecccc 100644 --- a/docs/explanation/git-family-object-storage.md +++ b/docs/explanation/git-family-object-storage.md @@ -382,14 +382,15 @@ View lifecycle locks remain responsible for “may this path be removed?” The family lock is responsible for “is this object inventory and manifest update atomic?” Neither lock grants authorization. -## Integrity and healing +## Storage and authorization integrity -Integrity is defined for an object-format/identifier family, not for one -owner path. A pass checks the family pack set and object graph, then checks -every owner and `/prs/` view with the same identifier for the correct alternate -and for refs whose targets are available in the family. Multiple independent -histories in one identifier are valid, and unreachable objects are not an -error: retained delete-state and rollback data intentionally remain present. +Storage integrity is defined for an object-format/identifier family, not for +one owner path. The storage pass checks the family pack set and object graph, +then checks every owner and `/prs/` view with the same identifier for the +correct alternate and for refs whose targets are available in the family. +Multiple independent histories in one identifier are valid, and unreachable +objects are not an error: retained delete-state and rollback data intentionally +remain present. The relay starts one non-blocking pass after migration, database initialization, and construction of the hardened outbound Git client. Broken @@ -401,6 +402,33 @@ integrity check again. An unresolved family produces an `ERROR` log containing bounded counts and OID/diagnostic samples; it does not make availability depend on remote servers. +Authorization integrity is the second, view-scoped layer. It derives each +owner's branch, tag, and `HEAD` set from the NIP-01-preferred accepted State +event published by that owner's confirmed maintainer set. It derives +`refs/nostr/` from accepted PR and PR Update events using the same +maintainer-overlap propagation rule as normal event processing. GRASP-06 views +also require the signer and this relay's clone URL to name that exact +submitter/identifier coordinate. Exactly scoped active purgatory entries are +recognized so an in-flight push is not mistaken for corruption. + +The pass creates or updates missing/wrong authorized refs once their objects +are available, deletes stale State-governed branch and tag refs, and repairs +`HEAD`. Before giving up on a missing expected object it uses the same accepted +clone URLs and hardened fetch path as storage healing. It does not delete an +unexplained PR or unknown-namespace ref: that ref may be evidence of a past +authorization bypass. Instead it emits a bounded, structured `ERROR` with +`manual_inspection=true`, the view and ref, and actual/expected targets. An +authorized State still in purgatory also defers State mutation and is reported +for inspection rather than racing event promotion. + +Both startup passes run in the background after database initialization. They +do not extend the offline migration window or make relay availability depend +on a remote Git server. Their stable terminal log messages are `Git +storage-integrity startup pass completed` and `Git authorization-integrity +startup pass completed`. Because the authorization pass is online, it refreshes +the accepted events immediately before mutation and refuses to overwrite any +ref that changed after its initial snapshot. + This is also the migration repair path. Migration remains a deterministic, offline conversion that preserves every Git-readable object and the exact legacy refs. Once those paths are thin family views, the ordinary family pass @@ -408,7 +436,7 @@ can heal pre-existing missing objects. Unindexed legacy packs are preserved under `.grasp/migration/unindexed-packs/` when their backup is retired; there is no separate legacy repair subsystem. -Operators can queue the same identifier-scoped check in the live process: +Operators can queue an identifier-scoped storage check in the live process: ```console ngit-grasp integrity-check --identifier example diff --git a/docs/how-to/upgrade-git-family-storage.md b/docs/how-to/upgrade-git-family-storage.md index ed7e996..df6d07f 100644 --- a/docs/how-to/upgrade-git-family-storage.md +++ b/docs/how-to/upgrade-git-family-storage.md @@ -51,18 +51,42 @@ crash-safe, and safe to resume by restarting the same v3 release. It: Progress is stored below `.grasp/migration/`; the completed layout is marked by `.grasp/storage-version`. Do not edit these files while the service is running. +A later v3 restart sees that completed marker and does not repeat the offline +v2 conversion; the production-scale 51-minute cost applies to the first +v2-to-v3 migration, not every v3 deployment. -Confirm that the service reaches its normal listening state, then check the -terminal integrity summary: +Confirm that the service reaches its normal listening state, then follow the +logs until both terminal integrity summaries appear: ```text -Git identifier-family integrity startup pass completed +Git storage-integrity startup pass completed +Git authorization-integrity startup pass completed ``` -An `unresolved` or `failed` count above zero has a corresponding `ERROR` naming -the identifier. The server has already attempted automatic repair. A retained -backup protects an unhealthy family's legacy data, and a legacy shallow view -continues serving at its pre-upgrade level rather than being made less usable. +These passes are non-blocking: the relay is online while they inspect and heal +the migrated views. In the storage summary, an `unresolved` or `failed` count +above zero has a corresponding `ERROR` naming the identifier. In the +authorization summary, a `manual_inspection` or `failed` count above zero has +a corresponding `ERROR` naming the exact view, ref, actual target, expected +target, and reason. The server has already attempted safe automatic repair. +It deliberately preserves unexplained `refs/nostr/*` and unknown-namespace +refs as evidence instead of guessing that deletion is safe. + +The authorization pass treats the accepted event database as authoritative: + +- the latest State event from the owner's confirmed maintainer set defines + all and only `refs/heads/*`, `refs/tags/*`, and `HEAD`; +- accepted PR and PR Update events define `refs/nostr/` in every + owner view selected by the existing maintainer-overlap rules; +- a GRASP-06 contributor view additionally requires the event signer, clone + URL, and repository identifier to match that exact `/prs/` coordinate; +- precisely scoped active purgatory entries are tolerated as in-flight state. + Legacy unscoped placeholders are preserved but reported because they cannot + prove which owner and identifier originally received the push. + +A retained backup protects an unhealthy family's legacy data, and a legacy +shallow view continues serving at its pre-upgrade level rather than being made +less usable. ## Roll back diff --git a/src/git/authorization_integrity.rs b/src/git/authorization_integrity.rs new file mode 100644 index 0000000..8ae7839 --- /dev/null +++ b/src/git/authorization_integrity.rs @@ -0,0 +1,1232 @@ +//! Reconcile served Git refs with the Nostr events that authorize them. +//! +//! This is intentionally separate from [`super::integrity`]. The storage pass +//! proves that objects, packs, alternates, and ref targets are structurally +//! sound. This pass proves that the refs exposed by each repository view are +//! exactly those justified by accepted State, PR, and PR Update events. + +use std::collections::{BTreeMap, BTreeSet, HashSet}; +use std::path::Path; +use std::process::Command; + +use anyhow::{anyhow, Context, Result}; +use nostr_sdk::prelude::{Event, EventId, FromBech32, PublicKey, ToBech32}; +use tracing::{error, info, warn}; + +use super::authorization::{compute_membership, extract_commit_tag, RepositoryData}; +use super::integrity::{discover_families, discover_views, list_refs, FamilyRepairSource}; +use super::storage::{FamilyKey, LocalGitStorage, ObjectFormat}; +use crate::nostr::events::RepositoryState; +use crate::purgatory::sync::RealSyncContext; +use crate::purgatory::{PrPurgatoryEntry, StatePurgatoryEntry}; + +const MAX_ISSUES_PER_VIEW: usize = 16; + +#[derive(Debug, Default)] +struct PassStats { + families: usize, + views_checked: usize, + healthy: usize, + repaired: usize, + manual_inspection: usize, + failed: usize, + refs_created: usize, + refs_updated: usize, + refs_deleted: usize, + heads_set: usize, +} + +#[derive(Debug, Default)] +struct RepairCounts { + created: usize, + updated: usize, + deleted: usize, + heads_set: usize, +} + +impl RepairCounts { + fn any(&self) -> bool { + self.created + self.updated + self.deleted + self.heads_set > 0 + } +} + +#[derive(Debug)] +struct IntegrityIssue { + category: &'static str, + reference: String, + actual_oid: Option, + expected_oid: Option, + detail: String, +} + +impl IntegrityIssue { + fn new( + category: &'static str, + reference: impl Into, + actual_oid: Option, + expected_oid: Option, + detail: impl Into, + ) -> Self { + Self { + category, + reference: reference.into(), + actual_oid, + expected_oid, + detail: detail.into(), + } + } +} + +#[derive(Debug, Clone)] +enum ViewIdentity { + Owner { pubkey: PublicKey }, + Prs { submitter: PublicKey }, +} + +impl ViewIdentity { + fn kind(&self) -> &'static str { + match self { + Self::Owner { .. } => "owner", + Self::Prs { .. } => "prs", + } + } + + fn pubkey_hex(&self) -> String { + match self { + Self::Owner { pubkey } => pubkey.to_hex(), + Self::Prs { submitter } => submitter.to_hex(), + } + } +} + +#[derive(Debug, Default)] +struct ViewExpectation { + state_refs: BTreeMap, + state_head: Option, + has_authoritative_state: bool, + state_deferred: bool, + pr_refs: BTreeMap, + pending_pr_refs: BTreeMap, + ambiguous_placeholders: BTreeMap, + issues: Vec, +} + +/// Run the authorization-aware pass after storage integrity has completed. +/// +/// The relay is already serving traffic while this runs. Missing objects are +/// fetched only through the hardened integrity-fetch path, and each family is +/// leased only for the short local ref reconciliation phase. +pub async fn run_startup_pass(storage: &LocalGitStorage, source: &RealSyncContext) { + let mut stats = PassStats::default(); + let families = match discover_families(storage) { + Ok(families) => families, + Err(error) => { + error!(%error, "Git authorization-integrity startup discovery failed"); + stats.failed = 1; + log_completion(&stats); + return; + } + }; + let accepted_prs = match source.accepted_pr_events_for_integrity().await { + Ok(events) => events, + Err(error) => { + error!(%error, "Git authorization-integrity authoritative event query failed"); + stats.families = families.len(); + stats.failed = 1; + log_completion(&stats); + return; + } + }; + let pending_prs = source.pending_pr_entries_for_integrity(); + + info!( + families = families.len(), + accepted_pr_events = accepted_prs.len(), + pending_pr_entries = pending_prs.len(), + "Git authorization-integrity startup pass started" + ); + + for key in families { + stats.families += 1; + if let Err(error) = check_family( + storage, + source, + &key, + &accepted_prs, + &pending_prs, + &mut stats, + ) + .await + { + stats.failed += 1; + error!( + identifier = %key.identifier, + object_format = %key.object_format, + %error, + manual_inspection = true, + "Git authorization-integrity family check failed" + ); + } + tokio::task::yield_now().await; + } + + log_completion(&stats); +} + +fn log_completion(stats: &PassStats) { + info!( + families = stats.families, + views_checked = stats.views_checked, + healthy = stats.healthy, + repaired = stats.repaired, + manual_inspection = stats.manual_inspection, + failed = stats.failed, + refs_created = stats.refs_created, + refs_updated = stats.refs_updated, + refs_deleted = stats.refs_deleted, + heads_set = stats.heads_set, + "Git authorization-integrity startup pass completed" + ); +} + +async fn check_family( + storage: &LocalGitStorage, + source: &RealSyncContext, + key: &FamilyKey, + accepted_prs: &[Event], + pending_prs: &[(String, PrPurgatoryEntry)], + stats: &mut PassStats, +) -> Result<()> { + let repo_data = source + .accepted_repository_data_for_integrity(&key.identifier) + .await + .with_context(|| format!("load accepted events for {}", key.identifier))?; + let pending_states = source.pending_state_events_for_integrity(&key.identifier); + let views = discover_views(storage, key)?; + let mut checks = Vec::with_capacity(views.len()); + + for view in views { + let identity = view_identity(storage, &view)?; + let baseline_refs: BTreeMap<_, _> = list_refs(&view)?.into_iter().collect(); + let baseline_head = symbolic_head(&view)?; + let expectation = build_expectation( + &identity, + &key.identifier, + &repo_data, + accepted_prs, + pending_prs, + &pending_states, + source.service_address_for_integrity(), + ); + checks.push((view, identity, expectation, baseline_refs, baseline_head)); + } + + // Repair missing expected objects before taking the family lease: the + // hardened fetch path takes that same lease while installing objects. + let family = storage.family_repo_path(key); + let mut missing = BTreeSet::new(); + for (_, _, expectation, _, _) in &checks { + for oid in expectation + .state_refs + .values() + .chain(expectation.pr_refs.values()) + { + if valid_oid(key.object_format, oid) && !oid_exists(&family, oid)? { + missing.insert(oid.clone()); + } + } + } + if !missing.is_empty() { + if let Some((target_view, _, _, _, _)) = checks.first() { + fetch_missing_expected_oids(source, &key.identifier, target_view, &missing).await; + } + } + + let _lease = storage.write_lease(key).await?; + // The pass is online and remote object recovery may have taken time. + // Revalidate State and every candidate PR ID while the family is leased, + // and refuse to mutate a ref that changed since the initial snapshot. + let repo_data = source + .accepted_repository_data_for_integrity(&key.identifier) + .await + .with_context(|| format!("refresh accepted events for {}", key.identifier))?; + let pending_states = source.pending_state_events_for_integrity(&key.identifier); + let mut candidate_pr_ids: BTreeSet = accepted_prs + .iter() + .filter(|event| event_mentions_identifier(event, &key.identifier)) + .map(|event| event.id) + .collect(); + for (view, _, _, _, _) in &checks { + for (reference, _) in list_refs(view)? { + if let Some(event_id) = reference.strip_prefix("refs/nostr/") { + if let Ok(event_id) = EventId::from_hex(event_id) { + candidate_pr_ids.insert(event_id); + } + } + } + } + let fresh_prs = source + .accepted_pr_events_by_id_for_integrity(&candidate_pr_ids.into_iter().collect::>()) + .await + .with_context(|| format!("refresh accepted PR events for {}", key.identifier))?; + let pending_prs = source.pending_pr_entries_for_integrity(); + + for (view, identity, _, baseline_refs, baseline_head) in checks { + let expectation = build_expectation( + &identity, + &key.identifier, + &repo_data, + &fresh_prs, + &pending_prs, + &pending_states, + source.service_address_for_integrity(), + ); + stats.views_checked += 1; + match reconcile_view( + &view, + key.object_format, + expectation, + &baseline_refs, + baseline_head.as_deref(), + ) { + Ok((repairs, issues)) => { + stats.refs_created += repairs.created; + stats.refs_updated += repairs.updated; + stats.refs_deleted += repairs.deleted; + stats.heads_set += repairs.heads_set; + if issues.is_empty() { + if repairs.any() { + stats.repaired += 1; + info!( + identifier = %key.identifier, + object_format = %key.object_format, + view = %view.display(), + view_type = identity.kind(), + pubkey = %identity.pubkey_hex(), + refs_created = repairs.created, + refs_updated = repairs.updated, + refs_deleted = repairs.deleted, + heads_set = repairs.heads_set, + "Git authorization-integrity repaired repository view" + ); + } else { + stats.healthy += 1; + } + } else { + stats.manual_inspection += 1; + log_issues(key, &view, &identity, &issues); + } + } + Err(error) => { + stats.failed += 1; + error!( + identifier = %key.identifier, + object_format = %key.object_format, + view = %view.display(), + view_type = identity.kind(), + pubkey = %identity.pubkey_hex(), + %error, + manual_inspection = true, + "Git authorization-integrity repository check failed" + ); + } + } + } + Ok(()) +} + +async fn fetch_missing_expected_oids( + source: &RealSyncContext, + identifier: &str, + target_view: &Path, + missing: &BTreeSet, +) { + let urls = match source.accepted_clone_urls(identifier).await { + Ok(urls) => urls, + Err(error) => { + warn!(%identifier, %error, "Authorization-integrity could not load repair sources"); + return; + } + }; + let mut remaining: Vec = missing.iter().cloned().collect(); + for url in urls { + if remaining.is_empty() { + break; + } + match source + .fetch_missing_oids(target_view, &url, &remaining) + .await + { + Ok(_) => remaining.retain(|oid| !super::oid_exists(target_view, oid)), + Err(error) => warn!( + %identifier, + %url, + %error, + "Authorization-integrity repair source failed" + ), + } + } +} + +fn build_expectation( + identity: &ViewIdentity, + identifier: &str, + repo_data: &RepositoryData, + accepted_prs: &[Event], + pending_prs: &[(String, PrPurgatoryEntry)], + pending_states: &[StatePurgatoryEntry], + service_address: Option<&str>, +) -> ViewExpectation { + let mut expected = ViewExpectation::default(); + let maintainers = match identity { + ViewIdentity::Owner { pubkey } => { + let membership = + compute_membership(&repo_data.announcements, &pubkey.to_hex(), identifier); + let mut maintainers: Vec<_> = membership.state_maintainers.into_iter().collect(); + maintainers.sort(); + if let Some(state) = + super::sync::latest_authorized_state(&maintainers, &repo_data.states) + { + expected.has_authoritative_state = true; + add_state_expectation(&mut expected, state); + } + expected.state_deferred = pending_states + .iter() + .any(|entry| maintainers.contains(&entry.author.to_hex())); + Some(maintainers) + } + ViewIdentity::Prs { .. } => None, + }; + + for event in accepted_prs { + if event_applies_to_view( + event, + identity, + identifier, + maintainers.as_deref(), + service_address, + ) { + add_pr_expectation(&mut expected, event); + } + } + + for (event_id, entry) in pending_prs { + let reference = format!("refs/nostr/{event_id}"); + match &entry.event { + Some(event) + if event_applies_to_view( + event, + identity, + identifier, + maintainers.as_deref(), + service_address, + ) => + { + expected + .pending_pr_refs + .insert(reference, entry.commit.clone()); + } + None if placeholder_applies_to_view(entry, identity, identifier) => { + expected + .pending_pr_refs + .insert(reference, entry.commit.clone()); + } + None if entry.prs_scope.is_none() && matches!(identity, ViewIdentity::Owner { .. }) => { + // Legacy standard-endpoint placeholders did not record their + // owner/identifier. Preserve a matching ref, but make the + // uncertainty visible rather than treating it as authority. + expected + .ambiguous_placeholders + .insert(reference, entry.commit.clone()); + } + _ => {} + } + } + + expected +} + +fn add_state_expectation(expected: &mut ViewExpectation, state: &RepositoryState) { + for branch in &state.branches { + insert_expected_ref( + &mut expected.state_refs, + format!("refs/heads/{}", branch.name), + branch.commit.clone(), + &mut expected.issues, + "conflicting_state_ref", + ); + } + for tag in &state.tags { + insert_expected_ref( + &mut expected.state_refs, + format!("refs/tags/{}", tag.name), + tag.commit.clone(), + &mut expected.issues, + "conflicting_state_ref", + ); + } + expected.state_head = state.head.clone(); +} + +fn add_pr_expectation(expected: &mut ViewExpectation, event: &Event) { + let reference = format!("refs/nostr/{}", event.id); + let Some(commit) = extract_commit_tag(event) else { + expected.issues.push(IntegrityIssue::new( + "accepted_pr_missing_commit", + reference, + None, + None, + format!("accepted event {} has no usable c tag", event.id), + )); + return; + }; + insert_expected_ref( + &mut expected.pr_refs, + reference, + commit, + &mut expected.issues, + "conflicting_pr_ref", + ); +} + +fn insert_expected_ref( + refs: &mut BTreeMap, + reference: String, + oid: String, + issues: &mut Vec, + category: &'static str, +) { + if let Some(previous) = refs.insert(reference.clone(), oid.clone()) { + if previous != oid { + issues.push(IntegrityIssue::new( + category, + reference, + Some(previous), + Some(oid), + "authoritative events specify conflicting targets", + )); + } + } +} + +fn event_applies_to_view( + event: &Event, + identity: &ViewIdentity, + identifier: &str, + maintainers: Option<&[String]>, + service_address: Option<&str>, +) -> bool { + match identity { + ViewIdentity::Owner { .. } => { + let tagged = tagged_owners(event, identifier); + maintainers.is_some_and(|maintainers| { + maintainers + .iter() + .any(|maintainer| tagged.contains(maintainer)) + }) + } + ViewIdentity::Prs { submitter } => { + event.pubkey == *submitter + && service_address.is_some_and(|domain| { + crate::grasp06::policy::prs_identifiers_named_by_event_clone_tags(event, domain) + .iter() + .any(|candidate| candidate == identifier) + }) + } + } +} + +fn tagged_owners(event: &Event, identifier: &str) -> HashSet { + event + .tags + .iter() + .filter_map(|tag| { + let values = tag.as_slice(); + if values.first().map(String::as_str) != Some("a") { + return None; + } + let coordinate = values.get(1)?; + let mut parts = coordinate.splitn(3, ':'); + if parts.next()? != "30617" { + return None; + } + let owner = PublicKey::parse(parts.next()?).ok()?; + (parts.next()? == identifier).then(|| owner.to_hex()) + }) + .collect() +} + +fn event_mentions_identifier(event: &Event, identifier: &str) -> bool { + !tagged_owners(event, identifier).is_empty() +} + +fn placeholder_applies_to_view( + entry: &PrPurgatoryEntry, + identity: &ViewIdentity, + identifier: &str, +) -> bool { + let (ViewIdentity::Prs { submitter }, Some(scope)) = (identity, &entry.prs_scope) else { + return false; + }; + scope.submitter == *submitter && scope.identifier == identifier +} + +fn reconcile_view( + view: &Path, + object_format: ObjectFormat, + mut expected: ViewExpectation, + baseline_refs: &BTreeMap, + baseline_head: Option<&str>, +) -> Result<(RepairCounts, Vec)> { + let current: BTreeMap<_, _> = list_refs(view)?.into_iter().collect(); + let mut repairs = RepairCounts::default(); + + if expected.state_deferred { + expected.issues.push(IntegrityIssue::new( + "pending_authorized_state", + "HEAD and refs/heads/*, refs/tags/*", + None, + None, + "an authorized State event is still in purgatory; automatic State reconciliation was deferred", + )); + } else if expected.has_authoritative_state { + for (reference, actual) in current.iter().filter(|(reference, _)| { + reference.starts_with("refs/heads/") || reference.starts_with("refs/tags/") + }) { + if !expected.state_refs.contains_key(reference) { + if baseline_refs.get(reference) != Some(actual) { + expected.issues.push(IntegrityIssue::new( + "concurrent_ref_change", + reference, + Some(actual.clone()), + None, + "preserved because the ref changed while the online integrity pass was running", + )); + continue; + } + match delete_ref_cas(view, reference, actual) { + Ok(()) => repairs.deleted += 1, + Err(error) => expected.issues.push(IntegrityIssue::new( + "stale_state_ref_repair_failed", + reference, + Some(actual.clone()), + None, + error.to_string(), + )), + } + } + } + reconcile_expected_refs( + view, + object_format, + ¤t, + baseline_refs, + &expected.state_refs, + &mut repairs, + &mut expected.issues, + ); + reconcile_head( + view, + expected.state_head.as_deref(), + &expected.state_refs, + baseline_head, + &mut repairs, + &mut expected.issues, + ); + } else { + for (reference, actual) in current.iter().filter(|(reference, _)| { + reference.starts_with("refs/heads/") || reference.starts_with("refs/tags/") + }) { + expected.issues.push(IntegrityIssue::new( + "state_ref_without_authoritative_state", + reference, + Some(actual.clone()), + None, + "preserved because no accepted State event authorizes an automatic decision", + )); + } + } + + reconcile_expected_refs( + view, + object_format, + ¤t, + baseline_refs, + &expected.pr_refs, + &mut repairs, + &mut expected.issues, + ); + + for (reference, actual) in current + .iter() + .filter(|(reference, _)| reference.starts_with("refs/nostr/")) + { + if expected.pr_refs.contains_key(reference) { + continue; + } + if expected.pending_pr_refs.get(reference) == Some(actual) { + continue; + } + if expected.ambiguous_placeholders.get(reference) == Some(actual) { + expected.issues.push(IntegrityIssue::new( + "unscoped_pending_placeholder", + reference, + Some(actual.clone()), + Some(actual.clone()), + "preserved, but the legacy placeholder does not record which owner and identifier received the push", + )); + continue; + } + expected.issues.push(IntegrityIssue::new( + "unauthorized_pr_ref", + reference, + Some(actual.clone()), + expected.pending_pr_refs.get(reference).cloned(), + "preserved for manual inspection because no accepted or exactly scoped pending PR event authorizes it", + )); + } + + for (reference, actual) in ¤t { + if reference.starts_with("refs/heads/") + || reference.starts_with("refs/tags/") + || reference.starts_with("refs/nostr/") + { + continue; + } + expected.issues.push(IntegrityIssue::new( + "unrecognized_served_ref", + reference, + Some(actual.clone()), + None, + "preserved for manual inspection because this ref namespace is not derived from State or PR events", + )); + } + + Ok((repairs, expected.issues)) +} + +fn reconcile_expected_refs( + view: &Path, + object_format: ObjectFormat, + current: &BTreeMap, + baseline: &BTreeMap, + expected: &BTreeMap, + repairs: &mut RepairCounts, + issues: &mut Vec, +) { + for (reference, target) in expected { + let actual = current.get(reference); + if actual == Some(target) { + continue; + } + if baseline.get(reference) != actual { + issues.push(IntegrityIssue::new( + "concurrent_ref_change", + reference, + actual.cloned(), + Some(target.clone()), + "preserved because the ref changed while the online integrity pass was running", + )); + continue; + } + if !valid_oid(object_format, target) { + issues.push(IntegrityIssue::new( + "invalid_authoritative_oid", + reference, + actual.cloned(), + Some(target.clone()), + "authoritative event target is not a valid object ID for this family", + )); + continue; + } + if !super::oid_exists(view, target) { + issues.push(IntegrityIssue::new( + "missing_authoritative_object", + reference, + actual.cloned(), + Some(target.clone()), + "the expected object remains unavailable after accepted-source repair attempts", + )); + continue; + } + match update_ref_cas( + view, + reference, + target, + actual.map(String::as_str), + object_format, + ) { + Ok(()) if actual.is_some() => repairs.updated += 1, + Ok(()) => repairs.created += 1, + Err(error) => issues.push(IntegrityIssue::new( + "authorized_ref_repair_failed", + reference, + actual.cloned(), + Some(target.clone()), + error.to_string(), + )), + } + } +} + +fn reconcile_head( + view: &Path, + expected_head: Option<&str>, + state_refs: &BTreeMap, + baseline_head: Option<&str>, + repairs: &mut RepairCounts, + issues: &mut Vec, +) { + let Some(expected_head) = expected_head else { + return; + }; + if !expected_head.starts_with("refs/heads/") || !state_refs.contains_key(expected_head) { + issues.push(IntegrityIssue::new( + "invalid_authoritative_head", + "HEAD", + symbolic_head(view).ok().flatten(), + Some(expected_head.to_owned()), + "State HEAD must name a branch declared by the same State event", + )); + return; + } + let actual = match symbolic_head(view) { + Ok(actual) => actual, + Err(error) => { + issues.push(IntegrityIssue::new( + "head_inspection_failed", + "HEAD", + None, + Some(expected_head.to_owned()), + error.to_string(), + )); + return; + } + }; + if actual.as_deref() == Some(expected_head) { + return; + } + if actual.as_deref() != baseline_head { + issues.push(IntegrityIssue::new( + "concurrent_head_change", + "HEAD", + actual, + Some(expected_head.to_owned()), + "preserved because HEAD changed while the online integrity pass was running", + )); + return; + } + match super::set_repository_head(view, expected_head) { + Ok(()) => repairs.heads_set += 1, + Err(error) => issues.push(IntegrityIssue::new( + "head_repair_failed", + "HEAD", + actual, + Some(expected_head.to_owned()), + error, + )), + } +} + +fn view_identity(storage: &LocalGitStorage, view: &Path) -> Result { + let relative = view + .strip_prefix(storage.git_data_path()) + .with_context(|| format!("view {} escapes Git data path", view.display()))?; + let parts: Vec<_> = relative.components().collect(); + match parts.as_slice() { + [owner, _repo] => { + let owner = owner.as_os_str().to_string_lossy(); + let pubkey = PublicKey::from_bech32(&owner) + .with_context(|| format!("invalid owner directory {owner}"))?; + if pubkey.to_bech32().ok().as_deref() != Some(owner.as_ref()) { + return Err(anyhow!("non-canonical owner directory {owner}")); + } + Ok(ViewIdentity::Owner { pubkey }) + } + [namespace, submitter, _repo] + if namespace.as_os_str() == crate::grasp06::paths::PRS_DISK_PREFIX => + { + let submitter = submitter.as_os_str().to_string_lossy(); + let pubkey = PublicKey::parse(submitter.as_ref()) + .with_context(|| format!("invalid PR submitter directory {submitter}"))?; + if pubkey.to_hex() != submitter { + return Err(anyhow!("non-canonical PR submitter directory {submitter}")); + } + Ok(ViewIdentity::Prs { submitter: pubkey }) + } + _ => Err(anyhow!("unrecognized repository view {}", view.display())), + } +} + +fn valid_oid(object_format: ObjectFormat, oid: &str) -> bool { + let length = match object_format { + ObjectFormat::Sha1 => 40, + ObjectFormat::Sha256 => 64, + }; + oid.len() == length && oid.chars().all(|character| character.is_ascii_hexdigit()) +} + +fn oid_exists(repo: &Path, oid: &str) -> Result { + let status = Command::new("git") + .args(["cat-file", "-e", oid]) + .current_dir(repo) + .status() + .with_context(|| format!("inspect object {oid} in {}", repo.display()))?; + Ok(status.success()) +} + +fn update_ref_cas( + repo: &Path, + reference: &str, + target: &str, + actual: Option<&str>, + object_format: ObjectFormat, +) -> Result<()> { + let zero = match object_format { + ObjectFormat::Sha1 => "0000000000000000000000000000000000000000", + ObjectFormat::Sha256 => "0000000000000000000000000000000000000000000000000000000000000000", + }; + let output = Command::new("git") + .args(["update-ref", reference, target, actual.unwrap_or(zero)]) + .current_dir(repo) + .output() + .with_context(|| format!("update {reference} in {}", repo.display()))?; + if output.status.success() { + Ok(()) + } else { + Err(anyhow!( + "git update-ref failed: {}", + String::from_utf8_lossy(&output.stderr).trim() + )) + } +} + +fn delete_ref_cas(repo: &Path, reference: &str, actual: &str) -> Result<()> { + let output = Command::new("git") + .args(["update-ref", "-d", reference, actual]) + .current_dir(repo) + .output() + .with_context(|| format!("delete {reference} in {}", repo.display()))?; + if output.status.success() { + Ok(()) + } else { + Err(anyhow!( + "git update-ref -d failed: {}", + String::from_utf8_lossy(&output.stderr).trim() + )) + } +} + +fn symbolic_head(repo: &Path) -> Result> { + let output = Command::new("git") + .args(["symbolic-ref", "-q", "HEAD"]) + .current_dir(repo) + .output() + .with_context(|| format!("inspect HEAD in {}", repo.display()))?; + if output.status.success() { + Ok(Some(String::from_utf8(output.stdout)?.trim().to_owned())) + } else if output.status.code() == Some(1) { + Ok(None) + } else { + Err(anyhow!( + "git symbolic-ref failed: {}", + String::from_utf8_lossy(&output.stderr).trim() + )) + } +} + +fn log_issues(key: &FamilyKey, view: &Path, identity: &ViewIdentity, issues: &[IntegrityIssue]) { + for issue in issues.iter().take(MAX_ISSUES_PER_VIEW) { + error!( + identifier = %key.identifier, + object_format = %key.object_format, + view = %view.display(), + view_type = identity.kind(), + pubkey = %identity.pubkey_hex(), + category = issue.category, + reference = %issue.reference, + actual_oid = issue.actual_oid.as_deref().unwrap_or(""), + expected_oid = issue.expected_oid.as_deref().unwrap_or(""), + detail = %issue.detail, + manual_inspection = true, + "Git authorization-integrity requires manual inspection" + ); + } + if issues.len() > MAX_ISSUES_PER_VIEW { + error!( + identifier = %key.identifier, + object_format = %key.object_format, + view = %view.display(), + suppressed = issues.len() - MAX_ISSUES_PER_VIEW, + manual_inspection = true, + "Additional Git authorization-integrity issues suppressed for this view" + ); + } +} + +#[cfg(test)] +mod tests { + use std::io::Write; + use std::path::PathBuf; + use std::process::Stdio; + + use nostr_sdk::prelude::{EventBuilder, FinalizeEvent, Keys, Kind, Tag}; + + use super::*; + use crate::nostr::events::RepositoryAnnouncement; + + fn git(repo: &Path, args: &[&str]) -> String { + let output = Command::new("git") + .args(args) + .current_dir(repo) + .env("GIT_CONFIG_NOSYSTEM", "1") + .output() + .unwrap(); + assert!( + output.status.success(), + "git {args:?}: {}", + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8_lossy(&output.stdout).trim().to_owned() + } + + fn ref_exists(repo: &Path, reference: &str) -> bool { + Command::new("git") + .args(["show-ref", "--verify", "--quiet", reference]) + .current_dir(repo) + .status() + .unwrap() + .success() + } + + fn write_commit(repo: &Path, message: &str) -> String { + let tree = git(repo, &["mktree"]); + let contents = format!( + "tree {tree}\nauthor Test 0 +0000\ncommitter Test 0 +0000\n\n{message}\n" + ); + let mut child = Command::new("git") + .args(["hash-object", "-t", "commit", "-w", "--stdin"]) + .current_dir(repo) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .unwrap(); + child + .stdin + .take() + .unwrap() + .write_all(contents.as_bytes()) + .unwrap(); + let output = child.wait_with_output().unwrap(); + assert!(output.status.success()); + String::from_utf8_lossy(&output.stdout).trim().to_owned() + } + + fn fixture() -> ( + tempfile::TempDir, + LocalGitStorage, + FamilyKey, + PathBuf, + Vec, + ) { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let key = FamilyKey::sha1("project").unwrap(); + let family = storage.ensure_family(&key).unwrap(); + let commits = ["old", "new", "pr", "unexpected"] + .into_iter() + .map(|message| write_commit(&family, message)) + .collect::>(); + let owner = Keys::generate().public_key().to_bech32().unwrap(); + let view = storage.git_data_path().join(owner).join("project.git"); + storage.create_thin_view(&key, &view).unwrap(); + (temp, storage, key, view, commits) + } + + fn reconcile_fixture( + view: &Path, + expectation: ViewExpectation, + ) -> (RepairCounts, Vec) { + let baseline_refs = list_refs(view).unwrap().into_iter().collect(); + let baseline_head = symbolic_head(view).unwrap(); + reconcile_view( + view, + ObjectFormat::Sha1, + expectation, + &baseline_refs, + baseline_head.as_deref(), + ) + .unwrap() + } + + #[test] + fn repairs_authoritative_refs_and_preserves_unexpected_pr_evidence() { + let (_temp, _storage, _key, view, commits) = fixture(); + let expected_pr = format!("refs/nostr/{}", "a".repeat(64)); + let unexpected_pr = format!("refs/nostr/{}", "b".repeat(64)); + git(&view, &["update-ref", "refs/heads/main", &commits[0]]); + git(&view, &["update-ref", "refs/heads/stale", &commits[0]]); + git(&view, &["update-ref", &unexpected_pr, &commits[3]]); + let expectation = ViewExpectation { + state_refs: BTreeMap::from([("refs/heads/main".to_owned(), commits[1].clone())]), + state_head: Some("refs/heads/main".to_owned()), + has_authoritative_state: true, + pr_refs: BTreeMap::from([(expected_pr.clone(), commits[2].clone())]), + ..ViewExpectation::default() + }; + + let (repairs, issues) = reconcile_fixture(&view, expectation); + + assert_eq!(repairs.updated, 1); + assert_eq!(repairs.created, 1); + assert_eq!(repairs.deleted, 1); + assert_eq!(git(&view, &["rev-parse", "refs/heads/main"]), commits[1]); + assert_eq!(git(&view, &["rev-parse", &expected_pr]), commits[2]); + assert_eq!(git(&view, &["rev-parse", &unexpected_pr]), commits[3]); + assert!(!ref_exists(&view, "refs/heads/stale")); + assert_eq!(issues.len(), 1); + assert_eq!(issues[0].category, "unauthorized_pr_ref"); + } + + #[test] + fn exact_authorized_view_is_healthy() { + let (_temp, _storage, _key, view, commits) = fixture(); + let pr_ref = format!("refs/nostr/{}", "c".repeat(64)); + git(&view, &["update-ref", "refs/heads/main", &commits[0]]); + git(&view, &["update-ref", &pr_ref, &commits[2]]); + let expectation = ViewExpectation { + state_refs: BTreeMap::from([("refs/heads/main".to_owned(), commits[0].clone())]), + state_head: Some("refs/heads/main".to_owned()), + has_authoritative_state: true, + pr_refs: BTreeMap::from([(pr_ref, commits[2].clone())]), + ..ViewExpectation::default() + }; + + let (repairs, issues) = reconcile_fixture(&view, expectation); + + assert!(!repairs.any()); + assert!(issues.is_empty(), "{issues:#?}"); + } + + #[test] + fn missing_authoritative_object_requires_manual_inspection() { + let (_temp, _storage, _key, view, commits) = fixture(); + git(&view, &["update-ref", "refs/heads/main", &commits[0]]); + let missing = "f".repeat(40); + let expectation = ViewExpectation { + state_refs: BTreeMap::from([("refs/heads/main".to_owned(), missing.clone())]), + state_head: Some("refs/heads/main".to_owned()), + has_authoritative_state: true, + ..ViewExpectation::default() + }; + + let (repairs, issues) = reconcile_fixture(&view, expectation); + + assert!(!repairs.any()); + assert_eq!(git(&view, &["rev-parse", "refs/heads/main"]), commits[0]); + assert!(issues.iter().any(|issue| { + issue.category == "missing_authoritative_object" + && issue.expected_oid.as_deref() == Some(missing.as_str()) + })); + } + + #[test] + fn online_pass_never_overwrites_a_ref_changed_after_its_snapshot() { + let (_temp, _storage, _key, view, commits) = fixture(); + git(&view, &["update-ref", "refs/heads/main", &commits[0]]); + let baseline_refs = list_refs(&view).unwrap().into_iter().collect(); + let baseline_head = symbolic_head(&view).unwrap(); + git(&view, &["update-ref", "refs/heads/main", &commits[1]]); + let expectation = ViewExpectation { + state_refs: BTreeMap::from([("refs/heads/main".to_owned(), commits[2].clone())]), + state_head: Some("refs/heads/main".to_owned()), + has_authoritative_state: true, + ..ViewExpectation::default() + }; + + let (repairs, issues) = reconcile_view( + &view, + ObjectFormat::Sha1, + expectation, + &baseline_refs, + baseline_head.as_deref(), + ) + .unwrap(); + + assert!(!repairs.any()); + assert_eq!(git(&view, &["rev-parse", "refs/heads/main"]), commits[1]); + assert!(issues + .iter() + .any(|issue| issue.category == "concurrent_ref_change")); + } + + #[test] + fn expectations_follow_authorized_state_and_exact_pr_coordinates() { + let owner = Keys::generate(); + let submitter = Keys::generate(); + let other_owner = Keys::generate(); + let commit = "1".repeat(40); + let pr_commit = "2".repeat(40); + let announcement_event = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags([Tag::identifier("project")]) + .finalize(&owner) + .unwrap(); + let state_event = EventBuilder::new(Kind::RepoState, "") + .tags([ + Tag::identifier("project"), + Tag::custom("refs/heads/main", [commit.clone()]), + Tag::custom("HEAD", ["ref: refs/heads/main"]), + ]) + .finalize(&owner) + .unwrap(); + let matching_pr = EventBuilder::new(Kind::GitPullRequest, "") + .tags([ + Tag::custom( + "a", + [format!("30617:{}:project", owner.public_key().to_hex())], + ), + Tag::custom("c", [pr_commit.clone()]), + ]) + .finalize(&submitter) + .unwrap(); + let other_identifier_pr = EventBuilder::new(Kind::GitPullRequestUpdate, "") + .tags([ + Tag::custom( + "a", + [format!( + "30617:{}:different-project", + other_owner.public_key().to_hex() + )], + ), + Tag::custom("c", ["3".repeat(40)]), + ]) + .finalize(&submitter) + .unwrap(); + let repo_data = RepositoryData { + announcements: vec![RepositoryAnnouncement::from_event(announcement_event).unwrap()], + states: vec![RepositoryState::from_event(state_event).unwrap()], + }; + + let expectation = build_expectation( + &ViewIdentity::Owner { + pubkey: owner.public_key(), + }, + "project", + &repo_data, + &[matching_pr.clone(), other_identifier_pr], + &[], + &[], + Some("relay.example"), + ); + + assert_eq!(expectation.state_refs.get("refs/heads/main"), Some(&commit)); + assert_eq!(expectation.state_head.as_deref(), Some("refs/heads/main")); + assert_eq!( + expectation + .pr_refs + .get(&format!("refs/nostr/{}", matching_pr.id)), + Some(&pr_commit) + ); + assert_eq!(expectation.pr_refs.len(), 1); + } +} diff --git a/src/git/integrity.rs b/src/git/integrity.rs index d1c3669..de4fad3 100644 --- a/src/git/integrity.rs +++ b/src/git/integrity.rs @@ -165,6 +165,7 @@ pub fn spawn_integrity_worker( ) -> JoinHandle<()> { tokio::spawn(async move { run_startup_pass(&storage, source.as_ref()).await; + super::authorization_integrity::run_startup_pass(&storage, source.as_ref()).await; let first = tokio::time::Instant::now() + REQUEST_POLL_INTERVAL; let mut interval = tokio::time::interval_at(first, REQUEST_POLL_INTERVAL); loop { @@ -178,13 +179,13 @@ async fn run_startup_pass(storage: &LocalGitStor let families = match discover_families(storage) { Ok(families) => families, Err(error) => { - error!(%error, "Git identifier-family integrity startup discovery failed"); + error!(%error, "Git storage-integrity startup discovery failed"); return; } }; info!( families = families.len(), - "Git identifier-family integrity startup pass started" + "Git storage-integrity startup pass started" ); let mut stats = PassStats::default(); for key in families { @@ -197,7 +198,7 @@ async fn run_startup_pass(storage: &LocalGitStor repaired = stats.repaired, unresolved = stats.unresolved, failed = stats.failed, - "Git identifier-family integrity startup pass completed" + "Git storage-integrity startup pass completed" ); } @@ -651,7 +652,7 @@ pub(crate) fn family_contains_closure(family: &Path, refs_source: &Path) -> Resu Ok(()) } -fn discover_views(storage: &LocalGitStorage, key: &FamilyKey) -> Result> { +pub(crate) fn discover_views(storage: &LocalGitStorage, key: &FamilyKey) -> Result> { let mut views = Vec::new(); let entries = match std::fs::read_dir(storage.git_data_path()) { Ok(entries) => entries, @@ -699,7 +700,7 @@ fn maybe_add_view(directory: &Path, key: &FamilyKey, views: &mut Vec) - Ok(()) } -fn list_refs(repo: &Path) -> Result> { +pub(crate) fn list_refs(repo: &Path) -> Result> { let output = Command::new("git") .args(["for-each-ref", "--format=%(refname)%00%(objectname)"]) .current_dir(repo) diff --git a/src/git/mod.rs b/src/git/mod.rs index 8c09327..d93a684 100644 --- a/src/git/mod.rs +++ b/src/git/mod.rs @@ -18,6 +18,7 @@ //! - `POST //.git/git-receive-pack` - Push operation pub mod authorization; +pub mod authorization_integrity; pub mod handlers; pub mod integrity; pub mod migration; diff --git a/src/git/sync.rs b/src/git/sync.rs index 89a2b1b..3f7673e 100644 --- a/src/git/sync.rs +++ b/src/git/sync.rs @@ -1351,25 +1351,23 @@ async fn process_purgatory_state_events( result } -/// Check if a state event is the latest authorized state for a given maintainer set. +/// Select the latest authorized State event for a given maintainer set. /// /// Only considers states already in the database, not other purgatory states. /// /// # Arguments -/// * `state` - The state event to check /// * `maintainers` - The set of authorized maintainers for the owner /// * `db_states` - State events from the database /// /// # Returns -/// true if this state is the latest (or equal latest) among all authorized states in the DB -fn is_latest_authorized_state( - state: &RepositoryState, +/// The NIP-01-preferred State: newest timestamp, then lowest event ID. +pub fn latest_authorized_state<'a>( maintainers: &[String], - db_states: &[RepositoryState], -) -> bool { + db_states: &'a [RepositoryState], +) -> Option<&'a RepositoryState> { // Find the NIP-01-preferred authorized state from the database: newest // timestamp wins, and the lowest event ID wins an equal-timestamp tie. - let latest_db_state = db_states + db_states .iter() .filter(|s| maintainers.contains(&s.event.pubkey.to_hex())) .max_by(|a, b| { @@ -1379,7 +1377,15 @@ fn is_latest_authorized_state( .created_at .cmp(&b.event.created_at) .then_with(|| b.event.id.cmp(&a.event.id)) - }); + }) +} + +fn is_latest_authorized_state( + state: &RepositoryState, + maintainers: &[String], + db_states: &[RepositoryState], +) -> bool { + let latest_db_state = latest_authorized_state(maintainers, db_states); match latest_db_state { None => true, // No other states exist in DB, this is the latest diff --git a/src/purgatory/mod.rs b/src/purgatory/mod.rs index 96c2334..4299505 100644 --- a/src/purgatory/mod.rs +++ b/src/purgatory/mod.rs @@ -1121,6 +1121,19 @@ impl Purgatory { .collect() } + /// Snapshot every active PR event and git-data-first placeholder. + /// + /// The authorization-integrity pass needs the event ID (the map key) as + /// well as the entry so it can distinguish a legitimate in-flight ref + /// from an unexplained served ref. Keeping this crate-private avoids + /// exposing the purgatory's internal indexing as public API. + pub(crate) fn pr_entries_for_integrity(&self) -> Vec<(String, PrPurgatoryEntry)> { + self.pr_events + .iter() + .map(|entry| (entry.key().clone(), entry.value().clone())) + .collect() + } + /// Remove expired entries from purgatory. /// /// Should be called periodically (every 60 seconds) by background task to clean up diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index 9cebcea..e9ac583 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -360,6 +360,71 @@ impl RealSyncContext { urls.dedup(); Ok(urls) } + + pub(crate) async fn accepted_repository_data_for_integrity( + &self, + identifier: &str, + ) -> Result { + crate::git::authorization::fetch_repository_data_excluding_purgatory( + &self.database, + identifier, + ) + .await + } + + /// Accepted PR and PR-update events used to reconstruct the authorized + /// `refs/nostr/*` surface during the startup integrity pass. + pub(crate) async fn accepted_pr_events_for_integrity( + &self, + ) -> Result> { + use nostr_sdk::prelude::{Filter, Kind}; + + self.database + .query(Filter::new().kinds([Kind::GitPullRequest, Kind::GitPullRequestUpdate])) + .await + .map(|events| events.into_iter().collect()) + .map_err(|error| anyhow::anyhow!("Database query failed: {error}")) + } + + /// Revalidate a bounded set of PR IDs immediately before ref mutation. + /// Deleted events disappear from this result, preventing a long-running + /// online pass from recreating refs from its older startup snapshot. + pub(crate) async fn accepted_pr_events_by_id_for_integrity( + &self, + ids: &[nostr_sdk::prelude::EventId], + ) -> Result> { + use nostr_sdk::prelude::{Filter, Kind}; + + if ids.is_empty() { + return Ok(Vec::new()); + } + self.database + .query( + Filter::new() + .ids(ids.iter().copied()) + .kinds([Kind::GitPullRequest, Kind::GitPullRequestUpdate]), + ) + .await + .map(|events| events.into_iter().collect()) + .map_err(|error| anyhow::anyhow!("Database query failed: {error}")) + } + + pub(crate) fn pending_pr_entries_for_integrity( + &self, + ) -> Vec<(String, crate::purgatory::PrPurgatoryEntry)> { + self.purgatory.pr_entries_for_integrity() + } + + pub(crate) fn pending_state_events_for_integrity( + &self, + identifier: &str, + ) -> Vec { + self.purgatory.find_state(identifier) + } + + pub(crate) fn service_address_for_integrity(&self) -> Option<&str> { + self.our_domain_value.as_deref() + } } const MISS_MEMO_TTL: Duration = Duration::from_secs(30 * 60); diff --git a/src/server.rs b/src/server.rs index 7d7a3e7..49ae032 100644 --- a/src/server.rs +++ b/src/server.rs @@ -456,7 +456,7 @@ impl RelayServer { git_storage, sync_ctx.clone(), )); - info!("Git identifier-family integrity worker started"); + info!("Git storage and authorization-integrity worker started"); // Create throttle manager for rate limiting remote git servers // Default: 5 concurrent requests per domain, 60 requests per minute per domain