1use std::collections::{BTreeMap, BTreeSet};
7
8use chrono::{DateTime, Utc};
9use serde::{Deserialize, Serialize};
10use serde_json::Value;
11use starweaver_core::{RunId, SessionId};
12use starweaver_stream::{DisplayMessage, ReplayCursor};
13
14use crate::{InputPart, RunStatus, SessionStatus};
15
16pub const LOCAL_SESSION_NAMESPACE: &str = "local";
18
19#[derive(Clone, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
21pub struct ManagedSessionTarget {
22 pub namespace_id: String,
24 pub session_id: SessionId,
26}
27
28impl ManagedSessionTarget {
29 #[must_use]
31 pub fn new(namespace_id: impl Into<String>, session_id: SessionId) -> Self {
32 Self {
33 namespace_id: namespace_id.into(),
34 session_id,
35 }
36 }
37}
38
39#[derive(Clone, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
41pub struct ManagedRunTarget {
42 pub namespace_id: String,
44 pub session_id: SessionId,
46 pub run_id: RunId,
48}
49
50impl ManagedRunTarget {
51 #[must_use]
53 pub fn new(namespace_id: impl Into<String>, session_id: SessionId, run_id: RunId) -> Self {
54 Self {
55 namespace_id: namespace_id.into(),
56 session_id,
57 run_id,
58 }
59 }
60}
61
62#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
64#[serde(rename_all = "snake_case")]
65pub enum AgentSessionOperation {
66 Read,
68 Search,
70 Create,
72 Update,
74 Control,
76 Delete,
78}
79
80#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
82pub struct AgentSessionScope {
83 pub namespace_id: String,
85 #[serde(default, skip_serializing_if = "Option::is_none")]
87 pub owner_id: Option<String>,
88 pub source_product: String,
90 #[serde(default, skip_serializing_if = "Option::is_none")]
92 pub source_session_id: Option<SessionId>,
93 #[serde(default, skip_serializing_if = "Option::is_none")]
95 pub source_run_id: Option<RunId>,
96 #[serde(default)]
98 pub operations: BTreeSet<AgentSessionOperation>,
99 #[serde(default)]
101 pub allowed_session_ids: BTreeSet<SessionId>,
102 #[serde(default = "default_true")]
104 pub allow_self_query: bool,
105 #[serde(default)]
107 pub allow_self_control: bool,
108 pub policy_fingerprint: String,
110 #[serde(default, skip_serializing_if = "Option::is_none")]
112 pub deadline: Option<DateTime<Utc>>,
113 #[serde(default = "default_page_limit")]
115 pub max_page_size: u32,
116}
117
118const fn default_true() -> bool {
119 true
120}
121
122const fn default_page_limit() -> u32 {
123 50
124}
125
126impl AgentSessionScope {
127 #[must_use]
129 pub fn allows(&self, operation: AgentSessionOperation) -> bool {
130 self.operations.contains(&operation)
131 }
132
133 #[must_use]
135 pub fn allows_session(&self, session_id: &SessionId) -> bool {
136 self.allowed_session_ids.is_empty() || self.allowed_session_ids.contains(session_id)
137 }
138
139 #[must_use]
141 pub fn is_self_run(&self, target: &ManagedRunTarget) -> bool {
142 self.namespace_id == target.namespace_id
143 && self.source_session_id.as_ref() == Some(&target.session_id)
144 && self.source_run_id.as_ref() == Some(&target.run_id)
145 }
146
147 #[must_use]
149 pub fn is_self_session(&self, target: &ManagedSessionTarget) -> bool {
150 self.namespace_id == target.namespace_id
151 && self.source_session_id.as_ref() == Some(&target.session_id)
152 }
153}
154
155#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
157#[serde(tag = "state", rename_all = "snake_case")]
158pub enum SessionDeletionFence {
159 #[default]
161 Stable,
162 Deleting {
164 fence_id: String,
166 expected_revision: u64,
168 requested_by: String,
170 started_at: DateTime<Utc>,
172 },
173 Deleted {
175 fence_id: String,
177 deleted_at: DateTime<Utc>,
179 },
180}
181
182impl SessionDeletionFence {
183 #[must_use]
185 pub const fn blocks_continuation(&self) -> bool {
186 !matches!(self, Self::Stable)
187 }
188}
189
190#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
192pub struct SessionContinuationFence {
193 pub target: ManagedSessionTarget,
195 pub revision: u64,
197 pub continuation_allowed: bool,
199 #[serde(default, skip_serializing_if = "Option::is_none")]
201 pub fence_id: Option<String>,
202}
203
204#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
206pub struct AgentSessionListQuery {
207 #[serde(default, skip_serializing_if = "Option::is_none")]
209 pub status: Option<SessionStatus>,
210 #[serde(default, skip_serializing_if = "Option::is_none")]
212 pub profile: Option<String>,
213 #[serde(default, skip_serializing_if = "Option::is_none")]
215 pub workspace: Option<String>,
216 #[serde(default = "default_page_limit")]
218 pub limit: u32,
219 #[serde(default, skip_serializing_if = "Option::is_none")]
221 pub page_token: Option<String>,
222}
223
224#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
226pub struct AgentSessionInclude {
227 #[serde(default)]
229 pub recent_runs: bool,
230 #[serde(default)]
232 pub trace: bool,
233}
234
235#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
237pub struct AgentRunListQuery {
238 #[serde(default = "default_page_limit")]
240 pub limit: u32,
241 #[serde(default, skip_serializing_if = "Option::is_none")]
243 pub page_token: Option<String>,
244}
245
246impl Default for AgentRunListQuery {
247 fn default() -> Self {
248 Self {
249 limit: default_page_limit(),
250 page_token: None,
251 }
252 }
253}
254
255#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
257pub struct AgentReplayQuery {
258 #[serde(default, skip_serializing_if = "Option::is_none")]
260 pub after: Option<ReplayCursor>,
261 #[serde(default = "default_page_limit")]
263 pub limit: u32,
264}
265
266#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
268pub struct AgentSessionView {
269 pub target: ManagedSessionTarget,
271 #[serde(default, skip_serializing_if = "Option::is_none")]
273 pub title: Option<String>,
274 pub status: SessionStatus,
276 #[serde(default, skip_serializing_if = "Option::is_none")]
278 pub profile: Option<String>,
279 #[serde(default, skip_serializing_if = "Option::is_none")]
281 pub workspace: Option<String>,
282 pub revision: u64,
284 #[serde(default, skip_serializing_if = "Option::is_none")]
286 pub head_run_id: Option<RunId>,
287 #[serde(default, skip_serializing_if = "Option::is_none")]
289 pub active_run_id: Option<RunId>,
290 pub resumable: bool,
292 pub controllable: bool,
294 #[serde(default, skip_serializing_if = "Vec::is_empty")]
296 pub recent_runs: Vec<AgentRunView>,
297 pub created_at: DateTime<Utc>,
299 pub updated_at: DateTime<Utc>,
301}
302
303#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
305pub struct AgentRunView {
306 pub target: ManagedRunTarget,
308 pub status: RunStatus,
310 pub sequence_no: usize,
312 #[serde(default, skip_serializing_if = "Option::is_none")]
314 pub input_preview: Option<String>,
315 #[serde(default, skip_serializing_if = "Option::is_none")]
317 pub output_preview: Option<String>,
318 #[serde(default, skip_serializing_if = "Option::is_none")]
320 pub error_category: Option<String>,
321 pub controllable: bool,
323 pub created_at: DateTime<Utc>,
325 pub updated_at: DateTime<Utc>,
327}
328
329#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
331pub struct AgentSessionPage {
332 pub sessions: Vec<AgentSessionView>,
334 #[serde(default, skip_serializing_if = "Option::is_none")]
336 pub next_page_token: Option<String>,
337}
338
339#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
341pub struct AgentRunPage {
342 pub runs: Vec<AgentRunView>,
344 #[serde(default, skip_serializing_if = "Option::is_none")]
346 pub next_page_token: Option<String>,
347}
348
349#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
351pub struct AgentDisplayPage {
352 pub messages: Vec<DisplayMessage>,
354 #[serde(default, skip_serializing_if = "Option::is_none")]
356 pub next_cursor: Option<ReplayCursor>,
357 #[serde(default = "untrusted_evidence_label")]
359 pub trust: String,
360}
361
362fn untrusted_evidence_label() -> String {
363 "untrusted_historical_evidence".to_string()
364}
365
366#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
368pub struct CreateManagedSession {
369 #[serde(default, skip_serializing_if = "Option::is_none")]
371 pub title: Option<String>,
372 #[serde(default, skip_serializing_if = "Option::is_none")]
374 pub profile: Option<String>,
375 #[serde(default, skip_serializing_if = "Option::is_none")]
377 pub workspace: Option<String>,
378 #[serde(default)]
380 pub metadata: BTreeMap<String, Value>,
381 pub idempotency_key: String,
383}
384
385#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
387pub struct ManagedSessionPatch {
388 #[serde(default, skip_serializing_if = "Option::is_none")]
390 pub title: Option<Option<String>>,
391 #[serde(default, skip_serializing_if = "Option::is_none")]
393 pub profile: Option<Option<String>>,
394 #[serde(default, skip_serializing_if = "Option::is_none")]
396 pub archived: Option<bool>,
397 #[serde(default)]
399 pub metadata: BTreeMap<String, Value>,
400}
401
402#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
404pub struct UpdateManagedSession {
405 pub session_id: SessionId,
407 pub expected_revision: u64,
409 pub patch: ManagedSessionPatch,
411 pub idempotency_key: String,
413}
414
415#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
417pub struct DeleteManagedSession {
418 pub session_id: SessionId,
420 pub expected_revision: u64,
422 pub idempotency_key: String,
424 #[serde(default, skip_serializing_if = "Option::is_none")]
426 pub approval_receipt_id: Option<String>,
427}
428
429#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
431pub struct StartManagedRun {
432 pub session_id: SessionId,
434 pub input: Vec<InputPart>,
436 #[serde(default, skip_serializing_if = "Option::is_none")]
438 pub profile: Option<String>,
439 #[serde(default)]
441 pub environment_refs: Vec<String>,
442 pub idempotency_key: String,
444}
445
446#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
448pub struct SteerManagedRun {
449 pub target: ManagedRunTarget,
451 pub steering_id: String,
453 pub text: String,
455 #[serde(default, skip_serializing_if = "Option::is_none")]
457 pub idempotency_key: Option<String>,
458}
459
460#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
462pub struct InterruptManagedRun {
463 pub target: ManagedRunTarget,
465 pub operation_id: String,
467 #[serde(default, skip_serializing_if = "Option::is_none")]
469 pub reason_category: Option<String>,
470 #[serde(default, skip_serializing_if = "Option::is_none")]
472 pub idempotency_key: Option<String>,
473}
474
475#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
477pub struct SessionMutationReceipt {
478 pub receipt_id: String,
480 pub session: AgentSessionView,
482 pub idempotent_replay: bool,
484}
485
486#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
488pub struct RunStartReceipt {
489 pub receipt_id: String,
491 pub target: ManagedRunTarget,
493 pub status: RunStatus,
495 pub fencing_generation: u64,
497 pub idempotent_replay: bool,
499}
500
501#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
503pub struct RunControlReceipt {
504 pub receipt_id: String,
506 pub target: ManagedRunTarget,
508 pub operation_id: String,
510 pub fencing_generation: u64,
512 pub accepted: bool,
514 pub idempotent_replay: bool,
516 pub created_at: DateTime<Utc>,
518}
519
520#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
522pub struct RunAdmissionLease {
523 pub target: ManagedRunTarget,
525 pub admission_id: String,
527 pub host_instance_id: String,
529 pub fencing_generation: u64,
531 pub lease_expires_at: DateTime<Utc>,
533 pub heartbeat_at: DateTime<Utc>,
535 pub command_fingerprint: String,
537 pub idempotency_key: String,
539}
540
541impl starweaver_core::VersionedRecord for RunAdmissionLease {
542 const SCHEMA: &'static str = "starweaver.session.run_admission_lease";
543}
544
545impl RunAdmissionLease {
546 #[must_use]
548 pub fn expired_at(&self, now: DateTime<Utc>) -> bool {
549 self.lease_expires_at <= now
550 }
551}
552
553#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
555pub struct AcquireRunAdmission {
556 pub run: crate::RunRecord,
558 pub namespace_id: String,
560 pub host_instance_id: String,
562 pub admission_id: String,
564 pub lease_expires_at: DateTime<Utc>,
566 pub idempotency_key: String,
568 pub command_fingerprint: String,
570}
571
572#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
574pub struct RunAdmissionReceipt {
575 pub run: crate::RunRecord,
577 pub lease: RunAdmissionLease,
579 pub idempotent_replay: bool,
581}
582
583impl starweaver_core::VersionedRecord for RunAdmissionReceipt {
584 const SCHEMA: &'static str = "starweaver.session.run_admission_receipt";
585}
586
587#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
589pub struct DurableControlReceipt {
590 pub receipt_id: String,
592 pub target: ManagedRunTarget,
594 pub operation_id: String,
596 pub operation: String,
598 pub idempotency_key: String,
600 pub command_fingerprint: String,
602 pub fencing_generation: u64,
604 pub state: String,
606 pub created_at: DateTime<Utc>,
608}
609
610impl starweaver_core::VersionedRecord for DurableControlReceipt {
611 const SCHEMA: &'static str = "starweaver.session.durable_control_receipt";
612}
613
614#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
616#[serde(rename_all = "snake_case")]
617pub enum AgentSessionQueryErrorCode {
618 InvalidQuery,
620 NotFound,
622 Unsupported,
624 Unavailable,
626 PermissionDenied,
628 InvalidCursor,
630 Failed,
632}
633
634#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
636pub struct AgentSessionQueryError {
637 pub code: AgentSessionQueryErrorCode,
639 pub message: String,
641}
642
643impl std::fmt::Display for AgentSessionQueryError {
644 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
645 write!(formatter, "{:?}: {}", self.code, self.message)
646 }
647}
648
649impl std::error::Error for AgentSessionQueryError {}
650
651#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
653#[serde(rename_all = "snake_case")]
654pub enum AgentSessionControlErrorCode {
655 InvalidCommand,
657 NotFound,
659 PermissionDenied,
661 ApprovalRequired,
663 Conflict,
665 IdempotencyConflict,
667 RunConflict,
669 NotActive,
671 Terminal,
673 StaleActive,
675 QuotaExceeded,
677 Unavailable,
679 Failed,
681}
682
683#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
685pub struct AgentSessionControlError {
686 pub code: AgentSessionControlErrorCode,
688 pub message: String,
690 #[serde(default, skip_serializing_if = "Option::is_none")]
692 pub current_revision: Option<u64>,
693}
694
695impl std::fmt::Display for AgentSessionControlError {
696 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
697 write!(formatter, "{:?}: {}", self.code, self.message)
698 }
699}
700
701impl std::error::Error for AgentSessionControlError {}