feat(persistence): P6b — câblage live du checkpoint conversationnel (best-effort)
ask_agent persiste désormais chaque paire dans son conversationId : tour Prompt à l'enqueue, tour Response au succès — best-effort (un échec de persistance ne dégrade jamais la délégation). - application : port RecordTurnProvider (matérialise un RecordTurn sur le bon project root — OrchestratorService est mono-instance multi-projets, le log est par root) ; wither with_record_turn(provider, Clock) ; helper record_turn_best_effort ; horodatage via domain::ports::Clock (pas d'horloge infra dans application) - app-tauri : AppRecordTurnProvider (Fs* sur le root du projet) câblé au composition root avec le SystemClock partagé - tests : 5 cas (paire Prompt→Response, fil A↔B vs User↔B, no-op sans provider, ask Ok même si record échoue) ; orchestrator_service 40, application + app-tauri verts Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@ -22,11 +22,15 @@ use tokio::sync::Mutex as AsyncMutex;
|
||||
use domain::conversation::{
|
||||
ConversationParty, ConversationRegistry, SessionRef, WaitForGraph,
|
||||
};
|
||||
use domain::input::InputMediator;
|
||||
use domain::conversation_log::{ConversationTurn, TurnId, TurnRole};
|
||||
use domain::input::{InputMediator, InputSource};
|
||||
use domain::mailbox::{Ticket, TicketId};
|
||||
use domain::ports::{EventBus, ProfileStore, PtyHandle};
|
||||
use domain::ports::{Clock, EventBus, ProfileStore, PtyHandle};
|
||||
use domain::project::ProjectPath;
|
||||
use domain::{AgentId, DomainEvent, OrchestratorCommand, OrchestratorVisibility, ProfileId, Project};
|
||||
|
||||
use crate::conversation::RecordTurn;
|
||||
|
||||
use crate::agent::{
|
||||
CreateAgentFromScratch, CreateAgentInput, LaunchAgent, LaunchAgentInput, ListAgents,
|
||||
ListAgentsInput, McpRuntime, ReattachDecision, UpdateAgentContext, UpdateAgentContextInput,
|
||||
@ -83,6 +87,27 @@ pub trait McpRuntimeProvider: Send + Sync {
|
||||
fn runtime_for(&self, project: &Project, agent_id: AgentId) -> Option<McpRuntime>;
|
||||
}
|
||||
|
||||
/// Fournit le use case [`RecordTurn`] **lié au project root** de la délégation en
|
||||
/// cours (lot P6b).
|
||||
///
|
||||
/// L'[`OrchestratorService`] est **unique et partagé par tous les projets ouverts**
|
||||
/// (un seul `Arc` au composition root), alors que le log/handoff conversationnel est
|
||||
/// **par project root** (`<root>/.ideai/conversations/`, comme la mémoire). Les
|
||||
/// adapters `Fs*` de P6a fixent leur racine à la construction et leur port ne porte
|
||||
/// pas le root par appel ; on ne peut donc pas figer un `RecordTurn` global. Ce
|
||||
/// **port** lève la tension : `ask_agent` connaît le `project.root` du tour et demande
|
||||
/// au provider de matérialiser un `RecordTurn` ciblant le **bon** dossier (les adapters
|
||||
/// `Fs*` ne font que des jointures de chemin, leur construction est triviale).
|
||||
///
|
||||
/// `None` ⇒ aucune persistance (best-effort absente) : zéro régression pour les call
|
||||
/// sites/tests qui ne le branchent pas. Implémenté dans app-tauri (seul détenteur des
|
||||
/// adapters `Fs*`).
|
||||
pub trait RecordTurnProvider: Send + Sync {
|
||||
/// Construit le [`RecordTurn`] dont le log/handoff ciblent `root`. Appelé une fois
|
||||
/// par checkpoint best-effort ; `None` ⇒ on saute silencieusement la persistance.
|
||||
fn record_turn_for(&self, root: &ProjectPath) -> Option<Arc<RecordTurn>>;
|
||||
}
|
||||
|
||||
/// Dispatches validated orchestrator commands to the agent/terminal use cases.
|
||||
pub struct OrchestratorService {
|
||||
create_agent: Arc<CreateAgentFromScratch>,
|
||||
@ -147,6 +172,16 @@ pub struct OrchestratorService {
|
||||
/// [`Self::with_context_guard`] ; `None` ⇒ les commandes `context.*`/`memory.*`
|
||||
/// renvoient une erreur typée (call sites/tests legacy restent verts).
|
||||
context_guard: Option<Arc<ContextGuardUseCases>>,
|
||||
/// Provider du use case [`RecordTurn`] **par project root** (lot P6b), pour
|
||||
/// persister **best-effort** le Prompt et la Response de chaque paire déléguée dans
|
||||
/// le bon `conversationId`. Injecté avec son horloge via [`Self::with_record_turn`] ;
|
||||
/// `None` ⇒ aucune persistance (zéro régression pour les call sites/tests legacy).
|
||||
/// Un échec de persistance ne transforme **jamais** un succès de délégation en erreur.
|
||||
record_turn: Option<Arc<dyn RecordTurnProvider>>,
|
||||
/// Horloge millis (port [`Clock`]) pour estampiller `at_ms` des tours persistés —
|
||||
/// injectée avec [`Self::record_turn`]. La couche `application` reste **pure** : pas
|
||||
/// de `SystemTime::now()` brut ici, le temps vient du port injecté au composition root.
|
||||
clock: Option<Arc<dyn Clock>>,
|
||||
}
|
||||
|
||||
/// Bundle des quatre use cases C7 sous [`domain::fileguard::FileGuard`], injectés
|
||||
@ -207,6 +242,8 @@ impl OrchestratorService {
|
||||
ask_locks: StdMutex::new(HashMap::new()),
|
||||
mcp_runtime_provider: None,
|
||||
context_guard: None,
|
||||
record_turn: None,
|
||||
clock: None,
|
||||
}
|
||||
}
|
||||
|
||||
@ -279,6 +316,58 @@ impl OrchestratorService {
|
||||
self
|
||||
}
|
||||
|
||||
/// Branche le provider de [`RecordTurn`] **par project root** + son horloge (lot
|
||||
/// P6b) pour persister **best-effort** le Prompt et la Response de chaque paire
|
||||
/// déléguée. Builder additif : signature de [`Self::new`] **inchangée** (les
|
||||
/// tests/call sites legacy restent verts ; `None` ⇒ aucune persistance, donc aucune
|
||||
/// régression). Un échec de persistance ne casse jamais la délégation live.
|
||||
#[must_use]
|
||||
pub fn with_record_turn(
|
||||
mut self,
|
||||
record_turn: Arc<dyn RecordTurnProvider>,
|
||||
clock: Arc<dyn Clock>,
|
||||
) -> Self {
|
||||
self.record_turn = Some(record_turn);
|
||||
self.clock = Some(clock);
|
||||
self
|
||||
}
|
||||
|
||||
/// Persiste **best-effort** un tour (`Prompt`/`Response`) dans `conversation`.
|
||||
///
|
||||
/// No-op silencieux quand le provider/horloge ne sont pas câblés, quand le provider
|
||||
/// ne rend pas de [`RecordTurn`] pour ce root, ou quand l'`append`/`save` échoue : la
|
||||
/// persistance ne doit **jamais** transformer un succès de délégation en erreur, ni
|
||||
/// paniquer (contrat P6b). N'ajoute pas de latence inutile (un seul `await` borné par
|
||||
/// les adapters `Fs*`, déjà sérialisés par le verrou de tour de la cible).
|
||||
async fn record_turn_best_effort(
|
||||
&self,
|
||||
root: &ProjectPath,
|
||||
conversation: domain::conversation::ConversationId,
|
||||
source: InputSource,
|
||||
role: TurnRole,
|
||||
text: String,
|
||||
) {
|
||||
let (Some(provider), Some(clock)) = (&self.record_turn, &self.clock) else {
|
||||
return;
|
||||
};
|
||||
let Some(record) = provider.record_turn_for(root) else {
|
||||
return;
|
||||
};
|
||||
let at_ms = u64::try_from(clock.now_millis()).unwrap_or(0);
|
||||
let turn = ConversationTurn::new(
|
||||
TurnId::new_random(),
|
||||
conversation,
|
||||
at_ms,
|
||||
source,
|
||||
role,
|
||||
text,
|
||||
);
|
||||
// Best-effort : un échec de persistance est avalé (le contrat P6b interdit qu'il
|
||||
// remonte). Pas de framework de log dans `application` ; on reste cohérent avec
|
||||
// les autres effets best-effort du service (cf. publication `AgentReplied`).
|
||||
let _ = record.record(conversation, turn).await;
|
||||
}
|
||||
|
||||
/// Dispatches a validated command against `project`.
|
||||
///
|
||||
/// # Errors
|
||||
@ -674,6 +763,21 @@ impl OrchestratorService {
|
||||
// porte la source (Human/Agent) et la conversation cible.
|
||||
let requester_label = self.requester_label(project, requester).await;
|
||||
let ticket_id = TicketId::new_random();
|
||||
// Checkpoint Prompt (P6b, best-effort) : persister l'invite AVANT que `task` ne
|
||||
// soit déplacé dans le `Ticket`. Source = origine de la requête (agent demandeur
|
||||
// `Some(from)` ⇒ Agent, sinon Humain) — **même** `InputSource` que le ticket.
|
||||
let prompt_source = match requester {
|
||||
Some(from) => InputSource::agent(from),
|
||||
None => InputSource::Human,
|
||||
};
|
||||
self.record_turn_best_effort(
|
||||
&project.root,
|
||||
conversation_id,
|
||||
prompt_source,
|
||||
TurnRole::Prompt,
|
||||
task.clone(),
|
||||
)
|
||||
.await;
|
||||
let ticket = match requester {
|
||||
Some(from) => {
|
||||
Ticket::from_agent(ticket_id, from, conversation_id, requester_label, task)
|
||||
@ -688,7 +792,20 @@ impl OrchestratorService {
|
||||
// 3. Attendre la réponse, bornée. Timeout/canal fermé ⇒ retirer le ticket
|
||||
// (cible laissée vivante) et renvoyer une erreur typée.
|
||||
match tokio::time::timeout(ASK_AGENT_TIMEOUT, pending).await {
|
||||
Ok(Ok(result)) => Ok(self.reply_outcome(agent_id, &target, result)),
|
||||
Ok(Ok(result)) => {
|
||||
// Checkpoint Response (P6b, best-effort) : persister la réponse AVANT de
|
||||
// déplacer `result` dans `reply_outcome`. Source = la **cible** (c'est
|
||||
// elle qui a rendu le tour) ; même conversation que le Prompt.
|
||||
self.record_turn_best_effort(
|
||||
&project.root,
|
||||
conversation_id,
|
||||
InputSource::agent(agent_id),
|
||||
TurnRole::Response,
|
||||
result.clone(),
|
||||
)
|
||||
.await;
|
||||
Ok(self.reply_outcome(agent_id, &target, result))
|
||||
}
|
||||
Ok(Err(_cancelled)) => {
|
||||
mailbox.cancel_head(agent_id, ticket_id);
|
||||
Err(AppError::Process(format!(
|
||||
|
||||
Reference in New Issue
Block a user