diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..5b89428 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,105 @@ +name: CI + +on: + push: + branches: + - main + pull_request: + branches: + - main + +permissions: + contents: read + +concurrency: + group: ci-${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }} + cancel-in-progress: true + +env: + CARGO_TERM_COLOR: always + +jobs: + quality: + name: Rust quality + runs-on: ubuntu-24.04 + steps: + - name: Checkout + uses: actions/checkout@v7 + + - name: Install Rust + uses: dtolnay/rust-toolchain@stable + with: + components: clippy,rustfmt + + - name: Cache Cargo + uses: Swatinem/rust-cache@v2 + + - name: Check formatting + run: cargo fmt --all -- --check + + - name: Run Clippy + run: cargo clippy --all-targets --all-features -- -D warnings + + - name: Test all features + run: cargo test --all-targets --all-features + + - name: Test minimal features + run: cargo test --all-targets --no-default-features + + - name: Build documentation + env: + RUSTDOCFLAGS: -D warnings + run: cargo doc --all-features --no-deps + + - name: Verify package contents + run: cargo package + + jetstream: + name: NATS JetStream 2.11.8 + runs-on: ubuntu-24.04 + env: + A3S_EVENT_REQUIRE_NATS: "1" + NATS_IMAGE: nats:2.11.8-alpine@sha256:71092f77d707a4a81b12aca5096d6b2d2e07ad16aa57c84066940a17af74f61a + steps: + - name: Checkout + uses: actions/checkout@v7 + + - name: Install Rust + uses: dtolnay/rust-toolchain@stable + + - name: Cache Cargo + uses: Swatinem/rust-cache@v2 + with: + key: jetstream + + - name: Start NATS with JetStream + run: | + docker run --detach \ + --name a3s-event-nats \ + --publish 4222:4222 \ + --publish 8222:8222 \ + "$NATS_IMAGE" \ + --jetstream \ + --http_port 8222 + + for attempt in $(seq 1 30); do + if curl --fail --silent \ + "http://127.0.0.1:8222/healthz?js-enabled-only=true" >/dev/null; then + exit 0 + fi + sleep 1 + done + + docker logs a3s-event-nats + exit 1 + + - name: Run real JetStream integration tests + run: cargo test --test nats_integration -- --nocapture --test-threads=1 + + - name: Show NATS logs after failure + if: failure() + run: docker logs a3s-event-nats + + - name: Stop NATS + if: always() + run: docker rm --force a3s-event-nats diff --git a/README.md b/README.md index b91d3dd..8d58100 100644 --- a/README.md +++ b/README.md @@ -322,7 +322,7 @@ A3S Event does NOT re-implement capabilities that providers already offer native ```bash just build # Build the project just test # Run all tests -just test-integration # NATS integration tests (requires nats-server -js) +just test-integration # NATS integration tests (fails if JetStream is unavailable) just bench # Performance benchmarks just lint # Run clippy just fmt # Format code @@ -332,7 +332,7 @@ just doc # Generate and open docs ### Test Coverage -176 unit tests + 29 memory integration tests + 9 NATS integration tests + 2 doc tests across 15 modules. +196 unit tests + 29 memory integration tests + 11 NATS integration tests + 2 doc tests across 15 modules. ```bash # Unit tests (no external dependencies) @@ -343,6 +343,11 @@ nats-server -js just test-integration ``` +`just test-integration` sets `A3S_EVENT_REQUIRE_NATS=1`, so a missing or +misconfigured JetStream server fails instead of silently skipping the suite. +Every pull request runs the same fail-closed suite against the pinned official +NATS 2.11.8 image. + ## Community Join us on [Discord](https://discord.gg/XVg6Hu6H) for questions, discussions, and updates. diff --git a/justfile b/justfile index 78e9afa..85bd3cf 100644 --- a/justfile +++ b/justfile @@ -151,9 +151,9 @@ test-nats: test-store: cargo test --lib -- store::tests -# NATS integration tests (requires running NATS server: nats-server -js) +# NATS integration tests; fail if JetStream is unavailable test-integration: - cargo test --test nats_integration + A3S_EVENT_REQUIRE_NATS=1 cargo test --test nats_integration -- --test-threads=1 # ============================================================================ # Benchmarks (requires: cargo install criterion) diff --git a/src/messaging.rs b/src/messaging.rs index e7c3ec8..bfcfa9f 100644 --- a/src/messaging.rs +++ b/src/messaging.rs @@ -125,9 +125,11 @@ pub trait MessageStream: Send + Sync { async fn next_timeout(&mut self, timeout: Duration) -> Result>; } +type SubscriberRegistry = std::sync::Arc)>>>; + /// In-memory message broker for single-process testing pub struct InMemoryMessaging { - subscribers: std::sync::Arc)>>>, + subscribers: SubscriberRegistry, } impl InMemoryMessaging { @@ -381,10 +383,14 @@ mod tests { "test".to_string(), serde_json::json!({}), ) - .to_session("session.123".to_string()); + .to_session("123".to_string()); messaging.send(&msg).await.unwrap(); - let received = stream.next().await.unwrap().unwrap(); + let received = stream + .next_timeout(std::time::Duration::from_secs(1)) + .await + .unwrap() + .expect("targeted message was not delivered"); assert_eq!(received.source_id, "session-1"); } @@ -404,7 +410,7 @@ mod tests { async fn test_message_handler_ref_none() { let handler = MessageHandlerRef::none(); let msg = Message::new("s1".to_string(), "test".to_string(), serde_json::json!({})); - handler.handle(msg); // Should not panic + handler.handle(msg).await; // Should not panic } #[test] diff --git a/src/metrics.rs b/src/metrics.rs index f8fd7f6..4077c39 100644 --- a/src/metrics.rs +++ b/src/metrics.rs @@ -110,11 +110,7 @@ impl EventMetrics { pub fn snapshot(&self) -> MetricsSnapshot { let publish_count = self.publish_count.load(Ordering::Relaxed); let total_latency = self.publish_latency_us.load(Ordering::Relaxed); - let avg_latency_us = if publish_count > 0 { - total_latency / publish_count - } else { - 0 - }; + let avg_latency_us = total_latency.checked_div(publish_count).unwrap_or(0); MetricsSnapshot { publish_count, diff --git a/src/provider/nats/client.rs b/src/provider/nats/client.rs index 40b74be..733be9c 100644 --- a/src/provider/nats/client.rs +++ b/src/provider/nats/client.rs @@ -6,9 +6,11 @@ use crate::error::{EventError, Result}; use crate::types::{DeliverPolicy, Event, PublishOptions, SubscribeOptions}; use async_nats::jetstream; use std::sync::Arc; -use std::time::Duration; +use std::time::{Duration, Instant}; use tokio::sync::Mutex; +const HISTORY_FETCH_BATCH_SIZE: usize = 256; + /// NATS JetStream client /// /// Low-level client for publishing and subscribing to events via NATS. @@ -289,8 +291,12 @@ impl NatsClient { /// Fetch historical events from the stream pub async fn history(&self, filter_subject: Option<&str>, limit: usize) -> Result> { + if limit == 0 { + return Ok(Vec::new()); + } + let mut config = jetstream::consumer::pull::Config { - deliver_policy: jetstream::consumer::DeliverPolicy::Last, + deliver_policy: jetstream::consumer::DeliverPolicy::All, ack_policy: jetstream::consumer::AckPolicy::None, ..Default::default() }; @@ -309,35 +315,63 @@ impl NatsClient { EventError::Consumer(format!("Failed to create history consumer: {}", e)) })?; - let mut events = Vec::with_capacity(limit); - let batch = consumer - .fetch() - .max_messages(limit) - .expires(Duration::from_secs(self.config.request_timeout_secs)) - .messages() - .await - .map_err(|e| EventError::JetStream(format!("Failed to fetch history: {}", e)))?; + let pending = usize::try_from(consumer.cached_info().num_pending).unwrap_or(usize::MAX); + if pending == 0 { + return Ok(Vec::new()); + } use futures_util::StreamExt; - let mut batch = std::pin::pin!(batch); - while let Some(msg) = batch.next().await { - match msg { - Ok(msg) => { - if let Ok(event) = serde_json::from_slice::(&msg.payload) { - events.push(event); + let mut events = std::collections::VecDeque::with_capacity(pending.min(limit)); + let timeout = Duration::from_secs(self.config.request_timeout_secs); + let started_at = Instant::now(); + let mut remaining = pending; + + 'history: while remaining > 0 { + let remaining_timeout = timeout.saturating_sub(started_at.elapsed()); + if remaining_timeout.is_zero() { + tracing::warn!( + remaining, + "History fetch reached its request timeout before scanning all pending events" + ); + break; + } + + let requested = remaining.min(HISTORY_FETCH_BATCH_SIZE); + let batch = consumer + .fetch() + .max_messages(requested) + .expires(remaining_timeout) + .messages() + .await + .map_err(|e| EventError::JetStream(format!("Failed to fetch history: {}", e)))?; + let mut batch = std::pin::pin!(batch); + let mut received = 0usize; + + while let Some(msg) = batch.next().await { + match msg { + Ok(msg) => { + received += 1; + if let Ok(event) = serde_json::from_slice::(&msg.payload) { + if events.len() == limit { + events.pop_front(); + } + events.push_back(event); + } } - if events.len() >= limit { - break; + Err(e) => { + tracing::warn!("Error fetching history message: {}", e); + break 'history; } } - Err(e) => { - tracing::warn!("Error fetching history message: {}", e); - break; - } + } + + remaining = remaining.saturating_sub(received); + if received < requested { + break; } } - Ok(events) + Ok(events.into_iter().rev().collect()) } /// Delete a durable consumer diff --git a/src/store.rs b/src/store.rs index ca2f455..0f92aac 100644 --- a/src/store.rs +++ b/src/store.rs @@ -401,7 +401,7 @@ impl EventBus { let mut subscribers = Vec::new(); for subject in &filter.subjects { - let consumer_name = format!("{}-{}", subscriber_id, subject.replace('.', "-")); + let consumer_name = subscription_consumer_name(subscriber_id, subject); let sub = match (&filter.options, filter.durable) { (Some(opts), true) => { self.provider @@ -436,7 +436,7 @@ impl EventBus { if let Some(filter) = filter { self.metrics.record_unsubscribe(); for subject in &filter.subjects { - let consumer_name = format!("{}-{}", subscriber_id, subject.replace('.', "-")); + let consumer_name = subscription_consumer_name(subscriber_id, subject); if let Err(e) = self.provider.unsubscribe(&consumer_name).await { tracing::warn!( consumer = %consumer_name, @@ -552,6 +552,22 @@ impl EventBus { } } +fn subscription_consumer_name(subscriber_id: &str, subject: &str) -> String { + format!("{subscriber_id}-{subject}") + .chars() + .map(|character| { + if character.is_whitespace() + || character.is_control() + || matches!(character, '.' | '*' | '>' | '/' | '\\') + { + '-' + } else { + character + } + }) + .collect() +} + #[cfg(test)] mod tests { use super::*; @@ -564,6 +580,18 @@ mod tests { EventBus::new(MemoryProvider::default()) } + #[test] + fn test_subscription_consumer_name_is_provider_safe() { + assert_eq!( + subscription_consumer_name("tenant 1", "events.market.*.>"), + "tenant-1-events-market----" + ); + assert_eq!( + subscription_consumer_name("analyst", "events.market.usd"), + "analyst-events-market-usd" + ); + } + #[tokio::test] async fn test_publish_and_list() { let bus = test_bus(); diff --git a/src/types.rs b/src/types.rs index 21f4eef..253185d 100644 --- a/src/types.rs +++ b/src/types.rs @@ -17,7 +17,7 @@ pub(crate) type BoxFuture<'a, T> = Pin + Send + 'a>>; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct Event { - /// Unique event identifier (evt-) + /// Unique event identifier (`evt-`) pub id: String, /// Subject this event was published to diff --git a/tests/nats_integration.rs b/tests/nats_integration.rs index 82bcc45..81016b7 100644 --- a/tests/nats_integration.rs +++ b/tests/nats_integration.rs @@ -4,7 +4,9 @@ //! These tests require a running NATS server with JetStream enabled: //! nats-server -js //! -//! Tests are skipped automatically if NATS is not available. +//! Tests are skipped automatically if NATS is not available, unless +//! `A3S_EVENT_REQUIRE_NATS` is set. CI sets this variable so an unavailable +//! or misconfigured server fails the suite instead of producing a false pass. use a3s_event::provider::nats::{NatsConfig, NatsProvider, StorageType}; use a3s_event::{ @@ -12,7 +14,8 @@ use a3s_event::{ SubscriptionFilter, }; -/// Try to connect to NATS. Returns None if server is unavailable. +/// Try to connect to NATS. Returns `None` if the server is unavailable and +/// fail-closed mode is not enabled. async fn try_nats_provider(stream_suffix: &str) -> Option { let config = NatsConfig { url: "nats://127.0.0.1:4222".to_string(), @@ -26,8 +29,11 @@ async fn try_nats_provider(stream_suffix: &str) -> Option { match NatsProvider::connect(config).await { Ok(provider) => Some(provider), - Err(_) => { - eprintln!("NATS not available, skipping integration test"); + Err(error) if std::env::var_os("A3S_EVENT_REQUIRE_NATS").is_some() => { + panic!("NATS JetStream is required but unavailable: {error}"); + } + Err(error) => { + eprintln!("NATS not available, skipping integration test: {error}"); None } } @@ -89,6 +95,39 @@ async fn test_nats_publish_multiple_categories() { assert!(all.len() >= 3); } +#[tokio::test] +async fn test_nats_history_returns_latest_events_first() { + let bus = nats_bus!("history_order"); + let mut published_ids = Vec::new(); + + // Exceed the client's fixed history-fetch batch size so this also proves + // that newest-first limiting remains correct across multiple bounded pulls. + for index in 0..260 { + let event = bus + .publish( + "history", + &format!("topic.{index}"), + &format!("Event {index}"), + "test", + serde_json::json!({"index": index}), + ) + .await + .unwrap(); + published_ids.push(event.id); + } + + let events = bus.list_events(None, 3).await.unwrap(); + let event_ids: Vec<&str> = events.iter().map(|event| event.id.as_str()).collect(); + let expected_ids: Vec<&str> = published_ids + .iter() + .rev() + .take(3) + .map(String::as_str) + .collect(); + + assert_eq!(event_ids, expected_ids); +} + #[tokio::test] async fn test_nats_publish_with_dedup() { let bus = nats_bus!("dedup"); @@ -144,15 +183,16 @@ async fn test_nats_durable_subscription() { // Try to receive (with timeout to avoid hanging) let sub = &mut subs[0]; - let result = tokio::time::timeout(std::time::Duration::from_secs(2), sub.next()).await; + let received = tokio::time::timeout(std::time::Duration::from_secs(5), sub.next()) + .await + .expect("timed out waiting for the durable subscription") + .expect("durable subscription returned an error") + .expect("durable subscription ended before delivering an event"); // Clean up bus.remove_subscription("test-analyst").await.unwrap(); - if let Ok(Ok(Some(received))) = result { - assert_eq!(received.event.category, "market"); - } - // If timeout, that's ok — the subscription was created successfully + assert_eq!(received.event.category, "market"); } #[tokio::test] @@ -253,14 +293,81 @@ async fn test_nats_manual_ack() { .unwrap(); // Receive with manual ack - let result = - tokio::time::timeout(std::time::Duration::from_secs(2), sub.next_manual_ack()).await; + let pending = tokio::time::timeout(std::time::Duration::from_secs(5), sub.next_manual_ack()) + .await + .expect("timed out waiting for a manually acknowledged event") + .expect("manual-ack subscription returned an error") + .expect("manual-ack subscription ended before delivering an event"); - if let Ok(Ok(Some(pending))) = result { - assert_eq!(pending.received.event.summary, "Ack test"); - pending.ack().await.unwrap(); - } + assert_eq!(pending.received.event.summary, "Ack test"); + pending.ack().await.unwrap(); // Clean up let _ = provider.unsubscribe("ack-test-consumer").await; } + +#[tokio::test] +async fn test_nats_unacked_message_is_redelivered() { + let suffix = "ack_redelivery"; + let provider = match try_nats_provider(suffix).await { + Some(provider) => provider, + None => return, + }; + + let mut subscription = provider + .subscribe_durable_with_options( + "ack-redelivery-consumer", + &format!("test.{suffix}.>"), + &SubscribeOptions { + max_deliver: Some(3), + ack_wait_secs: Some(1), + ..Default::default() + }, + ) + .await + .unwrap(); + + let event = Event::new( + format!("test.{suffix}.topic"), + "test", + "Ack redelivery test", + "test", + serde_json::json!({}), + ); + provider.publish(&event).await.unwrap(); + + let first = tokio::time::timeout( + std::time::Duration::from_secs(5), + subscription.next_manual_ack(), + ) + .await + .expect("timed out waiting for the initial delivery") + .expect("initial delivery returned an error") + .expect("subscription ended before the initial delivery"); + + assert_eq!(first.received.event.id, event.id); + assert_eq!(first.received.num_delivered, 1); + drop(first); + + let redelivered = tokio::time::timeout( + std::time::Duration::from_secs(5), + subscription.next_manual_ack(), + ) + .await + .expect("timed out waiting for the unacknowledged message to be redelivered") + .expect("redelivery returned an error") + .expect("subscription ended before redelivery"); + + assert_eq!(redelivered.received.event.id, event.id); + assert!( + redelivered.received.num_delivered >= 2, + "expected a redelivery count of at least 2, got {}", + redelivered.received.num_delivered + ); + redelivered.ack().await.unwrap(); + + provider + .unsubscribe("ack-redelivery-consumer") + .await + .unwrap(); +}