fix(session-limit): émettre AgentRateLimited sur le chemin délégué (#7)
Sur le rendez-vous délégué (idea_ask_agent), quand la cible touchait sa limite de session, l'événement AgentRateLimited n'était pas émis : l'écart backend documenté (ticket #7/F2) laissait la surface UI sans signal de limite pour la cible déléguée. Le service orchestrateur relaie désormais la limite de la cible vers le service de limite de session, fermant l'écart en cohérence avec le chemin direct (#30). Couverture QA (sortie réelle) : orchestrator_service 63 passed, session_limit_service 15 passed, session_limit_t4 7 passed, cargo test -p application 0 failed. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@ -2232,6 +2232,10 @@ impl AppState {
|
||||
Arc::clone(&background_tasks_port),
|
||||
Arc::clone(&clock) as Arc<dyn Clock>,
|
||||
)
|
||||
// Limites de session sur le chemin délégué headless : quand le drain
|
||||
// structuré retourne `RateLimited`, l'orchestrateur arme la reprise pour
|
||||
// la cible limitée via le même service que les chemins directs.
|
||||
.with_session_limits(Arc::clone(&session_limit_service))
|
||||
// Producteur B8 côté agent MCP : `idea_run_in_background` passe par le
|
||||
// même use case que la commande Tauri, sans aller-retour frontend.
|
||||
.with_spawn_background_command(Arc::clone(&spawn_background_command)),
|
||||
|
||||
@ -20,8 +20,8 @@ pub(crate) use lifecycle::ReattachDecision;
|
||||
pub use session_limit::{AgentResumer, SessionLimitService, RESUME_PROMPT};
|
||||
pub use structured::{
|
||||
drain_reply_stream_with_readiness, drain_with_readiness,
|
||||
drain_with_readiness_and_announcements, drain_with_readiness_outcome, send_blocking,
|
||||
AnnouncementPublisher, TurnOutcome,
|
||||
drain_with_readiness_and_announcements, drain_with_readiness_and_announcements_outcome,
|
||||
drain_with_readiness_outcome, send_blocking, AnnouncementPublisher, TurnOutcome,
|
||||
};
|
||||
|
||||
pub use catalogue::{
|
||||
|
||||
@ -167,8 +167,35 @@ pub async fn drain_with_readiness_and_announcements(
|
||||
agent: AgentId,
|
||||
announcements: Option<AnnouncementPublisher>,
|
||||
) -> Result<String, AgentSessionError> {
|
||||
match drain_with_readiness_and_announcements_outcome(
|
||||
session,
|
||||
prompt,
|
||||
timeout,
|
||||
mediator,
|
||||
agent,
|
||||
announcements,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
TurnOutcome::Completed(content) => Ok(content),
|
||||
TurnOutcome::RateLimited { .. } => Err(AgentSessionError::Io(
|
||||
"le tour s'est clos en limite de session, sans contenu Final".to_string(),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
/// Comme [`drain_with_readiness_and_announcements`], mais conserve l'issue riche
|
||||
/// [`TurnOutcome::RateLimited`] pour les appelants capables d'armer la reprise.
|
||||
pub async fn drain_with_readiness_and_announcements_outcome(
|
||||
session: &dyn AgentSession,
|
||||
prompt: &str,
|
||||
timeout: Option<Duration>,
|
||||
mediator: &dyn InputMediator,
|
||||
agent: AgentId,
|
||||
announcements: Option<AnnouncementPublisher>,
|
||||
) -> Result<TurnOutcome, AgentSessionError> {
|
||||
let Some(publisher) = announcements else {
|
||||
return drain_with_readiness(session, prompt, timeout, mediator, agent).await;
|
||||
return drain_with_readiness_outcome(session, prompt, timeout, mediator, agent).await;
|
||||
};
|
||||
|
||||
let (tap_tx, tap_rx) = std::sync::mpsc::channel();
|
||||
@ -188,7 +215,7 @@ pub async fn drain_with_readiness_and_announcements(
|
||||
let _ = live_pump.join();
|
||||
let stream = stream_result?;
|
||||
|
||||
match drain_stream_to_final(
|
||||
drain_stream_to_final(
|
||||
stream,
|
||||
|event| {
|
||||
if !matches!(event, ReplyEvent::Final { .. }) {
|
||||
@ -201,12 +228,7 @@ pub async fn drain_with_readiness_and_announcements(
|
||||
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(),
|
||||
)),
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
/// Draine un flux de réponse déjà ouvert, avec la même politique de readiness que
|
||||
|
||||
@ -36,9 +36,10 @@ use crate::conversation::RecordTurn;
|
||||
use crate::memory::HarvestMemoryFromTurn;
|
||||
|
||||
use crate::agent::{
|
||||
drain_with_readiness_and_announcements, AnnouncementPublisher, CreateAgentFromScratch,
|
||||
drain_with_readiness_and_announcements_outcome, AnnouncementPublisher, CreateAgentFromScratch,
|
||||
CreateAgentInput, LaunchAgent, LaunchAgentInput, ListAgents, ListAgentsInput, McpRuntime,
|
||||
ReattachDecision, UpdateAgentContext, UpdateAgentContextInput, CODEX_SUBMIT_DELAY_MS,
|
||||
ReattachDecision, SessionLimitService, TurnOutcome, UpdateAgentContext,
|
||||
UpdateAgentContextInput, CODEX_SUBMIT_DELAY_MS,
|
||||
};
|
||||
use crate::background::{SpawnBackgroundCommand, SpawnBackgroundCommandInput};
|
||||
use crate::error::AppError;
|
||||
@ -484,6 +485,10 @@ pub struct OrchestratorService {
|
||||
/// `idea_ask_agent` comme [`BackgroundTaskKind::HeadlessRendezvous`]. `None` ⇒
|
||||
/// comportement historique sans persistance de tâche (call sites non câblés).
|
||||
background_tasks: Option<Arc<dyn BackgroundTaskStore>>,
|
||||
/// Service applicatif des limites de session, injecté depuis la composition root.
|
||||
/// `None` conserve les call sites/tests legacy ; les chemins conscients des limites
|
||||
/// publient alors seulement leur erreur de tour historique.
|
||||
session_limits: Option<Arc<SessionLimitService>>,
|
||||
/// Use case B8 qui produit une tâche de fond command-backed depuis la surface
|
||||
/// agent MCP `idea_run_in_background`. `None` ⇒ outil non configuré.
|
||||
spawn_background_command: Option<Arc<SpawnBackgroundCommand>>,
|
||||
@ -558,10 +563,18 @@ impl OrchestratorService {
|
||||
ask_liveness_probe: None,
|
||||
ask_ceiling: ASK_AGENT_CEILING,
|
||||
background_tasks: None,
|
||||
session_limits: None,
|
||||
spawn_background_command: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Branche le service de reprise automatique des limites de session.
|
||||
#[must_use]
|
||||
pub fn with_session_limits(mut self, service: Arc<SessionLimitService>) -> Self {
|
||||
self.session_limits = Some(service);
|
||||
self
|
||||
}
|
||||
|
||||
/// Branche le producteur B8 command-backed pour `idea_run_in_background`.
|
||||
#[must_use]
|
||||
pub fn with_spawn_background_command(mut self, spawn: Arc<SpawnBackgroundCommand>) -> Self {
|
||||
@ -1822,7 +1835,7 @@ impl OrchestratorService {
|
||||
target: agent_id,
|
||||
ticket: ticket_id,
|
||||
});
|
||||
let drain = drain_with_readiness_and_announcements(
|
||||
let drain = drain_with_readiness_and_announcements_outcome(
|
||||
session,
|
||||
&task,
|
||||
None,
|
||||
@ -1836,13 +1849,40 @@ impl OrchestratorService {
|
||||
// au no-reply du chemin headless : la cible a rendu la main sans réponse
|
||||
// exploitable, on expose donc l'erreur métier retryable plutôt qu'un PROCESS.
|
||||
let wait = async {
|
||||
drain.await.map_err(|err| {
|
||||
if structured_no_reply_error(&err) {
|
||||
AppError::TargetReturnedNoReply(target.to_owned())
|
||||
} else {
|
||||
AppError::from(err)
|
||||
match drain.await {
|
||||
Ok(TurnOutcome::Completed(content)) => Ok(content),
|
||||
Ok(TurnOutcome::RateLimited { resets_at_ms }) => {
|
||||
let conversation_id = session.conversation_id();
|
||||
let (conversation_id, resets_at_ms) = match (conversation_id, resets_at_ms) {
|
||||
(Some(conversation_id), resets_at_ms) => {
|
||||
(Some(conversation_id), resets_at_ms)
|
||||
}
|
||||
(None, None) => (None, None),
|
||||
// Reset connu mais conversation moteur absente : ne pas armer
|
||||
// une reprise non reprenable ; surfacer le fallback humain.
|
||||
(None, Some(_)) => (None, None),
|
||||
};
|
||||
if let (Some(service), Some(structured)) =
|
||||
(&self.session_limits, &self.structured)
|
||||
{
|
||||
if let Some(node_id) = structured.node_for_agent(&agent_id) {
|
||||
service.on_rate_limited(
|
||||
agent_id,
|
||||
node_id,
|
||||
conversation_id,
|
||||
resets_at_ms,
|
||||
);
|
||||
}
|
||||
}
|
||||
Err(AppError::Process(
|
||||
"le tour s'est clos en limite de session, sans contenu Final".to_owned(),
|
||||
))
|
||||
}
|
||||
Err(err) if structured_no_reply_error(&err) => {
|
||||
Err(AppError::TargetReturnedNoReply(target.to_owned()))
|
||||
}
|
||||
Err(err) => Err(AppError::from(err)),
|
||||
}
|
||||
})
|
||||
};
|
||||
|
||||
// Borne par la fenêtre d'inactivité (réarmée sur signe de vie) sous plafond absolu.
|
||||
|
||||
@ -45,7 +45,8 @@ use uuid::Uuid;
|
||||
use application::{
|
||||
CloseTerminal, CreateAgentFromScratch, CreateSkill, GetLiveStateLean, HarvestMemoryFromTurn,
|
||||
LaunchAgent, ListAgents, LiveStateProvider, LiveStateReadProvider, OrchestratorService,
|
||||
SpawnBackgroundCommand, TerminalSessions, UpdateAgentContext, UpdateLiveState,
|
||||
SessionLimitService, SpawnBackgroundCommand, TerminalSessions, UpdateAgentContext,
|
||||
UpdateLiveState,
|
||||
};
|
||||
use domain::live_state::{LiveEntry, LiveState, WorkStatus};
|
||||
use domain::ports::LiveStateStore;
|
||||
@ -4024,9 +4025,11 @@ async fn p6b_failing_record_does_not_degrade_ask() {
|
||||
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
use application::StructuredSessions;
|
||||
use application::{AgentResumer, StructuredSessions};
|
||||
use domain::ids::ScheduleId;
|
||||
use domain::ports::{
|
||||
AgentSession, AgentSessionError, AgentSessionFactory, ReplyEvent, ReplyStream,
|
||||
AgentSession, AgentSessionError, AgentSessionFactory, ReplyEvent, ReplyStream, ScheduledTask,
|
||||
Scheduler,
|
||||
};
|
||||
|
||||
/// Session structurée fake : `send()` rend un unique [`ReplyEvent::Final`] qui **écho**
|
||||
@ -4034,6 +4037,7 @@ use domain::ports::{
|
||||
/// du tour **et** en `mark_idle`. Aucun `idea_reply` n'est nécessaire.
|
||||
struct EchoSession {
|
||||
id: SessionId,
|
||||
conversation_id: Option<String>,
|
||||
}
|
||||
#[async_trait]
|
||||
impl AgentSession for EchoSession {
|
||||
@ -4041,7 +4045,7 @@ impl AgentSession for EchoSession {
|
||||
self.id
|
||||
}
|
||||
fn conversation_id(&self) -> Option<String> {
|
||||
None
|
||||
self.conversation_id.clone()
|
||||
}
|
||||
async fn send(&self, prompt: &str) -> Result<ReplyStream, AgentSessionError> {
|
||||
let stream: ReplyStream = Box::new(
|
||||
@ -4096,7 +4100,72 @@ impl AgentSessionFactory for CountingFactory {
|
||||
*n += 1;
|
||||
id
|
||||
};
|
||||
Ok(Arc::new(EchoSession { id }))
|
||||
Ok(Arc::new(EchoSession {
|
||||
id,
|
||||
conversation_id: Some(format!("engine-{id}")),
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
struct RateLimitedSession {
|
||||
id: SessionId,
|
||||
conversation_id: Option<String>,
|
||||
resets_at_ms: i64,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl AgentSession for RateLimitedSession {
|
||||
fn id(&self) -> SessionId {
|
||||
self.id
|
||||
}
|
||||
|
||||
fn conversation_id(&self) -> Option<String> {
|
||||
self.conversation_id.clone()
|
||||
}
|
||||
|
||||
async fn send(&self, _prompt: &str) -> Result<ReplyStream, AgentSessionError> {
|
||||
Ok(Box::new(
|
||||
vec![ReplyEvent::RateLimited {
|
||||
resets_at_ms: Some(self.resets_at_ms),
|
||||
}]
|
||||
.into_iter(),
|
||||
))
|
||||
}
|
||||
|
||||
async fn shutdown(&self) -> Result<(), AgentSessionError> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct NoopScheduler {
|
||||
next: Mutex<u128>,
|
||||
}
|
||||
|
||||
impl Scheduler for NoopScheduler {
|
||||
fn arm(&self, _fire_at_ms: i64, _task: ScheduledTask) -> ScheduleId {
|
||||
let mut next = self.next.lock().unwrap();
|
||||
*next += 1;
|
||||
ScheduleId::from_uuid(Uuid::from_u128(*next))
|
||||
}
|
||||
|
||||
fn cancel(&self, _id: ScheduleId) -> bool {
|
||||
true
|
||||
}
|
||||
}
|
||||
|
||||
struct NoopResumer;
|
||||
|
||||
#[async_trait]
|
||||
impl AgentResumer for NoopResumer {
|
||||
async fn resume(
|
||||
&self,
|
||||
_agent_id: AgentId,
|
||||
_node_id: NodeId,
|
||||
_conversation_id: Option<String>,
|
||||
_resume_prompt: &str,
|
||||
) -> Result<(), application::AppError> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
@ -4108,6 +4177,8 @@ struct StructuredAskFixture {
|
||||
sessions: Arc<TerminalSessions>,
|
||||
factory: CountingFactory,
|
||||
pty: FakePty,
|
||||
bus: SpyBus,
|
||||
session_limits: Arc<SessionLimitService>,
|
||||
}
|
||||
|
||||
fn structured_ask_fixture(contexts: FakeContexts) -> StructuredAskFixture {
|
||||
@ -4161,6 +4232,12 @@ fn structured_ask_fixture(contexts: FakeContexts) -> StructuredAskFixture {
|
||||
Arc::new(pty.clone()) as Arc<dyn PtyPort>,
|
||||
));
|
||||
let conversations = Arc::new(TestConversations::new()) as Arc<dyn ConversationRegistry>;
|
||||
let session_limits = Arc::new(SessionLimitService::new(
|
||||
Arc::new(FixedMillisClock(1_000)) as Arc<dyn Clock>,
|
||||
Arc::new(NoopScheduler::default()) as Arc<dyn Scheduler>,
|
||||
Arc::new(bus.clone()),
|
||||
Arc::new(NoopResumer),
|
||||
));
|
||||
|
||||
let service = Arc::new(
|
||||
OrchestratorService::new(
|
||||
@ -4179,6 +4256,7 @@ fn structured_ask_fixture(contexts: FakeContexts) -> StructuredAskFixture {
|
||||
)
|
||||
.with_conversations(conversations)
|
||||
.with_events(Arc::new(bus.clone()))
|
||||
.with_session_limits(Arc::clone(&session_limits))
|
||||
// Même registre que le launcher : `ask_agent` y relit la session fraîchement
|
||||
// insérée par `launch_structured`.
|
||||
.with_structured(Arc::clone(&structured)),
|
||||
@ -4190,6 +4268,8 @@ fn structured_ask_fixture(contexts: FakeContexts) -> StructuredAskFixture {
|
||||
sessions,
|
||||
factory,
|
||||
pty,
|
||||
bus,
|
||||
session_limits,
|
||||
}
|
||||
}
|
||||
|
||||
@ -4295,3 +4375,110 @@ async fn ask_structured_target_live_as_pty_keeps_visible_terminal() {
|
||||
"target still has its visible PTY session"
|
||||
);
|
||||
}
|
||||
|
||||
/// Sous-lot adjacent #7/F2 : une limite détectée pendant un rendez-vous structuré
|
||||
/// doit aussi passer par `SessionLimitService::on_rate_limited` pour la cible, afin
|
||||
/// de publier le flux mono-agent `AgentRateLimited` puis `AgentResumeScheduled`.
|
||||
#[tokio::test]
|
||||
async fn ask_structured_rate_limit_arms_target_resume_events() {
|
||||
let agent = scratch_agent(aid(1), "architect", "agents/architect.md");
|
||||
let fx = structured_ask_fixture(FakeContexts::with_agent(&agent, "# persona"));
|
||||
let resets_at_ms = 9_999_999;
|
||||
fx.structured.insert(
|
||||
Arc::new(RateLimitedSession {
|
||||
id: sid(9700),
|
||||
conversation_id: Some("engine-target".to_owned()),
|
||||
resets_at_ms,
|
||||
}),
|
||||
aid(1),
|
||||
nid(700),
|
||||
);
|
||||
|
||||
let err = fx
|
||||
.service
|
||||
.dispatch(&project(), cmd(ASK_JSON))
|
||||
.await
|
||||
.expect_err("a rate-limited delegated turn has no Final reply");
|
||||
assert_eq!(err.code(), "PROCESS");
|
||||
|
||||
let events = fx.bus.events();
|
||||
assert!(
|
||||
events.iter().any(|event| matches!(
|
||||
event,
|
||||
DomainEvent::AgentRateLimited {
|
||||
agent_id,
|
||||
resets_at_ms: Some(t),
|
||||
} if *agent_id == aid(1) && *t == resets_at_ms
|
||||
)),
|
||||
"target AgentRateLimited must be published, got {events:?}"
|
||||
);
|
||||
assert!(
|
||||
events.iter().any(|event| matches!(
|
||||
event,
|
||||
DomainEvent::AgentResumeScheduled {
|
||||
agent_id,
|
||||
fire_at_ms,
|
||||
} if *agent_id == aid(1) && *fire_at_ms == resets_at_ms
|
||||
)),
|
||||
"target AgentResumeScheduled must be published, got {events:?}"
|
||||
);
|
||||
assert!(
|
||||
fx.session_limits.cancel_resume(aid(1)),
|
||||
"delegated target resume must be cancellable"
|
||||
);
|
||||
}
|
||||
|
||||
/// Sous-lot adjacent #7/F2 : si la cible structurée est limitée avec un reset connu
|
||||
/// mais sans `conversation_id` moteur, l'orchestrateur doit exposer le fallback humain
|
||||
/// et ne jamais armer une reprise non reprenable.
|
||||
#[tokio::test]
|
||||
async fn ask_structured_rate_limit_without_conversation_uses_human_fallback_no_resume() {
|
||||
let agent = scratch_agent(aid(1), "architect", "agents/architect.md");
|
||||
let fx = structured_ask_fixture(FakeContexts::with_agent(&agent, "# persona"));
|
||||
fx.structured.insert(
|
||||
Arc::new(RateLimitedSession {
|
||||
id: sid(9701),
|
||||
conversation_id: None,
|
||||
resets_at_ms: 9_999_999,
|
||||
}),
|
||||
aid(1),
|
||||
nid(700),
|
||||
);
|
||||
|
||||
let err = fx
|
||||
.service
|
||||
.dispatch(&project(), cmd(ASK_JSON))
|
||||
.await
|
||||
.expect_err("a rate-limited delegated turn has no Final reply");
|
||||
assert_eq!(err.code(), "PROCESS");
|
||||
|
||||
let events = fx.bus.events();
|
||||
assert!(
|
||||
events.iter().any(|event| matches!(
|
||||
event,
|
||||
DomainEvent::AgentRateLimited {
|
||||
agent_id,
|
||||
resets_at_ms: None,
|
||||
} if *agent_id == aid(1)
|
||||
)),
|
||||
"target AgentRateLimited(None) must be published, got {events:?}"
|
||||
);
|
||||
assert!(
|
||||
events.iter().any(|event| matches!(
|
||||
event,
|
||||
DomainEvent::AgentRateLimitSuspected { agent_id, .. } if *agent_id == aid(1)
|
||||
)),
|
||||
"target AgentRateLimitSuspected must be published, got {events:?}"
|
||||
);
|
||||
assert!(
|
||||
!events.iter().any(|event| matches!(
|
||||
event,
|
||||
DomainEvent::AgentResumeScheduled { agent_id, .. } if *agent_id == aid(1)
|
||||
)),
|
||||
"target resume must not be scheduled without engine conversation_id, got {events:?}"
|
||||
);
|
||||
assert!(
|
||||
!fx.session_limits.cancel_resume(aid(1)),
|
||||
"no delegated resume should be armed without engine conversation_id"
|
||||
);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user