harn_vm/stdlib/
session_change.rs1use std::path::Path;
13use std::sync::atomic::{AtomicU64, Ordering};
14use std::sync::{Arc, RwLock, RwLockWriteGuard};
15
16use harn_session_store::{SessionMeta, SharedSessionChangeObserver};
17
18use super::session_wal_watch;
19
20static OBSERVERS: RwLock<Vec<(u64, SharedSessionChangeObserver)>> = RwLock::new(Vec::new());
33static NEXT_SUBSCRIPTION: AtomicU64 = AtomicU64::new(1);
34type TitleFingerprint = (Option<String>, bool);
35type RememberedTitles = Vec<(String, TitleFingerprint)>;
36static TITLES: RwLock<RememberedTitles> = RwLock::new(Vec::new());
37
38fn observers() -> RwLockWriteGuard<'static, Vec<(u64, SharedSessionChangeObserver)>> {
39 OBSERVERS
40 .write()
41 .unwrap_or_else(|poisoned| poisoned.into_inner())
42}
43
44fn titles() -> RwLockWriteGuard<'static, RememberedTitles> {
45 TITLES
46 .write()
47 .unwrap_or_else(|poisoned| poisoned.into_inner())
48}
49
50#[must_use = "dropping the subscription immediately unregisters the observer"]
54pub struct SessionChangeSubscription {
55 id: u64,
56}
57
58impl Drop for SessionChangeSubscription {
59 fn drop(&mut self) {
60 observers().retain(|(id, _)| *id != self.id);
61 let live = subscriber_count() > 0;
62 if !live {
63 titles().clear();
64 }
65 session_wal_watch::sync_watchers(live);
66 }
67}
68
69pub fn subscribe(observer: SharedSessionChangeObserver) -> SessionChangeSubscription {
77 let id = NEXT_SUBSCRIPTION.fetch_add(1, Ordering::Relaxed);
78 observers().push((id, observer));
79 session_wal_watch::sync_watchers(true);
80 SessionChangeSubscription { id }
81}
82
83pub(crate) fn watch_store(path: &Path) {
85 session_wal_watch::register_store_path(path);
86 if subscriber_count() > 0 {
87 session_wal_watch::sync_watchers(true);
88 }
89}
90
91fn subscriber_count() -> usize {
92 OBSERVERS
93 .read()
94 .unwrap_or_else(|poisoned| poisoned.into_inner())
95 .len()
96}
97
98#[derive(Clone, Copy, Debug, PartialEq, Eq)]
100pub(super) enum TitleMemory {
101 New,
102 Unchanged,
103 Changed,
104}
105
106pub(super) fn remember_title(
107 session_id: &str,
108 title: Option<&str>,
109 title_pinned: bool,
110) -> TitleMemory {
111 let next = (title.map(str::to_string), title_pinned);
112 let mut titles = titles();
113 if let Some((_, previous)) = titles.iter_mut().find(|(id, _)| id == session_id) {
114 if *previous == next {
115 return TitleMemory::Unchanged;
116 }
117 *previous = next;
118 return TitleMemory::Changed;
119 }
120 titles.push((session_id.to_string(), next));
121 TitleMemory::New
122}
123
124pub(super) fn dispatch(meta: &SessionMeta) {
126 let observers: Vec<SharedSessionChangeObserver> = OBSERVERS
127 .read()
128 .unwrap_or_else(|poisoned| poisoned.into_inner())
129 .iter()
130 .map(|(_, observer)| Arc::clone(observer))
131 .collect();
132 for observer in observers {
135 observer.session_updated(meta);
136 }
137}
138
139struct SessionChangeFanout;
140
141impl harn_session_store::SessionChangeObserver for SessionChangeFanout {
142 fn session_updated(&self, meta: &SessionMeta) {
143 remember_title(&meta.id, meta.title.as_deref(), meta.title_pinned);
144 dispatch(meta);
145 }
146}
147
148pub(crate) fn current_observer() -> Option<SharedSessionChangeObserver> {
149 if subscriber_count() == 0 {
150 return None;
151 }
152 Some(Arc::new(SessionChangeFanout))
153}