use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::{
DurableHostEventClass, DurableHostEventScope, MutationReceipt, PendingHostEventPublication,
SessionStoreError, SessionStoreResult,
};
pub const MAX_ENVIRONMENT_PAGE_SIZE: u32 = 200;
pub const MAX_ENVIRONMENT_MOUNTS_PER_RUN: u32 = 128;
pub const ENVIRONMENT_ATTACH_OPERATION: &str = "environment.attach";
pub const ENVIRONMENT_DETACH_OPERATION: &str = "environment.detach";
pub const ENVIRONMENT_MOUNT_OPERATION: &str = "environment.mount";
pub const ENVIRONMENT_UNMOUNT_OPERATION: &str = "environment.unmount";
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum DurableEnvironmentScope {
Connection {
connection_id: String,
},
Session {
session_id: String,
},
Run {
session_id: String,
run_id: String,
},
}
impl DurableEnvironmentScope {
pub fn validate(&self) -> SessionStoreResult<()> {
match self {
Self::Connection { connection_id } => {
require_non_empty("environment scope connection", connection_id)
}
Self::Session { session_id } => {
require_non_empty("environment scope session", session_id)
}
Self::Run { session_id, run_id } => {
require_non_empty("environment scope session", session_id)?;
require_non_empty("environment scope run", run_id)
}
}
}
#[must_use]
pub fn permits_run(&self, connection_id: Option<&str>, session_id: &str, run_id: &str) -> bool {
match self {
Self::Connection {
connection_id: owner,
} => connection_id == Some(owner.as_str()),
Self::Session { session_id: owner } => owner == session_id,
Self::Run {
session_id: owner_session,
run_id: owner_run,
} => owner_session == session_id && owner_run == run_id,
}
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum DurableEnvironmentStatus {
Attaching,
Ready,
Degraded,
Detached,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct DurableEnvironmentAttachment {
pub authority_binding: String,
pub attachment_id: String,
pub environment_id: String,
pub display_name: Option<String>,
pub scope: DurableEnvironmentScope,
pub status: DurableEnvironmentStatus,
pub revision: u64,
pub updated_at: DateTime<Utc>,
}
impl starweaver_core::VersionedRecord for DurableEnvironmentAttachment {
const SCHEMA: &'static str = "starweaver.session.environment_attachment";
}
impl DurableEnvironmentAttachment {
pub fn validate(&self) -> SessionStoreResult<()> {
require_non_empty("environment authority binding", &self.authority_binding)?;
require_non_empty("environment attachment id", &self.attachment_id)?;
require_non_empty("public environment id", &self.environment_id)?;
if self.display_name.as_ref().is_some_and(String::is_empty) {
return Err(SessionStoreError::Failed(
"environment display name cannot be empty".to_string(),
));
}
self.scope.validate()?;
require_revision(self.revision, "environment attachment")
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum DurableEnvironmentMountStatus {
Mounted,
Unmounted,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct DurableEnvironmentMount {
pub authority_binding: String,
pub mount_id: String,
pub attachment_id: String,
pub session_id: String,
pub run_id: String,
pub resource_label: String,
pub status: DurableEnvironmentMountStatus,
pub revision: u64,
pub updated_at: DateTime<Utc>,
}
impl starweaver_core::VersionedRecord for DurableEnvironmentMount {
const SCHEMA: &'static str = "starweaver.session.environment_mount";
}
impl DurableEnvironmentMount {
pub fn validate(&self) -> SessionStoreResult<()> {
require_non_empty("environment authority binding", &self.authority_binding)?;
require_non_empty("environment mount id", &self.mount_id)?;
require_non_empty("environment attachment id", &self.attachment_id)?;
require_non_empty("environment mount session", &self.session_id)?;
require_non_empty("environment mount run", &self.run_id)?;
require_non_empty("environment resource label", &self.resource_label)?;
require_revision(self.revision, "environment mount")
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct EnvironmentHostEventContext {
pub transition_identity: String,
pub scope: DurableHostEventScope,
}
impl EnvironmentHostEventContext {
pub fn validate(&self) -> SessionStoreResult<()> {
require_non_empty(
"environment host event transition identity",
&self.transition_identity,
)?;
match &self.scope {
DurableHostEventScope::Global => Ok(()),
DurableHostEventScope::Session { session_id } => {
require_non_empty("environment host event session", session_id.as_str())
}
DurableHostEventScope::Run { session_id, run_id } => {
require_non_empty("environment host event session", session_id.as_str())?;
require_non_empty("environment host event run", run_id.as_str())
}
}
}
pub fn publication(
&self,
attachment: &DurableEnvironmentAttachment,
) -> SessionStoreResult<PendingHostEventPublication> {
self.validate()?;
attachment.validate()?;
let attachment_scope = match &attachment.scope {
DurableEnvironmentScope::Connection { .. } => {
serde_json::json!({"kind": "connection"})
}
DurableEnvironmentScope::Session { session_id } => {
serde_json::json!({"kind": "session", "sessionId": session_id})
}
DurableEnvironmentScope::Run { session_id, run_id } => serde_json::json!({
"kind": "run",
"sessionId": session_id,
"runId": run_id,
}),
};
let status = match attachment.status {
DurableEnvironmentStatus::Attaching => "attaching",
DurableEnvironmentStatus::Ready => "ready",
DurableEnvironmentStatus::Degraded => "degraded",
DurableEnvironmentStatus::Detached => "detached",
};
let mut projected_attachment = serde_json::Map::from_iter([
(
"attachmentId".to_string(),
serde_json::Value::String(attachment.attachment_id.clone()),
),
(
"environmentId".to_string(),
serde_json::Value::String(attachment.environment_id.clone()),
),
(
"revision".to_string(),
serde_json::Value::String(attachment.revision.to_string()),
),
("scope".to_string(), attachment_scope),
(
"status".to_string(),
serde_json::Value::String(status.to_string()),
),
]);
if let Some(display_name) = &attachment.display_name {
projected_attachment.insert(
"displayName".to_string(),
serde_json::Value::String(display_name.clone()),
);
}
PendingHostEventPublication::new(
&self.transition_identity,
0,
self.scope.clone(),
DurableHostEventClass::EnvironmentChanged,
serde_json::json!({
"kind": "environment_changed",
"attachment": serde_json::Value::Object(projected_attachment),
}),
attachment.updated_at,
)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct EnvironmentMutationContext {
pub authority_binding: String,
pub idempotency_key: String,
pub command_fingerprint: String,
pub occurred_at: DateTime<Utc>,
pub host_event: Option<EnvironmentHostEventContext>,
}
impl EnvironmentMutationContext {
pub fn validate(&self) -> SessionStoreResult<()> {
require_non_empty("environment authority binding", &self.authority_binding)?;
require_non_empty("environment idempotency key", &self.idempotency_key)?;
require_non_empty("environment command fingerprint", &self.command_fingerprint)?;
if let Some(host_event) = &self.host_event {
host_event.validate()?;
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct AttachEnvironment {
pub context: EnvironmentMutationContext,
pub attachment_id: String,
pub environment_id: String,
pub display_name: Option<String>,
pub scope: DurableEnvironmentScope,
pub status: DurableEnvironmentStatus,
}
impl AttachEnvironment {
pub fn validate(&self) -> SessionStoreResult<()> {
self.context.validate()?;
require_non_empty("environment attachment id", &self.attachment_id)?;
require_non_empty("public environment id", &self.environment_id)?;
if self.display_name.as_ref().is_some_and(String::is_empty) {
return Err(SessionStoreError::Failed(
"environment display name cannot be empty".to_string(),
));
}
if self.status == DurableEnvironmentStatus::Detached {
return Err(SessionStoreError::Failed(
"new environment attachment cannot be detached".to_string(),
));
}
self.scope.validate()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct DetachEnvironment {
pub context: EnvironmentMutationContext,
pub attachment_id: String,
}
impl DetachEnvironment {
pub fn validate(&self) -> SessionStoreResult<()> {
self.context.validate()?;
require_non_empty("environment attachment id", &self.attachment_id)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct MountEnvironmentResource {
pub context: EnvironmentMutationContext,
pub mount_id: String,
pub attachment_id: String,
pub session_id: String,
pub run_id: String,
pub connection_id: Option<String>,
pub resource_label: String,
}
impl MountEnvironmentResource {
pub fn validate(&self) -> SessionStoreResult<()> {
self.context.validate()?;
require_non_empty("environment mount id", &self.mount_id)?;
require_non_empty("environment attachment id", &self.attachment_id)?;
require_non_empty("environment mount session", &self.session_id)?;
require_non_empty("environment mount run", &self.run_id)?;
if self.connection_id.as_ref().is_some_and(String::is_empty) {
return Err(SessionStoreError::Failed(
"environment mount connection identity cannot be empty".to_string(),
));
}
require_non_empty("environment resource label", &self.resource_label)
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct UnmountEnvironmentResource {
pub context: EnvironmentMutationContext,
pub mount_id: String,
}
impl UnmountEnvironmentResource {
pub fn validate(&self) -> SessionStoreResult<()> {
self.context.validate()?;
require_non_empty("environment mount id", &self.mount_id)
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct EnvironmentAttachmentMutationResult {
pub attachment: DurableEnvironmentAttachment,
pub receipt: MutationReceipt,
}
impl starweaver_core::VersionedRecord for EnvironmentAttachmentMutationResult {
const SCHEMA: &'static str = "starweaver.session.environment_attachment_mutation_result";
}
impl EnvironmentAttachmentMutationResult {
pub fn validate(&self) -> SessionStoreResult<()> {
self.attachment.validate()?;
self.receipt.validate()?;
if self.receipt.target_ref != self.attachment.attachment_id {
return Err(SessionStoreError::Conflict(
"environment receipt target does not match attachment".to_string(),
));
}
Ok(())
}
#[must_use]
pub fn replayed_projection(&self) -> Self {
Self {
attachment: self.attachment.clone(),
receipt: self.receipt.replayed_projection(),
}
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct EnvironmentMountMutationResult {
pub mount: DurableEnvironmentMount,
pub receipt: MutationReceipt,
}
impl starweaver_core::VersionedRecord for EnvironmentMountMutationResult {
const SCHEMA: &'static str = "starweaver.session.environment_mount_mutation_result";
}
impl EnvironmentMountMutationResult {
pub fn validate(&self) -> SessionStoreResult<()> {
self.mount.validate()?;
self.receipt.validate()?;
if self.receipt.target_ref != self.mount.mount_id {
return Err(SessionStoreError::Conflict(
"environment receipt target does not match mount".to_string(),
));
}
Ok(())
}
#[must_use]
pub fn replayed_projection(&self) -> Self {
Self {
mount: self.mount.clone(),
receipt: self.receipt.replayed_projection(),
}
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "result_kind", rename_all = "snake_case")]
pub enum EnvironmentMutationResult {
Attachment(EnvironmentAttachmentMutationResult),
Mount(EnvironmentMountMutationResult),
}
impl starweaver_core::VersionedRecord for EnvironmentMutationResult {
const SCHEMA: &'static str = "starweaver.session.environment_mutation_result";
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct EnvironmentAttachmentPageKey {
pub updated_at: DateTime<Utc>,
pub attachment_id: String,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct EnvironmentAttachmentQuery {
pub authority_binding: String,
pub scope: Option<DurableEnvironmentScope>,
pub connection_id: Option<String>,
pub limit: u32,
pub after: Option<EnvironmentAttachmentPageKey>,
}
impl EnvironmentAttachmentQuery {
pub fn validate(&self) -> SessionStoreResult<()> {
require_non_empty("environment authority binding", &self.authority_binding)?;
require_page_limit(self.limit)?;
if let Some(scope) = &self.scope {
scope.validate()?;
if let DurableEnvironmentScope::Connection { connection_id } = scope
&& self.connection_id.as_deref() != Some(connection_id)
{
return Err(SessionStoreError::NotFound(
"environment connection scope".to_string(),
));
}
}
if self.connection_id.as_ref().is_some_and(String::is_empty) {
return Err(SessionStoreError::Failed(
"environment viewer connection identity cannot be empty".to_string(),
));
}
if let Some(after) = &self.after {
require_non_empty("environment attachment page key", &after.attachment_id)?;
}
Ok(())
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct EnvironmentAttachmentPage {
pub items: Vec<DurableEnvironmentAttachment>,
pub next: Option<EnvironmentAttachmentPageKeyProjection>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct EnvironmentAttachmentPageKeyProjection {
pub updated_at: DateTime<Utc>,
pub attachment_id: String,
}
impl From<EnvironmentAttachmentPageKeyProjection> for EnvironmentAttachmentPageKey {
fn from(value: EnvironmentAttachmentPageKeyProjection) -> Self {
Self {
updated_at: value.updated_at,
attachment_id: value.attachment_id,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct EnvironmentMountQuery {
pub authority_binding: String,
pub session_id: String,
pub run_id: String,
pub limit: u32,
}
impl EnvironmentMountQuery {
pub fn validate(&self) -> SessionStoreResult<()> {
require_non_empty("environment authority binding", &self.authority_binding)?;
require_non_empty("environment mount session", &self.session_id)?;
require_non_empty("environment mount run", &self.run_id)?;
require_page_limit(self.limit)
}
}
fn require_page_limit(limit: u32) -> SessionStoreResult<()> {
if !(1..=MAX_ENVIRONMENT_PAGE_SIZE).contains(&limit) {
return Err(SessionStoreError::Failed(format!(
"environment page limit must be between 1 and {MAX_ENVIRONMENT_PAGE_SIZE}"
)));
}
Ok(())
}
fn require_revision(revision: u64, label: &str) -> SessionStoreResult<()> {
if revision == 0 {
return Err(SessionStoreError::Failed(format!(
"{label} revision must be greater than zero"
)));
}
Ok(())
}
fn require_non_empty(label: &str, value: &str) -> SessionStoreResult<()> {
if value.is_empty() {
return Err(SessionStoreError::Failed(format!(
"{label} cannot be empty"
)));
}
Ok(())
}