1use std::cell::RefCell;
7use std::collections::VecDeque;
8use std::panic::{catch_unwind, AssertUnwindSafe};
9use std::rc::Rc;
10
11use crate::graph::{Graph, GraphObserver, ObserveEvent};
12use crate::storage::{
13 observe_event_frame, AppendLogStorageTier, ObserveEventFrame, ObserveEventFrameOptions,
14 StorageError, StorageResult,
15};
16
17#[derive(Clone, Debug, Eq, PartialEq)]
18pub enum ObserveEventLogErrorPhase {
20 Map,
22 Write,
24 Flush,
26 Rollback,
28 Dispose,
30}
31
32#[derive(Clone)]
33pub struct ObserveEventLogErrorContext<T> {
35 pub phase: ObserveEventLogErrorPhase,
37 pub event: Option<ObserveEvent>,
39 pub value: Option<ObserveEventFrame<T>>,
41}
42
43pub type ObserveEventLogMap<T> = dyn Fn(&ObserveEvent) -> Option<T>;
45pub type ObserveEventLogErrorFn<T> = dyn Fn(StorageError, ObserveEventLogErrorContext<T>);
47
48#[derive(Clone)]
49pub struct AttachObserveEventLogOptions<T: Clone> {
51 pub path: Option<String>,
53 pub stream: Option<String>,
55 pub map: Rc<ObserveEventLogMap<T>>,
57 pub on_error: Option<Rc<ObserveEventLogErrorFn<T>>>,
59}
60
61impl<T> AttachObserveEventLogOptions<T>
62where
63 T: Clone + From<ObserveEvent> + 'static,
64{
65 pub fn new() -> Self {
67 Self {
68 path: None,
69 stream: None,
70 map: Rc::new(|event| Some(event.clone().into())),
71 on_error: None,
72 }
73 }
74}
75
76impl<T> Default for AttachObserveEventLogOptions<T>
77where
78 T: Clone + From<ObserveEvent> + 'static,
79{
80 fn default() -> Self {
81 Self::new()
82 }
83}
84
85impl<T: Clone> AttachObserveEventLogOptions<T> {
86 pub fn from_map(map: impl Fn(&ObserveEvent) -> Option<T> + 'static) -> Self {
88 Self {
89 path: None,
90 stream: None,
91 map: Rc::new(map),
92 on_error: None,
93 }
94 }
95
96 pub fn with_path(mut self, path: impl Into<String>) -> Self {
98 self.path = Some(path.into());
99 self
100 }
101
102 pub fn with_stream(mut self, stream: impl Into<String>) -> Self {
104 self.stream = Some(stream.into());
105 self
106 }
107
108 pub fn with_map(mut self, map: impl Fn(&ObserveEvent) -> Option<T> + 'static) -> Self {
110 self.map = Rc::new(map);
111 self
112 }
113
114 pub fn with_on_error(
116 mut self,
117 on_error: impl Fn(StorageError, ObserveEventLogErrorContext<T>) + 'static,
118 ) -> Self {
119 self.on_error = Some(Rc::new(on_error));
120 self
121 }
122}
123
124pub struct ObserveEventLogHandle {
126 observer: Option<GraphObserver>,
127 flush: Rc<dyn Fn() -> StorageResult<()>>,
128 rollback: Rc<dyn Fn() -> StorageResult<()>>,
129}
130
131impl ObserveEventLogHandle {
132 pub fn flush(&self) -> StorageResult<()> {
134 (self.flush)()
135 }
136
137 pub fn rollback(&self) -> StorageResult<()> {
139 (self.rollback)()
140 }
141
142 pub fn dispose(&mut self) -> StorageResult<()> {
144 self.observer.take();
145 (self.flush)()
146 }
147}
148
149impl Drop for ObserveEventLogHandle {
150 fn drop(&mut self) {
151 let _ = self.dispose();
152 }
153}
154
155pub fn attach_observe_event_log<T: Clone + 'static>(
157 graph: &Graph,
158 log: Rc<dyn AppendLogStorageTier<ObserveEventFrame<T>>>,
159 opts: AttachObserveEventLogOptions<T>,
160) -> ObserveEventLogHandle {
161 let stream = match &opts.path {
162 Some(path) => graph.observe_path(path),
163 None => graph.observe(),
164 };
165 let map = opts.map.clone();
166 let on_error = opts.on_error.clone();
167 let frame_opts = ObserveEventFrameOptions {
168 stream: opts.stream.clone(),
169 };
170 let pending = Rc::new(RefCell::new(
171 VecDeque::<(ObserveEvent, ObserveEventFrame<T>)>::new(),
172 ));
173 let flush_pending = {
174 let pending = pending.clone();
175 let log = log.clone();
176 let on_error = on_error.clone();
177 Rc::new(move || {
178 let mut first_error = None;
179 loop {
180 let Some((event, frame)) = pending.borrow().front().cloned() else {
181 break;
182 };
183 if let Err(error) = log.append(frame.clone()) {
184 report_observe_event_log_error(
185 &on_error,
186 error.clone(),
187 ObserveEventLogErrorContext {
188 phase: ObserveEventLogErrorPhase::Write,
189 event: Some(event),
190 value: Some(frame),
191 },
192 );
193 if first_error.is_none() {
194 first_error = Some(error);
195 }
196 break;
197 }
198 pending.borrow_mut().pop_front();
199 }
200 match first_error {
201 Some(error) => Err(error),
202 None => Ok(()),
203 }
204 }) as Rc<dyn Fn() -> StorageResult<()>>
205 };
206 let rollback_pending = {
207 let pending = pending.clone();
208 Rc::new(move || {
209 pending.borrow_mut().clear();
210 Ok(())
211 }) as Rc<dyn Fn() -> StorageResult<()>>
212 };
213 let observer = stream.subscribe(move |event| {
214 let mapped = match catch_unwind(AssertUnwindSafe(|| (map)(&event))) {
215 Ok(mapped) => mapped,
216 Err(_) => {
217 report_observe_event_log_error(
218 &on_error,
219 StorageError::backend("attach_observe_event_log: map panicked"),
220 ObserveEventLogErrorContext {
221 phase: ObserveEventLogErrorPhase::Map,
222 event: Some(event),
223 value: None,
224 },
225 );
226 return;
227 }
228 };
229 let Some(value) = mapped else {
230 return;
231 };
232 let frame = match observe_event_frame(
233 event.seq,
234 event.path.clone(),
235 value,
236 ObserveEventFrameOptions {
237 stream: frame_opts.stream.clone(),
238 },
239 ) {
240 Ok(frame) => frame,
241 Err(error) => {
242 report_observe_event_log_error(
243 &on_error,
244 error,
245 ObserveEventLogErrorContext {
246 phase: ObserveEventLogErrorPhase::Map,
247 event: Some(event),
248 value: None,
249 },
250 );
251 return;
252 }
253 };
254 pending.borrow_mut().push_back((event, frame));
255 });
256 ObserveEventLogHandle {
257 observer: Some(observer),
258 flush: flush_pending,
259 rollback: rollback_pending,
260 }
261}
262
263fn report_observe_event_log_error<T: Clone>(
264 on_error: &Option<Rc<ObserveEventLogErrorFn<T>>>,
265 error: StorageError,
266 ctx: ObserveEventLogErrorContext<T>,
267) {
268 if let Some(on_error) = on_error {
269 let _ = catch_unwind(AssertUnwindSafe(|| on_error(error, ctx)));
270 }
271}