use std::collections::{BTreeMap, BTreeSet};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use starweaver_core::{Metadata, RunId, SessionId};
use crate::{
ApprovalDecision, ApprovalRecord, ApprovalStatus, DeferredToolRecord, ExecutionStatus,
MutationReceipt, PendingHostEventPublication, SessionStoreError, SessionStoreResult,
};
pub const APPROVAL_DECIDE_OPERATION: &str = "approval.decide";
pub const DEFERRED_COMPLETE_OPERATION: &str = "deferred.complete";
pub const DEFERRED_FAIL_OPERATION: &str = "deferred.fail";
pub const CLARIFICATION_RESOLVE_OPERATION: &str = "clarification.resolve";
pub const ASK_USER_QUESTION_ACTION: &str = "ask_user_question";
pub const CLARIFICATION_ANSWERS_METADATA_KEY: &str = "clarification_answers";
pub const CLARIFICATION_RESPONSE_METADATA_KEY: &str = "clarification_response";
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct InteractionMutationContext {
pub authority_binding: String,
pub expected_revision: u64,
pub idempotency_key: String,
pub command_fingerprint: String,
pub occurred_at: DateTime<Utc>,
pub host_event_publication: Option<PendingHostEventPublication>,
}
impl InteractionMutationContext {
pub fn validate(&self) -> SessionStoreResult<()> {
require_non_empty("interaction authority binding", &self.authority_binding)?;
require_non_empty("interaction idempotency key", &self.idempotency_key)?;
require_non_empty("interaction command fingerprint", &self.command_fingerprint)?;
if self.expected_revision == 0 {
return Err(SessionStoreError::Failed(
"interaction expected revision must be greater than zero".to_string(),
));
}
if let Some(publication) = &self.host_event_publication {
publication.validate()?;
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct DecideApproval {
pub context: InteractionMutationContext,
pub session_id: SessionId,
pub run_id: RunId,
pub approval_id: String,
pub decision: ApprovalDecision,
}
impl DecideApproval {
pub fn validate(&self) -> SessionStoreResult<()> {
self.context.validate()?;
require_non_empty("approval id", &self.approval_id)?;
if !matches!(
self.decision.status,
ApprovalStatus::Approved | ApprovalStatus::Denied
) {
return Err(SessionStoreError::Failed(
"approval decision must be approved or denied".to_string(),
));
}
if self.decision.decided_at != self.context.occurred_at {
return Err(SessionStoreError::Failed(
"approval decision time must match mutation time".to_string(),
));
}
Ok(())
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ApprovalMutationResult {
pub approval: ApprovalRecord,
pub receipt: MutationReceipt,
}
impl starweaver_core::VersionedRecord for ApprovalMutationResult {
const SCHEMA: &'static str = "starweaver.session.approval_mutation_result";
}
impl ApprovalMutationResult {
#[must_use]
pub fn replayed_projection(&self) -> Self {
Self {
approval: self.approval.clone(),
receipt: self.receipt.replayed_projection(),
}
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum DeferredMutationOutcome {
Completed {
response: Value,
#[serde(default, skip_serializing_if = "Metadata::is_empty")]
metadata: Metadata,
},
Failed {
response: Value,
#[serde(default, skip_serializing_if = "Metadata::is_empty")]
metadata: Metadata,
},
}
impl DeferredMutationOutcome {
#[must_use]
pub const fn operation(&self) -> &'static str {
match self {
Self::Completed { .. } => DEFERRED_COMPLETE_OPERATION,
Self::Failed { .. } => DEFERRED_FAIL_OPERATION,
}
}
#[must_use]
pub const fn status(&self) -> ExecutionStatus {
match self {
Self::Completed { .. } => ExecutionStatus::Completed,
Self::Failed { .. } => ExecutionStatus::Failed,
}
}
#[must_use]
pub const fn parts(&self) -> (&Value, &Metadata) {
match self {
Self::Completed { response, metadata } | Self::Failed { response, metadata } => {
(response, metadata)
}
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ResolveDeferredTool {
pub context: InteractionMutationContext,
pub session_id: SessionId,
pub run_id: RunId,
pub deferred_id: String,
pub outcome: DeferredMutationOutcome,
}
impl ResolveDeferredTool {
pub fn validate(&self) -> SessionStoreResult<()> {
self.context.validate()?;
require_non_empty("deferred id", &self.deferred_id)
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct DeferredMutationResult {
pub deferred: DeferredToolRecord,
pub receipt: MutationReceipt,
}
impl starweaver_core::VersionedRecord for DeferredMutationResult {
const SCHEMA: &'static str = "starweaver.session.deferred_mutation_result";
}
impl DeferredMutationResult {
#[must_use]
pub fn replayed_projection(&self) -> Self {
Self {
deferred: self.deferred.clone(),
receipt: self.receipt.replayed_projection(),
}
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ClarificationOption {
pub label: String,
#[serde(default)]
pub description: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub preview: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ClarificationQuestion {
#[serde(default)]
pub header: String,
pub question: String,
#[serde(default, alias = "multiSelect")]
pub multi_select: bool,
#[serde(default)]
pub options: Vec<ClarificationOption>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ClarificationAnswer {
pub question: String,
#[serde(default, alias = "selectedOptions")]
pub selected_options: Vec<String>,
#[serde(default, alias = "freeText", skip_serializing_if = "Option::is_none")]
pub free_text: Option<String>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ResolveClarification {
pub context: InteractionMutationContext,
pub session_id: SessionId,
pub run_id: RunId,
pub clarification_id: String,
pub answers: Vec<ClarificationAnswer>,
pub response: Option<String>,
pub resolved_by: Option<String>,
}
impl ResolveClarification {
pub fn validate(&self) -> SessionStoreResult<()> {
self.context.validate()?;
require_non_empty("clarification id", &self.clarification_id)?;
if self.response.as_deref().is_some_and(str::is_empty) {
return Err(SessionStoreError::Failed(
"clarification response cannot be empty".to_string(),
));
}
if self.resolved_by.as_deref().is_some_and(str::is_empty) {
return Err(SessionStoreError::Failed(
"clarification resolver cannot be empty".to_string(),
));
}
Ok(())
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ClarificationResolution {
pub clarification_id: String,
pub session_id: SessionId,
pub run_id: RunId,
pub questions: Vec<ClarificationQuestion>,
pub answers: Vec<ClarificationAnswer>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub response: Option<String>,
pub revision: u64,
pub resolved_at: DateTime<Utc>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ClarificationMutationResult {
pub clarification: ClarificationResolution,
pub approval: ApprovalRecord,
pub receipt: MutationReceipt,
}
impl starweaver_core::VersionedRecord for ClarificationMutationResult {
const SCHEMA: &'static str = "starweaver.session.clarification_mutation_result";
}
impl ClarificationMutationResult {
#[must_use]
pub fn replayed_projection(&self) -> Self {
Self {
clarification: self.clarification.clone(),
approval: self.approval.clone(),
receipt: self.receipt.replayed_projection(),
}
}
}
#[allow(clippy::too_many_lines)]
pub fn validate_clarification_answers(
request: &Value,
answers: &[ClarificationAnswer],
) -> SessionStoreResult<(Vec<ClarificationQuestion>, Vec<ClarificationAnswer>)> {
let questions_value = request
.as_object()
.and_then(|object| object.get("questions"))
.ok_or_else(|| {
SessionStoreError::Failed(
"clarification request must contain a questions array".to_string(),
)
})?;
let questions = serde_json::from_value::<Vec<ClarificationQuestion>>(questions_value.clone())
.map_err(|error| {
SessionStoreError::Failed(format!("malformed clarification questions: {error}"))
})?;
if questions.is_empty() {
return Err(SessionStoreError::Failed(
"clarification request must contain at least one question".to_string(),
));
}
let mut by_question = BTreeMap::new();
for question in &questions {
require_non_empty("clarification question", &question.question)?;
let mut labels = BTreeSet::new();
for option in &question.options {
require_non_empty("clarification option label", &option.label)?;
if !labels.insert(option.label.as_str()) {
return Err(SessionStoreError::Failed(format!(
"duplicate clarification option {} for question {}",
option.label, question.question
)));
}
}
if by_question
.insert(question.question.as_str(), question)
.is_some()
{
return Err(SessionStoreError::Failed(format!(
"duplicate clarification question {}",
question.question
)));
}
}
if answers.len() != questions.len() {
return Err(SessionStoreError::Failed(
"clarification answers must match every durable question exactly once".to_string(),
));
}
let mut answer_by_question = BTreeMap::new();
for answer in answers {
let question = by_question.get(answer.question.as_str()).ok_or_else(|| {
SessionStoreError::Failed(format!(
"clarification answer does not match durable question {}",
answer.question
))
})?;
if answer_by_question
.insert(answer.question.as_str(), answer)
.is_some()
{
return Err(SessionStoreError::Failed(format!(
"duplicate clarification answer for {}",
answer.question
)));
}
if answer.free_text.as_deref().is_some_and(str::is_empty) {
return Err(SessionStoreError::Failed(format!(
"clarification free text cannot be empty for {}",
answer.question
)));
}
if answer.selected_options.is_empty() && answer.free_text.is_none() {
return Err(SessionStoreError::Failed(format!(
"clarification answer is empty for {}",
answer.question
)));
}
if !question.multi_select && answer.selected_options.len() > 1 {
return Err(SessionStoreError::Failed(format!(
"clarification question {} does not allow multiple selections",
answer.question
)));
}
let allowed = question
.options
.iter()
.map(|option| option.label.as_str())
.collect::<BTreeSet<_>>();
let mut selected = BTreeSet::new();
for label in &answer.selected_options {
if !selected.insert(label.as_str()) {
return Err(SessionStoreError::Failed(format!(
"duplicate clarification selection {label} for {}",
answer.question
)));
}
if !allowed.contains(label.as_str()) {
return Err(SessionStoreError::Failed(format!(
"clarification selection {label} is not an option for {}",
answer.question
)));
}
}
}
let ordered_answers = questions
.iter()
.map(|question| {
answer_by_question
.get(question.question.as_str())
.map(|answer| (*answer).clone())
.ok_or_else(|| {
SessionStoreError::Failed(format!(
"missing clarification answer for {}",
question.question
))
})
})
.collect::<SessionStoreResult<Vec<_>>>()?;
Ok((questions, ordered_answers))
}
fn require_non_empty(label: &str, value: &str) -> SessionStoreResult<()> {
if value.is_empty() {
return Err(SessionStoreError::Failed(format!(
"{label} cannot be empty"
)));
}
Ok(())
}