Skip to main content

locus_sdk/application/
memory_evict.rs

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}