51 lines
1.5 KiB
Rust
51 lines
1.5 KiB
Rust
use mem_core::Record;
|
|||
|
|
use futures::stream::Stream;
|
||
|
|
|
||
|
|
/// A source of records, shaped as a stream from day one.
|
||
|
|
/// Sources decide how to produce records; the chunker never learns
|
||
|
|
/// whether they came from pi, claude, or a socket.
|
||
|
|
pub trait RecordSource {
|
||
|
|
fn records(self) -> Box<dyn Stream<Item = Result<Record, String>> + Unpin>;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// A test vector source that produces records from a Vec.
|
||
|
|
pub struct VecSource(pub Vec<Record>);
|
||
|
|
|
||
|
|
impl RecordSource for VecSource {
|
||
|
|
fn records(self) -> Box<dyn Stream<Item = Result<Record, String>> + Unpin> {
|
||
|
|
Box::new(futures::stream::iter(self.0.into_iter().map(Ok)))
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
#[cfg(test)]
|
||
|
|
mod tests {
|
||
|
|
use super::*;
|
||
|
|
use mem_core::{Provenance, Role};
|
||
|
|
use time::macros::datetime;
|
||
|
|
use futures::StreamExt;
|
||
|
|
|
||
|
|
#[tokio::test]
|
||
|
|
async fn test_vec_source() {
|
||
|
|
let records = vec![
|
||
|
|
Record {
|
||
|
|
role: Role::User,
|
||
|
|
text: "Hello".to_string(),
|
||
|
|
timestamp: datetime!(2024-08-20 12:00:00 UTC),
|
||
|
|
provenance: Provenance {
|
||
|
|
source_id: "session1".to_string(),
|
||
|
|
offset: 0,
|
||
|
|
},
|
||
|
|
},
|
||
|
|
];
|
||
|
|
|
||
|
|
let source = VecSource(records.clone());
|
||
|
|
let mut stream = source.records();
|
||
|
|
|
||
|
|
let result = stream.next().await;
|
||
|
|
assert!(result.is_some());
|
||
|
|
let record = result.unwrap().unwrap();
|
||
|
|
assert_eq!(record.role, Role::User);
|
||
|
|
assert_eq!(record.text, "Hello");
|
||
|
|
}
|
||
|
|
}
|