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::models::Event;
15use crate::publish_routing::resolve_publish_target;
16use crate::registry::TopicRegistry;
17use crate::storage::{InProcStoragePort, StoragePort};
18
19pub struct GenericPhotonBackend {
25 port: Arc<dyn StoragePort>,
26 registry: TopicRegistry,
27 capabilities: BackendCapabilities,
28}
29
30impl GenericPhotonBackend {
31 #[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 #[must_use]
44 pub const fn capabilities(&self) -> BackendCapabilities {
45 self.capabilities
46 }
47
48 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 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}