Skip to main content

photon_backend/backend/
generic.rs

1//! Single [`PhotonBackend`] implementation wrapping any [`StoragePort`].
2
3use std::pin::Pin;
4use std::sync::Arc;
5
6use async_trait::async_trait;
7use futures::stream::Stream;
8use serde_json::Value;
9
10use super::capabilities::BackendCapabilities;
11use super::context::BackendContext;
12use super::photon_backend::PhotonBackend;
13use crate::error::Result;
14use crate::input::{validate_payload_size, validate_topic_name};
15use crate::models::Event;
16use crate::publish_routing::resolve_publish_target;
17use crate::registry::TopicRegistry;
18use crate::storage::{InProcStoragePort, StoragePort};
19
20/// Unified backend — all storage adapters install this wrapping their port.
21///
22/// Prefer selecting a [`StoragePort`](crate::storage::StoragePort) and letting
23/// [`PhotonBuilder`](https://docs.rs/uf-photon/latest/photon/struct.PhotonBuilder.html) install this
24/// automatically. See also the [`EmbeddedBackend`](crate::EmbeddedBackend) alias.
25pub struct GenericPhotonBackend {
26    port: Arc<dyn StoragePort>,
27    registry: TopicRegistry,
28    capabilities: BackendCapabilities,
29}
30
31impl GenericPhotonBackend {
32    /// Wrap an existing storage port and topic registry.
33    #[must_use]
34    pub fn new(port: Arc<dyn StoragePort>, registry: TopicRegistry) -> Self {
35        let capabilities = port.capabilities();
36        Self {
37            port,
38            registry,
39            capabilities,
40        }
41    }
42
43    /// Capabilities for contract tests and host introspection.
44    #[must_use]
45    pub const fn capabilities(&self) -> BackendCapabilities {
46        self.capabilities
47    }
48
49    /// Install the default `mem` tier from builder context.
50    ///
51    /// Loads `PHOTON_TRANSPORT_KEY` via [`TransportCrypto::from_env`](crate::event::TransportCrypto::from_env).
52    ///
53    /// # Errors
54    ///
55    /// Returns an error if the transport key cannot be loaded from the environment.
56    pub fn install_mem(ctx: BackendContext) -> Result<Arc<dyn PhotonBackend>> {
57        let port = Arc::new(InProcStoragePort::new(
58            crate::event::TransportCrypto::from_env()?,
59        ));
60        Ok(Arc::new(Self::new(port, ctx.registry)))
61    }
62
63    /// Install a custom storage port from builder context.
64    ///
65    /// # Errors
66    ///
67    /// Returns an error if the operation fails.
68    pub fn install_with_port(
69        ctx: BackendContext,
70        port: Arc<dyn StoragePort>,
71    ) -> Result<Arc<dyn PhotonBackend>> {
72        Ok(Arc::new(Self::new(port, ctx.registry)))
73    }
74}
75
76#[async_trait]
77impl PhotonBackend for GenericPhotonBackend {
78    fn telemetry_label(&self) -> &'static str {
79        self.capabilities.telemetry_label
80    }
81
82    fn capabilities(&self) -> BackendCapabilities {
83        self.capabilities
84    }
85
86    async fn publish(
87        &self,
88        topic_name: &str,
89        topic_key: Option<&str>,
90        actor_json: Value,
91        payload_json: Value,
92    ) -> Result<String> {
93        validate_topic_name(topic_name)?;
94        validate_payload_size(&payload_json)?;
95        let target = resolve_publish_target(&self.registry, topic_name, topic_key, &payload_json);
96        let event = self
97            .port
98            .append(
99                topic_name,
100                target.topic_key.as_deref(),
101                actor_json,
102                payload_json,
103            )
104            .await?;
105        Ok(event.event_id)
106    }
107
108    fn subscribe(
109        &self,
110        topic_name: String,
111        topic_key_filter: Option<String>,
112        after_seq: Option<i64>,
113    ) -> Pin<Box<dyn Stream<Item = Result<Event>> + Send>> {
114        if let Err(error) = validate_topic_name(&topic_name) {
115            return Box::pin(futures::stream::once(async move { Err(error) }));
116        }
117        self.port.subscribe(topic_name, topic_key_filter, after_seq)
118    }
119
120    async fn get_event(&self, event_id: &str) -> Result<Option<Event>> {
121        if !self.capabilities.supports_get_event {
122            return Ok(None);
123        }
124        self.port.get_event(event_id).await
125    }
126
127    fn registry(&self) -> &TopicRegistry {
128        &self.registry
129    }
130
131    async fn get_checkpoint_seq(
132        &self,
133        subscription_name: &str,
134        topic_name: &str,
135        topic_key: Option<&str>,
136    ) -> Result<Option<i64>> {
137        self.port
138            .load_checkpoint(subscription_name, topic_name, topic_key)
139            .await
140    }
141
142    async fn set_checkpoint(
143        &self,
144        subscription_name: &str,
145        topic_name: &str,
146        topic_key: Option<&str>,
147        last_seq: i64,
148    ) -> Result<()> {
149        self.port
150            .commit_checkpoint(subscription_name, topic_name, topic_key, last_seq)
151            .await
152    }
153}