Skip to main content

bamboo_server/schedule_app/
ledger_bridge.rs

1//! Ledger → schedule bridge.
2//!
3//! Keeps real [`ScheduleStore`] entries in step with a ledger record's
4//! `remind_at` / `recurrence` times, implementing the
5//! [`LedgerScheduleBridge`] port the `ledger` tool talks to. Managed
6//! schedules are tagged by name (`ledger:<record_id>:…`) and their ids live on
7//! the record's `schedule_ids`; the invariant is that a terminal record owns
8//! no schedules.
9//!
10//! The `ledger` tool sits in the base tool layer, which is assembled before
11//! the schedule store exists, so the wiring goes through
12//! [`LateBoundLedgerBridge`]: a handle installed at tool-build time and bound
13//! to the real bridge once the scheduler is up.
14
15use std::sync::Arc;
16
17use async_trait::async_trait;
18use chrono::Utc;
19use tokio::sync::RwLock;
20
21use bamboo_domain::ledger::LedgerRecord;
22use bamboo_domain::schedule::{ScheduleRunConfig, ScheduleTrigger};
23use bamboo_server_tools::LedgerScheduleBridge;
24
25use super::store::ScheduleStore;
26
27pub struct ScheduleLedgerBridge {
28    store: Arc<ScheduleStore>,
29}
30
31impl ScheduleLedgerBridge {
32    pub fn new(store: Arc<ScheduleStore>) -> Self {
33        Self { store }
34    }
35
36    fn reminder_run_config(record: &LedgerRecord) -> ScheduleRunConfig {
37        ScheduleRunConfig {
38            task_message: Some(format!(
39                "Reminder fired for ledger record `{id}`: {title}.\n\
40                 Load it with the `ledger` tool (action=get, id={id}) and check its \
41                 current status. If it is already done or cancelled, stop. Otherwise \
42                 do whatever preparation is genuinely useful, then alert the user \
43                 with a concise notify message that names the record and when it is due.",
44                id = record.id,
45                title = record.title,
46            )),
47            auto_execute: true,
48            ..ScheduleRunConfig::default()
49        }
50    }
51}
52
53#[async_trait]
54impl LedgerScheduleBridge for ScheduleLedgerBridge {
55    async fn sync_record_schedules(&self, record: &LedgerRecord) -> Result<Vec<String>, String> {
56        // Reconcile by replacement: drop the old managed set, create the new
57        // one. Reminder counts are tiny and this stays deterministic across
58        // partial edits (changed times, removed reminders, added recurrence).
59        self.release_schedules(&record.schedule_ids).await?;
60
61        let now = Utc::now();
62        let mut created = Vec::new();
63        let mut triggers: Vec<(String, ScheduleTrigger)> = Vec::new();
64        for (index, remind_at) in record.time.remind_at.iter().enumerate() {
65            // A past reminder can never fire; creating it would be rejected.
66            if *remind_at <= now {
67                continue;
68            }
69            triggers.push((
70                format!("ledger:{}:remind:{}", record.id, index),
71                ScheduleTrigger::Once { at: *remind_at },
72            ));
73        }
74        if let Some(recurrence) = &record.time.recurrence {
75            triggers.push((
76                format!("ledger:{}:recurrence", record.id),
77                recurrence.clone(),
78            ));
79        }
80
81        for (name, trigger) in triggers {
82            let entry = self
83                .store
84                .create_schedule(name, trigger, true, Self::reminder_run_config(record))
85                .await
86                .map_err(|error| format!("failed to create reminder schedule: {error}"))?;
87            created.push(entry.id);
88        }
89        Ok(created)
90    }
91
92    async fn release_schedules(&self, schedule_ids: &[String]) -> Result<(), String> {
93        for schedule_id in schedule_ids {
94            // Already-gone schedules (fired Once entries prune, manual deletes)
95            // are fine — the goal is absence.
96            self.store
97                .delete_schedule(schedule_id)
98                .await
99                .map_err(|error| format!("failed to delete schedule {schedule_id}: {error}"))?;
100        }
101        Ok(())
102    }
103}
104
105/// A bridge handle that can be installed before the scheduler exists and bound
106/// afterwards. Until bound, mutations that need scheduling degrade to a
107/// warning on the tool result (the record itself still persists).
108#[derive(Default)]
109pub struct LateBoundLedgerBridge {
110    inner: RwLock<Option<Arc<dyn LedgerScheduleBridge>>>,
111}
112
113impl LateBoundLedgerBridge {
114    pub async fn bind(&self, bridge: Arc<dyn LedgerScheduleBridge>) {
115        *self.inner.write().await = Some(bridge);
116    }
117}
118
119#[async_trait]
120impl LedgerScheduleBridge for LateBoundLedgerBridge {
121    async fn sync_record_schedules(&self, record: &LedgerRecord) -> Result<Vec<String>, String> {
122        match self.inner.read().await.as_ref() {
123            Some(bridge) => bridge.sync_record_schedules(record).await,
124            None => Err("the scheduler is not initialized yet".to_string()),
125        }
126    }
127
128    async fn release_schedules(&self, schedule_ids: &[String]) -> Result<(), String> {
129        match self.inner.read().await.as_ref() {
130            Some(bridge) => bridge.release_schedules(schedule_ids).await,
131            None => Err("the scheduler is not initialized yet".to_string()),
132        }
133    }
134}
135
136#[cfg(test)]
137mod tests {
138    use super::*;
139    use bamboo_domain::ledger::{RecordKind, RecordStatus};
140    use chrono::Duration;
141
142    async fn test_store() -> (Arc<ScheduleStore>, tempfile::TempDir) {
143        let dir = tempfile::tempdir().unwrap();
144        let store = Arc::new(ScheduleStore::new(dir.path().to_path_buf()).await.unwrap());
145        (store, dir)
146    }
147
148    #[tokio::test]
149    async fn sync_creates_once_schedules_for_future_reminders_and_skips_past() {
150        let (store, _dir) = test_store().await;
151        let bridge = ScheduleLedgerBridge::new(store.clone());
152
153        let mut record = LedgerRecord::new("rec_1", RecordKind::Reminder, "Take medication");
154        record.time.remind_at = vec![
155            Utc::now() - Duration::hours(1), // past: skipped
156            Utc::now() + Duration::hours(3),
157        ];
158        let ids = bridge.sync_record_schedules(&record).await.unwrap();
159        assert_eq!(ids.len(), 1);
160
161        let entry = store.get_schedule(&ids[0]).await.unwrap();
162        assert!(entry.name.starts_with("ledger:rec_1:remind:"));
163        assert!(matches!(entry.trigger, ScheduleTrigger::Once { .. }));
164        assert!(entry.run_config.auto_execute);
165        assert!(entry
166            .run_config
167            .task_message
168            .as_deref()
169            .unwrap()
170            .contains("rec_1"));
171    }
172
173    #[tokio::test]
174    async fn sync_replaces_previous_managed_schedules() {
175        let (store, _dir) = test_store().await;
176        let bridge = ScheduleLedgerBridge::new(store.clone());
177
178        let mut record = LedgerRecord::new("rec_1", RecordKind::Reminder, "Standup");
179        record.time.remind_at = vec![Utc::now() + Duration::hours(1)];
180        let first = bridge.sync_record_schedules(&record).await.unwrap();
181        record.schedule_ids = first.clone();
182
183        // The reminder moves; the old schedule must not survive.
184        record.time.remind_at = vec![Utc::now() + Duration::hours(2)];
185        let second = bridge.sync_record_schedules(&record).await.unwrap();
186        assert_eq!(second.len(), 1);
187        assert!(store.get_schedule(&first[0]).await.is_none());
188        assert!(store.get_schedule(&second[0]).await.is_some());
189    }
190
191    #[tokio::test]
192    async fn recurrence_maps_to_a_recurring_schedule_and_release_deletes() {
193        let (store, _dir) = test_store().await;
194        let bridge = ScheduleLedgerBridge::new(store.clone());
195
196        let mut record = LedgerRecord::new("rec_habit", RecordKind::Habit, "Daily review");
197        record.time.recurrence = Some(ScheduleTrigger::Daily {
198            hour: 9,
199            minute: 0,
200            second: 0,
201        });
202        let ids = bridge.sync_record_schedules(&record).await.unwrap();
203        assert_eq!(ids.len(), 1);
204        assert!(matches!(
205            store.get_schedule(&ids[0]).await.unwrap().trigger,
206            ScheduleTrigger::Daily { .. }
207        ));
208
209        record.transition_to(RecordStatus::Done, None);
210        bridge.release_schedules(&ids).await.unwrap();
211        assert!(store.get_schedule(&ids[0]).await.is_none());
212        // Releasing again is a no-op, not an error.
213        bridge.release_schedules(&ids).await.unwrap();
214    }
215
216    #[tokio::test]
217    async fn late_bound_bridge_errors_until_bound_then_delegates() {
218        let (store, _dir) = test_store().await;
219        let handle = LateBoundLedgerBridge::default();
220
221        let mut record = LedgerRecord::new("rec_1", RecordKind::Reminder, "Ping");
222        record.time.remind_at = vec![Utc::now() + Duration::hours(1)];
223        assert!(handle.sync_record_schedules(&record).await.is_err());
224
225        handle
226            .bind(Arc::new(ScheduleLedgerBridge::new(store)))
227            .await;
228        assert_eq!(
229            handle.sync_record_schedules(&record).await.unwrap().len(),
230            1
231        );
232    }
233}