use crate::backend::*;
use crate::content::ContentWriter;
use onlyne_acp::{ContentBlock, PromptOutcome};
use super::journal::{Journal, append, completion_head, drain};
use super::state::SessionEntry;
const REFUSAL_LIST_LIMIT: usize = 3;
fn settle_for(
stop_reason: &str,
answer: Option<&str>,
) -> (Option<TaskState>, Option<String>, String) {
let named = if stop_reason.is_empty() {
"(absent)".to_string()
} else {
stop_reason.to_string()
};
match stop_reason {
PromptOutcome::END_TURN => (None, None, named),
PromptOutcome::CANCELLED => (
Some(TaskState::Cancelled),
Some("the client cancelled this session".to_string()),
named,
),
PromptOutcome::REFUSAL | PromptOutcome::MAX_TOKENS | PromptOutcome::MAX_TURN_REQUESTS => (
Some(TaskState::Failed),
Some(format!("agent stopped the turn: {named}")),
named,
),
_ => (
Some(TaskState::Failed),
Some(match answer {
Some(_) => format!("agent ended the turn with stopReason {named:?}"),
None => format!("agent ended the turn with stopReason {named:?} and no answer"),
}),
named,
),
}
}
pub(super) fn run_turn(
entry: Arc<SessionEntry>,
sink: OutcomeSink,
content: ContentWriter,
prompt: String,
record: &'static str,
policy: &'static str,
) {
let task_id = entry.current_task();
let journal = Journal::new(&entry.workdir, &task_id, &entry.id, content);
journal.record(
record,
vec![
("task_id", Value::from(task_id.clone())),
("prompt", Value::from(prompt.clone())),
],
);
let events = entry.agent.subscribe();
let turn = entry
.agent
.prompt(&entry.id, vec![ContentBlock::text(&prompt)]);
let drained = drain(&events, &entry.id);
let refusals = take_refusals(&entry, policy);
for record in drained.lines {
journal.raw(record);
}
if !drained.log.is_empty() {
append(&journal.log, &drained.log);
}
let head = completion_head(&drained.message);
let (settled, note, stop_reason) = match &turn {
Ok(outcome) => settle_for(&outcome.stop_reason, head.as_deref()),
Err(error) => (
Some(TaskState::Failed),
Some(death_note(error, drained.exited.as_deref())),
"(error)".to_string(),
),
};
journal.record(
"turn",
vec![
("task_id", Value::from(task_id.clone())),
("stop_reason", Value::from(stop_reason)),
("head", head.clone().map(Value::from).unwrap_or(Value::Null)),
],
);
entry.turn.finish();
sink.push(SessionOutcome {
task_id,
outcome: settled,
head,
note,
refusals,
});
}
fn take_refusals(entry: &SessionEntry, policy: &str) -> Option<String> {
let mut refused = std::mem::take(&mut *entry.refusals.lock());
if refused.is_empty() {
return None;
}
let count = refused.len();
refused.truncate(REFUSAL_LIST_LIMIT);
Some(format!(
"{count} permission ask(s) refused (policy={policy}): {}",
refused.join("; ")
))
}
fn death_note(error: &anyhow::Error, exited: Option<&str>) -> String {
let detail = error.to_string();
match exited {
Some(exited) if !detail.contains(exited) => format!("{detail}; {exited}"),
_ => detail,
}
}