diff --git a/crates/app-tauri/src/commands.rs b/crates/app-tauri/src/commands.rs index 2b982dc..0d2c4b7 100644 --- a/crates/app-tauri/src/commands.rs +++ b/crates/app-tauri/src/commands.rs @@ -31,28 +31,27 @@ use crate::dto::{ parse_agent_id, parse_close_terminal, parse_delete_profile, parse_layout_id, parse_memory_slug, parse_node_id, parse_profile_id, parse_project_id, parse_session_id, parse_skill_id, parse_task_id, parse_template_id, parse_ticket_id, AgentDriftListDto, AgentDto, AgentListDto, - BackgroundTaskDto, AssignSkillRequestDto, AttachLiveAgentRequestDto, AttachLiveAgentResponseDto, - ChangeAgentProfileDto, ChangeAgentProfileRequestDto, ConfigureProfilesRequestDto, - ConversationDetailsDto, CreateAgentFromTemplateRequestDto, CreateAgentRequestDto, - CreateLayoutRequestDto, CreateLayoutResultDto, CreateMemoryRequestDto, CreateProjectRequestDto, - CreateSkillRequestDto, CreateTemplateRequestDto, DeleteLayoutRequestDto, DeleteLayoutResultDto, - DeliveredDelegationRequestDto, DetectProfilesRequestDto, DetectProfilesResponseDto, - EffectivePermissionsDto, EmbedderEnginesDto, EmbedderProfileDto, EmbedderProfileListDto, - ErrorDto, FirstRunStateDto, FrontAttachedRequestDto, GitBranchesDto, GitCheckoutRequestDto, - GitCommitDto, GitCommitListDto, GitCommitRequestDto, GitStageRequestDto, GitStatusListDto, - GraphCommitListDto, HealthRequestDto, HealthResponseDto, InspectConversationRequestDto, - InterruptAgentRequestDto, LaunchAgentRequestDto, LayoutDto, LayoutOperationDto, ListLayoutsDto, - LiveAgentListDto, MemoryDto, MemoryIndexDto, MemoryLinksDto, MemoryListDto, - OpenTerminalRequestDto, ProfileDto, ProfileListDto, ProjectDto, ProjectListDto, - ProjectPermissionsDto, ProjectWorkStateDto, ReadAgentContextResponseDto, - ReadConversationPageRequestDto, ReattachChatDto, ReattachResultDto, RecallMemoryRequestDto, - RenameLayoutRequestDto, ReplyChunk, ResizeTerminalRequestDto, - ResolveAgentPermissionsRequestDto, ResumableAgentListDto, SaveEmbedderProfileRequestDto, - SaveProfileRequestDto, SetActiveLayoutRequestDto, SetActiveLayoutResultDto, SkillDto, - SkillListDto, StopLiveAgentRequestDto, StopLiveAgentResponseDto, - SyncAgentWithTemplateRequestDto, SyncResultDto, TemplateDto, TemplateListDto, - TerminalClosedDto, TerminalSessionDto, TurnPageDto, UnassignSkillRequestDto, + BackgroundTaskDto, ChangeAgentProfileDto, ChangeAgentProfileRequestDto, + ConfigureProfilesRequestDto, ConversationDetailsDto, CreateAgentFromTemplateRequestDto, + CreateAgentRequestDto, CreateLayoutRequestDto, CreateLayoutResultDto, CreateMemoryRequestDto, + CreateProjectRequestDto, CreateSkillRequestDto, CreateTemplateRequestDto, + DeleteLayoutRequestDto, DeleteLayoutResultDto, DeliveredDelegationRequestDto, + DetectProfilesRequestDto, DetectProfilesResponseDto, EffectivePermissionsDto, + EmbedderEnginesDto, EmbedderProfileDto, EmbedderProfileListDto, ErrorDto, FirstRunStateDto, + FrontAttachedRequestDto, GitBranchesDto, GitCheckoutRequestDto, GitCommitDto, GitCommitListDto, + GitCommitRequestDto, GitStageRequestDto, GitStatusListDto, GraphCommitListDto, + HealthRequestDto, HealthResponseDto, InspectConversationRequestDto, InterruptAgentRequestDto, + LaunchAgentRequestDto, LayoutDto, LayoutOperationDto, ListLayoutsDto, LiveAgentListDto, + MemoryDto, MemoryIndexDto, MemoryLinksDto, MemoryListDto, OpenTerminalRequestDto, ProfileDto, + ProfileListDto, ProjectDto, ProjectListDto, ProjectPermissionsDto, ProjectWorkStateDto, + ReadAgentContextResponseDto, ReadConversationPageRequestDto, ReattachChatDto, + ReattachResultDto, RecallMemoryRequestDto, RenameLayoutRequestDto, ReplyChunk, + ResizeTerminalRequestDto, ResolveAgentPermissionsRequestDto, ResumableAgentListDto, + SaveEmbedderProfileRequestDto, SaveProfileRequestDto, SetActiveLayoutRequestDto, + SetActiveLayoutResultDto, SkillDto, SkillListDto, StopLiveAgentRequestDto, + StopLiveAgentResponseDto, SyncAgentWithTemplateRequestDto, SyncResultDto, TemplateDto, + TemplateListDto, TerminalClosedDto, TerminalSessionDto, TurnPageDto, UnassignSkillRequestDto, UpdateAgentContextRequestDto, UpdateAgentPermissionsRequestDto, UpdateMemoryRequestDto, UpdateProjectContextRequestDto, UpdateProjectPermissionsRequestDto, UpdateSkillRequestDto, UpdateTemplateRequestDto, WriteTerminalRequestDto, diff --git a/crates/app-tauri/src/dto.rs b/crates/app-tauri/src/dto.rs index 55f4849..0bb2aea 100644 --- a/crates/app-tauri/src/dto.rs +++ b/crates/app-tauri/src/dto.rs @@ -9,11 +9,12 @@ use serde::{Deserialize, Serialize}; use application::{ - AgentTicketState, AppError, AttachLiveAgentOutput, ConversationPreviewStatus, - ConversationTurnWorkPreview, ConversationWorkSummary, CreateProjectInput, CreateProjectOutput, - GitGraphOutput, HealthInput, HealthReport, LayoutKind, ListProjectsOutput, LiveSessionKind, - LiveSessionSnapshot, OpenProjectOutput, ProjectWorkState, StopLiveAgentOutput, - TicketWorkSource, TicketWorkStatus, TurnPage, TurnSource, TurnView, + AgentBackgroundTaskState, AgentTicketState, AppError, AttachLiveAgentOutput, + BackgroundTaskKindLabel, ConversationPreviewStatus, ConversationTurnWorkPreview, + ConversationWorkSummary, CreateProjectInput, CreateProjectOutput, GitGraphOutput, HealthInput, + HealthReport, LayoutKind, ListProjectsOutput, LiveSessionKind, LiveSessionSnapshot, + OpenProjectOutput, ProjectWorkState, StopLiveAgentOutput, TicketWorkSource, TicketWorkStatus, + TurnPage, TurnSource, TurnView, }; use domain::{AgentBusyState, PageCursor, PageDirection, Project, ProjectId, TurnRole}; @@ -1607,6 +1608,50 @@ impl From for AgentTicketStateDto { } } +/// One background task owned by an agent in the Work panel read model. +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct AgentBackgroundTaskStateDto { + /// Stable task id (UUID string). + pub task_id: String, + /// Kind discriminant. + pub kind: String, + /// Lifecycle state. + pub state: String, + /// Process exit code, when the terminal result carries one. + #[serde(skip_serializing_if = "Option::is_none")] + pub exit_code: Option, + /// Human-readable summary / error / reason of the terminal result. + #[serde(skip_serializing_if = "Option::is_none")] + pub summary: Option, + /// Bounded stdout tail. + #[serde(skip_serializing_if = "Option::is_none")] + pub stdout_tail: Option, + /// Bounded stderr tail. + #[serde(skip_serializing_if = "Option::is_none")] + pub stderr_tail: Option, + /// Creation timestamp, epoch milliseconds. + pub created_at_ms: u64, + /// Last update timestamp, epoch milliseconds. + pub updated_at_ms: u64, +} + +impl From for AgentBackgroundTaskStateDto { + fn from(task: AgentBackgroundTaskState) -> Self { + Self { + task_id: task.task_id.to_string(), + kind: background_kind_label_from_work_state(task.kind).to_owned(), + state: background_state_label(task.state).to_owned(), + exit_code: task.exit_code, + summary: task.summary, + stdout_tail: task.stdout_tail, + stderr_tail: task.stderr_tail, + created_at_ms: task.created_at_ms, + updated_at_ms: task.updated_at_ms, + } + } +} + /// One manifest agent's current live/busy state. #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "camelCase")] @@ -1623,6 +1668,8 @@ pub struct AgentWorkStateDto { pub busy: AgentBusyState, /// Pending/in-progress delegation tickets, in FIFO order. pub tickets: Vec, + /// Best-effort first-class background tasks owned by this agent. + pub background_tasks: Vec, } /// How much of a [`ConversationWorkSummaryDto`] could be derived, best-effort. @@ -1748,6 +1795,11 @@ impl From for ProjectWorkStateDto { .into_iter() .map(AgentTicketStateDto::from) .collect(), + background_tasks: agent + .background_tasks + .into_iter() + .map(AgentBackgroundTaskStateDto::from) + .collect(), }) .collect(), conversations: state @@ -2924,6 +2976,16 @@ fn background_kind_label(kind: &BackgroundTaskKind) -> &'static str { } } +/// Kind discriminant string for a work-state [`BackgroundTaskKindLabel`]. +fn background_kind_label_from_work_state(kind: BackgroundTaskKindLabel) -> &'static str { + match kind { + BackgroundTaskKindLabel::Command => "command", + BackgroundTaskKindLabel::HeadlessRendezvous => "headlessRendezvous", + BackgroundTaskKindLabel::SessionResume => "sessionResume", + BackgroundTaskKindLabel::Maintenance => "maintenance", + } +} + /// Lifecycle state string for a [`BackgroundTaskState`]. fn background_state_label(state: BackgroundTaskState) -> &'static str { match state { diff --git a/crates/app-tauri/src/state.rs b/crates/app-tauri/src/state.rs index 51754ca..ace4b4b 100644 --- a/crates/app-tauri/src/state.rs +++ b/crates/app-tauri/src/state.rs @@ -503,6 +503,48 @@ impl AppReconcileBackgroundTasks { } } +fn schedule_background_ready_retry( + ready_tx: tokio::sync::mpsc::UnboundedSender, + ready: BackgroundTaskReadyToDeliver, +) { + tauri::async_runtime::spawn(async move { + tokio::time::sleep(std::time::Duration::from_secs(1)).await; + let _ = ready_tx.send(ready); + }); +} + +fn schedule_background_wake_retry( + wake: Arc, + project: Project, + owner_agent_id: AgentId, + task_id: TaskId, +) { + tauri::async_runtime::spawn(async move { + loop { + tokio::time::sleep(std::time::Duration::from_secs(1)).await; + match wake + .wake_agent( + &project, + owner_agent_id, + WakeReason::BackgroundCompletion { task_id }, + ) + .await + { + Ok(()) => break, + Err(WakeError::AgentBusy { .. }) => continue, + Err(err) => { + application::diag!( + "[background-task] completion wake retry failed: task={} owner={} err={err}", + task_id, + owner_agent_id + ); + break; + } + } + } + }); +} + struct AppWakeSessionProvider { launch_agent: Arc, structured_sessions: Arc, @@ -1690,6 +1732,7 @@ impl AppState { let wake = Arc::clone(&background_wake); let projects = Arc::clone(&store_port); let clock_for_items = Arc::clone(&clock) as Arc; + let retry_ready = background_ready_tx.clone(); tauri::async_runtime::spawn(async move { while let Some(ready) = background_ready_rx.recv().await { let item = InboxItem { @@ -1712,6 +1755,7 @@ impl AppState { ready.owner_agent_id, receipt.depth ); + schedule_background_ready_retry(retry_ready.clone(), ready); continue; } Ok(_) => {} @@ -1722,6 +1766,7 @@ impl AppState { ready.task_id, ready.owner_agent_id ); + schedule_background_ready_retry(retry_ready.clone(), ready); continue; } Err(err) => { @@ -1758,6 +1803,20 @@ impl AppState { ) .await { + if matches!(err, WakeError::AgentBusy { .. }) { + application::diag!( + "[background-task] completion wake postponed: task={} owner={}", + ready.task_id, + ready.owner_agent_id + ); + schedule_background_wake_retry( + Arc::clone(&wake), + project, + ready.owner_agent_id, + ready.task_id, + ); + continue; + } application::diag!( "[background-task] completion wake failed: task={} owner={} err={err}", ready.task_id, @@ -1801,7 +1860,8 @@ impl AppState { Arc::new(AppHandoffProvider) as Arc, Arc::new(AppConversationLogProvider) as Arc, - ), + ) + .with_background_tasks(Arc::clone(&background_tasks_port)), ); // Lot LS6 — rotation (hors chemin chaud) + lecture humaine paginée. Tous deux // composent le provider d'archive par root ; la rotation lit aussi le handoff diff --git a/crates/app-tauri/tests/dto_agents.rs b/crates/app-tauri/tests/dto_agents.rs index fa1cc59..e5c59a5 100644 --- a/crates/app-tauri/tests/dto_agents.rs +++ b/crates/app-tauri/tests/dto_agents.rs @@ -227,6 +227,7 @@ fn project_work_state_dto_serialises_live_and_busy_camelcase() { since_ms: 1_234, }, tickets: vec![], + background_tasks: vec![], }], conversations: vec![], }); @@ -285,6 +286,7 @@ fn project_work_state_dto_serialises_tickets_camelcase() { task_len: 8, }, ], + background_tasks: vec![], }], conversations: vec![], }); diff --git a/crates/application/src/agent/mod.rs b/crates/application/src/agent/mod.rs index 95b0d98..130500a 100644 --- a/crates/application/src/agent/mod.rs +++ b/crates/application/src/agent/mod.rs @@ -19,7 +19,8 @@ pub(crate) use lifecycle::ReattachDecision; pub use session_limit::{AgentResumer, SessionLimitService, RESUME_PROMPT}; pub use structured::{ - drain_with_readiness, drain_with_readiness_outcome, send_blocking, TurnOutcome, + drain_reply_stream_with_readiness, drain_with_readiness, drain_with_readiness_outcome, + send_blocking, TurnOutcome, }; pub use catalogue::{ diff --git a/crates/application/src/agent/structured.rs b/crates/application/src/agent/structured.rs index 37aacd5..ab8085e 100644 --- a/crates/application/src/agent/structured.rs +++ b/crates/application/src/agent/structured.rs @@ -16,7 +16,7 @@ use std::time::Duration; use domain::ids::AgentId; use domain::input::InputMediator; -use domain::ports::{AgentSession, AgentSessionError, ReplyEvent}; +use domain::ports::{AgentSession, AgentSessionError, ReplyEvent, ReplyStream}; use domain::readiness::{ReadinessPolicy, ReadinessSignal}; /// Issue d'un tour drainé (ARCHITECTURE §21.2-T4, réconciliation de la limite de @@ -113,6 +113,36 @@ pub async fn drain_with_readiness( } } +/// Draine un flux de réponse déjà ouvert, avec la même politique de readiness que +/// [`drain_with_readiness`]. +/// +/// Ce point d'entrée sert aux appelants qui doivent effectuer une action durable +/// après `AgentSession::send` réussi mais avant le drain complet du tour. +pub async fn drain_reply_stream_with_readiness( + stream: ReplyStream, + mediator: &dyn InputMediator, + agent: AgentId, +) -> Result { + match drain_stream_to_final( + stream, + |event| { + if !matches!(event, ReplyEvent::Final { .. }) { + mediator.mark_alive(agent); + } + }, + |signal| { + if signal == ReadinessSignal::TurnEnded { + mediator.mark_idle(agent); + } + }, + )? { + TurnOutcome::Completed(content) => Ok(content), + TurnOutcome::RateLimited { .. } => Err(AgentSessionError::Io( + "le tour s'est clos en limite de session, sans contenu Final".to_string(), + )), + } +} + /// Variante **riche** de [`drain_with_readiness`] : même branchement readiness, mais /// retourne le [`TurnOutcome`] complet au lieu de réduire la limite à une erreur. /// @@ -209,6 +239,15 @@ async fn drain_bounded_events( async fn drain_to_final( session: &dyn AgentSession, prompt: &str, + on_event: impl FnMut(&ReplyEvent), + on_signal: impl FnMut(ReadinessSignal), +) -> Result { + let stream = session.send(prompt).await?; + drain_stream_to_final(stream, on_event, on_signal) +} + +fn drain_stream_to_final( + stream: ReplyStream, mut on_event: impl FnMut(&ReplyEvent), mut on_signal: impl FnMut(ReadinessSignal), ) -> Result { @@ -216,7 +255,6 @@ async fn drain_to_final( // `RateLimited` traverse le flux. Sert UNIQUEMENT au cas « clos sans Final » — // un `Final` ultérieur l'emporte toujours (le tour a réellement abouti). let mut last_rate_limit: Option> = None; - let stream = session.send(prompt).await?; for event in stream { // Battement de vivacité (lot 2) : notifié pour CHAQUE événement brut, avant le // classement readiness. Le sink décide (les non-terminaux prouvent la vivacité). diff --git a/crates/application/src/background/mod.rs b/crates/application/src/background/mod.rs index 32f6a97..9617edf 100644 --- a/crates/application/src/background/mod.rs +++ b/crates/application/src/background/mod.rs @@ -193,10 +193,7 @@ impl CancelBackgroundTask { /// /// # Errors /// [`AppError`] if the runner or store fails. - pub async fn execute( - &self, - task_id: TaskId, - ) -> Result { + pub async fn execute(&self, task_id: TaskId) -> Result { self.runner.cancel(task_id).await.map_err(map_port_err)?; let task = self.store.get(task_id).await.map_err(map_port_err)?; Ok(CancelBackgroundTaskOutput { task }) @@ -256,10 +253,7 @@ impl RetryBackgroundTask { /// # Errors /// [`AppError::NotFound`] if the task is unknown, [`AppError::Invalid`] if it /// is not a terminal command task or its command is no longer available. - pub async fn execute( - &self, - task_id: TaskId, - ) -> Result { + pub async fn execute(&self, task_id: TaskId) -> Result { let old = self .store .get(task_id) diff --git a/crates/application/src/lib.rs b/crates/application/src/lib.rs index 5ebfd6a..f9ecc17 100644 --- a/crates/application/src/lib.rs +++ b/crates/application/src/lib.rs @@ -33,25 +33,27 @@ pub mod window; pub mod workstate; pub use agent::{ - drain_with_readiness, drain_with_readiness_outcome, reference_profile_id, reference_profiles, - selectable_reference_profiles, send_blocking, AgentResumer, ChangeAgentProfile, - ChangeAgentProfileInput, ChangeAgentProfileOutput, ConfigureProfiles, ConfigureProfilesInput, - ConfigureProfilesOutput, CreateAgentFromScratch, CreateAgentInput, CreateAgentOutput, - DeleteAgent, DeleteAgentInput, DeleteProfile, DeleteProfileInput, DetectProfiles, - DetectProfilesInput, DetectProfilesOutput, FirstRunState, FirstRunStateOutput, HandoffProvider, - InjectedLiveRow, InspectConversation, InspectConversationInput, InspectConversationOutput, - LaunchAgent, LaunchAgentInput, LaunchAgentOutput, ListAgents, ListAgentsInput, - ListAgentsOutput, ListProfiles, ListProfilesOutput, ListResumableAgents, - ListResumableAgentsInput, ListResumableAgentsOutput, LiveStateLeanProvider, McpRuntime, - PermissionProjectorRegistry, ProfileAvailability, ProviderSessionProvider, ReadAgentContext, - ReadAgentContextInput, ReadAgentContextOutput, ReferenceProfiles, ReferenceProfilesOutput, - ResumableAgent, SaveProfile, SaveProfileInput, SaveProfileOutput, SessionLimitService, - StructuredSessionDescriptor, TurnOutcome, UpdateAgentContext, UpdateAgentContextInput, - AGENT_MEMORY_RECALL_BUDGET, CODEX_SUBMIT_DELAY_MS, LIVE_STATE_INJECT_MAX, RESUME_PROMPT, + drain_reply_stream_with_readiness, drain_with_readiness, drain_with_readiness_outcome, + reference_profile_id, reference_profiles, selectable_reference_profiles, send_blocking, + AgentResumer, ChangeAgentProfile, ChangeAgentProfileInput, ChangeAgentProfileOutput, + ConfigureProfiles, ConfigureProfilesInput, ConfigureProfilesOutput, CreateAgentFromScratch, + CreateAgentInput, CreateAgentOutput, DeleteAgent, DeleteAgentInput, DeleteProfile, + DeleteProfileInput, DetectProfiles, DetectProfilesInput, DetectProfilesOutput, FirstRunState, + FirstRunStateOutput, HandoffProvider, InjectedLiveRow, InspectConversation, + InspectConversationInput, InspectConversationOutput, LaunchAgent, LaunchAgentInput, + LaunchAgentOutput, ListAgents, ListAgentsInput, ListAgentsOutput, ListProfiles, + ListProfilesOutput, ListResumableAgents, ListResumableAgentsInput, ListResumableAgentsOutput, + LiveStateLeanProvider, McpRuntime, PermissionProjectorRegistry, ProfileAvailability, + ProviderSessionProvider, ReadAgentContext, ReadAgentContextInput, ReadAgentContextOutput, + ReferenceProfiles, ReferenceProfilesOutput, ResumableAgent, SaveProfile, SaveProfileInput, + SaveProfileOutput, SessionLimitService, StructuredSessionDescriptor, TurnOutcome, + UpdateAgentContext, UpdateAgentContextInput, AGENT_MEMORY_RECALL_BUDGET, CODEX_SUBMIT_DELAY_MS, + LIVE_STATE_INJECT_MAX, RESUME_PROMPT, }; pub use background::{ - BackgroundCommandArchive, CancelBackgroundTask, CancelBackgroundTaskOutput, RetryBackgroundTask, - SpawnBackgroundCommand, SpawnBackgroundCommandInput, SpawnBackgroundCommandOutput, + BackgroundCommandArchive, CancelBackgroundTask, CancelBackgroundTaskOutput, + RetryBackgroundTask, SpawnBackgroundCommand, SpawnBackgroundCommandInput, + SpawnBackgroundCommandOutput, }; pub use conversation::{ ConversationArchiveProvider, ReadConversationPage, ReadConversationPageInput, RecordTurn, @@ -140,11 +142,11 @@ pub use terminal::{ }; pub use window::{MoveTabToNewWindow, MoveTabToNewWindowInput, MoveTabToNewWindowOutput}; pub use workstate::{ - AgentTicketState, AgentWorkState, AttachLiveAgent, AttachLiveAgentInput, AttachLiveAgentOutput, - ConversationLogProvider, ConversationPreviewStatus, ConversationTurnWorkPreview, - ConversationWorkSummary, GetLiveStateLean, GetProjectWorkState, GetProjectWorkStateInput, - LeanLiveEntry, LeanLiveState, LiveWorkSession, ProjectWorkState, ReconcileLiveState, - ReconcileLiveStateInput, StopLiveAgent, StopLiveAgentInput, StopLiveAgentOutput, - TicketWorkSource, TicketWorkStatus, UpdateLiveState, UpdateLiveStateInput, - LIVE_STATE_MAX_ENTRIES, LIVE_STATE_TTL_MS, + AgentBackgroundTaskState, AgentTicketState, AgentWorkState, AttachLiveAgent, + AttachLiveAgentInput, AttachLiveAgentOutput, BackgroundTaskKindLabel, ConversationLogProvider, + ConversationPreviewStatus, ConversationTurnWorkPreview, ConversationWorkSummary, + GetLiveStateLean, GetProjectWorkState, GetProjectWorkStateInput, LeanLiveEntry, LeanLiveState, + LiveWorkSession, ProjectWorkState, ReconcileLiveState, ReconcileLiveStateInput, StopLiveAgent, + StopLiveAgentInput, StopLiveAgentOutput, TicketWorkSource, TicketWorkStatus, UpdateLiveState, + UpdateLiveStateInput, LIVE_STATE_MAX_ENTRIES, LIVE_STATE_TTL_MS, }; diff --git a/crates/application/src/orchestrator/wake.rs b/crates/application/src/orchestrator/wake.rs index f785ca5..b48f1fb 100644 --- a/crates/application/src/orchestrator/wake.rs +++ b/crates/application/src/orchestrator/wake.rs @@ -19,7 +19,7 @@ use domain::ports::{ }; use domain::project::Project; -use crate::agent::drain_with_readiness; +use crate::agent::drain_reply_stream_with_readiness; /// Provides a structured/headless session for a wake turn. /// @@ -96,7 +96,7 @@ impl AgentWakeService { }); if self.input.busy_state(agent).is_busy() { - return Ok(()); + return Err(WakeError::AgentBusy { agent_id: agent }); } let Some(item) = self.inbox.dequeue_next(agent) else { @@ -117,26 +117,29 @@ impl AgentWakeService { agent_id: agent, }); let session = self.sessions.session_for_wake(project, agent).await?; - drain_with_readiness( - session.as_ref(), - &delivery.prompt, - None, - self.input.as_ref(), - agent, - ) - .await - .map_err(|err| WakeError::Session(err.to_string()))?; - - self.mailbox.cancel_head(agent, delivery.ticket_id); - self.input.mark_idle(agent); - guard.disarm(); + let stream = session + .send(&delivery.prompt) + .await + .map_err(|err| WakeError::Session(err.to_string()))?; if let Some(task_id) = delivery.delivered_task_id { self.tasks .mark_completion_delivered(task_id) .await .map_err(|err| WakeError::Store(err.to_string()))?; + self.publish(DomainEvent::BackgroundTaskCompletionDelivered { + project_id: project.id, + task_id, + owner_agent_id: agent, + }); } + drain_reply_stream_with_readiness(stream, self.input.as_ref(), agent) + .await + .map_err(|err| WakeError::Session(err.to_string()))?; + + self.mailbox.cancel_head(agent, delivery.ticket_id); + self.input.mark_idle(agent); + guard.disarm(); Ok(()) } diff --git a/crates/application/src/workstate/mod.rs b/crates/application/src/workstate/mod.rs index 459ef9f..2e5071d 100644 --- a/crates/application/src/workstate/mod.rs +++ b/crates/application/src/workstate/mod.rs @@ -23,11 +23,12 @@ use std::collections::{HashMap, HashSet}; use std::sync::Arc; use domain::input::{AgentBusyState, InputSource}; -use domain::ports::AgentContextStore; +use domain::ports::{AgentContextStore, BackgroundTaskPortError, BackgroundTaskStore}; use domain::{ - AgentId, AgentQueueSnapshot, ConversationId, ConversationLog, ConversationTurn, Handoff, - HandoffStore, InputMediator, NodeId, ProfileId, Project, ProjectPath, QueuedTicketSnapshot, - SessionId, TicketId, TurnId, TurnRole, + AgentId, AgentQueueSnapshot, BackgroundTask, BackgroundTaskKind, BackgroundTaskResult, + BackgroundTaskState, ConversationId, ConversationLog, ConversationTurn, Handoff, HandoffStore, + InputMediator, NodeId, ProfileId, Project, ProjectPath, QueuedTicketSnapshot, SessionId, + TaskId, TicketId, TurnId, TurnRole, }; use crate::error::AppError; @@ -150,6 +151,48 @@ pub struct AgentWorkState { pub busy: AgentBusyState, /// Pending/in-progress delegation tickets, in FIFO order. pub tickets: Vec, + /// Best-effort first-class background tasks owned by this agent. + /// + /// V1 bound: terminal tasks already delivered to the owner are no longer + /// enumerable by the background-task store, so they disappear from Work; + /// retry is available only while the terminal completion is still undelivered. + pub background_tasks: Vec, +} + +/// Background task projected into an agent's work state. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct AgentBackgroundTaskState { + /// Stable task id. + pub task_id: TaskId, + /// Kind discriminant. + pub kind: BackgroundTaskKindLabel, + /// Current lifecycle state. + pub state: BackgroundTaskState, + /// Process exit code, when the terminal result carries one. + pub exit_code: Option, + /// Human-readable summary / error / reason of the terminal result. + pub summary: Option, + /// Bounded stdout tail. + pub stdout_tail: Option, + /// Bounded stderr tail. + pub stderr_tail: Option, + /// Creation timestamp, epoch milliseconds. + pub created_at_ms: u64, + /// Last update timestamp, epoch milliseconds. + pub updated_at_ms: u64, +} + +/// Background task kind label used by the work-state read model. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum BackgroundTaskKindLabel { + /// A non-interactive command. + Command, + /// A headless inter-agent rendezvous. + HeadlessRendezvous, + /// Resume of a structured agent session. + SessionResume, + /// Internal maintenance work. + Maintenance, } /// One queued delegation ticket, projected for the work-state read model. @@ -229,6 +272,8 @@ pub struct GetProjectWorkState { handoffs: Option>, /// Per-root log source (fallback), wired best-effort; `None` ⇒ no fallback turns. logs: Option>, + /// Background-task read model source, wired best-effort; `None` ⇒ no tasks. + background_tasks: Option>, } impl GetProjectWorkState { @@ -250,6 +295,7 @@ impl GetProjectWorkState { queue, handoffs: None, logs: None, + background_tasks: None, } } @@ -267,6 +313,15 @@ impl GetProjectWorkState { self } + /// Wires the best-effort background-task source for the Work panel. Builder so + /// existing call sites/tests stay unchanged; absent or failing store ⇒ empty + /// `background_tasks` without affecting live/busy/tickets. + #[must_use] + pub fn with_background_tasks(mut self, store: Arc) -> Self { + self.background_tasks = Some(store); + self + } + /// Executes the read-only aggregation. /// /// # Errors @@ -278,40 +333,47 @@ impl GetProjectWorkState { ) -> Result { let manifest = self.contexts.load_manifest(&input.project).await?; let live_by_agent = live_by_agent(self.live.live_agent_snapshots()); - let agents = manifest - .entries - .iter() - .map(|entry| { - let agent = entry - .to_agent() - .map_err(|err| AppError::Invalid(err.to_string()))?; - let live = live_by_agent - .get(&agent.id) - .map(|snapshot| LiveWorkSession { - node_id: snapshot.node_id, - session_id: snapshot.session_id, - kind: snapshot.kind, - }); - // Tickets are crossed with the busy state: only an agent absent from - // the manifest is dropped (this loop only visits manifest entries), so - // the manifest boundary is naturally preserved. - let busy_ticket = self.input.busy_state(agent.id).ticket(); - let tickets = self - .queue - .queue_for(agent.id) - .into_iter() - .map(|snapshot| ticket_state(snapshot, busy_ticket)) - .collect(); - Ok(AgentWorkState { - agent_id: agent.id, - name: agent.name, - profile_id: agent.profile_id, - live, - busy: self.input.busy_state(agent.id), - tickets, - }) - }) - .collect::, AppError>>()?; + let undelivered_completions = match &self.background_tasks { + Some(store) => Some(store.list_undelivered_completions().await), + None => None, + }; + + let mut agents = Vec::with_capacity(manifest.entries.len()); + for entry in &manifest.entries { + let agent = entry + .to_agent() + .map_err(|err| AppError::Invalid(err.to_string()))?; + let live = live_by_agent + .get(&agent.id) + .map(|snapshot| LiveWorkSession { + node_id: snapshot.node_id, + session_id: snapshot.session_id, + kind: snapshot.kind, + }); + // Tickets are crossed with the busy state: only an agent absent from + // the manifest is dropped (this loop only visits manifest entries), so + // the manifest boundary is naturally preserved. + let busy = self.input.busy_state(agent.id); + let busy_ticket = busy.ticket(); + let tickets = self + .queue + .queue_for(agent.id) + .into_iter() + .map(|snapshot| ticket_state(snapshot, busy_ticket)) + .collect(); + let background_tasks = self + .background_tasks_for_agent(&input.project, agent.id, &undelivered_completions) + .await; + agents.push(AgentWorkState { + agent_id: agent.id, + name: agent.name, + profile_id: agent.profile_id, + live, + busy, + tickets, + background_tasks, + }); + } // Lot C — best-effort conversation summaries. Only the conversations // visible through the projected tickets are summarised, deduplicated in @@ -354,6 +416,106 @@ impl GetProjectWorkState { } summaries } + + async fn background_tasks_for_agent( + &self, + project: &Project, + agent_id: AgentId, + undelivered_completions: &Option, BackgroundTaskPortError>>, + ) -> Vec { + let Some(store) = &self.background_tasks else { + return Vec::new(); + }; + let Some(Ok(undelivered_completions)) = undelivered_completions.as_ref() else { + return Vec::new(); + }; + let Ok(open) = store.list_open_for_agent(agent_id).await else { + return Vec::new(); + }; + + let mut by_id: HashMap = HashMap::new(); + for task in open { + by_id.insert(task.id, task); + } + for task in undelivered_completions + .iter() + .filter(|task| task.owner_agent_id == agent_id) + { + by_id.entry(task.id).or_insert_with(|| task.clone()); + } + + let mut tasks: Vec<_> = by_id + .into_values() + .filter(|task| task.project_id == project.id) + .map(AgentBackgroundTaskState::from) + .collect(); + tasks.sort_by_key(|task| (task.created_at_ms, task.task_id)); + tasks + } +} + +impl From for AgentBackgroundTaskState { + fn from(task: BackgroundTask) -> Self { + let (exit_code, summary, stdout_tail, stderr_tail) = flatten_background_result(&task); + Self { + task_id: task.id, + kind: BackgroundTaskKindLabel::from(&task.kind), + state: task.state, + exit_code, + summary, + stdout_tail, + stderr_tail, + created_at_ms: task.created_at_ms, + updated_at_ms: task.updated_at_ms, + } + } +} + +impl From<&BackgroundTaskKind> for BackgroundTaskKindLabel { + fn from(kind: &BackgroundTaskKind) -> Self { + match kind { + BackgroundTaskKind::Command { .. } => Self::Command, + BackgroundTaskKind::HeadlessRendezvous { .. } => Self::HeadlessRendezvous, + BackgroundTaskKind::SessionResume { .. } => Self::SessionResume, + BackgroundTaskKind::Maintenance { .. } => Self::Maintenance, + } + } +} + +fn flatten_background_result( + task: &BackgroundTask, +) -> (Option, Option, Option, Option) { + match &task.result { + Some(BackgroundTaskResult::Success { + exit_code, + summary, + stdout_tail, + stderr_tail, + .. + }) => ( + *exit_code, + Some(summary.clone()), + stdout_tail.clone(), + stderr_tail.clone(), + ), + Some(BackgroundTaskResult::Failure { + exit_code, + error, + stdout_tail, + stderr_tail, + .. + }) => ( + *exit_code, + Some(error.clone()), + stdout_tail.clone(), + stderr_tail.clone(), + ), + Some(BackgroundTaskResult::Cancelled { reason, .. }) + | Some(BackgroundTaskResult::Expired { reason, .. }) => { + (None, Some(reason.clone()), None, None) + } + None => (None, None, None, None), + } } /// Collects the conversation ids referenced by the agents' tickets, deduplicated in diff --git a/crates/application/tests/agent_wake.rs b/crates/application/tests/agent_wake.rs index d5cd519..02a9ee2 100644 --- a/crates/application/tests/agent_wake.rs +++ b/crates/application/tests/agent_wake.rs @@ -298,6 +298,16 @@ impl BackgroundTaskStore for FakeTaskStore { #[derive(Default)] struct FakeSession { prompts: Mutex>, + events: Mutex>>, +} + +impl FakeSession { + fn with_events(events: Vec) -> Self { + Self { + prompts: Mutex::new(Vec::new()), + events: Mutex::new(Some(events)), + } + } } #[async_trait] @@ -312,12 +322,12 @@ impl AgentSession for FakeSession { async fn send(&self, prompt: &str) -> Result { self.prompts.lock().unwrap().push(prompt.to_owned()); - Ok(Box::new( + let events = self.events.lock().unwrap().take().unwrap_or_else(|| { vec![ReplyEvent::Final { content: "ack".to_owned(), }] - .into_iter(), - )) + }); + Ok(Box::new(events.into_iter())) } async fn shutdown(&self) -> Result<(), AgentSessionError> { @@ -421,15 +431,16 @@ async fn owner_busy_does_not_start_concurrent_wake_and_keeps_item_queued() { .unwrap(); turns.force_busy(owner, ticket(99)); - service(inbox.clone(), turns, tasks, sessions) + let err = service(inbox.clone(), turns, tasks, sessions) .wake_agent( &project, owner, WakeReason::BackgroundCompletion { task_id }, ) .await - .unwrap(); + .unwrap_err(); + assert_eq!(err, WakeError::AgentBusy { agent_id: owner }); assert!(session.prompts.lock().unwrap().is_empty()); assert_eq!(inbox.snapshot(owner).depth, 1); } @@ -489,6 +500,35 @@ async fn completion_is_marked_delivered_after_successful_wake() { assert_eq!(tasks.delivered(), vec![task_id]); } +#[tokio::test] +async fn completion_is_marked_delivered_once_send_is_accepted_even_if_drain_fails() { + let project = project(); + let owner = agent(1); + let task_id = task_id(10); + let inbox = Arc::new(FakeInbox::default()); + let turns = Arc::new(SharedTurnState::default()); + let tasks = Arc::new(FakeTaskStore::default()); + let session = Arc::new(FakeSession::with_events(Vec::new())); + let sessions = Arc::new(FakeSessionProvider::with_session(session.clone())); + tasks.insert(completed_task(&project, owner, task_id, "done")); + inbox + .enqueue_message(owner, completion_item(owner, task_id, ticket(20))) + .unwrap(); + + let err = service(inbox, turns, tasks.clone(), sessions) + .wake_agent( + &project, + owner, + WakeReason::BackgroundCompletion { task_id }, + ) + .await + .unwrap_err(); + + assert!(matches!(err, WakeError::Session(_))); + assert_eq!(session.prompts.lock().unwrap().len(), 1); + assert_eq!(tasks.delivered(), vec![task_id]); +} + #[tokio::test] async fn wake_drains_exactly_one_item_per_turn() { let project = project(); diff --git a/crates/application/tests/workstate.rs b/crates/application/tests/workstate.rs index 2ddf00b..a153123 100644 --- a/crates/application/tests/workstate.rs +++ b/crates/application/tests/workstate.rs @@ -16,13 +16,15 @@ use domain::mailbox::{ AgentQueueSnapshot, MailboxError, PendingReply, QueuedTicketSnapshot, Ticket, TurnResolution, }; use domain::ports::{ - AgentContextStore, AgentSession, AgentSessionError, PtyHandle, ReplyStream, StoreError, + AgentContextStore, AgentSession, AgentSessionError, BackgroundTaskPortError, + BackgroundTaskStore, PtyHandle, ReplyStream, StoreError, }; use domain::{ - Agent, AgentBusyState, AgentId, AgentManifest, AgentOrigin, ConversationId, ConversationLog, - ConversationTurn, Handoff, HandoffStore, InputMediator, InputSource, ManifestEntry, - MarkdownDoc, NodeId, ProfileId, Project, ProjectId, ProjectPath, PtySize, RemoteRef, SessionId, - SessionKind, TerminalSession, TicketId, TurnId, TurnRole, + 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; @@ -46,6 +48,10 @@ 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)), @@ -147,6 +153,106 @@ struct FakeQueue { queues: Mutex>>, } +#[derive(Default)] +struct FakeBackgroundTaskStore { + tasks: Mutex>, + fail_open: Mutex, + fail_undelivered: Mutex, +} + +impl FakeBackgroundTaskStore { + fn set_tasks(&self, tasks: Vec) { + *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, 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, 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, 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) { self.queues.lock().unwrap().insert(agent, tickets); @@ -400,6 +506,106 @@ struct ConvFixture { project: Project, } +struct BackgroundFixture { + usecase: GetProjectWorkState, + queue: Arc, + store: Arc, + 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, + Arc::clone(&queue) as Arc, + ) + .with_background_tasks(Arc::clone(&store) as Arc); + 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 { @@ -539,6 +745,123 @@ async fn workstate_agent_without_queue_has_no_tickets() { 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"); diff --git a/crates/domain/src/lib.rs b/crates/domain/src/lib.rs index 2669905..5823b0a 100644 --- a/crates/domain/src/lib.rs +++ b/crates/domain/src/lib.rs @@ -181,8 +181,7 @@ pub use ports::{ EmbedderPromptStore, EventBus, EventStream, ExitStatus, FileSystem, FsError, GitCommitInfo, GitError, GitFileStatus, GitPort, GraphCommit, IdGenerator, IssueNumberAllocator, IssueStore, IssueStoreError, LiveStateStore, MemoryError, MemoryQuery, MemoryRecall, MemoryStore, Output, - OutputStream, PermissionStore, PreparedContext, - ProcessError, ProcessSpawner, ProfileStore, ProjectStore, PtyError, PtyHandle, PtyPort, - RemoteError, RemoteHost, RemotePath, RuntimeError, ScheduledTask, Scheduler, SpawnSpec, - StoreError, TemplateStore, + OutputStream, PermissionStore, PreparedContext, ProcessError, ProcessSpawner, ProfileStore, + ProjectStore, PtyError, PtyHandle, PtyPort, RemoteError, RemoteHost, RemotePath, RuntimeError, + ScheduledTask, Scheduler, SpawnSpec, StoreError, TemplateStore, }; diff --git a/crates/infrastructure/src/background_task/mod.rs b/crates/infrastructure/src/background_task/mod.rs index a95c3c0..8624656 100644 --- a/crates/infrastructure/src/background_task/mod.rs +++ b/crates/infrastructure/src/background_task/mod.rs @@ -17,6 +17,7 @@ mod tail; pub use runner::CommandBackgroundRunner; pub use sink::{ start_background_ready_inbox_bridge, BackgroundCompletionSink, BackgroundCompletionSinkError, - BackgroundCompletionSinkOutcome, BackgroundReadyInboxBridgeHandle, BackgroundTaskReadyToDeliver, + BackgroundCompletionSinkOutcome, BackgroundReadyInboxBridgeHandle, + BackgroundTaskReadyToDeliver, }; pub use tail::{bounded_tail, tail_cap_bytes, BoundedTail}; diff --git a/crates/infrastructure/src/background_task/runner.rs b/crates/infrastructure/src/background_task/runner.rs index 0c7c4b9..d5b22a3 100644 --- a/crates/infrastructure/src/background_task/runner.rs +++ b/crates/infrastructure/src/background_task/runner.rs @@ -170,7 +170,10 @@ impl CommandBackgroundRunner { } }; - running.lock().expect("runner registry poisoned").remove(&task_id); + running + .lock() + .expect("runner registry poisoned") + .remove(&task_id); let _ = completion_tx.send(BackgroundTaskCompletion { task_id, result }); } @@ -213,10 +216,7 @@ impl CommandBackgroundRunner { let Ok(stream) = pty.subscribe_output(&handle) else { return; }; - let _ = tokio::task::spawn_blocking(move || { - for _chunk in stream {} - }) - .await; + let _ = tokio::task::spawn_blocking(move || for _chunk in stream {}).await; } } diff --git a/crates/infrastructure/src/input/mod.rs b/crates/infrastructure/src/input/mod.rs index ab1b44b..eac3343 100644 --- a/crates/infrastructure/src/input/mod.rs +++ b/crates/infrastructure/src/input/mod.rs @@ -1077,10 +1077,7 @@ impl AgentInbox for MediatedInbox { let item_id = item.id; let ticket = inbox_item_to_ticket(&item); self.inbox_items().insert(item_id, item); - let _pending = match ticket.source { - domain::input::InputSource::Human => self.enqueue(agent, ticket), - domain::input::InputSource::Agent { .. } => self.enqueue_silent(agent, ticket), - }; + let _pending = self.mailbox.enqueue(agent, ticket); let depth = self.mailbox.pending(&agent); if let Some(events) = &self.tracker.events { events.publish(DomainEvent::AgentInboxQueued { @@ -1347,6 +1344,36 @@ mod tests { assert_eq!(inbox.mailbox.pending(&a), 1); } + #[test] + fn enqueue_message_queues_inbox_item_without_starting_turn() { + let bus = Arc::new(RecordingBus::default()); + let inbox = MediatedInbox::new(Arc::new(InMemoryMailbox::new()), Arc::new(FixedClock(1))) + .with_events(Arc::clone(&bus) as Arc); + let a = agent(1); + let task_id = domain::TaskId::from_uuid(uuid::Uuid::from_u128(42)); + let item = domain::InboxItem { + id: TicketId::from_uuid(uuid::Uuid::from_u128(10)), + agent_id: a, + source: InboxSource::BackgroundTask { task_id }, + kind: domain::InboxItemKind::BackgroundCompletion, + body: "Background task completed.".to_owned(), + created_at_ms: 1, + correlation_id: Some(task_id.to_string()), + }; + + let receipt = inbox.enqueue_message(a, item.clone()).unwrap(); + + assert_eq!(receipt.status, InboxReceiptStatus::Queued); + assert_eq!(inbox.busy_state(a), AgentBusyState::Idle); + assert_eq!(inbox.mailbox.pending(&a), 1); + assert!( + bus.busy_events().is_empty(), + "inbox enqueue must not start a mediated turn" + ); + assert_eq!(inbox.dequeue_next(a), Some(item)); + assert_eq!(inbox.mailbox.pending(&a), 0); + } + #[test] fn delegation_ready_carries_bound_submit_config() { let bus = Arc::new(RecordingBus::default()); diff --git a/crates/infrastructure/src/lib.rs b/crates/infrastructure/src/lib.rs index 1abf876..1a08f0e 100644 --- a/crates/infrastructure/src/lib.rs +++ b/crates/infrastructure/src/lib.rs @@ -40,8 +40,9 @@ pub mod timeparse; pub use background_task::{ bounded_tail, start_background_ready_inbox_bridge, tail_cap_bytes, BackgroundCompletionSink, - BackgroundCompletionSinkError, BackgroundCompletionSinkOutcome, BackgroundReadyInboxBridgeHandle, - BackgroundTaskReadyToDeliver, BoundedTail, CommandBackgroundRunner, + BackgroundCompletionSinkError, BackgroundCompletionSinkOutcome, + BackgroundReadyInboxBridgeHandle, BackgroundTaskReadyToDeliver, BoundedTail, + CommandBackgroundRunner, }; pub use clock::SystemClock; pub use conversation::InMemoryConversationRegistry; diff --git a/frontend/src/adapters/workStateNormalization.ts b/frontend/src/adapters/workStateNormalization.ts index 49e0803..40f832f 100644 --- a/frontend/src/adapters/workStateNormalization.ts +++ b/frontend/src/adapters/workStateNormalization.ts @@ -89,6 +89,8 @@ function normalizeBackgroundStatus( if (raw === "cancelled" || raw === "canceled" || raw === "expired") return "cancelled"; if (raw === "running" || raw === "started") return "running"; if (raw === "delivered") return "delivered"; + // Backend `queued`/`waiting` collapse to the cancellable "pending" bucket. + if (raw === "queued" || raw === "waiting") return "pending"; return "pending"; } @@ -107,7 +109,7 @@ function normalizeBackgroundTask( exitCode: nullableNumber(task.exitCode), stdoutTail: nullableString(task.stdoutTail), stderrTail: nullableString(task.stderrTail), - finishedAtMs: nullableNumber(task.finishedAtMs), + updatedAtMs: numberValue(task.updatedAtMs), }; } diff --git a/frontend/src/domain/index.ts b/frontend/src/domain/index.ts index a103a2a..dff739e 100644 --- a/frontend/src/domain/index.ts +++ b/frontend/src/domain/index.ts @@ -252,7 +252,8 @@ export interface BackgroundCompletion { exitCode: number | null; stdoutTail: string | null; stderrTail: string | null; - finishedAtMs: number | null; + /** Last-update timestamp (epoch ms); chronological ordering key. */ + updatedAtMs: number; } /** One item currently pending in an agent's inbox. */ diff --git a/frontend/src/features/workstate/ProjectWorkStatePanel.tsx b/frontend/src/features/workstate/ProjectWorkStatePanel.tsx index 85788f6..85a539e 100644 --- a/frontend/src/features/workstate/ProjectWorkStatePanel.tsx +++ b/frontend/src/features/workstate/ProjectWorkStatePanel.tsx @@ -451,11 +451,9 @@ function AgentRow({ const busy = busyState.state === "busy"; const tickets = [...(agent.tickets ?? [])].sort((a, b) => a.position - b.position); const inbox = [...(agent.inbox ?? [])].sort((a, b) => a.createdAtMs - b.createdAtMs); - const backgroundTasks = [...(agent.backgroundTasks ?? [])].sort((a, b) => { - const aDone = a.finishedAtMs ?? 0; - const bDone = b.finishedAtMs ?? 0; - return bDone - aDone || a.taskId.localeCompare(b.taskId); - }); + const backgroundTasks = [...(agent.backgroundTasks ?? [])].sort( + (a, b) => b.updatedAtMs - a.updatedAtMs || a.taskId.localeCompare(b.taskId), + ); const inboxDepth = Math.max(agent.inboxDepth ?? 0, inbox.length, tickets.length); const visibleNodeIds = new Set(visibleLeaves.map((leaf) => leaf.id)); const liveVisible = live && visibleNodeIds.has(agent.live!.nodeId); diff --git a/frontend/src/features/workstate/workstate.test.tsx b/frontend/src/features/workstate/workstate.test.tsx index 44bb65a..9d18a3c 100644 --- a/frontend/src/features/workstate/workstate.test.tsx +++ b/frontend/src/features/workstate/workstate.test.tsx @@ -176,6 +176,89 @@ describe("ProjectWorkStatePanel", () => { ); }); + it("normalizes the real backend background-task payload shape", async () => { + // Mirrors AgentBackgroundTaskStateDto (crates/app-tauri/src/dto.rs): camelCase, + // `state` (not `status`), optional exitCode/summary/stdoutTail/stderrTail omitted + // when absent, createdAtMs/updatedAtMs present, NO ownerAgentId/projectId/finishedAtMs. + const workState = new MockWorkStateGateway(); + workState._setProjectWorkState(PROJECT_ID, { + agents: [ + { + agentId: "agent-bg", + name: "Runner", + profileId: "codex", + busy: { state: "idle" }, + backgroundTasks: [ + { + taskId: "task-waiting-1", + kind: "command", + state: "waiting", + createdAtMs: 10, + updatedAtMs: 11, + }, + { + taskId: "task-queued-1", + kind: "headlessRendezvous", + state: "queued", + createdAtMs: 12, + updatedAtMs: 13, + }, + { + taskId: "task-failed-1", + kind: "command", + state: "failed", + exitCode: 1, + stderrTail: "boom", + createdAtMs: 14, + updatedAtMs: 15, + }, + { + taskId: "task-expired-1", + kind: "sessionResume", + state: "expired", + createdAtMs: 16, + updatedAtMs: 17, + }, + ], + }, + ], + }); + + // Adapter-level assertion: locks the mapping against the real payload. + const state = await workState.getProjectWorkState(PROJECT_ID); + const tasks = state.agents[0]?.backgroundTasks ?? []; + expect(tasks.map((t) => [t.taskId, t.status])).toEqual([ + ["task-waiting-1", "pending"], + ["task-queued-1", "pending"], + ["task-failed-1", "failed"], + ["task-expired-1", "cancelled"], + ]); + // Optional fields omitted by the backend normalize to null (not undefined/NaN). + expect(tasks[0]?.exitCode).toBeNull(); + expect(tasks[0]?.stdoutTail).toBeNull(); + expect(tasks[2]?.stderrTail).toBe("boom"); + // updatedAtMs (backend-emitted) is carried through for chronological ordering. + expect(tasks.map((t) => t.updatedAtMs)).toEqual([11, 13, 15, 17]); + + renderPanel(workState); + + const list = await screen.findByLabelText("Runner background tasks"); + const cancelButtons = within(list).getAllByRole("button", { name: "Cancel" }); + const retryButtons = within(list).getAllByRole("button", { name: "Retry" }); + // waiting + queued are cancellable; failed + expired are retryable. + expect(cancelButtons.filter((b) => !(b as HTMLButtonElement).disabled)).toHaveLength(2); + expect(retryButtons.filter((b) => !(b as HTMLButtonElement).disabled)).toHaveLength(2); + + // Rows are ordered most-recently-updated first (updatedAtMs desc), taskId as + // the deterministic tiebreak — locks the fix for the dead finishedAtMs sort. + const text = list.textContent ?? ""; + const order = ["task-exp", "task-fai", "task-que", "task-wai"].map((id) => + text.indexOf(id), + ); + expect(order).toEqual([...order].sort((a, b) => a - b)); + expect(order.every((i) => i >= 0)).toBe(true); + }); + it("renders a legacy agent without tickets", async () => { const workState = new MockWorkStateGateway(); workState._setProjectWorkState(PROJECT_ID, {