use turnframe_core::flow::WritingStage;
use turnframe_core::knowledge::{KnowledgeChunk, KnowledgeError, KnowledgeRequest};
use turnframe_core::locale::LocalizedText;
use turnframe_core::plan::AnswerBasis;
use turnframe_core::reduce::{AnswerTask, SourcePolicy};
use turnframe_core::response::{
AnswerStatus, AnsweredEnumeration, AnsweredValue, GeneratedAnswer, NarratableFact,
};
use super::{Composer, CompositionInput};
use crate::narrate::Narrator;
use crate::narrate::tasks::{AnswerInput, AnswerKind, Answered};
pub(super) struct Answers<'a> {
pub composer: &'a Composer,
pub narrator: &'a Narrator<'a>,
pub input: &'a CompositionInput<'a>,
}
fn block_id_of(task: &AnswerTask) -> turnframe_core::ids::BlockId {
turnframe_core::ids::BlockId::from(format!("answer:{}", task.question_id))
}
impl Answers<'_> {
pub async fn all(&self) -> Vec<GeneratedAnswer> {
futures::future::join_all(self.input.answer_tasks.iter().map(|task| self.one(task))).await
}
async fn one(&self, task: &AnswerTask) -> GeneratedAnswer {
let copy = &self.composer.copy;
if !task.enumerations.is_empty() && task.basis == AnswerBasis::GeneralDomainKnowledge {
return self.enumerated(task);
}
if !self.composer.narrates(self.input) {
return self.unanswered(task, AnswerStatus::Unsupported, ©.answer_unsupported);
}
let chunks = match self.sources(task).await {
Ok(chunks) => chunks,
Err(status) => {
let text = match status {
AnswerStatus::SourceUnavailable => ©.answer_source_unavailable,
_ => ©.answer_unsupported,
};
return self.unanswered(task, status, text);
}
};
let facts = self.facts(task, &chunks);
let guidance = self.composer.guidance(
self.input,
&task
.case_refs
.iter()
.map(|case_ref| case_ref.key())
.collect::<Vec<_>>(),
WritingStage::Answer,
);
let answer = AnswerInput {
question: &task.question,
locale: self.input.turn.locale.as_str(),
tone: self.composer.narration.tone,
asked_before: self
.input
.recent
.iter()
.rev()
.find(|message| message.role == crate::conversation::TranscriptRole::User)
.map(|message| message.text.as_str())
.filter(|_| task.continues_previous && self.input.preceding_reply.is_some()),
previous: self
.input
.preceding_reply
.filter(|_| task.continues_previous),
facts: &facts,
guidance: &guidance,
attachments: &self.input.attachments,
};
match self
.narrator
.answer(task.question_id.as_str(), &answer)
.await
{
Ok(Answered {
kind: AnswerKind::Answered,
text,
}) => GeneratedAnswer {
block_id: block_id_of(task),
question_id: Some(task.question_id.clone()),
text,
basis: task.basis,
status: AnswerStatus::Answered,
facts_used: facts,
citations: chunks.into_iter().map(|chunk| chunk.citation).collect(),
enumerations: Vec::new(),
},
Ok(Answered {
kind: AnswerKind::CannotAnswer,
..
}) => self.unanswered(task, AnswerStatus::Unsupported, ©.answer_unsupported),
Err(code) if code == "too_long" => {
self.unanswered(task, AnswerStatus::Withheld, ©.answer_too_long)
}
Err(_) if !task.capabilities.is_empty() => self.capabilities(task, facts),
Err(_) => self.unanswered(task, AnswerStatus::NotWritten, ©.answer_not_written),
}
}
fn capabilities(&self, task: &AnswerTask, facts: Vec<NarratableFact>) -> GeneratedAnswer {
let text = task
.capabilities
.iter()
.map(|capability| format!("- {}", capability.summary))
.collect::<Vec<_>>()
.join("\n");
GeneratedAnswer {
block_id: block_id_of(task),
question_id: Some(task.question_id.clone()),
text,
basis: task.basis,
status: AnswerStatus::Answered,
facts_used: facts,
citations: Vec::new(),
enumerations: Vec::new(),
}
}
async fn sources(&self, task: &AnswerTask) -> Result<Vec<KnowledgeChunk>, AnswerStatus> {
let sources_are_the_answer =
task.basis == AnswerBasis::GeneralDomainKnowledge && task.capabilities.is_empty();
if task.required_sources == SourcePolicy::NoRetrieval
|| !(sources_are_the_answer || self.composer.knowledge.is_some())
{
return Ok(Vec::new());
}
match self.retrieve(task).await {
Ok(found) if found.is_empty() && sources_are_the_answer => {
Err(AnswerStatus::Unsupported)
}
Ok(found) => Ok(found),
Err(error) => {
tracing::warn!(
target: "turnframe.compose",
error = %error,
"knowledge retrieval failed"
);
if sources_are_the_answer {
Err(AnswerStatus::SourceUnavailable)
} else {
Ok(Vec::new())
}
}
}
}
async fn retrieve(&self, task: &AnswerTask) -> Result<Vec<KnowledgeChunk>, KnowledgeError> {
let Some(knowledge) = self.composer.knowledge.as_ref() else {
return Err(KnowledgeError::Unavailable {
source_id: "none_configured".to_owned(),
});
};
let turn = self.input.turn;
knowledge
.retrieve(KnowledgeRequest {
account_id: turn.actor.account_id.clone(),
query: task.question.clone(),
locale: turn.locale.clone(),
case_refs: task.case_refs.clone(),
source_policy: task.required_sources,
max_chunks: self.composer.max_chunks,
as_of: None,
})
.await
}
fn facts(&self, task: &AnswerTask, chunks: &[KnowledgeChunk]) -> Vec<NarratableFact> {
let locale = &self.input.turn.locale;
let mut facts = Vec::new();
for view in self.input.views {
if !task
.case_refs
.iter()
.any(|case_ref| case_ref.key() == view.case_ref.key())
{
continue;
}
facts.push(NarratableFact::Record {
case_ref: view.case_ref.clone(),
label: self
.input
.case_labels
.iter()
.find(|named| named.case_ref.key() == view.case_ref.key())
.map(|named| named.label.clone()),
phase: view.phase.clone(),
outcome: view.outcome.clone(),
});
facts.extend(view.state.iter().map(|held| NarratableFact::StateValue {
case_ref: view.case_ref.clone(),
field: held.field.clone(),
value: held.value.clone(),
}));
facts.extend(view.obligations.iter().map(|obligation| {
NarratableFact::ObligationOpen {
case_ref: view.case_ref.clone(),
obligation: obligation.value.clone(),
sentence: obligation
.sentence
.as_ref()
.map(|text| text.resolve(locale).to_owned()),
relevance: turnframe_core::response::FactRelevance::ThisTurn,
}
}));
}
facts.extend(task.capabilities.iter().map(|capability| {
NarratableFact::OperationAvailable {
workflow: capability.workflow.clone(),
operation: capability.operation.clone(),
summary: capability.summary.clone(),
}
}));
facts.extend(chunks.iter().map(|chunk| NarratableFact::Knowledge {
chunk_id: chunk.chunk_id.clone(),
source_id: chunk.source_id.clone(),
text: chunk.text.clone(),
}));
facts
}
fn unanswered(
&self,
task: &AnswerTask,
status: AnswerStatus,
text: &LocalizedText,
) -> GeneratedAnswer {
GeneratedAnswer {
block_id: block_id_of(task),
question_id: Some(task.question_id.clone()),
text: text.resolve(&self.input.turn.locale).to_owned(),
basis: task.basis,
status,
facts_used: Vec::new(),
citations: Vec::new(),
enumerations: Vec::new(),
}
}
fn enumerated(&self, task: &AnswerTask) -> GeneratedAnswer {
let locale = &self.input.turn.locale;
let enumerations: Vec<AnsweredEnumeration> = task
.enumerations
.iter()
.map(|enumeration| AnsweredEnumeration {
subject: enumeration.subject.as_str().to_owned(),
preamble: enumeration
.preamble
.as_ref()
.map(|text| text.resolve(locale).to_owned()),
values: enumeration
.values
.iter()
.map(|value| AnsweredValue {
id: value.id.clone(),
label: value.label.resolve(locale).to_owned(),
})
.collect(),
})
.collect();
let text = enumerations
.iter()
.map(|enumeration| {
let labels: Vec<&str> = enumeration
.values
.iter()
.map(|value| value.label.as_str())
.collect();
match enumeration.preamble.as_deref() {
Some(preamble) => format!(
"{}: {}.",
preamble.trim_end().trim_end_matches(['.', ':']),
labels.join(", ")
),
None => format!("{}.", labels.join(", ")),
}
})
.collect::<Vec<_>>()
.join("\n");
GeneratedAnswer {
block_id: block_id_of(task),
question_id: Some(task.question_id.clone()),
text,
basis: task.basis,
status: AnswerStatus::Answered,
facts_used: Vec::new(),
citations: Vec::new(),
enumerations,
}
}
}