tangle


git clone https://radroots.dev/git/tangle.git
Log | Files | Refs | README | LICENSE

commit 8568b59960664a87ae1e15d257f68aa96187b5ca
parent 216ed05e5de4c0624e9e912a25066297c5cf32b2
Author: triesap <tyson@radroots.org>
Date:   Fri, 26 Jun 2026 09:41:47 +0000

runtime: add neutral relay output projection hooks

- add projection contexts for historical query and live fanout output
- route session query and stored-offset fanout through projection decisions
- keep replacement output limited to stored events behind group read gates
- cover no-op, suppression, replacement, and private-group gate behavior

Diffstat:
MAGENTS.md | 1+
Mcrates/tangle_runtime/src/relay/core.rs | 4++--
Mcrates/tangle_runtime/src/runtime.rs | 690+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcrates/tangle_runtime/src/session.rs | 68++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcrates/tangle_runtime/tests/phase2_acceptance_targets.rs | 2+-
5 files changed, 741 insertions(+), 24 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -3,6 +3,7 @@ - this repo owns the `tangle` Nostr relay, including protocol handling, event validation, session admission, rate limiting, storage, indexing or projection boundaries, subscription fanout, group policy, relay-generated state, and relay-owned signing when explicitly part of the relay contract - do not make this repo responsible for user private-key custody, client-side signing, signer approval UX, remote-signer session control, wallet flows, account key management, platform-wide artifacts, publication, promotion, or deployment transport unless an approved spec assigns that responsibility here - product behavior should follow approved relay specs and repo-local adopted decisions when they conflict with legacy implementation details +- virtual relay tenancy is required for `tangle_v1_mvp`; host and tenant relay boundaries must remain explicit runtime behavior - prefer the smallest coherent change that fully addresses the request; do not mix unrelated cleanup, speculative refactors, compatibility scaffolding, or roadmap work into the same change - inspect the relevant implementation, tests, manifests, specs, and docs before changing behavior - do not invent requirements, APIs, dependencies, release processes, or external integration behavior diff --git a/crates/tangle_runtime/src/relay/core.rs b/crates/tangle_runtime/src/relay/core.rs @@ -143,7 +143,7 @@ struct BaseRelayGroupCountQuery<'a> { } impl BaseRelayQueryReport { - fn new( + pub(crate) fn new( messages: Vec<RuntimeRelayMessage>, group_read_denied: bool, query_metrics: BaseRelayQueryMetrics, @@ -248,7 +248,7 @@ impl BaseRelayQueryMetrics { } } - fn with_returned_events(self, returned_events: usize) -> Self { + pub(crate) fn with_returned_events(self, returned_events: usize) -> Self { Self { returned_events: u64::try_from(returned_events).expect("returned events fit in u64"), ..self diff --git a/crates/tangle_runtime/src/runtime.rs b/crates/tangle_runtime/src/runtime.rs @@ -90,6 +90,13 @@ pub trait RelayRuntimeHooks: Send + Sync { } fn event_stored(&self, _context: &RelayEventStoredContext) {} + + fn project_event( + &self, + _context: &RelayEventProjectionContext, + ) -> RelayEventProjectionDecision { + RelayEventProjectionDecision::Emit + } } #[derive(Debug, Default)] @@ -111,6 +118,62 @@ impl EventAdmissionDecision { } } +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct RelayProjectionContext { + identifier: Option<String>, +} + +impl RelayProjectionContext { + pub fn named(identifier: impl Into<String>) -> Result<Self, BaseRelayError> { + let identifier = identifier.into(); + if identifier.is_empty() { + return Err(BaseRelayError::invalid( + "relay projection identifier must not be empty", + )); + } + if identifier.chars().any(char::is_control) { + return Err(BaseRelayError::invalid( + "relay projection identifier must not contain control characters", + )); + } + Ok(Self { + identifier: Some(identifier), + }) + } + + pub fn identifier(&self) -> Option<&str> { + self.identifier.as_deref() + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RelayEventProjectionSource { + HistoricalQuery, + LiveFanout { store_offset: u64 }, +} + +impl RelayEventProjectionSource { + pub fn store_offset(self) -> Option<u64> { + match self { + Self::HistoricalQuery => None, + Self::LiveFanout { store_offset } => Some(store_offset), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RelayEventProjectionDecision { + Emit, + Suppress, + ReplaceWithStoredOffset { store_offset: u64 }, +} + +impl RelayEventProjectionDecision { + pub fn replace_with_stored_offset(store_offset: u64) -> Self { + Self::ReplaceWithStoredOffset { store_offset } + } +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct RelayEventContext { event_id: String, @@ -252,6 +315,46 @@ pub struct RelayEventStoredContext { store_offsets: Vec<u64>, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RelayEventProjectionContext { + subscription_id: SubscriptionId, + projection: RelayProjectionContext, + source: RelayEventProjectionSource, + event: RelayEventContext, +} + +impl RelayEventProjectionContext { + pub fn new( + subscription_id: SubscriptionId, + projection: RelayProjectionContext, + source: RelayEventProjectionSource, + event: RelayEventContext, + ) -> Self { + Self { + subscription_id, + projection, + source, + event, + } + } + + pub fn subscription_id(&self) -> &SubscriptionId { + &self.subscription_id + } + + pub fn projection(&self) -> &RelayProjectionContext { + &self.projection + } + + pub fn source(&self) -> RelayEventProjectionSource { + self.source + } + + pub fn event(&self) -> &RelayEventContext { + &self.event + } +} + impl RelayEventStoredContext { pub fn new(event: RelayEventContext, store_offsets: Vec<u64>) -> Self { Self { @@ -780,6 +883,92 @@ impl RelayRuntimeShared { ) } + fn project_query_report( + &self, + report: BaseRelayQueryReport, + projection: &RelayProjectionContext, + auth: &BaseAuthState, + ) -> Result<BaseRelayQueryReport, BaseRelayError> { + let group_read_denied = report.group_read_denied(); + let query_metrics = report.query_metrics(); + let messages = self.project_runtime_messages(report.into_messages(), projection, auth)?; + let returned_events = messages + .iter() + .filter(|message| matches!(message, RuntimeRelayMessage::Event { .. })) + .count(); + let query_metrics = query_metrics.with_returned_events(returned_events); + Ok(BaseRelayQueryReport::new( + messages, + group_read_denied, + query_metrics, + )) + } + + fn project_runtime_messages( + &self, + messages: Vec<RuntimeRelayMessage>, + projection: &RelayProjectionContext, + auth: &BaseAuthState, + ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { + let mut output = Vec::with_capacity(messages.len()); + for message in messages { + match message { + RuntimeRelayMessage::Event { + subscription_id, + event, + } => { + if let Some(projected) = self.project_event_output( + &subscription_id, + projection, + RelayEventProjectionSource::HistoricalQuery, + &event, + auth, + )? { + output.push(RuntimeRelayMessage::event(subscription_id, projected)); + } + } + message => output.push(message), + } + } + Ok(output) + } + + fn project_event_output( + &self, + subscription_id: &SubscriptionId, + projection: &RelayProjectionContext, + source: RelayEventProjectionSource, + event: &PocketEvent, + auth: &BaseAuthState, + ) -> Result<Option<PocketOwnedEvent>, BaseRelayError> { + let context = RelayEventProjectionContext::new( + subscription_id.clone(), + projection.clone(), + source, + RelayEventContext::from_pocket_event(event)?, + ); + match self.hooks.project_event(&context) { + RelayEventProjectionDecision::Emit => Ok(Some(event.to_owned())), + RelayEventProjectionDecision::Suppress => Ok(None), + RelayEventProjectionDecision::ReplaceWithStoredOffset { store_offset } => { + let Ok(replacement) = self.store.event_by_offset(store_offset) else { + return Ok(None); + }; + let group_auth = + GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); + if BaseRelay::group_read_gate_visible_to_auth( + self.groups.as_ref(), + &replacement, + &group_auth, + )? { + Ok(Some(replacement)) + } else { + Ok(None) + } + } + } + } + fn handle_count_with_auth_report( &self, subscription_id: SubscriptionId, @@ -1228,6 +1417,11 @@ impl RelayRuntimeHandle { search_present, auth, )?; + let report = self.inner.project_query_report( + report, + &RelayProjectionContext::default(), + auth, + )?; self.inner .metrics .record_query_metrics(report.query_metrics()); @@ -1395,12 +1589,13 @@ impl RelayRuntimeHandle { .rate_limit_req_pocket(subscription_id, filters, auth, rate_limit_context, now) } - pub(crate) async fn query_req_with_auth_report( + pub(crate) async fn query_req_with_auth_report_with_projection_context( &self, subscription_id: SubscriptionId, filters: Vec<PocketOwnedFilter>, search_present: bool, auth: &BaseAuthState, + projection: &RelayProjectionContext, ) -> Result<BaseRelayQueryReport, BaseRelayError> { let started_at = Instant::now(); let report = self.inner.query_req_with_auth_report( @@ -1409,6 +1604,7 @@ impl RelayRuntimeHandle { search_present, auth, )?; + let report = self.inner.project_query_report(report, projection, auth)?; if report.group_read_denied() { self.inner.metrics.record_group_read_denial(); } @@ -1437,11 +1633,12 @@ impl RelayRuntimeHandle { Ok(Some(pocket_event)) } - pub(crate) async fn fanout_event_offset( + pub(crate) async fn fanout_event_offset_with_projection_context( &self, offset: StoreOffset, subscriptions: &mut LiveSubscriptionSet, auth: &BaseAuthState, + projection: &RelayProjectionContext, ) -> Result<Vec<RuntimeRelayMessage>, BaseRelayError> { let pocket_event = self.inner.store.event_by_offset(offset.as_u64())?; let group_auth = GroupAuthContext::new(auth.authenticated_pubkeys().iter().cloned()); @@ -1449,12 +1646,21 @@ impl RelayRuntimeHandle { BaseRelay::group_read_gate_visible_to_auth(self.inner.groups.as_ref(), event, auth) .unwrap_or(false) })?; - Ok(subscriptions - .into_iter() - .map(|subscription_id| { - RuntimeRelayMessage::event(subscription_id, pocket_event.clone()) - }) - .collect()) + let mut messages = Vec::with_capacity(subscriptions.len()); + for subscription_id in subscriptions { + if let Some(projected) = self.inner.project_event_output( + &subscription_id, + projection, + RelayEventProjectionSource::LiveFanout { + store_offset: offset.as_u64(), + }, + &pocket_event, + auth, + )? { + messages.push(RuntimeRelayMessage::event(subscription_id, projected)); + } + } + Ok(messages) } pub async fn shutdown(&self) -> Result<BaseRelayShutdownReport, BaseRelayError> { @@ -2327,9 +2533,11 @@ impl Default for TangleShutdownSignal { mod tests { use super::{ BROAD_QUERY_TIME_WINDOW_SECONDS, EventAdmissionDecision, RelayEventAdmissionContext, - RelayEventStoredContext, RelayRuntime, RelayRuntimeHandle, RelayRuntimeHooks, - RuntimeClientMessage, TangleBroadQueryReason, TangleClientRateLimitContext, - TangleQueryClassification, TangleQueryClassifier, TangleRuntimeLimits, + RelayEventProjectionContext, RelayEventProjectionDecision, RelayEventProjectionSource, + RelayEventStoredContext, RelayProjectionContext, RelayRuntime, RelayRuntimeHandle, + RelayRuntimeHooks, RuntimeClientMessage, TangleBroadQueryReason, + TangleClientRateLimitContext, TangleQueryClassification, TangleQueryClassifier, + TangleRuntimeLimits, }; use crate::config::{BaseRelayRuntimeConfig, parse_base_relay_runtime_config_json}; use crate::event_bus::{TangleEventBus, TangleEventReceiveError, TangleEventReceiver}; @@ -2582,7 +2790,12 @@ mod tests { let offset = offsets.try_recv().expect("offset"); assert!(matches!( handle - .fanout_event_offset(offset, &mut subscriptions, &auth) + .fanout_event_offset_with_projection_context( + offset, + &mut subscriptions, + &auth, + &RelayProjectionContext::default(), + ) .await .expect("fanout") .as_slice(), @@ -2705,6 +2918,359 @@ mod tests { } #[tokio::test] + async fn runtime_projection_default_keeps_query_and_live_output() { + let root = temp_root("runtime-projection-default"); + let _ = std::fs::remove_dir_all(&root); + let handle = + RelayRuntimeHandle::new(RelayRuntime::open(runtime_config(&root, 8)).expect("runtime")); + let mut offsets = handle.subscribe_events().await; + let mut auth = handle.auth_state().await.expect("auth"); + let event = tangle_v2_event( + FixtureKey::Member, + 1_714_124_433, + 1, + Vec::new(), + "default projection", + ) + .expect("event"); + + assert_eq!( + handle + .handle_protocol_client_message_for_test( + ClientMessage::Event(event.clone()), + &mut auth, + UnixTimestamp::new(1_714_124_433) + ) + .await + .expect("event"), + vec![RelayMessage::Ok { + event_id: event.id().clone(), + accepted: true, + message: String::new() + }] + ); + let offset = offsets.try_recv().expect("offset"); + let query_sub = SubscriptionId::new("projection-default-query").expect("subscription"); + let report = handle + .query_req_with_auth_report_with_projection_context( + query_sub.clone(), + vec![pocket_filter(json!({"ids": [event.id().as_str()]}))], + false, + &auth, + &RelayProjectionContext::default(), + ) + .await + .expect("query"); + assert!(matches!( + report.into_messages().as_slice(), + [ + RuntimeRelayMessage::Event { + subscription_id, + event: found + }, + RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) + ] if subscription_id == &query_sub + && found.id().as_hex_string() == event.id().as_str() + && eose == &query_sub + )); + + let live_sub = SubscriptionId::new("projection-default-live").expect("subscription"); + let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); + subscriptions + .subscribe(live_sub.clone(), vec![pocket_filter(json!({"kinds": [1]}))]) + .expect("subscribe"); + assert!(matches!( + handle + .fanout_event_offset_with_projection_context( + offset, + &mut subscriptions, + &auth, + &RelayProjectionContext::default() + ) + .await + .expect("fanout") + .as_slice(), + [RuntimeRelayMessage::Event { + subscription_id, + event: found + }] if subscription_id == &live_sub && found.id().as_hex_string() == event.id().as_str() + )); + + let _ = std::fs::remove_dir_all(root); + } + + #[tokio::test] + async fn runtime_projection_can_suppress_historical_query_events() { + let root = temp_root("runtime-projection-query-suppress"); + let _ = std::fs::remove_dir_all(&root); + let hooks = Arc::new(ProjectingHooks::new( + "quiet", + ProjectionHookScope::Historical, + None, + RelayEventProjectionDecision::Suppress, + )); + let handle = RelayRuntimeHandle::new( + RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) + .expect("runtime"), + ); + let mut auth = handle.auth_state().await.expect("auth"); + let event = tangle_v2_event( + FixtureKey::Member, + 1_714_124_433, + 1, + Vec::new(), + "suppress query", + ) + .expect("event"); + assert_accepted_reply( + runtime_event_reply(&handle, event.clone(), &mut auth, 1_714_124_433).await, + &event, + ); + let subscription_id = SubscriptionId::new("query-suppressed").expect("subscription"); + + let report = handle + .query_req_with_auth_report_with_projection_context( + subscription_id.clone(), + vec![pocket_filter(json!({"ids": [event.id().as_str()]}))], + false, + &auth, + &RelayProjectionContext::named("quiet").expect("projection"), + ) + .await + .expect("query"); + assert_eq!( + report.into_messages(), + vec![RuntimeRelayMessage::from(RelayMessage::Eose( + subscription_id + ))] + ); + let contexts = hooks.contexts(); + assert_eq!(contexts.len(), 1); + assert_eq!(contexts[0].projection().identifier(), Some("quiet")); + assert_eq!( + contexts[0].source(), + RelayEventProjectionSource::HistoricalQuery + ); + assert_eq!(contexts[0].event().event_id(), event.id().as_str()); + + let _ = std::fs::remove_dir_all(root); + } + + #[tokio::test] + async fn runtime_projection_can_suppress_live_fanout_events() { + let root = temp_root("runtime-projection-live-suppress"); + let _ = std::fs::remove_dir_all(&root); + let hooks = Arc::new(ProjectingHooks::new( + "quiet-live", + ProjectionHookScope::Live, + None, + RelayEventProjectionDecision::Suppress, + )); + let handle = RelayRuntimeHandle::new( + RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) + .expect("runtime"), + ); + let mut offsets = handle.subscribe_events().await; + let mut auth = handle.auth_state().await.expect("auth"); + let event = tangle_v2_event( + FixtureKey::Member, + 1_714_124_433, + 1, + Vec::new(), + "suppress live", + ) + .expect("event"); + assert_accepted_reply( + runtime_event_reply(&handle, event.clone(), &mut auth, 1_714_124_433).await, + &event, + ); + let offset = offsets.try_recv().expect("offset"); + let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); + subscriptions + .subscribe( + SubscriptionId::new("live-suppressed").expect("subscription"), + vec![pocket_filter(json!({"kinds": [1]}))], + ) + .expect("subscribe"); + + assert!( + handle + .fanout_event_offset_with_projection_context( + offset, + &mut subscriptions, + &auth, + &RelayProjectionContext::named("quiet-live").expect("projection") + ) + .await + .expect("fanout") + .is_empty() + ); + let contexts = hooks.contexts(); + assert_eq!(contexts.len(), 1); + assert_eq!(contexts[0].projection().identifier(), Some("quiet-live")); + assert_eq!( + contexts[0].source(), + RelayEventProjectionSource::LiveFanout { + store_offset: offset.as_u64() + } + ); + assert_eq!(contexts[0].event().event_id(), event.id().as_str()); + + let _ = std::fs::remove_dir_all(root); + } + + #[tokio::test] + async fn runtime_projection_replaces_with_existing_stored_events_only() { + let root = temp_root("runtime-projection-replace"); + let _ = std::fs::remove_dir_all(&root); + let hooks = Arc::new(ProjectingHooks::new( + "replace", + ProjectionHookScope::Historical, + Some("source"), + RelayEventProjectionDecision::Emit, + )); + let handle = RelayRuntimeHandle::new( + RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) + .expect("runtime"), + ); + let mut offsets = handle.subscribe_events().await; + let mut auth = handle.auth_state().await.expect("auth"); + let source = tangle_v2_event(FixtureKey::Member, 1_714_124_433, 1, Vec::new(), "source") + .expect("source"); + let replacement = tangle_v2_event( + FixtureKey::Admin, + 1_714_124_434, + 1, + Vec::new(), + "replacement", + ) + .expect("replacement"); + assert_accepted_reply( + runtime_event_reply(&handle, source.clone(), &mut auth, 1_714_124_433).await, + &source, + ); + let _source_offset = offsets.try_recv().expect("source offset"); + assert_accepted_reply( + runtime_event_reply(&handle, replacement.clone(), &mut auth, 1_714_124_434).await, + &replacement, + ); + let replacement_offset = offsets.try_recv().expect("replacement offset"); + hooks.set_decision(RelayEventProjectionDecision::replace_with_stored_offset( + replacement_offset.as_u64(), + )); + let subscription_id = SubscriptionId::new("replace-existing").expect("subscription"); + + let report = handle + .query_req_with_auth_report_with_projection_context( + subscription_id.clone(), + vec![pocket_filter(json!({"ids": [source.id().as_str()]}))], + false, + &auth, + &RelayProjectionContext::named("replace").expect("projection"), + ) + .await + .expect("query"); + assert!(matches!( + report.into_messages().as_slice(), + [ + RuntimeRelayMessage::Event { + subscription_id: delivered, + event + }, + RuntimeRelayMessage::Protocol(RelayMessage::Eose(eose)) + ] if delivered == &subscription_id + && event.id().as_hex_string() == replacement.id().as_str() + && eose == &subscription_id + )); + + hooks.set_decision(RelayEventProjectionDecision::replace_with_stored_offset( + u64::MAX, + )); + let missing_subscription = SubscriptionId::new("replace-missing").expect("subscription"); + let report = handle + .query_req_with_auth_report_with_projection_context( + missing_subscription.clone(), + vec![pocket_filter(json!({"ids": [source.id().as_str()]}))], + false, + &auth, + &RelayProjectionContext::named("replace").expect("projection"), + ) + .await + .expect("missing query"); + assert_eq!( + report.into_messages(), + vec![RuntimeRelayMessage::from(RelayMessage::Eose( + missing_subscription + ))] + ); + + let _ = std::fs::remove_dir_all(root); + } + + #[tokio::test] + async fn runtime_projection_runs_after_group_read_gates() { + let root = temp_root("runtime-projection-group-gate"); + let _ = std::fs::remove_dir_all(&root); + let hooks = Arc::new(ProjectingHooks::new( + "group-gate", + ProjectionHookScope::Historical, + None, + RelayEventProjectionDecision::Suppress, + )); + let handle = RelayRuntimeHandle::new( + RelayRuntime::open_with_hooks(runtime_config(&root, 8), hooks.clone()) + .expect("runtime"), + ); + let mut owner_auth = + authenticated_runtime_state(&handle, FixtureKey::Owner, "group-gate-owner", 120).await; + let public_auth = handle.auth_state().await.expect("public auth"); + let create = + tangle_v2_group_create_event(FixtureKey::Owner, "ProjectionPrivate", 121, &["private"]) + .expect("create"); + assert_accepted_reply( + runtime_event_reply(&handle, create.clone(), &mut owner_auth, 121).await, + &create, + ); + let private_event = tangle_v2_group_event( + FixtureKey::Owner, + "ProjectionPrivate", + 122, + 1, + "private projection", + ) + .expect("private event"); + assert_accepted_reply( + runtime_event_reply(&handle, private_event.clone(), &mut owner_auth, 122).await, + &private_event, + ); + let subscription_id = SubscriptionId::new("group-gate").expect("subscription"); + + let report = handle + .query_req_with_auth_report_with_projection_context( + subscription_id.clone(), + vec![pocket_filter(json!({ + "kinds": [1], + "#h": ["ProjectionPrivate"] + }))], + false, + &public_auth, + &RelayProjectionContext::named("group-gate").expect("projection"), + ) + .await + .expect("query"); + assert_eq!( + report.into_messages(), + vec![RuntimeRelayMessage::from(RelayMessage::Closed { + subscription_id, + message: "auth-required: authentication required to read group events".to_owned() + })] + ); + assert!(hooks.contexts().is_empty()); + + let _ = std::fs::remove_dir_all(root); + } + + #[tokio::test] async fn runtime_rate_limits_event_pubkeys_before_storage() { let root = temp_root("runtime-event-rate-limit"); let _ = std::fs::remove_dir_all(&root); @@ -3969,7 +4535,12 @@ mod tests { let mut generated_kinds = BTreeSet::new(); for offset in generated_offsets { let messages = handle - .fanout_event_offset(offset, &mut subscriptions, &auth) + .fanout_event_offset_with_projection_context( + offset, + &mut subscriptions, + &auth, + &RelayProjectionContext::default(), + ) .await .expect("fanout"); assert!(matches!( @@ -4653,12 +5224,22 @@ mod tests { let mut member_fanout_count = 0; for offset in &published_offsets { let public_replies = handle - .fanout_event_offset(*offset, &mut public_subscriptions, &public_auth) + .fanout_event_offset_with_projection_context( + *offset, + &mut public_subscriptions, + &public_auth, + &RelayProjectionContext::default(), + ) .await .expect("public fanout"); assert!(public_replies.is_empty()); let member_replies = handle - .fanout_event_offset(*offset, &mut member_subscriptions, &member_auth) + .fanout_event_offset_with_projection_context( + *offset, + &mut member_subscriptions, + &member_auth, + &RelayProjectionContext::default(), + ) .await .expect("member fanout"); for reply in member_replies { @@ -4870,6 +5451,85 @@ mod tests { } } + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + enum ProjectionHookScope { + Historical, + Live, + } + + struct ProjectingHooks { + projection_identifier: &'static str, + scope: ProjectionHookScope, + source_content: Option<&'static str>, + decision: Mutex<RelayEventProjectionDecision>, + contexts: Mutex<Vec<RelayEventProjectionContext>>, + } + + impl ProjectingHooks { + fn new( + projection_identifier: &'static str, + scope: ProjectionHookScope, + source_content: Option<&'static str>, + decision: RelayEventProjectionDecision, + ) -> Self { + Self { + projection_identifier, + scope, + source_content, + decision: Mutex::new(decision), + contexts: Mutex::new(Vec::new()), + } + } + + fn set_decision(&self, decision: RelayEventProjectionDecision) { + *self.decision.lock().expect("decision") = decision; + } + + fn contexts(&self) -> Vec<RelayEventProjectionContext> { + self.contexts.lock().expect("contexts").clone() + } + + fn scope_matches(&self, source: RelayEventProjectionSource) -> bool { + matches!( + (self.scope, source), + ( + ProjectionHookScope::Historical, + RelayEventProjectionSource::HistoricalQuery + ) | ( + ProjectionHookScope::Live, + RelayEventProjectionSource::LiveFanout { .. } + ) + ) + } + + fn content_matches(&self, context: &RelayEventProjectionContext) -> bool { + match self.source_content { + Some(content) => context.event().content() == content, + None => true, + } + } + } + + impl RelayRuntimeHooks for ProjectingHooks { + fn project_event( + &self, + context: &RelayEventProjectionContext, + ) -> RelayEventProjectionDecision { + self.contexts + .lock() + .expect("contexts") + .push(context.clone()); + if context.projection().identifier() == Some(self.projection_identifier) + && self.scope_matches(context.source()) + && self.content_matches(context) + { + *self.decision.lock().expect("decision") + } else { + RelayEventProjectionDecision::Emit + } + } + } + fn runtime_config_with_group_policy( root: &Path, per_connection_outbound_queue: usize, diff --git a/crates/tangle_runtime/src/session.rs b/crates/tangle_runtime/src/session.rs @@ -13,8 +13,8 @@ use crate::{ }, resource_limits::{RelayResourceLimiter, RelaySubscriptionPermit}, runtime::{ - RelayRuntimeHandle, TangleClientMessageMetricKind, TangleClientRateLimitContext, - TangleRuntimeLimits, + RelayProjectionContext, RelayRuntimeHandle, TangleClientMessageMetricKind, + TangleClientRateLimitContext, TangleRuntimeLimits, }, }; use axum::extract::ws::{CloseFrame, Message, Utf8Bytes, WebSocket}; @@ -43,10 +43,39 @@ pub struct TangleWebSocketSession { resource_limiter: Option<RelayResourceLimiter>, subscription_permits: BTreeMap<SubscriptionId, RelaySubscriptionPermit>, events: TangleEventReceiver, + projection: RelayProjectionContext, } static NEXT_TANGLE_CONNECTION_ID: AtomicU64 = AtomicU64::new(1); +#[derive(Debug, Clone, Default)] +pub struct TangleWebSocketSessionOptions { + peer_ip: Option<IpAddr>, + resource_limiter: Option<RelayResourceLimiter>, + projection: RelayProjectionContext, +} + +impl TangleWebSocketSessionOptions { + pub fn new() -> Self { + Self::default() + } + + pub fn with_peer_ip(mut self, peer_ip: Option<IpAddr>) -> Self { + self.peer_ip = peer_ip; + self + } + + pub fn with_resource_limiter(mut self, resource_limiter: Option<RelayResourceLimiter>) -> Self { + self.resource_limiter = resource_limiter; + self + } + + pub fn with_projection_context(mut self, projection: RelayProjectionContext) -> Self { + self.projection = projection; + self + } +} + impl TangleWebSocketSession { pub fn new( limits: TangleRuntimeLimits, @@ -78,6 +107,26 @@ impl TangleWebSocketSession { peer_ip: Option<IpAddr>, resource_limiter: Option<RelayResourceLimiter>, ) -> Result<Self, BaseRelayError> { + Self::new_with_options( + limits, + shutdown, + runtime, + auth, + events, + TangleWebSocketSessionOptions::new() + .with_peer_ip(peer_ip) + .with_resource_limiter(resource_limiter), + ) + } + + pub fn new_with_options( + limits: TangleRuntimeLimits, + shutdown: watch::Receiver<bool>, + runtime: RelayRuntimeHandle, + auth: BaseAuthState, + events: TangleEventReceiver, + options: TangleWebSocketSessionOptions, + ) -> Result<Self, BaseRelayError> { let outbound_queue_capacity = limits.outbound_queue_capacity(); let (sender, receiver) = mpsc::channel(outbound_queue_capacity); let subscriptions = LiveSubscriptionSet::new( @@ -86,7 +135,7 @@ impl TangleWebSocketSession { )?; Ok(Self { connection_id: NEXT_TANGLE_CONNECTION_ID.fetch_add(1, Ordering::Relaxed), - peer_ip, + peer_ip: options.peer_ip, connected_at: Instant::now(), outbound: TangleOutboundSender { sender, @@ -98,9 +147,10 @@ impl TangleWebSocketSession { limits, auth, subscriptions, - resource_limiter, + resource_limiter: options.resource_limiter, subscription_permits: BTreeMap::new(), events, + projection: options.projection, }) } @@ -219,7 +269,12 @@ impl TangleWebSocketSession { let runtime = self.runtime.clone(); let auth = self.auth.clone(); let replies = match runtime - .fanout_event_offset(offset, &mut self.subscriptions, &auth) + .fanout_event_offset_with_projection_context( + offset, + &mut self.subscriptions, + &auth, + &self.projection, + ) .await { Ok(replies) => replies, @@ -384,11 +439,12 @@ impl TangleWebSocketSession { } let report = self .runtime - .query_req_with_auth_report( + .query_req_with_auth_report_with_projection_context( subscription_id.clone(), filters.clone(), search_present, &self.auth, + &self.projection, ) .await?; let closes_subscription = report.group_read_denied(); diff --git a/crates/tangle_runtime/tests/phase2_acceptance_targets.rs b/crates/tangle_runtime/tests/phase2_acceptance_targets.rs @@ -1875,7 +1875,7 @@ fn req_count_and_live_fanout_share_one_group_read_gate() { runtime .matches("BaseRelay::group_read_gate_visible_to_auth") .count(), - 2 + 3 ); assert!(!relay_core.contains("fn event_visible_to_auth(")); assert!(!relay_core.contains("fn pocket_event_visible_to_auth("));