pub struct InProcStoragePort { /* private fields */ }Expand description
In-process storage for the mem adapter tier.
Single-process only: live fanout and a bounded replay buffer live in memory. Photon installs
this by default when you omit
PhotonBuilder::storage_port.
When not to use: multiple binaries or hosts must share a topic log — pick NATS/Kafka/Fluvio
(Mode 2) or SQLite for durable single-host Mode 1.
Getting started: Mode 1.
§Examples
§Mode 1 host (publish + handlers)
One binary owns both publish and #[subscribe] dispatch via start_executor.
ⓘ
use std::sync::Arc;
use photon_backend::{InProcStoragePort, StoragePort, TransportCrypto};
use photon_core::JsonIdentityFactory;
use photon_runtime::Photon;
let port: Arc<dyn StoragePort> = Arc::new(InProcStoragePort::new(
TransportCrypto::from_env()?,
));
let photon = Photon::builder()
.storage_port(port)
.auto_registry()
.build()?;
photon.start_executor(Arc::new(JsonIdentityFactory))?;
// EventType { … }.publish_on(&photon).await?;Default path (same mem port, no explicit storage_port):
Photon::builder().auto_registry().build()?.
Implementations§
Source§impl InProcStoragePort
impl InProcStoragePort
Sourcepub fn new(crypto: TransportCrypto) -> InProcStoragePort
pub fn new(crypto: TransportCrypto) -> InProcStoragePort
Create a new in-process port with default replay buffer size.
Trait Implementations§
Source§impl StoragePort for InProcStoragePort
impl StoragePort for InProcStoragePort
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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: '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,
InProcStoragePort: 'async_trait,
Last delivered seq pin for retention watermarks (optional). Read more
Auto Trait Implementations§
impl !RefUnwindSafe for InProcStoragePort
impl !UnwindSafe for InProcStoragePort
impl Freeze for InProcStoragePort
impl Send for InProcStoragePort
impl Sync for InProcStoragePort
impl Unpin for InProcStoragePort
impl UnsafeUnpin for InProcStoragePort
Blanket Implementations§
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