use std::collections::{BTreeMap, BTreeSet};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use starweaver_core::{RunId, SessionId};
use starweaver_stream::{DisplayMessage, ReplayCursor};
use crate::{InputPart, RunStatus, SessionStatus};
pub const LOCAL_SESSION_NAMESPACE: &str = "local";
#[derive(Clone, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
pub struct ManagedSessionTarget {
pub namespace_id: String,
pub session_id: SessionId,
}
impl ManagedSessionTarget {
#[must_use]
pub fn new(namespace_id: impl Into<String>, session_id: SessionId) -> Self {
Self {
namespace_id: namespace_id.into(),
session_id,
}
}
}
#[derive(Clone, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
pub struct ManagedRunTarget {
pub namespace_id: String,
pub session_id: SessionId,
pub run_id: RunId,
}
impl ManagedRunTarget {
#[must_use]
pub fn new(namespace_id: impl Into<String>, session_id: SessionId, run_id: RunId) -> Self {
Self {
namespace_id: namespace_id.into(),
session_id,
run_id,
}
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum AgentSessionOperation {
Read,
Search,
Create,
Update,
Control,
Delete,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentSessionScope {
pub namespace_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub owner_id: Option<String>,
pub source_product: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source_session_id: Option<SessionId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source_run_id: Option<RunId>,
#[serde(default)]
pub operations: BTreeSet<AgentSessionOperation>,
#[serde(default)]
pub allowed_session_ids: BTreeSet<SessionId>,
#[serde(default = "default_true")]
pub allow_self_query: bool,
#[serde(default)]
pub allow_self_control: bool,
pub policy_fingerprint: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deadline: Option<DateTime<Utc>>,
#[serde(default = "default_page_limit")]
pub max_page_size: u32,
}
const fn default_true() -> bool {
true
}
const fn default_page_limit() -> u32 {
50
}
impl AgentSessionScope {
#[must_use]
pub fn allows(&self, operation: AgentSessionOperation) -> bool {
self.operations.contains(&operation)
}
#[must_use]
pub fn allows_session(&self, session_id: &SessionId) -> bool {
self.allowed_session_ids.is_empty() || self.allowed_session_ids.contains(session_id)
}
#[must_use]
pub fn is_self_run(&self, target: &ManagedRunTarget) -> bool {
self.namespace_id == target.namespace_id
&& self.source_session_id.as_ref() == Some(&target.session_id)
&& self.source_run_id.as_ref() == Some(&target.run_id)
}
#[must_use]
pub fn is_self_session(&self, target: &ManagedSessionTarget) -> bool {
self.namespace_id == target.namespace_id
&& self.source_session_id.as_ref() == Some(&target.session_id)
}
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "state", rename_all = "snake_case")]
pub enum SessionDeletionFence {
#[default]
Stable,
Deleting {
fence_id: String,
expected_revision: u64,
requested_by: String,
started_at: DateTime<Utc>,
},
Deleted {
fence_id: String,
deleted_at: DateTime<Utc>,
},
}
impl SessionDeletionFence {
#[must_use]
pub const fn blocks_continuation(&self) -> bool {
!matches!(self, Self::Stable)
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct SessionContinuationFence {
pub target: ManagedSessionTarget,
pub revision: u64,
pub continuation_allowed: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub fence_id: Option<String>,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentSessionListQuery {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub status: Option<SessionStatus>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workspace: Option<String>,
#[serde(default = "default_page_limit")]
pub limit: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub page_token: Option<String>,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentSessionInclude {
#[serde(default)]
pub recent_runs: bool,
#[serde(default)]
pub trace: bool,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentRunListQuery {
#[serde(default = "default_page_limit")]
pub limit: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub page_token: Option<String>,
}
impl Default for AgentRunListQuery {
fn default() -> Self {
Self {
limit: default_page_limit(),
page_token: None,
}
}
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentReplayQuery {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub after: Option<ReplayCursor>,
#[serde(default = "default_page_limit")]
pub limit: u32,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentSessionView {
pub target: ManagedSessionTarget,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
pub status: SessionStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workspace: Option<String>,
pub revision: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub head_run_id: Option<RunId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub active_run_id: Option<RunId>,
pub resumable: bool,
pub controllable: bool,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub recent_runs: Vec<AgentRunView>,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentRunView {
pub target: ManagedRunTarget,
pub status: RunStatus,
pub sequence_no: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input_preview: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output_preview: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error_category: Option<String>,
pub controllable: bool,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentSessionPage {
pub sessions: Vec<AgentSessionView>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_page_token: Option<String>,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentRunPage {
pub runs: Vec<AgentRunView>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_page_token: Option<String>,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentDisplayPage {
pub messages: Vec<DisplayMessage>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_cursor: Option<ReplayCursor>,
#[serde(default = "untrusted_evidence_label")]
pub trust: String,
}
fn untrusted_evidence_label() -> String {
"untrusted_historical_evidence".to_string()
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct CreateManagedSession {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workspace: Option<String>,
#[serde(default)]
pub metadata: BTreeMap<String, Value>,
pub idempotency_key: String,
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
pub struct ManagedSessionPatch {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub title: Option<Option<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile: Option<Option<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub archived: Option<bool>,
#[serde(default)]
pub metadata: BTreeMap<String, Value>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct UpdateManagedSession {
pub session_id: SessionId,
pub expected_revision: u64,
pub patch: ManagedSessionPatch,
pub idempotency_key: String,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct DeleteManagedSession {
pub session_id: SessionId,
pub expected_revision: u64,
pub idempotency_key: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub approval_receipt_id: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct StartManagedRun {
pub session_id: SessionId,
pub input: Vec<InputPart>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub profile: Option<String>,
#[serde(default)]
pub environment_refs: Vec<String>,
pub idempotency_key: String,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct SteerManagedRun {
pub target: ManagedRunTarget,
pub steering_id: String,
pub text: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct InterruptManagedRun {
pub target: ManagedRunTarget,
pub operation_id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reason_category: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub idempotency_key: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct SessionMutationReceipt {
pub receipt_id: String,
pub session: AgentSessionView,
pub idempotent_replay: bool,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct RunStartReceipt {
pub receipt_id: String,
pub target: ManagedRunTarget,
pub status: RunStatus,
pub fencing_generation: u64,
pub idempotent_replay: bool,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct RunControlReceipt {
pub receipt_id: String,
pub target: ManagedRunTarget,
pub operation_id: String,
pub fencing_generation: u64,
pub accepted: bool,
pub idempotent_replay: bool,
pub created_at: DateTime<Utc>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct RunAdmissionLease {
pub target: ManagedRunTarget,
pub admission_id: String,
pub host_instance_id: String,
pub fencing_generation: u64,
pub lease_expires_at: DateTime<Utc>,
pub heartbeat_at: DateTime<Utc>,
pub command_fingerprint: String,
pub idempotency_key: String,
}
impl starweaver_core::VersionedRecord for RunAdmissionLease {
const SCHEMA: &'static str = "starweaver.session.run_admission_lease";
}
impl RunAdmissionLease {
#[must_use]
pub fn expired_at(&self, now: DateTime<Utc>) -> bool {
self.lease_expires_at <= now
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AcquireRunAdmission {
pub run: crate::RunRecord,
pub namespace_id: String,
pub host_instance_id: String,
pub admission_id: String,
pub lease_expires_at: DateTime<Utc>,
pub idempotency_key: String,
pub command_fingerprint: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub replaces_waiting_run_id: Option<RunId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hitl_resume_claim_id: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct RunAdmissionReceipt {
pub run: crate::RunRecord,
pub lease: RunAdmissionLease,
pub idempotent_replay: bool,
}
impl starweaver_core::VersionedRecord for RunAdmissionReceipt {
const SCHEMA: &'static str = "starweaver.session.run_admission_receipt";
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct DurableControlReceipt {
pub receipt_id: String,
pub target: ManagedRunTarget,
pub operation_id: String,
pub operation: String,
pub idempotency_key: String,
pub command_fingerprint: String,
pub fencing_generation: u64,
pub state: String,
pub created_at: DateTime<Utc>,
}
impl starweaver_core::VersionedRecord for DurableControlReceipt {
const SCHEMA: &'static str = "starweaver.session.durable_control_receipt";
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum AgentSessionQueryErrorCode {
InvalidQuery,
NotFound,
Unsupported,
Unavailable,
PermissionDenied,
InvalidCursor,
Failed,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentSessionQueryError {
pub code: AgentSessionQueryErrorCode,
pub message: String,
}
impl std::fmt::Display for AgentSessionQueryError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "{:?}: {}", self.code, self.message)
}
}
impl std::error::Error for AgentSessionQueryError {}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum AgentSessionControlErrorCode {
InvalidCommand,
NotFound,
PermissionDenied,
ApprovalRequired,
Conflict,
IdempotencyConflict,
RunConflict,
NotActive,
Terminal,
StaleActive,
QuotaExceeded,
Unavailable,
Failed,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AgentSessionControlError {
pub code: AgentSessionControlErrorCode,
pub message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub current_revision: Option<u64>,
}
impl std::fmt::Display for AgentSessionControlError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "{:?}: {}", self.code, self.message)
}
}
impl std::error::Error for AgentSessionControlError {}