1use std::collections::{BTreeSet, HashSet};
2use std::sync::Arc;
3
4use anyhow::{Result, anyhow};
5use locus_core_rs::domain::contracts::{NodeStore, SemanticIndexStore};
6use locus_core_rs::domain::models::{
7 NodeDeleteRecord, NodeDeleteRequest, NodeDeleteResult, NodeDeleteStatus, NodeQuery,
8 SessionPurgeRequest, SttpNode,
9};
10use locus_core_rs::storage::derive_tenant_id_from_session;
11
12use crate::application::memory_filters::{
13 build_session_filter, node_matches_common_filters, resolve_indexed_sync_keys,
14};
15use crate::application::memory_graph::graph_node_id;
16use crate::domain::evict::{
17 InboundReferencesPreview, MemoryEvictMode, MemoryEvictRecord, MemoryEvictRequest,
18 MemoryEvictResult,
19};
20use crate::domain::memory::{clamp_batch_size, clamp_nodes};
21
22const TENANT_SCAN_LIMIT: usize = 5000;
23
24#[derive(Debug, Clone)]
25struct EvictCandidate {
26 store_node_id: String,
27 sync_key: String,
28 graph_id: String,
29}
30
31pub struct MemoryEvictService {
32 store: Arc<dyn NodeStore>,
33 semantic_index: Option<Arc<dyn SemanticIndexStore>>,
34}
35
36impl MemoryEvictService {
37 pub fn new(store: Arc<dyn NodeStore>) -> Self {
38 Self {
39 store,
40 semantic_index: None,
41 }
42 }
43
44 pub fn with_semantic_index(
45 mut self,
46 semantic_index: Arc<dyn SemanticIndexStore>,
47 ) -> Self {
48 self.semantic_index = Some(semantic_index);
49 self
50 }
51
52 pub async fn execute(&self, request: &MemoryEvictRequest) -> Result<MemoryEvictResult> {
53 if matches!(request.mode, MemoryEvictMode::PurgeSession) {
54 return self.execute_purge(request).await;
55 }
56
57 let session_id = single_session_id(&request.scope)?;
58 let tenant_id = resolve_tenant_id(&request.scope, &session_id);
59 let max_nodes = clamp_nodes(if request.max_nodes == 0 {
60 5000
61 } else {
62 request.max_nodes
63 });
64
65 let candidates = self
66 .resolve_candidates(request, &session_id, &tenant_id, max_nodes)
67 .await?;
68
69 if candidates.is_empty() {
70 return Ok(MemoryEvictResult::default());
71 }
72
73 let session_nodes = self
74 .store
75 .query_nodes_async(NodeQuery {
76 limit: TENANT_SCAN_LIMIT,
77 session_id: Some(session_id.clone()),
78 from_utc: request.scope.from_utc,
79 to_utc: request.scope.to_utc,
80 tiers: request.scope.tiers.clone(),
81 })
82 .await?;
83
84 let candidate_sync_keys = candidates
85 .iter()
86 .map(|candidate| candidate.sync_key.clone())
87 .collect::<HashSet<_>>();
88
89 let mut records = Vec::new();
90 let mut to_delete_sync_keys = Vec::new();
91 let mut to_delete_node_ids = Vec::new();
92
93 for candidate in candidates {
94 if !request.force {
95 let inbound = collect_inbound_refs(&candidate, &session_nodes, &candidate_sync_keys);
96 if !inbound.child_parent_links.is_empty()
97 || !inbound.incoming_semantic_refs.is_empty()
98 {
99 records.push(MemoryEvictRecord {
100 node_id: candidate.store_node_id.clone(),
101 sync_key: candidate.sync_key.clone(),
102 status: "blocked".to_string(),
103 reason: Some("inbound references exist".to_string()),
104 inbound_references: Some(InboundReferencesPreview {
105 child_parent_links: inbound.child_parent_links,
106 incoming_semantic_refs: inbound.incoming_semantic_refs,
107 }),
108 });
109 continue;
110 }
111 }
112
113 if request.dry_run {
114 records.push(MemoryEvictRecord {
115 node_id: candidate.store_node_id.clone(),
116 sync_key: candidate.sync_key.clone(),
117 status: "would_delete".to_string(),
118 reason: None,
119 inbound_references: None,
120 });
121 } else {
122 if !candidate.store_node_id.is_empty() {
123 to_delete_node_ids.push(candidate.store_node_id.clone());
124 }
125 if !candidate.sync_key.is_empty() {
126 to_delete_sync_keys.push(candidate.sync_key.clone());
127 }
128 }
129 }
130
131 let mut core_result = NodeDeleteResult::default();
132 if !request.dry_run && (!to_delete_sync_keys.is_empty() || !to_delete_node_ids.is_empty()) {
133 let batch_size = clamp_batch_size(500);
134 for chunk in to_delete_sync_keys.chunks(batch_size) {
135 let chunk_result = self
136 .store
137 .delete_nodes_async(NodeDeleteRequest {
138 tenant_id: tenant_id.clone(),
139 session_id: session_id.clone(),
140 sync_keys: chunk.to_vec(),
141 node_ids: Vec::new(),
142 dry_run: false,
143 })
144 .await?;
145 merge_delete_result(&mut core_result, chunk_result);
146 }
147
148 for chunk in to_delete_node_ids.chunks(batch_size) {
149 let chunk_result = self
150 .store
151 .delete_nodes_async(NodeDeleteRequest {
152 tenant_id: tenant_id.clone(),
153 session_id: session_id.clone(),
154 sync_keys: Vec::new(),
155 node_ids: chunk.to_vec(),
156 dry_run: false,
157 })
158 .await?;
159 merge_delete_result(&mut core_result, chunk_result);
160 }
161
162 for record in &core_result.records {
163 if record.status == NodeDeleteStatus::Deleted && !record.sync_key.is_empty() {
164 self.delete_tag_index_rows(&tenant_id, &record.sync_key).await?;
165 }
166 }
167
168 for record in core_result.records.iter().filter(|record| {
169 record.status == NodeDeleteStatus::Deleted
170 || record.status == NodeDeleteStatus::NotFound
171 }) {
172 records.push(map_core_record(record));
173 }
174 }
175
176 let blocked = records.iter().filter(|record| record.status == "blocked").count();
177 let would_delete = records
178 .iter()
179 .filter(|record| record.status == "would_delete")
180 .map(|record| record.sync_key.clone())
181 .collect::<Vec<_>>();
182 let deleted = if request.dry_run {
183 would_delete.len()
184 } else {
185 core_result.deleted
186 };
187 let not_found = records
188 .iter()
189 .filter(|record| record.status == "not_found")
190 .count()
191 + if request.dry_run {
192 0
193 } else {
194 core_result.not_found
195 };
196
197 Ok(MemoryEvictResult {
198 dry_run: request.dry_run,
199 deleted,
200 blocked,
201 not_found,
202 skipped: core_result.skipped,
203 would_delete,
204 calibrations_deleted: 0,
205 checkpoints_deleted: 0,
206 records,
207 })
208 }
209
210 async fn execute_purge(&self, request: &MemoryEvictRequest) -> Result<MemoryEvictResult> {
211 let session_id = single_session_id(&request.scope)?;
212 let tenant_id = resolve_tenant_id(&request.scope, &session_id);
213
214 let core_result = self
215 .store
216 .purge_session_async(SessionPurgeRequest {
217 tenant_id: tenant_id.clone(),
218 session_id: session_id.clone(),
219 tiers: request.scope.tiers.clone(),
220 dry_run: request.dry_run,
221 include_calibration: request.include_calibration,
222 include_checkpoints: request.include_checkpoints,
223 })
224 .await?;
225
226 let records = core_result
227 .records
228 .iter()
229 .map(map_core_record)
230 .collect::<Vec<_>>();
231 let would_delete = records
232 .iter()
233 .filter(|record| record.status == "would_delete" || record.status == "skipped")
234 .map(|record| record.sync_key.clone())
235 .filter(|sync_key| !sync_key.is_empty())
236 .collect::<Vec<_>>();
237
238 Ok(MemoryEvictResult {
239 dry_run: core_result.dry_run,
240 deleted: core_result.deleted,
241 blocked: core_result.blocked,
242 not_found: core_result.not_found,
243 skipped: core_result.skipped,
244 would_delete,
245 calibrations_deleted: core_result.calibrations_deleted,
246 checkpoints_deleted: core_result.checkpoints_deleted,
247 records,
248 })
249 }
250
251 async fn resolve_candidates(
252 &self,
253 request: &MemoryEvictRequest,
254 session_id: &str,
255 tenant_id: &str,
256 max_nodes: usize,
257 ) -> Result<Vec<EvictCandidate>> {
258 match request.mode {
259 MemoryEvictMode::BySyncKeys => {
260 let keys = request
261 .sync_keys
262 .as_ref()
263 .ok_or_else(|| anyhow!("sync_keys are required for by_sync_keys mode"))?;
264 if keys.is_empty() {
265 return Err(anyhow!("at least one sync key is required"));
266 }
267
268 let session_nodes = self.load_session_nodes(request, session_id).await?;
269 let mut candidates = Vec::new();
270 for sync_key in keys.iter().map(|value| value.trim()).filter(|value| !value.is_empty()) {
271 if let Some(node) = session_nodes.iter().find(|node| node.sync_key == sync_key) {
272 candidates.push(build_candidate(node, String::new()));
273 } else {
274 candidates.push(EvictCandidate {
275 store_node_id: String::new(),
276 sync_key: sync_key.to_string(),
277 graph_id: String::new(),
278 });
279 }
280 }
281 Ok(candidates)
282 }
283 MemoryEvictMode::ByNodeIds => {
284 let node_ids = request
285 .node_ids
286 .as_ref()
287 .ok_or_else(|| anyhow!("node_ids are required for by_node_ids mode"))?;
288 if node_ids.is_empty() {
289 return Err(anyhow!("at least one node id is required"));
290 }
291
292 Ok(node_ids
293 .iter()
294 .map(|node_id| EvictCandidate {
295 store_node_id: node_id.trim().to_string(),
296 sync_key: String::new(),
297 graph_id: String::new(),
298 })
299 .filter(|candidate| !candidate.store_node_id.is_empty())
300 .collect())
301 }
302 MemoryEvictMode::ByFilter => {
303 let nodes = self
304 .load_filtered_nodes(request, session_id, tenant_id, max_nodes)
305 .await?;
306 Ok(nodes
307 .iter()
308 .map(|node| build_candidate(node, String::new()))
309 .collect())
310 }
311 MemoryEvictMode::PurgeSession => Err(anyhow!("purge session uses dedicated path")),
312 }
313 }
314
315 async fn load_session_nodes(
316 &self,
317 request: &MemoryEvictRequest,
318 session_id: &str,
319 ) -> Result<Vec<SttpNode>> {
320 self.store
321 .query_nodes_async(NodeQuery {
322 limit: TENANT_SCAN_LIMIT,
323 session_id: Some(session_id.to_string()),
324 from_utc: request.scope.from_utc,
325 to_utc: request.scope.to_utc,
326 tiers: request.scope.tiers.clone(),
327 })
328 .await
329 }
330
331 async fn load_filtered_nodes(
332 &self,
333 request: &MemoryEvictRequest,
334 session_id: &str,
335 tenant_id: &str,
336 max_nodes: usize,
337 ) -> Result<Vec<SttpNode>> {
338 let mut nodes = self.load_session_nodes(request, session_id).await?;
339 let session_filter = build_session_filter(&request.scope);
340 let indexed_sync_keys = if let Some(index) = self.semantic_index.as_ref() {
341 resolve_indexed_sync_keys(
342 index.as_ref(),
343 tenant_id,
344 &request.filter,
345 Some(session_id),
346 TENANT_SCAN_LIMIT,
347 )
348 .await?
349 } else {
350 None
351 };
352
353 nodes.retain(|node| {
354 if let Some(keys) = &indexed_sync_keys
355 && !keys.contains(&node.sync_key)
356 {
357 return false;
358 }
359
360 node_matches_common_filters(
361 node,
362 &request.scope,
363 &request.filter,
364 session_filter.as_ref(),
365 )
366 });
367 nodes.truncate(max_nodes);
368 Ok(nodes)
369 }
370
371 async fn delete_tag_index_rows(&self, tenant_id: &str, sync_key: &str) -> Result<()> {
372 if let Some(index) = self.semantic_index.as_ref() {
373 index
374 .delete_node_tags_async(tenant_id, sync_key)
375 .await?;
376 }
377 Ok(())
378 }
379}
380
381fn single_session_id(scope: &crate::domain::memory::MemoryScope) -> Result<String> {
382 scope
383 .session_ids
384 .as_ref()
385 .and_then(|sessions| sessions.first().cloned())
386 .filter(|session| !session.trim().is_empty())
387 .ok_or_else(|| anyhow!("exactly one session id is required"))
388}
389
390fn resolve_tenant_id(
391 scope: &crate::domain::memory::MemoryScope,
392 session_id: &str,
393) -> String {
394 scope
395 .tenant_id
396 .clone()
397 .or_else(|| Some(derive_tenant_id_from_session(session_id)))
398 .unwrap_or_else(|| "default".to_string())
399}
400
401fn build_candidate(node: &SttpNode, store_node_id: String) -> EvictCandidate {
402 EvictCandidate {
403 store_node_id,
404 sync_key: node.sync_key.clone(),
405 graph_id: graph_node_id(node),
406 }
407}
408
409fn collect_inbound_refs(
410 candidate: &EvictCandidate,
411 nodes: &[SttpNode],
412 candidate_sync_keys: &HashSet<String>,
413) -> locus_core_rs::domain::models::InboundNodeReferences {
414 let mut child_parent_links = BTreeSet::new();
415 let mut incoming_semantic_refs = BTreeSet::new();
416
417 for node in nodes {
418 if candidate_sync_keys.contains(&node.sync_key) {
419 continue;
420 }
421
422 if let Some(parent) = node.parent_node_id.as_ref() {
423 if parent == &candidate.graph_id
424 || (!candidate.store_node_id.is_empty() && parent == &candidate.store_node_id)
425 {
426 child_parent_links.insert(node.sync_key.clone());
427 }
428 }
429
430 if let Some(links) = &node.semantic_links {
431 for link in links {
432 let target = link.target.trim();
433 if target_matches_candidate(target, candidate) {
434 incoming_semantic_refs.insert(node.sync_key.clone());
435 }
436 }
437 }
438 }
439
440 locus_core_rs::domain::models::InboundNodeReferences {
441 child_parent_links: child_parent_links.into_iter().collect(),
442 incoming_semantic_refs: incoming_semantic_refs.into_iter().collect(),
443 }
444}
445
446fn target_matches_candidate(target: &str, candidate: &EvictCandidate) -> bool {
447 let lower = target.to_ascii_lowercase();
448 if !candidate.sync_key.is_empty()
449 && lower == format!("ref:{}", candidate.sync_key.to_ascii_lowercase())
450 {
451 return true;
452 }
453 if !candidate.graph_id.is_empty()
454 && lower == format!("ref:{}", candidate.graph_id.to_ascii_lowercase())
455 {
456 return true;
457 }
458 if !candidate.store_node_id.is_empty()
459 && lower == format!("ref:{}", candidate.store_node_id.to_ascii_lowercase())
460 {
461 return true;
462 }
463 false
464}
465
466fn merge_delete_result(target: &mut NodeDeleteResult, source: NodeDeleteResult) {
467 target.deleted += source.deleted;
468 target.blocked += source.blocked;
469 target.not_found += source.not_found;
470 target.skipped += source.skipped;
471 target.records.extend(source.records);
472}
473
474fn map_core_record(record: &NodeDeleteRecord) -> MemoryEvictRecord {
475 let status = match record.status {
476 NodeDeleteStatus::Deleted => "deleted",
477 NodeDeleteStatus::NotFound => "not_found",
478 NodeDeleteStatus::Blocked => "blocked",
479 NodeDeleteStatus::Skipped => {
480 if record.reason.as_deref() == Some("would delete") {
481 "would_delete"
482 } else {
483 "skipped"
484 }
485 }
486 };
487
488 MemoryEvictRecord {
489 node_id: record.node_id.clone(),
490 sync_key: record.sync_key.clone(),
491 status: status.to_string(),
492 reason: record.reason.clone(),
493 inbound_references: None,
494 }
495}
496
497#[cfg(test)]
498mod tests {
499 use std::sync::Arc;
500
501 use chrono::Utc;
502 use locus_core_rs::domain::contracts::{NodeStore, SemanticIndexStore};
503 use locus_core_rs::domain::models::{AvecState, NodeUpsertStatus, SemanticLink, SttpNode};
504 use locus_core_rs::storage::{InMemoryNodeStore, InMemorySemanticIndexStore};
505 use locus_core_rs::SemanticIndexStoreInitializer;
506
507 use super::*;
508 use crate::domain::memory::{MemoryFilter, MemoryScope};
509
510 fn build_node(session_id: &str, sync_key: &str, parent: Option<&str>) -> SttpNode {
511 SttpNode {
512 raw: "raw".to_string(),
513 session_id: session_id.to_string(),
514 tier: "raw".to_string(),
515 timestamp: Utc::now(),
516 compression_depth: 1,
517 parent_node_id: parent.map(str::to_string),
518 sync_key: sync_key.to_string(),
519 updated_at: Utc::now(),
520 source_metadata: None,
521 context_summary: None,
522 semantic_tags: None,
523 semantic_links: None,
524 embedding: None,
525 embedding_model: None,
526 embedding_dimensions: None,
527 embedded_at: None,
528 user_avec: AvecState::zero(),
529 model_avec: AvecState::zero(),
530 compression_avec: None,
531 rho: 0.9,
532 kappa: 0.9,
533 psi: 2.0,
534 }
535 }
536
537 async fn seed_parent_child(store: &InMemoryNodeStore, session: &str) -> (String, String) {
538 let parent = build_node(session, "parent-sync", None);
539 let parent_graph = graph_node_id(&parent);
540 let parent_result = store
541 .upsert_node_async(parent)
542 .await
543 .expect("parent upsert");
544 assert_eq!(parent_result.status, NodeUpsertStatus::Created);
545
546 let child = build_node(session, "child-sync", Some(&parent_graph));
547 store.upsert_node_async(child).await.expect("child upsert");
548 (parent_result.node_id, parent_result.sync_key)
549 }
550
551 #[tokio::test(flavor = "current_thread")]
552 async fn blocked_when_child_parent_link_exists() {
553 let store = Arc::new(InMemoryNodeStore::new());
554 let session = "demo";
555 let (_parent_id, parent_sync) = seed_parent_child(store.as_ref(), session).await;
556
557 let service = MemoryEvictService::new(store.clone());
558 let result = service
559 .execute(&MemoryEvictRequest {
560 mode: MemoryEvictMode::BySyncKeys,
561 scope: MemoryScope {
562 session_ids: Some(vec![session.to_string()]),
563 ..Default::default()
564 },
565 sync_keys: Some(vec![parent_sync.clone()]),
566 dry_run: false,
567 force: false,
568 ..Default::default()
569 })
570 .await
571 .expect("evict");
572
573 assert_eq!(result.blocked, 1);
574 assert_eq!(result.deleted, 0);
575 }
576
577 #[tokio::test(flavor = "current_thread")]
578 async fn force_deletes_despite_lineage() {
579 let store = Arc::new(InMemoryNodeStore::new());
580 let session = "demo";
581 let (_parent_id, parent_sync) = seed_parent_child(store.as_ref(), session).await;
582
583 let service = MemoryEvictService::new(store.clone());
584 let result = service
585 .execute(&MemoryEvictRequest {
586 mode: MemoryEvictMode::BySyncKeys,
587 scope: MemoryScope {
588 session_ids: Some(vec![session.to_string()]),
589 ..Default::default()
590 },
591 sync_keys: Some(vec![parent_sync]),
592 dry_run: false,
593 force: true,
594 ..Default::default()
595 })
596 .await
597 .expect("evict");
598
599 assert_eq!(result.deleted, 1);
600 }
601
602 #[tokio::test(flavor = "current_thread")]
603 async fn dry_run_does_not_mutate_store() {
604 let store = Arc::new(InMemoryNodeStore::new());
605 let session = "demo";
606 store
607 .upsert_node_async(build_node(session, "sync-42", None))
608 .await
609 .expect("upsert");
610
611 let service = MemoryEvictService::new(store.clone());
612 let result = service
613 .execute(&MemoryEvictRequest {
614 mode: MemoryEvictMode::BySyncKeys,
615 scope: MemoryScope {
616 session_ids: Some(vec![session.to_string()]),
617 ..Default::default()
618 },
619 sync_keys: Some(vec!["sync-42".to_string()]),
620 dry_run: true,
621 ..Default::default()
622 })
623 .await
624 .expect("evict");
625
626 assert_eq!(result.deleted, 1);
627 assert_eq!(result.would_delete, vec!["sync-42".to_string()]);
628
629 let remaining = store
630 .query_nodes_async(NodeQuery {
631 limit: 10,
632 session_id: Some(session.to_string()),
633 ..Default::default()
634 })
635 .await
636 .expect("query");
637 assert_eq!(remaining.len(), 1);
638 }
639
640 #[tokio::test(flavor = "current_thread")]
641 async fn blocked_by_inbound_semantic_ref() {
642 let store = Arc::new(InMemoryNodeStore::new());
643 let session = "demo";
644 let target = build_node(session, "target-sync", None);
645 let target_graph = graph_node_id(&target);
646 store.upsert_node_async(target).await.expect("target");
647
648 let mut referrer = build_node(session, "referrer-sync", None);
649 referrer.semantic_links = Some(vec![SemanticLink {
650 rel: "evaluates".to_string(),
651 target: format!("ref:{target_graph}"),
652 confidence: Some(0.9),
653 }]);
654 store.upsert_node_async(referrer).await.expect("referrer");
655
656 let service = MemoryEvictService::new(store.clone());
657 let result = service
658 .execute(&MemoryEvictRequest {
659 mode: MemoryEvictMode::BySyncKeys,
660 scope: MemoryScope {
661 session_ids: Some(vec![session.to_string()]),
662 ..Default::default()
663 },
664 sync_keys: Some(vec!["target-sync".to_string()]),
665 dry_run: false,
666 force: false,
667 ..Default::default()
668 })
669 .await
670 .expect("evict");
671
672 assert_eq!(result.blocked, 1);
673 }
674
675 #[tokio::test(flavor = "current_thread")]
676 async fn delete_removes_semantic_tag_index_rows() {
677 let store = Arc::new(InMemoryNodeStore::new());
678 let index = Arc::new(InMemorySemanticIndexStore::new());
679 index.initialize_async().await.expect("init");
680 let session = "demo";
681 let mut node = build_node(session, "tagged-sync", None);
682 node.semantic_tags = Some(vec!["stale".to_string(), "debug".to_string()]);
683 store.upsert_node_async(node).await.expect("upsert");
684
685 index
686 .sync_node_tags_async(
687 locus_core_rs::domain::models::SemanticTagNodeRef {
688 tenant_id: "default".to_string(),
689 session_id: session.to_string(),
690 node_id: "node".to_string(),
691 sync_key: "tagged-sync".to_string(),
692 },
693 &["stale".to_string(), "debug".to_string()],
694 None,
695 )
696 .await
697 .expect("sync tags");
698
699 let service = MemoryEvictService::new(store.clone()).with_semantic_index(index.clone());
700 let result = service
701 .execute(&MemoryEvictRequest {
702 mode: MemoryEvictMode::BySyncKeys,
703 scope: MemoryScope {
704 session_ids: Some(vec![session.to_string()]),
705 ..Default::default()
706 },
707 sync_keys: Some(vec!["tagged-sync".to_string()]),
708 dry_run: false,
709 force: true,
710 ..Default::default()
711 })
712 .await
713 .expect("evict");
714
715 assert_eq!(result.deleted, 1);
716
717 let tags = index
718 .query_tag_records_async(locus_core_rs::domain::models::SemanticTagQueryFilter {
719 tenant_id: Some("default".to_string()),
720 session_id: Some(session.to_string()),
721 ..Default::default()
722 })
723 .await
724 .expect("query tags");
725 assert!(tags.is_empty());
726 }
727
728 #[tokio::test(flavor = "current_thread")]
729 async fn session_purge_removes_nodes_and_checkpoints() {
730 let store = Arc::new(InMemoryNodeStore::new());
731 let session = "temp-import";
732 store
733 .upsert_node_async(build_node(session, "one", None))
734 .await
735 .expect("upsert");
736 store
737 .upsert_node_async(build_node(session, "two", None))
738 .await
739 .expect("upsert");
740 store
741 .put_checkpoint_async(locus_core_rs::domain::models::SyncCheckpoint {
742 session_id: session.to_string(),
743 connector_id: "demo".to_string(),
744 cursor: None,
745 metadata: None,
746 updated_at: Utc::now(),
747 })
748 .await
749 .expect("checkpoint");
750
751 let service = MemoryEvictService::new(store.clone());
752 let preview = service
753 .execute(&MemoryEvictRequest {
754 mode: MemoryEvictMode::PurgeSession,
755 scope: MemoryScope {
756 session_ids: Some(vec![session.to_string()]),
757 ..Default::default()
758 },
759 dry_run: true,
760 include_checkpoints: true,
761 ..Default::default()
762 })
763 .await
764 .expect("preview");
765 assert_eq!(preview.deleted, 2);
766 assert_eq!(preview.checkpoints_deleted, 1);
767
768 let applied = service
769 .execute(&MemoryEvictRequest {
770 mode: MemoryEvictMode::PurgeSession,
771 scope: MemoryScope {
772 session_ids: Some(vec![session.to_string()]),
773 ..Default::default()
774 },
775 include_checkpoints: true,
776 include_calibration: true,
777 ..Default::default()
778 })
779 .await
780 .expect("purge");
781 assert_eq!(applied.deleted, 2);
782
783 let remaining = store
784 .query_nodes_async(NodeQuery {
785 limit: 10,
786 session_id: Some(session.to_string()),
787 ..Default::default()
788 })
789 .await
790 .expect("query");
791 assert!(remaining.is_empty());
792 }
793
794 #[tokio::test(flavor = "current_thread")]
795 async fn filter_mode_evicts_matching_nodes() {
796 let store = Arc::new(InMemoryNodeStore::new());
797 let session = "demo";
798 let mut stale = build_node(session, "stale-node", None);
799 stale.semantic_tags = Some(vec!["stale".to_string()]);
800 store.upsert_node_async(stale).await.expect("stale");
801 store
802 .upsert_node_async(build_node(session, "fresh-node", None))
803 .await
804 .expect("fresh");
805
806 let service = MemoryEvictService::new(store.clone());
807 let result = service
808 .execute(&MemoryEvictRequest {
809 mode: MemoryEvictMode::ByFilter,
810 scope: MemoryScope {
811 session_ids: Some(vec![session.to_string()]),
812 ..Default::default()
813 },
814 filter: MemoryFilter {
815 tags_contains: Some(vec!["stale".to_string()]),
816 ..Default::default()
817 },
818 force: true,
819 ..Default::default()
820 })
821 .await
822 .expect("evict");
823
824 assert_eq!(result.deleted, 1);
825 let remaining = store
826 .query_nodes_async(NodeQuery {
827 limit: 10,
828 session_id: Some(session.to_string()),
829 ..Default::default()
830 })
831 .await
832 .expect("query");
833 assert_eq!(remaining.len(), 1);
834 assert_eq!(remaining[0].sync_key, "fresh-node");
835 }
836}