Skip to main content

myko/client/
view_map.rs

1use std::sync::Arc;
2
3use hyphae::{Cell, CellImmutable, CellMap, CellMutable, Mutable as _, Watchable as _};
4use serde::de::DeserializeOwned;
5use serde_json::Value;
6use tracing::{debug, error, trace};
7
8use super::{
9    ConnectionStatus, MykoClient,
10    map_response::{MapSequence, decode_map_upserts},
11    query_map::apply_incremental_map_update,
12};
13use crate::{
14    common::with_id::WithId,
15    core::{
16        item::Eventable,
17        view::{ViewParams, ViewRequest},
18    },
19    wire::{message::MykoMessage, wrap_view},
20};
21
22/// Fine-grained view data together with explicit initial-response readiness.
23#[derive(Clone)]
24pub struct ViewMapWatch<T: hyphae::CellValue> {
25    map: CellMap<Arc<str>, Arc<T>, CellImmutable>,
26    ready: Cell<bool, CellImmutable>,
27}
28
29impl<T: hyphae::CellValue> ViewMapWatch<T> {
30    #[must_use]
31    pub const fn map(&self) -> &CellMap<Arc<str>, Arc<T>, CellImmutable> {
32        &self.map
33    }
34
35    #[must_use]
36    pub const fn ready(&self) -> &Cell<bool, CellImmutable> {
37        &self.ready
38    }
39
40    #[must_use]
41    pub fn into_map(self) -> CellMap<Arc<str>, Arc<T>, CellImmutable> {
42        self.map
43    }
44}
45
46impl MykoClient {
47    /// Watch a view with stable, independently reactive item cells.
48    /// Identical view parameters share one decoded map and wire subscription.
49    pub fn watch_view_map<V>(
50        &self,
51        view: impl Into<ViewRequest<V>>,
52    ) -> CellMap<Arc<str>, Arc<V::Item>, CellImmutable>
53    where
54        V: ViewParams + Clone,
55        V::Item: Eventable + WithId + DeserializeOwned + Clone + std::fmt::Debug + 'static,
56    {
57        self.watch_view_map_state(view).into_map()
58    }
59
60    /// Watch a fine-grained view map and retain an initial-response signal.
61    #[allow(clippy::too_many_lines)]
62    pub fn watch_view_map_state<V>(&self, view: impl Into<ViewRequest<V>>) -> ViewMapWatch<V::Item>
63    where
64        V: ViewParams + Clone,
65        V::Item: Eventable + WithId + DeserializeOwned + Clone + std::fmt::Debug + 'static,
66    {
67        let supplied: ViewRequest<V> = view.into();
68        let view_id = supplied.view.view_id();
69        let cache_key = format!(
70            "view-map:{view_id}:{}:{:016x}",
71            std::any::type_name::<V::Item>(),
72            supplied.view.cache_key_hash()
73        );
74        let _cache_gate = self
75            .inner
76            .map_watch_cache_gate
77            .lock()
78            .unwrap_or_else(std::sync::PoisonError::into_inner);
79        if let Some((map, ready)) = self.cached_map_watch(&cache_key) {
80            debug!("watch_view_map_state: cache hit for {cache_key}");
81            return ViewMapWatch { map, ready };
82        }
83        self.inner.map_watch_cache.remove(&cache_key);
84
85        let view = ViewRequest::with_tx(supplied.view, super::next_subscription_tx());
86        let tx = view.tx.clone();
87        let map: CellMap<Arc<str>, Arc<V::Item>> =
88            CellMap::new().with_name(format!("view_map:{view_id}"));
89        let map_weak = map.downgrade();
90        let ready =
91            Cell::<bool, CellMutable>::new(false).with_name(format!("view_map_ready:{view_id}"));
92        let ready_weak = ready.downgrade();
93        let ready_read = ready.clone().lock();
94
95        let Ok(wrapped) = wrap_view(tx.clone(), &view.view) else {
96            error!("Could not serialize view map request for {view_id}");
97            return ViewMapWatch {
98                map: map.lock(),
99                ready: ready_read,
100            };
101        };
102        let Ok(frame) = self.encode_message(&MykoMessage::View(wrapped)) else {
103            error!("Could not encode view map request for {view_id}");
104            return ViewMapWatch {
105                map: map.lock(),
106                ready: ready_read,
107            };
108        };
109
110        let tx_for_handler = tx.clone();
111        let view_id_for_handler = view_id.clone();
112        let sequences = Arc::new(MapSequence::new());
113        let sequences_for_handler = Arc::clone(&sequences);
114        let handler: super::QueryHandler = Box::new(move |response_value: Value| {
115            let Some(map_writer) = map_weak.upgrade() else {
116                return;
117            };
118            let response =
119                match serde_json::from_value::<crate::wire::ClientQueryResponse>(response_value) {
120                    Ok(response) => response,
121                    Err(error) => {
122                        error!(
123                            "Rejected view '{}' malformed response: {}",
124                            view_id_for_handler, error
125                        );
126                        return;
127                    }
128                };
129            if response.tx != tx_for_handler {
130                return;
131            }
132
133            let upserts = match decode_map_upserts::<V::Item, _>(response.upserts, WithId::id) {
134                Ok(upserts) => upserts,
135                Err(error) => {
136                    error!(
137                        "Rejected view '{}' response: invalid {} upsert: {}",
138                        view_id_for_handler,
139                        std::any::type_name::<V::Item>(),
140                        error
141                    );
142                    return;
143                }
144            };
145            if !sequences_for_handler.accept(response.sequence) {
146                error!(
147                    "Rejected view '{}' out-of-order sequence {}",
148                    view_id_for_handler, response.sequence
149                );
150                return;
151            }
152            let is_initial_response = response.sequence == 0;
153            if is_initial_response {
154                trace!("Sequence reset: replacing {} view map", view_id_for_handler);
155                map_writer.replace_all(upserts);
156            } else {
157                apply_incremental_map_update(&map_writer, response.deletes, upserts);
158            }
159            if is_initial_response && let Some(ready_writer) = ready_weak.upgrade() {
160                ready_writer.set(true);
161            }
162        });
163        if !self.try_register_query_handler(tx.clone(), handler) {
164            error!("Refusing duplicate view map transaction {tx}");
165            return ViewMapWatch {
166                map: map.lock(),
167                ready: ready_read,
168            };
169        }
170
171        let socket = self.inner.socket.clone();
172        let ready_for_status = ready.downgrade();
173        let sequences_for_status = sequences;
174        let status_cell = self.connection_status();
175        let send_view_id = view_id;
176        let status_guard = status_cell.subscribe(move |signal| {
177            if let hyphae::Signal::Value(status) = signal {
178                if let ConnectionStatus::Connected(_) = &**status {
179                    match socket.send(frame.clone()) {
180                        Ok(()) => debug!("Watching view map {send_view_id}"),
181                        Err(error) => error!("Could not send view: {error:?}"),
182                    }
183                } else {
184                    sequences_for_status.reset_epoch();
185                    if let Some(ready_writer) = ready_for_status.upgrade() {
186                        ready_writer.set(false);
187                    }
188                    debug!("View map {send_view_id} disconnected");
189                }
190            }
191        });
192        map.own(status_guard);
193        map.own(super::view_cancel_guard(tx.clone(), self.inner.clone()));
194        map.own(super::retain_cell_guard(ready_read.clone()));
195        map.own(super::map_watch_cache_guard(
196            cache_key.clone(),
197            tx.clone(),
198            self.inner.clone(),
199        ));
200        let watch = ViewMapWatch {
201            map: map.lock(),
202            ready: ready_read,
203        };
204        self.cache_map_watch(cache_key, tx, &watch.map, &watch.ready);
205        watch
206    }
207}