pub struct NatsStoragePort { /* private fields */ }Expand description
NATS JetStream-backed storage port.
Implementations§
Source§impl NatsStoragePort
impl NatsStoragePort
Sourcepub fn builder() -> NatsStoragePortBuilder
pub fn builder() -> NatsStoragePortBuilder
Start a builder for explicit host wiring.
Sourcepub async fn from_env() -> Result<NatsStoragePort, PhotonError>
pub async fn from_env() -> Result<NatsStoragePort, PhotonError>
Connect using env (PHOTON_NATS_* defaults via builder).
§Errors
Returns an error when env is missing or connection fails.
Sourcepub async fn connect(
url: &str,
stream_name: &str,
) -> Result<NatsStoragePort, PhotonError>
pub async fn connect( url: &str, stream_name: &str, ) -> Result<NatsStoragePort, PhotonError>
Connect to NATS with explicit URL and stream name (legacy; uses env defaults for firehose options).
§Errors
Returns an error when connection or stream setup fails.
Sourcepub const fn config(&self) -> &NatsConfig
pub const fn config(&self) -> &NatsConfig
Resolved adapter configuration.
Trait Implementations§
Source§impl StoragePort for NatsStoragePort
impl StoragePort for NatsStoragePort
Source§fn capabilities(&self) -> StorageCapabilities
fn capabilities(&self) -> StorageCapabilities
Adapter capabilities for contract tests and telemetry.
Source§fn append<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
topic_name: &'life1 str,
topic_key: Option<&'life2 str>,
actor_json: Value,
payload_json: Value,
) -> Pin<Box<dyn Future<Output = Result<Event, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
NatsStoragePort: 'async_trait,
fn append<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
topic_name: &'life1 str,
topic_key: Option<&'life2 str>,
actor_json: Value,
payload_json: Value,
) -> Pin<Box<dyn Future<Output = Result<Event, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
NatsStoragePort: 'async_trait,
Append one event to a topic partition. Read more
Source§fn subscribe(
&self,
topic_name: String,
topic_key_filter: Option<String>,
after_seq: Option<i64>,
) -> Pin<Box<dyn Stream<Item = Result<Event, PhotonError>> + Send>>
fn subscribe( &self, topic_name: String, topic_key_filter: Option<String>, after_seq: Option<i64>, ) -> Pin<Box<dyn Stream<Item = Result<Event, PhotonError>> + Send>>
Stream events for a topic partition, optionally replaying after
after_seq. Read moreSource§fn get_event<'life0, 'life1, 'async_trait>(
&'life0 self,
_event_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Event>, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
NatsStoragePort: 'async_trait,
fn get_event<'life0, 'life1, 'async_trait>(
&'life0 self,
_event_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<Event>, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
NatsStoragePort: 'async_trait,
Point lookup by event id (optional — see
StorageCapabilities::supports_get_event). Read moreSource§fn load_checkpoint<'life0, 'life1, 'life2, 'life3, 'async_trait>(
&'life0 self,
subscription_name: &'life1 str,
topic_name: &'life2 str,
topic_key: Option<&'life3 str>,
) -> Pin<Box<dyn Future<Output = Result<Option<i64>, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
NatsStoragePort: 'async_trait,
fn load_checkpoint<'life0, 'life1, 'life2, 'life3, 'async_trait>(
&'life0 self,
subscription_name: &'life1 str,
topic_name: &'life2 str,
topic_key: Option<&'life3 str>,
) -> Pin<Box<dyn Future<Output = Result<Option<i64>, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
NatsStoragePort: 'async_trait,
Load durable subscription high-water seq. Read more
Source§fn commit_checkpoint<'life0, 'life1, 'life2, 'life3, 'async_trait>(
&'life0 self,
subscription_name: &'life1 str,
topic_name: &'life2 str,
topic_key: Option<&'life3 str>,
last_seq: i64,
) -> Pin<Box<dyn Future<Output = Result<(), PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
NatsStoragePort: 'async_trait,
fn commit_checkpoint<'life0, 'life1, 'life2, 'life3, 'async_trait>(
&'life0 self,
subscription_name: &'life1 str,
topic_name: &'life2 str,
topic_key: Option<&'life3 str>,
last_seq: i64,
) -> Pin<Box<dyn Future<Output = Result<(), PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
NatsStoragePort: 'async_trait,
Persist durable subscription high-water seq. Read more
Source§fn truncate_before<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
topic_name: &'life1 str,
topic_key: Option<&'life2 str>,
truncate_bound: i64,
) -> Pin<Box<dyn Future<Output = Result<u64, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Self: 'async_trait,
fn truncate_before<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
topic_name: &'life1 str,
topic_key: Option<&'life2 str>,
truncate_bound: i64,
) -> Pin<Box<dyn Future<Output = Result<u64, PhotonError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Self: 'async_trait,
Trim events before
truncate_bound for a partition (no-op when unsupported). Read moreSource§fn delivery_seq_pin<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
topic_name: &'life1 str,
topic_key: Option<&'life2 str>,
) -> Pin<Box<dyn Future<Output = Option<i64>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Self: 'async_trait,
fn delivery_seq_pin<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
topic_name: &'life1 str,
topic_key: Option<&'life2 str>,
) -> Pin<Box<dyn Future<Output = Option<i64>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Self: 'async_trait,
Last delivered seq pin for retention watermarks (optional). Read more
Auto Trait Implementations§
impl !RefUnwindSafe for NatsStoragePort
impl !UnwindSafe for NatsStoragePort
impl Freeze for NatsStoragePort
impl Send for NatsStoragePort
impl Sync for NatsStoragePort
impl Unpin for NatsStoragePort
impl UnsafeUnpin for NatsStoragePort
Blanket Implementations§
impl<T> AsyncConnector for T
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
Converts
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
Converts
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more