Skip to main content

stasis/application/runtime/
memory_rollup_job_handler.rs

1use std::sync::Arc;
2
3use async_trait::async_trait;
4use serde_json::json;
5
6use crate::application::orchestration::runtime_job_payloads::MemoryRollupJobPayload;
7use crate::application::runtime::in_memory_runtime::{JobExecutionOutcome, JobHandler};
8use crate::application::runtime::memory_operation_job_outcome_helpers::{
9    operation_failure, operation_success, policy_violation_failure,
10};
11use crate::domain::errors::Result;
12use crate::domain::runtime::job::Job;
13use crate::ports::outbound::memory::memory_models::{MemoryRollupRequest, MemoryScope};
14use crate::ports::outbound::memory::memory_operations::MemoryOperations;
15
16pub struct MemoryRollupJobHandler {
17    operations: Arc<dyn MemoryOperations>,
18}
19
20impl MemoryRollupJobHandler {
21    pub fn new(operations: Arc<dyn MemoryOperations>) -> Self {
22        Self { operations }
23    }
24
25    fn parse_payload(raw: &str) -> std::result::Result<MemoryRollupJobPayload, String> {
26        serde_json::from_str(raw)
27            .map_err(|err| format!("policy violation: invalid memory-rollup payload json: {err}"))
28    }
29}
30
31#[async_trait]
32impl JobHandler for MemoryRollupJobHandler {
33    fn job_type(&self) -> &'static str {
34        "workflow.stasis.memory.rollup"
35    }
36
37    async fn execute(&self, job: &Job) -> Result<JobExecutionOutcome> {
38        let payload = match Self::parse_payload(&job.payload_ref) {
39            Ok(payload) => payload,
40            Err(message) => return Ok(policy_violation_failure("stasis-memory-rollup", message)),
41        };
42
43        let request = MemoryRollupRequest {
44            scope: MemoryScope {
45                session_ids: payload.session_ids,
46                tiers: payload.tiers,
47                from_utc: payload.from_utc,
48                to_utc: payload.to_utc,
49            },
50            max_days: payload.max_days.unwrap_or(30),
51            max_nodes: payload.max_nodes.unwrap_or(5000),
52        };
53
54        match self.operations.rollup(&request).await {
55            Ok(result) => Ok(operation_success(
56                "stasis-memory-rollup",
57                "memory-rollup",
58                &job.id,
59                json!({
60                    "total_groups": result.total_groups,
61                    "scanned_nodes": result.scanned_nodes,
62                }),
63            )),
64            Err(err) => Ok(operation_failure("stasis-memory-rollup", err.to_string())),
65        }
66    }
67}