use std::sync::Arc;
use turnframe_core::error::{AuthorizationError, InteractionError, OrchestratorError, StoreError};
use turnframe_core::interaction::InteractionRejection;
use turnframe_core::observe::{Signal, SignalLabels};
use turnframe_core::replay::WorkflowVersionRecord;
use turnframe_core::response::ResponseBlock;
use turnframe_store::conversation::StoredUserTurn;
use super::Session;
use crate::attachments::TurnAttachments;
use crate::interactions::{AcceptedInteraction, ResponseAdmission, ResponseContext};
use crate::orchestrator::CaseCandidate;
use crate::resolve::TargetResolver;
use crate::turn::{
DirectoryTerms, addressable_cards, admit_card_case, blocking_summary, build_resolver,
candidates_of_open_cards, project_case, recent_messages, unaddressable_card,
};
impl Session<'_> {
pub(super) async fn accept(&mut self) -> Result<(), OrchestratorError> {
let conversations = self.runtime.stores.conversations();
conversations
.load_conversation(self.account(), &self.input.conversation_id)
.await
.map_err(|error| match error {
StoreError::NotFound => {
OrchestratorError::Unauthorized(AuthorizationError::ConversationNotAccessible {
conversation_id: self.input.conversation_id,
})
}
other => OrchestratorError::Store(other),
})?;
match conversations
.append_user_turn(StoredUserTurn::new(self.input.clone(), self.now))
.await
{
Ok(()) | Err(StoreError::Conflict) => {}
Err(error) => return Err(OrchestratorError::Store(error)),
}
self.runtime
.stores
.replay()
.put(self.record.clone())
.await
.map_err(OrchestratorError::Store)?;
self.open_interactions = self
.runtime
.interactions
.open_for_conversation(self.account(), &self.input.conversation_id)
.await
.map_err(OrchestratorError::Interaction)?;
Ok(())
}
pub(super) async fn load_cases(&mut self) -> Result<(), OrchestratorError> {
let directory = Arc::clone(&self.runtime.directory);
let mut candidates = directory
.candidates(&self.input.actor, &self.input.conversation_id)
.await
.map_err(OrchestratorError::Store)?;
if let Some(origin) = self.input.origin.as_ref()
&& let Some(candidate) = directory
.resolve_origin(&self.input.actor, &origin.origin_token)
.await
.map_err(OrchestratorError::Store)?
{
self.origin_case = Some(candidate.key.clone());
candidates.push(candidate);
}
if let Some(confirmed) = self.confirmed_case.clone()
&& !candidates
.iter()
.any(|candidate| candidate.key == confirmed.key)
{
candidates.push(confirmed);
}
let from_cards = candidates_of_open_cards(&candidates, &self.open_interactions);
let proposed = candidates
.into_iter()
.map(|candidate| (candidate, false))
.chain(from_cards.into_iter().map(|candidate| (candidate, true)));
for (candidate, from_card) in proposed {
self.load_candidate(candidate, from_card).await?;
}
self.open_interactions =
addressable_cards(&self.cases, std::mem::take(&mut self.open_interactions));
Ok(())
}
pub(super) async fn load_candidate(
&mut self,
candidate: CaseCandidate,
from_card: bool,
) -> Result<(), OrchestratorError> {
if self.cases.contains_key(&candidate.key) {
return Ok(());
}
let Ok(registered) = self.runtime.workflows.require(&candidate.key.workflow) else {
return Ok(());
};
let loaded = registered
.executor
.load(self.account(), &candidate.key.case_id)
.await
.map_err(OrchestratorError::Store)?;
let candidate = if from_card {
let admitted = admit_card_case(
self.runtime.observer.as_ref(),
self.runtime.directory.as_ref(),
&self.input.actor,
&self.input.conversation_id,
candidate,
loaded.value.is_some(),
)
.await
.map_err(OrchestratorError::Store)?;
let Some(admitted) = admitted else {
return Ok(());
};
admitted
} else {
candidate
};
let case_ref = candidate.key.clone().at(loaded.revision);
let terms = DirectoryTerms {
confirm_every_write: candidate.confirm_every_write,
subject_only_when_named: candidate.subject_only_when_named,
};
let case = project_case(
self.runtime.observer.as_ref(),
®istered.definition,
case_ref.clone(),
candidate.label,
loaded.value,
terms,
)
.inspect_err(|error| {
if matches!(error, OrchestratorError::InvariantViolation(_)) {
self.runtime.observer.observe_labeled(
&Signal::WorkflowInvariantViolation,
&SignalLabels::workflow(candidate.key.workflow.clone()),
);
}
})?;
self.record.loaded_cases.push(case_ref);
self.record.workflow_versions.push(WorkflowVersionRecord {
key: registered.key.clone(),
version: registered.version.clone(),
});
self.cases.insert(candidate.key, case);
Ok(())
}
pub(super) async fn reload_cases(&mut self) -> Result<(), OrchestratorError> {
self.cases.clear();
self.open_interactions = self
.runtime
.interactions
.open_for_conversation(self.account(), &self.input.conversation_id)
.await
.map_err(OrchestratorError::Interaction)?;
self.load_cases().await
}
pub(super) async fn admit_response(
&mut self,
) -> Result<Option<AcceptedInteraction>, OrchestratorError> {
let Some(response) = self.input.interaction_response.clone() else {
return Ok(None);
};
let record = self
.runtime
.interactions
.get(self.account(), &response.interaction_id)
.await
.map_err(OrchestratorError::Interaction)?;
if let Some(rejection) = unaddressable_card(&self.cases, &record) {
return Err(OrchestratorError::Interaction(InteractionError::Rejected(
rejection,
)));
}
let current = self
.cases
.get(&record.interaction.case_ref.key())
.map_or(record.interaction.case_ref.expected_revision, |case| {
case.case_ref.expected_revision
});
let context = ResponseContext::click(
&self.input.actor,
&self.input.conversation_id,
self.input.turn_id,
current,
self.now,
);
let admission = self
.runtime
.interactions
.accept(context, &response)
.await
.map_err(|error| self.observed_rejection(error))?;
match admission {
ResponseAdmission::Accepted(accepted) => {
self.resolving_card = Some(accepted.interaction_id());
Ok(Some(*accepted))
}
ResponseAdmission::AlreadyAnswered(record) => {
self.replayed = Some(record);
Ok(None)
}
}
}
pub(super) fn observed_rejection(&self, error: InteractionError) -> OrchestratorError {
if let InteractionError::Rejected(InteractionRejection::Stale { .. }) = &error {
self.runtime
.observer
.observe_labeled(&Signal::InteractionStale, &SignalLabels::none());
}
OrchestratorError::Interaction(error)
}
pub(super) async fn load_conversation(&mut self) {
let window = self
.runtime
.config
.understanding
.transcript_turns
.unwrap_or(usize::MAX);
let Ok(turns) = self
.runtime
.stores
.conversations()
.load_recent_turns(self.account(), &self.input.conversation_id, window)
.await
else {
return;
};
for assistant in turns.iter().filter_map(|turn| turn.assistant.as_ref()) {
self.artifacts_shown
.extend(assistant.blocks.iter().filter_map(|block| match block {
ResponseBlock::Artifact(view) => Some(view.block_id.clone()),
_ => None,
}));
if !assistant.subjects.is_empty() {
self.carried_subjects = assistant.subjects.clone();
}
}
(self.recent, self.previous) = recent_messages(turns, self.input.turn_id);
}
pub(super) async fn gather_attachments(&mut self) {
let composer = &self.runtime.composer;
self.attachments = TurnAttachments::gather(
self.runtime.attachment_source.as_ref(),
self.runtime.config.attachments,
&self.input.turn_id,
&self.input.attachments,
|part| composer.carries(part),
)
.await;
}
pub(super) fn resolver(&self, answered: Option<&AcceptedInteraction>) -> TargetResolver {
build_resolver(
self.account(),
self.input.turn_id,
&self.runtime.workflows.definitions(),
&self.runtime.case_ids,
&self.cases,
self.input
.origin
.as_ref()
.map(|origin| &origin.origin_token)
.zip(self.origin_case.as_ref()),
blocking_summary(answered, &self.open_interactions),
)
}
}