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#[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 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 #[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}