feat(persistence): couche conversationnelle — cadrage §18/§19 + briques P1→P4
Resync ARCHITECTURE.md (état livré + cadrage persistance/handoff) et premières briques de la couche de persistance conversationnelle (log canonique par paire + handoff incrémental), indépendante du provider — prépare reprise fiable et handoff cross-profile Claude↔Codex. ARCHITECTURE.md - §14.3.2/§17 : M5 marqué livré, verrou « ouvert » périmé, §17 réconcilié (vue = terminal de sortie, pas d'UI chat) ; §18 état livré 2026-06-12 ; §19 cadrage persistance/handoff (log par paire + handoff, 10 lots P1→P10) Domaine (conversation_log.rs, pur) - P1 : ConversationTurn / TurnId / TurnRole + port ConversationLog - P3 : Handoff + port HandoffStore - P4 : port HandoffSummarizer (async, seam OCP pour adapter LLM futur) Infrastructure (conversation_log/) - P2 : FsConversationLog — JSONL append-only par paire, sync_all (durabilité crash), skip ligne corrompue, fichier absent ⇒ vide - P3 : FsHandoffStore — handoff.md front-matter, write atomique tmp+rename - P4 : HeuristicHandoffSummarizer — incrémental, zéro modèle/I/O, fenêtre WINDOW Tests : domaine 12 + infra 24 (conversation_log) verts, suites complètes sans régression. Cycle dev/test : le binôme a débusqué et corrigé un bug de durabilité (append sans flush) au passage. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
559
crates/domain/src/conversation_log.rs
Normal file
559
crates/domain/src/conversation_log.rs
Normal file
@ -0,0 +1,559 @@
|
||||
//! Log canonique de conversation (cadrage « persistance conversationnelle », lot P1).
|
||||
//!
|
||||
//! Aujourd'hui la continuité d'un fil repose sur le `resumable_id` CLI du provider
|
||||
//! ([`crate::conversation::Conversation::resumable_id`]) : fragile (perdu si le run
|
||||
//! est nettoyé) et **non portable** d'un provider à l'autre (un swap Claude→Codex
|
||||
//! repart de zéro). Pour y remédier, IdeA tient sa **propre mémoire de conversation**,
|
||||
//! indépendante du provider, dont la **source de vérité durable** est ce **log
|
||||
//! canonique append-only, par conversation (paire)** (ARCHITECTURE §19, décision
|
||||
//! D19-1a).
|
||||
//!
|
||||
//! Ce module est **pur** (règle de dépendance du domaine) : ni `tokio`, ni `std::fs`,
|
||||
//! ni I/O. Il possède le value object [`ConversationTurn`] (un tour de conversation),
|
||||
//! ses newtypes/énumérés ([`TurnId`], [`TurnRole`]) et le port driven
|
||||
//! [`ConversationLog`]. L'écriture/lecture réelle du `log.jsonl` est une affaire
|
||||
//! d'infrastructure (l'adapter `FsConversationLog`, lot P2).
|
||||
//!
|
||||
//! ## Frontière (D19-2)
|
||||
//!
|
||||
//! Ce log est **distinct** de la mémoire durable `.ideai/memory/` (savoir projet
|
||||
//! stable, peu bruité) et de la live-state (busy, sessions en cours). Le log est
|
||||
//! **volumineux et bruité** par nature : il ne **pollue jamais** `memory/`.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::conversation::ConversationId;
|
||||
use crate::input::InputSource;
|
||||
use crate::ports::StoreError;
|
||||
|
||||
/// Identifie un [`ConversationTurn`] dans le log d'une conversation.
|
||||
///
|
||||
/// Newtype autour d'[`uuid::Uuid`], calqué sur [`crate::conversation::ConversationId`]
|
||||
/// / [`crate::mailbox::TicketId`]. Sa monotonie (un id frappé à l'`append`) sert de
|
||||
/// **curseur** à la relecture incrémentale ([`ConversationLog::read`] avec `since`).
|
||||
#[derive(
|
||||
Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize,
|
||||
)]
|
||||
#[serde(transparent)]
|
||||
pub struct TurnId(pub uuid::Uuid);
|
||||
|
||||
impl TurnId {
|
||||
/// Enrobe un [`uuid::Uuid`] existant.
|
||||
#[must_use]
|
||||
pub const fn from_uuid(id: uuid::Uuid) -> Self {
|
||||
Self(id)
|
||||
}
|
||||
|
||||
/// Frappe un nouvel identifiant de tour aléatoire.
|
||||
#[must_use]
|
||||
pub fn new_random() -> Self {
|
||||
Self(uuid::Uuid::new_v4())
|
||||
}
|
||||
|
||||
/// Renvoie l'[`uuid::Uuid`] interne.
|
||||
#[must_use]
|
||||
pub const fn as_uuid(&self) -> uuid::Uuid {
|
||||
self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl std::fmt::Display for TurnId {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
write!(f, "{}", self.0)
|
||||
}
|
||||
}
|
||||
|
||||
/// Nature d'un tour dans le fil (qui parle / ce qui se passe).
|
||||
///
|
||||
/// Sépare l'invite émise vers l'agent ([`TurnRole::Prompt`]), sa réponse
|
||||
/// ([`TurnRole::Response`]) et l'activité outillée intermédiaire
|
||||
/// ([`TurnRole::ToolActivity`]) — assez pour rejouer/résumer fidèlement le travail
|
||||
/// utile sans dépendre du format d'un provider donné (D19-3).
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub enum TurnRole {
|
||||
/// Une invite adressée à l'agent (humaine ou déléguée par un autre agent).
|
||||
Prompt,
|
||||
/// Une réponse rendue par l'agent (le `result`/`final` d'un tour).
|
||||
Response,
|
||||
/// Une activité outillée intermédiaire (appel d'outil, trace) au sein d'un tour.
|
||||
ToolActivity,
|
||||
}
|
||||
|
||||
/// Un tour de conversation : qui a parlé, quand, dans quel fil, et quoi.
|
||||
///
|
||||
/// Value object pur (aucun comportement, aucune I/O), sérialisable : c'est l'unité
|
||||
/// **append-only** persistée (une ligne JSON par tour dans `log.jsonl`, lot P2). La
|
||||
/// [`InputSource`] est la *source de vérité* de l'origine (Humain ou Agent), réutilisée
|
||||
/// telle quelle depuis le domaine de l'entrée médiée.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct ConversationTurn {
|
||||
/// Identifiant stable de ce tour (frappé à l'`append`) ; sert aussi de curseur.
|
||||
pub id: TurnId,
|
||||
/// La conversation (paire) à laquelle ce tour appartient.
|
||||
pub conversation: ConversationId,
|
||||
/// Horodatage (epoch millisecondes) du tour.
|
||||
pub at_ms: u64,
|
||||
/// L'origine de l'entrée (Humain ou Agent délégant).
|
||||
pub source: InputSource,
|
||||
/// La nature du tour (invite, réponse, activité outillée).
|
||||
pub role: TurnRole,
|
||||
/// Le contenu textuel du tour.
|
||||
pub text: String,
|
||||
}
|
||||
|
||||
impl ConversationTurn {
|
||||
/// Construit un tour à partir de ses composants.
|
||||
#[must_use]
|
||||
pub fn new(
|
||||
id: TurnId,
|
||||
conversation: ConversationId,
|
||||
at_ms: u64,
|
||||
source: InputSource,
|
||||
role: TurnRole,
|
||||
text: impl Into<String>,
|
||||
) -> Self {
|
||||
Self {
|
||||
id,
|
||||
conversation,
|
||||
at_ms,
|
||||
source,
|
||||
role,
|
||||
text: text.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Le log canonique append-only des conversations, par paire (port driven, lot P1).
|
||||
///
|
||||
/// **Source de vérité durable** (D19-1a) : on **ajoute** un tour à chaque checkpoint
|
||||
/// (fin de tour, D19-5) et on **relit** à la reprise ou pour recalculer le handoff. Le
|
||||
/// log est gardé **par conversation** : deux fils distincts n'interfèrent jamais.
|
||||
///
|
||||
/// `#[async_trait]` et erreur [`StoreError`] comme les autres ports driven de
|
||||
/// persistance ([`crate::ports::TemplateStore`]…) : injecté en
|
||||
/// `Arc<dyn ConversationLog>` au composition root, donc gardé object-safe (les ports
|
||||
/// async dyn-compatibles passent par `async_trait` qui box le futur — cf. note
|
||||
/// d'en-tête de [`crate::ports`]).
|
||||
#[async_trait::async_trait]
|
||||
pub trait ConversationLog: Send + Sync {
|
||||
/// Ajoute `turn` à la fin du log canonique de `conversation`.
|
||||
///
|
||||
/// # Errors
|
||||
/// [`StoreError`] en cas d'échec d'écriture ou de sérialisation.
|
||||
async fn append(
|
||||
&self,
|
||||
conversation: ConversationId,
|
||||
turn: ConversationTurn,
|
||||
) -> Result<(), StoreError>;
|
||||
|
||||
/// Relit les tours de `conversation`, dans l'ordre d'ajout.
|
||||
///
|
||||
/// `since` est un curseur **exclusif** : quand il vaut `Some(id)`, seuls les tours
|
||||
/// **postérieurs** au tour `id` sont renvoyés (relecture incrémentale pour le
|
||||
/// recalcul de handoff) ; `None` relit tout le fil depuis le début.
|
||||
///
|
||||
/// # Errors
|
||||
/// [`StoreError`] en cas d'échec de lecture ou de désérialisation.
|
||||
async fn read(
|
||||
&self,
|
||||
conversation: ConversationId,
|
||||
since: Option<TurnId>,
|
||||
) -> Result<Vec<ConversationTurn>, StoreError>;
|
||||
|
||||
/// Renvoie les `n` derniers tours de `conversation`, dans l'ordre d'ajout.
|
||||
///
|
||||
/// Renvoie moins de `n` éléments si le fil est plus court (et un `Vec` vide pour
|
||||
/// `n == 0` ou un fil vide). Sert au résumé incrémental (les N derniers tours).
|
||||
///
|
||||
/// # Errors
|
||||
/// [`StoreError`] en cas d'échec de lecture ou de désérialisation.
|
||||
async fn last(
|
||||
&self,
|
||||
conversation: ConversationId,
|
||||
n: usize,
|
||||
) -> Result<Vec<ConversationTurn>, StoreError>;
|
||||
}
|
||||
|
||||
/// Résumé cumulatif de reprise d'une conversation (handoff, lot P3).
|
||||
///
|
||||
/// Value object pur (aucune I/O, aucun comportement) : c'est le **point de reprise**
|
||||
/// portable d'une conversation (ARCHITECTURE §19.2/§19.3). À la reprise d'un fil — ou
|
||||
/// au swap d'un provider à l'autre — IdeA injecte ce `summary_md` plutôt que de
|
||||
/// dépendre du `resumable_id` CLI d'un provider donné. Le champ [`Handoff::up_to`]
|
||||
/// est le **curseur** ([`TurnId`]) jusqu'auquel le résumé a été calculé : un recalcul
|
||||
/// incrémental relit le log à partir de ce curseur ([`ConversationLog::read`] avec
|
||||
/// `since`) pour étendre le résumé aux tours postérieurs.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct Handoff {
|
||||
/// Le résumé cumulatif, en Markdown : ce qu'un agent doit savoir pour reprendre.
|
||||
pub summary_md: String,
|
||||
/// Le curseur : dernier [`TurnId`] couvert par ce résumé (borne du recalcul).
|
||||
pub up_to: TurnId,
|
||||
/// L'objectif courant de la conversation, le cas échéant (fil sans but explicite).
|
||||
pub objective: Option<String>,
|
||||
}
|
||||
|
||||
impl Handoff {
|
||||
/// Construit un handoff à partir de ses composants.
|
||||
#[must_use]
|
||||
pub fn new(
|
||||
summary_md: impl Into<String>,
|
||||
up_to: TurnId,
|
||||
objective: Option<String>,
|
||||
) -> Self {
|
||||
Self {
|
||||
summary_md: summary_md.into(),
|
||||
up_to,
|
||||
objective,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Le store du résumé de reprise (handoff) d'une conversation (port driven, lot P3).
|
||||
///
|
||||
/// Distinct du [`ConversationLog`] append-only : ici on garde **un seul** [`Handoff`]
|
||||
/// par conversation (le dernier point de reprise), écrasé à chaque recalcul. La
|
||||
/// persistance réelle (`handoff.md`) est une affaire d'infrastructure (l'adapter
|
||||
/// `FsHandoffStore`, lot P3).
|
||||
///
|
||||
/// `#[async_trait]` et erreur [`StoreError`] comme le port voisin [`ConversationLog`] :
|
||||
/// injecté en `Arc<dyn HandoffStore>` au composition root, gardé object-safe.
|
||||
#[async_trait::async_trait]
|
||||
pub trait HandoffStore: Send + Sync {
|
||||
/// Charge le handoff de `conversation`, ou `None` si aucun n'a encore été écrit.
|
||||
///
|
||||
/// L'absence d'un handoff n'est **jamais** une erreur (`Ok(None)`).
|
||||
///
|
||||
/// # Errors
|
||||
/// [`StoreError`] en cas d'échec de lecture ou de désérialisation.
|
||||
async fn load(&self, conversation: ConversationId) -> Result<Option<Handoff>, StoreError>;
|
||||
|
||||
/// Écrit (en écrasant) le handoff de `conversation`.
|
||||
///
|
||||
/// # Errors
|
||||
/// [`StoreError`] en cas d'échec d'écriture ou de sérialisation.
|
||||
async fn save(
|
||||
&self,
|
||||
conversation: ConversationId,
|
||||
handoff: Handoff,
|
||||
) -> Result<(), StoreError>;
|
||||
}
|
||||
|
||||
/// Le résumeur incrémental de handoff (port driving, lot P4).
|
||||
///
|
||||
/// **Replie** (au sens d'un *fold*) un [`Handoff`] : on part du résumé précédent (s'il
|
||||
/// existe) et on n'intègre que **l'incrément** de tours nouveaux ([`new_turns`]), sans
|
||||
/// jamais relire tout le fil — c'est le seul appelant (l'application) qui relit le log à
|
||||
/// partir de [`Handoff::up_to`] ([`ConversationLog::read`] avec `since`) pour fournir cet
|
||||
/// incrément.
|
||||
///
|
||||
/// ## Pourquoi `async` (OCP, décision tranchée par Main)
|
||||
///
|
||||
/// Le port est **async** même si l'implémentation heuristique zéro-dépendance
|
||||
/// ([`HeuristicHandoffSummarizer`] côté infra) n'`await` rien : on fige **maintenant** le
|
||||
/// seam pour qu'un futur `LlmHandoffSummarizer` (P10, qui fera de l'I/O réseau async) se
|
||||
/// **substitue sans modifier ni le trait ni l'application** (open/closed). `#[async_trait]`
|
||||
/// pour rester object-safe (injecté en `Arc<dyn HandoffSummarizer>` au composition root).
|
||||
///
|
||||
/// ## Pas de `Result` (best-effort, D19-6)
|
||||
///
|
||||
/// Le repli est **best-effort** : il ne doit **jamais** bloquer la persistance du log
|
||||
/// canonique (la source de vérité durable). Une heuristique ne peut pas échouer ; un futur
|
||||
/// adapter LLM qui échouerait renverra simplement le `prev` (ou un repli heuristique) plutôt
|
||||
/// que de propager une erreur. Le port ne porte donc **aucune** erreur.
|
||||
///
|
||||
/// [`new_turns`]: HandoffSummarizer::fold
|
||||
#[async_trait::async_trait]
|
||||
pub trait HandoffSummarizer: Send + Sync {
|
||||
/// Replie `prev` (s'il existe) avec **uniquement** `new_turns` et renvoie le handoff
|
||||
/// étendu.
|
||||
///
|
||||
/// `prev` est le dernier point de reprise connu (ou `None` pour un premier calcul) ;
|
||||
/// `new_turns` est l'incrément des tours postérieurs à `prev.up_to`, dans l'ordre
|
||||
/// d'ajout. Le résultat couvre jusqu'au dernier tour vu (cf. [`Handoff::up_to`]).
|
||||
/// Best-effort : ne renvoie jamais d'erreur (cf. note du trait).
|
||||
async fn fold(&self, prev: Option<Handoff>, new_turns: &[ConversationTurn]) -> Handoff;
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Mutex;
|
||||
|
||||
// ---------------------------------------------------------------------
|
||||
// Deterministic constructors (calqués sur conversation.rs / mailbox.rs)
|
||||
// ---------------------------------------------------------------------
|
||||
|
||||
fn conv_id(n: u128) -> ConversationId {
|
||||
ConversationId::from_uuid(uuid::Uuid::from_u128(n))
|
||||
}
|
||||
|
||||
fn turn_id(n: u128) -> TurnId {
|
||||
TurnId::from_uuid(uuid::Uuid::from_u128(n))
|
||||
}
|
||||
|
||||
fn turn(conv: ConversationId, id: TurnId, role: TurnRole, text: &str) -> ConversationTurn {
|
||||
ConversationTurn::new(id, conv, 1_000, InputSource::Human, role, text)
|
||||
}
|
||||
|
||||
fn texts(turns: &[ConversationTurn]) -> Vec<String> {
|
||||
turns.iter().map(|t| t.text.clone()).collect()
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------
|
||||
// In-memory double of the port (the normal way to lock a pure port's
|
||||
// contract before the Fs adapter lands in P2). Vec-backed, per conversation.
|
||||
//
|
||||
// The trait is async but the double's state is purely in-memory: the sync
|
||||
// Mutex is locked briefly and **never** held across an `.await` point.
|
||||
// ---------------------------------------------------------------------
|
||||
|
||||
#[derive(Default)]
|
||||
struct InMemoryConversationLog {
|
||||
threads: Mutex<HashMap<ConversationId, Vec<ConversationTurn>>>,
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl ConversationLog for InMemoryConversationLog {
|
||||
async fn append(
|
||||
&self,
|
||||
conversation: ConversationId,
|
||||
turn: ConversationTurn,
|
||||
) -> Result<(), StoreError> {
|
||||
self.threads
|
||||
.lock()
|
||||
.unwrap()
|
||||
.entry(conversation)
|
||||
.or_default()
|
||||
.push(turn);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn read(
|
||||
&self,
|
||||
conversation: ConversationId,
|
||||
since: Option<TurnId>,
|
||||
) -> Result<Vec<ConversationTurn>, StoreError> {
|
||||
let guard = self.threads.lock().unwrap();
|
||||
let Some(thread) = guard.get(&conversation) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
let out = match since {
|
||||
None => thread.clone(),
|
||||
Some(cursor) => {
|
||||
// Strictly-after semantics: everything past the cursor's position.
|
||||
match thread.iter().position(|t| t.id == cursor) {
|
||||
Some(idx) => thread[idx + 1..].to_vec(),
|
||||
None => Vec::new(),
|
||||
}
|
||||
}
|
||||
};
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
async fn last(
|
||||
&self,
|
||||
conversation: ConversationId,
|
||||
n: usize,
|
||||
) -> Result<Vec<ConversationTurn>, StoreError> {
|
||||
let guard = self.threads.lock().unwrap();
|
||||
let Some(thread) = guard.get(&conversation) else {
|
||||
return Ok(Vec::new());
|
||||
};
|
||||
let start = thread.len().saturating_sub(n);
|
||||
Ok(thread[start..].to_vec())
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------
|
||||
// Serde round-trips (verrouille camelCase / transparent)
|
||||
// ---------------------------------------------------------------------
|
||||
|
||||
#[test]
|
||||
fn turn_id_serializes_transparently() {
|
||||
let u = uuid::Uuid::from_u128(7);
|
||||
let id = TurnId::from_uuid(u);
|
||||
let json = serde_json::to_string(&id).unwrap();
|
||||
// `#[serde(transparent)]` => bare quoted uuid, no wrapper object.
|
||||
assert_eq!(json, format!("\"{u}\""));
|
||||
let back: TurnId = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(back, id);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn turn_role_serializes_camel_case() {
|
||||
assert_eq!(serde_json::to_string(&TurnRole::Prompt).unwrap(), "\"prompt\"");
|
||||
assert_eq!(
|
||||
serde_json::to_string(&TurnRole::Response).unwrap(),
|
||||
"\"response\""
|
||||
);
|
||||
assert_eq!(
|
||||
serde_json::to_string(&TurnRole::ToolActivity).unwrap(),
|
||||
"\"toolActivity\""
|
||||
);
|
||||
for role in [TurnRole::Prompt, TurnRole::Response, TurnRole::ToolActivity] {
|
||||
let json = serde_json::to_string(&role).unwrap();
|
||||
let back: TurnRole = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(back, role);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn conversation_turn_round_trips_in_camel_case() {
|
||||
let t = ConversationTurn::new(
|
||||
turn_id(3),
|
||||
conv_id(9),
|
||||
1_700_000_000_123,
|
||||
InputSource::Human,
|
||||
TurnRole::Prompt,
|
||||
"hello",
|
||||
);
|
||||
let json = serde_json::to_string(&t).unwrap();
|
||||
// camelCase field renaming on the struct.
|
||||
assert!(json.contains("\"atMs\":1700000000123"), "got: {json}");
|
||||
assert!(!json.contains("at_ms"), "snake_case leaked: {json}");
|
||||
let back: ConversationTurn = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(back, t);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn conversation_turn_round_trips_agent_source() {
|
||||
let from = crate::ids::AgentId::from_uuid(uuid::Uuid::from_u128(11));
|
||||
let t = ConversationTurn::new(
|
||||
turn_id(4),
|
||||
conv_id(9),
|
||||
42,
|
||||
InputSource::agent(from),
|
||||
TurnRole::ToolActivity,
|
||||
"ran tool",
|
||||
);
|
||||
let json = serde_json::to_string(&t).unwrap();
|
||||
let back: ConversationTurn = serde_json::from_str(&json).unwrap();
|
||||
assert_eq!(back, t);
|
||||
assert_eq!(back.source.as_agent(), Some(from));
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------
|
||||
// Port contract — via the in-memory double
|
||||
// ---------------------------------------------------------------------
|
||||
|
||||
#[tokio::test]
|
||||
async fn append_then_read_all_preserves_insertion_order() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
let c = conv_id(1);
|
||||
log.append(c, turn(c, turn_id(1), TurnRole::Prompt, "a"))
|
||||
.await
|
||||
.unwrap();
|
||||
log.append(c, turn(c, turn_id(2), TurnRole::Response, "b"))
|
||||
.await
|
||||
.unwrap();
|
||||
log.append(c, turn(c, turn_id(3), TurnRole::Prompt, "c"))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let all = log.read(c, None).await.unwrap();
|
||||
assert_eq!(texts(&all), vec!["a", "b", "c"]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn read_with_cursor_is_strictly_exclusive() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
let c = conv_id(1);
|
||||
for (id, txt) in [(1, "a"), (2, "b"), (3, "c")] {
|
||||
log.append(c, turn(c, turn_id(id), TurnRole::Prompt, txt))
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
// since = id of "a" => only the strictly-posterior tours, "a" excluded.
|
||||
let after_a = log.read(c, Some(turn_id(1))).await.unwrap();
|
||||
assert_eq!(texts(&after_a), vec!["b", "c"]);
|
||||
|
||||
let after_b = log.read(c, Some(turn_id(2))).await.unwrap();
|
||||
assert_eq!(texts(&after_b), vec!["c"]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn read_cursor_on_last_id_yields_empty() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
let c = conv_id(1);
|
||||
log.append(c, turn(c, turn_id(1), TurnRole::Prompt, "a"))
|
||||
.await
|
||||
.unwrap();
|
||||
log.append(c, turn(c, turn_id(2), TurnRole::Response, "b"))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let after_last = log.read(c, Some(turn_id(2))).await.unwrap();
|
||||
assert!(after_last.is_empty(), "cursor on last id => nothing after");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn last_zero_is_empty() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
let c = conv_id(1);
|
||||
log.append(c, turn(c, turn_id(1), TurnRole::Prompt, "a"))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(log.last(c, 0).await.unwrap().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn last_on_empty_thread_is_empty() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
// Never appended to this conversation.
|
||||
assert!(log.last(conv_id(99), 5).await.unwrap().is_empty());
|
||||
assert!(log.read(conv_id(99), None).await.unwrap().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn last_n_greater_than_len_returns_everything() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
let c = conv_id(1);
|
||||
for (id, txt) in [(1, "a"), (2, "b")] {
|
||||
log.append(c, turn(c, turn_id(id), TurnRole::Prompt, txt))
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
let last = log.last(c, 10).await.unwrap();
|
||||
assert_eq!(texts(&last), vec!["a", "b"]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn last_n_less_than_len_returns_the_n_last_in_order() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
let c = conv_id(1);
|
||||
for (id, txt) in [(1, "a"), (2, "b"), (3, "c"), (4, "d")] {
|
||||
log.append(c, turn(c, turn_id(id), TurnRole::Prompt, txt))
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
let last2 = log.last(c, 2).await.unwrap();
|
||||
assert_eq!(texts(&last2), vec!["c", "d"], "the 2 last, in insertion order");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn conversations_are_disjoint_threads() {
|
||||
let log = InMemoryConversationLog::default();
|
||||
let c1 = conv_id(1);
|
||||
let c2 = conv_id(2);
|
||||
log.append(c1, turn(c1, turn_id(1), TurnRole::Prompt, "c1-a"))
|
||||
.await
|
||||
.unwrap();
|
||||
log.append(c2, turn(c2, turn_id(2), TurnRole::Prompt, "c2-a"))
|
||||
.await
|
||||
.unwrap();
|
||||
log.append(c1, turn(c1, turn_id(3), TurnRole::Response, "c1-b"))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(texts(&log.read(c1, None).await.unwrap()), vec!["c1-a", "c1-b"]);
|
||||
assert_eq!(texts(&log.read(c2, None).await.unwrap()), vec!["c2-a"]);
|
||||
// A cursor from one thread never leaks tours from another.
|
||||
assert_eq!(texts(&log.last(c2, 10).await.unwrap()), vec!["c2-a"]);
|
||||
}
|
||||
}
|
||||
@ -32,6 +32,7 @@
|
||||
|
||||
pub mod agent;
|
||||
pub mod conversation;
|
||||
pub mod conversation_log;
|
||||
pub mod error;
|
||||
pub mod events;
|
||||
pub mod fileguard;
|
||||
@ -85,6 +86,10 @@ pub use conversation::{
|
||||
|
||||
pub use input::{AgentBusyState, InputMediator, InputSource};
|
||||
|
||||
pub use conversation_log::{
|
||||
ConversationLog, ConversationTurn, Handoff, HandoffStore, HandoffSummarizer, TurnId, TurnRole,
|
||||
};
|
||||
|
||||
pub use fileguard::{
|
||||
is_orchestrator, may_write_directly, FileGuard, GuardError, GuardedResource, ReadLease,
|
||||
WriteLease,
|
||||
|
||||
Reference in New Issue
Block a user