Skip to main content

stasis/application/orchestration/
runtime_job_payloads.rs

1use chrono::{DateTime, Utc};
2use serde::{Deserialize, Serialize};
3use serde_json::Value;
4
5use crate::domain::errors::{Result, StasisError};
6
7#[derive(Clone, Debug, Serialize, Deserialize)]
8#[serde(rename_all = "snake_case")]
9pub enum AgentToolCallMode {
10    Auto,
11    Strict,
12}
13
14#[derive(Clone, Debug, Serialize, Deserialize)]
15#[serde(rename_all = "snake_case")]
16pub enum MemoryFallbackPolicyPayload {
17    Never,
18    OnEmpty,
19    Always,
20}
21
22#[derive(Clone, Debug, Serialize, Deserialize)]
23#[serde(rename_all = "snake_case")]
24pub enum MemoryStrictnessModePayload {
25    Precision,
26    Balanced,
27    Recall,
28}
29
30#[derive(Clone, Debug, Serialize, Deserialize)]
31#[serde(rename_all = "snake_case")]
32pub enum MemoryStoreModePayload {
33    Disabled,
34    SummaryOnly,
35    Full,
36}
37
38#[derive(Clone, Debug, Serialize, Deserialize)]
39#[serde(rename_all = "camelCase")]
40pub struct MemoryPolicyPayload {
41    pub session_ids: Option<Vec<String>>,
42    pub tiers: Option<Vec<String>>,
43    pub from_utc: Option<DateTime<Utc>>,
44    pub to_utc: Option<DateTime<Utc>>,
45    pub limit: Option<usize>,
46    pub alpha: Option<f32>,
47    pub beta: Option<f32>,
48    pub fallback_policy: Option<MemoryFallbackPolicyPayload>,
49    pub strictness: Option<MemoryStrictnessModePayload>,
50    pub query_text: Option<String>,
51    pub include_explain: Option<bool>,
52    pub store_mode: Option<MemoryStoreModePayload>,
53}
54
55#[derive(Clone, Debug, Serialize, Deserialize)]
56pub struct AgentSessionParticipantPayload {
57    pub agent_id: String,
58    pub system_prompt: Option<String>,
59    pub tool_name: String,
60    pub tool_input: Option<Value>,
61}
62
63#[derive(Clone, Debug, Serialize, Deserialize)]
64pub struct AgentSessionJobPayload {
65    pub thread_id: Option<String>,
66    pub initial_user_prompt: String,
67    pub participants: Vec<AgentSessionParticipantPayload>,
68    pub policy_profile: Option<String>,
69    pub model_hint: Option<String>,
70    pub max_turns: Option<usize>,
71    pub tool_call_mode: Option<AgentToolCallMode>,
72    pub memory_policy: Option<MemoryPolicyPayload>,
73}
74
75impl AgentSessionJobPayload {
76    pub fn to_payload_ref(&self) -> Result<String> {
77        serde_json::to_string(self).map_err(|err| {
78            StasisError::PortFailure(format!("failed to encode agent-session payload: {err}"))
79        })
80    }
81}
82
83#[derive(Clone, Debug, Serialize, Deserialize)]
84pub struct AgentTurnJobPayload {
85    pub agent_id: String,
86    pub thread_id: Option<String>,
87    pub user_prompt: String,
88    pub system_prompt: Option<String>,
89    pub policy_profile: Option<String>,
90    pub model_hint: Option<String>,
91    pub tool_name: String,
92    pub tool_input: Option<Value>,
93    pub tool_call_mode: Option<AgentToolCallMode>,
94    pub memory_policy: Option<MemoryPolicyPayload>,
95}
96
97impl AgentTurnJobPayload {
98    pub fn to_payload_ref(&self) -> Result<String> {
99        serde_json::to_string(self).map_err(|err| {
100            StasisError::PortFailure(format!("failed to encode agent-turn payload: {err}"))
101        })
102    }
103}
104
105#[derive(Clone, Debug, Serialize, Deserialize)]
106pub struct ToolLoopJobPayload {
107    pub user_prompt: String,
108    pub system_prompt: Option<String>,
109    pub policy_profile: Option<String>,
110    pub model_hint: Option<String>,
111    pub tool_name: String,
112    pub tool_input: Option<Value>,
113    pub tool_call_mode: Option<AgentToolCallMode>,
114    pub memory_policy: Option<MemoryPolicyPayload>,
115}
116
117impl ToolLoopJobPayload {
118    pub fn to_payload_ref(&self) -> Result<String> {
119        serde_json::to_string(self).map_err(|err| {
120            StasisError::PortFailure(format!("failed to encode tool-loop payload: {err}"))
121        })
122    }
123}
124
125#[derive(Clone, Debug, Serialize, Deserialize)]
126pub struct PromptJobPayload {
127    pub user_prompt: String,
128    pub system_prompt: Option<String>,
129    pub policy_profile: Option<String>,
130    pub model_hint: Option<String>,
131    pub memory_policy: Option<MemoryPolicyPayload>,
132}
133
134impl PromptJobPayload {
135    pub fn to_payload_ref(&self) -> Result<String> {
136        serde_json::to_string(self).map_err(|err| {
137            StasisError::PortFailure(format!("failed to encode prompt payload: {err}"))
138        })
139    }
140}
141
142#[derive(Clone, Debug, Serialize, Deserialize)]
143pub struct SequentialStageJobPayload {
144    pub stage_id: String,
145    pub user_prompt_template: String,
146    pub system_prompt: Option<String>,
147    pub policy_profile: Option<String>,
148    pub model_hint: Option<String>,
149}
150
151#[derive(Clone, Debug, Serialize, Deserialize)]
152pub struct SequentialPatternJobPayload {
153    pub thread_id: Option<String>,
154    pub initial_user_prompt: String,
155    pub policy_profile: Option<String>,
156    pub model_hint: Option<String>,
157    pub stages: Vec<SequentialStageJobPayload>,
158}
159
160impl SequentialPatternJobPayload {
161    pub fn to_payload_ref(&self) -> Result<String> {
162        serde_json::to_string(self).map_err(|err| {
163            StasisError::PortFailure(format!(
164                "failed to encode sequential-pattern payload: {err}"
165            ))
166        })
167    }
168}
169
170#[derive(Clone, Debug, Serialize, Deserialize)]
171pub struct ConcurrentBranchJobPayload {
172    pub branch_id: String,
173    pub user_prompt_template: String,
174    pub system_prompt: Option<String>,
175    pub policy_profile: Option<String>,
176    pub model_hint: Option<String>,
177}
178
179#[derive(Clone, Debug, Serialize, Deserialize)]
180pub struct ConcurrentPatternJobPayload {
181    pub thread_id: Option<String>,
182    pub initial_user_prompt: String,
183    pub policy_profile: Option<String>,
184    pub model_hint: Option<String>,
185    pub merge_strategy: Option<String>,
186    pub branches: Vec<ConcurrentBranchJobPayload>,
187}
188
189impl ConcurrentPatternJobPayload {
190    pub fn to_payload_ref(&self) -> Result<String> {
191        serde_json::to_string(self).map_err(|err| {
192            StasisError::PortFailure(format!(
193                "failed to encode concurrent-pattern payload: {err}"
194            ))
195        })
196    }
197}
198
199#[derive(Clone, Debug, Serialize, Deserialize)]
200pub struct HandoffTurnJobPayload {
201    pub actor_id: String,
202    pub user_prompt_template: String,
203    pub system_prompt: Option<String>,
204    pub policy_profile: Option<String>,
205    pub model_hint: Option<String>,
206}
207
208#[derive(Clone, Debug, Serialize, Deserialize)]
209pub struct HandoffPatternJobPayload {
210    pub thread_id: Option<String>,
211    pub initial_user_prompt: String,
212    pub policy_profile: Option<String>,
213    pub model_hint: Option<String>,
214    pub turns: Vec<HandoffTurnJobPayload>,
215}
216
217impl HandoffPatternJobPayload {
218    pub fn to_payload_ref(&self) -> Result<String> {
219        serde_json::to_string(self).map_err(|err| {
220            StasisError::PortFailure(format!("failed to encode handoff-pattern payload: {err}"))
221        })
222    }
223}
224
225#[derive(Clone, Debug, Serialize, Deserialize)]
226pub struct OrchestratorRouteJobPayload {
227    pub route_id: String,
228    pub selector_keywords: Vec<String>,
229    pub user_prompt_template: String,
230    pub system_prompt: Option<String>,
231    pub policy_profile: Option<String>,
232    pub model_hint: Option<String>,
233}
234
235#[derive(Clone, Debug, Serialize, Deserialize)]
236pub struct OrchestratorPatternJobPayload {
237    pub thread_id: Option<String>,
238    pub initial_user_prompt: String,
239    pub policy_profile: Option<String>,
240    pub model_hint: Option<String>,
241    pub routes: Vec<OrchestratorRouteJobPayload>,
242}
243
244impl OrchestratorPatternJobPayload {
245    pub fn to_payload_ref(&self) -> Result<String> {
246        serde_json::to_string(self).map_err(|err| {
247            StasisError::PortFailure(format!(
248                "failed to encode orchestrator-pattern payload: {err}"
249            ))
250        })
251    }
252}
253
254#[derive(Clone, Debug, Serialize, Deserialize)]
255pub struct MemoryRecallJobPayload {
256    pub memory_policy: Option<MemoryPolicyPayload>,
257}
258
259impl MemoryRecallJobPayload {
260    pub fn to_payload_ref(&self) -> Result<String> {
261        serde_json::to_string(self).map_err(|err| {
262            StasisError::PortFailure(format!("failed to encode memory-recall payload: {err}"))
263        })
264    }
265}
266
267#[derive(Clone, Debug, Serialize, Deserialize)]
268#[serde(rename_all = "camelCase")]
269pub struct MemoryFindJobPayload {
270    pub session_ids: Option<Vec<String>>,
271    pub tiers: Option<Vec<String>>,
272    pub from_utc: Option<DateTime<Utc>>,
273    pub to_utc: Option<DateTime<Utc>>,
274    pub limit: Option<usize>,
275    pub cursor: Option<String>,
276    pub text_contains: Option<String>,
277    pub sort_field: Option<String>,
278    pub sort_direction: Option<String>,
279}
280
281impl MemoryFindJobPayload {
282    pub fn to_payload_ref(&self) -> Result<String> {
283        serde_json::to_string(self).map_err(|err| {
284            StasisError::PortFailure(format!("failed to encode memory-find payload: {err}"))
285        })
286    }
287}
288
289#[derive(Clone, Debug, Serialize, Deserialize)]
290#[serde(rename_all = "camelCase")]
291pub struct MemoryAggregateJobPayload {
292    pub session_ids: Option<Vec<String>>,
293    pub tiers: Option<Vec<String>>,
294    pub from_utc: Option<DateTime<Utc>>,
295    pub to_utc: Option<DateTime<Utc>>,
296    pub max_groups: Option<usize>,
297    pub max_nodes: Option<usize>,
298}
299
300impl MemoryAggregateJobPayload {
301    pub fn to_payload_ref(&self) -> Result<String> {
302        serde_json::to_string(self).map_err(|err| {
303            StasisError::PortFailure(format!("failed to encode memory-aggregate payload: {err}"))
304        })
305    }
306}
307
308#[derive(Clone, Debug, Serialize, Deserialize)]
309#[serde(rename_all = "snake_case")]
310pub enum MemoryTransformOperationPayload {
311    EmbedBackfill,
312    ReindexEmbeddings,
313}
314
315#[derive(Clone, Debug, Serialize, Deserialize)]
316#[serde(rename_all = "camelCase")]
317pub struct MemoryTransformJobPayload {
318    pub session_ids: Option<Vec<String>>,
319    pub tiers: Option<Vec<String>>,
320    pub from_utc: Option<DateTime<Utc>>,
321    pub to_utc: Option<DateTime<Utc>>,
322    pub operation: Option<MemoryTransformOperationPayload>,
323    pub dry_run: Option<bool>,
324    pub batch_size: Option<usize>,
325    pub max_nodes: Option<usize>,
326    pub provider_id: Option<String>,
327    pub model: Option<String>,
328}
329
330impl MemoryTransformJobPayload {
331    pub fn to_payload_ref(&self) -> Result<String> {
332        serde_json::to_string(self).map_err(|err| {
333            StasisError::PortFailure(format!("failed to encode memory-transform payload: {err}"))
334        })
335    }
336}
337
338#[derive(Clone, Debug, Serialize, Deserialize)]
339#[serde(rename_all = "camelCase")]
340pub struct MemoryRollupJobPayload {
341    pub session_ids: Option<Vec<String>>,
342    pub tiers: Option<Vec<String>>,
343    pub from_utc: Option<DateTime<Utc>>,
344    pub to_utc: Option<DateTime<Utc>>,
345    pub max_days: Option<usize>,
346    pub max_nodes: Option<usize>,
347}
348
349impl MemoryRollupJobPayload {
350    pub fn to_payload_ref(&self) -> Result<String> {
351        serde_json::to_string(self).map_err(|err| {
352            StasisError::PortFailure(format!("failed to encode memory-rollup payload: {err}"))
353        })
354    }
355}
356
357#[derive(Clone, Debug, Default, Serialize, Deserialize)]
358pub struct MemorySchemaJobPayload {}
359
360impl MemorySchemaJobPayload {
361    pub fn to_payload_ref(&self) -> Result<String> {
362        serde_json::to_string(self).map_err(|err| {
363            StasisError::PortFailure(format!("failed to encode memory-schema payload: {err}"))
364        })
365    }
366}