stasis/application/runtime/
memory_graph_job_handler.rs1use std::sync::Arc;
2
3use async_trait::async_trait;
4use serde_json::json;
5
6use crate::application::orchestration::runtime_job_payloads::MemoryGraphJobPayload;
7use crate::application::runtime::in_memory_runtime::{JobExecutionOutcome, JobHandler};
8use crate::application::runtime::memory_job_request_helpers::memory_scope_from_fields;
9use crate::application::runtime::memory_operation_job_outcome_helpers::{
10 operation_failure, operation_success, policy_violation_failure,
11};
12use crate::application::runtime::memory_recall_request_builder::memory_filter_from_payload;
13use crate::domain::errors::Result;
14use crate::domain::runtime::job::Job;
15use crate::ports::outbound::memory::memory_context_reader::MemoryContextReader;
16use crate::ports::outbound::memory::memory_models::MemoryGraphRequest;
17
18pub struct MemoryGraphJobHandler {
19 reader: Arc<dyn MemoryContextReader>,
20}
21
22impl MemoryGraphJobHandler {
23 pub fn new(reader: Arc<dyn MemoryContextReader>) -> Self {
24 Self { reader }
25 }
26
27 fn parse_payload(raw: &str) -> std::result::Result<MemoryGraphJobPayload, String> {
28 serde_json::from_str(raw)
29 .map_err(|err| format!("policy violation: invalid memory-graph payload json: {err}"))
30 }
31}
32
33#[async_trait]
34impl JobHandler for MemoryGraphJobHandler {
35 fn job_type(&self) -> &'static str {
36 "workflow.stasis.memory.graph"
37 }
38
39 async fn execute(&self, job: &Job) -> Result<JobExecutionOutcome> {
40 let payload = match Self::parse_payload(&job.payload_ref) {
41 Ok(payload) => payload,
42 Err(message) => return Ok(policy_violation_failure("stasis-memory-graph", message)),
43 };
44
45 let request = MemoryGraphRequest {
46 scope: memory_scope_from_fields(
47 payload.tenant_id,
48 payload.session_ids,
49 payload.tiers,
50 payload.from_utc,
51 payload.to_utc,
52 ),
53 filter: memory_filter_from_payload(&payload.filter),
54 include_lineage: payload.include_lineage.unwrap_or(true),
55 include_semantic: payload.include_semantic.unwrap_or(true),
56 include_session_topology: payload.include_session_topology.unwrap_or(true),
57 rel: payload.rel,
58 target_prefix: payload.target_prefix,
59 limit: payload.limit.unwrap_or(200),
60 };
61
62 match self.reader.graph(&request).await {
63 Ok(result) => Ok(operation_success(
64 "stasis-memory-graph",
65 "memory-graph",
66 &job.id,
67 json!({
68 "retrieved": result.retrieved,
69 "sessions": result.sessions,
70 "nodes": result.nodes,
71 "edges": result.edges,
72 }),
73 )),
74 Err(err) => Ok(operation_failure("stasis-memory-graph", err.to_string())),
75 }
76 }
77}