T1 — Wake automatique du propriétaire à la complétion. La complétion d'une tâche de fond est désormais livrée à l'agent propriétaire dès que la session accepte l'envoi (wake.rs : mark_completion_delivered au send accepté), via un drain de flux dédié (structured.rs : drain_reply_stream_with_readiness). L'inbox médiée enfile l'item sans démarrer de tour ni marquer l'agent busy (input/mod.rs : enqueue FIFO silencieux). Régression couverte (tests/agent_wake.rs, tests input/mod.rs). T3 — Tâches de fond projetées dans le read-model du panneau Work. AgentWorkState porte désormais background_tasks (VO AgentBackgroundTaskState), alimenté par un builder best-effort with_background_tasks(store) : union list_open_for_agent + dispatch des completions non livrées par owner_agent_id, erreur store => Vec vide (aucune régression live/busy/tickets). DTO Tauri backgroundTasks et wiring du BackgroundTaskStore côté state.rs. Le frontend, déjà câblé, affiche Cancel/Retry (mapping queued/waiting -> pending, tri sur updatedAtMs). Borne V1 : une tâche terminale déjà livrée n'est plus énumérable (Retry limité à la fenêtre non livrée). Tests : cargo build --workspace OK ; cargo test -p application / -p app-tauri / -p infrastructure verts ; frontend build + vitest verts. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1305 lines
38 KiB
Rust
1305 lines
38 KiB
Rust
//! Tests for the project work-state read model.
|
|
|
|
use std::collections::HashMap;
|
|
use std::future::Future;
|
|
use std::pin::Pin;
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
use async_trait::async_trait;
|
|
|
|
use application::{
|
|
ConversationLogProvider, ConversationPreviewStatus, GetProjectWorkState,
|
|
GetProjectWorkStateInput, HandoffProvider, LiveSessionKind, LiveSessions, StructuredSessions,
|
|
TerminalSessions, TicketWorkSource, TicketWorkStatus,
|
|
};
|
|
use domain::mailbox::{
|
|
AgentQueueSnapshot, MailboxError, PendingReply, QueuedTicketSnapshot, Ticket, TurnResolution,
|
|
};
|
|
use domain::ports::{
|
|
AgentContextStore, AgentSession, AgentSessionError, BackgroundTaskPortError,
|
|
BackgroundTaskStore, PtyHandle, ReplyStream, StoreError,
|
|
};
|
|
use domain::{
|
|
Agent, AgentBusyState, AgentId, AgentManifest, AgentOrigin, BackgroundTask, BackgroundTaskKind,
|
|
BackgroundTaskResult, BackgroundTaskState, BackgroundTaskWakePolicy, ConversationId,
|
|
ConversationLog, ConversationTurn, Handoff, HandoffStore, InputMediator, InputSource,
|
|
ManifestEntry, MarkdownDoc, NodeId, ProfileId, Project, ProjectId, ProjectPath, PtySize,
|
|
RemoteRef, SessionId, SessionKind, TaskId, TerminalSession, TicketId, TurnId, TurnRole,
|
|
};
|
|
use uuid::Uuid;
|
|
|
|
fn aid(n: u128) -> AgentId {
|
|
AgentId::from_uuid(Uuid::from_u128(n))
|
|
}
|
|
|
|
fn pid(n: u128) -> ProfileId {
|
|
ProfileId::from_uuid(Uuid::from_u128(n))
|
|
}
|
|
|
|
fn sid(n: u128) -> SessionId {
|
|
SessionId::from_uuid(Uuid::from_u128(n))
|
|
}
|
|
|
|
fn nid(n: u128) -> NodeId {
|
|
NodeId::from_uuid(Uuid::from_u128(n))
|
|
}
|
|
|
|
fn ticket_id(n: u128) -> domain::TicketId {
|
|
domain::TicketId::from_uuid(Uuid::from_u128(n))
|
|
}
|
|
|
|
fn task_id(n: u128) -> TaskId {
|
|
TaskId::from_uuid(Uuid::from_u128(n))
|
|
}
|
|
|
|
fn project() -> Project {
|
|
Project::new(
|
|
ProjectId::from_uuid(Uuid::from_u128(1)),
|
|
"demo",
|
|
ProjectPath::new("/tmp/idea-workstate-test").unwrap(),
|
|
RemoteRef::local(),
|
|
1_700_000_000_000,
|
|
)
|
|
.unwrap()
|
|
}
|
|
|
|
fn agent(n: u128, name: &str) -> Agent {
|
|
Agent::new(
|
|
aid(n),
|
|
name,
|
|
format!("agents/{name}.md"),
|
|
pid(100 + n),
|
|
AgentOrigin::Scratch,
|
|
false,
|
|
)
|
|
.unwrap()
|
|
}
|
|
|
|
fn manifest(agents: &[Agent]) -> AgentManifest {
|
|
AgentManifest::new(1, agents.iter().map(ManifestEntry::from_agent).collect()).unwrap()
|
|
}
|
|
|
|
#[derive(Clone)]
|
|
struct FakeContexts {
|
|
manifest: AgentManifest,
|
|
}
|
|
|
|
#[async_trait]
|
|
impl AgentContextStore for FakeContexts {
|
|
async fn read_context(
|
|
&self,
|
|
_project: &Project,
|
|
_agent: &AgentId,
|
|
) -> Result<MarkdownDoc, StoreError> {
|
|
Err(StoreError::NotFound)
|
|
}
|
|
|
|
async fn write_context(
|
|
&self,
|
|
_project: &Project,
|
|
_agent: &AgentId,
|
|
_md: &MarkdownDoc,
|
|
) -> Result<(), StoreError> {
|
|
Ok(())
|
|
}
|
|
|
|
async fn load_manifest(&self, _project: &Project) -> Result<AgentManifest, StoreError> {
|
|
Ok(self.manifest.clone())
|
|
}
|
|
|
|
async fn save_manifest(
|
|
&self,
|
|
_project: &Project,
|
|
_manifest: &AgentManifest,
|
|
) -> Result<(), StoreError> {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct FakeInput {
|
|
busy: Mutex<HashMap<AgentId, AgentBusyState>>,
|
|
}
|
|
|
|
impl FakeInput {
|
|
fn set_busy(&self, agent: AgentId, busy: AgentBusyState) {
|
|
self.busy.lock().unwrap().insert(agent, busy);
|
|
}
|
|
}
|
|
|
|
impl InputMediator for FakeInput {
|
|
fn enqueue(&self, _agent: AgentId, _ticket: Ticket) -> PendingReply {
|
|
let fut: Pin<Box<dyn Future<Output = Result<TurnResolution, MailboxError>> + Send>> =
|
|
Box::pin(async { Err(MailboxError::Cancelled) });
|
|
PendingReply::new(fut)
|
|
}
|
|
|
|
fn preempt(&self, _agent: AgentId) {}
|
|
|
|
fn mark_idle(&self, _agent: AgentId) {}
|
|
|
|
fn busy_state(&self, agent: AgentId) -> AgentBusyState {
|
|
self.busy
|
|
.lock()
|
|
.unwrap()
|
|
.get(&agent)
|
|
.copied()
|
|
.unwrap_or(AgentBusyState::Idle)
|
|
}
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct FakeQueue {
|
|
queues: Mutex<HashMap<AgentId, Vec<QueuedTicketSnapshot>>>,
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct FakeBackgroundTaskStore {
|
|
tasks: Mutex<Vec<BackgroundTask>>,
|
|
fail_open: Mutex<bool>,
|
|
fail_undelivered: Mutex<bool>,
|
|
}
|
|
|
|
impl FakeBackgroundTaskStore {
|
|
fn set_tasks(&self, tasks: Vec<BackgroundTask>) {
|
|
*self.tasks.lock().unwrap() = tasks;
|
|
}
|
|
|
|
fn fail_open(&self) {
|
|
*self.fail_open.lock().unwrap() = true;
|
|
}
|
|
|
|
fn fail_undelivered(&self) {
|
|
*self.fail_undelivered.lock().unwrap() = true;
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl BackgroundTaskStore for FakeBackgroundTaskStore {
|
|
async fn create(&self, task: &BackgroundTask) -> Result<(), BackgroundTaskPortError> {
|
|
self.tasks.lock().unwrap().push(task.clone());
|
|
Ok(())
|
|
}
|
|
|
|
async fn get(&self, id: TaskId) -> Result<Option<BackgroundTask>, BackgroundTaskPortError> {
|
|
Ok(self
|
|
.tasks
|
|
.lock()
|
|
.unwrap()
|
|
.iter()
|
|
.find(|task| task.id == id)
|
|
.cloned())
|
|
}
|
|
|
|
async fn save(&self, task: &BackgroundTask) -> Result<(), BackgroundTaskPortError> {
|
|
let mut tasks = self.tasks.lock().unwrap();
|
|
if let Some(existing) = tasks.iter_mut().find(|existing| existing.id == task.id) {
|
|
*existing = task.clone();
|
|
} else {
|
|
tasks.push(task.clone());
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
async fn list_open_for_agent(
|
|
&self,
|
|
agent_id: AgentId,
|
|
) -> Result<Vec<BackgroundTask>, BackgroundTaskPortError> {
|
|
if *self.fail_open.lock().unwrap() {
|
|
return Err(BackgroundTaskPortError::Store("open failed".to_owned()));
|
|
}
|
|
Ok(self
|
|
.tasks
|
|
.lock()
|
|
.unwrap()
|
|
.iter()
|
|
.filter(|task| task.owner_agent_id == agent_id && !task.is_terminal())
|
|
.cloned()
|
|
.collect())
|
|
}
|
|
|
|
async fn list_undelivered_completions(
|
|
&self,
|
|
) -> Result<Vec<BackgroundTask>, BackgroundTaskPortError> {
|
|
if *self.fail_undelivered.lock().unwrap() {
|
|
return Err(BackgroundTaskPortError::Store(
|
|
"undelivered failed".to_owned(),
|
|
));
|
|
}
|
|
Ok(self
|
|
.tasks
|
|
.lock()
|
|
.unwrap()
|
|
.iter()
|
|
.filter(|task| task.has_pending_completion_delivery())
|
|
.cloned()
|
|
.collect())
|
|
}
|
|
|
|
async fn mark_completion_delivered(
|
|
&self,
|
|
task_id: TaskId,
|
|
) -> Result<(), BackgroundTaskPortError> {
|
|
if let Some(task) = self
|
|
.tasks
|
|
.lock()
|
|
.unwrap()
|
|
.iter_mut()
|
|
.find(|task| task.id == task_id)
|
|
{
|
|
task.completion_delivered = true;
|
|
}
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
impl FakeQueue {
|
|
fn set(&self, agent: AgentId, tickets: Vec<QueuedTicketSnapshot>) {
|
|
self.queues.lock().unwrap().insert(agent, tickets);
|
|
}
|
|
}
|
|
|
|
impl AgentQueueSnapshot for FakeQueue {
|
|
fn queue_for(&self, agent: AgentId) -> Vec<QueuedTicketSnapshot> {
|
|
self.queues
|
|
.lock()
|
|
.unwrap()
|
|
.get(&agent)
|
|
.cloned()
|
|
.unwrap_or_default()
|
|
}
|
|
}
|
|
|
|
/// Configurable per-conversation outcome for the fake handoff store.
|
|
///
|
|
/// Absence is modelled by simply not configuring a conversation (see the `_` arm
|
|
/// of [`FakeHandoffStore::load`]), so it needs no dedicated variant.
|
|
#[derive(Clone)]
|
|
enum HandoffOutcome {
|
|
/// A readable handoff (status Ready).
|
|
Present(Handoff),
|
|
/// An unreadable handoff (status Partial/Unavailable depending on the log).
|
|
Error,
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct FakeHandoffStore {
|
|
outcomes: Mutex<HashMap<ConversationId, HandoffOutcome>>,
|
|
}
|
|
|
|
impl FakeHandoffStore {
|
|
fn set(&self, conversation: ConversationId, outcome: HandoffOutcome) {
|
|
self.outcomes.lock().unwrap().insert(conversation, outcome);
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl HandoffStore for FakeHandoffStore {
|
|
async fn load(&self, conversation: ConversationId) -> Result<Option<Handoff>, StoreError> {
|
|
match self.outcomes.lock().unwrap().get(&conversation) {
|
|
Some(HandoffOutcome::Present(handoff)) => Ok(Some(handoff.clone())),
|
|
Some(HandoffOutcome::Error) => Err(StoreError::Io("handoff unreadable".to_owned())),
|
|
// Unconfigured conversation ⇒ no handoff yet (never an error).
|
|
None => Ok(None),
|
|
}
|
|
}
|
|
|
|
async fn save(
|
|
&self,
|
|
_conversation: ConversationId,
|
|
_handoff: Handoff,
|
|
) -> Result<(), StoreError> {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
/// Configurable per-conversation outcome for the fake conversation log.
|
|
#[derive(Clone)]
|
|
enum LogOutcome {
|
|
/// A readable thread (possibly empty).
|
|
Turns(Vec<ConversationTurn>),
|
|
/// An unreadable thread.
|
|
Error,
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct FakeLog {
|
|
outcomes: Mutex<HashMap<ConversationId, LogOutcome>>,
|
|
}
|
|
|
|
impl FakeLog {
|
|
fn set(&self, conversation: ConversationId, outcome: LogOutcome) {
|
|
self.outcomes.lock().unwrap().insert(conversation, outcome);
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl ConversationLog for FakeLog {
|
|
async fn append(
|
|
&self,
|
|
_conversation: ConversationId,
|
|
_turn: ConversationTurn,
|
|
) -> Result<(), StoreError> {
|
|
Ok(())
|
|
}
|
|
|
|
async fn read(
|
|
&self,
|
|
_conversation: ConversationId,
|
|
_since: Option<TurnId>,
|
|
) -> Result<Vec<ConversationTurn>, StoreError> {
|
|
Ok(Vec::new())
|
|
}
|
|
|
|
async fn last(
|
|
&self,
|
|
conversation: ConversationId,
|
|
n: usize,
|
|
) -> Result<Vec<ConversationTurn>, StoreError> {
|
|
match self.outcomes.lock().unwrap().get(&conversation) {
|
|
Some(LogOutcome::Turns(turns)) => {
|
|
let start = turns.len().saturating_sub(n);
|
|
Ok(turns[start..].to_vec())
|
|
}
|
|
Some(LogOutcome::Error) => Err(StoreError::Io("log unreadable".to_owned())),
|
|
None => Ok(Vec::new()),
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Stateless providers handing the same fake stores back regardless of root.
|
|
struct FakeHandoffProvider(Arc<FakeHandoffStore>);
|
|
|
|
impl HandoffProvider for FakeHandoffProvider {
|
|
fn handoff_store_for(&self, _root: &ProjectPath) -> Option<Arc<dyn HandoffStore>> {
|
|
Some(Arc::clone(&self.0) as Arc<dyn HandoffStore>)
|
|
}
|
|
}
|
|
|
|
struct FakeLogProvider(Arc<FakeLog>);
|
|
|
|
impl ConversationLogProvider for FakeLogProvider {
|
|
fn conversation_log_for(&self, _root: &ProjectPath) -> Option<Arc<dyn ConversationLog>> {
|
|
Some(Arc::clone(&self.0) as Arc<dyn ConversationLog>)
|
|
}
|
|
}
|
|
|
|
fn turn(id: u128, conversation: ConversationId, role: TurnRole, text: &str) -> ConversationTurn {
|
|
ConversationTurn::new(
|
|
TurnId::from_uuid(Uuid::from_u128(id)),
|
|
conversation,
|
|
1_700,
|
|
InputSource::Human,
|
|
role,
|
|
text,
|
|
)
|
|
}
|
|
|
|
fn snapshot(
|
|
id: u128,
|
|
position: u32,
|
|
source: InputSource,
|
|
requester: &str,
|
|
task: &str,
|
|
) -> QueuedTicketSnapshot {
|
|
QueuedTicketSnapshot {
|
|
id: TicketId::from_uuid(Uuid::from_u128(id)),
|
|
source,
|
|
conversation: ConversationId::from_uuid(Uuid::from_u128(id + 1000)),
|
|
requester: requester.to_owned(),
|
|
task: task.to_owned(),
|
|
position,
|
|
}
|
|
}
|
|
|
|
struct FakeSession {
|
|
id: SessionId,
|
|
}
|
|
|
|
#[async_trait]
|
|
impl AgentSession for FakeSession {
|
|
fn id(&self) -> SessionId {
|
|
self.id
|
|
}
|
|
|
|
fn conversation_id(&self) -> Option<String> {
|
|
None
|
|
}
|
|
|
|
async fn send(&self, _prompt: &str) -> Result<ReplyStream, AgentSessionError> {
|
|
Ok(Box::new(std::iter::empty()))
|
|
}
|
|
|
|
async fn shutdown(&self) -> Result<(), AgentSessionError> {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
fn fake_session(id: SessionId) -> Arc<dyn AgentSession> {
|
|
Arc::new(FakeSession { id })
|
|
}
|
|
|
|
fn insert_pty(
|
|
sessions: &TerminalSessions,
|
|
session_id: SessionId,
|
|
agent_id: AgentId,
|
|
node_id: NodeId,
|
|
) {
|
|
sessions.insert(
|
|
PtyHandle { session_id },
|
|
TerminalSession::starting(
|
|
session_id,
|
|
node_id,
|
|
ProjectPath::new("/tmp/idea-workstate-test").unwrap(),
|
|
SessionKind::Agent { agent_id },
|
|
PtySize { rows: 24, cols: 80 },
|
|
),
|
|
);
|
|
}
|
|
|
|
struct Fixture {
|
|
usecase: GetProjectWorkState,
|
|
pty: Arc<TerminalSessions>,
|
|
structured: Arc<StructuredSessions>,
|
|
input: Arc<FakeInput>,
|
|
queue: Arc<FakeQueue>,
|
|
project: Project,
|
|
}
|
|
|
|
fn fixture(agents: &[Agent]) -> Fixture {
|
|
let pty = Arc::new(TerminalSessions::new());
|
|
let structured = Arc::new(StructuredSessions::new());
|
|
let live = Arc::new(LiveSessions::new(Arc::clone(&pty), Arc::clone(&structured)));
|
|
let input = Arc::new(FakeInput::default());
|
|
let input_port = Arc::clone(&input) as Arc<dyn InputMediator>;
|
|
let queue = Arc::new(FakeQueue::default());
|
|
let queue_port = Arc::clone(&queue) as Arc<dyn AgentQueueSnapshot>;
|
|
let usecase = GetProjectWorkState::new(
|
|
Arc::new(FakeContexts {
|
|
manifest: manifest(agents),
|
|
}),
|
|
live,
|
|
input_port,
|
|
queue_port,
|
|
);
|
|
Fixture {
|
|
usecase,
|
|
pty,
|
|
structured,
|
|
input,
|
|
queue,
|
|
project: project(),
|
|
}
|
|
}
|
|
|
|
/// Conversation id minted by [`snapshot`] for a ticket of the given `id`.
|
|
fn conv_of(id: u128) -> ConversationId {
|
|
ConversationId::from_uuid(Uuid::from_u128(id + 1000))
|
|
}
|
|
|
|
struct ConvFixture {
|
|
usecase: GetProjectWorkState,
|
|
queue: Arc<FakeQueue>,
|
|
input: Arc<FakeInput>,
|
|
handoffs: Arc<FakeHandoffStore>,
|
|
logs: Arc<FakeLog>,
|
|
project: Project,
|
|
}
|
|
|
|
struct BackgroundFixture {
|
|
usecase: GetProjectWorkState,
|
|
queue: Arc<FakeQueue>,
|
|
store: Arc<FakeBackgroundTaskStore>,
|
|
project: Project,
|
|
}
|
|
|
|
fn background_fixture(agents: &[Agent]) -> BackgroundFixture {
|
|
let pty = Arc::new(TerminalSessions::new());
|
|
let structured = Arc::new(StructuredSessions::new());
|
|
let live = Arc::new(LiveSessions::new(pty, structured));
|
|
let input = Arc::new(FakeInput::default());
|
|
let queue = Arc::new(FakeQueue::default());
|
|
let store = Arc::new(FakeBackgroundTaskStore::default());
|
|
let usecase = GetProjectWorkState::new(
|
|
Arc::new(FakeContexts {
|
|
manifest: manifest(agents),
|
|
}),
|
|
live,
|
|
input as Arc<dyn InputMediator>,
|
|
Arc::clone(&queue) as Arc<dyn AgentQueueSnapshot>,
|
|
)
|
|
.with_background_tasks(Arc::clone(&store) as Arc<dyn BackgroundTaskStore>);
|
|
BackgroundFixture {
|
|
usecase,
|
|
queue,
|
|
store,
|
|
project: project(),
|
|
}
|
|
}
|
|
|
|
fn background_task(
|
|
id: u128,
|
|
project_id: ProjectId,
|
|
owner: AgentId,
|
|
created_at_ms: u64,
|
|
state: BackgroundTaskState,
|
|
delivered: bool,
|
|
) -> BackgroundTask {
|
|
let base = BackgroundTask::new(
|
|
task_id(id),
|
|
project_id,
|
|
owner,
|
|
BackgroundTaskKind::Command {
|
|
label: format!("task-{id}"),
|
|
},
|
|
BackgroundTaskWakePolicy::WakeOwner,
|
|
created_at_ms,
|
|
None,
|
|
)
|
|
.unwrap();
|
|
match state {
|
|
BackgroundTaskState::Queued => base,
|
|
BackgroundTaskState::Running | BackgroundTaskState::Waiting => {
|
|
base.transition(state, created_at_ms + 10).unwrap()
|
|
}
|
|
BackgroundTaskState::Completed
|
|
| BackgroundTaskState::Failed
|
|
| BackgroundTaskState::Cancelled
|
|
| BackgroundTaskState::Expired => {
|
|
let result = match state {
|
|
BackgroundTaskState::Completed => BackgroundTaskResult::Success {
|
|
finished_at_ms: created_at_ms + 20,
|
|
exit_code: Some(0),
|
|
summary: "done".to_owned(),
|
|
stdout_tail: Some("stdout".to_owned()),
|
|
stderr_tail: None,
|
|
},
|
|
BackgroundTaskState::Failed => BackgroundTaskResult::Failure {
|
|
finished_at_ms: created_at_ms + 20,
|
|
exit_code: Some(2),
|
|
error: "failed".to_owned(),
|
|
stdout_tail: Some("out".to_owned()),
|
|
stderr_tail: Some("err".to_owned()),
|
|
},
|
|
BackgroundTaskState::Cancelled => BackgroundTaskResult::Cancelled {
|
|
finished_at_ms: created_at_ms + 20,
|
|
reason: "cancelled".to_owned(),
|
|
},
|
|
BackgroundTaskState::Expired => BackgroundTaskResult::Expired {
|
|
finished_at_ms: created_at_ms + 20,
|
|
reason: "expired".to_owned(),
|
|
},
|
|
BackgroundTaskState::Queued
|
|
| BackgroundTaskState::Running
|
|
| BackgroundTaskState::Waiting => unreachable!(),
|
|
};
|
|
let running = base
|
|
.transition(BackgroundTaskState::Running, created_at_ms + 10)
|
|
.unwrap();
|
|
let completed = running.complete(result).unwrap();
|
|
if delivered {
|
|
completed.mark_completion_delivered().unwrap()
|
|
} else {
|
|
completed
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Fixture wiring the best-effort conversation sources (handoff + log) onto the
|
|
/// read model, exposing the fake stores so each test configures their outcomes.
|
|
fn conv_fixture(agents: &[Agent]) -> ConvFixture {
|
|
let pty = Arc::new(TerminalSessions::new());
|
|
let structured = Arc::new(StructuredSessions::new());
|
|
let live = Arc::new(LiveSessions::new(pty, structured));
|
|
let input = Arc::new(FakeInput::default());
|
|
let queue = Arc::new(FakeQueue::default());
|
|
let handoffs = Arc::new(FakeHandoffStore::default());
|
|
let logs = Arc::new(FakeLog::default());
|
|
let usecase = GetProjectWorkState::new(
|
|
Arc::new(FakeContexts {
|
|
manifest: manifest(agents),
|
|
}),
|
|
live,
|
|
Arc::clone(&input) as Arc<dyn InputMediator>,
|
|
Arc::clone(&queue) as Arc<dyn AgentQueueSnapshot>,
|
|
)
|
|
.with_conversation_sources(
|
|
Arc::new(FakeHandoffProvider(Arc::clone(&handoffs))) as Arc<dyn HandoffProvider>,
|
|
Arc::new(FakeLogProvider(Arc::clone(&logs))) as Arc<dyn ConversationLogProvider>,
|
|
);
|
|
ConvFixture {
|
|
usecase,
|
|
queue,
|
|
input,
|
|
handoffs,
|
|
logs,
|
|
project: project(),
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_lists_manifest_agents_idle_without_live_sessions() {
|
|
let a = agent(10, "alpha");
|
|
let b = agent(20, "beta");
|
|
let f = fixture(&[a.clone(), b.clone()]);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(out.agents.len(), 2);
|
|
assert_eq!(out.agents[0].agent_id, a.id);
|
|
assert_eq!(out.agents[0].name, "alpha");
|
|
assert_eq!(out.agents[0].profile_id, a.profile_id);
|
|
assert_eq!(out.agents[0].live, None);
|
|
assert_eq!(out.agents[0].busy, AgentBusyState::Idle);
|
|
assert_eq!(out.agents[1].agent_id, b.id);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_attaches_live_pty_session_to_manifest_agent() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
insert_pty(&f.pty, sid(1), a.id, nid(100));
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let live = out.agents[0].live.unwrap();
|
|
assert_eq!(live.session_id, sid(1));
|
|
assert_eq!(live.node_id, nid(100));
|
|
assert_eq!(live.kind, LiveSessionKind::Pty);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_attaches_live_structured_session_to_manifest_agent() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
f.structured.insert(fake_session(sid(2)), a.id, nid(200));
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let live = out.agents[0].live.unwrap();
|
|
assert_eq!(live.session_id, sid(2));
|
|
assert_eq!(live.node_id, nid(200));
|
|
assert_eq!(live.kind, LiveSessionKind::Structured);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_includes_busy_state_from_input_mediator() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
let busy = AgentBusyState::Busy {
|
|
ticket: ticket_id(7),
|
|
since_ms: 1_234,
|
|
};
|
|
f.input.set_busy(a.id, busy);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(out.agents[0].busy, busy);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_ignores_live_agents_absent_from_manifest() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
insert_pty(&f.pty, sid(99), aid(999), nid(999));
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(out.agents.len(), 1);
|
|
assert_eq!(out.agents[0].agent_id, a.id);
|
|
assert_eq!(out.agents[0].live, None);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_agent_without_queue_has_no_tickets() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert!(out.agents[0].tickets.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_projects_open_and_undelivered_background_tasks_by_agent() {
|
|
let a = agent(10, "alpha");
|
|
let b = agent(20, "beta");
|
|
let f = background_fixture(&[a.clone(), b.clone()]);
|
|
let other_project = ProjectId::from_uuid(Uuid::from_u128(999));
|
|
f.store.set_tasks(vec![
|
|
background_task(
|
|
2,
|
|
f.project.id,
|
|
a.id,
|
|
2_000,
|
|
BackgroundTaskState::Running,
|
|
false,
|
|
),
|
|
background_task(
|
|
1,
|
|
f.project.id,
|
|
a.id,
|
|
1_000,
|
|
BackgroundTaskState::Failed,
|
|
false,
|
|
),
|
|
background_task(
|
|
3,
|
|
f.project.id,
|
|
a.id,
|
|
3_000,
|
|
BackgroundTaskState::Completed,
|
|
true,
|
|
),
|
|
background_task(
|
|
4,
|
|
other_project,
|
|
a.id,
|
|
4_000,
|
|
BackgroundTaskState::Running,
|
|
false,
|
|
),
|
|
]);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let alpha_tasks = &out.agents[0].background_tasks;
|
|
assert_eq!(alpha_tasks.len(), 2);
|
|
assert_eq!(alpha_tasks[0].task_id, task_id(1));
|
|
assert_eq!(alpha_tasks[0].state, BackgroundTaskState::Failed);
|
|
assert_eq!(alpha_tasks[0].exit_code, Some(2));
|
|
assert_eq!(alpha_tasks[0].summary.as_deref(), Some("failed"));
|
|
assert_eq!(alpha_tasks[0].stdout_tail.as_deref(), Some("out"));
|
|
assert_eq!(alpha_tasks[0].stderr_tail.as_deref(), Some("err"));
|
|
assert_eq!(alpha_tasks[1].task_id, task_id(2));
|
|
assert_eq!(alpha_tasks[1].state, BackgroundTaskState::Running);
|
|
assert!(out.agents[1].background_tasks.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_background_task_store_error_degrades_without_dropping_tickets() {
|
|
let a = agent(10, "alpha");
|
|
let f = background_fixture(std::slice::from_ref(&a));
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::Human, "Human", "queued")],
|
|
);
|
|
f.store.set_tasks(vec![background_task(
|
|
1,
|
|
f.project.id,
|
|
a.id,
|
|
1_000,
|
|
BackgroundTaskState::Running,
|
|
false,
|
|
)]);
|
|
f.store.fail_open();
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert!(out.agents[0].background_tasks.is_empty());
|
|
assert_eq!(out.agents[0].tickets.len(), 1);
|
|
assert_eq!(out.agents[0].tickets[0].ticket_id, ticket_id(1));
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_undelivered_completion_error_degrades_background_tasks_only() {
|
|
let a = agent(10, "alpha");
|
|
let f = background_fixture(std::slice::from_ref(&a));
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::Human, "Human", "queued")],
|
|
);
|
|
f.store.set_tasks(vec![background_task(
|
|
1,
|
|
f.project.id,
|
|
a.id,
|
|
1_000,
|
|
BackgroundTaskState::Running,
|
|
false,
|
|
)]);
|
|
f.store.fail_undelivered();
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert!(out.agents[0].background_tasks.is_empty());
|
|
assert_eq!(out.agents[0].tickets.len(), 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_lists_two_tickets_in_fifo_order() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![
|
|
snapshot(1, 0, InputSource::agent(from), "Main", "first"),
|
|
snapshot(2, 1, InputSource::agent(from), "Main", "second"),
|
|
],
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let tickets = &out.agents[0].tickets;
|
|
assert_eq!(tickets.len(), 2);
|
|
assert_eq!(tickets[0].ticket_id, ticket_id(1));
|
|
assert_eq!(tickets[0].position, 0);
|
|
assert_eq!(tickets[1].ticket_id, ticket_id(2));
|
|
assert_eq!(tickets[1].position, 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_marks_busy_head_in_progress_and_rest_queued() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![
|
|
snapshot(1, 0, InputSource::agent(from), "Main", "first"),
|
|
snapshot(2, 1, InputSource::agent(from), "Main", "second"),
|
|
],
|
|
);
|
|
f.input.set_busy(
|
|
a.id,
|
|
AgentBusyState::Busy {
|
|
ticket: ticket_id(1),
|
|
since_ms: 5,
|
|
},
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let tickets = &out.agents[0].tickets;
|
|
assert_eq!(tickets[0].status, TicketWorkStatus::InProgress);
|
|
assert_eq!(tickets[1].status, TicketWorkStatus::Queued);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_marks_all_queued_when_agent_idle() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "first")],
|
|
);
|
|
// Agent left Idle (default): even the head is only queued, not in-progress.
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(out.agents[0].tickets[0].status, TicketWorkStatus::Queued);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_ignores_queue_for_agent_absent_from_manifest() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
// A queue exists for an agent that is not in the manifest: it must not surface.
|
|
f.queue.set(
|
|
aid(999),
|
|
vec![snapshot(1, 0, InputSource::Human, "User", "ghost")],
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(out.agents.len(), 1);
|
|
assert!(out.agents[0].tickets.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_truncates_task_preview_and_keeps_original_len() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
// 400 chars with runs of whitespace to normalise; well over the 160 cap.
|
|
let long_task = format!("start{}end", " x ".repeat(80));
|
|
let original_len = long_task.chars().count();
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::Human, "User", &long_task)],
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let ticket = &out.agents[0].tickets[0];
|
|
assert_eq!(ticket.task_preview.chars().count(), 160);
|
|
assert!(!ticket.task_preview.contains(" "), "whitespace normalised");
|
|
assert_eq!(ticket.task_len, original_len);
|
|
assert!(ticket.task_len > 160);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_maps_human_and_agent_ticket_sources() {
|
|
let a = agent(10, "alpha");
|
|
let f = fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![
|
|
snapshot(1, 0, InputSource::Human, "User", "from human"),
|
|
snapshot(2, 1, InputSource::agent(from), "Main", "from agent"),
|
|
],
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let tickets = &out.agents[0].tickets;
|
|
assert_eq!(tickets[0].source, TicketWorkSource::Human);
|
|
assert_eq!(tickets[0].requester_label, "User");
|
|
assert_eq!(
|
|
tickets[1].source,
|
|
TicketWorkSource::Agent { agent_id: from }
|
|
);
|
|
assert_eq!(tickets[1].requester_label, "Main");
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Lot C — conversation summaries (best-effort, read-only)
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[tokio::test]
|
|
async fn workstate_has_no_conversations_without_tickets() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert!(out.conversations.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_dedups_conversation_ids_from_tickets() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
let shared = ConversationId::from_uuid(Uuid::from_u128(7777));
|
|
// Two tickets pointing at the *same* conversation must yield one summary.
|
|
let mut t1 = snapshot(1, 0, InputSource::agent(from), "Main", "first");
|
|
t1.conversation = shared;
|
|
let mut t2 = snapshot(2, 1, InputSource::agent(from), "Main", "second");
|
|
t2.conversation = shared;
|
|
f.queue.set(a.id, vec![t1, t2]);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
assert_eq!(out.conversations.len(), 1);
|
|
assert_eq!(out.conversations[0].conversation_id, shared);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_conversation_ready_from_handoff() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "task")],
|
|
);
|
|
let conv = conv_of(1);
|
|
let up_to = TurnId::from_uuid(Uuid::from_u128(500));
|
|
f.handoffs.set(
|
|
conv,
|
|
HandoffOutcome::Present(Handoff::new(
|
|
"Résumé du fil",
|
|
up_to,
|
|
Some("Livrer le lot C".to_owned()),
|
|
)),
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let summary = &out.conversations[0];
|
|
assert_eq!(summary.conversation_id, conv);
|
|
assert_eq!(summary.status, ConversationPreviewStatus::Ready);
|
|
assert_eq!(summary.summary_preview.as_deref(), Some("Résumé du fil"));
|
|
assert_eq!(summary.summary_len, "Résumé du fil".chars().count());
|
|
assert_eq!(
|
|
summary.objective_preview.as_deref(),
|
|
Some("Livrer le lot C")
|
|
);
|
|
assert_eq!(summary.up_to, Some(up_to));
|
|
assert!(
|
|
summary.recent_turns.is_empty(),
|
|
"Ready uses the handoff, not turns"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_conversation_missing_with_bounded_recent_turns() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "task")],
|
|
);
|
|
let conv = conv_of(1);
|
|
// No handoff (Absent) but a readable log of 5 turns ⇒ Missing + last 3.
|
|
f.logs.set(
|
|
conv,
|
|
LogOutcome::Turns(vec![
|
|
turn(1, conv, TurnRole::Prompt, "t1"),
|
|
turn(2, conv, TurnRole::Response, "t2"),
|
|
turn(3, conv, TurnRole::Prompt, "t3"),
|
|
turn(4, conv, TurnRole::Response, "t4"),
|
|
turn(5, conv, TurnRole::Prompt, "t5"),
|
|
]),
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let summary = &out.conversations[0];
|
|
assert_eq!(summary.status, ConversationPreviewStatus::Missing);
|
|
assert!(summary.summary_preview.is_none());
|
|
assert_eq!(summary.summary_len, 0);
|
|
assert_eq!(summary.recent_turns.len(), 3, "bounded to RECENT_TURNS_MAX");
|
|
// The *last* three turns, in order.
|
|
assert_eq!(summary.recent_turns[0].text_preview, "t3");
|
|
assert_eq!(summary.recent_turns[2].text_preview, "t5");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_conversation_partial_when_handoff_errors_but_log_ok() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "task")],
|
|
);
|
|
let conv = conv_of(1);
|
|
f.handoffs.set(conv, HandoffOutcome::Error);
|
|
f.logs.set(
|
|
conv,
|
|
LogOutcome::Turns(vec![turn(1, conv, TurnRole::Response, "only turn")]),
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let summary = &out.conversations[0];
|
|
assert_eq!(summary.status, ConversationPreviewStatus::Partial);
|
|
assert!(summary.summary_preview.is_none(), "no summary on partial");
|
|
assert_eq!(summary.recent_turns.len(), 1);
|
|
assert_eq!(summary.recent_turns[0].text_preview, "only turn");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_conversation_unavailable_when_handoff_and_log_fail() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "task")],
|
|
);
|
|
let conv = conv_of(1);
|
|
f.handoffs.set(conv, HandoffOutcome::Error);
|
|
f.logs.set(conv, LogOutcome::Error);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let summary = &out.conversations[0];
|
|
assert_eq!(summary.status, ConversationPreviewStatus::Unavailable);
|
|
assert!(summary.summary_preview.is_none());
|
|
assert!(summary.recent_turns.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_preview_failure_preserves_agents_live_busy_tickets() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "task")],
|
|
);
|
|
let conv = conv_of(1);
|
|
// Both sources KO for this conversation ⇒ Unavailable, but the rest stands.
|
|
f.handoffs.set(conv, HandoffOutcome::Error);
|
|
f.logs.set(conv, LogOutcome::Error);
|
|
let busy = AgentBusyState::Busy {
|
|
ticket: ticket_id(1),
|
|
since_ms: 9,
|
|
};
|
|
f.input.set_busy(a.id, busy);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
// Agents/busy/tickets fully intact despite the failed preview.
|
|
assert_eq!(out.agents.len(), 1);
|
|
assert_eq!(out.agents[0].busy, busy);
|
|
assert_eq!(out.agents[0].tickets.len(), 1);
|
|
assert_eq!(
|
|
out.agents[0].tickets[0].status,
|
|
TicketWorkStatus::InProgress
|
|
);
|
|
assert_eq!(
|
|
out.conversations[0].status,
|
|
ConversationPreviewStatus::Unavailable
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_conversation_previews_truncated_and_normalised() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "task")],
|
|
);
|
|
let conv = conv_of(1);
|
|
// Whitespace runs to normalise; lengths well over each cap (480 / 160).
|
|
// Each " s " collapses to one " s" (2 chars) ⇒ pick counts past the caps.
|
|
let long_summary = format!("start{}end", " s ".repeat(300));
|
|
let long_objective = format!("goal{}done", " o ".repeat(120));
|
|
f.handoffs.set(
|
|
conv,
|
|
HandoffOutcome::Present(Handoff::new(
|
|
long_summary.clone(),
|
|
TurnId::from_uuid(Uuid::from_u128(1)),
|
|
Some(long_objective),
|
|
)),
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let summary = &out.conversations[0];
|
|
let preview = summary.summary_preview.as_deref().unwrap();
|
|
assert_eq!(preview.chars().count(), 480, "summary capped");
|
|
assert!(!preview.contains(" "), "summary whitespace normalised");
|
|
assert_eq!(
|
|
summary.summary_len,
|
|
long_summary.chars().count(),
|
|
"original length preserved"
|
|
);
|
|
let objective = summary.objective_preview.as_deref().unwrap();
|
|
assert_eq!(objective.chars().count(), 160, "objective capped");
|
|
assert!(!objective.contains(" "), "objective whitespace normalised");
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn workstate_conversation_turn_text_preview_truncated() {
|
|
let a = agent(10, "alpha");
|
|
let f = conv_fixture(std::slice::from_ref(&a));
|
|
let from = aid(20);
|
|
f.queue.set(
|
|
a.id,
|
|
vec![snapshot(1, 0, InputSource::agent(from), "Main", "task")],
|
|
);
|
|
let conv = conv_of(1);
|
|
let long_text = format!("begin{}fin", " w ".repeat(120));
|
|
f.logs.set(
|
|
conv,
|
|
LogOutcome::Turns(vec![turn(1, conv, TurnRole::Prompt, &long_text)]),
|
|
);
|
|
|
|
let out = f
|
|
.usecase
|
|
.execute(GetProjectWorkStateInput { project: f.project })
|
|
.await
|
|
.unwrap();
|
|
|
|
let turn_preview = &out.conversations[0].recent_turns[0];
|
|
assert_eq!(
|
|
turn_preview.text_preview.chars().count(),
|
|
220,
|
|
"turn capped"
|
|
);
|
|
assert!(!turn_preview.text_preview.contains(" "), "normalised");
|
|
assert_eq!(turn_preview.text_len, long_text.chars().count());
|
|
}
|