CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s
Rename all `corex-cache`, `corex-storage`, `corex-event`, `corex-config`, and `corex-crypto` crates to `mytheclipse-cache`, `mytheclipse-storage`, `mytheclipse-event`, `mytheclipse-config`, and `mytheclipse-crypto` respectively, aligning with the `mytheclipse` root crate naming convention. - Convert root `Cargo.toml` from a package manifest to a workspace manifest with explicit member paths - Update crate names in each `Cargo.toml` and docs.rs URLs - Update all import paths, module doc comments, and error/panic messages from `corex-*` to `mytheclipse-*` across source and README files - Update CI workflow and publish script to reference the new crate names<tool_call></think>chore: rename corex-* crates to mytheclipse-* namespace Rename all `corex-cache`, `corex-storage`, `corex-event`, `corex-config`, and `corex-crypto` crates to `mytheclipse-cache`, `mytheclipse-storage`, `mytheclipse-event`, `mytheclipse-config`, and `mytheclipse-crypto` respectively, aligning with the `mytheclipse` root crate naming convention. - Convert root `Cargo.toml` from a package manifest to a workspace manifest with explicit member paths - Update crate names in each `Cargo.toml` and docs.rs URLs - Update all import paths, module doc comments, and error/panic messages from `corex-*` to `mytheclipse-*` across source and README files - Update CI workflow and publish script to reference the new crate names
96 lines
2.9 KiB
Rust
96 lines
2.9 KiB
Rust
//! A NATS-backed [`EventBus`] (feature `nats`), via the pure-Rust
|
|
//! `async-nats` client.
|
|
//!
|
|
//! Topics map directly onto NATS subjects.
|
|
|
|
use async_trait::async_trait;
|
|
use futures_util::StreamExt;
|
|
|
|
use crate::traits::{EventBus, EventError, Subscription, SubscriptionInner};
|
|
|
|
fn map_connect_err(e: async_nats::ConnectError) -> EventError {
|
|
EventError::Io(e.to_string())
|
|
}
|
|
|
|
fn map_publish_err(e: async_nats::PublishError) -> EventError {
|
|
EventError::Io(e.to_string())
|
|
}
|
|
|
|
fn map_subscribe_err(e: async_nats::SubscribeError) -> EventError {
|
|
EventError::Io(e.to_string())
|
|
}
|
|
|
|
/// An [`EventBus`] backed by a NATS client.
|
|
#[derive(Clone)]
|
|
pub struct NatsEventBus {
|
|
client: async_nats::Client,
|
|
}
|
|
|
|
impl NatsEventBus {
|
|
/// Connects to a NATS server at `addr` (e.g. `nats://127.0.0.1:4222`).
|
|
pub async fn connect(addr: &str) -> Result<Self, EventError> {
|
|
let client = async_nats::connect(addr).await.map_err(map_connect_err)?;
|
|
Ok(Self { client })
|
|
}
|
|
|
|
/// Wraps an already-connected client.
|
|
pub fn from_client(client: async_nats::Client) -> Self {
|
|
Self { client }
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl EventBus for NatsEventBus {
|
|
async fn publish(&self, topic: &str, payload: Vec<u8>) -> Result<(), EventError> {
|
|
self.client
|
|
.publish(topic.to_string(), payload.into())
|
|
.await
|
|
.map_err(map_publish_err)
|
|
}
|
|
|
|
async fn subscribe(&self, topic: &str) -> Result<Subscription, EventError> {
|
|
let subscriber = self
|
|
.client
|
|
.subscribe(topic.to_string())
|
|
.await
|
|
.map_err(map_subscribe_err)?;
|
|
Ok(Subscription::new(Box::new(NatsSubscription { subscriber })))
|
|
}
|
|
}
|
|
|
|
struct NatsSubscription {
|
|
subscriber: async_nats::Subscriber,
|
|
}
|
|
|
|
#[async_trait]
|
|
impl SubscriptionInner for NatsSubscription {
|
|
async fn recv(&mut self) -> Result<Vec<u8>, EventError> {
|
|
match self.subscriber.next().await {
|
|
Some(msg) => Ok(msg.payload.to_vec()),
|
|
None => Err(EventError::Closed),
|
|
}
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
/// Requires a live NATS server at `NATS_URL` (e.g. `nats://127.0.0.1:4222`).
|
|
/// Run with: `NATS_URL=... cargo test -p mytheclipse-event --features nats -- --ignored`.
|
|
#[tokio::test]
|
|
#[ignore = "requires a live NATS instance (NATS_URL)"]
|
|
async fn publish_subscribe_roundtrip_live() {
|
|
let url = std::env::var("NATS_URL").expect("set NATS_URL");
|
|
let bus = NatsEventBus::connect(&url).await.unwrap();
|
|
let mut sub = bus.subscribe("orders.created").await.unwrap();
|
|
// Give the subscription a moment to register server-side.
|
|
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
|
|
bus.publish("orders.created", b"hello".to_vec())
|
|
.await
|
|
.unwrap();
|
|
let msg = sub.recv().await.unwrap();
|
|
assert_eq!(msg, b"hello");
|
|
}
|
|
}
|