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, Default, Serialize, Deserialize)]
39#[serde(rename_all = "camelCase")]
40pub struct MemoryMetricRangePayload {
41    pub min: Option<f32>,
42    pub max: Option<f32>,
43}
44
45#[derive(Clone, Debug, Default, Serialize, Deserialize)]
46#[serde(rename_all = "camelCase")]
47pub struct MemoryFilterPayload {
48    pub has_embedding: Option<bool>,
49    pub embedding_model: Option<String>,
50    pub psi: Option<MemoryMetricRangePayload>,
51    pub rho: Option<MemoryMetricRangePayload>,
52    pub kappa: Option<MemoryMetricRangePayload>,
53    pub text_contains: Option<String>,
54    pub tags_contains: Option<Vec<String>>,
55    pub has_tag: Option<String>,
56    pub indexed_tags: Option<Vec<String>>,
57    pub tag_prefix: Option<String>,
58    pub has_semantic_links: Option<bool>,
59    pub link_rel: Option<String>,
60    pub link_target: Option<String>,
61    pub links_to_ref: Option<String>,
62}
63
64#[derive(Clone, Debug, Serialize, Deserialize)]
65#[serde(rename_all = "camelCase")]
66pub struct MemoryPolicyPayload {
67    pub tenant_id: Option<String>,
68    pub session_ids: Option<Vec<String>>,
69    pub tiers: Option<Vec<String>>,
70    pub from_utc: Option<DateTime<Utc>>,
71    pub to_utc: Option<DateTime<Utc>>,
72    pub limit: Option<usize>,
73    pub alpha: Option<f32>,
74    pub beta: Option<f32>,
75    pub gamma: Option<f32>,
76    pub fallback_policy: Option<MemoryFallbackPolicyPayload>,
77    pub strictness: Option<MemoryStrictnessModePayload>,
78    pub query_text: Option<String>,
79    pub include_explain: Option<bool>,
80    pub store_mode: Option<MemoryStoreModePayload>,
81    #[serde(default)]
82    pub filter: MemoryFilterPayload,
83}
84
85#[derive(Clone, Debug, Serialize, Deserialize)]
86pub struct AgentSessionParticipantPayload {
87    pub agent_id: String,
88    pub system_prompt: Option<String>,
89    pub tool_name: String,
90    pub tool_input: Option<Value>,
91}
92
93#[derive(Clone, Debug, Serialize, Deserialize)]
94pub struct AgentSessionJobPayload {
95    pub thread_id: Option<String>,
96    pub initial_user_prompt: String,
97    pub participants: Vec<AgentSessionParticipantPayload>,
98    pub policy_profile: Option<String>,
99    pub model_hint: Option<String>,
100    pub reasoning_effort: Option<String>,
101    pub max_turns: Option<usize>,
102    pub tool_call_mode: Option<AgentToolCallMode>,
103    pub memory_policy: Option<MemoryPolicyPayload>,
104}
105
106impl AgentSessionJobPayload {
107    pub fn to_payload_ref(&self) -> Result<String> {
108        serde_json::to_string(self).map_err(|err| {
109            StasisError::PortFailure(format!("failed to encode agent-session payload: {err}"))
110        })
111    }
112}
113
114#[derive(Clone, Debug, Serialize, Deserialize)]
115pub struct AgentTurnJobPayload {
116    pub agent_id: String,
117    pub thread_id: Option<String>,
118    pub user_prompt: String,
119    pub system_prompt: Option<String>,
120    pub policy_profile: Option<String>,
121    pub model_hint: Option<String>,
122    pub reasoning_effort: Option<String>,
123    pub tool_name: String,
124    pub tool_input: Option<Value>,
125    pub tool_call_mode: Option<AgentToolCallMode>,
126    pub memory_policy: Option<MemoryPolicyPayload>,
127}
128
129impl AgentTurnJobPayload {
130    pub fn to_payload_ref(&self) -> Result<String> {
131        serde_json::to_string(self).map_err(|err| {
132            StasisError::PortFailure(format!("failed to encode agent-turn payload: {err}"))
133        })
134    }
135}
136
137#[derive(Clone, Debug, Serialize, Deserialize)]
138pub struct ToolLoopJobPayload {
139    pub user_prompt: String,
140    pub system_prompt: Option<String>,
141    pub policy_profile: Option<String>,
142    pub model_hint: Option<String>,
143    pub reasoning_effort: Option<String>,
144    pub tool_name: String,
145    pub tool_input: Option<Value>,
146    pub tool_call_mode: Option<AgentToolCallMode>,
147    pub memory_policy: Option<MemoryPolicyPayload>,
148}
149
150impl ToolLoopJobPayload {
151    pub fn to_payload_ref(&self) -> Result<String> {
152        serde_json::to_string(self).map_err(|err| {
153            StasisError::PortFailure(format!("failed to encode tool-loop payload: {err}"))
154        })
155    }
156}
157
158#[derive(Clone, Debug, Serialize, Deserialize)]
159pub struct PromptJobPayload {
160    pub user_prompt: String,
161    pub system_prompt: Option<String>,
162    pub policy_profile: Option<String>,
163    pub model_hint: Option<String>,
164    pub reasoning_effort: Option<String>,
165    pub memory_policy: Option<MemoryPolicyPayload>,
166}
167
168impl PromptJobPayload {
169    pub fn to_payload_ref(&self) -> Result<String> {
170        serde_json::to_string(self).map_err(|err| {
171            StasisError::PortFailure(format!("failed to encode prompt payload: {err}"))
172        })
173    }
174}
175
176#[derive(Clone, Debug, Serialize, Deserialize)]
177pub struct SequentialStageJobPayload {
178    pub stage_id: String,
179    pub user_prompt_template: String,
180    pub system_prompt: Option<String>,
181    pub policy_profile: Option<String>,
182    pub model_hint: Option<String>,
183    pub reasoning_effort: Option<String>,
184}
185
186#[derive(Clone, Debug, Serialize, Deserialize)]
187pub struct SequentialPatternJobPayload {
188    pub thread_id: Option<String>,
189    pub initial_user_prompt: String,
190    pub policy_profile: Option<String>,
191    pub model_hint: Option<String>,
192    pub reasoning_effort: Option<String>,
193    pub stages: Vec<SequentialStageJobPayload>,
194}
195
196impl SequentialPatternJobPayload {
197    pub fn to_payload_ref(&self) -> Result<String> {
198        serde_json::to_string(self).map_err(|err| {
199            StasisError::PortFailure(format!(
200                "failed to encode sequential-pattern payload: {err}"
201            ))
202        })
203    }
204}
205
206#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
207#[serde(rename_all = "snake_case")]
208pub enum ConcurrentBranchExecutionMode {
209    #[default]
210    Prompt,
211    ToolLoop,
212}
213
214#[derive(Clone, Debug, Serialize, Deserialize)]
215pub struct ConcurrentBranchJobPayload {
216    pub branch_id: String,
217    pub user_prompt_template: String,
218    pub system_prompt: Option<String>,
219    pub policy_profile: Option<String>,
220    pub model_hint: Option<String>,
221    pub reasoning_effort: Option<String>,
222    #[serde(default)]
223    pub execution_mode: ConcurrentBranchExecutionMode,
224    pub tool_name: Option<String>,
225    pub tool_input: Option<Value>,
226    pub tool_call_mode: Option<AgentToolCallMode>,
227    pub memory_policy: Option<MemoryPolicyPayload>,
228}
229
230#[derive(Clone, Debug, Serialize, Deserialize)]
231pub struct ConcurrentPatternJobPayload {
232    pub thread_id: Option<String>,
233    pub initial_user_prompt: String,
234    pub policy_profile: Option<String>,
235    pub model_hint: Option<String>,
236    pub reasoning_effort: Option<String>,
237    pub merge_strategy: Option<String>,
238    pub tool_call_mode: Option<AgentToolCallMode>,
239    pub memory_policy: Option<MemoryPolicyPayload>,
240    pub branches: Vec<ConcurrentBranchJobPayload>,
241}
242
243impl ConcurrentBranchJobPayload {
244    pub fn prompt(
245        branch_id: impl Into<String>,
246        user_prompt_template: impl Into<String>,
247    ) -> Self {
248        Self {
249            branch_id: branch_id.into(),
250            user_prompt_template: user_prompt_template.into(),
251            system_prompt: None,
252            policy_profile: None,
253            model_hint: None,
254            reasoning_effort: None,
255            execution_mode: ConcurrentBranchExecutionMode::Prompt,
256            tool_name: None,
257            tool_input: None,
258            tool_call_mode: None,
259            memory_policy: None,
260        }
261    }
262
263    pub fn tool_loop(
264        branch_id: impl Into<String>,
265        user_prompt_template: impl Into<String>,
266        tool_name: impl Into<String>,
267        tool_input: Option<Value>,
268    ) -> Self {
269        Self {
270            branch_id: branch_id.into(),
271            user_prompt_template: user_prompt_template.into(),
272            system_prompt: None,
273            policy_profile: None,
274            model_hint: None,
275            reasoning_effort: None,
276            execution_mode: ConcurrentBranchExecutionMode::ToolLoop,
277            tool_name: Some(tool_name.into()),
278            tool_input,
279            tool_call_mode: None,
280            memory_policy: None,
281        }
282    }
283}
284
285impl ConcurrentPatternJobPayload {
286    pub fn to_payload_ref(&self) -> Result<String> {
287        serde_json::to_string(self).map_err(|err| {
288            StasisError::PortFailure(format!(
289                "failed to encode concurrent-pattern payload: {err}"
290            ))
291        })
292    }
293}
294
295#[derive(Clone, Debug, Serialize, Deserialize)]
296pub struct HandoffTurnJobPayload {
297    pub actor_id: String,
298    pub user_prompt_template: String,
299    pub system_prompt: Option<String>,
300    pub policy_profile: Option<String>,
301    pub model_hint: Option<String>,
302    pub reasoning_effort: Option<String>,
303}
304
305#[derive(Clone, Debug, Serialize, Deserialize)]
306pub struct HandoffPatternJobPayload {
307    pub thread_id: Option<String>,
308    pub initial_user_prompt: String,
309    pub policy_profile: Option<String>,
310    pub model_hint: Option<String>,
311    pub reasoning_effort: Option<String>,
312    pub turns: Vec<HandoffTurnJobPayload>,
313}
314
315impl HandoffPatternJobPayload {
316    pub fn to_payload_ref(&self) -> Result<String> {
317        serde_json::to_string(self).map_err(|err| {
318            StasisError::PortFailure(format!("failed to encode handoff-pattern payload: {err}"))
319        })
320    }
321}
322
323#[derive(Clone, Debug, Serialize, Deserialize)]
324pub struct OrchestratorRouteJobPayload {
325    pub route_id: String,
326    pub selector_keywords: Vec<String>,
327    pub user_prompt_template: String,
328    pub system_prompt: Option<String>,
329    pub policy_profile: Option<String>,
330    pub model_hint: Option<String>,
331    pub reasoning_effort: Option<String>,
332}
333
334#[derive(Clone, Debug, Serialize, Deserialize)]
335pub struct OrchestratorPatternJobPayload {
336    pub thread_id: Option<String>,
337    pub initial_user_prompt: String,
338    pub policy_profile: Option<String>,
339    pub model_hint: Option<String>,
340    pub reasoning_effort: Option<String>,
341    pub routes: Vec<OrchestratorRouteJobPayload>,
342}
343
344impl OrchestratorPatternJobPayload {
345    pub fn to_payload_ref(&self) -> Result<String> {
346        serde_json::to_string(self).map_err(|err| {
347            StasisError::PortFailure(format!(
348                "failed to encode orchestrator-pattern payload: {err}"
349            ))
350        })
351    }
352}
353
354#[derive(Clone, Debug, Serialize, Deserialize)]
355pub struct MemoryRecallJobPayload {
356    pub memory_policy: Option<MemoryPolicyPayload>,
357}
358
359impl MemoryRecallJobPayload {
360    pub fn to_payload_ref(&self) -> Result<String> {
361        serde_json::to_string(self).map_err(|err| {
362            StasisError::PortFailure(format!("failed to encode memory-recall payload: {err}"))
363        })
364    }
365}
366
367#[derive(Clone, Debug, Serialize, Deserialize)]
368#[serde(rename_all = "camelCase")]
369pub struct MemoryFindJobPayload {
370    pub tenant_id: Option<String>,
371    pub session_ids: Option<Vec<String>>,
372    pub tiers: Option<Vec<String>>,
373    pub from_utc: Option<DateTime<Utc>>,
374    pub to_utc: Option<DateTime<Utc>>,
375    pub limit: Option<usize>,
376    pub cursor: Option<String>,
377    pub text_contains: Option<String>,
378    pub tags_contains: Option<Vec<String>>,
379    pub has_tag: Option<String>,
380    pub indexed_tags: Option<Vec<String>>,
381    pub tag_prefix: Option<String>,
382    pub has_semantic_links: Option<bool>,
383    pub link_rel: Option<String>,
384    pub link_target: Option<String>,
385    pub links_to_ref: Option<String>,
386    pub sort_field: Option<String>,
387    pub sort_direction: Option<String>,
388}
389
390impl MemoryFindJobPayload {
391    pub fn to_payload_ref(&self) -> Result<String> {
392        serde_json::to_string(self).map_err(|err| {
393            StasisError::PortFailure(format!("failed to encode memory-find payload: {err}"))
394        })
395    }
396}
397
398#[derive(Clone, Debug, Serialize, Deserialize)]
399#[serde(rename_all = "camelCase")]
400pub struct MemoryAggregateJobPayload {
401    pub session_ids: Option<Vec<String>>,
402    pub tiers: Option<Vec<String>>,
403    pub from_utc: Option<DateTime<Utc>>,
404    pub to_utc: Option<DateTime<Utc>>,
405    pub max_groups: Option<usize>,
406    pub max_nodes: Option<usize>,
407}
408
409impl MemoryAggregateJobPayload {
410    pub fn to_payload_ref(&self) -> Result<String> {
411        serde_json::to_string(self).map_err(|err| {
412            StasisError::PortFailure(format!("failed to encode memory-aggregate payload: {err}"))
413        })
414    }
415}
416
417#[derive(Clone, Debug, Serialize, Deserialize)]
418#[serde(rename_all = "snake_case")]
419pub enum MemoryTransformOperationPayload {
420    EmbedBackfill,
421    ReindexEmbeddings,
422    EmbedTagBackfill,
423    ReindexTagEmbeddings,
424}
425
426#[derive(Clone, Debug, Serialize, Deserialize)]
427#[serde(rename_all = "camelCase")]
428pub struct MemoryTransformJobPayload {
429    pub session_ids: Option<Vec<String>>,
430    pub tiers: Option<Vec<String>>,
431    pub from_utc: Option<DateTime<Utc>>,
432    pub to_utc: Option<DateTime<Utc>>,
433    pub operation: Option<MemoryTransformOperationPayload>,
434    pub dry_run: Option<bool>,
435    pub batch_size: Option<usize>,
436    pub max_nodes: Option<usize>,
437    pub provider_id: Option<String>,
438    pub model: Option<String>,
439}
440
441impl MemoryTransformJobPayload {
442    pub fn to_payload_ref(&self) -> Result<String> {
443        serde_json::to_string(self).map_err(|err| {
444            StasisError::PortFailure(format!("failed to encode memory-transform payload: {err}"))
445        })
446    }
447}
448
449#[derive(Clone, Debug, Serialize, Deserialize)]
450#[serde(rename_all = "camelCase")]
451pub struct MemoryRollupJobPayload {
452    pub session_ids: Option<Vec<String>>,
453    pub tiers: Option<Vec<String>>,
454    pub from_utc: Option<DateTime<Utc>>,
455    pub to_utc: Option<DateTime<Utc>>,
456    pub max_days: Option<usize>,
457    pub max_nodes: Option<usize>,
458}
459
460impl MemoryRollupJobPayload {
461    pub fn to_payload_ref(&self) -> Result<String> {
462        serde_json::to_string(self).map_err(|err| {
463            StasisError::PortFailure(format!("failed to encode memory-rollup payload: {err}"))
464        })
465    }
466}
467
468#[derive(Clone, Debug, Default, Serialize, Deserialize)]
469pub struct MemorySchemaJobPayload {}
470
471impl MemorySchemaJobPayload {
472    pub fn to_payload_ref(&self) -> Result<String> {
473        serde_json::to_string(self).map_err(|err| {
474            StasisError::PortFailure(format!("failed to encode memory-schema payload: {err}"))
475        })
476    }
477}
478
479#[derive(Clone, Debug, Serialize, Deserialize)]
480#[serde(rename_all = "snake_case")]
481pub enum MemoryEvictModePayload {
482    BySyncKeys,
483    ByNodeIds,
484    ByFilter,
485    PurgeSession,
486}
487
488#[derive(Clone, Debug, Serialize, Deserialize)]
489#[serde(rename_all = "camelCase")]
490pub struct MemoryEvictJobPayload {
491    pub mode: Option<MemoryEvictModePayload>,
492    pub tenant_id: Option<String>,
493    pub session_ids: Option<Vec<String>>,
494    pub tiers: Option<Vec<String>>,
495    pub from_utc: Option<DateTime<Utc>>,
496    pub to_utc: Option<DateTime<Utc>>,
497    #[serde(default)]
498    pub filter: MemoryFilterPayload,
499    pub sync_keys: Option<Vec<String>>,
500    pub node_ids: Option<Vec<String>>,
501    pub dry_run: Option<bool>,
502    pub force: Option<bool>,
503    pub max_nodes: Option<usize>,
504    pub include_calibration: Option<bool>,
505    pub include_checkpoints: Option<bool>,
506}
507
508impl MemoryEvictJobPayload {
509    pub fn to_payload_ref(&self) -> Result<String> {
510        serde_json::to_string(self).map_err(|err| {
511            StasisError::PortFailure(format!("failed to encode memory-evict payload: {err}"))
512        })
513    }
514}
515
516#[derive(Clone, Debug, Serialize, Deserialize)]
517#[serde(rename_all = "camelCase")]
518pub struct MemoryGraphJobPayload {
519    pub tenant_id: Option<String>,
520    pub session_ids: Option<Vec<String>>,
521    pub tiers: Option<Vec<String>>,
522    pub from_utc: Option<DateTime<Utc>>,
523    pub to_utc: Option<DateTime<Utc>>,
524    #[serde(default)]
525    pub filter: MemoryFilterPayload,
526    pub include_lineage: Option<bool>,
527    pub include_semantic: Option<bool>,
528    pub include_session_topology: Option<bool>,
529    pub rel: Option<String>,
530    pub target_prefix: Option<String>,
531    pub limit: Option<usize>,
532}
533
534impl MemoryGraphJobPayload {
535    pub fn to_payload_ref(&self) -> Result<String> {
536        serde_json::to_string(self).map_err(|err| {
537            StasisError::PortFailure(format!("failed to encode memory-graph payload: {err}"))
538        })
539    }
540}