Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
105 changes: 105 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
@@ -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
9 changes: 7 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand All @@ -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.
Expand Down
4 changes: 2 additions & 2 deletions justfile
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
14 changes: 10 additions & 4 deletions src/messaging.rs
Original file line number Diff line number Diff line change
Expand Up @@ -125,9 +125,11 @@ pub trait MessageStream: Send + Sync {
async fn next_timeout(&mut self, timeout: Duration) -> Result<Option<Message>>;
}

type SubscriberRegistry = std::sync::Arc<RwLock<Vec<(String, flume::Sender<Message>)>>>;

/// In-memory message broker for single-process testing
pub struct InMemoryMessaging {
subscribers: std::sync::Arc<RwLock<Vec<(String, flume::Sender<Message>)>>>,
subscribers: SubscriberRegistry,
}

impl InMemoryMessaging {
Expand Down Expand Up @@ -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");
}

Expand All @@ -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]
Expand Down
6 changes: 1 addition & 5 deletions src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
80 changes: 57 additions & 23 deletions src/provider/nats/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -289,8 +291,12 @@ impl NatsClient {

/// Fetch historical events from the stream
pub async fn history(&self, filter_subject: Option<&str>, limit: usize) -> Result<Vec<Event>> {
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()
};
Expand All @@ -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::<Event>(&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::<Event>(&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
Expand Down
32 changes: 30 additions & 2 deletions src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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::*;
Expand All @@ -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();
Expand Down
2 changes: 1 addition & 1 deletion src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ pub(crate) type BoxFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct Event {
/// Unique event identifier (evt-<uuid>)
/// Unique event identifier (`evt-<uuid>`)
pub id: String,

/// Subject this event was published to
Expand Down
Loading