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::models::Event;
15use crate::publish_routing::resolve_publish_target;
16use crate::registry::TopicRegistry;
17use crate::storage::{InProcStoragePort, StoragePort};
18
19/// Unified backend — all storage adapters install this wrapping their port.
20///
21/// Prefer selecting a [`StoragePort`](crate::storage::StoragePort) and letting
22/// [`PhotonBuilder`](https://docs.rs/uf-photon/latest/photon/struct.PhotonBuilder.html) install this
23/// automatically. See also the [`EmbeddedBackend`](crate::EmbeddedBackend) alias.
24pub struct GenericPhotonBackend {
25    port: Arc<dyn StoragePort>,
26    registry: TopicRegistry,
27    capabilities: BackendCapabilities,
28}
29
30impl GenericPhotonBackend {
31    /// Wrap an existing storage port and topic registry.
32    #[must_use]
33    pub fn new(port: Arc<dyn StoragePort>, registry: TopicRegistry) -> Self {
34        let capabilities = port.capabilities();
35        Self {
36            port,
37            registry,
38            capabilities,
39        }
40    }
41
42    /// Capabilities for contract tests and host introspection.
43    #[must_use]
44    pub const fn capabilities(&self) -> BackendCapabilities {
45        self.capabilities
46    }
47
48    /// Install the default `mem` tier from builder context.
49    ///
50    /// Loads `PHOTON_TRANSPORT_KEY` via [`TransportCrypto::from_env`](crate::event::TransportCrypto::from_env).
51    ///
52    /// # Errors
53    ///
54    /// Returns an error if the transport key cannot be loaded from the environment.
55    pub fn install_mem(ctx: BackendContext) -> Result<Arc<dyn PhotonBackend>> {
56        let port = Arc::new(InProcStoragePort::new(
57            crate::event::TransportCrypto::from_env()?,
58        ));
59        Ok(Arc::new(Self::new(port, ctx.registry)))
60    }
61
62    /// Install a custom storage port from builder context.
63    ///
64    /// # Errors
65    ///
66    /// Returns an error if the operation fails.
67    pub fn install_with_port(
68        ctx: BackendContext,
69        port: Arc<dyn StoragePort>,
70    ) -> Result<Arc<dyn PhotonBackend>> {
71        Ok(Arc::new(Self::new(port, ctx.registry)))
72    }
73}
74
75#[async_trait]
76impl PhotonBackend for GenericPhotonBackend {
77    fn telemetry_label(&self) -> &'static str {
78        self.capabilities.telemetry_label
79    }
80
81    fn capabilities(&self) -> BackendCapabilities {
82        self.capabilities
83    }
84
85    async fn publish(
86        &self,
87        topic_name: &str,
88        topic_key: Option<&str>,
89        actor_json: Value,
90        payload_json: Value,
91    ) -> Result<String> {
92        let target = resolve_publish_target(
93            &self.registry,
94            topic_name,
95            topic_key,
96            &payload_json,
97        );
98        let event = self
99            .port
100            .append(
101                topic_name,
102                target.topic_key.as_deref(),
103                actor_json,
104                payload_json,
105            )
106            .await?;
107        Ok(event.event_id)
108    }
109
110    fn subscribe(
111        &self,
112        topic_name: String,
113        topic_key_filter: Option<String>,
114        after_seq: Option<i64>,
115    ) -> Pin<Box<dyn Stream<Item = Result<Event>> + Send>> {
116        self.port
117            .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}