1use std::{
21 collections::{HashMap, HashSet},
22 sync::Arc,
23 time::{Duration, Instant},
24};
25
26use chrono::Utc;
27use hyphae::Gettable;
28use marshal_entities::{AutoSource, Message, Room, RoomKind, RoomMember, Session};
29use myko::{core::item::Eventable, server::CellServerCtx, utils::downcast_item};
30
31pub const STALE_AFTER: Duration = Duration::from_secs(60);
41
42pub const MESSAGE_TTL: Duration = Duration::from_secs(14 * 24 * 60 * 60);
48
49const MESSAGE_SWEEP_EVERY: u64 = 100;
53
54pub const HOOK_BACKSTOP: Duration = Duration::from_secs(60 * 60);
60
61pub const TICK_INTERVAL: Duration = Duration::from_secs(3);
65
66pub async fn run_sweeper(ctx: CellServerCtx) {
68 let mut disconnected_since: HashMap<Arc<str>, Instant> = HashMap::new();
69 let mut interval = tokio::time::interval(TICK_INTERVAL);
70 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
71 let mut tick: u64 = 0;
72
73 loop {
74 interval.tick().await;
75 sweep_once(&ctx, &mut disconnected_since);
76 sweep_rooms(&ctx);
77 if tick.is_multiple_of(MESSAGE_SWEEP_EVERY) {
81 sweep_messages(&ctx);
82 }
83 tick = tick.wrapping_add(1);
84 }
85}
86
87fn sweep_messages(ctx: &CellServerCtx) {
92 let Some(store) = ctx.registry.get(Message::ENTITY_NAME_STATIC) else {
93 return;
94 };
95 let cutoff = Utc::now().timestamp_millis() - MESSAGE_TTL.as_millis() as i64;
96 let mut to_delete: Vec<Arc<str>> = Vec::new();
97 for (id, item) in store.entries().get() {
98 if let Some(m) = downcast_item::<Message>(&item)
99 && m.sent_at < cutoff
100 {
101 to_delete.push(id);
102 }
103 }
104 if to_delete.is_empty() {
105 return;
106 }
107 log::info!(
108 "[cleanup] pruning {} message(s) older than {} days",
109 to_delete.len(),
110 MESSAGE_TTL.as_secs() / 86_400,
111 );
112 for id in to_delete {
113 if let Err(e) = ctx.del_by_id(Message::ENTITY_NAME_STATIC, &id) {
114 log::warn!("[cleanup] prune message {} failed: {}", id, e);
115 }
116 }
117}
118
119fn sweep_rooms(ctx: &CellServerCtx) {
127 let Some(room_store) = ctx.registry.get(Room::ENTITY_NAME_STATIC) else {
128 return;
129 };
130 let Some(member_store) = ctx.registry.get(RoomMember::ENTITY_NAME_STATIC) else {
131 return;
132 };
133
134 let mut member_counts: HashMap<Arc<str>, usize> = HashMap::new();
136 for (_id, item) in member_store.entries().get() {
137 if let Some(m) = downcast_item::<RoomMember>(&item) {
138 *member_counts.entry(m.room_id.0.clone()).or_default() += 1;
139 }
140 }
141
142 let mut to_delete: Vec<Arc<str>> = Vec::new();
143 for (id, item) in room_store.entries().get() {
144 let Some(room) = downcast_item::<Room>(&item) else {
145 continue;
146 };
147 let RoomKind::Auto { source } = &room.kind else {
149 continue;
150 };
151 if matches!(source, AutoSource::Everyone) {
153 continue;
154 }
155 let is_host = matches!(source, AutoSource::Host { .. });
156 let empty = member_counts.get(&room.id.0).copied().unwrap_or(0) == 0;
157 if is_host || empty {
158 to_delete.push(id);
159 }
160 }
161
162 for id in to_delete {
163 log::info!("[cleanup] DELing stale auto-room {}", id);
164 if let Err(e) = ctx.del_by_id(Room::ENTITY_NAME_STATIC, &id) {
165 log::warn!("[cleanup] del room {} failed: {}", id, e);
166 }
167 }
168}
169
170fn sweep_once(ctx: &CellServerCtx, disconnected_since: &mut HashMap<Arc<str>, Instant>) {
171 let Some(session_store) = ctx.registry.get(Session::ENTITY_NAME_STATIC) else {
172 return;
173 };
174 let Some(client_store) = ctx.registry.get("Client") else {
175 return;
179 };
180
181 let live_client_ids: HashSet<Arc<str>> = client_store
182 .entries()
183 .get()
184 .into_iter()
185 .map(|(id, _)| id)
186 .collect();
187
188 let now = Instant::now();
189 let now_ms = Utc::now().timestamp_millis();
190 let backstop_ms = HOOK_BACKSTOP.as_millis() as i64;
191 let mut to_delete: Vec<Arc<str>> = Vec::new();
192 let mut still_disconnected: HashSet<Arc<str>> = HashSet::new();
193
194 for (id, item) in session_store.entries().get() {
195 let Some(session) = downcast_item::<Session>(&item) else {
196 continue;
197 };
198
199 match session.client_id.as_ref() {
200 None => {
204 let last = session.last_activity_at.unwrap_or(session.connected_at);
205 if now_ms.saturating_sub(last) >= backstop_ms {
206 to_delete.push(id);
207 }
208 }
209 Some(cid) if live_client_ids.contains(&cid.0) => {
211 disconnected_since.remove(&id);
212 }
213 Some(_) => {
215 still_disconnected.insert(id.clone());
216 let first_seen = *disconnected_since.entry(id.clone()).or_insert(now);
217 if now.duration_since(first_seen) >= STALE_AFTER {
218 to_delete.push(id);
219 }
220 }
221 }
222 }
223
224 disconnected_since.retain(|id, _| still_disconnected.contains(id));
227
228 for id in to_delete {
229 log::info!("[cleanup] DELing abandoned session {}", id);
230 if let Err(e) = ctx.del_by_id(Session::ENTITY_NAME_STATIC, &id) {
231 log::warn!("[cleanup] del session {} failed: {}", id, e);
232 continue;
233 }
234 disconnected_since.remove(&id);
235 }
236}
237
238#[cfg(test)]
239mod tests {
240 use super::*;
241 use marshal_entities::{Message, MessageId, RoomId, RoomMemberId, SessionId};
242 use myko::{
243 server::Persister,
244 wire::{MEvent, MEventType},
245 };
246 use myko_server::{BlackholePersister, CellServer};
247 use std::collections::HashSet;
248 use uuid::Uuid;
249
250 fn setup() -> CellServerCtx {
251 marshal_entities::link();
252 crate::link();
253 let blackhole: Arc<dyn Persister> = Arc::new(BlackholePersister);
254 let server = CellServer::builder()
255 .with_default_persister(blackhole)
256 .build();
257 let ctx = server.ctx();
258 Box::leak(Box::new(server));
259 ctx
260 }
261
262 fn set_room(ctx: &CellServerCtx, id: &str, kind: RoomKind) {
263 let room = Room {
264 id: RoomId(Arc::from(id)),
265 name: id.to_string(),
266 description: None,
267 kind,
268 created_at: 0,
269 };
270 let ev = MEvent::from_item(&room, MEventType::SET, &Uuid::new_v4().to_string());
271 ctx.apply_event_batch(vec![ev]).expect("apply Room SET");
272 }
273
274 fn set_member(ctx: &CellServerCtx, room_id: &str, session_id: &str) {
275 let member = RoomMember {
276 id: RoomMemberId(Arc::from(RoomMember::make_id(room_id, session_id).as_str())),
277 room_id: RoomId(Arc::from(room_id)),
278 session_id: SessionId(Arc::from(session_id)),
279 joined_at: 0,
280 };
281 let ev = MEvent::from_item(&member, MEventType::SET, &Uuid::new_v4().to_string());
282 ctx.apply_event_batch(vec![ev])
283 .expect("apply RoomMember SET");
284 }
285
286 fn room_ids(ctx: &CellServerCtx) -> HashSet<String> {
287 ctx.registry
288 .get(Room::ENTITY_NAME_STATIC)
289 .map(|s| {
290 s.entries()
291 .get()
292 .into_iter()
293 .map(|(id, _)| id.to_string())
294 .collect()
295 })
296 .unwrap_or_default()
297 }
298
299 #[test]
300 fn sweep_reaps_host_and_empty_auto_rooms_and_spares_the_rest() {
301 let ctx = setup();
302
303 set_room(
305 &ctx,
306 "everyone",
307 RoomKind::Auto {
308 source: AutoSource::Everyone,
309 },
310 );
311 set_member(&ctx, "everyone", "sess-a");
312 set_room(
314 &ctx,
315 "host:node1",
316 RoomKind::Auto {
317 source: AutoSource::Host {
318 name: "node1".into(),
319 },
320 },
321 );
322 set_member(&ctx, "host:node1", "sess-a");
323 set_room(
325 &ctx,
326 "project:live",
327 RoomKind::Auto {
328 source: AutoSource::Project {
329 basename: "live".into(),
330 },
331 },
332 );
333 set_member(&ctx, "project:live", "sess-a");
334 set_room(
336 &ctx,
337 "project:stale",
338 RoomKind::Auto {
339 source: AutoSource::Project {
340 basename: "stale".into(),
341 },
342 },
343 );
344 set_room(&ctx, "design-sync", RoomKind::Adhoc);
346
347 sweep_rooms(&ctx);
348
349 let ids = room_ids(&ctx);
350 assert!(ids.contains("everyone"), "global room must survive");
351 assert!(
352 ids.contains("project:live"),
353 "populated project room must survive"
354 );
355 assert!(ids.contains("design-sync"), "empty adhoc room must survive");
356 assert!(!ids.contains("host:node1"), "host room must be reaped");
357 assert!(
358 !ids.contains("project:stale"),
359 "empty auto-room must be reaped"
360 );
361 }
362
363 fn set_message(ctx: &CellServerCtx, id: &str, sent_at: i64) {
364 let msg = Message {
365 id: MessageId(Arc::from(id)),
366 from_session_id: SessionId(Arc::from("sender")),
367 to_session_id: Some(SessionId(Arc::from("recipient"))),
368 to_room_id: None,
369 to_operator: None,
370 body: "hi".into(),
371 sent_at,
372 };
373 let ev = MEvent::from_item(&msg, MEventType::SET, &Uuid::new_v4().to_string());
374 ctx.apply_event_batch(vec![ev]).expect("apply Message SET");
375 }
376
377 fn message_ids(ctx: &CellServerCtx) -> HashSet<String> {
378 ctx.registry
379 .get(Message::ENTITY_NAME_STATIC)
380 .map(|s| {
381 s.entries()
382 .get()
383 .into_iter()
384 .map(|(id, _)| id.to_string())
385 .collect()
386 })
387 .unwrap_or_default()
388 }
389
390 #[test]
391 fn sweep_messages_prunes_old_and_keeps_recent() {
392 let ctx = setup();
393 let now = chrono::Utc::now().timestamp_millis();
394 set_message(&ctx, "old", now - MESSAGE_TTL.as_millis() as i64 - 1);
396 set_message(&ctx, "recent", now - 1_000);
398
399 sweep_messages(&ctx);
400
401 let ids = message_ids(&ctx);
402 assert!(
403 !ids.contains("old"),
404 "message older than the TTL must be pruned"
405 );
406 assert!(ids.contains("recent"), "recent message must survive");
407 }
408}