commit f7e947c2e5404b14974eef2ed11cae31c6b55061
parent 09065a610d95e57acdc895a14c07580fa099e7c3
Author: triesap <tyson@radroots.org>
Date: Fri, 7 Aug 2026 06:36:07 +0000
fix(storage): rebuild current event visibility
- validate admission-selected contracts without weakening generic classification
- reduce current heads, deletion cutoffs, and ephemeral visibility deterministically
- preserve memory and SQLite parity across atomic ingest, pagination, and reopen
- harden downstream recovery semantics and corrupt-row rejection coverage
Diffstat:
11 files changed, 1155 insertions(+), 93 deletions(-)
diff --git a/crates/event/src/contract/registry_v7.rs b/crates/event/src/contract/registry_v7.rs
@@ -3926,6 +3926,7 @@ fn validate_event_contract_in_registry(
event.content(),
contract,
event_contracts,
+ false,
)?;
Ok(contract)
}
@@ -3938,6 +3939,31 @@ pub fn validate_event_contract_shape(
validate_event_contract_parts(event.kind_u32(), &tags, event.content(), contract_id)
}
+/// Validates a contract selected explicitly by an admission boundary.
+///
+/// This is the required path for contracts whose discriminator is
+/// [`EventDiscriminator::AdmissionOnly`]. The selected contract's kind,
+/// content, tags, and custom invariants are still validated in full.
+pub fn validate_event_contract_for_admission(
+ event: &EventEnvelope,
+ contract_id: &str,
+) -> Result<&'static EventContract, ContractValidationError> {
+ let contract =
+ event_contract(contract_id).ok_or_else(|| ContractValidationError::UnknownContract {
+ contract_id: contract_id.to_owned(),
+ })?;
+ let tags = event.tags_as_vec();
+ validate_event_contract_parts_in_registry(
+ event.kind_u32(),
+ &tags,
+ event.content(),
+ contract,
+ EVENT_CONTRACTS_REGISTRY_V7,
+ true,
+ )?;
+ Ok(contract)
+}
+
pub fn validate_event_contract_parts(
kind: u32,
tags: &[Vec<String>],
@@ -3954,6 +3980,7 @@ pub fn validate_event_contract_parts(
content,
contract,
EVENT_CONTRACTS_REGISTRY_V7,
+ false,
)
}
@@ -3963,6 +3990,7 @@ fn validate_event_contract_parts_in_registry(
content: &str,
contract: &EventContract,
event_contracts: &'static [EventContract],
+ admission_selected: bool,
) -> Result<(), ContractValidationError> {
crate::require_invariant(kind == contract.kind, &|| {
ContractValidationError::KindMismatch {
@@ -3971,7 +3999,7 @@ fn validate_event_contract_parts_in_registry(
}
})?;
crate::require_invariant(
- !matches!(contract.discriminator, EventDiscriminator::AdmissionOnly),
+ admission_selected || !matches!(contract.discriminator, EventDiscriminator::AdmissionOnly),
&|| ContractValidationError::AdmissionRequired {
contract_id: contract.id,
},
@@ -3979,12 +4007,33 @@ fn validate_event_contract_parts_in_registry(
validate_classified_listing_partition_parts(tags, contract)?;
validate_content_shape_parts(content, contract)?;
validate_contract_tags_parts_in_registry(tags, contract, event_contracts)?;
- validate_discriminator_parts(content, contract)?;
+ if admission_selected {
+ validate_admission_selected_contract_parts(kind, tags, contract)?;
+ }
+ validate_discriminator_parts(content, contract, admission_selected)?;
validate_custom_calendar_contract_parts(tags, contract)?;
validate_custom_knowledge_contract_parts(content, contract)?;
Ok(())
}
+fn validate_admission_selected_contract_parts(
+ kind: u32,
+ tags: &[Vec<String>],
+ contract: &EventContract,
+) -> Result<(), ContractValidationError> {
+ if contract.id == "radroots.social.deletion_request.v1"
+ && !tags.iter().any(|tag| {
+ tag.first()
+ .is_some_and(|name| matches!(name.as_str(), "e" | "a"))
+ })
+ {
+ return Err(ContractValidationError::ContractMatch {
+ error: ContractMatchError::UnsupportedShape(kind),
+ });
+ }
+ Ok(())
+}
+
fn validate_classified_listing_partition_parts(
tags: &[Vec<String>],
contract: &EventContract,
@@ -4939,9 +4988,10 @@ fn validate_custom_knowledge_contract_parts(
fn validate_discriminator_parts(
content: &str,
contract: &EventContract,
+ admission_selected: bool,
) -> Result<(), ContractValidationError> {
crate::require_invariant(
- !matches!(contract.discriminator, EventDiscriminator::AdmissionOnly),
+ admission_selected || !matches!(contract.discriminator, EventDiscriminator::AdmissionOnly),
&|| ContractValidationError::AdmissionRequired {
contract_id: contract.id,
},
diff --git a/crates/event/src/contract/registry_v7/tests.rs b/crates/event/src/contract/registry_v7/tests.rs
@@ -1088,6 +1088,27 @@ fn nip09_deletion_request_contract_is_typed_and_admission_only() {
contract_id: "radroots.social.deletion_request.v1",
})
);
+ let selected = unsigned_event(
+ KIND_DELETION_REQUEST,
+ vec![vec![
+ "e",
+ "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
+ ]],
+ "superseded",
+ );
+ assert_eq!(
+ validate_event_contract_for_admission(&selected, contract.id),
+ Ok(contract)
+ );
+ assert_eq!(
+ validate_event_contract_for_admission(
+ &unsigned_event(KIND_DELETION_REQUEST, vec![], "superseded"),
+ contract.id,
+ ),
+ Err(ContractValidationError::ContractMatch {
+ error: ContractMatchError::UnsupportedShape(KIND_DELETION_REQUEST),
+ })
+ );
}
#[test]
@@ -1165,13 +1186,13 @@ fn supports_content_field_discriminators() {
discriminator,
..base
};
- assert!(validate_discriminator_parts(r#"{"type":"proposal"}"#, &contract).is_ok());
+ assert!(validate_discriminator_parts(r#"{"type":"proposal"}"#, &contract, false).is_ok());
assert!(matches!(
- validate_discriminator_parts(r#"{"type":"decision"}"#, &contract),
+ validate_discriminator_parts(r#"{"type":"decision"}"#, &contract, false),
Err(ContractValidationError::ContentFieldMismatch { .. })
));
assert!(matches!(
- validate_discriminator_parts("{}", &contract),
+ validate_discriminator_parts("{}", &contract, false),
Err(ContractValidationError::MissingContentField { .. })
));
}
@@ -1179,7 +1200,7 @@ fn supports_content_field_discriminators() {
discriminator: EventDiscriminator::KindOnly,
..base
};
- assert!(validate_discriminator_parts("not-json", &contract).is_ok());
+ assert!(validate_discriminator_parts("not-json", &contract, false).is_ok());
}
#[test]
diff --git a/crates/event/src/verification.rs b/crates/event/src/verification.rs
@@ -9,7 +9,8 @@ use core::fmt;
use crate::{
contract::registry_v7::{
- ContractValidationError, EventContract, validate_event_contract_registry_v7,
+ ContractValidationError, EventContract, validate_event_contract_for_admission,
+ validate_event_contract_registry_v7,
},
envelope::EventEnvelope,
id::EventId,
@@ -137,6 +138,23 @@ impl SignatureVerifiedEvent {
contract,
})
}
+
+ /// Validates a contract selected explicitly by the admission boundary.
+ ///
+ /// This transition supports profiles such as root updates, replies,
+ /// comments, and deletion requests that intentionally cannot be selected
+ /// solely from their public Nostr wire shape.
+ pub fn validate_contract_for_admission(
+ self,
+ contract_id: &str,
+ ) -> Result<ContractValidatedEvent, Error> {
+ let contract = validate_event_contract_for_admission(&self.0, contract_id)
+ .map_err(Error::ContractValidation)?;
+ Ok(ContractValidatedEvent {
+ event: self,
+ contract,
+ })
+ }
}
/// A signature-verified event whose registry-selected contract is valid.
diff --git a/crates/storage/src/event.rs b/crates/storage/src/event.rs
@@ -12,6 +12,10 @@ use std::collections::BTreeSet;
use crate::{Error, status::EventStoreStatus};
+mod visibility;
+#[doc(hidden)]
+pub use visibility::{VisibilityEvaluation, VisibilityInput, evaluate_visibility};
+
/// Maximum events returned by one storage query.
pub const EVENT_QUERY_LIMIT_MAX: u16 = 1_000;
/// Maximum explicit event identifiers in one storage query.
@@ -401,6 +405,87 @@ pub struct StoredVisibleEvent {
event: SignedEvent,
}
+/// Deterministic digest of one complete visibility rebuild.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct VisibilityDigest([u8; 32]);
+
+impl VisibilityDigest {
+ /// Returns the canonical SHA-256 bytes for the rebuilt visibility state.
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ &self.0
+ }
+
+ pub(crate) const fn new(bytes: [u8; 32]) -> Self {
+ Self(bytes)
+ }
+}
+
+/// Complete deterministic result of rebuilding current event visibility.
+///
+/// Current heads remain listed even when their selected event is suppressed.
+/// This prevents a deleted head from resurrecting an older revision.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct VisibilitySnapshot {
+ generation: SourceGeneration,
+ current_heads: Vec<radroots_event::envelope::event_head::CurrentEventHead>,
+ deletion_request_ids: Vec<EventId>,
+ visible_event_ids: Vec<EventId>,
+ suppressed_event_ids: Vec<EventId>,
+ superseded_event_ids: Vec<EventId>,
+ digest: VisibilityDigest,
+}
+
+impl VisibilitySnapshot {
+ pub(crate) const fn new(
+ generation: SourceGeneration,
+ current_heads: Vec<radroots_event::envelope::event_head::CurrentEventHead>,
+ deletion_request_ids: Vec<EventId>,
+ visible_event_ids: Vec<EventId>,
+ suppressed_event_ids: Vec<EventId>,
+ superseded_event_ids: Vec<EventId>,
+ digest: VisibilityDigest,
+ ) -> Self {
+ Self {
+ generation,
+ current_heads,
+ deletion_request_ids,
+ visible_event_ids,
+ suppressed_event_ids,
+ superseded_event_ids,
+ digest,
+ }
+ }
+
+ pub const fn generation(&self) -> SourceGeneration {
+ self.generation
+ }
+
+ pub fn current_heads(&self) -> &[radroots_event::envelope::event_head::CurrentEventHead] {
+ self.current_heads.as_slice()
+ }
+
+ pub fn deletion_request_ids(&self) -> &[EventId] {
+ self.deletion_request_ids.as_slice()
+ }
+
+ pub fn visible_event_ids(&self) -> &[EventId] {
+ self.visible_event_ids.as_slice()
+ }
+
+ pub fn suppressed_event_ids(&self) -> &[EventId] {
+ self.suppressed_event_ids.as_slice()
+ }
+
+ pub fn superseded_event_ids(&self) -> &[EventId] {
+ self.superseded_event_ids.as_slice()
+ }
+
+ pub const fn digest(&self) -> VisibilityDigest {
+ self.digest
+ }
+}
+
impl StoredVisibleEvent {
pub const fn new(position: EventPosition, event: SignedEvent) -> Self {
Self { position, event }
@@ -535,6 +620,11 @@ pub trait EventStore: Send + Sync {
query: EventQuery,
) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>>;
+ /// Rebuilds current visibility from immutable retained event truth.
+ ///
+ /// Implementations must use the same reducer as [`Self::query_visible`].
+ fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>>;
+
/// Queries bounded provenance for one event.
fn query_provenance(
&self,
diff --git a/crates/storage/src/event/visibility.rs b/crates/storage/src/event/visibility.rs
@@ -0,0 +1,499 @@
+//! Deterministic current-visibility reduction over immutable event truth.
+
+use std::collections::{BTreeMap, BTreeSet};
+
+use radroots_event::{
+ EventId, SignedEvent,
+ envelope::event_head::{
+ CurrentEventHead, EventHeadCandidateResult, EventHeadCoordinate, EventHeadDecision,
+ event_head_candidate_for_nip01_event, select_event_head,
+ },
+ envelope::kind::KIND_DELETION_REQUEST,
+ id::Nip01Coordinate,
+};
+use sha2::{Digest, Sha256};
+
+use super::{
+ AdmissionStage, EventPosition, SourceGeneration, VisibilityDigest, VisibilitySnapshot,
+};
+use crate::Error;
+
+#[doc(hidden)]
+pub struct VisibilityInput<'a> {
+ position: EventPosition,
+ event: &'a SignedEvent,
+ stage: AdmissionStage,
+}
+
+impl<'a> VisibilityInput<'a> {
+ #[doc(hidden)]
+ pub const fn new(
+ position: EventPosition,
+ event: &'a SignedEvent,
+ stage: AdmissionStage,
+ ) -> Self {
+ Self {
+ position,
+ event,
+ stage,
+ }
+ }
+}
+
+#[doc(hidden)]
+pub struct VisibilityEvaluation {
+ snapshot: VisibilitySnapshot,
+ visible: BTreeSet<EventId>,
+}
+
+impl VisibilityEvaluation {
+ #[doc(hidden)]
+ pub fn is_visible(&self, event_id: &EventId) -> bool {
+ self.visible.contains(event_id)
+ }
+
+ #[doc(hidden)]
+ pub const fn snapshot(&self) -> &VisibilitySnapshot {
+ &self.snapshot
+ }
+
+ #[doc(hidden)]
+ pub fn into_snapshot(self) -> VisibilitySnapshot {
+ self.snapshot
+ }
+}
+
+#[derive(Clone)]
+struct Candidate<'a> {
+ position: EventPosition,
+ event: &'a SignedEvent,
+ coordinate: Option<EventHeadCoordinate>,
+ ephemeral: bool,
+}
+
+struct DeletionRequest {
+ request_id: EventId,
+ author: [u8; 32],
+ created_at: u64,
+ event_targets: BTreeSet<EventId>,
+ address_targets: BTreeSet<Nip01Coordinate>,
+}
+
+#[doc(hidden)]
+pub fn evaluate_visibility<'a>(
+ generation: SourceGeneration,
+ inputs: impl IntoIterator<Item = VisibilityInput<'a>>,
+) -> Result<VisibilityEvaluation, Error> {
+ let mut candidates = Vec::new();
+ let mut deletion_requests = Vec::new();
+ let mut heads = BTreeMap::<EventHeadCoordinate, CurrentEventHead>::new();
+
+ for input in inputs {
+ if input.position.generation() != generation {
+ return Err(Error::CorruptStoredEvent);
+ }
+ if input.stage != AdmissionStage::Visible {
+ continue;
+ }
+ let envelope = input.event.envelope();
+ if envelope.kind_u32() == KIND_DELETION_REQUEST {
+ deletion_requests.push(parse_deletion_request(input.event)?);
+ }
+ let (coordinate, ephemeral) = match event_head_candidate_for_nip01_event(envelope) {
+ EventHeadCandidateResult::Candidate(candidate) => {
+ let coordinate = candidate.coordinate.clone();
+ match select_event_head(candidate, heads.get(&coordinate)) {
+ EventHeadDecision::Applied(head) => {
+ heads.insert(coordinate.clone(), head);
+ }
+ EventHeadDecision::SkippedDuplicate
+ | EventHeadDecision::SkippedOlder
+ | EventHeadDecision::SkippedSameTimestampHigherEventId => {}
+ EventHeadDecision::CoordinateMismatch => {
+ return Err(Error::CorruptStoredEvent);
+ }
+ }
+ (Some(coordinate), false)
+ }
+ EventHeadCandidateResult::NotHeadSelected => (None, false),
+ EventHeadCandidateResult::NotPersisted => (None, true),
+ EventHeadCandidateResult::Malformed(_) => return Err(Error::CorruptStoredEvent),
+ };
+ candidates.push(Candidate {
+ position: input.position,
+ event: input.event,
+ coordinate,
+ ephemeral,
+ });
+ }
+
+ candidates.sort_by_key(|candidate| candidate.position.sequence());
+ deletion_requests.sort_by_key(|request| request.request_id);
+
+ let mut visible = BTreeSet::new();
+ let mut suppressed = BTreeSet::new();
+ let mut superseded = BTreeSet::new();
+ for candidate in candidates {
+ let event_id = *candidate.event.id();
+ if candidate.ephemeral {
+ suppressed.insert(event_id);
+ continue;
+ }
+ if candidate.coordinate.as_ref().is_some_and(|coordinate| {
+ heads
+ .get(coordinate)
+ .is_none_or(|head| head.event_id != event_id)
+ }) {
+ superseded.insert(event_id);
+ continue;
+ }
+ if is_suppressed(candidate.event, &deletion_requests) {
+ suppressed.insert(event_id);
+ } else {
+ visible.insert(event_id);
+ }
+ }
+
+ let current_heads = heads.into_values().collect::<Vec<_>>();
+ let deletion_request_ids = deletion_requests
+ .iter()
+ .map(|request| request.request_id)
+ .collect::<Vec<_>>();
+ let visible_event_ids = visible.iter().copied().collect::<Vec<_>>();
+ let suppressed_event_ids = suppressed.into_iter().collect::<Vec<_>>();
+ let superseded_event_ids = superseded.into_iter().collect::<Vec<_>>();
+ let digest = visibility_digest(
+ generation,
+ current_heads.as_slice(),
+ deletion_request_ids.as_slice(),
+ visible_event_ids.as_slice(),
+ suppressed_event_ids.as_slice(),
+ superseded_event_ids.as_slice(),
+ )?;
+ Ok(VisibilityEvaluation {
+ snapshot: VisibilitySnapshot::new(
+ generation,
+ current_heads,
+ deletion_request_ids,
+ visible_event_ids,
+ suppressed_event_ids,
+ superseded_event_ids,
+ digest,
+ ),
+ visible,
+ })
+}
+
+fn parse_deletion_request(event: &SignedEvent) -> Result<DeletionRequest, Error> {
+ let mut event_targets = BTreeSet::new();
+ let mut address_targets = BTreeSet::new();
+ for tag in event.envelope().tag_slices() {
+ match tag.as_slice().first().map(String::as_str) {
+ Some("e") => {
+ let value = tag.as_slice().get(1).ok_or(Error::CorruptStoredEvent)?;
+ event_targets.insert(EventId::parse(value).map_err(|_| Error::CorruptStoredEvent)?);
+ }
+ Some("a") => {
+ let value = tag.as_slice().get(1).ok_or(Error::CorruptStoredEvent)?;
+ address_targets
+ .insert(Nip01Coordinate::parse(value).map_err(|_| Error::CorruptStoredEvent)?);
+ }
+ _ => {}
+ }
+ }
+ if event_targets.is_empty() && address_targets.is_empty() {
+ return Err(Error::CorruptStoredEvent);
+ }
+ Ok(DeletionRequest {
+ request_id: *event.id(),
+ author: *event.envelope().author().as_bytes(),
+ created_at: event.envelope().created_at_u64(),
+ event_targets,
+ address_targets,
+ })
+}
+
+fn is_suppressed(target: &SignedEvent, requests: &[DeletionRequest]) -> bool {
+ let target_event = target.envelope();
+ if target_event.kind_u32() == KIND_DELETION_REQUEST {
+ return false;
+ }
+ let coordinate = event_coordinate(target_event);
+ requests.iter().any(|request| {
+ let request_event_id_match = request.event_targets.contains(target_event.id());
+ let request_coordinate_match = coordinate.as_ref().is_some_and(|coordinate| {
+ request.address_targets.contains(coordinate)
+ && target_event.created_at_u64() <= request.created_at
+ });
+ (request_event_id_match || request_coordinate_match)
+ && request.author == *target_event.author().as_bytes()
+ })
+}
+
+fn event_coordinate(event: &radroots_event::Event) -> Option<Nip01Coordinate> {
+ let kind = event.kind_u32();
+ let identifier = match event.kind_class() {
+ radroots_event::envelope::EventKindClass::Replaceable => "",
+ radroots_event::envelope::EventKindClass::Addressable => event
+ .tag_slices()
+ .iter()
+ .find(|tag| tag.as_slice().first().is_some_and(|name| name == "d"))?
+ .as_slice()
+ .get(1)?
+ .as_str(),
+ radroots_event::envelope::EventKindClass::Regular
+ | radroots_event::envelope::EventKindClass::Ephemeral => return None,
+ };
+ Nip01Coordinate::parse(format!("{kind}:{}:{identifier}", event.author())).ok()
+}
+
+fn visibility_digest(
+ generation: SourceGeneration,
+ heads: &[CurrentEventHead],
+ deletion_requests: &[EventId],
+ visible: &[EventId],
+ suppressed: &[EventId],
+ superseded: &[EventId],
+) -> Result<VisibilityDigest, Error> {
+ let mut digest = Sha256::new();
+ digest.update(b"radroots.storage.visibility.v1\0");
+ digest.update(generation.as_bytes());
+ for head in heads {
+ digest.update(b"head\0");
+ match &head.coordinate {
+ EventHeadCoordinate::Replaceable { kind, pubkey } => {
+ digest.update(b"replaceable\0");
+ digest.update(kind.to_be_bytes());
+ digest.update(pubkey.as_bytes());
+ }
+ EventHeadCoordinate::Addressable {
+ kind,
+ pubkey,
+ d_tag,
+ } => {
+ digest.update(b"addressable\0");
+ digest.update(kind.to_be_bytes());
+ digest.update(pubkey.as_bytes());
+ let d_tag_length =
+ u64::try_from(d_tag.len()).map_err(|_| Error::CorruptStoredEvent)?;
+ digest.update(d_tag_length.to_be_bytes());
+ digest.update(d_tag.as_bytes());
+ }
+ }
+ digest.update(head.event_id.as_bytes());
+ digest.update(head.created_at.to_be_bytes());
+ }
+ update_event_ids(&mut digest, b"deletion\0", deletion_requests);
+ update_event_ids(&mut digest, b"visible\0", visible);
+ update_event_ids(&mut digest, b"suppressed\0", suppressed);
+ update_event_ids(&mut digest, b"superseded\0", superseded);
+ Ok(VisibilityDigest::new(digest.finalize().into()))
+}
+
+fn update_event_ids(digest: &mut Sha256, prefix: &[u8], event_ids: &[EventId]) {
+ for event_id in event_ids {
+ digest.update(prefix);
+ digest.update(event_id.as_bytes());
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use radroots_event::wire::Nip01EventWire;
+
+ const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+ const OTHER_AUTHOR: &str = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798";
+
+ fn signed_event(
+ author: &str,
+ created_at: u64,
+ kind: u32,
+ tags: Vec<Vec<&str>>,
+ content: &str,
+ ) -> SignedEvent {
+ let tags = tags
+ .into_iter()
+ .map(|tag| tag.into_iter().map(str::to_owned).collect::<Vec<_>>())
+ .collect::<Vec<_>>();
+ let mut wire = Nip01EventWire {
+ id: "0".repeat(64),
+ pubkey: author.to_owned(),
+ created_at,
+ kind,
+ tags,
+ content: content.to_owned(),
+ sig: "42".repeat(64),
+ extra: Default::default(),
+ };
+ wire.id = wire.computed_event_id().expect("event id").to_hex();
+ let raw_json = serde_json::json!({
+ "id": &wire.id,
+ "pubkey": &wire.pubkey,
+ "created_at": wire.created_at,
+ "kind": wire.kind,
+ "tags": &wire.tags,
+ "content": &wire.content,
+ "sig": &wire.sig,
+ })
+ .to_string();
+ SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event")
+ }
+
+ fn evaluate<'a>(
+ events: impl IntoIterator<Item = (&'a SignedEvent, AdmissionStage)>,
+ ) -> VisibilityEvaluation {
+ let generation = SourceGeneration::new([7; 32]).expect("generation");
+ evaluate_visibility(
+ generation,
+ events
+ .into_iter()
+ .enumerate()
+ .map(|(index, (event, stage))| {
+ let sequence = u64::try_from(index)
+ .expect("index")
+ .checked_add(1)
+ .and_then(|value| super::super::EventSequence::new(value).ok())
+ .expect("sequence");
+ VisibilityInput::new(EventPosition::new(generation, sequence), event, stage)
+ }),
+ )
+ .expect("visibility")
+ }
+
+ #[test]
+ fn selected_head_is_stable_and_a_deleted_head_does_not_resurrect() {
+ let old = signed_event(AUTHOR, 10, 0, vec![], r#"{"name":"old"}"#);
+ let current = signed_event(AUTHOR, 20, 0, vec![], r#"{"name":"current"}"#);
+ let deletion = signed_event(
+ AUTHOR,
+ 30,
+ KIND_DELETION_REQUEST,
+ vec![vec!["e", current.id().to_hex().as_str()]],
+ "",
+ );
+ let evaluation = evaluate([
+ (&old, AdmissionStage::Visible),
+ (¤t, AdmissionStage::Visible),
+ (&deletion, AdmissionStage::Visible),
+ ]);
+
+ assert!(!evaluation.is_visible(old.id()));
+ assert!(!evaluation.is_visible(current.id()));
+ assert!(evaluation.is_visible(deletion.id()));
+ assert_eq!(evaluation.snapshot().superseded_event_ids(), &[*old.id()]);
+ assert_eq!(
+ evaluation.snapshot().suppressed_event_ids(),
+ &[*current.id()]
+ );
+ assert_eq!(
+ evaluation.snapshot().current_heads()[0].event_id,
+ *current.id()
+ );
+ }
+
+ #[test]
+ fn address_cutoff_allows_a_later_replacement_and_wrong_author_is_ineffective() {
+ let old = signed_event(AUTHOR, 10, 30_023, vec![vec!["d", "farm-update"]], "old");
+ let wrong_author_deletion = signed_event(
+ OTHER_AUTHOR,
+ 15,
+ KIND_DELETION_REQUEST,
+ vec![vec!["a", format!("30023:{AUTHOR}:farm-update").as_str()]],
+ "",
+ );
+ let cutoff = signed_event(
+ AUTHOR,
+ 20,
+ KIND_DELETION_REQUEST,
+ vec![vec!["a", format!("30023:{AUTHOR}:farm-update").as_str()]],
+ "",
+ );
+ let later = signed_event(AUTHOR, 30, 30_023, vec![vec!["d", "farm-update"]], "later");
+ let evaluation = evaluate([
+ (&old, AdmissionStage::Visible),
+ (&wrong_author_deletion, AdmissionStage::Visible),
+ (&cutoff, AdmissionStage::Visible),
+ (&later, AdmissionStage::Visible),
+ ]);
+
+ assert!(evaluation.is_visible(later.id()));
+ assert!(evaluation.is_visible(wrong_author_deletion.id()));
+ assert!(evaluation.is_visible(cutoff.id()));
+ assert!(!evaluation.is_visible(old.id()));
+ assert_eq!(
+ evaluation.snapshot().current_heads()[0].event_id,
+ *later.id()
+ );
+ }
+
+ #[test]
+ fn rebuild_is_order_independent_and_ignores_nonvisible_admissions() {
+ let visible = signed_event(AUTHOR, 10, 1, vec![], "visible");
+ let excluded = signed_event(AUTHOR, 20, 1, vec![], "excluded");
+ let ephemeral = signed_event(AUTHOR, 30, 20_000, vec![], "ephemeral");
+ let first = evaluate([
+ (&visible, AdmissionStage::Visible),
+ (&excluded, AdmissionStage::Verified),
+ (&ephemeral, AdmissionStage::Visible),
+ ]);
+ let second = evaluate([
+ (&ephemeral, AdmissionStage::Visible),
+ (&excluded, AdmissionStage::Verified),
+ (&visible, AdmissionStage::Visible),
+ ]);
+
+ assert!(first.is_visible(visible.id()));
+ assert!(!first.is_visible(excluded.id()));
+ assert!(!first.is_visible(ephemeral.id()));
+ assert_eq!(first.snapshot().digest(), second.snapshot().digest());
+ assert_eq!(first.snapshot(), second.snapshot());
+ }
+
+ #[test]
+ fn deletion_requests_without_targets_fail_closed() {
+ let deletion = signed_event(AUTHOR, 10, KIND_DELETION_REQUEST, vec![], "");
+ let generation = SourceGeneration::new([7; 32]).expect("generation");
+ let position = EventPosition::new(
+ generation,
+ super::super::EventSequence::new(1).expect("sequence"),
+ );
+
+ let error = evaluate_visibility(
+ generation,
+ [VisibilityInput::new(
+ position,
+ &deletion,
+ AdmissionStage::Visible,
+ )],
+ )
+ .err()
+ .expect("targetless deletion must fail");
+
+ assert_eq!(error, Error::CorruptStoredEvent);
+ }
+
+ #[test]
+ fn events_from_another_generation_fail_closed() {
+ let event = signed_event(AUTHOR, 10, 1, vec![], "event");
+ let requested_generation = SourceGeneration::new([7; 32]).expect("generation");
+ let stored_generation = SourceGeneration::new([8; 32]).expect("generation");
+ let position = EventPosition::new(
+ stored_generation,
+ super::super::EventSequence::new(1).expect("sequence"),
+ );
+
+ let error = evaluate_visibility(
+ requested_generation,
+ [VisibilityInput::new(
+ position,
+ &event,
+ AdmissionStage::Visible,
+ )],
+ )
+ .err()
+ .expect("cross-generation visibility input must fail");
+
+ assert_eq!(error, Error::CorruptStoredEvent);
+ }
+}
diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs
@@ -25,6 +25,7 @@ use crate::{
AdmissionDisposition, AdmissionReceipt, AdmissionStage, EventAdmission, EventPage,
EventPosition, EventQuery, EventQueryBounds, EventSequence, SourceGeneration,
StoredEventProvenance, StoredRawEvent, StoredVerifiedEvent, StoredVisibleEvent,
+ VisibilityEvaluation, VisibilityInput, VisibilitySnapshot, evaluate_visibility,
},
journal::{
IdempotencyKey, JournalStage, JournalTransition, OperationInstanceId, OperationRecord,
@@ -130,7 +131,12 @@ impl MemoryStorage {
self.state.lock().map_err(|_| Error::BackendUnavailable)
}
- fn selected(&self, state: &State, query: &EventQuery) -> Result<Vec<EventEntry>, Error> {
+ fn selected(
+ &self,
+ state: &State,
+ query: &EventQuery,
+ mut eligible: impl FnMut(&EventEntry) -> bool,
+ ) -> Result<(Vec<EventEntry>, Option<EventPosition>), Error> {
if query
.bounds()
.cursor()
@@ -142,15 +148,37 @@ impl MemoryStorage {
.bounds()
.cursor()
.map_or(0, |cursor| cursor.sequence().get());
- Ok(state
+ let mut selected = state
.events
.iter()
.filter(|entry| {
- entry.position.sequence().get() > after && query.selects(entry.admission.event_id())
+ entry.position.sequence().get() > after
+ && query.selects(entry.admission.event_id())
+ && eligible(entry)
})
- .take(usize::from(query.bounds().limit()))
+ .take(usize::from(query.bounds().limit()) + 1)
.cloned()
- .collect())
+ .collect::<Vec<_>>();
+ let next = if selected.len() > usize::from(query.bounds().limit()) {
+ selected.truncate(usize::from(query.bounds().limit()));
+ selected.last().map(|entry| entry.position)
+ } else {
+ None
+ };
+ Ok((selected, next))
+ }
+
+ fn visibility_locked(&self, state: &State) -> Result<VisibilityEvaluation, Error> {
+ evaluate_visibility(
+ self.generation,
+ state.events.iter().map(|entry| {
+ VisibilityInput::new(
+ entry.position,
+ entry.admission.event(),
+ entry.admission.stage(),
+ )
+ }),
+ )
}
fn admit_locked(
@@ -374,11 +402,10 @@ impl EventStore for MemoryStorage {
)
.map_err(|_| Error::CorruptStoredEvent)?;
let visible = u64::try_from(
- state
- .events
- .iter()
- .filter(|entry| entry.admission.stage() == AdmissionStage::Visible)
- .count(),
+ self.visibility_locked(&state)?
+ .snapshot()
+ .visible_event_ids()
+ .len(),
)
.map_err(|_| Error::CorruptStoredEvent)?;
EventStoreStatus::new(
@@ -405,8 +432,8 @@ impl EventStore for MemoryStorage {
) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>> {
Box::pin(async move {
let state = self.state()?;
- let items = self
- .selected(&state, &query)?
+ let (entries, next) = self.selected(&state, &query, |_| true)?;
+ let items = entries
.into_iter()
.map(|entry| {
StoredRawEvent::new(
@@ -416,7 +443,7 @@ impl EventStore for MemoryStorage {
)
})
.collect();
- EventPage::new(self.generation, items, None, query.bounds())
+ EventPage::new(self.generation, items, next, query.bounds())
})
}
@@ -426,16 +453,16 @@ impl EventStore for MemoryStorage {
) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>> {
Box::pin(async move {
let state = self.state()?;
- let items = self
- .selected(&state, &query)?
+ let (entries, next) = self.selected(&state, &query, |entry| {
+ entry.admission.stage() >= AdmissionStage::Verified
+ })?;
+ let items = entries
.into_iter()
- .filter_map(|entry| {
- (entry.admission.stage() >= AdmissionStage::Verified).then(|| {
- StoredVerifiedEvent::new(entry.position, entry.admission.event().clone())
- })
+ .map(|entry| {
+ StoredVerifiedEvent::new(entry.position, entry.admission.event().clone())
})
.collect();
- EventPage::new(self.generation, items, None, query.bounds())
+ EventPage::new(self.generation, items, next, query.bounds())
})
}
@@ -445,16 +472,24 @@ impl EventStore for MemoryStorage {
) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> {
Box::pin(async move {
let state = self.state()?;
- let items = self
- .selected(&state, &query)?
+ let visibility = self.visibility_locked(&state)?;
+ let (entries, next) = self.selected(&state, &query, |entry| {
+ visibility.is_visible(entry.admission.event_id())
+ })?;
+ let items = entries
.into_iter()
- .filter_map(|entry| {
- (entry.admission.stage() == AdmissionStage::Visible).then(|| {
- StoredVisibleEvent::new(entry.position, entry.admission.event().clone())
- })
+ .map(|entry| {
+ StoredVisibleEvent::new(entry.position, entry.admission.event().clone())
})
.collect();
- EventPage::new(self.generation, items, None, query.bounds())
+ EventPage::new(self.generation, items, next, query.bounds())
+ })
+ }
+
+ fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>> {
+ Box::pin(async move {
+ let state = self.state()?;
+ Ok(self.visibility_locked(&state)?.into_snapshot())
})
}
diff --git a/crates/storage/tests/event_store.rs b/crates/storage/tests/event_store.rs
@@ -10,6 +10,7 @@ use radroots_storage::{
AdmissionDisposition, AdmissionReceipt, AdmissionStage, EventAdmission, EventPage,
EventPosition, EventQuery, EventQueryBounds, EventSequence, SourceGeneration,
StoredEventProvenance, StoredRawEvent, StoredVerifiedEvent, StoredVisibleEvent,
+ VisibilityInput, VisibilitySnapshot, evaluate_visibility,
},
status::{EventStoreHealth, EventStoreMode, EventStoreStatus},
};
@@ -191,6 +192,23 @@ impl EventStore for MemoryEventStore {
})
}
+ fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>> {
+ Box::pin(async move {
+ let entries = self.entries.lock().expect("test store lock");
+ evaluate_visibility(
+ self.generation,
+ entries.iter().map(|entry| {
+ VisibilityInput::new(
+ entry.position,
+ entry.admission.event(),
+ entry.admission.stage(),
+ )
+ }),
+ )
+ .map(|evaluation| evaluation.into_snapshot())
+ })
+ }
+
fn query_provenance(
&self,
event_id: EventId,
diff --git a/crates/storage/tests/memory.rs b/crates/storage/tests/memory.rs
@@ -1,7 +1,12 @@
#![cfg(feature = "memory")]
use futures_executor::block_on;
-use radroots_event::{SignedEvent, wire::Nip01EventWire};
+use radroots_event::{
+ SignedEvent,
+ admission::{AdmissionPolicy, RawEvent, SignatureVerifier, VisibilityPolicy},
+ envelope::EventEnvelope,
+ wire::Nip01EventWire,
+};
use radroots_protocol::runtime::v1::OperationId;
use radroots_storage::{
Error, EventStore, Journal, Outbox, ProjectionStore,
@@ -41,13 +46,22 @@ use radroots_transport::{
};
fn signed_event() -> SignedEvent {
+ signed_event_with(1_800_000_100, 0, vec![], "memory-backend")
+}
+
+fn signed_event_with(
+ created_at: u64,
+ kind: u32,
+ tags: Vec<Vec<String>>,
+ content: &str,
+) -> SignedEvent {
let mut wire = Nip01EventWire {
id: "0".repeat(64),
pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
- created_at: 1_800_000_100,
- kind: 0,
- tags: vec![],
- content: "memory-backend".to_owned(),
+ created_at,
+ kind,
+ tags,
+ content: content.to_owned(),
sig: "42".repeat(64),
extra: Default::default(),
};
@@ -65,6 +79,44 @@ fn signed_event() -> SignedEvent {
SignedEvent::from_wire_verified_id(wire, raw).expect("signed event")
}
+struct Allow;
+
+impl SignatureVerifier for Allow {
+ fn verify_signature(&self, _event: &EventEnvelope) -> Result<(), radroots_event::Error> {
+ Ok(())
+ }
+}
+
+impl AdmissionPolicy for Allow {
+ type Error = core::convert::Infallible;
+
+ fn policy_id(&self) -> &'static str {
+ "test.storage-memory.admission.v1"
+ }
+
+ fn admit(
+ &self,
+ _event: &radroots_event::admission::ContractValidatedEvent,
+ ) -> Result<(), Self::Error> {
+ Ok(())
+ }
+}
+
+impl VisibilityPolicy for Allow {
+ type Error = core::convert::Infallible;
+
+ fn policy_id(&self) -> &'static str {
+ "test.storage-memory.visibility.v1"
+ }
+
+ fn make_visible(
+ &self,
+ _event: &radroots_event::admission::AdmittedEvent,
+ ) -> Result<(), Self::Error> {
+ Ok(())
+ }
+}
+
fn admission(event: SignedEvent, at: u64) -> EventAdmission {
let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
let provenance = EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), at)
@@ -72,6 +124,95 @@ fn admission(event: SignedEvent, at: u64) -> EventAdmission {
EventAdmission::raw(ObservedEvent::new(event, provenance))
}
+fn visible_admission(event: SignedEvent, at: u64) -> EventAdmission {
+ let verified = RawEvent::new(event.envelope().clone())
+ .verify_id()
+ .expect("event id")
+ .verify_signature(&Allow)
+ .expect("signature");
+ let validated = if event.envelope().kind_u32() == 5 {
+ verified
+ .validate_contract_for_admission("radroots.social.deletion_request.v1")
+ .expect("admission-selected contract")
+ } else {
+ verified.validate_contract().expect("contract")
+ };
+ let visible = validated
+ .admit_with(&Allow)
+ .expect("admission")
+ .make_visible_with(&Allow)
+ .expect("visibility");
+ let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
+ let provenance = EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), at)
+ .expect("provenance");
+ EventAdmission::visible(ObservedEvent::new(event, provenance), visible)
+ .expect("visible admission")
+}
+
+#[test]
+fn memory_visibility_rebuild_is_current_delete_aware_and_atomic_parity_safe() {
+ let generation = SourceGeneration::new([17; 32]).expect("generation");
+ let direct = MemoryStorage::new(generation);
+ let atomic_store = MemoryStorage::new(generation);
+ let old = signed_event_with(
+ 1_800_000_100,
+ 0,
+ vec![],
+ r#"{"display_name":"Old Farm","bot":false}"#,
+ );
+ let current = signed_event_with(
+ 1_800_000_200,
+ 0,
+ vec![],
+ r#"{"display_name":"Current Farm","bot":false}"#,
+ );
+ let deletion = signed_event_with(
+ 1_800_000_300,
+ 5,
+ vec![vec!["e".to_owned(), current.id().to_hex()]],
+ "retired profile",
+ );
+ let admissions = [
+ visible_admission(old.clone(), 100),
+ visible_admission(current.clone(), 200),
+ visible_admission(deletion.clone(), 300),
+ ];
+
+ for admission in admissions.clone() {
+ block_on(direct.admit(admission)).expect("direct admission");
+ }
+ for (index, admission) in admissions.into_iter().enumerate() {
+ let identity = u8::try_from(index).expect("commit identity") + 20;
+ block_on(atomic_store.commit(atomic(
+ identity,
+ identity,
+ AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission, None))),
+ )))
+ .expect("atomic admission");
+ }
+
+ let snapshot = block_on(direct.rebuild_visibility()).expect("visibility rebuild");
+ let atomic_snapshot =
+ block_on(atomic_store.rebuild_visibility()).expect("atomic visibility rebuild");
+ assert_eq!(snapshot, atomic_snapshot);
+ assert_eq!(snapshot.current_heads()[0].event_id, *current.id());
+ assert_eq!(snapshot.visible_event_ids(), &[*deletion.id()]);
+ assert_eq!(snapshot.suppressed_event_ids(), &[*current.id()]);
+ assert_eq!(snapshot.superseded_event_ids(), &[*old.id()]);
+ let page = block_on(direct.query_visible(EventQuery::all(
+ EventQueryBounds::first(10).expect("bounds"),
+ )))
+ .expect("visible page");
+ assert_eq!(page.items().len(), 1);
+ assert_eq!(page.items()[0].event().id(), deletion.id());
+ assert_eq!(
+ block_on(EventStore::status(&direct))
+ .expect("status")
+ .visible_events(),
+ 1
+ );
+}
+
fn prepare(instance: OperationInstanceId) -> PrepareOperation {
PrepareOperation::new(
instance,
diff --git a/crates/storage_sqlite/src/event/mod.rs b/crates/storage_sqlite/src/event/mod.rs
@@ -1,3 +1,4 @@
+use radroots_event::SignedEvent;
use radroots_event_codec::Codec;
use radroots_storage::{
Error, EventStore,
@@ -5,7 +6,8 @@ use radroots_storage::{
AdmissionDisposition, AdmissionReceipt, AdmissionStage, BoxFuture, EventAdmission,
EventCursor, EventId, EventPage, EventPosition, EventQuery, EventQueryBounds,
EventSequence, SourceGeneration, StoredEventProvenance, StoredRawEvent,
- StoredVerifiedEvent, StoredVisibleEvent,
+ StoredVerifiedEvent, StoredVisibleEvent, VisibilityEvaluation, VisibilityInput,
+ VisibilitySnapshot, evaluate_visibility,
},
status::{EventStoreHealth, EventStoreMode, EventStoreStatus},
};
@@ -33,6 +35,7 @@ pub struct SqliteStorage {
struct StoredEventRow {
position: EventPosition,
raw_json: String,
+ event: SignedEvent,
stage: AdmissionStage,
}
@@ -130,8 +133,10 @@ impl SqliteStorage {
let mut builder = QueryBuilder::<Sqlite>::new(
"SELECT source_generation, source_sequence, signed_event, admission_stage, \
admitted_contract_id, admitted_registry_version \
- FROM radroots_runtime_events WHERE source_sequence > ",
+ FROM radroots_runtime_events WHERE source_generation = ",
);
+ builder.push_bind(self.generation.as_bytes().as_slice());
+ builder.push(" AND source_sequence > ");
builder.push_bind(i64_from_u64(after)?);
if minimum_stage == AdmissionStage::Verified {
builder.push(" AND admission_stage IN ('verified', 'visible')");
@@ -167,6 +172,29 @@ impl SqliteStorage {
Ok((decoded, next))
}
+ async fn current_event_rows(&self) -> Result<Vec<StoredEventRow>, Error> {
+ let rows = sqlx::query(
+ "SELECT source_generation, source_sequence, signed_event, admission_stage,
+ admitted_contract_id, admitted_registry_version
+ FROM radroots_runtime_events
+ WHERE source_generation = ?
+ ORDER BY source_sequence",
+ )
+ .bind(self.generation.as_bytes().as_slice())
+ .fetch_all(&self.pool)
+ .await
+ .map_err(map_backend)?;
+ rows.iter().map(|row| self.decode_event_row(row)).collect()
+ }
+
+ fn visibility_for_rows(&self, rows: &[StoredEventRow]) -> Result<VisibilityEvaluation, Error> {
+ evaluate_visibility(
+ self.generation,
+ rows.iter()
+ .map(|row| VisibilityInput::new(row.position, &row.event, row.stage)),
+ )
+ }
+
fn validate_cursor(&self, query: &EventQuery) -> Result<(), Error> {
if query
.bounds()
@@ -198,10 +226,15 @@ impl SqliteStorage {
.map(|value| u32::try_from(value).map_err(|_| Error::CorruptStoredEvent))
.transpose()?;
if contract_id.is_some() != registry_version.is_some()
+ || registry_version.is_some_and(|version| {
+ version != radroots_event::contract::RegistryVersion::CURRENT.get()
+ })
|| contract_id.as_deref().is_some_and(|stored| {
- contract_metadata(&event).is_none_or(|(contract, version)| {
- stored != contract || Some(version) != registry_version
- })
+ radroots_event::contract::registry_v7::validate_event_contract_for_admission(
+ event.envelope(),
+ stored,
+ )
+ .is_err()
})
{
return Err(Error::CorruptStoredEvent);
@@ -209,6 +242,7 @@ impl SqliteStorage {
Ok(StoredEventRow {
position: EventPosition::new(generation, sequence),
raw_json,
+ event,
stage,
})
}
@@ -252,7 +286,7 @@ impl SqliteStorage {
let (position, disposition) = if let Some(row) = existing {
let stored = self.decode_event_row(&row)?;
- let metadata = contract_metadata(admission.event());
+ let metadata = admission_contract_metadata(&admission);
if stored.raw_json.as_bytes() != admission.event().raw_json().as_bytes() {
return Err(Error::EventConflict);
}
@@ -281,7 +315,7 @@ impl SqliteStorage {
};
(stored.position, disposition)
} else {
- let metadata = contract_metadata(admission.event());
+ let metadata = admission_contract_metadata(&admission);
let next = sqlx::query_scalar::<_, i64>(
"UPDATE radroots_runtime_source_generations
SET sequence_head = sequence_head + 1
@@ -329,15 +363,13 @@ impl SqliteStorage {
}
}
-fn contract_metadata(event: &radroots_event::SignedEvent) -> Option<(&'static str, u32)> {
- radroots_event::contract::registry_v7::validate_event_contract_registry_v7(event.envelope())
- .ok()
- .map(|contract| {
- (
- contract.id,
- radroots_event::contract::RegistryVersion::CURRENT.get(),
- )
- })
+fn admission_contract_metadata(admission: &EventAdmission) -> Option<(&'static str, u32)> {
+ admission.visible_event().map(|event| {
+ (
+ event.admitted_event().validated_event().contract_id(),
+ radroots_event::contract::RegistryVersion::CURRENT.get(),
+ )
+ })
}
#[cfg_attr(coverage_nightly, coverage(off))]
@@ -348,22 +380,28 @@ impl EventStore for SqliteStorage {
"SELECT
COUNT(*) AS raw_events,
COALESCE(SUM(CASE WHEN admission_stage IN ('verified', 'visible') THEN 1 ELSE 0 END), 0)
- AS verified_events,
- COALESCE(SUM(CASE WHEN admission_stage = 'visible' THEN 1 ELSE 0 END), 0)
- AS visible_events
+ AS verified_events
FROM radroots_runtime_events WHERE source_generation = ?",
)
.bind(self.generation.as_bytes().as_slice())
.fetch_one(&self.pool)
.await
.map_err(map_backend)?;
+ let visibility_rows = self.current_event_rows().await?;
+ let visible_events = u64::try_from(
+ self.visibility_for_rows(&visibility_rows)?
+ .snapshot()
+ .visible_event_ids()
+ .len(),
+ )
+ .map_err(|_| Error::CorruptStoredEvent)?;
EventStoreStatus::new(
self.generation,
self.mode,
EventStoreHealth::Available,
u64_from_i64(row.try_get("raw_events").map_err(map_corrupt)?)?,
u64_from_i64(row.try_get("verified_events").map_err(map_corrupt)?)?,
- u64_from_i64(row.try_get("visible_events").map_err(map_corrupt)?)?,
+ visible_events,
)
})
}
@@ -392,12 +430,8 @@ impl EventStore for SqliteStorage {
let (rows, next) = self.selected(&query, AdmissionStage::Raw).await?;
let items = rows
.into_iter()
- .map(|row| {
- let event = Codec::decode_signed_event(row.raw_json.as_str())
- .map_err(|_| Error::CorruptStoredEvent)?;
- Ok(StoredRawEvent::new(row.position, event, row.stage))
- })
- .collect::<Result<Vec<_>, Error>>()?;
+ .map(|row| StoredRawEvent::new(row.position, row.event, row.stage))
+ .collect();
EventPage::new(self.generation, items, next, query.bounds())
})
}
@@ -410,12 +444,8 @@ impl EventStore for SqliteStorage {
let (rows, next) = self.selected(&query, AdmissionStage::Verified).await?;
let items = rows
.into_iter()
- .map(|row| {
- let event = Codec::decode_signed_event(row.raw_json.as_str())
- .map_err(|_| Error::CorruptStoredEvent)?;
- Ok(StoredVerifiedEvent::new(row.position, event))
- })
- .collect::<Result<Vec<_>, Error>>()?;
+ .map(|row| StoredVerifiedEvent::new(row.position, row.event))
+ .collect();
EventPage::new(self.generation, items, next, query.bounds())
})
}
@@ -425,19 +455,43 @@ impl EventStore for SqliteStorage {
query: EventQuery,
) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> {
Box::pin(async move {
- let (rows, next) = self.selected(&query, AdmissionStage::Visible).await?;
- let items = rows
+ self.validate_cursor(&query)?;
+ let after = query
+ .bounds()
+ .cursor()
+ .map_or(0, |cursor| cursor.sequence().get());
+ let rows = self.current_event_rows().await?;
+ let visibility = self.visibility_for_rows(&rows)?;
+ let mut selected = rows
.into_iter()
- .map(|row| {
- let event = Codec::decode_signed_event(row.raw_json.as_str())
- .map_err(|_| Error::CorruptStoredEvent)?;
- Ok(StoredVisibleEvent::new(row.position, event))
+ .filter(|row| {
+ row.position.sequence().get() > after
+ && query.selects(row.event.id())
+ && visibility.is_visible(row.event.id())
})
- .collect::<Result<Vec<_>, Error>>()?;
+ .take(usize::from(query.bounds().limit()) + 1)
+ .collect::<Vec<_>>();
+ let next = if selected.len() > usize::from(query.bounds().limit()) {
+ selected.truncate(usize::from(query.bounds().limit()));
+ selected.last().map(|row| row.position)
+ } else {
+ None
+ };
+ let items = selected
+ .into_iter()
+ .map(|row| StoredVisibleEvent::new(row.position, row.event))
+ .collect();
EventPage::new(self.generation, items, next, query.bounds())
})
}
+ fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>> {
+ Box::pin(async move {
+ let rows = self.current_event_rows().await?;
+ Ok(self.visibility_for_rows(&rows)?.into_snapshot())
+ })
+ }
+
fn query_provenance(
&self,
event_id: EventId,
@@ -638,12 +692,22 @@ mod tests {
}
fn signed_event(content: &str, pretty: bool) -> SignedEvent {
+ signed_event_with(content, pretty, 1_800_000_100, 0, vec![])
+ }
+
+ fn signed_event_with(
+ content: &str,
+ pretty: bool,
+ created_at: u64,
+ kind: u32,
+ tags: Vec<Vec<String>>,
+ ) -> SignedEvent {
let mut wire = Nip01EventWire {
id: "0".repeat(64),
pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
- created_at: 1_800_000_100,
- kind: 0,
- tags: vec![],
+ created_at,
+ kind,
+ tags,
content: content.to_owned(),
sig: "42".repeat(64),
extra: Default::default(),
@@ -689,9 +753,15 @@ mod tests {
}
fn visible(event: &SignedEvent) -> VisibleEvent {
- verified(event)
- .validate_contract()
- .expect("contract")
+ let verified = verified(event);
+ let validated = if event.envelope().kind_u32() == 5 {
+ verified
+ .validate_contract_for_admission("radroots.social.deletion_request.v1")
+ .expect("admission-selected contract")
+ } else {
+ verified.validate_contract().expect("contract")
+ };
+ validated
.admit_with(&Allow)
.expect("admission")
.make_visible_with(&Allow)
@@ -876,6 +946,74 @@ mod tests {
}
#[tokio::test]
+ async fn visibility_rebuild_survives_reopen_and_matches_current_head_deletion_queries() {
+ let generation = SourceGeneration::new([18; 32]).expect("generation");
+ let store = store(generation).await;
+ let old = signed_event_with(
+ r#"{"display_name":"Old Farm","bot":false}"#,
+ false,
+ 1_800_000_100,
+ 0,
+ vec![],
+ );
+ let current = signed_event_with(
+ r#"{"display_name":"Current Farm","bot":false}"#,
+ false,
+ 1_800_000_200,
+ 0,
+ vec![],
+ );
+ let deletion = signed_event_with(
+ "retired profile",
+ false,
+ 1_800_000_300,
+ 5,
+ vec![vec!["e".to_owned(), current.id().to_hex()]],
+ );
+ for (event, observed_at) in [
+ (old.clone(), 100),
+ (current.clone(), 200),
+ (deletion.clone(), 300),
+ ] {
+ store
+ .admit(
+ EventAdmission::visible(
+ observed(event.clone(), observed_at, None),
+ visible(&event),
+ )
+ .expect("visible admission"),
+ )
+ .await
+ .expect("admit visible event");
+ }
+
+ let before = store
+ .rebuild_visibility()
+ .await
+ .expect("visibility rebuild");
+ let reopened =
+ SqliteStorage::new(store.pool.clone(), generation, EventStoreMode::ReadWrite);
+ let after = reopened
+ .rebuild_visibility()
+ .await
+ .expect("reopened visibility rebuild");
+ assert_eq!(before, after);
+ assert_eq!(after.current_heads()[0].event_id, *current.id());
+ assert_eq!(after.visible_event_ids(), &[*deletion.id()]);
+ assert_eq!(after.suppressed_event_ids(), &[*current.id()]);
+ assert_eq!(after.superseded_event_ids(), &[*old.id()]);
+ let page = reopened
+ .query_visible(EventQuery::all(
+ EventQueryBounds::first(10).expect("bounds"),
+ ))
+ .await
+ .expect("visible page");
+ assert_eq!(page.items().len(), 1);
+ assert_eq!(page.items()[0].event().id(), deletion.id());
+ assert_eq!(reopened.status().await.expect("status").visible_events(), 1);
+ }
+
+ #[tokio::test]
async fn corrupt_rows_fail_closed_and_source_history_is_immutable() {
let generation = SourceGeneration::new([10; 32]).expect("generation");
let store = store(generation).await;
@@ -884,10 +1022,23 @@ mod tests {
false,
);
store
- .admit(EventAdmission::raw(observed(event, 30, None)))
+ .admit(EventAdmission::raw(observed(event.clone(), 30, None)))
.await
.expect("event");
+ let wrong_generation = SourceGeneration::new([11; 32]).expect("generation");
+ let mismatched_store = SqliteStorage::new(
+ store.pool.clone(),
+ wrong_generation,
+ EventStoreMode::ReadWrite,
+ );
+ assert_eq!(
+ mismatched_store
+ .admit(EventAdmission::raw(observed(event, 31, None)))
+ .await,
+ Err(Error::CorruptStoredEvent)
+ );
+
assert!(
sqlx::query("DELETE FROM radroots_runtime_events")
.execute(&store.pool)
@@ -917,6 +1068,33 @@ mod tests {
.await,
Err(Error::CorruptStoredEvent)
);
+ sqlx::query(
+ "UPDATE radroots_runtime_events
+ SET admitted_contract_id = 'radroots.social.geochat.v1',
+ admitted_registry_version = 7",
+ )
+ .execute(&store.pool)
+ .await
+ .expect("forge mismatched selected contract");
+ assert_eq!(
+ store
+ .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),))
+ .await,
+ Err(Error::CorruptStoredEvent)
+ );
+ sqlx::query(
+ "UPDATE radroots_runtime_events
+ SET admitted_registry_version = 6",
+ )
+ .execute(&store.pool)
+ .await
+ .expect("forge registry version");
+ assert_eq!(
+ store
+ .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),))
+ .await,
+ Err(Error::CorruptStoredEvent)
+ );
sqlx::query("PRAGMA ignore_check_constraints = ON")
.execute(&store.pool)
.await
diff --git a/crates/sync/tests/projection.rs b/crates/sync/tests/projection.rs
@@ -173,19 +173,19 @@ fn setup() -> (Engine, Arc<MemoryStorage>, ProjectionId) {
fn signed_event(created_at: u64) -> SignedEvent {
let tags: Vec<Vec<String>> = vec![];
- let id = compute_canonical_nip01_event_id(PUBKEY, created_at, 0, &tags, CONTENT)
+ let id = compute_canonical_nip01_event_id(PUBKEY, created_at, 1, &tags, CONTENT)
.expect("event id")
.to_hex();
let signature = "42".repeat(64);
let raw_json = format!(
- "{{\"id\":\"{id}\",\"pubkey\":\"{PUBKEY}\",\"created_at\":{created_at},\"kind\":0,\"tags\":[],\"content\":{content:?},\"sig\":\"{signature}\"}}",
+ "{{\"id\":\"{id}\",\"pubkey\":\"{PUBKEY}\",\"created_at\":{created_at},\"kind\":1,\"tags\":[],\"content\":{content:?},\"sig\":\"{signature}\"}}",
content = CONTENT,
);
SignedEvent::new(SignedEventParts {
id,
pubkey: PUBKEY.to_owned(),
created_at,
- kind: 0,
+ kind: 1,
tags,
content: CONTENT.to_owned(),
sig: signature,
diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs
@@ -154,6 +154,14 @@ impl EventStore for FaultStorage {
> {
EventStore::query_visible(self.inner.as_ref(), value)
}
+ fn rebuild_visibility(
+ &self,
+ ) -> radroots_transport::BoxFuture<
+ '_,
+ Result<radroots_storage::event::VisibilitySnapshot, radroots_storage::Error>,
+ > {
+ EventStore::rebuild_visibility(self.inner.as_ref())
+ }
fn query_provenance(
&self,
id: radroots_storage::event::EventId,
@@ -1427,13 +1435,17 @@ async fn sqlite_signing_and_admission_recover_across_every_reopen_boundary() {
.expect("admission replay");
assert!(replay.is_replay());
assert_eq!(signer.calls.load(Ordering::Relaxed), 0);
- let visible = store
- .query_visible(EventQuery::all(
+ let admitted_events = store
+ .query_raw(EventQuery::all(
EventQueryBounds::first(10).expect("bounds"),
))
.await
- .expect("visible events");
- assert_eq!(visible.items().len(), 1);
+ .expect("admitted events");
+ assert_eq!(admitted_events.items().len(), 1);
+ assert_eq!(
+ admitted_events.items()[0].stage(),
+ radroots_storage::event::AdmissionStage::Visible
+ );
}
let store = Arc::new(