nodedb_lite/engine/htap/
bridge.rs1use std::collections::HashMap;
14use std::sync::Mutex;
15use std::time::{SystemTime, UNIX_EPOCH};
16
17use nodedb_types::value::Value;
18
19use crate::engine::columnar::ColumnarEngine;
20use crate::storage::engine::StorageEngine;
21
22#[derive(Debug, Clone)]
24pub struct MaterializedView {
25 pub source: String,
27 pub target: String,
29 pub last_replicated_ms: u64,
31 pub rows_replicated: u64,
33}
34
35pub struct HtapBridge {
41 views: HashMap<String, Vec<MaterializedView>>,
43}
44
45impl HtapBridge {
46 pub fn new() -> Self {
48 Self {
49 views: HashMap::new(),
50 }
51 }
52
53 pub fn register_view(&mut self, source: &str, target: &str) {
57 let view = MaterializedView {
58 source: source.to_string(),
59 target: target.to_string(),
60 last_replicated_ms: now_ms(),
61 rows_replicated: 0,
62 };
63 self.views.entry(source.to_string()).or_default().push(view);
64 }
65
66 pub fn remove_view(&mut self, target: &str) {
68 for views in self.views.values_mut() {
69 views.retain(|v| v.target != target);
70 }
71 self.views.retain(|_, views| !views.is_empty());
72 }
73
74 pub fn views_for_source(&self, source: &str) -> &[MaterializedView] {
76 self.views.get(source).map(|v| v.as_slice()).unwrap_or(&[])
77 }
78
79 pub fn view_by_target(&self, target: &str) -> Option<&MaterializedView> {
81 self.views.values().flatten().find(|v| v.target == target)
82 }
83
84 pub fn all_targets(&self) -> Vec<&str> {
86 self.views
87 .values()
88 .flatten()
89 .map(|v| v.target.as_str())
90 .collect()
91 }
92
93 pub fn replicate_insert<S: StorageEngine>(
99 &mut self,
100 source: &str,
101 values: &[Value],
102 columnar: &Mutex<ColumnarEngine<S>>,
103 ) {
104 let Some(views) = self.views.get_mut(source) else {
105 return;
106 };
107
108 let mut engine = match columnar.lock() {
109 Ok(e) => e,
110 Err(p) => p.into_inner(),
111 };
112
113 for view in views.iter_mut() {
114 if engine.insert(&view.target, values).is_ok() {
115 view.rows_replicated += 1;
116 view.last_replicated_ms = now_ms();
117 }
118 }
119 }
120
121 pub fn replicate_delete<S: StorageEngine>(
124 &mut self,
125 source: &str,
126 pk: &Value,
127 columnar: &Mutex<ColumnarEngine<S>>,
128 ) {
129 let Some(views) = self.views.get_mut(source) else {
130 return;
131 };
132
133 let mut engine = match columnar.lock() {
134 Ok(e) => e,
135 Err(p) => p.into_inner(),
136 };
137
138 for view in views.iter_mut() {
139 if engine.delete(&view.target, pk).unwrap_or(false) {
140 view.last_replicated_ms = now_ms();
141 }
142 }
143 }
144
145 pub fn lag_ms(&self, target: &str) -> u64 {
147 self.view_by_target(target)
148 .map(|v| now_ms().saturating_sub(v.last_replicated_ms))
149 .unwrap_or(0)
150 }
151
152 pub fn is_empty(&self) -> bool {
154 self.views.is_empty()
155 }
156}
157
158impl Default for HtapBridge {
159 fn default() -> Self {
160 Self::new()
161 }
162}
163
164fn now_ms() -> u64 {
165 SystemTime::now()
166 .duration_since(UNIX_EPOCH)
167 .map(|d| d.as_millis() as u64)
168 .unwrap_or(0)
169}
170
171#[cfg(test)]
172mod tests {
173 use super::*;
174
175 #[test]
176 fn register_and_lookup_view() {
177 let mut bridge = HtapBridge::new();
178 bridge.register_view("customers", "customer_analytics");
179
180 assert!(!bridge.is_empty());
181 assert_eq!(bridge.views_for_source("customers").len(), 1);
182 assert_eq!(
183 bridge.views_for_source("customers")[0].target,
184 "customer_analytics"
185 );
186 assert!(bridge.view_by_target("customer_analytics").is_some());
187 assert!(bridge.view_by_target("nonexistent").is_none());
188 }
189
190 #[test]
191 fn remove_view() {
192 let mut bridge = HtapBridge::new();
193 bridge.register_view("customers", "analytics_1");
194 bridge.register_view("customers", "analytics_2");
195
196 assert_eq!(bridge.views_for_source("customers").len(), 2);
197
198 bridge.remove_view("analytics_1");
199 assert_eq!(bridge.views_for_source("customers").len(), 1);
200 assert_eq!(
201 bridge.views_for_source("customers")[0].target,
202 "analytics_2"
203 );
204 }
205
206 #[test]
207 fn multiple_sources() {
208 let mut bridge = HtapBridge::new();
209 bridge.register_view("orders", "order_analytics");
210 bridge.register_view("customers", "customer_analytics");
211
212 assert_eq!(bridge.all_targets().len(), 2);
213 }
214}