1use chrono::{DateTime, Utc};
4use serde::{Deserialize, Deserializer, Serialize, Serializer, de::Error as _};
5use serde_json::Value;
6use starweaver_context::ResumableState;
7use starweaver_core::{
8 CheckpointId, ConversationId, Metadata, RunId, RunLifecycle, SessionId, TaskId, TraceContext,
9};
10use starweaver_stream::{ReplayCursor, ReplayCursorFamily, ReplayScope};
11
12use crate::{input::InputPart, management::SessionDeletionFence};
13
14#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
16#[serde(rename_all = "snake_case")]
17pub enum SessionStatus {
18 #[default]
20 Active,
21 Archived,
23 Failed,
25 Deleted,
27}
28
29#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
31pub struct DurableRunStatus(Option<RunLifecycle>);
32
33pub type RunStatus = DurableRunStatus;
35
36#[allow(non_upper_case_globals)]
37impl DurableRunStatus {
38 pub const Queued: Self = Self(None);
40 pub const Starting: Self = Self(Some(RunLifecycle::Starting));
42 pub const Running: Self = Self(Some(RunLifecycle::Running));
44 pub const Waiting: Self = Self(Some(RunLifecycle::Waiting));
46 pub const Completed: Self = Self(Some(RunLifecycle::Completed));
48 pub const Failed: Self = Self(Some(RunLifecycle::Failed));
50 pub const Cancelled: Self = Self(Some(RunLifecycle::Cancelled));
52
53 #[must_use]
55 pub const fn lifecycle(self) -> Option<RunLifecycle> {
56 self.0
57 }
58
59 #[must_use]
61 pub const fn as_str(self) -> &'static str {
62 match self.0 {
63 None => "queued",
64 Some(lifecycle) => lifecycle.as_str(),
65 }
66 }
67
68 #[must_use]
70 pub const fn is_active(self) -> bool {
71 matches!(
72 self.0,
73 None | Some(RunLifecycle::Starting | RunLifecycle::Running | RunLifecycle::Waiting)
74 )
75 }
76
77 #[must_use]
79 pub const fn is_terminal(self) -> bool {
80 match self.0 {
81 Some(lifecycle) => lifecycle.is_terminal(),
82 None => false,
83 }
84 }
85}
86
87impl From<RunLifecycle> for DurableRunStatus {
88 fn from(value: RunLifecycle) -> Self {
89 Self(Some(value))
90 }
91}
92
93impl TryFrom<DurableRunStatus> for RunLifecycle {
94 type Error = QueuedRunStatus;
95
96 fn try_from(value: DurableRunStatus) -> Result<Self, Self::Error> {
97 value.0.ok_or(QueuedRunStatus)
98 }
99}
100
101impl Serialize for DurableRunStatus {
102 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
103 where
104 S: Serializer,
105 {
106 serializer.serialize_str(self.as_str())
107 }
108}
109
110impl<'de> Deserialize<'de> for DurableRunStatus {
111 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
112 where
113 D: Deserializer<'de>,
114 {
115 match String::deserialize(deserializer)?.as_str() {
116 "queued" => Ok(Self::Queued),
117 "starting" => Ok(Self::Starting),
118 "running" => Ok(Self::Running),
119 "waiting" => Ok(Self::Waiting),
120 "completed" => Ok(Self::Completed),
121 "failed" => Ok(Self::Failed),
122 "cancelled" => Ok(Self::Cancelled),
123 other => Err(D::Error::custom(format!("unknown run status: {other}"))),
124 }
125 }
126}
127
128#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
130pub struct RunTerminalError {
131 pub code: String,
133 pub message: String,
135}
136
137impl RunTerminalError {
138 #[must_use]
140 pub fn new(code: impl Into<String>, message: impl Into<String>) -> Self {
141 Self {
142 code: code.into(),
143 message: message.into(),
144 }
145 }
146}
147
148#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
150pub struct RunTerminalProjection {
151 pub status: RunStatus,
153 #[serde(default, skip_serializing_if = "Option::is_none")]
155 pub output_preview: Option<String>,
156 #[serde(default, skip_serializing_if = "Option::is_none")]
158 pub error: Option<RunTerminalError>,
159}
160
161impl RunTerminalProjection {
162 pub fn try_new(
169 status: RunStatus,
170 output_preview: Option<String>,
171 error: Option<RunTerminalError>,
172 ) -> Result<Self, RunTerminalProjectionError> {
173 let projection = Self {
174 status,
175 output_preview,
176 error,
177 };
178 projection.validate()?;
179 Ok(projection)
180 }
181
182 #[must_use]
184 pub const fn completed(output_preview: Option<String>) -> Self {
185 Self {
186 status: RunStatus::Completed,
187 output_preview,
188 error: None,
189 }
190 }
191
192 #[must_use]
194 pub const fn failed(error: RunTerminalError) -> Self {
195 Self {
196 status: RunStatus::Failed,
197 output_preview: None,
198 error: Some(error),
199 }
200 }
201
202 #[must_use]
204 pub const fn cancelled(error: Option<RunTerminalError>) -> Self {
205 Self {
206 status: RunStatus::Cancelled,
207 output_preview: None,
208 error,
209 }
210 }
211
212 pub fn validate(&self) -> Result<(), RunTerminalProjectionError> {
221 if !self.status.is_terminal() {
222 return Err(RunTerminalProjectionError::NonTerminalStatus(self.status));
223 }
224 if self.status == RunStatus::Failed && self.error.is_none() {
225 return Err(RunTerminalProjectionError::MissingFailureDiagnostic);
226 }
227 if self.status == RunStatus::Completed && self.error.is_some() {
228 return Err(RunTerminalProjectionError::UnexpectedSuccessDiagnostic);
229 }
230 if let Some(error) = self.error.as_ref() {
231 if error.code.is_empty() {
232 return Err(RunTerminalProjectionError::EmptyDiagnosticCode);
233 }
234 if error.message.is_empty() {
235 return Err(RunTerminalProjectionError::EmptyDiagnosticMessage);
236 }
237 }
238 Ok(())
239 }
240
241 #[must_use]
243 pub fn matches(&self, run: &RunRecord) -> bool {
244 (run.status, &run.output_preview, &run.terminal_error)
245 == (self.status, &self.output_preview, &self.error)
246 }
247
248 pub fn apply_to(&self, run: &mut RunRecord) {
250 run.status = self.status;
251 run.output_preview.clone_from(&self.output_preview);
252 run.terminal_error.clone_from(&self.error);
253 }
254}
255
256#[derive(Clone, Debug, Eq, PartialEq)]
258pub enum RunTerminalProjectionError {
259 NonTerminalStatus(RunStatus),
261 MissingFailureDiagnostic,
263 UnexpectedSuccessDiagnostic,
265 UnexpectedNonTerminalDiagnostic(RunStatus),
267 EmptyDiagnosticCode,
269 EmptyDiagnosticMessage,
271}
272
273impl std::fmt::Display for RunTerminalProjectionError {
274 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
275 match self {
276 Self::NonTerminalStatus(status) => {
277 write!(formatter, "run status {} is not terminal", status.as_str())
278 }
279 Self::MissingFailureDiagnostic => {
280 formatter.write_str("failed run requires a terminal diagnostic")
281 }
282 Self::UnexpectedSuccessDiagnostic => {
283 formatter.write_str("completed run cannot carry a terminal diagnostic")
284 }
285 Self::UnexpectedNonTerminalDiagnostic(status) => write!(
286 formatter,
287 "non-terminal run status {} cannot carry a terminal diagnostic",
288 status.as_str()
289 ),
290 Self::EmptyDiagnosticCode => formatter.write_str("terminal diagnostic code is empty"),
291 Self::EmptyDiagnosticMessage => {
292 formatter.write_str("terminal diagnostic message is empty")
293 }
294 }
295 }
296}
297
298impl std::error::Error for RunTerminalProjectionError {}
299
300#[derive(Clone, Copy, Debug, Eq, PartialEq)]
302pub struct QueuedRunStatus;
303
304impl std::fmt::Display for QueuedRunStatus {
305 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
306 formatter.write_str("queued run has no runtime lifecycle")
307 }
308}
309
310impl std::error::Error for QueuedRunStatus {}
311
312#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
314#[serde(rename_all = "snake_case")]
315pub enum ExecutionStatus {
316 Pending,
318 Running,
320 Waiting,
322 Completed,
324 Failed,
326 Cancelled,
328}
329
330#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
332pub struct EnvironmentStateRef {
333 pub provider: String,
335 pub reference: String,
337 #[serde(default, skip_serializing_if = "Option::is_none")]
339 pub revision: Option<String>,
340 #[serde(default, skip_serializing_if = "Metadata::is_empty")]
342 pub metadata: Metadata,
343}
344
345#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
347pub struct CheckpointRef {
348 pub checkpoint_id: CheckpointId,
350 pub run_id: RunId,
352 pub sequence: usize,
354 pub node: String,
356 #[serde(default, skip_serializing_if = "Option::is_none")]
358 pub storage_ref: Option<String>,
359 #[serde(default, skip_serializing_if = "Option::is_none")]
361 pub stream_cursor: Option<usize>,
362 pub created_at: DateTime<Utc>,
364 #[serde(default, skip_serializing_if = "Metadata::is_empty")]
366 pub metadata: Metadata,
367}
368
369#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
371pub struct StreamCursorRef {
372 pub position: ReplayCursor,
374 pub created_at: DateTime<Utc>,
376 #[serde(default, skip_serializing_if = "Metadata::is_empty")]
378 pub metadata: Metadata,
379}
380
381impl StreamCursorRef {
382 #[must_use]
384 pub fn new(position: ReplayCursor) -> Self {
385 Self {
386 position,
387 created_at: Utc::now(),
388 metadata: Metadata::default(),
389 }
390 }
391
392 #[must_use]
394 pub const fn family(&self) -> ReplayCursorFamily {
395 self.position.family
396 }
397
398 #[must_use]
400 pub const fn scope(&self) -> &ReplayScope {
401 &self.position.scope
402 }
403
404 #[must_use]
406 pub const fn sequence(&self) -> usize {
407 self.position.sequence
408 }
409
410 #[must_use]
412 pub fn same_stream(&self, other: &Self) -> bool {
413 self.family() == other.family() && self.scope() == other.scope()
414 }
415
416 pub fn validate_for_run(&self, run_id: &RunId) -> Result<(), StreamCursorRefError> {
422 let expected = ReplayScope::run(run_id.as_str());
423 if self.scope() != &expected {
424 return Err(StreamCursorRefError::WrongScope {
425 expected: expected.as_str().to_string(),
426 actual: self.scope().as_str().to_string(),
427 });
428 }
429 Ok(())
430 }
431
432 pub fn validate_progression(&self, current: &Self) -> Result<(), StreamCursorRefError> {
438 if self.same_stream(current) && self.sequence() < current.sequence() {
439 return Err(StreamCursorRefError::SequenceRegression {
440 family: self.family(),
441 scope: self.scope().as_str().to_string(),
442 current: current.sequence(),
443 proposed: self.sequence(),
444 });
445 }
446 Ok(())
447 }
448}
449
450#[derive(Clone, Debug, Eq, PartialEq)]
452pub enum StreamCursorRefError {
453 WrongScope {
455 expected: String,
457 actual: String,
459 },
460 SequenceRegression {
462 family: ReplayCursorFamily,
464 scope: String,
466 current: usize,
468 proposed: usize,
470 },
471}
472
473impl std::fmt::Display for StreamCursorRefError {
474 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
475 match self {
476 Self::WrongScope { expected, actual } => {
477 write!(
478 formatter,
479 "expected cursor scope {expected}, received {actual}"
480 )
481 }
482 Self::SequenceRegression {
483 family,
484 scope,
485 current,
486 proposed,
487 } => write!(
488 formatter,
489 "{} cursor for {scope} regressed from {current} to {proposed}",
490 family.as_str()
491 ),
492 }
493 }
494}
495
496impl std::error::Error for StreamCursorRefError {}
497
498#[derive(Deserialize)]
499#[serde(deny_unknown_fields)]
500struct CurrentStreamCursorRefWire {
501 position: ReplayCursor,
502 created_at: DateTime<Utc>,
503 #[serde(default)]
504 metadata: Metadata,
505}
506
507#[derive(Deserialize)]
508#[serde(deny_unknown_fields)]
509struct LegacyStreamCursorRefWire {
510 family: String,
511 scope: String,
512 sequence: usize,
513 #[serde(default)]
514 cursor: Option<String>,
515 created_at: DateTime<Utc>,
516 #[serde(default)]
517 metadata: Metadata,
518}
519
520#[derive(Deserialize)]
521#[serde(untagged)]
522enum StreamCursorRefWire {
523 Current(CurrentStreamCursorRefWire),
524 Legacy(LegacyStreamCursorRefWire),
525}
526
527impl<'de> Deserialize<'de> for StreamCursorRef {
528 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
529 where
530 D: Deserializer<'de>,
531 {
532 match StreamCursorRefWire::deserialize(deserializer)? {
533 StreamCursorRefWire::Current(current) => Ok(Self {
534 position: current.position,
535 created_at: current.created_at,
536 metadata: current.metadata,
537 }),
538 StreamCursorRefWire::Legacy(legacy) => {
539 let LegacyStreamCursorRefWire {
540 family,
541 scope,
542 sequence,
543 cursor,
544 created_at,
545 metadata,
546 } = legacy;
547 let family = match family.as_str() {
548 "raw_runtime" => ReplayCursorFamily::RawRuntime,
549 "display" => ReplayCursorFamily::Display,
550 "replay_event" => ReplayCursorFamily::ReplayEvent,
551 other => {
552 return Err(D::Error::custom(format!(
553 "unknown stream cursor family: {other}"
554 )));
555 }
556 };
557 let mut position =
558 ReplayCursor::for_family(family, ReplayScope::from_string(scope), sequence);
559 position.backend_cursor = cursor;
560 Ok(Self {
561 position,
562 created_at,
563 metadata,
564 })
565 }
566 }
567 }
568}
569
570#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
572pub struct SessionRecord {
573 pub session_id: SessionId,
575 #[serde(default = "default_session_namespace")]
577 pub namespace_id: String,
578 #[serde(default, skip_serializing_if = "Option::is_none")]
580 pub owner_id: Option<String>,
581 #[serde(default = "initial_session_revision")]
583 pub revision: u64,
584 #[serde(default)]
586 pub deletion_fence: SessionDeletionFence,
587 #[serde(default, skip_serializing_if = "Option::is_none")]
589 pub title: Option<String>,
590 #[serde(default, skip_serializing_if = "Option::is_none")]
592 pub workspace: Option<String>,
593 #[serde(default, skip_serializing_if = "Option::is_none")]
595 pub profile: Option<String>,
596 #[serde(default)]
598 pub status: SessionStatus,
599 #[serde(default)]
601 pub state: ResumableState,
602 #[serde(default, skip_serializing_if = "Option::is_none")]
604 pub environment_state: Option<EnvironmentStateRef>,
605 #[serde(default, skip_serializing_if = "Vec::is_empty")]
607 pub stream_cursors: Vec<StreamCursorRef>,
608 #[serde(default, skip_serializing_if = "TraceContext::is_empty")]
610 pub trace_context: TraceContext,
611 #[serde(default, skip_serializing_if = "Option::is_none")]
613 pub parent_session_id: Option<SessionId>,
614 #[serde(default, skip_serializing_if = "Option::is_none")]
616 pub head_run_id: Option<RunId>,
617 #[serde(default, skip_serializing_if = "Option::is_none")]
619 pub head_success_run_id: Option<RunId>,
620 #[serde(default, skip_serializing_if = "Option::is_none")]
622 pub active_run_id: Option<RunId>,
623 pub created_at: DateTime<Utc>,
625 pub updated_at: DateTime<Utc>,
627 #[serde(default, skip_serializing_if = "Metadata::is_empty")]
629 pub metadata: Metadata,
630}
631
632fn default_session_namespace() -> String {
633 crate::LOCAL_SESSION_NAMESPACE.to_string()
634}
635
636const fn initial_session_revision() -> u64 {
637 1
638}
639
640const fn initial_run_revision() -> u64 {
641 1
642}
643
644impl starweaver_core::VersionedRecord for SessionRecord {
645 const SCHEMA: &'static str = "starweaver.session.session_record";
646 const ALLOW_BARE_V0: bool = true;
647}
648
649impl SessionRecord {
650 #[must_use]
652 pub fn new(session_id: SessionId) -> Self {
653 let now = Utc::now();
654 Self {
655 session_id,
656 namespace_id: default_session_namespace(),
657 owner_id: None,
658 revision: initial_session_revision(),
659 deletion_fence: SessionDeletionFence::Stable,
660 title: None,
661 workspace: None,
662 profile: None,
663 status: SessionStatus::Active,
664 state: ResumableState::default(),
665 environment_state: None,
666 stream_cursors: Vec::new(),
667 trace_context: TraceContext::default(),
668 parent_session_id: None,
669 head_run_id: None,
670 head_success_run_id: None,
671 active_run_id: None,
672 created_at: now,
673 updated_at: now,
674 metadata: Metadata::default(),
675 }
676 }
677}
678
679#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
681pub struct RunRecord {
682 pub session_id: SessionId,
684 pub run_id: RunId,
686 #[serde(default = "initial_run_revision")]
688 pub revision: u64,
689 pub conversation_id: ConversationId,
691 #[serde(default, skip_serializing_if = "Vec::is_empty")]
693 pub input: Vec<InputPart>,
694 #[serde(default)]
696 pub status: RunStatus,
697 #[serde(default, skip_serializing_if = "Option::is_none")]
699 pub output_preview: Option<String>,
700 #[serde(default, skip_serializing_if = "Option::is_none")]
702 pub terminal_error: Option<RunTerminalError>,
703 #[serde(default, skip_serializing_if = "Value::is_null")]
705 pub structured_output: Value,
706 #[serde(default, skip_serializing_if = "Option::is_none")]
708 pub latest_checkpoint: Option<CheckpointRef>,
709 #[serde(default, skip_serializing_if = "Option::is_none")]
711 pub environment_state: Option<EnvironmentStateRef>,
712 #[serde(default, skip_serializing_if = "Vec::is_empty")]
714 pub stream_cursors: Vec<StreamCursorRef>,
715 #[serde(default, skip_serializing_if = "TraceContext::is_empty")]
717 pub trace_context: TraceContext,
718 #[serde(default)]
720 pub sequence_no: usize,
721 #[serde(default, skip_serializing_if = "Option::is_none")]
723 pub restore_from_run_id: Option<RunId>,
724 #[serde(default, skip_serializing_if = "Option::is_none")]
726 pub parent_run_id: Option<RunId>,
727 #[serde(default, skip_serializing_if = "Option::is_none")]
729 pub parent_task_id: Option<TaskId>,
730 #[serde(default, skip_serializing_if = "Option::is_none")]
732 pub trigger_type: Option<String>,
733 #[serde(default, skip_serializing_if = "Option::is_none")]
735 pub profile: Option<String>,
736 pub created_at: DateTime<Utc>,
738 pub updated_at: DateTime<Utc>,
740 #[serde(default, skip_serializing_if = "Metadata::is_empty")]
742 pub metadata: Metadata,
743}
744
745impl starweaver_core::VersionedRecord for RunRecord {
746 const SCHEMA: &'static str = "starweaver.session.run_record";
747 const ALLOW_BARE_V0: bool = true;
748}
749
750impl RunRecord {
751 #[must_use]
755 pub fn terminal_projection(&self) -> Option<RunTerminalProjection> {
756 self.status.is_terminal().then(|| RunTerminalProjection {
757 status: self.status,
758 output_preview: self.output_preview.clone(),
759 error: self.terminal_error.clone(),
760 })
761 }
762
763 pub fn validate_new_write(&self) -> Result<(), RunTerminalProjectionError> {
772 self.terminal_projection().map_or_else(
773 || {
774 if self.terminal_error.is_some() {
775 Err(RunTerminalProjectionError::UnexpectedNonTerminalDiagnostic(
776 self.status,
777 ))
778 } else {
779 Ok(())
780 }
781 },
782 |terminal| terminal.validate(),
783 )
784 }
785
786 pub fn normalize_for_admission(&mut self) {
791 self.status = RunStatus::Queued;
792 self.output_preview = None;
793 self.terminal_error = None;
794 }
795
796 pub fn apply_legacy_status_update(
801 &mut self,
802 status: RunStatus,
803 output_preview: Option<String>,
804 ) {
805 self.status = status;
806 match status {
807 RunStatus::Failed => {
808 self.output_preview = None;
809 self.terminal_error = Some(RunTerminalError::new(
810 "legacy_status_update_failed",
811 "run failed",
812 ));
813 }
814 RunStatus::Cancelled => {
815 self.output_preview = None;
816 self.terminal_error = Some(RunTerminalError::new(
817 "legacy_status_update_cancelled",
818 "run cancelled",
819 ));
820 }
821 _ => {
822 self.output_preview = output_preview;
823 self.terminal_error = None;
824 }
825 }
826 }
827
828 #[must_use]
830 pub fn new(session_id: SessionId, run_id: RunId, conversation_id: ConversationId) -> Self {
831 let now = Utc::now();
832 Self {
833 session_id,
834 run_id,
835 revision: initial_run_revision(),
836 conversation_id,
837 input: Vec::new(),
838 status: RunStatus::Queued,
839 output_preview: None,
840 terminal_error: None,
841 structured_output: Value::Null,
842 latest_checkpoint: None,
843 environment_state: None,
844 stream_cursors: Vec::new(),
845 trace_context: TraceContext::default(),
846 sequence_no: 0,
847 restore_from_run_id: None,
848 parent_run_id: None,
849 parent_task_id: None,
850 trigger_type: None,
851 profile: None,
852 created_at: now,
853 updated_at: now,
854 metadata: Metadata::default(),
855 }
856 }
857}