stasis/application/runtime/
memory_evict_job_handler.rs1use std::sync::Arc;
2
3use async_trait::async_trait;
4use serde_json::json;
5
6use crate::application::orchestration::runtime_job_payloads::{
7 MemoryEvictJobPayload, MemoryEvictModePayload,
8};
9use crate::application::runtime::in_memory_runtime::{JobExecutionOutcome, JobHandler};
10use crate::application::runtime::memory_job_request_helpers::memory_scope_from_fields;
11use crate::application::runtime::memory_operation_job_outcome_helpers::{
12 operation_failure, operation_success, policy_violation_failure,
13};
14use crate::application::runtime::memory_recall_request_builder::memory_filter_from_payload;
15use crate::domain::errors::Result;
16use crate::domain::runtime::job::Job;
17use crate::ports::outbound::memory::memory_models::{MemoryEvictMode, MemoryEvictRequest};
18use crate::ports::outbound::memory::memory_operations::MemoryOperations;
19
20pub struct MemoryEvictJobHandler {
21 operations: Arc<dyn MemoryOperations>,
22}
23
24impl MemoryEvictJobHandler {
25 pub fn new(operations: Arc<dyn MemoryOperations>) -> Self {
26 Self { operations }
27 }
28
29 fn parse_payload(raw: &str) -> std::result::Result<MemoryEvictJobPayload, String> {
30 serde_json::from_str(raw)
31 .map_err(|err| format!("policy violation: invalid memory-evict payload json: {err}"))
32 }
33
34 fn map_mode(value: Option<MemoryEvictModePayload>) -> MemoryEvictMode {
35 match value {
36 Some(MemoryEvictModePayload::ByNodeIds) => MemoryEvictMode::ByNodeIds,
37 Some(MemoryEvictModePayload::ByFilter) => MemoryEvictMode::ByFilter,
38 Some(MemoryEvictModePayload::PurgeSession) => MemoryEvictMode::PurgeSession,
39 _ => MemoryEvictMode::BySyncKeys,
40 }
41 }
42}
43
44#[async_trait]
45impl JobHandler for MemoryEvictJobHandler {
46 fn job_type(&self) -> &'static str {
47 "workflow.stasis.memory.evict"
48 }
49
50 async fn execute(&self, job: &Job) -> Result<JobExecutionOutcome> {
51 let payload = match Self::parse_payload(&job.payload_ref) {
52 Ok(payload) => payload,
53 Err(message) => return Ok(policy_violation_failure("stasis-memory-evict", message)),
54 };
55
56 let request = MemoryEvictRequest {
57 mode: Self::map_mode(payload.mode),
58 scope: memory_scope_from_fields(
59 payload.tenant_id,
60 payload.session_ids,
61 payload.tiers,
62 payload.from_utc,
63 payload.to_utc,
64 ),
65 filter: memory_filter_from_payload(&payload.filter),
66 sync_keys: payload.sync_keys,
67 node_ids: payload.node_ids,
68 dry_run: payload.dry_run.unwrap_or(true),
69 force: payload.force.unwrap_or(false),
70 max_nodes: payload.max_nodes.unwrap_or(5000),
71 include_calibration: payload.include_calibration.unwrap_or(false),
72 include_checkpoints: payload.include_checkpoints.unwrap_or(false),
73 };
74
75 match self.operations.evict(&request).await {
76 Ok(result) => Ok(operation_success(
77 "stasis-memory-evict",
78 "memory-evict",
79 &job.id,
80 json!({
81 "dry_run": result.dry_run,
82 "deleted": result.deleted,
83 "blocked": result.blocked,
84 "not_found": result.not_found,
85 "skipped": result.skipped,
86 "would_delete": result.would_delete,
87 "calibrations_deleted": result.calibrations_deleted,
88 "checkpoints_deleted": result.checkpoints_deleted,
89 "records": result.records,
90 }),
91 )),
92 Err(err) => Ok(operation_failure("stasis-memory-evict", err.to_string())),
93 }
94 }
95}