use serde::{Deserialize, Serialize};
use basemyai_core::libsql;
use super::{Memory, MemoryLayer};
use crate::{MemoryError, Result, now_unix};
fn storage(e: libsql::Error) -> MemoryError {
basemyai_core::CoreError::Storage(e.to_string()).into()
}
fn to_vec_literal(v: &[f32]) -> String {
let mut s = String::with_capacity(v.len() * 8 + 2);
s.push('[');
for (i, x) in v.iter().enumerate() {
if i > 0 {
s.push(',');
}
s.push_str(&x.to_string());
}
s.push(']');
s
}
const FORMAT: &str = "basemyai-export";
const VERSION: u32 = 1;
const EMBED_CHUNK: usize = 128;
#[derive(Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
enum ExportLine {
Header {
format: String,
version: u32,
agent_id: String,
embedding_model: String,
embedding_dim: usize,
exported_at: i64,
},
Memory {
id: String,
layer: String,
content: String,
valid_from: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
valid_until: Option<i64>,
#[serde(default)]
importance: f64,
#[serde(default, skip_serializing_if = "Option::is_none")]
last_access: Option<i64>,
},
Entity {
id: String,
kind: String,
label: String,
valid_from: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
valid_until: Option<i64>,
#[serde(default)]
importance: f64,
},
Edge {
src: String,
dst: String,
relation: String,
#[serde(default = "default_weight")]
weight: f64,
valid_from: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
valid_until: Option<i64>,
},
}
fn default_weight() -> f64 {
1.0
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
pub struct ImportReport {
pub memories: usize,
pub memories_skipped: usize,
pub entities: usize,
pub entities_skipped: usize,
pub edges: usize,
pub edges_skipped: usize,
}
impl Memory {
pub async fn export_jsonl(&self) -> Result<String> {
let conn = self.libsql_engine().store().connect();
let mut out = String::new();
push_line(
&mut out,
&ExportLine::Header {
format: FORMAT.to_string(),
version: VERSION,
agent_id: self.agent.as_str().to_string(),
embedding_model: self.embedder.model_id().to_string(),
embedding_dim: self.embedder.dim(),
exported_at: now_unix(),
},
)?;
let mut rows = conn
.query(
"SELECT id, layer, content, valid_from, valid_until, importance, last_access \
FROM memory WHERE agent_id = ?1 ORDER BY valid_from, id",
libsql::params![self.agent.as_str()],
)
.await
.map_err(storage)?;
while let Some(row) = rows.next().await.map_err(storage)? {
push_line(
&mut out,
&ExportLine::Memory {
id: text(&row, 0)?,
layer: text(&row, 1)?,
content: text(&row, 2)?,
valid_from: integer(&row, 3)?,
valid_until: integer_opt(&row, 4)?,
importance: real(&row, 5)?,
last_access: integer_opt(&row, 6)?,
},
)?;
}
let mut rows = conn
.query(
"SELECT id, kind, label, valid_from, valid_until, importance \
FROM entity WHERE agent_id = ?1 ORDER BY id",
libsql::params![self.agent.as_str()],
)
.await
.map_err(storage)?;
while let Some(row) = rows.next().await.map_err(storage)? {
push_line(
&mut out,
&ExportLine::Entity {
id: text(&row, 0)?,
kind: text(&row, 1)?,
label: text(&row, 2)?,
valid_from: integer(&row, 3)?,
valid_until: integer_opt(&row, 4)?,
importance: real(&row, 5)?,
},
)?;
}
let mut rows = conn
.query(
"SELECT src, dst, relation, weight, valid_from, valid_until \
FROM edge WHERE agent_id = ?1 ORDER BY src, dst, relation",
libsql::params![self.agent.as_str()],
)
.await
.map_err(storage)?;
while let Some(row) = rows.next().await.map_err(storage)? {
push_line(
&mut out,
&ExportLine::Edge {
src: text(&row, 0)?,
dst: text(&row, 1)?,
relation: text(&row, 2)?,
weight: real(&row, 3)?,
valid_from: integer(&row, 4)?,
valid_until: integer_opt(&row, 5)?,
},
)?;
}
Ok(out)
}
pub async fn import_jsonl(&self, jsonl: &str) -> Result<ImportReport> {
let mut header_seen = false;
let mut memories = Vec::new();
let mut entities = Vec::new();
let mut edges = Vec::new();
for (i, raw) in jsonl.lines().enumerate() {
let line = raw.trim();
if line.is_empty() {
continue;
}
let parsed: ExportLine = serde_json::from_str(line)
.map_err(|e| MemoryError::Porting(format!("ligne {} malformée : {e}", i + 1)))?;
if !header_seen && !matches!(parsed, ExportLine::Header { .. }) {
return Err(MemoryError::Porting(format!(
"ligne {} : données avant l'en-tête (fichier tronqué ?)",
i + 1
)));
}
match parsed {
ExportLine::Header { format, version, .. } => {
if header_seen {
return Err(MemoryError::Porting(format!("en-tête dupliqué (ligne {})", i + 1)));
}
if format != FORMAT {
return Err(MemoryError::Porting(format!("format inconnu : {format:?}")));
}
if version != VERSION {
return Err(MemoryError::Porting(format!(
"version d'export {version} non supportée (max {VERSION})"
)));
}
header_seen = true;
}
ExportLine::Memory {
id,
layer,
content,
valid_from,
valid_until,
importance,
last_access,
} => {
let layer = MemoryLayer::from_table(&layer)?;
memories.push((id, layer, content, valid_from, valid_until, importance, last_access));
}
ExportLine::Entity {
id,
kind,
label,
valid_from,
valid_until,
importance,
} => entities.push((id, kind, label, valid_from, valid_until, importance)),
ExportLine::Edge {
src,
dst,
relation,
weight,
valid_from,
valid_until,
} => edges.push((src, dst, relation, weight, valid_from, valid_until)),
}
}
if !header_seen {
return Err(MemoryError::Porting(
"en-tête absent : ce n'est pas un export BaseMyAI".into(),
));
}
let contents: Vec<String> = memories.iter().map(|m| m.2.clone()).collect();
let mut vectors = Vec::with_capacity(contents.len());
for chunk in contents.chunks(EMBED_CHUNK) {
vectors.extend(self.embedder.embed_batch(chunk)?);
}
let mut report = ImportReport::default();
let txn = self.libsql_engine().store().begin_write().await?;
for ((id, layer, content, valid_from, valid_until, importance, last_access), vector) in
memories.iter().zip(&vectors)
{
let inserted = txn
.execute(
"INSERT OR IGNORE INTO memory \
(id, agent_id, layer, content, valid_from, valid_until, importance, last_access, emb) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, vector(?9))",
libsql::params![
id.as_str(),
self.agent.as_str(),
layer.table(),
content.as_str(),
*valid_from,
*valid_until,
*importance,
*last_access,
to_vec_literal(vector),
],
)
.await
.map_err(storage)?;
if inserted > 0 {
txn.execute(
"INSERT INTO memory_fts (id, agent_id, content) VALUES (?1, ?2, ?3)",
libsql::params![id.as_str(), self.agent.as_str(), content.as_str()],
)
.await
.map_err(storage)?;
report.memories += 1;
} else {
report.memories_skipped += 1;
}
}
for (id, kind, label, valid_from, valid_until, importance) in &entities {
let inserted = txn
.execute(
"INSERT OR IGNORE INTO entity (id, agent_id, kind, label, valid_from, valid_until, importance) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
libsql::params![
id.as_str(),
self.agent.as_str(),
kind.as_str(),
label.as_str(),
*valid_from,
*valid_until,
*importance,
],
)
.await
.map_err(storage)?;
if inserted > 0 {
report.entities += 1;
} else {
report.entities_skipped += 1;
}
}
for (src, dst, relation, weight, valid_from, valid_until) in &edges {
let inserted = txn
.execute(
"INSERT OR IGNORE INTO edge (src, dst, agent_id, relation, weight, valid_from, valid_until) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
libsql::params![
src.as_str(),
dst.as_str(),
self.agent.as_str(),
relation.as_str(),
*weight,
*valid_from,
*valid_until,
],
)
.await
.map_err(storage)?;
if inserted > 0 {
report.edges += 1;
} else {
report.edges_skipped += 1;
}
}
txn.commit().await?;
Ok(report)
}
}
fn push_line(out: &mut String, line: &ExportLine) -> Result<()> {
let json = serde_json::to_string(line).map_err(|e| MemoryError::Porting(format!("sérialisation : {e}")))?;
out.push_str(&json);
out.push('\n');
Ok(())
}
fn text(row: &libsql::Row, idx: i32) -> Result<String> {
row.get::<String>(idx).map_err(storage)
}
fn integer(row: &libsql::Row, idx: i32) -> Result<i64> {
row.get::<i64>(idx).map_err(storage)
}
fn integer_opt(row: &libsql::Row, idx: i32) -> Result<Option<i64>> {
match row.get_value(idx).map_err(storage)? {
libsql::Value::Null => Ok(None),
libsql::Value::Integer(i) => Ok(Some(i)),
other => {
Err(basemyai_core::CoreError::Storage(format!("colonne {idx} : entier attendu, reçu {other:?}")).into())
}
}
}
fn real(row: &libsql::Row, idx: i32) -> Result<f64> {
match row.get_value(idx).map_err(storage)? {
libsql::Value::Real(r) => Ok(r),
#[allow(clippy::cast_precision_loss)]
libsql::Value::Integer(i) => Ok(i as f64),
other => Err(basemyai_core::CoreError::Storage(format!("colonne {idx} : réel attendu, reçu {other:?}")).into()),
}
}