backbone_core/
integration.rs1use async_trait::async_trait;
25use std::marker::PhantomData;
26
27#[async_trait]
37pub trait ModuleAdapter<External, Internal>: Send + Sync
38where
39 External: Send + 'static,
40 Internal: Send + 'static,
41{
42 type Error: std::error::Error + Send + Sync;
44
45 fn source_module() -> &'static str
49 where
50 Self: Sized,
51 {
52 "unknown"
53 }
54
55 fn target_module() -> &'static str
59 where
60 Self: Sized,
61 {
62 "unknown"
63 }
64
65 async fn to_internal(&self, external: External) -> Result<Internal, Self::Error>;
67
68 async fn to_external(&self, _internal: Internal) -> Result<External, Self::Error> {
72 Err(self.not_implemented("to_external"))
73 }
74
75 fn not_implemented(&self, method: &str) -> Self::Error;
76}
77
78pub trait ProjectionAdapter<External, Internal>: Send + Sync {
82 fn project(&self, external: &External) -> Internal;
83}
84
85#[derive(Debug, thiserror::Error)]
89pub enum IntegrationError {
90 #[error("mapping failed: {0}")]
91 MappingFailed(String),
92
93 #[error("external service unavailable: {0}")]
94 ServiceUnavailable(String),
95
96 #[error("schema mismatch: {0}")]
97 SchemaMismatch(String),
98
99 #[error("not implemented: {0}")]
100 NotImplemented(String),
101}
102
103pub struct EventBridge<External, Internal, A>
109where
110 External: Send + 'static,
111 Internal: Send + 'static,
112 A: ModuleAdapter<External, Internal>,
113{
114 adapter: A,
115 _phantom: PhantomData<(External, Internal)>,
116}
117
118impl<External, Internal, A> EventBridge<External, Internal, A>
119where
120 External: Send + 'static,
121 Internal: Send + 'static,
122 A: ModuleAdapter<External, Internal>,
123{
124 pub fn new(adapter: A) -> Self {
125 Self {
126 adapter,
127 _phantom: PhantomData,
128 }
129 }
130
131 pub async fn process(
133 &self,
134 event: External,
135 ) -> Result<Internal, <A as ModuleAdapter<External, Internal>>::Error> {
136 self.adapter.to_internal(event).await
137 }
138}
139
140pub struct IdentityAdapter<T>(PhantomData<T>);
145
146impl<T> IdentityAdapter<T> {
147 pub fn new() -> Self {
148 Self(PhantomData)
149 }
150}
151
152impl<T> Default for IdentityAdapter<T> {
153 fn default() -> Self {
154 Self::new()
155 }
156}
157
158#[async_trait]
159impl<T: Clone + Send + Sync + 'static> ModuleAdapter<T, T> for IdentityAdapter<T>
160{
161 type Error = IntegrationError;
162
163 async fn to_internal(&self, external: T) -> Result<T, Self::Error> {
164 Ok(external)
165 }
166
167 async fn to_external(&self, internal: T) -> Result<T, Self::Error> {
168 Ok(internal)
169 }
170
171 fn not_implemented(&self, method: &str) -> Self::Error {
172 IntegrationError::NotImplemented(method.into())
173 }
174}
175
176pub fn identity_adapter<T: Clone + Send + Sync + 'static>() -> IdentityAdapter<T> {
178 IdentityAdapter::new()
179}
180
181#[cfg(test)]
182mod tests {
183 use super::*;
184
185 #[derive(Debug, Clone, PartialEq)]
186 struct ExternalUserEvent {
187 user_id: String,
188 email: String,
189 }
190
191 #[derive(Debug, Clone, PartialEq)]
192 struct InternalCustomerEvent {
193 customer_id: String,
194 contact_email: String,
195 }
196
197 struct UserToCustomerAdapter;
198
199 #[async_trait]
200 impl ModuleAdapter<ExternalUserEvent, InternalCustomerEvent> for UserToCustomerAdapter
201 {
202 type Error = IntegrationError;
203
204 fn source_module() -> &'static str {
205 "sapiens"
206 }
207
208 fn target_module() -> &'static str {
209 "corpus"
210 }
211
212 async fn to_internal(
213 &self,
214 external: ExternalUserEvent,
215 ) -> Result<InternalCustomerEvent, Self::Error> {
216 Ok(InternalCustomerEvent {
217 customer_id: external.user_id,
218 contact_email: external.email,
219 })
220 }
221
222 fn not_implemented(&self, method: &str) -> Self::Error {
223 IntegrationError::NotImplemented(method.into())
224 }
225 }
226
227 #[tokio::test]
228 async fn adapter_maps_fields_correctly() {
229 let bridge = EventBridge::new(UserToCustomerAdapter);
230 let external = ExternalUserEvent {
231 user_id: "u-1".into(),
232 email: "user@example.com".into(),
233 };
234
235 let internal = bridge.process(external).await.unwrap();
236 assert_eq!(internal.customer_id, "u-1");
237 assert_eq!(internal.contact_email, "user@example.com");
238 }
239
240 #[tokio::test]
241 async fn identity_adapter_roundtrips() {
242 let adapter: IdentityAdapter<String> = IdentityAdapter::new();
243 let value = "hello".to_string();
244 let out = adapter.to_internal(value.clone()).await.unwrap();
245 assert_eq!(out, value);
246 }
247
248 #[test]
249 fn module_names_are_exposed() {
250 assert_eq!(UserToCustomerAdapter::source_module(), "sapiens");
251 assert_eq!(UserToCustomerAdapter::target_module(), "corpus");
252 }
253}