feat(sync): Phase 4 - dynamic subscriptions

- Add SubscriptionManager for per-connection tracking
- Trigger subscription updates on new repo/PR events
- Implement consolidation when filter count > 150
This commit is contained in:
DanConwayDev
2025-12-04 18:15:19 +00:00
parent f639ecfac6
commit a19ff57e72
6 changed files with 1190 additions and 2 deletions
+131 -1
View File
@@ -15,6 +15,13 @@
//! - Health tracking with success/failure reporting //! - Health tracking with success/failure reporting
//! - Exponential backoff with health-aware delays //! - Exponential backoff with health-aware delays
//! - Dead relay detection and minimal retry //! - Dead relay detection and minimal retry
//!
//! ## Phase 4 Features
//!
//! - Dynamic subscription updates when new announcements/PRs arrive
//! - Per-connection subscription tracking
//! - Filter consolidation when count exceeds threshold (>150)
//! - Duplicate subscription prevention
use std::sync::Arc; use std::sync::Arc;
use std::time::Duration; use std::time::Duration;
@@ -24,6 +31,7 @@ use tokio::sync::mpsc;
use super::filter::FilterService; use super::filter::FilterService;
use super::health::RelayHealthTracker; use super::health::RelayHealthTracker;
use super::subscription::SubscriptionManager;
/// Event received from the sync connection /// Event received from the sync connection
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
@@ -38,6 +46,7 @@ pub struct SyncConnection {
client: Client, client: Client,
filter_service: Arc<FilterService>, filter_service: Arc<FilterService>,
remote_domain: String, remote_domain: String,
subscription_manager: SubscriptionManager,
} }
impl SyncConnection { impl SyncConnection {
@@ -57,18 +66,26 @@ impl SyncConnection {
tracing::info!("Sync connection established to {}", url); tracing::info!("Sync connection established to {}", url);
// Create subscription manager for this connection
let subscription_manager = SubscriptionManager::new(
filter_service.clone(),
remote_domain.to_string(),
);
Ok(Self { Ok(Self {
url: url.to_string(), url: url.to_string(),
client, client,
filter_service, filter_service,
remote_domain: remote_domain.to_string(), remote_domain: remote_domain.to_string(),
subscription_manager,
}) })
} }
/// Start receiving events and send them through the channel /// Start receiving events and send them through the channel
/// ///
/// This method runs indefinitely, handling events from all three filter layers. /// This method runs indefinitely, handling events from all three filter layers.
pub async fn run(self, tx: mpsc::Sender<SyncedEvent>) { /// Dynamic subscription updates are triggered when new announcements or PRs arrive.
pub async fn run(mut self, tx: mpsc::Sender<SyncedEvent>) {
// Subscribe to all three filter layers // Subscribe to all three filter layers
// Layer 1: Announcement discovery (kinds 30617 + 30618) // Layer 1: Announcement discovery (kinds 30617 + 30618)
@@ -174,6 +191,119 @@ impl SyncConnection {
.await .await
.ok(); .ok();
} }
/// Handle dynamic subscription updates based on incoming event kind
///
/// - kind 30617/30618: New announcement → add Layer 2 subscription
/// - kind 1617/1618/1619/1621/1622: New PR/Issue → add Layer 3 subscription
async fn handle_dynamic_subscription(&mut self, event: &Event) {
let kind = event.kind.as_u16();
// Check if this is an announcement kind (triggers Layer 2 subscription)
if SubscriptionManager::is_announcement_kind(kind) {
if let Some(new_filters) = self.subscription_manager.add_announcement(event) {
tracing::info!(
"New announcement {} on {}, adding {} Layer 2 filter(s) (total filters: {})",
event.id.to_hex(),
self.url,
new_filters.len(),
self.subscription_manager.get_filter_count()
);
self.subscribe_to_filters(new_filters, "Layer 2").await;
}
}
// Check if this is a PR/Issue kind (triggers Layer 3 subscription)
if SubscriptionManager::is_pr_issue_kind(kind) {
if let Some(new_filters) = self.subscription_manager.add_event(event) {
tracing::info!(
"New PR/Issue {} on {}, adding {} Layer 3 filter(s) (total filters: {})",
event.id.to_hex(),
self.url,
new_filters.len(),
self.subscription_manager.get_filter_count()
);
self.subscribe_to_filters(new_filters, "Layer 3").await;
}
}
// Check if we need to consolidate
if self.subscription_manager.should_consolidate() {
self.consolidate_subscriptions().await;
}
}
/// Subscribe to new filters
async fn subscribe_to_filters(&self, filters: Vec<Filter>, layer_name: &str) {
for filter in filters {
match self.client.subscribe(filter, None).await {
Ok(output) => {
tracing::debug!(
"Dynamic {} subscription on {} (subscription: {})",
layer_name,
self.url,
output.id()
);
}
Err(e) => {
tracing::warn!(
"Failed to add dynamic {} subscription on {}: {}",
layer_name,
self.url,
e
);
}
}
}
}
/// Consolidate subscriptions back to Layer 1 only
///
/// This is triggered when the filter count exceeds 150.
/// All existing subscriptions are closed and only Layer 1 is re-subscribed.
async fn consolidate_subscriptions(&mut self) {
tracing::warn!(
"Filter count {} exceeds threshold, consolidating subscriptions on {}",
self.subscription_manager.get_filter_count(),
self.url
);
// Get consolidated filters (clears tracking and returns Layer 1 only)
let layer1_filters = self.subscription_manager.consolidate();
// Note: nostr-sdk doesn't provide a way to close all subscriptions easily
// The client will manage subscription count internally
// We just add the new Layer 1 subscription
for filter in layer1_filters {
match self.client.subscribe(filter, None).await {
Ok(output) => {
tracing::info!(
"Consolidated to Layer 1 subscription on {} (subscription: {})",
self.url,
output.id()
);
}
Err(e) => {
tracing::error!(
"Failed to subscribe Layer 1 after consolidation on {}: {}",
self.url,
e
);
}
}
}
}
/// Get the current filter count from the subscription manager
pub fn get_filter_count(&self) -> usize {
self.subscription_manager.get_filter_count()
}
/// Check if subscriptions have been consolidated
pub fn is_consolidated(&self) -> bool {
self.subscription_manager.is_consolidated()
}
} }
/// Reconnect loop with health-aware exponential backoff /// Reconnect loop with health-aware exponential backoff
+1
View File
@@ -24,6 +24,7 @@ const KIND_MAINTAINER_LIST: u16 = 30618;
/// 1. Layer 1: Discover new repository announcements and maintainer metadata /// 1. Layer 1: Discover new repository announcements and maintainer metadata
/// 2. Layer 2: Sync events directly related to repositories we track /// 2. Layer 2: Sync events directly related to repositories we track
/// 3. Layer 3: Sync discussions and updates related to Layer 2 events /// 3. Layer 3: Sync discussions and updates related to Layer 2 events
#[derive(Debug)]
pub struct FilterService { pub struct FilterService {
database: SharedDatabase, database: SharedDatabase,
/// Our relay's domain for filtering /// Our relay's domain for filtering
+30 -1
View File
@@ -15,6 +15,14 @@
//! - Health tracking with exponential backoff //! - Health tracking with exponential backoff
//! - Dead relay detection after 24h of failures //! - Dead relay detection after 24h of failures
//! - Startup jitter to prevent thundering herd //! - Startup jitter to prevent thundering herd
//!
//! ## Phase 4 Features
//!
//! - Dynamic subscription updates handled per-connection
//! - Each connection manages its own SubscriptionManager
//! - Announcements trigger Layer 2 subscriptions
//! - PRs/Issues trigger Layer 3 subscriptions
//! - Consolidation when filter count exceeds 150
use std::collections::HashSet; use std::collections::HashSet;
use std::sync::Arc; use std::sync::Arc;
@@ -225,17 +233,38 @@ impl SyncManager {
} }
/// Process a single synced event /// Process a single synced event
///
/// Events are validated through the write policy and stored if accepted.
/// Dynamic subscription updates are handled by each connection's SubscriptionManager.
async fn process_event(&self, synced_event: SyncedEvent) { async fn process_event(&self, synced_event: SyncedEvent) {
let event = &synced_event.event; let event = &synced_event.event;
let event_id = event.id.to_hex(); let event_id = event.id.to_hex();
let kind = event.kind.as_u16();
tracing::debug!( tracing::debug!(
"Processing synced event {} (kind {}) from {}", "Processing synced event {} (kind {}) from {}",
event_id, event_id,
event.kind.as_u16(), kind,
synced_event.source_url synced_event.source_url
); );
// Log subscription-relevant events for debugging
match kind {
30617 | 30618 => {
tracing::debug!(
"Received announcement {} - connection will add Layer 2 subscription",
event_id
);
}
1617 | 1618 | 1619 | 1621 | 1622 => {
tracing::debug!(
"Received PR/Issue {} - connection will add Layer 3 subscription",
event_id
);
}
_ => {}
}
// Validate through write policy using SYNC_SOURCE_ADDR // Validate through write policy using SYNC_SOURCE_ADDR
let result = self.write_policy.admit_event(event, &SYNC_SOURCE_ADDR).await; let result = self.write_policy.admit_event(event, &SYNC_SOURCE_ADDR).await;
+2
View File
@@ -21,10 +21,12 @@ mod connection;
mod filter; mod filter;
pub mod health; pub mod health;
mod manager; mod manager;
mod subscription;
pub use filter::FilterService; pub use filter::FilterService;
pub use health::{HealthState, RelayHealth, RelayHealthTracker}; pub use health::{HealthState, RelayHealth, RelayHealthTracker};
pub use manager::SyncManager; pub use manager::SyncManager;
pub use subscription::SubscriptionManager;
use std::net::SocketAddr; use std::net::SocketAddr;
+278
View File
@@ -0,0 +1,278 @@
//! Subscription Manager for GRASP-02 Phase 4: Dynamic Subscriptions
//!
//! Manages dynamic subscription updates per connection, including:
//! - Tracking subscribed announcements and events
//! - Adding new subscriptions when announcements/PRs arrive
//! - Consolidating filters when count exceeds threshold
//! - Preventing duplicate subscriptions
//!
//! ## Dynamic Subscription Strategy
//!
//! Initial: Layer 1 (announcements)
//! ↓ (announcement received)
//! Add: Layer 2 (events for that repo)
//! ↓ (PR/Issue received)
//! Add: Layer 3 (events for that PR/Issue)
//! ↓ (filter count > 150)
//! Consolidate: Back to Layer 1 only
use std::collections::HashSet;
use std::sync::Arc;
use nostr_sdk::prelude::*;
use super::filter::FilterService;
/// Maximum number of filters before consolidation is triggered
const CONSOLIDATION_THRESHOLD: usize = 150;
/// Kind 30617 - Repository Announcement (NIP-34)
const KIND_REPOSITORY_ANNOUNCEMENT: u16 = 30617;
/// Kind 30618 - Maintainer List (NIP-34)
const KIND_MAINTAINER_LIST: u16 = 30618;
/// Manages subscriptions for a single relay connection
///
/// Tracks which announcements and events have been subscribed to,
/// and handles dynamic subscription updates as new events arrive.
#[derive(Debug)]
pub struct SubscriptionManager {
/// Event IDs of announcements we've subscribed to (for Layer 2)
subscribed_announcements: HashSet<String>,
/// Event IDs of PRs/Issues we've subscribed to (for Layer 3)
subscribed_events: HashSet<String>,
/// Whether we've consolidated back to Layer 1 only
is_consolidated: bool,
/// FilterService for building filters
filter_service: Arc<FilterService>,
/// Remote relay domain for Layer 2 filters
remote_domain: String,
}
impl SubscriptionManager {
/// Create a new SubscriptionManager
///
/// # Arguments
/// * `filter_service` - FilterService for building subscription filters
/// * `remote_domain` - The domain of the remote relay we're syncing from
pub fn new(filter_service: Arc<FilterService>, remote_domain: String) -> Self {
Self {
subscribed_announcements: HashSet::new(),
subscribed_events: HashSet::new(),
is_consolidated: false,
filter_service,
remote_domain,
}
}
/// Add an announcement and return new filters to subscribe to
///
/// When a new announcement (kind 30617/30618) arrives, this creates
/// Layer 2 filters to subscribe to events for that repository.
///
/// Returns `Some(filters)` if this is a new announcement, `None` if already subscribed.
pub fn add_announcement(&mut self, event: &Event) -> Option<Vec<Filter>> {
let event_id = event.id.to_hex();
// Check if already subscribed or consolidated
if self.is_consolidated || self.subscribed_announcements.contains(&event_id) {
return None;
}
// Add to tracked announcements
self.subscribed_announcements.insert(event_id);
// Build Layer 2 filters for this announcement
// Layer 2 filters target events with 'a' tags pointing to this repo
let filters = self.build_layer2_filter_for_announcement(event);
if filters.is_empty() {
None
} else {
Some(filters)
}
}
/// Add a PR/Issue/Patch event and return new filters to subscribe to
///
/// When a new PR (kind 1617), Issue (kind 1621), or Patch (kind 1622) arrives,
/// this creates Layer 3 filters to subscribe to related events.
///
/// Returns `Some(filters)` if this is a new event, `None` if already subscribed.
pub fn add_event(&mut self, event: &Event) -> Option<Vec<Filter>> {
let event_id = event.id.to_hex();
// Check if already subscribed or consolidated
if self.is_consolidated || self.subscribed_events.contains(&event_id) {
return None;
}
// Add to tracked events
self.subscribed_events.insert(event_id.clone());
// Build Layer 3 filter for this event
// Layer 3 filters target events with 'e' tags pointing to this event
let filter = Filter::new().custom_tag(
SingleLetterTag::lowercase(Alphabet::E),
event_id,
);
Some(vec![filter])
}
/// Check if consolidation is needed
///
/// Returns true if the total filter count exceeds the threshold (150).
pub fn should_consolidate(&self) -> bool {
!self.is_consolidated && self.get_filter_count() > CONSOLIDATION_THRESHOLD
}
/// Consolidate all subscriptions back to Layer 1 only
///
/// Clears all tracked announcements and events, marks as consolidated,
/// and returns the Layer 1 filters to re-subscribe to.
pub fn consolidate(&mut self) -> Vec<Filter> {
tracing::info!(
"Consolidating subscriptions: {} announcements, {} events -> Layer 1 only",
self.subscribed_announcements.len(),
self.subscribed_events.len()
);
// Clear tracked subscriptions
self.subscribed_announcements.clear();
self.subscribed_events.clear();
self.is_consolidated = true;
// Return Layer 1 filters
self.filter_service.get_layer1_filters()
}
/// Get the total count of active filters
///
/// Counts 1 filter per announcement (Layer 2) + 1 filter per event (Layer 3),
/// plus the base Layer 1 filter count.
pub fn get_filter_count(&self) -> usize {
if self.is_consolidated {
// When consolidated, we only have Layer 1 filters
1
} else {
// Layer 1 (1) + Layer 2 (announcements) + Layer 3 (events)
1 + self.subscribed_announcements.len() + self.subscribed_events.len()
}
}
/// Check if an announcement has been subscribed to
pub fn has_announcement(&self, event_id: &str) -> bool {
self.subscribed_announcements.contains(event_id)
}
/// Check if an event has been subscribed to
pub fn has_event(&self, event_id: &str) -> bool {
self.subscribed_events.contains(event_id)
}
/// Check if subscriptions have been consolidated
pub fn is_consolidated(&self) -> bool {
self.is_consolidated
}
/// Get the count of subscribed announcements
pub fn announcement_count(&self) -> usize {
self.subscribed_announcements.len()
}
/// Get the count of subscribed events
pub fn event_count(&self) -> usize {
self.subscribed_events.len()
}
/// Build Layer 2 filter for a specific announcement event
///
/// Creates a filter with an 'a' tag pointing to the announcement's coordinates.
fn build_layer2_filter_for_announcement(&self, event: &Event) -> Vec<Filter> {
// Extract the d tag (identifier) from the event
let identifier = event.tags.iter().find_map(|tag| {
let tag_vec = tag.clone().to_vec();
if tag_vec.len() >= 2 && tag_vec[0] == "d" {
Some(tag_vec[1].clone())
} else {
None
}
});
let identifier = match identifier {
Some(id) => id,
None => {
tracing::warn!(
"Announcement {} has no 'd' tag, cannot build Layer 2 filter",
event.id.to_hex()
);
return Vec::new();
}
};
// Determine the kind for the coordinate
let kind = event.kind.as_u16();
if kind != KIND_REPOSITORY_ANNOUNCEMENT && kind != KIND_MAINTAINER_LIST {
tracing::warn!(
"Event {} is not an announcement (kind {}), cannot build Layer 2 filter",
event.id.to_hex(),
kind
);
return Vec::new();
}
// Build the addressable coordinate: kind:pubkey:identifier
let coord = format!("{}:{}:{}", kind, event.pubkey.to_hex(), identifier);
// Create filter with 'a' tag for this coordinate
let filter = Filter::new().custom_tag(
SingleLetterTag::lowercase(Alphabet::A),
coord,
);
vec![filter]
}
/// Check if an event kind is an announcement kind
pub fn is_announcement_kind(kind: u16) -> bool {
kind == KIND_REPOSITORY_ANNOUNCEMENT || kind == KIND_MAINTAINER_LIST
}
/// Check if an event kind is a PR/Issue/Patch kind that should trigger Layer 3
pub fn is_pr_issue_kind(kind: u16) -> bool {
matches!(
kind,
1617 | // Patch proposal (NIP-34)
1618 | // PR
1619 | // PR Update
1621 | // Issue
1622 // Reply
)
}
}
#[cfg(test)]
mod tests {
use super::SubscriptionManager;
#[test]
fn test_is_announcement_kind() {
assert!(SubscriptionManager::is_announcement_kind(30617));
assert!(SubscriptionManager::is_announcement_kind(30618));
assert!(!SubscriptionManager::is_announcement_kind(1));
assert!(!SubscriptionManager::is_announcement_kind(1617));
}
#[test]
fn test_is_pr_issue_kind() {
assert!(SubscriptionManager::is_pr_issue_kind(1617));
assert!(SubscriptionManager::is_pr_issue_kind(1618));
assert!(SubscriptionManager::is_pr_issue_kind(1619));
assert!(SubscriptionManager::is_pr_issue_kind(1621));
assert!(SubscriptionManager::is_pr_issue_kind(1622));
assert!(!SubscriptionManager::is_pr_issue_kind(30617));
assert!(!SubscriptionManager::is_pr_issue_kind(1));
}
}
+748
View File
@@ -0,0 +1,748 @@
//! GRASP-02 Phase 4: Dynamic Subscription Integration Tests
//!
//! Tests verify dynamic subscription management:
//! - New announcement triggers Layer 2 subscription
//! - New PR/Issue triggers Layer 3 subscription
//! - Subscription count tracking per connection
//! - Consolidation at filter count > 150
//! - No duplicate subscriptions
//!
//! # Running Tests
//!
//! ```bash
//! cargo test --test proactive_sync_dynamic
//! cargo test --test proactive_sync_dynamic -- --nocapture
//! ```
use std::collections::HashSet;
use ngit_grasp::sync::SubscriptionManager;
use nostr_sdk::prelude::*;
/// Kind 30617 - Repository Announcement (NIP-34)
const KIND_REPOSITORY_ANNOUNCEMENT: u16 = 30617;
/// Kind 30618 - Maintainer List (NIP-34)
const KIND_MAINTAINER_LIST: u16 = 30618;
/// Maximum filters before consolidation (from spec)
const CONSOLIDATION_THRESHOLD: usize = 150;
/// Helper to create a test announcement event
fn create_test_announcement(keys: &Keys, identifier: &str) -> Event {
let tags = vec![
Tag::identifier(identifier),
Tag::custom(
TagKind::custom("clone"),
vec![format!("http://test.example.com/{}", identifier)],
),
Tag::custom(
TagKind::custom("relays"),
vec!["ws://test.example.com".to_string()],
),
];
EventBuilder::new(Kind::Custom(KIND_REPOSITORY_ANNOUNCEMENT), "Test repo")
.tags(tags)
.sign_with_keys(keys)
.expect("Failed to sign event")
}
/// Helper to create a test maintainer list event
fn create_test_maintainer_list(keys: &Keys, identifier: &str) -> Event {
let tags = vec![
Tag::identifier(identifier),
Tag::custom(
TagKind::custom("relays"),
vec!["ws://test.example.com".to_string()],
),
];
EventBuilder::new(Kind::Custom(KIND_MAINTAINER_LIST), "Maintainer list")
.tags(tags)
.sign_with_keys(keys)
.expect("Failed to sign event")
}
/// Helper to create a test PR event (kind 1617)
fn create_test_pr_event(keys: &Keys, repo_coord: &str) -> Event {
let tags = vec![Tag::custom(
TagKind::custom("a"),
vec![repo_coord.to_string()],
)];
EventBuilder::new(Kind::Custom(1617), "Test patch proposal")
.tags(tags)
.sign_with_keys(keys)
.expect("Failed to sign event")
}
/// Helper to create a test PR event (kind 1618)
fn create_test_pr_1618_event(keys: &Keys, repo_coord: &str) -> Event {
let tags = vec![Tag::custom(
TagKind::custom("a"),
vec![repo_coord.to_string()],
)];
EventBuilder::new(Kind::Custom(1618), "Test PR")
.tags(tags)
.sign_with_keys(keys)
.expect("Failed to sign event")
}
/// Helper to create a test Issue event (kind 1621)
fn create_test_issue_event(keys: &Keys, repo_coord: &str) -> Event {
let tags = vec![Tag::custom(
TagKind::custom("a"),
vec![repo_coord.to_string()],
)];
EventBuilder::new(Kind::Custom(1621), "Test issue")
.tags(tags)
.sign_with_keys(keys)
.expect("Failed to sign event")
}
/// Helper to create a test Reply event (kind 1622)
fn create_test_reply_event(keys: &Keys, event_id: &str) -> Event {
let tags = vec![Tag::custom(
TagKind::custom("e"),
vec![event_id.to_string()],
)];
EventBuilder::new(Kind::Custom(1622), "Test reply")
.tags(tags)
.sign_with_keys(keys)
.expect("Failed to sign event")
}
// ============================================================================
// Kind Detection Tests
// ============================================================================
/// Test that announcement kinds are correctly identified
#[test]
fn test_is_announcement_kind_30617() {
assert!(SubscriptionManager::is_announcement_kind(30617));
}
/// Test that maintainer list kind is correctly identified
#[test]
fn test_is_announcement_kind_30618() {
assert!(SubscriptionManager::is_announcement_kind(30618));
}
/// Test that non-announcement kinds are not identified as announcements
#[test]
fn test_is_announcement_kind_negative() {
assert!(!SubscriptionManager::is_announcement_kind(1)); // Text note
assert!(!SubscriptionManager::is_announcement_kind(1617)); // PR
assert!(!SubscriptionManager::is_announcement_kind(1621)); // Issue
assert!(!SubscriptionManager::is_announcement_kind(0)); // Unknown
}
/// Test that PR/Issue kinds are correctly identified
#[test]
fn test_is_pr_issue_kind_1617() {
assert!(SubscriptionManager::is_pr_issue_kind(1617)); // Patch proposal
}
/// Test that PR kind 1618 is correctly identified
#[test]
fn test_is_pr_issue_kind_1618() {
assert!(SubscriptionManager::is_pr_issue_kind(1618)); // PR
}
/// Test that PR update kind is correctly identified
#[test]
fn test_is_pr_issue_kind_1619() {
assert!(SubscriptionManager::is_pr_issue_kind(1619)); // PR Update
}
/// Test that Issue kind is correctly identified
#[test]
fn test_is_pr_issue_kind_1621() {
assert!(SubscriptionManager::is_pr_issue_kind(1621)); // Issue
}
/// Test that Reply kind is correctly identified
#[test]
fn test_is_pr_issue_kind_1622() {
assert!(SubscriptionManager::is_pr_issue_kind(1622)); // Reply
}
/// Test that non-PR/Issue kinds are not identified
#[test]
fn test_is_pr_issue_kind_negative() {
assert!(!SubscriptionManager::is_pr_issue_kind(30617)); // Announcement
assert!(!SubscriptionManager::is_pr_issue_kind(1)); // Text note
assert!(!SubscriptionManager::is_pr_issue_kind(0)); // Unknown
}
// ============================================================================
// Filter Count Tests
// ============================================================================
/// Test initial filter count is 1 (Layer 1 only)
#[test]
fn test_initial_filter_count() {
// Create a minimal SubscriptionManager-like state for testing
// We test the logic without needing a full FilterService
// Initial state: 0 announcements, 0 events, not consolidated
// Filter count should be: 1 (Layer 1) + 0 + 0 = 1
let announcement_count = 0;
let event_count = 0;
let is_consolidated = false;
let filter_count = if is_consolidated {
1
} else {
1 + announcement_count + event_count
};
assert_eq!(filter_count, 1);
}
/// Test filter count increases with announcements
#[test]
fn test_filter_count_with_announcements() {
let announcement_count = 5;
let event_count = 0;
let is_consolidated = false;
let filter_count = if is_consolidated {
1
} else {
1 + announcement_count + event_count
};
// 1 (Layer 1) + 5 (announcements) = 6
assert_eq!(filter_count, 6);
}
/// Test filter count increases with events
#[test]
fn test_filter_count_with_events() {
let announcement_count = 0;
let event_count = 10;
let is_consolidated = false;
let filter_count = if is_consolidated {
1
} else {
1 + announcement_count + event_count
};
// 1 (Layer 1) + 10 (events) = 11
assert_eq!(filter_count, 11);
}
/// Test filter count with both announcements and events
#[test]
fn test_filter_count_mixed() {
let announcement_count = 50;
let event_count = 30;
let is_consolidated = false;
let filter_count = if is_consolidated {
1
} else {
1 + announcement_count + event_count
};
// 1 + 50 + 30 = 81
assert_eq!(filter_count, 81);
}
/// Test filter count is 1 when consolidated
#[test]
fn test_filter_count_consolidated() {
let announcement_count = 100; // These would be cleared on consolidation
let event_count = 100;
let is_consolidated = true;
let filter_count = if is_consolidated {
1
} else {
1 + announcement_count + event_count
};
assert_eq!(filter_count, 1);
}
// ============================================================================
// Consolidation Threshold Tests
// ============================================================================
/// Test consolidation is not triggered below threshold
#[test]
fn test_should_consolidate_below_threshold() {
let filter_count = 100;
let is_consolidated = false;
let should_consolidate = !is_consolidated && filter_count > CONSOLIDATION_THRESHOLD;
assert!(!should_consolidate);
}
/// Test consolidation is triggered at threshold
#[test]
fn test_should_consolidate_at_threshold() {
let filter_count = 151; // > 150
let is_consolidated = false;
let should_consolidate = !is_consolidated && filter_count > CONSOLIDATION_THRESHOLD;
assert!(should_consolidate);
}
/// Test consolidation is triggered well above threshold
#[test]
fn test_should_consolidate_above_threshold() {
let filter_count = 200;
let is_consolidated = false;
let should_consolidate = !is_consolidated && filter_count > CONSOLIDATION_THRESHOLD;
assert!(should_consolidate);
}
/// Test consolidation is not triggered if already consolidated
#[test]
fn test_should_consolidate_already_consolidated() {
let filter_count = 200; // Would trigger, but already consolidated
let is_consolidated = true;
let should_consolidate = !is_consolidated && filter_count > CONSOLIDATION_THRESHOLD;
assert!(!should_consolidate);
}
/// Test exact threshold boundary (150 should NOT trigger, 151 should)
#[test]
fn test_consolidation_threshold_boundary() {
let is_consolidated = false;
// 150 should NOT trigger (> 150, not >= 150)
let should_consolidate_at_150 = !is_consolidated && 150 > CONSOLIDATION_THRESHOLD;
assert!(!should_consolidate_at_150);
// 151 should trigger
let should_consolidate_at_151 = !is_consolidated && 151 > CONSOLIDATION_THRESHOLD;
assert!(should_consolidate_at_151);
}
// ============================================================================
// Duplicate Prevention Tests
// ============================================================================
/// Test duplicate announcement detection
#[test]
fn test_duplicate_announcement_prevention() {
let mut subscribed_announcements: HashSet<String> = HashSet::new();
let event_id = "abc123".to_string();
// First add should succeed
let is_new = !subscribed_announcements.contains(&event_id);
assert!(is_new);
subscribed_announcements.insert(event_id.clone());
// Second add should fail (duplicate)
let is_new_again = !subscribed_announcements.contains(&event_id);
assert!(!is_new_again);
}
/// Test duplicate event detection
#[test]
fn test_duplicate_event_prevention() {
let mut subscribed_events: HashSet<String> = HashSet::new();
let event_id = "def456".to_string();
// First add should succeed
let is_new = !subscribed_events.contains(&event_id);
assert!(is_new);
subscribed_events.insert(event_id.clone());
// Second add should fail (duplicate)
let is_new_again = !subscribed_events.contains(&event_id);
assert!(!is_new_again);
}
/// Test multiple unique items are tracked correctly
#[test]
fn test_multiple_unique_items_tracked() {
let mut subscribed_announcements: HashSet<String> = HashSet::new();
// Add multiple unique announcements
for i in 0..10 {
let id = format!("announcement_{}", i);
assert!(!subscribed_announcements.contains(&id));
subscribed_announcements.insert(id);
}
assert_eq!(subscribed_announcements.len(), 10);
}
// ============================================================================
// Event Creation and Validation Tests
// ============================================================================
/// Test announcement event has required d tag
#[test]
fn test_announcement_has_d_tag() {
let keys = Keys::generate();
let event = create_test_announcement(&keys, "my-repo");
let has_d_tag = event.tags.iter().any(|tag| {
let tag_vec = tag.clone().to_vec();
tag_vec.len() >= 2 && tag_vec[0] == "d"
});
assert!(has_d_tag);
}
/// Test announcement event has correct kind
#[test]
fn test_announcement_correct_kind() {
let keys = Keys::generate();
let event = create_test_announcement(&keys, "my-repo");
assert_eq!(event.kind.as_u16(), KIND_REPOSITORY_ANNOUNCEMENT);
}
/// Test maintainer list event has correct kind
#[test]
fn test_maintainer_list_correct_kind() {
let keys = Keys::generate();
let event = create_test_maintainer_list(&keys, "maintainers");
assert_eq!(event.kind.as_u16(), KIND_MAINTAINER_LIST);
}
/// Test PR event has a tag
#[test]
fn test_pr_event_has_a_tag() {
let keys = Keys::generate();
let coord = "30617:pubkey123:my-repo";
let event = create_test_pr_event(&keys, coord);
let has_a_tag = event.tags.iter().any(|tag| {
let tag_vec = tag.clone().to_vec();
tag_vec.len() >= 2 && tag_vec[0] == "a"
});
assert!(has_a_tag);
}
/// Test issue event has a tag
#[test]
fn test_issue_event_has_a_tag() {
let keys = Keys::generate();
let coord = "30617:pubkey123:my-repo";
let event = create_test_issue_event(&keys, coord);
let has_a_tag = event.tags.iter().any(|tag| {
let tag_vec = tag.clone().to_vec();
tag_vec.len() >= 2 && tag_vec[0] == "a"
});
assert!(has_a_tag);
}
/// Test reply event has e tag
#[test]
fn test_reply_event_has_e_tag() {
let keys = Keys::generate();
let event_id = "abc123def456";
let event = create_test_reply_event(&keys, event_id);
let has_e_tag = event.tags.iter().any(|tag| {
let tag_vec = tag.clone().to_vec();
tag_vec.len() >= 2 && tag_vec[0] == "e"
});
assert!(has_e_tag);
}
// ============================================================================
// Subscription Lifecycle Tests
// ============================================================================
/// Test subscription lifecycle: initial -> add announcements -> add events -> consolidate
#[test]
fn test_subscription_lifecycle() {
let mut subscribed_announcements: HashSet<String> = HashSet::new();
let mut subscribed_events: HashSet<String> = HashSet::new();
let mut is_consolidated = false;
// Initial state
let initial_count = 1 + subscribed_announcements.len() + subscribed_events.len();
assert_eq!(initial_count, 1);
// Add some announcements
for i in 0..50 {
subscribed_announcements.insert(format!("ann_{}", i));
}
let after_announcements = 1 + subscribed_announcements.len() + subscribed_events.len();
assert_eq!(after_announcements, 51);
// Add some events
for i in 0..50 {
subscribed_events.insert(format!("evt_{}", i));
}
let after_events = 1 + subscribed_announcements.len() + subscribed_events.len();
assert_eq!(after_events, 101);
// Add more to exceed threshold
for i in 50..100 {
subscribed_announcements.insert(format!("ann_{}", i));
}
let before_consolidation = 1 + subscribed_announcements.len() + subscribed_events.len();
assert_eq!(before_consolidation, 151);
// Should trigger consolidation
let should_consolidate = !is_consolidated && before_consolidation > CONSOLIDATION_THRESHOLD;
assert!(should_consolidate);
// Consolidate
subscribed_announcements.clear();
subscribed_events.clear();
is_consolidated = true;
// After consolidation
let after_consolidation = if is_consolidated { 1 } else { 1 + subscribed_announcements.len() + subscribed_events.len() };
assert_eq!(after_consolidation, 1);
// Should not trigger consolidation again
let should_consolidate_again = !is_consolidated && after_consolidation > CONSOLIDATION_THRESHOLD;
assert!(!should_consolidate_again);
}
/// Test that consolidated state blocks new additions
#[test]
fn test_consolidated_blocks_additions() {
let is_consolidated = true;
// When consolidated, add_announcement should return None (simulated)
// The logic is: if is_consolidated, return None
let should_add = !is_consolidated;
assert!(!should_add);
}
/// Test that non-consolidated state allows additions
#[test]
fn test_non_consolidated_allows_additions() {
let is_consolidated = false;
let mut subscribed_announcements: HashSet<String> = HashSet::new();
let event_id = "new_announcement";
// When not consolidated and event not in set, should add
let should_add = !is_consolidated && !subscribed_announcements.contains(event_id);
assert!(should_add);
subscribed_announcements.insert(event_id.to_string());
assert!(subscribed_announcements.contains(event_id));
}
// ============================================================================
// Filter Building Tests (coordinate format)
// ============================================================================
/// Test announcement coordinate format
#[test]
fn test_announcement_coordinate_format() {
let keys = Keys::generate();
let identifier = "my-repo";
let event = create_test_announcement(&keys, identifier);
// Extract d tag
let d_tag = event.tags.iter().find_map(|tag| {
let tag_vec = tag.clone().to_vec();
if tag_vec.len() >= 2 && tag_vec[0] == "d" {
Some(tag_vec[1].clone())
} else {
None
}
});
assert!(d_tag.is_some());
assert_eq!(d_tag.unwrap(), identifier);
// Build coordinate: kind:pubkey:identifier
let coord = format!("{}:{}:{}", KIND_REPOSITORY_ANNOUNCEMENT, event.pubkey.to_hex(), identifier);
// Verify format
let parts: Vec<&str> = coord.split(':').collect();
assert_eq!(parts.len(), 3);
assert_eq!(parts[0], "30617");
assert_eq!(parts[2], identifier);
}
/// Test multiple announcement coordinates are unique
#[test]
fn test_multiple_announcement_coordinates_unique() {
let keys = Keys::generate();
let identifiers = vec!["repo1", "repo2", "repo3"];
let mut coords: HashSet<String> = HashSet::new();
for id in identifiers {
let event = create_test_announcement(&keys, id);
let coord = format!("{}:{}:{}", KIND_REPOSITORY_ANNOUNCEMENT, event.pubkey.to_hex(), id);
coords.insert(coord);
}
assert_eq!(coords.len(), 3);
}
// ============================================================================
// Integration-style Tests
// ============================================================================
/// Test simulated workflow: announcement received, then PR received
#[test]
fn test_workflow_announcement_then_pr() {
let keys = Keys::generate();
let mut subscribed_announcements: HashSet<String> = HashSet::new();
let mut subscribed_events: HashSet<String> = HashSet::new();
let is_consolidated = false;
// Step 1: Receive announcement
let announcement = create_test_announcement(&keys, "my-repo");
let ann_id = announcement.id.to_hex();
// Should add to tracking (simulating add_announcement)
let should_add_ann = !is_consolidated && !subscribed_announcements.contains(&ann_id);
assert!(should_add_ann);
subscribed_announcements.insert(ann_id.clone());
// Filter count should increase
let filter_count = 1 + subscribed_announcements.len() + subscribed_events.len();
assert_eq!(filter_count, 2);
// Step 2: Receive PR for that repo
let coord = format!("{}:{}:my-repo", KIND_REPOSITORY_ANNOUNCEMENT, keys.public_key().to_hex());
let pr = create_test_pr_event(&keys, &coord);
let pr_id = pr.id.to_hex();
// Should add to tracking (simulating add_event)
let should_add_pr = !is_consolidated && !subscribed_events.contains(&pr_id);
assert!(should_add_pr);
subscribed_events.insert(pr_id.clone());
// Filter count should increase again
let filter_count = 1 + subscribed_announcements.len() + subscribed_events.len();
assert_eq!(filter_count, 3);
}
/// Test stress: adding many items triggers consolidation
#[test]
fn test_stress_many_items_triggers_consolidation() {
let keys = Keys::generate();
let mut subscribed_announcements: HashSet<String> = HashSet::new();
let mut subscribed_events: HashSet<String> = HashSet::new();
let mut is_consolidated = false;
let mut consolidation_triggered = false;
// Add 100 announcements
for i in 0..100 {
let event = create_test_announcement(&keys, &format!("repo-{}", i));
let event_id = event.id.to_hex();
if !is_consolidated && !subscribed_announcements.contains(&event_id) {
subscribed_announcements.insert(event_id);
}
// Check consolidation after each add
let filter_count = 1 + subscribed_announcements.len() + subscribed_events.len();
if !is_consolidated && filter_count > CONSOLIDATION_THRESHOLD {
consolidation_triggered = true;
subscribed_announcements.clear();
subscribed_events.clear();
is_consolidated = true;
break;
}
}
// If we didn't consolidate yet, add events
if !consolidation_triggered {
for i in 0..100 {
let coord = format!("30617:pubkey:repo-{}", i);
let event = create_test_pr_event(&keys, &coord);
let event_id = event.id.to_hex();
if !is_consolidated && !subscribed_events.contains(&event_id) {
subscribed_events.insert(event_id);
}
// Check consolidation after each add
let filter_count = 1 + subscribed_announcements.len() + subscribed_events.len();
if !is_consolidated && filter_count > CONSOLIDATION_THRESHOLD {
consolidation_triggered = true;
subscribed_announcements.clear();
subscribed_events.clear();
is_consolidated = true;
break;
}
}
}
// Consolidation should have been triggered
assert!(consolidation_triggered);
assert!(is_consolidated);
// After consolidation, counts should be reset
assert_eq!(subscribed_announcements.len(), 0);
assert_eq!(subscribed_events.len(), 0);
}
/// Test that all PR/Issue kinds are handled consistently
#[test]
fn test_all_pr_issue_kinds_handled() {
let keys = Keys::generate();
let coord = "30617:pubkey:repo";
// All these kinds should be identified as PR/Issue
let pr_kinds = vec![1617, 1618, 1619, 1621, 1622];
for kind in pr_kinds {
assert!(
SubscriptionManager::is_pr_issue_kind(kind),
"Kind {} should be identified as PR/Issue",
kind
);
}
}
/// Test that announcement and PR/Issue kinds are mutually exclusive
#[test]
fn test_kind_categories_mutually_exclusive() {
let announcement_kinds = vec![30617, 30618];
let pr_issue_kinds = vec![1617, 1618, 1619, 1621, 1622];
// No announcement kind should be a PR/Issue kind
for kind in &announcement_kinds {
assert!(
!SubscriptionManager::is_pr_issue_kind(*kind),
"Announcement kind {} should not be PR/Issue",
kind
);
}
// No PR/Issue kind should be an announcement kind
for kind in &pr_issue_kinds {
assert!(
!SubscriptionManager::is_announcement_kind(*kind),
"PR/Issue kind {} should not be announcement",
kind
);
}
}