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