photon_backend/backend/
generic.rs1use 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
20pub struct GenericPhotonBackend {
26 port: Arc<dyn StoragePort>,
27 registry: TopicRegistry,
28 capabilities: BackendCapabilities,
29}
30
31impl GenericPhotonBackend {
32 #[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 #[must_use]
45 pub const fn capabilities(&self) -> BackendCapabilities {
46 self.capabilities
47 }
48
49 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 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}