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
78 changes: 78 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,9 @@ jobs:
run: cargo doc --all-features --no-deps

- name: Verify package contents
run: cargo clippy --features iggy --no-default-features --all-targets -- -D warnings

- name: Package
run: cargo package

jetstream:
Expand Down Expand Up @@ -103,3 +106,78 @@ jobs:
- name: Stop NATS
if: always()
run: docker rm --force a3s-event-nats

iggy:
name: Apache Iggy 0.9.0
runs-on: ubuntu-24.04
env:
A3S_EVENT_REQUIRE_IGGY: "1"
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: iggy

# io_uring needs seccomp=unconfined; single shard + 2s rebalancing
# keep the e2e suite fast (server contract documented in
# tests/iggy_integration.rs and tests/e2e_chaos_resilience.rs).
- name: Start Iggy
run: |
docker run --detach --name a3s-event-iggy --publish 5102:5102 --security-opt seccomp=unconfined -e RUST_LOG=info -e IGGY_ROOT_USERNAME=iggy -e IGGY_ROOT_PASSWORD=iggy -e IGGY_TCP_ADDRESS=0.0.0.0:5102 -e IGGY_NODE_ADVERTISED_ADDRESS=127.0.0.1 -e IGGY_SHARDING_CPU_ALLOCATION=1 -e IGGY_SHARDING_PIN_CORES=false -e IGGY_CONSUMER_GROUP_REBALANCING_TIMEOUT=2s apache/iggy:0.9.0

for attempt in $(seq 1 30); do
if (echo > /dev/tcp/127.0.0.1/5102) 2>/dev/null; then
exit 0
fi
sleep 1
done

docker logs a3s-event-iggy
exit 1

- name: Cross-provider conformance + iggy e2e
# --tests (not --all-targets): the criterion bench target does not
# accept libtest flags like --test-threads.
run: cargo test --features nats,iggy --no-default-features --tests -- --test-threads=1

# Known upstream flake (apache/iggy#4361): iggy 0.9.0 can panic on
# restart/boot replay, killing the server mid-suite. One bounded
# retry with a FRESH container keeps CI signal honest without hiding
# genuine failures (a real regression fails twice).
- name: Retry once on fresh Iggy (upstream #4361 flake)
if: failure()
run: |
docker logs a3s-event-iggy || true
docker rm --force a3s-event-iggy
docker run --detach \
--name a3s-event-iggy \
--publish 5102:5102 \
--security-opt seccomp=unconfined \
-e RUST_LOG=info \
-e IGGY_ROOT_USERNAME=iggy \
-e IGGY_ROOT_PASSWORD=iggy \
-e IGGY_TCP_ADDRESS=0.0.0.0:5102 \
-e IGGY_NODE_ADVERTISED_ADDRESS=127.0.0.1 \
-e IGGY_SHARDING_CPU_ALLOCATION=1 \
-e IGGY_SHARDING_PIN_CORES=false \
-e IGGY_CONSUMER_GROUP_REBALANCING_TIMEOUT=2s \
apache/iggy:0.9.0
for attempt in $(seq 1 30); do
if (echo > /dev/tcp/127.0.0.1/5102) 2>/dev/null; then break; fi
sleep 1
done
cargo test --features nats,iggy --no-default-features --tests -- --test-threads=1

- name: Show Iggy logs after failure
if: failure()
run: docker logs a3s-event-iggy

- name: Stop Iggy
if: always()
run: docker rm --force a3s-event-iggy
81 changes: 81 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
# Changelog

All notable changes to this project are documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [0.4.0] — 2026-09-30

### ⚠️ Operational migrations

- **NATS durable consumer names are now sanitized.** `EventBus` used to build
consumer names as `{subscriber}-{subject with '.'→'-'}`; subjects also
contain `*` and `>`, which JetStream rejects outright (`error 10103`). The
name is now built by collapsing every character outside `[a-zA-Z0-9_-]` to
`-`. **Deployed consumers subscribed under the old naming will see new,
empty consumers on upgrade** — either drain/retire old subscriptions before
upgrading, or accept a one-time redelivery from the deliver policy's start.
- **NATS `history()` now scans forward from the start of the retained
stream** (`DeliverPolicy::All`) instead of `Last`, which returned at most
one message. Read-side behavior change: `list_events`/`counts` now actually
return history.

### Added

- **Apache Iggy provider** (`iggy` feature, SDK `iggy` 0.11 / server 0.9):
stream→topic mapping where each subject category is one topic (full subject
preserved in payload + `a3s-subject` user header; subscription filters
narrowed client-side), durable subscriptions as consumer groups with
explicitly stored offsets (at-least-once, last-consumed convention),
ephemeral subscriptions, subscribe-time head probing for `New`/`Last`
positioning, bounded connect timeout owned by the provider, PAT or
username/password login. Fail-closed where the broker cannot honor the
contract: `expected_sequence`, `DeliverPolicy::LastPerSubject`. Accepted
but ignored (per the trait contract): `max_deliver`, `backoff_secs`,
`max_ack_pending`, `ack_wait_secs`. Single-partition topics only
(`IggyPartitioning::Balanced` is a documented no-op in this version).
- `EventBus::from_provider(Arc<dyn EventProvider>)` — share one provider
handle between the bus and its owner.
- `EventBus::set_schema_registry` — setter symmetry with the other optional
capabilities (`with_schema_registry` was previously the only path).
- **Broker routing failures now reach the DLQ.** `EventBus`'s documented
"routes failed events to a DlqHandler" contract is actually implemented:
failed sink deliveries produce a `DeadLetterEvent` (reason
`broker routing: n of m sink deliveries failed`) and advance the
`dlq_count` metric.
- Deep end-to-end test assets: a cross-provider conformance suite (tier 0 on
every provider × 7 scenarios; tier 1 on persistent providers × 5), feature
e2e suites (pipeline, routing/bridge, cron source, crypto, CloudEvents,
messaging, DLQ/schema/sinks completion, chaos), and opt-in chaos tests
(broker restart mid-stream, PAT login) driven by environment variables.

### Known issues (upstream)

- **Iggy server 0.9.0 has an intermittent restart-path panic**: boot replay can
hit `client_id 0 is reserved for internal use` (`core/consensus/src/client_table.rs`)
when the persisted client table contains certain sessions, killing the shard
and the server. Discovered by this crate's opt-in chaos suite (broker restart
mid-stream). Repro: connect clients, publish, `docker restart` the container;
sometimes the server exits (1) during boot. A fresh container (recreate, not
restart) boots clean. Until fixed upstream, Iggy restarts in production need
a supervisor plus a readiness gate — and this is a reason the `iggy` feature
should not be considered GA-hardened even when this crate is.

### Fixed

- `InMemoryMessaging::send` prefixed targeted patterns with `session.`,
producing `session.session.<id>` — no documented filter could ever match a
targeted send (and `test_subscribe_and_send_to_specific_session` hung
every full `cargo test` run). Patterns are now the target id itself.
- `matches_pattern` checked wildcards on the pattern side only, so
subscriber filters like `session.*` never matched targeted messages;
wildcards are now honored symmetrically.
- Iggy `DeliverPolicy::ByStartTime` positioning was consumed on the first
poll even when it returned nothing, falling back to `offset(0)` and
delivering pre-cutoff events; a timestamp position now sticks until a poll
actually returns messages.
- Several `clippy -D warnings` violations across the crate.

## [0.3.0] — prior release

See git history.
9 changes: 7 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "a3s-event"
version = "0.3.0"
version = "0.4.0"
edition = "2021"
authors = ["A3S Lab"]
license = "MIT"
Expand All @@ -18,10 +18,11 @@ path = "src/lib.rs"
[features]
default = ["nats", "encryption", "cloudevents", "routing"]
nats = ["dep:async-nats", "dep:time", "dep:futures-util"]
iggy = ["dep:iggy", "dep:bytes"]
encryption = ["dep:aes-gcm", "dep:base64"]
cloudevents = ["dep:chrono"]
routing = []
full = ["nats", "encryption", "cloudevents", "routing"]
full = ["nats", "iggy", "encryption", "cloudevents", "routing"]

[dependencies]
serde = { version = "1", features = ["derive"] }
Expand All @@ -38,6 +39,10 @@ async-nats = { version = "0.38", optional = true }
futures-util = { version = "0.3", optional = true }
time = { version = "0.3", optional = true }

# Optional: Apache Iggy provider
iggy = { version = "0.11", optional = true }
bytes = { version = "1", optional = true }

# Optional: AES-256-GCM payload encryption
aes-gcm = { version = "0.10", optional = true }
base64 = { version = "0.22", optional = true }
Expand Down
28 changes: 28 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ All optional modules are behind feature gates. The minimal core (types, memory p
| Feature | Default | Description |
|---------|---------|-------------|
| `nats` | ✅ | NATS JetStream provider (`async-nats`, `futures-util`, `time`) |
| `iggy` | — | Apache Iggy provider (`iggy`, `bytes`) |
| `encryption` | ✅ | AES-256-GCM payload encryption (`aes-gcm`, `base64`) |
| `cloudevents` | ✅ | CloudEvents v1.0 conversion (`chrono`) |
| `routing` | ✅ | Broker/Trigger event routing + Sink DLQ |
Expand All @@ -83,6 +84,7 @@ a3s-event = { version = "0.3", default-features = false, features = ["nats", "en
|----------|----------|-------------|--------------|
| `MemoryProvider` | Testing, development, single-process | In-process only | Single process |
| `NatsProvider` | Production, multi-service | JetStream (file/memory) | Distributed |
| `IggyProvider` | Production, multi-service, Rust-native broker | Iggy stream (per-topic log) | Distributed |

### Memory Provider

Expand Down Expand Up @@ -119,6 +121,32 @@ let provider = NatsProvider::connect(NatsConfig {
}).await?;
```

### Apache Iggy Provider

Requires `iggy` feature. Rust-native message streaming (stream → topic → partition). Subjects map onto Iggy with one rule: the stream holds every topic, and each subject **category** becomes a topic; subscription filters narrow client-side via `subject_matches`. Durable subscriptions are consumer groups with explicitly stored offsets (at-least-once); ordering is total within a category.

```rust
use a3s_event::provider::iggy::{IggyConfig, IggyProvider};

let provider = IggyProvider::connect(IggyConfig {
server_address: "127.0.0.1:5102".to_string(),
stream_name: "a3s_events".to_string(),
subject_prefix: "events".to_string(),
max_age_secs: 604_800, // 7 days
..Default::default()
}).await?;
```

Known limitations of the current version (fail-closed, not silently ignored): `expected_sequence` and `DeliverPolicy::LastPerSubject` are rejected; `max_deliver`/`backoff_secs`/`max_ack_pending`/`ack_wait_secs` are accepted and ignored (Iggy's low-level polling has no per-group redelivery controls); topics are single-partition so `IggyPartitioning::Balanced` currently behaves like `Single`.

## Operations

- **Resilience (verified by opt-in chaos tests)**: an Iggy broker restart mid-stream preserves stream/topic/consumer-offset state; consumers reconnecting under the same name resume from their committed offset without replaying acked events. Dead group members are evicted after the server's `consumer_group.rebalancing_timeout` (default 30s). Run the chaos suite locally with `A3S_EVENT_IGGY_RESTART="docker restart <container>" cargo test --test e2e_chaos_resilience`.
- **Timestamp positioning** (`DeliverPolicy::ByStartTime`) compares against the broker's server-side receive stamps; allow for clock skew between publishers and the broker when choosing cutoffs.
- **Coverage discipline**: unit coverage is measured with `cargo llvm-cov --lib` (~83% lines; broker provider bodies are exercised by the live-server e2e suites instead). `cargo clippy --all-targets -- -D warnings` runs against both the default and the `nats,iggy` feature sets in CI; the minimal core cross-compiles cleanly for Linux x64/arm64 and Windows.
- **Baseline performance** (memory provider, crate release profile `opt-level=z` + LTO, criterion, Apple Silicon): publish ~203 µs per 100-event batch (~2.0 µs/event) and ~1.70 ms per 1000-event batch (~1.7 µs/event); `history(limit 100)` ~12.4 µs unfiltered / ~20.8 µs subject-filtered. Broker-backed providers are dominated by network round-trips, not this crate's envelope handling; run `cargo bench --bench publish` for your own baseline.
- **Migration notes live in [CHANGELOG.md](CHANGELOG.md)** — the 0.4.0 release changes NATS durable consumer naming and NATS `history()` semantics.

## Architecture

```text
Expand Down
113 changes: 113 additions & 0 deletions docs/ga-readiness.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
# a3s-event 0.4.0 — GA readiness audit

Date: 2026-09-30 · Branch: `feat/iggy-provider-ga` (local, not pushed)

This document is the auditable evidence trail for calling this release
production-ready. It separates what is **verified** from what is **gated on
external action**, and names every known limitation. It is intentionally not
a marketing document.

## 1. Verification evidence

### 1.1 Test matrices (all green, re-verified on a live broker)

| Matrix | Result |
|---|---|
| Default features (nats/encryption/cloudevents/routing) | **280 passed / 0 failed** |
| `nats,iggy` (no default) | **246 passed / 0 failed** — against a live `apache/iggy:0.9.0` container |
| Chaos (opt-in env vars) | iggy restart-resume **passed live**; nats restart-recovery **passed live**; both skip cleanly when unset |

Caveat recorded in the chaos suite: **always check the broker is alive after
running chaos** — suites skip-pass against a dead server (skip-if-unavailable
is the harness contract). The upstream restart bug (§3.1) makes this a real
operational footgun, not a theoretical one.

### 1.2 Depth of coverage

- **Cross-provider conformance** (`tests/conformance.rs`): tier-0 (every
provider: envelope fidelity across 3 categories × 4 versions, fan-out
isolation with a 3-subscriber overlap matrix, per-category total order,
8×10 concurrent publish no-loss/no-dup, tail filters, options plumbing,
counts/info/health) and tier-1 (persistent providers: unacked redelivery,
resume across a NEW connection, group-rebuild replay, late-subscriber
ordered replay, competing consumers exactly-once across connections).
- **Iggy contract matrix**: 32 rows, all implemented — including
poison-message tolerance (foreign non-JSON frame skipped without wedging)
and the two fail-closed surfaces (`expected_sequence`,
`LastPerSubject`).
- **Feature e2e**: EventBus full pipeline (schema gate → encryption at rest
→ broker routing → DLQ capture → state persistence across "restart" →
metrics audit), routing/bridge (filter matrix + cross-bus TopicSink),
CronSource lifecycle, crypto key-rotation/tamper, CloudEvents fidelity,
messaging isolation, DLQ capacity/predicate/SinkDlqHandler notification
contract, schema compatibility matrix (stepwise), error paths
(unwritable state store, always-failing provider).
- **Bugs the depth bought** (all fixed, regression-covered): DLQ contract
unwired, wildcard matching dead code, NATS durable-name rejection,
NATS history policy, ByStartTime fallback-to-zero, messaging target
prefix, plus test-semantics fixes (clock-skew midpoint, per-connection
group identity, ack-wait redelivery windows).

### 1.3 Static quality gates

| Gate | Result |
|---|---|
| `cargo clippy --all-targets -- -D warnings` | clean × 4 feature sets (default, nats, iggy, nats+iggy) |
| `cargo fmt --check` | clean |
| Unit coverage (`cargo llvm-cov --lib`) | **82.9% lines** (broker provider bodies exercised by live e2e, not counted) |
| Cross-compile, minimal core | linux x64/arm64 + windows msvc clean (TLS deps need native or C cross-toolchain — CI runs native per-OS) |

### 1.4 Performance baseline (criterion, crate release profile, Apple Silicon)

Memory provider: publish ~203 µs/100-event batch (~2.0 µs/event), ~1.70 ms
per 1000 (~1.7 µs/event); `history(100)` 12.4 µs / 20.8 µs filtered. Broker
round-trips dominate all networked paths.

### 1.5 Packaging

`cargo publish --dry-run` verifies the packaged crate builds standalone and
ships README/LICENSE/CHANGELOG/docs. The stray root-monorepo `.gitmodules`
is excluded via `.gitignore` (never shipped).

## 2. Known limitations (documented, not hidden)

- iggy provider: no broker-side dedup (`msg_id` is header-only); redelivery
controls (`max_deliver`/`backoff`/`max_ack_pending`/`ack_wait`) accepted
and ignored; single-partition topics only (`Balanced` is a documented
no-op); group membership is per client connection.
- History ordering across providers is not part of the contract (memory is
newest-first, brokers oldest-first).
- TLS paths compile but have no e2e coverage.

## 3. Gated on external action (the honest remainder)

### 3.1 Upstream defect gating iggy-GA

Iggy server 0.9.0 intermittently panics on restart boot replay
(`client_id 0 is reserved for internal use`,
`core/consensus/src/client_table.rs`), leaving the server unbootable with
the same data directory. Observed twice locally. Issue draft:
`docs/upstream-iggy-restart-panic.md` (not filed — needs authorization).
**Until fixed upstream, the `iggy` feature must not be called
GA-hardened**, regardless of this crate's own quality.

### 3.2 Owner decisions

- Push `feat/iggy-provider-ga`, wire CI to the remote, first remote run.
- `cargo publish` for real (0.4.0; migration warnings in CHANGELOG §0.4.0).
- How this crate rejoins the a3s monorepo (in-tree vs submodule re-pin) —
the root `.gitmodules` deletion is mid-conversion and is an owners' call.
- Filing the upstream iggy issue.

### 3.3 Time-gated

- Production soak: weeks of real load. No substitute exists.

## 4. Verdict

For the **memory and nats** feature sets: engineering GA criteria are met
(tests, gates, coverage, packaging, docs, migration notes) pending §3.2's
publish/CI wiring. For the **iggy** feature set: same crate-level criteria
are met, but the feature is explicitly **not GA** until §3.1 is resolved
upstream — this is stated in the CHANGELOG and is not negotiable by test
count.
Loading
Loading