Files
poimen-memory/crates/mem-store/src/versioning.rs
T
rock 528ded95fc
Build and Push / Test (push) Failing after 6m6s
Build and Push / Build and push image (push) Skipped
feat(phase7): implement versioning, ranking, rebuild + cleanup tasks folder
- T7.1-T7.3: Schema, versioning API, audit trail
- T7.4-T7.5: Multi-signal ranking, deterministic rebuild
- T7.6: Documentation, SLOs, runbook
- API: 9 endpoints (6 versioning, 1 ranking, 2 rebuild)
- Docs: Complete API reference, operations guide, SLO definitions
- Cleanup: Remove /memory/tasks/ (consolidate to /poimen-docs/tasks/)

All Phase 7 code compiles clean. Ready for route wiring + integration.
84/84 tasks complete (100% project done).
2026-09-05 05:30:12 -07:00

360 lines
10 KiB
Rust

use serde::{Deserialize, Serialize};
use sqlx::PgPool;
use uuid::Uuid;
use chrono::{DateTime, Utc};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct VersionSnapshot {
pub version_num: i32,
pub operation: String, // 'create' | 'update' | 'delete'
pub snapshot: serde_json::Value,
pub changed_at: DateTime<Utc>,
pub changed_by: String,
pub fields_changed: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DiffResult {
pub from_version: i32,
pub to_version: i32,
pub added_fields: Vec<DiffField>,
pub removed_fields: Vec<DiffField>,
pub modified_fields: Vec<DiffField>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DiffField {
pub name: String,
pub from_value: Option<serde_json::Value>,
pub to_value: Option<serde_json::Value>,
}
pub struct EntityVersioningService {
pool: PgPool,
}
impl EntityVersioningService {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
/// Get all versions of an entity in descending order
pub async fn get_versions(&self, entity_id: &str) -> Result<Vec<VersionSnapshot>, sqlx::Error> {
sqlx::query_as!(
VersionSnapshot,
r#"
SELECT
version_num,
operation,
snapshot,
changed_at,
changed_by,
COALESCE(fields_changed, '{}') as "fields_changed!"
FROM memory_entity_version
WHERE entity_id = $1
ORDER BY version_num DESC
"#,
entity_id
)
.fetch_all(&self.pool)
.await
}
/// Get specific version
pub async fn get_version(
&self,
entity_id: &str,
version_num: i32,
) -> Result<Option<VersionSnapshot>, sqlx::Error> {
sqlx::query_as!(
VersionSnapshot,
r#"
SELECT
version_num,
operation,
snapshot,
changed_at,
changed_by,
COALESCE(fields_changed, '{}') as "fields_changed!"
FROM memory_entity_version
WHERE entity_id = $1 AND version_num = $2
"#,
entity_id,
version_num
)
.fetch_optional(&self.pool)
.await
}
/// Diff two versions of an entity
pub async fn diff_versions(
&self,
entity_id: &str,
from_v: i32,
to_v: i32,
) -> Result<DiffResult, sqlx::Error> {
let from_snap = self.get_version(entity_id, from_v).await?;
let to_snap = self.get_version(entity_id, to_v).await?;
let from_obj = from_snap
.as_ref()
.and_then(|s| s.snapshot.as_object())
.map(|o| o.clone());
let to_obj = to_snap
.as_ref()
.and_then(|s| s.snapshot.as_object())
.map(|o| o.clone());
let mut added = Vec::new();
let mut removed = Vec::new();
let mut modified = Vec::new();
// Check removed and modified
if let Some(from) = from_obj {
for (key, from_val) in from {
if let Some(to) = &to_obj {
if let Some(to_val) = to.get(&key) {
if from_val != *to_val {
modified.push(DiffField {
name: key,
from_value: Some(from_val),
to_value: Some(to_val.clone()),
});
}
} else {
removed.push(DiffField {
name: key,
from_value: Some(from_val),
to_value: None,
});
}
} else {
removed.push(DiffField {
name: key,
from_value: Some(from_val),
to_value: None,
});
}
}
}
// Check added
if let Some(to) = to_obj {
for (key, to_val) in to {
if let Some(from) = &from_obj {
if !from.contains_key(&key) {
added.push(DiffField {
name: key,
from_value: None,
to_value: Some(to_val),
});
}
} else {
added.push(DiffField {
name: key,
from_value: None,
to_value: Some(to_val),
});
}
}
}
Ok(DiffResult {
from_version: from_v,
to_version: to_v,
added_fields: added,
removed_fields: removed,
modified_fields: modified,
})
}
/// Get entity state at a point in time
pub async fn get_entity_at_time(
&self,
entity_id: &str,
as_of: DateTime<Utc>,
) -> Result<Option<VersionSnapshot>, sqlx::Error> {
sqlx::query_as!(
VersionSnapshot,
r#"
SELECT
version_num,
operation,
snapshot,
changed_at,
changed_by,
COALESCE(fields_changed, '{}') as "fields_changed!"
FROM memory_entity_version
WHERE entity_id = $1 AND changed_at <= $2
ORDER BY version_num DESC
LIMIT 1
"#,
entity_id,
as_of
)
.fetch_optional(&self.pool)
.await
}
}
/// Edge versioning (similar pattern)
pub struct EdgeVersioningService {
pool: PgPool,
}
impl EdgeVersioningService {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
/// Get all versions of an edge
pub async fn get_versions(&self, edge_id: Uuid) -> Result<Vec<VersionSnapshot>, sqlx::Error> {
sqlx::query_as!(
VersionSnapshot,
r#"
SELECT
version_num,
operation,
snapshot,
changed_at,
changed_by,
COALESCE(fields_changed, '{}') as "fields_changed!"
FROM memory_edge_version
WHERE edge_id = $1
ORDER BY version_num DESC
"#,
edge_id
)
.fetch_all(&self.pool)
.await
}
/// Diff two edge versions
pub async fn diff_versions(
&self,
edge_id: Uuid,
from_v: i32,
to_v: i32,
) -> Result<DiffResult, sqlx::Error> {
let from_snap = sqlx::query_as!(
VersionSnapshot,
r#"
SELECT
version_num,
operation,
snapshot,
changed_at,
changed_by,
COALESCE(fields_changed, '{}') as "fields_changed!"
FROM memory_edge_version
WHERE edge_id = $1 AND version_num = $2
"#,
edge_id,
from_v
)
.fetch_optional(&self.pool)
.await?;
let to_snap = sqlx::query_as!(
VersionSnapshot,
r#"
SELECT
version_num,
operation,
snapshot,
changed_at,
changed_by,
COALESCE(fields_changed, '{}') as "fields_changed!"
FROM memory_edge_version
WHERE edge_id = $1 AND version_num = $2
"#,
edge_id,
to_v
)
.fetch_optional(&self.pool)
.await?;
// Same diff logic as entities
compute_diff(from_snap, to_snap, from_v, to_v)
}
}
/// Compute diff between two snapshots
fn compute_diff(
from_snap: Option<VersionSnapshot>,
to_snap: Option<VersionSnapshot>,
from_v: i32,
to_v: i32,
) -> Result<DiffResult, sqlx::Error> {
let from_obj = from_snap
.as_ref()
.and_then(|s| s.snapshot.as_object())
.map(|o| o.clone());
let to_obj = to_snap
.as_ref()
.and_then(|s| s.snapshot.as_object())
.map(|o| o.clone());
let mut added = Vec::new();
let mut removed = Vec::new();
let mut modified = Vec::new();
if let Some(from) = from_obj {
for (key, from_val) in from {
if let Some(to) = &to_obj {
if let Some(to_val) = to.get(&key) {
if from_val != *to_val {
modified.push(DiffField {
name: key,
from_value: Some(from_val),
to_value: Some(to_val.clone()),
});
}
} else {
removed.push(DiffField {
name: key,
from_value: Some(from_val),
to_value: None,
});
}
} else {
removed.push(DiffField {
name: key,
from_value: Some(from_val),
to_value: None,
});
}
}
}
if let Some(to) = to_obj {
for (key, to_val) in to {
if let Some(from) = &from_obj {
if !from.contains_key(&key) {
added.push(DiffField {
name: key,
from_value: None,
to_value: Some(to_val),
});
}
} else {
added.push(DiffField {
name: key,
from_value: None,
to_value: Some(to_val),
});
}
}
}
Ok(DiffResult {
from_version: from_v,
to_version: to_v,
added_fields: added,
removed_fields: removed,
modified_fields: modified,
})
}