Skip to main content

graphrefly/adapters/
observe_storage.rs

1//! Graph.observe() -> passive storage adapter (D57/D74/D125).
2//!
3//! Storage frame codecs and append-log paging stay in [`crate::storage`]. This module owns the
4//! graph-bound observe subscription that writes those passive frames.
5
6use 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)]
18/// `ObserveEventLogErrorPhase` variants.
19pub enum ObserveEventLogErrorPhase {
20    /// `Map` variant.
21    Map,
22    /// `Write` variant.
23    Write,
24    /// `Flush` variant.
25    Flush,
26    /// `Rollback` variant.
27    Rollback,
28    /// `Dispose` variant.
29    Dispose,
30}
31
32#[derive(Clone)]
33/// `ObserveEventLogErrorContext` data container.
34pub struct ObserveEventLogErrorContext<T> {
35    /// `phase` field for phase.
36    pub phase: ObserveEventLogErrorPhase,
37    /// `event` field for event.
38    pub event: Option<ObserveEvent>,
39    /// `value` field for value.
40    pub value: Option<ObserveEventFrame<T>>,
41}
42
43/// `ObserveEventLogMap` type alias.
44pub type ObserveEventLogMap<T> = dyn Fn(&ObserveEvent) -> Option<T>;
45/// `ObserveEventLogErrorFn` type alias.
46pub type ObserveEventLogErrorFn<T> = dyn Fn(StorageError, ObserveEventLogErrorContext<T>);
47
48#[derive(Clone)]
49/// `AttachObserveEventLogOptions` data container.
50pub struct AttachObserveEventLogOptions<T: Clone> {
51    /// `path` field for path.
52    pub path: Option<String>,
53    /// `stream` field for stream.
54    pub stream: Option<String>,
55    /// `map` field for map.
56    pub map: Rc<ObserveEventLogMap<T>>,
57    /// `on_error` field for on error.
58    pub on_error: Option<Rc<ObserveEventLogErrorFn<T>>>,
59}
60
61impl<T> AttachObserveEventLogOptions<T>
62where
63    T: Clone + From<ObserveEvent> + 'static,
64{
65    /// Creates or computes `new`.
66    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    /// Creates or computes `from_map`.
87    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    /// Updates or reads `with_path`.
97    pub fn with_path(mut self, path: impl Into<String>) -> Self {
98        self.path = Some(path.into());
99        self
100    }
101
102    /// Updates or reads `with_stream`.
103    pub fn with_stream(mut self, stream: impl Into<String>) -> Self {
104        self.stream = Some(stream.into());
105        self
106    }
107
108    /// Updates or reads `with_map`.
109    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    /// Updates or reads `with_on_error`.
115    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
124/// `ObserveEventLogHandle` data container.
125pub struct ObserveEventLogHandle {
126    observer: Option<GraphObserver>,
127    flush: Rc<dyn Fn() -> StorageResult<()>>,
128    rollback: Rc<dyn Fn() -> StorageResult<()>>,
129}
130
131impl ObserveEventLogHandle {
132    /// Updates or reads `flush`.
133    pub fn flush(&self) -> StorageResult<()> {
134        (self.flush)()
135    }
136
137    /// Updates or reads `rollback`.
138    pub fn rollback(&self) -> StorageResult<()> {
139        (self.rollback)()
140    }
141
142    /// Updates or reads `dispose`.
143    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
155/// Creates or computes `attach_observe_event_log`.
156pub 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}