commit f2afcab4d8fcf6bb529649b89138318fd0ed3cdb parent f9281482f2fb4ceb5d524044821cdcedfef8c540 Author: triesap <tyson@radroots.org> Date: Mon, 6 Jul 2026 20:05:04 +0000 event-store: add transport observations - replace relay-only event-store observations with transport-kind endpoint rows - add typed observation models, endpoint lookup, and transport kind parsing - update Nostr relay ingest and trade projection count call sites - validate event-store, downstream transport/outbox/trade tests, and contracts Diffstat:
17 files changed, 358 insertions(+), 134 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock @@ -4287,6 +4287,7 @@ version = "0.1.0-alpha.2" dependencies = [ "radroots_events", "radroots_nostr", + "radroots_transport", "serde", "serde_json", "sqlx", @@ -4562,6 +4563,7 @@ dependencies = [ "radroots_events", "radroots_nostr", "radroots_outbox", + "radroots_transport", "serde", "serde_json", "thiserror 1.0.69", @@ -4854,6 +4856,7 @@ dependencies = [ "radroots_events", "radroots_events_codec", "radroots_nostr", + "radroots_transport", "serde", "serde_json", "sha2", diff --git a/crates/event_store/Cargo.toml b/crates/event_store/Cargo.toml @@ -25,6 +25,7 @@ radroots_nostr = { workspace = true, default-features = false, features = [ "std", "events", ] } +radroots_transport = { workspace = true, default-features = false } serde = { workspace = true, features = ["std"] } serde_json = { workspace = true, features = ["std"] } sqlx = { workspace = true, optional = true, features = ["derive"] } diff --git a/crates/event_store/migrations/0001_event_store.down.sql b/crates/event_store/migrations/0001_event_store.down.sql @@ -3,6 +3,6 @@ DROP TABLE trade_projection; DROP TABLE listing_projection; DROP TABLE projection_cursor; DROP TABLE nostr_event_head; -DROP TABLE relay_event_seen; +DROP TABLE event_transport_observation; DROP TABLE nostr_event_tags; DROP TABLE nostr_events; diff --git a/crates/event_store/migrations/0001_event_store.up.sql b/crates/event_store/migrations/0001_event_store.up.sql @@ -38,18 +38,21 @@ CREATE TABLE IF NOT EXISTS nostr_event_tags ( CREATE INDEX IF NOT EXISTS nostr_event_tag_lookup_idx ON nostr_event_tags(tag_name, tag_value, event_id); CREATE INDEX IF NOT EXISTS nostr_event_tag_relay_idx ON nostr_event_tags(relay_indexed, tag_name, tag_value, event_id); -CREATE TABLE IF NOT EXISTS relay_event_seen ( +CREATE TABLE IF NOT EXISTS event_transport_observation ( event_id TEXT NOT NULL REFERENCES nostr_events(event_id) ON DELETE CASCADE, - relay_url TEXT NOT NULL, + transport_kind TEXT NOT NULL, + endpoint_uri TEXT NOT NULL, + endpoint_fingerprint TEXT NOT NULL, observation_type TEXT NOT NULL, - first_seen_at_ms INTEGER NOT NULL, - last_seen_at_ms INTEGER NOT NULL, + first_observed_at_ms INTEGER NOT NULL, + last_observed_at_ms INTEGER NOT NULL, observation_count INTEGER NOT NULL, - last_message TEXT, - PRIMARY KEY (event_id, relay_url, observation_type) + redacted_message TEXT, + PRIMARY KEY (event_id, transport_kind, endpoint_fingerprint, observation_type) ); -CREATE INDEX IF NOT EXISTS relay_event_seen_relay_idx ON relay_event_seen(relay_url, last_seen_at_ms, event_id); +CREATE INDEX IF NOT EXISTS event_transport_observation_endpoint_idx +ON event_transport_observation(transport_kind, endpoint_fingerprint, last_observed_at_ms, event_id); CREATE TABLE IF NOT EXISTS nostr_event_head ( coordinate_type TEXT NOT NULL, @@ -152,7 +155,7 @@ CREATE TABLE IF NOT EXISTS trade_projection ( issues_json TEXT NOT NULL, issue_count INTEGER NOT NULL, source_event_count INTEGER NOT NULL, - relay_observation_count INTEGER NOT NULL, + transport_observation_count INTEGER NOT NULL, evidence_hash TEXT NOT NULL, last_source_event_seq INTEGER, updated_at_ms INTEGER NOT NULL, diff --git a/crates/event_store/src/error.rs b/crates/event_store/src/error.rs @@ -1,6 +1,7 @@ use radroots_events::contract::RadrootsContractMatchError; use radroots_events::event_head::RadrootsEventHeadMalformed; use radroots_events::ids::RadrootsIdParseError; +use radroots_transport::RadrootsTransportError; #[derive(Debug, thiserror::Error)] pub enum RadrootsEventStoreError { @@ -14,6 +15,8 @@ pub enum RadrootsEventStoreError { EventHeadMalformed(RadrootsEventHeadMalformed), #[error("identifier parse error: {0}")] IdParse(#[from] RadrootsIdParseError), + #[error("transport contract error: {0}")] + Transport(RadrootsTransportError), #[error("stored event `{0}` was not found")] MissingEvent(String), #[error("event-store tag query tag name cannot be empty")] @@ -26,6 +29,21 @@ pub enum RadrootsEventStoreError { QueryLimitOutOfRange { min: u32, max: u32, actual: u32 }, #[error("invalid stored enum value `{value}` for {field}")] InvalidStoredEnum { field: &'static str, value: String }, + #[error( + "stored transport observation fingerprint `{endpoint_fingerprint}` does not match `{transport_kind}` endpoint `{endpoint_uri}` for event `{event_id}`" + )] + InvalidStoredTransportEndpointFingerprint { + event_id: String, + transport_kind: String, + endpoint_uri: String, + endpoint_fingerprint: String, + }, #[error("integer value `{value}` is outside {field} range")] IntegerRange { field: &'static str, value: i64 }, } + +impl From<RadrootsTransportError> for RadrootsEventStoreError { + fn from(value: RadrootsTransportError) -> Self { + Self::Transport(value) + } +} diff --git a/crates/event_store/src/lib.rs b/crates/event_store/src/lib.rs @@ -18,11 +18,11 @@ pub use migrations::{EVENT_STORE_MIGRATION_DOWN, EVENT_STORE_MIGRATION_UP}; pub use model::{ RadrootsEventContractStatus, RadrootsEventHeadStoreDecision, RadrootsEventIngest, RadrootsEventIngestReceipt, RadrootsEventStoreStatusSummary, RadrootsEventVerificationStatus, - RadrootsProjectionCursor, RadrootsRelayObservation, RadrootsRelayObservationType, - RadrootsStoredEvent, RadrootsStoredEventHead, RadrootsStoredEventTag, StoredEventClass, + RadrootsProjectionCursor, RadrootsStoredEvent, RadrootsStoredEventHead, RadrootsStoredEventTag, + RadrootsTransportObservation, RadrootsTransportObservationType, StoredEventClass, }; #[cfg(feature = "sqlite")] pub use store::{ RADROOTS_EVENT_STORE_CONTRACT_QUERY_LIMIT_MAX, RADROOTS_EVENT_STORE_QUERY_LIMIT_MAX, - RadrootsEventStore, + RadrootsEventStore, RadrootsTransportObservationRow, }; diff --git a/crates/event_store/src/model.rs b/crates/event_store/src/model.rs @@ -4,6 +4,9 @@ use radroots_events::contract::{ RadrootsContractMatchError, RadrootsEventClass, RadrootsTagSemantic, RadrootsTagValueType, }; use radroots_events::event_head::RadrootsEventHeadDecision; +use radroots_transport::{ + RadrootsTransportKind, RadrootsTransportTargetFingerprint, RadrootsTransportTargetUri, +}; #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum RadrootsEventVerificationStatus { @@ -125,48 +128,87 @@ impl StoredEventClass { } #[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum RadrootsRelayObservationType { - Fetch, - Subscription, - PublishAck, - Import, +pub enum RadrootsTransportObservationType { + NostrFetch, + NostrSubscription, + NostrPublishAck, + LocalImport, + MeshHeard, + MeshForwarded, + GatewayStored, + GatewayRepublished, + DeliveryAck, + Other, } -impl RadrootsRelayObservationType { +impl RadrootsTransportObservationType { pub fn as_str(self) -> &'static str { match self { - Self::Fetch => "fetch", - Self::Subscription => "subscription", - Self::PublishAck => "publish_ack", - Self::Import => "import", + Self::NostrFetch => "nostr_fetch", + Self::NostrSubscription => "nostr_subscription", + Self::NostrPublishAck => "nostr_publish_ack", + Self::LocalImport => "local_import", + Self::MeshHeard => "mesh_heard", + Self::MeshForwarded => "mesh_forwarded", + Self::GatewayStored => "gateway_stored", + Self::GatewayRepublished => "gateway_republished", + Self::DeliveryAck => "delivery_ack", + Self::Other => "other", + } + } + + pub fn parse(value: &str) -> Result<Self, RadrootsEventStoreError> { + match value { + "nostr_fetch" => Ok(Self::NostrFetch), + "nostr_subscription" => Ok(Self::NostrSubscription), + "nostr_publish_ack" => Ok(Self::NostrPublishAck), + "local_import" => Ok(Self::LocalImport), + "mesh_heard" => Ok(Self::MeshHeard), + "mesh_forwarded" => Ok(Self::MeshForwarded), + "gateway_stored" => Ok(Self::GatewayStored), + "gateway_republished" => Ok(Self::GatewayRepublished), + "delivery_ack" => Ok(Self::DeliveryAck), + "other" => Ok(Self::Other), + _ => Err(RadrootsEventStoreError::InvalidStoredEnum { + field: "observation_type", + value: value.to_owned(), + }), } } } #[derive(Clone, Debug, PartialEq, Eq)] -pub struct RadrootsRelayObservation { - pub relay_url: String, - pub observation_type: RadrootsRelayObservationType, +pub struct RadrootsTransportObservation { + pub transport_kind: RadrootsTransportKind, + pub endpoint_uri: RadrootsTransportTargetUri, + pub endpoint_fingerprint: RadrootsTransportTargetFingerprint, + pub observation_type: RadrootsTransportObservationType, pub observed_at_ms: i64, - pub message: Option<String>, + pub redacted_message: Option<String>, } -impl RadrootsRelayObservation { +impl RadrootsTransportObservation { pub fn new( - relay_url: impl Into<String>, - observation_type: RadrootsRelayObservationType, + transport_kind: RadrootsTransportKind, + endpoint_uri: impl AsRef<str>, + observation_type: RadrootsTransportObservationType, observed_at_ms: i64, - ) -> Self { - Self { - relay_url: relay_url.into(), + ) -> Result<Self, RadrootsEventStoreError> { + let endpoint_uri = RadrootsTransportTargetUri::parse(endpoint_uri)?; + let endpoint_fingerprint = + RadrootsTransportTargetFingerprint::from_target(&transport_kind, &endpoint_uri); + Ok(Self { + transport_kind, + endpoint_uri, + endpoint_fingerprint, observation_type, observed_at_ms, - message: None, - } + redacted_message: None, + }) } - pub fn with_message(mut self, message: impl Into<String>) -> Self { - self.message = Some(message.into()); + pub fn with_redacted_message(mut self, message: impl Into<String>) -> Self { + self.redacted_message = Some(message.into()); self } } @@ -176,7 +218,7 @@ pub struct RadrootsEventIngest { pub event: RadrootsNostrEvent, pub raw_json: Option<String>, pub observed_at_ms: i64, - pub relay_observation: Option<RadrootsRelayObservation>, + pub transport_observation: Option<RadrootsTransportObservation>, } impl RadrootsEventIngest { @@ -185,7 +227,7 @@ impl RadrootsEventIngest { event, raw_json: None, observed_at_ms, - relay_observation: None, + transport_observation: None, } } @@ -194,8 +236,8 @@ impl RadrootsEventIngest { self } - pub fn with_observation(mut self, observation: RadrootsRelayObservation) -> Self { - self.relay_observation = Some(observation); + pub fn with_observation(mut self, observation: RadrootsTransportObservation) -> Self { + self.transport_observation = Some(observation); self } } @@ -243,7 +285,7 @@ pub struct RadrootsEventIngestReceipt { pub struct RadrootsEventStoreStatusSummary { pub total_events: i64, pub projection_eligible_events: i64, - pub relay_observations: i64, + pub transport_observations: i64, pub last_event_seq: Option<i64>, pub last_event_updated_at_ms: Option<i64>, } @@ -446,20 +488,37 @@ mod tests { assert!(StoredEventClass::parse("bad").is_err()); for observation_type in [ - RadrootsRelayObservationType::Fetch, - RadrootsRelayObservationType::Subscription, - RadrootsRelayObservationType::PublishAck, - RadrootsRelayObservationType::Import, + RadrootsTransportObservationType::NostrFetch, + RadrootsTransportObservationType::NostrSubscription, + RadrootsTransportObservationType::NostrPublishAck, + RadrootsTransportObservationType::LocalImport, + RadrootsTransportObservationType::MeshHeard, + RadrootsTransportObservationType::MeshForwarded, + RadrootsTransportObservationType::GatewayStored, + RadrootsTransportObservationType::GatewayRepublished, + RadrootsTransportObservationType::DeliveryAck, + RadrootsTransportObservationType::Other, ] { assert!(!observation_type.as_str().is_empty()); + assert_eq!( + RadrootsTransportObservationType::parse(observation_type.as_str()) + .expect("observation type"), + observation_type + ); } - let observation = RadrootsRelayObservation::new( + let observation = RadrootsTransportObservation::new( + RadrootsTransportKind::Nostr, "wss://relay.example.test", - RadrootsRelayObservationType::Fetch, + RadrootsTransportObservationType::NostrFetch, 1, ) - .with_message("seen"); - assert_eq!(observation.message.as_deref(), Some("seen")); + .expect("observation") + .with_redacted_message("seen"); + assert_eq!(observation.redacted_message.as_deref(), Some("seen")); + assert_eq!( + observation.endpoint_uri.as_str(), + "wss://relay.example.test" + ); } #[test] diff --git a/crates/event_store/src/store.rs b/crates/event_store/src/store.rs @@ -3,9 +3,9 @@ use crate::migrations::{EVENT_STORE_MIGRATION_DOWN, EVENT_STORE_MIGRATION_UP}; use crate::model::{ RadrootsEventContractStatus, RadrootsEventHeadStoreDecision, RadrootsEventIngest, RadrootsEventIngestReceipt, RadrootsEventStoreStatusSummary, RadrootsEventVerificationStatus, - RadrootsProjectionCursor, RadrootsRelayObservation, RadrootsStoredEvent, - RadrootsStoredEventHead, RadrootsStoredEventTag, StoredEventClass, tag_semantic_name, - tag_value_type_name, + RadrootsProjectionCursor, RadrootsStoredEvent, RadrootsStoredEventHead, RadrootsStoredEventTag, + RadrootsTransportObservation, RadrootsTransportObservationType, StoredEventClass, + tag_semantic_name, tag_value_type_name, }; use radroots_events::RadrootsNostrEvent; use radroots_events::contract::{ @@ -18,6 +18,9 @@ use radroots_events::event_head::{ }; use radroots_events::ids::{RadrootsEventId, RadrootsEventSignature, RadrootsPublicKey}; use radroots_nostr::prelude::{RadrootsNostrEventVerification, radroots_nostr_verify_event}; +use radroots_transport::{ + RadrootsTransportKind, RadrootsTransportTargetFingerprint, RadrootsTransportTargetUri, +}; use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; use sqlx::{Row, SqlitePool}; use std::path::Path; @@ -84,12 +87,15 @@ impl RadrootsEventStore { ) .fetch_one(&self.pool) .await?; - let relay_observations = - query_i64(&self.pool, "SELECT COUNT(*) FROM relay_event_seen").await?; + let transport_observations = query_i64( + &self.pool, + "SELECT COUNT(*) FROM event_transport_observation", + ) + .await?; Ok(RadrootsEventStoreStatusSummary { total_events: row.try_get("total_events")?, projection_eligible_events: row.try_get("projection_eligible_events")?, - relay_observations, + transport_observations, last_event_seq: row.try_get("last_event_seq")?, last_event_updated_at_ms: row.try_get("last_event_updated_at_ms")?, }) @@ -148,7 +154,7 @@ impl RadrootsEventStore { projection_eligible = false; } - if let Some(observation) = ingest.relay_observation.as_ref() { + if let Some(observation) = ingest.transport_observation.as_ref() { upsert_observation(&mut tx, ingest.event.id.as_str(), observation).await?; } @@ -197,14 +203,36 @@ impl RadrootsEventStore { pub async fn observations_for_event( &self, event_id: &str, - ) -> Result<Vec<RadrootsRelayObservationRow>, RadrootsEventStoreError> { + ) -> Result<Vec<RadrootsTransportObservationRow>, RadrootsEventStoreError> { let rows = sqlx::query( - "SELECT event_id, relay_url, observation_type, first_seen_at_ms, last_seen_at_ms, observation_count, last_message FROM relay_event_seen WHERE event_id = ? ORDER BY relay_url, observation_type", + "SELECT event_id, transport_kind, endpoint_uri, endpoint_fingerprint, observation_type, first_observed_at_ms, last_observed_at_ms, observation_count, redacted_message FROM event_transport_observation WHERE event_id = ? ORDER BY transport_kind, endpoint_uri, observation_type", ) .bind(event_id) .fetch_all(&self.pool) .await?; - rows.into_iter().map(relay_observation_from_row).collect() + rows.into_iter() + .map(transport_observation_from_row) + .collect() + } + + pub async fn observations_for_endpoint( + &self, + transport_kind: RadrootsTransportKind, + endpoint_uri: impl AsRef<str>, + ) -> Result<Vec<RadrootsTransportObservationRow>, RadrootsEventStoreError> { + let endpoint_uri = RadrootsTransportTargetUri::parse(endpoint_uri)?; + let endpoint_fingerprint = + RadrootsTransportTargetFingerprint::from_target(&transport_kind, &endpoint_uri); + let rows = sqlx::query( + "SELECT event_id, transport_kind, endpoint_uri, endpoint_fingerprint, observation_type, first_observed_at_ms, last_observed_at_ms, observation_count, redacted_message FROM event_transport_observation WHERE transport_kind = ? AND endpoint_fingerprint = ? ORDER BY last_observed_at_ms, event_id, observation_type", + ) + .bind(transport_kind.canonical_label()) + .bind(endpoint_fingerprint.as_str()) + .fetch_all(&self.pool) + .await?; + rows.into_iter() + .map(transport_observation_from_row) + .collect() } pub async fn event_head( @@ -338,14 +366,16 @@ impl RadrootsEventStore { } #[derive(Clone, Debug, PartialEq, Eq)] -pub struct RadrootsRelayObservationRow { +pub struct RadrootsTransportObservationRow { pub event_id: String, - pub relay_url: String, - pub observation_type: String, - pub first_seen_at_ms: i64, - pub last_seen_at_ms: i64, + pub transport_kind: RadrootsTransportKind, + pub endpoint_uri: RadrootsTransportTargetUri, + pub endpoint_fingerprint: RadrootsTransportTargetFingerprint, + pub observation_type: RadrootsTransportObservationType, + pub first_observed_at_ms: i64, + pub last_observed_at_ms: i64, pub observation_count: i64, - pub last_message: Option<String>, + pub redacted_message: Option<String>, } struct EventClassification { @@ -553,17 +583,19 @@ async fn insert_tags( async fn upsert_observation( tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, event_id: &str, - observation: &RadrootsRelayObservation, + observation: &RadrootsTransportObservation, ) -> Result<(), RadrootsEventStoreError> { sqlx::query( - "INSERT INTO relay_event_seen(event_id, relay_url, observation_type, first_seen_at_ms, last_seen_at_ms, observation_count, last_message) VALUES (?, ?, ?, ?, ?, 1, ?) ON CONFLICT(event_id, relay_url, observation_type) DO UPDATE SET last_seen_at_ms = excluded.last_seen_at_ms, observation_count = relay_event_seen.observation_count + 1, last_message = excluded.last_message", + "INSERT INTO event_transport_observation(event_id, transport_kind, endpoint_uri, endpoint_fingerprint, observation_type, first_observed_at_ms, last_observed_at_ms, observation_count, redacted_message) VALUES (?, ?, ?, ?, ?, ?, ?, 1, ?) ON CONFLICT(event_id, transport_kind, endpoint_fingerprint, observation_type) DO UPDATE SET endpoint_uri = excluded.endpoint_uri, last_observed_at_ms = excluded.last_observed_at_ms, observation_count = event_transport_observation.observation_count + 1, redacted_message = excluded.redacted_message", ) .bind(event_id) - .bind(observation.relay_url.as_str()) + .bind(observation.transport_kind.canonical_label()) + .bind(observation.endpoint_uri.as_str()) + .bind(observation.endpoint_fingerprint.as_str()) .bind(observation.observation_type.as_str()) .bind(observation.observed_at_ms) .bind(observation.observed_at_ms) - .bind(observation.message.as_deref()) + .bind(observation.redacted_message.as_deref()) .execute(&mut **tx) .await?; Ok(()) @@ -784,17 +816,41 @@ fn projection_cursor_from_row( } #[cfg_attr(coverage_nightly, coverage(off))] -fn relay_observation_from_row( +fn transport_observation_from_row( row: sqlx::sqlite::SqliteRow, -) -> Result<RadrootsRelayObservationRow, RadrootsEventStoreError> { - Ok(RadrootsRelayObservationRow { - event_id: row.try_get("event_id")?, - relay_url: row.try_get("relay_url")?, - observation_type: row.try_get("observation_type")?, - first_seen_at_ms: row.try_get("first_seen_at_ms")?, - last_seen_at_ms: row.try_get("last_seen_at_ms")?, +) -> Result<RadrootsTransportObservationRow, RadrootsEventStoreError> { + let event_id: String = row.try_get("event_id")?; + let transport_kind_label: String = row.try_get("transport_kind")?; + let endpoint_uri_raw: String = row.try_get("endpoint_uri")?; + let endpoint_fingerprint_raw: String = row.try_get("endpoint_fingerprint")?; + let transport_kind = RadrootsTransportKind::parse(&transport_kind_label)?; + let endpoint_uri = RadrootsTransportTargetUri::parse(&endpoint_uri_raw)?; + let endpoint_fingerprint = + RadrootsTransportTargetFingerprint::parse(&endpoint_fingerprint_raw)?; + let expected_fingerprint = + RadrootsTransportTargetFingerprint::from_target(&transport_kind, &endpoint_uri); + if endpoint_fingerprint != expected_fingerprint { + return Err( + RadrootsEventStoreError::InvalidStoredTransportEndpointFingerprint { + event_id, + transport_kind: transport_kind_label, + endpoint_uri: endpoint_uri_raw, + endpoint_fingerprint: endpoint_fingerprint_raw, + }, + ); + } + Ok(RadrootsTransportObservationRow { + event_id, + transport_kind, + endpoint_uri, + endpoint_fingerprint, + observation_type: RadrootsTransportObservationType::parse( + row.try_get("observation_type")?, + )?, + first_observed_at_ms: row.try_get("first_observed_at_ms")?, + last_observed_at_ms: row.try_get("last_observed_at_ms")?, observation_count: row.try_get("observation_count")?, - last_message: row.try_get("last_message")?, + redacted_message: row.try_get("redacted_message")?, }) } @@ -965,13 +1021,13 @@ mod tests { } #[tokio::test] - async fn status_summary_counts_events_projections_and_relay_observations() { + async fn status_summary_counts_events_projections_and_transport_observations() { let store = RadrootsEventStore::open_memory().await.expect("open"); let empty = store.status_summary().await.expect("empty status"); assert_eq!(empty.total_events, 0); assert_eq!(empty.projection_eligible_events, 0); - assert_eq!(empty.relay_observations, 0); + assert_eq!(empty.transport_observations, 0); assert_eq!(empty.last_event_seq, None); assert_eq!(empty.last_event_updated_at_ms, None); @@ -981,11 +1037,13 @@ mod tests { vec![vec!["t".to_owned(), "soil".to_owned()]], "hello", ); - let observation = RadrootsRelayObservation::new( + let observation = RadrootsTransportObservation::new( + RadrootsTransportKind::Nostr, "wss://relay.example.com", - crate::RadrootsRelayObservationType::PublishAck, + crate::RadrootsTransportObservationType::NostrPublishAck, 1_100, - ); + ) + .expect("observation"); store .ingest_event(RadrootsEventIngest::new(event.clone(), 1_000)) .await @@ -998,7 +1056,7 @@ mod tests { let status = store.status_summary().await.expect("status"); assert_eq!(status.total_events, 1); assert_eq!(status.projection_eligible_events, 1); - assert_eq!(status.relay_observations, 1); + assert_eq!(status.transport_observations, 1); assert_eq!(status.last_event_seq, Some(1)); assert_eq!(status.last_event_updated_at_ms, Some(1_000)); } @@ -1772,22 +1830,26 @@ mod tests { } #[tokio::test] - async fn relay_observations_upsert_separately_from_event_identity() { + async fn transport_observations_upsert_and_query_by_endpoint() { let store = RadrootsEventStore::open_memory().await.expect("open"); let event = signed_event(KIND_POST, 15, Vec::new(), "hello"); - let observation = RadrootsRelayObservation::new( + let observation = RadrootsTransportObservation::new( + RadrootsTransportKind::Nostr, "wss://relay.local", - crate::RadrootsRelayObservationType::Subscription, + crate::RadrootsTransportObservationType::NostrSubscription, 4_000, - ); + ) + .expect("observation"); let ingest = RadrootsEventIngest::new(event.clone(), 4_000).with_observation(observation); store.ingest_event(ingest).await.expect("first"); - let observation = RadrootsRelayObservation::new( + let observation = RadrootsTransportObservation::new( + RadrootsTransportKind::Nostr, "wss://relay.local", - crate::RadrootsRelayObservationType::Subscription, + crate::RadrootsTransportObservationType::NostrSubscription, 4_100, ) - .with_message("duplicate accepted"); + .expect("observation") + .with_redacted_message("duplicate accepted"); let ingest = RadrootsEventIngest::new(event.clone(), 4_100).with_observation(observation); store.ingest_event(ingest).await.expect("second"); @@ -1796,12 +1858,25 @@ mod tests { .await .expect("observations"); assert_eq!(observations.len(), 1); + assert_eq!(observations[0].transport_kind, RadrootsTransportKind::Nostr); + assert_eq!(observations[0].endpoint_uri.as_str(), "wss://relay.local"); + assert_eq!( + observations[0].observation_type, + crate::RadrootsTransportObservationType::NostrSubscription + ); assert_eq!(observations[0].observation_count, 2); - assert_eq!(observations[0].last_seen_at_ms, 4_100); + assert_eq!(observations[0].first_observed_at_ms, 4_000); + assert_eq!(observations[0].last_observed_at_ms, 4_100); assert_eq!( - observations[0].last_message.as_deref(), + observations[0].redacted_message.as_deref(), Some("duplicate accepted") ); + + let endpoint_observations = store + .observations_for_endpoint(RadrootsTransportKind::Nostr, "WSS://RELAY.LOCAL") + .await + .expect("endpoint observations"); + assert_eq!(endpoint_observations, observations); } #[tokio::test] diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -2698,7 +2698,7 @@ mod tests { } #[tokio::test] - async fn local_signed_event_ingest_is_idempotent_without_relay_observation() { + async fn local_signed_event_ingest_is_idempotent_without_transport_observation() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); let event_store = RadrootsEventStore::open_memory() .await diff --git a/crates/relay_transport/Cargo.toml b/crates/relay_transport/Cargo.toml @@ -46,6 +46,7 @@ radroots_outbox = { workspace = true, optional = true, default-features = false, "sqlite", "runtime-tokio", ] } +radroots_transport = { workspace = true, default-features = false } futures = { workspace = true } nostr = { workspace = true } serde = { workspace = true, features = ["derive", "std"] } diff --git a/crates/relay_transport/src/fetch.rs b/crates/relay_transport/src/fetch.rs @@ -5,12 +5,13 @@ use core::time::Duration; use futures::future::BoxFuture; use nostr::{JsonUtil, filter::MatchEventOptions}; use radroots_event_store::{ - RadrootsEventContractStatus, RadrootsEventIngest, RadrootsEventStore, RadrootsRelayObservation, - RadrootsRelayObservationType, + RadrootsEventContractStatus, RadrootsEventIngest, RadrootsEventStore, + RadrootsTransportObservation, RadrootsTransportObservationType, }; use radroots_nostr::prelude::{ RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrFilter, radroots_event_from_nostr, }; +use radroots_transport::RadrootsTransportKind; use serde::{Deserialize, Serialize}; use std::sync::{Arc, Mutex, PoisonError}; @@ -361,18 +362,20 @@ where }) => { let event = radroots_event_from_nostr(&raw_event); let observation_type = match mode { - RadrootsRelayFetchMode::Fetch => RadrootsRelayObservationType::Fetch, + RadrootsRelayFetchMode::Fetch => RadrootsTransportObservationType::NostrFetch, RadrootsRelayFetchMode::Subscription => { - RadrootsRelayObservationType::Subscription + RadrootsTransportObservationType::NostrSubscription } }; + let observation = RadrootsTransportObservation::new( + RadrootsTransportKind::Nostr, + relay_url.clone(), + observation_type, + observed_at_ms, + )?; let ingest = RadrootsEventIngest::new(event, observed_at_ms) .with_raw_json(raw_json) - .with_observation(RadrootsRelayObservation::new( - relay_url.clone(), - observation_type, - observed_at_ms, - )); + .with_observation(observation); match event_store.ingest_event(ingest).await { Ok(store_receipt) => { let unsupported = diff --git a/crates/relay_transport/src/outbox.rs b/crates/relay_transport/src/outbox.rs @@ -7,7 +7,8 @@ use crate::{ publish_signed_event, }; use radroots_event_store::{ - RadrootsEventIngest, RadrootsEventStore, RadrootsRelayObservation, RadrootsRelayObservationType, + RadrootsEventIngest, RadrootsEventStore, RadrootsTransportObservation, + RadrootsTransportObservationType, }; use radroots_events::RadrootsNostrEvent; use radroots_events::draft::RadrootsSignedNostrEvent; @@ -15,6 +16,7 @@ use radroots_outbox::{ RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxEventStoreIngestReceipt, RadrootsOutboxRelayStatus, }; +use radroots_transport::RadrootsTransportKind; #[derive(Clone, Debug, PartialEq, Eq)] pub struct RadrootsOutboxPublishPolicy { @@ -274,13 +276,14 @@ async fn ingest_publish_observation( message: Option<&str>, observed_at_ms: i64, ) -> Result<(), RadrootsRelayTransportError> { - let mut observation = RadrootsRelayObservation::new( + let mut observation = RadrootsTransportObservation::new( + RadrootsTransportKind::Nostr, relay_url, - RadrootsRelayObservationType::PublishAck, + RadrootsTransportObservationType::NostrPublishAck, observed_at_ms, - ); + )?; if let Some(message) = message { - observation = observation.with_message(message); + observation = observation.with_redacted_message(message); } let ingest = RadrootsEventIngest::new(event_from_signed(signed_event), observed_at_ms) .with_raw_json(signed_event.raw_json.clone()) diff --git a/crates/relay_transport/tests/transport.rs b/crates/relay_transport/tests/transport.rs @@ -1,6 +1,8 @@ use futures::future::BoxFuture; use nostr::JsonUtil; -use radroots_event_store::{RadrootsEventStore, RadrootsEventVerificationStatus}; +use radroots_event_store::{ + RadrootsEventStore, RadrootsEventVerificationStatus, RadrootsTransportObservationType, +}; use radroots_events::draft::{RadrootsFrozenEventDraft, RadrootsSignedNostrEvent}; use radroots_events::kinds::KIND_POST; use radroots_nostr::prelude::{ @@ -21,6 +23,7 @@ use radroots_relay_transport::{ RadrootsRelayUrl, RadrootsRelayUrlPolicy, fetch_and_ingest_relay_events, fetch_relay_events, fetch_relay_events_blocking, publish_claimed_outbox_event, publish_signed_event, }; +use radroots_transport::RadrootsTransportKind; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; const FIXTURE_ALICE_SECRET_KEY_HEX: &str = @@ -605,7 +608,7 @@ fn fetch_blocking_facade_runs_mock_adapter() { } #[tokio::test] -async fn fetch_ingests_events_and_records_relay_observations() { +async fn fetch_ingests_events_and_records_transport_observations() { let signed = signed_post("hello"); let store = RadrootsEventStore::open_memory().await.expect("store"); let adapter = RadrootsMockRelayFetchAdapter::new(vec![ @@ -721,7 +724,12 @@ async fn fetch_ingests_events_and_records_relay_observations() { .await .expect("observations"); assert_eq!(observations.len(), 1); - assert_eq!(observations[0].relay_url, RELAY_PRIMARY_WSS); + assert_eq!(observations[0].transport_kind, RadrootsTransportKind::Nostr); + assert_eq!(observations[0].endpoint_uri.as_str(), RELAY_PRIMARY_WSS); + assert_eq!( + observations[0].observation_type, + RadrootsTransportObservationType::NostrFetch + ); assert_eq!(observations[0].observation_count, 2); } @@ -1059,7 +1067,10 @@ async fn fetch_subscription_mode_and_store_errors_are_reported() { .await .expect("observations"); assert_eq!(observations.len(), 1); - assert_eq!(observations[0].observation_type, "subscription"); + assert_eq!( + observations[0].observation_type, + RadrootsTransportObservationType::NostrSubscription + ); let closed_store = RadrootsEventStore::open_memory().await.expect("store"); closed_store.pool().close().await; diff --git a/crates/trade/Cargo.toml b/crates/trade/Cargo.toml @@ -80,6 +80,7 @@ radroots_nostr = { workspace = true, default-features = false, features = [ "std", "events", ] } +radroots_transport = { workspace = true, default-features = false } sqlx = { workspace = true, default-features = false, features = [ "runtime-tokio", "sqlite", diff --git a/crates/trade/src/projection.rs b/crates/trade/src/projection.rs @@ -145,7 +145,7 @@ pub struct RadrootsProjectionRefreshReceipt { pub listing_upserts: usize, pub trade_upserts: usize, pub validation_receipts: usize, - pub relay_observations: i64, + pub transport_observations: i64, pub last_event_seq: Option<i64>, } @@ -215,8 +215,8 @@ pub async fn refresh_product_projections( for stored_event in &events { receipt.last_event_seq = Some(stored_event.seq); - receipt.relay_observations += - relay_observation_count_for_event(store, &stored_event.event_id).await?; + receipt.transport_observations += + transport_observation_count_for_event(store, &stored_event.event_id).await?; if is_listing_kind(stored_event.kind) { let event = stored_event_to_nostr_event(stored_event)?; let listing = validate_listing_event(&event).map_err(|source| { @@ -449,7 +449,8 @@ async fn upsert_trade_projection( }; let event_ids = projection_source_event_ids(&root_event_id, &projection); let source_event_count = event_ids.len(); - let relay_observation_count = relay_observation_count_for_events(store, &event_ids).await?; + let transport_observation_count = + transport_observation_count_for_events(store, &event_ids).await?; let evidence_hash = projection_evidence_hash(&event_ids); upsert_trade_projection_row( store, @@ -459,7 +460,7 @@ async fn upsert_trade_projection( expected_listing_event_id.as_ref(), current_listing_event_id.as_ref(), source_event_count, - relay_observation_count, + transport_observation_count, &evidence_hash, inputs.last_source_event_seq, updated_at_ms, @@ -479,7 +480,7 @@ async fn upsert_trade_projection_row( expected_listing_event_id: Option<&RadrootsEventId>, current_listing_event_id: Option<&RadrootsEventId>, source_event_count: usize, - relay_observation_count: i64, + transport_observation_count: i64, evidence_hash: &str, last_source_event_seq: Option<i64>, updated_at_ms: i64, @@ -521,7 +522,7 @@ async fn upsert_trade_projection_row( })?; sqlx::query( - "INSERT INTO trade_projection(order_id, root_event_id, projection_version, status, lifecycle_terminal, rhi_state, listing_addr, buyer_pubkey, seller_pubkey, request_event_id, decision_event_id, agreement_event_id, pending_revision_event_id, cancellation_event_id, validation_receipt_event_id, last_event_id, expected_listing_event_id, current_listing_event_id, economics_json, pending_inventory_json, committed_inventory_json, issues_json, issue_count, source_event_count, relay_observation_count, evidence_hash, last_source_event_seq, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(order_id, root_event_id, projection_version) DO UPDATE SET status = excluded.status, lifecycle_terminal = excluded.lifecycle_terminal, rhi_state = excluded.rhi_state, listing_addr = excluded.listing_addr, buyer_pubkey = excluded.buyer_pubkey, seller_pubkey = excluded.seller_pubkey, request_event_id = excluded.request_event_id, decision_event_id = excluded.decision_event_id, agreement_event_id = excluded.agreement_event_id, pending_revision_event_id = excluded.pending_revision_event_id, cancellation_event_id = excluded.cancellation_event_id, validation_receipt_event_id = excluded.validation_receipt_event_id, last_event_id = excluded.last_event_id, expected_listing_event_id = excluded.expected_listing_event_id, current_listing_event_id = excluded.current_listing_event_id, economics_json = excluded.economics_json, pending_inventory_json = excluded.pending_inventory_json, committed_inventory_json = excluded.committed_inventory_json, issues_json = excluded.issues_json, issue_count = excluded.issue_count, source_event_count = excluded.source_event_count, relay_observation_count = excluded.relay_observation_count, evidence_hash = excluded.evidence_hash, last_source_event_seq = excluded.last_source_event_seq, updated_at_ms = excluded.updated_at_ms", + "INSERT INTO trade_projection(order_id, root_event_id, projection_version, status, lifecycle_terminal, rhi_state, listing_addr, buyer_pubkey, seller_pubkey, request_event_id, decision_event_id, agreement_event_id, pending_revision_event_id, cancellation_event_id, validation_receipt_event_id, last_event_id, expected_listing_event_id, current_listing_event_id, economics_json, pending_inventory_json, committed_inventory_json, issues_json, issue_count, source_event_count, transport_observation_count, evidence_hash, last_source_event_seq, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(order_id, root_event_id, projection_version) DO UPDATE SET status = excluded.status, lifecycle_terminal = excluded.lifecycle_terminal, rhi_state = excluded.rhi_state, listing_addr = excluded.listing_addr, buyer_pubkey = excluded.buyer_pubkey, seller_pubkey = excluded.seller_pubkey, request_event_id = excluded.request_event_id, decision_event_id = excluded.decision_event_id, agreement_event_id = excluded.agreement_event_id, pending_revision_event_id = excluded.pending_revision_event_id, cancellation_event_id = excluded.cancellation_event_id, validation_receipt_event_id = excluded.validation_receipt_event_id, last_event_id = excluded.last_event_id, expected_listing_event_id = excluded.expected_listing_event_id, current_listing_event_id = excluded.current_listing_event_id, economics_json = excluded.economics_json, pending_inventory_json = excluded.pending_inventory_json, committed_inventory_json = excluded.committed_inventory_json, issues_json = excluded.issues_json, issue_count = excluded.issue_count, source_event_count = excluded.source_event_count, transport_observation_count = excluded.transport_observation_count, evidence_hash = excluded.evidence_hash, last_source_event_seq = excluded.last_source_event_seq, updated_at_ms = excluded.updated_at_ms", ) .bind(order_id.as_str()) .bind(root_event_id.as_str()) @@ -547,7 +548,7 @@ async fn upsert_trade_projection_row( .bind(issues_json) .bind(i64::try_from(projection.issues.len()).unwrap_or(i64::MAX)) .bind(i64::try_from(source_event_count).unwrap_or(i64::MAX)) - .bind(relay_observation_count) + .bind(transport_observation_count) .bind(evidence_hash) .bind(last_source_event_seq) .bind(updated_at_ms) @@ -769,25 +770,26 @@ fn stored_event_to_nostr_event( }) } -async fn relay_observation_count_for_events( +async fn transport_observation_count_for_events( store: &RadrootsEventStore, event_ids: &[RadrootsEventId], ) -> Result<i64, RadrootsTradeProjectionError> { let mut count = 0; for event_id in event_ids { - count += relay_observation_count_for_event(store, event_id.as_str()).await?; + count += transport_observation_count_for_event(store, event_id.as_str()).await?; } Ok(count) } -async fn relay_observation_count_for_event( +async fn transport_observation_count_for_event( store: &RadrootsEventStore, event_id: &str, ) -> Result<i64, RadrootsTradeProjectionError> { - let row = sqlx::query("SELECT COUNT(*) AS count FROM relay_event_seen WHERE event_id = ?") - .bind(event_id) - .fetch_one(store.pool()) - .await?; + let row = + sqlx::query("SELECT COUNT(*) AS count FROM event_transport_observation WHERE event_id = ?") + .bind(event_id) + .fetch_one(store.pool()) + .await?; Ok(row.try_get("count")?) } @@ -906,7 +908,9 @@ mod tests { RadrootsCoreCurrency, RadrootsCoreDecimal, RadrootsCoreMoney, RadrootsCoreQuantity, RadrootsCoreQuantityPrice, RadrootsCoreUnit, }; - use radroots_event_store::{RadrootsEventIngest, RadrootsRelayObservation}; + use radroots_event_store::{ + RadrootsEventIngest, RadrootsTransportObservation, RadrootsTransportObservationType, + }; use radroots_events::{ RadrootsNostrEventPtr, farm::RadrootsFarmRef, @@ -927,6 +931,7 @@ mod tests { RadrootsNostrKeys, RadrootsNostrSecretKey, RadrootsNostrTimestamp, radroots_event_from_nostr, radroots_nostr_build_event, }; + use radroots_transport::RadrootsTransportKind; use crate::validation_receipt::{ RadrootsTradeValidationReceipt, RadrootsValidationReceiptProof, @@ -1213,11 +1218,13 @@ mod tests { store .ingest_event( RadrootsEventIngest::new(request_event.clone(), 20).with_observation( - RadrootsRelayObservation::new( + RadrootsTransportObservation::new( + RadrootsTransportKind::Nostr, "wss://relay.example.test", - radroots_event_store::RadrootsRelayObservationType::Import, + RadrootsTransportObservationType::LocalImport, 20, - ), + ) + .expect("observation"), ), ) .await @@ -1239,7 +1246,7 @@ mod tests { assert_eq!(refresh.listing_upserts, 1); assert_eq!(refresh.trade_upserts, 1); assert_eq!(refresh.validation_receipts, 1); - assert_eq!(refresh.relay_observations, 1); + assert_eq!(refresh.transport_observations, 1); let rows = search_listing_projection( &store, @@ -1265,7 +1272,7 @@ mod tests { let root_event_id = RadrootsEventId::parse(request_event.id).expect("request"); let trade_row = sqlx::query( - "SELECT root_event_id, projection_version, status, rhi_state, relay_observation_count, source_event_count, evidence_hash FROM trade_projection WHERE order_id = ? AND root_event_id = ? AND projection_version = ?", + "SELECT root_event_id, projection_version, status, rhi_state, transport_observation_count, source_event_count, evidence_hash FROM trade_projection WHERE order_id = ? AND root_event_id = ? AND projection_version = ?", ) .bind(order_id().as_str()) .bind(root_event_id.as_str()) @@ -1291,7 +1298,7 @@ mod tests { ); assert_eq!( trade_row - .try_get::<i64, _>("relay_observation_count") + .try_get::<i64, _>("transport_observation_count") .unwrap(), 1 ); diff --git a/crates/transport/src/kind.rs b/crates/transport/src/kind.rs @@ -12,6 +12,17 @@ pub enum RadrootsTransportKind { } impl RadrootsTransportKind { + pub fn parse(value: impl AsRef<str>) -> Result<Self, RadrootsTransportError> { + let canonical = value.as_ref().trim().to_ascii_lowercase(); + match canonical.as_str() { + "nostr" => Ok(Self::Nostr), + "reticulum" => Ok(Self::Reticulum), + "mesh" => Ok(Self::Mesh), + "local" => Ok(Self::Local), + _ => Self::custom(canonical), + } + } + pub fn custom(value: impl Into<String>) -> Result<Self, RadrootsTransportError> { let value = value.into(); let canonical = value.trim().to_ascii_lowercase(); diff --git a/crates/transport/tests/transport.rs b/crates/transport/tests/transport.rs @@ -28,6 +28,34 @@ fn target_fingerprints_are_stable_and_transport_scoped() { } #[test] +fn transport_kind_parser_round_trips_canonical_labels_and_custom_values() { + assert_eq!( + RadrootsTransportKind::parse(" NOSTR ").expect("nostr kind"), + RadrootsTransportKind::Nostr + ); + assert_eq!( + RadrootsTransportKind::parse("reticulum").expect("reticulum kind"), + RadrootsTransportKind::Reticulum + ); + assert_eq!( + RadrootsTransportKind::parse("mesh").expect("mesh kind"), + RadrootsTransportKind::Mesh + ); + assert_eq!( + RadrootsTransportKind::parse("local").expect("local kind"), + RadrootsTransportKind::Local + ); + assert_eq!( + RadrootsTransportKind::parse("fieldbus").expect("custom kind"), + RadrootsTransportKind::Custom("fieldbus".to_owned()) + ); + assert_eq!( + RadrootsTransportKind::parse("bad kind").expect_err("invalid kind"), + RadrootsTransportError::InvalidTransportKind + ); +} + +#[test] fn target_set_rejects_duplicate_fingerprints() { let first = RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, "wss://relay.example/a") .expect("first target");