bamboo_server/schedule_app/
ledger_bridge.rs1use 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 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 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 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#[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), 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 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 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}