tangle


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

commit 70f9b47b8e91b11ab4fc5a85a000d4337a0d2e42
parent 1915c2e9946d860ccf9ccd405ec6762c21c76c69
Author: triesap <tyson@radroots.org>
Date:   Fri, 26 Jun 2026 19:34:24 +0000

runtime: add live projection candidates

Diffstat:
Mcrates/tangle_runtime/src/runtime.rs | 267++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
1 file changed, 246 insertions(+), 21 deletions(-)

diff --git a/crates/tangle_runtime/src/runtime.rs b/crates/tangle_runtime/src/runtime.rs @@ -97,6 +97,13 @@ pub trait RelayRuntimeHooks: Send + Sync { RelayProjectionQueryPlan::default() } + fn live_projection_candidates( + &self, + _context: &RelayLiveProjectionContext, + ) -> Vec<RelayLiveProjectionCandidate> { + Vec::new() + } + fn project_event( &self, _context: &RelayEventProjectionContext, @@ -275,6 +282,54 @@ pub enum RelayRequestedKinds { } #[derive(Debug, Clone, PartialEq, Eq)] +pub struct RelayLiveProjectionContext { + projection: RelayProjectionContext, + source_store_offset: u64, + event: RelayEventContext, +} + +impl RelayLiveProjectionContext { + pub fn new( + projection: RelayProjectionContext, + source_store_offset: u64, + event: RelayEventContext, + ) -> Self { + Self { + projection, + source_store_offset, + event, + } + } + + pub fn projection(&self) -> &RelayProjectionContext { + &self.projection + } + + pub fn source_store_offset(&self) -> u64 { + self.source_store_offset + } + + pub fn event(&self) -> &RelayEventContext { + &self.event + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct RelayLiveProjectionCandidate { + store_offset: u64, +} + +impl RelayLiveProjectionCandidate { + pub fn stored_offset(store_offset: u64) -> Self { + Self { store_offset } + } + + pub fn store_offset(self) -> u64 { + self.store_offset + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] pub struct RelayEventContext { event_id: String, pubkey: String, @@ -425,6 +480,15 @@ struct RelayEventProjectionRequest<'a> { filter: &'a PocketFilter, } +struct RelayLiveProjectionDelivery<'a> { + subscriptions: &'a LiveSubscriptionSet, + auth: &'a BaseAuthState, + projection: &'a RelayProjectionContext, + group_auth: &'a GroupAuthContext, + delivered: &'a mut BTreeSet<(SubscriptionId, String)>, + messages: &'a mut Vec<RuntimeRelayMessage>, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct RelayEventProjectionContext { subscription_id: SubscriptionId, @@ -1908,34 +1972,80 @@ impl RelayRuntimeHandle { ) -> 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()); - let subscriptions = subscriptions.fanout(&pocket_event, &group_auth, |event, auth| { - BaseRelay::group_read_gate_visible_to_auth(self.inner.groups.as_ref(), event, auth) - .unwrap_or(false) - })?; - let mut messages = Vec::with_capacity(subscriptions.len()); + let mut messages = Vec::new(); + let mut delivered = BTreeSet::new(); + let mut delivery = RelayLiveProjectionDelivery { + subscriptions, + auth, + projection, + group_auth: &group_auth, + delivered: &mut delivered, + messages: &mut messages, + }; + self.fanout_projected_live_event(&pocket_event, offset.as_u64(), &mut delivery)?; + let context = RelayLiveProjectionContext::new( + projection.clone(), + offset.as_u64(), + RelayEventContext::from_pocket_event(&pocket_event)?, + ); + for candidate in self.inner.hooks.live_projection_candidates(&context) { + let Ok(candidate_event) = self.inner.store.event_by_offset(candidate.store_offset()) + else { + continue; + }; + self.fanout_projected_live_event( + &candidate_event, + candidate.store_offset(), + &mut delivery, + )?; + } + Ok(messages) + } + + fn fanout_projected_live_event( + &self, + event: &PocketEvent, + store_offset: u64, + delivery: &mut RelayLiveProjectionDelivery<'_>, + ) -> Result<(), BaseRelayError> { + let subscriptions = + delivery + .subscriptions + .fanout(event, delivery.group_auth, |event, auth| { + BaseRelay::group_read_gate_visible_to_auth( + self.inner.groups.as_ref(), + event, + auth, + ) + .unwrap_or(false) + })?; for matched in subscriptions { + let subscription_id = matched.subscription_id().clone(); let matched_filter = matched.matched_filter_context(); if let Some(projected) = self.inner .project_event_output(RelayEventProjectionRequest { subscription_id: matched.subscription_id(), - projection, - source: RelayEventProjectionSource::LiveFanout { - store_offset: offset.as_u64(), - }, - event: &pocket_event, - auth, + projection: delivery.projection, + source: RelayEventProjectionSource::LiveFanout { store_offset }, + event, + auth: delivery.auth, matched_filter: &matched_filter, filter: matched.filter(), })? { - messages.push(RuntimeRelayMessage::event( - matched.into_subscription_id(), - projected, - )); + let event_id = projected.id().as_hex_string(); + if delivery + .delivered + .insert((subscription_id.clone(), event_id)) + { + delivery + .messages + .push(RuntimeRelayMessage::event(subscription_id, projected)); + } } } - Ok(messages) + Ok(()) } pub async fn shutdown(&self) -> Result<BaseRelayShutdownReport, BaseRelayError> { @@ -2809,11 +2919,11 @@ mod tests { use super::{ BROAD_QUERY_TIME_WINDOW_SECONDS, EventAdmissionDecision, RelayEventAdmissionContext, RelayEventProjectionContext, RelayEventProjectionDecision, RelayEventProjectionSource, - RelayEventStoredContext, RelayProjectionContext, RelayProjectionQueryPlan, - RelayQueryProjectionContext, RelayRequestedKinds, RelayRuntime, RelayRuntimeHandle, - RelayRuntimeHooks, RuntimeClientMessage, TangleBroadQueryReason, - TangleClientRateLimitContext, TangleQueryClassification, TangleQueryClassifier, - TangleRuntimeLimits, + RelayEventStoredContext, RelayLiveProjectionCandidate, RelayLiveProjectionContext, + RelayProjectionContext, RelayProjectionQueryPlan, RelayQueryProjectionContext, + RelayRequestedKinds, 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}; @@ -3649,6 +3759,91 @@ mod tests { } #[tokio::test] + async fn runtime_projection_live_candidates_match_candidate_filter() { + let root = temp_root("runtime-projection-live-candidate"); + let _ = std::fs::remove_dir_all(&root); + let hooks = Arc::new(ProjectingHooks::new( + "live-candidate", + ProjectionHookScope::Live, + None, + 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, + 30078, + Vec::new(), + "source", + ) + .expect("source"); + let candidate = + tangle_v2_event(FixtureKey::Admin, 1_714_124_434, 1, Vec::new(), "candidate") + .expect("candidate"); + 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, candidate.clone(), &mut auth, 1_714_124_434).await, + &candidate, + ); + let candidate_offset = offsets.try_recv().expect("candidate offset"); + hooks.set_live_candidates(vec![RelayLiveProjectionCandidate::stored_offset( + candidate_offset.as_u64(), + )]); + let mut subscriptions = LiveSubscriptionSet::new(8, 64).expect("subscriptions"); + let subscription_id = SubscriptionId::new("candidate-note").expect("subscription"); + subscriptions + .subscribe( + subscription_id.clone(), + vec![pocket_filter(json!({"kinds": [1]}))], + ) + .expect("subscribe"); + + let messages = handle + .fanout_event_offset_with_projection_context( + source_offset, + &mut subscriptions, + &auth, + &RelayProjectionContext::named("live-candidate").expect("projection"), + ) + .await + .expect("fanout"); + + assert!(matches!( + messages.as_slice(), + [RuntimeRelayMessage::Event { + subscription_id: delivered, + event + }] if delivered == &subscription_id + && event.id().as_hex_string() == candidate.id().as_str() + )); + let live_contexts = hooks.live_contexts(); + assert_eq!(live_contexts.len(), 1); + assert_eq!( + live_contexts[0].source_store_offset(), + source_offset.as_u64() + ); + assert_eq!(live_contexts[0].event().event_id(), source.id().as_str()); + let contexts = hooks.contexts(); + assert_eq!(contexts.len(), 1); + assert_eq!(contexts[0].event().event_id(), candidate.id().as_str()); + assert_eq!( + contexts[0].matched_filter().requested_kinds(), + &RelayRequestedKinds::Explicit(BTreeSet::from([1])) + ); + + 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); @@ -5904,7 +6099,9 @@ mod tests { source_content: Option<&'static str>, decision: Mutex<RelayEventProjectionDecision>, query_plan: Mutex<RelayProjectionQueryPlan>, + live_candidates: Mutex<Vec<RelayLiveProjectionCandidate>>, contexts: Mutex<Vec<RelayEventProjectionContext>>, + live_contexts: Mutex<Vec<RelayLiveProjectionContext>>, query_contexts: Mutex<Vec<RelayQueryProjectionContext>>, } @@ -5921,7 +6118,9 @@ mod tests { source_content, decision: Mutex::new(decision), query_plan: Mutex::new(RelayProjectionQueryPlan::default()), + live_candidates: Mutex::new(Vec::new()), contexts: Mutex::new(Vec::new()), + live_contexts: Mutex::new(Vec::new()), query_contexts: Mutex::new(Vec::new()), } } @@ -5934,10 +6133,18 @@ mod tests { *self.query_plan.lock().expect("query plan") = plan; } + fn set_live_candidates(&self, candidates: Vec<RelayLiveProjectionCandidate>) { + *self.live_candidates.lock().expect("live candidates") = candidates; + } + fn contexts(&self) -> Vec<RelayEventProjectionContext> { self.contexts.lock().expect("contexts").clone() } + fn live_contexts(&self) -> Vec<RelayLiveProjectionContext> { + self.live_contexts.lock().expect("live contexts").clone() + } + fn query_contexts(&self) -> Vec<RelayQueryProjectionContext> { self.query_contexts.lock().expect("query contexts").clone() } @@ -5972,6 +6179,24 @@ mod tests { *self.query_plan.lock().expect("query plan") } + fn live_projection_candidates( + &self, + context: &RelayLiveProjectionContext, + ) -> Vec<RelayLiveProjectionCandidate> { + self.live_contexts + .lock() + .expect("live contexts") + .push(context.clone()); + if context.projection().identifier() == Some(self.projection_identifier) { + self.live_candidates + .lock() + .expect("live candidates") + .clone() + } else { + Vec::new() + } + } + fn project_event( &self, context: &RelayEventProjectionContext,