Skip to main content

notedthat_write/
sinks.rs

1//! Where a committed write is announced: the indexing queue and, when one is
2//! configured, the object change event log.
3//!
4//! Bundled so every write path takes one argument and every surface names its
5//! [`EventSource`] once, where it builds the bundle, rather than at each call.
6
7use notedthat_core::{EventPublisher, EventSource};
8use notedthat_indexer::{IndexEvent, IndexHealth};
9use tokio::sync::mpsc::Sender;
10
11/// The two places a durable write is reported to, and on whose behalf.
12#[derive(Clone, Copy)]
13pub struct WriteSinks<'a> {
14    /// The in-process indexing queue (D38).
15    pub indexer_tx: &'a Sender<IndexEvent>,
16    /// The event log, when `NOTEDTHAT_EVENTS_BACKEND` selects one.
17    pub events: Option<&'a dyn EventPublisher>,
18    /// Where what happened at the queue — enqueued, refused, or a queue with
19    /// no worker behind it — is recorded for the health view (#97). `None`
20    /// only in tests that have no health view to keep.
21    pub index_health: Option<&'a IndexHealth>,
22    /// The surface making the write, stamped on every event it publishes.
23    pub source: EventSource,
24}
25
26impl<'a> WriteSinks<'a> {
27    /// Both sinks, recording on `index_health`.
28    #[must_use]
29    pub fn new(
30        indexer_tx: &'a Sender<IndexEvent>,
31        events: Option<&'a dyn EventPublisher>,
32        index_health: &'a IndexHealth,
33        source: EventSource,
34    ) -> Self {
35        Self {
36            indexer_tx,
37            events,
38            index_health: Some(index_health),
39            source,
40        }
41    }
42
43    /// The indexing queue alone, as every deployment without an events backend
44    /// runs, attributed to the HTTP API, with nothing keeping a health record.
45    #[must_use]
46    pub fn indexer_only(indexer_tx: &'a Sender<IndexEvent>) -> Self {
47        Self {
48            indexer_tx,
49            events: None,
50            index_health: None,
51            source: EventSource::Http,
52        }
53    }
54}