use lunaris_core::storage::types::{Lsn, WriteOp};
use lunaris_core::{
BiTemporal, Hlc, HlcClock, LunarisError, StorageError, StoragePort, ValidateError,
};
use serde::{Deserialize, Serialize};
use tracing::Instrument;
use ulid::Ulid;
use crate::audit::{AuditEvent, publish_audit_event};
use crate::handle::Lunaris;
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum ForgetTarget {
Id(Ulid),
Scope(ScopeSpec),
Before(Hlc),
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum ScopeSpec {
BySource(String),
ByMetadata(String, String),
ByEpisode(Ulid),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[non_exhaustive]
pub enum IndexKind {
Kv,
Vector,
Graph,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct ForgetReceipt {
pub target: ForgetTarget,
pub indices_affected: Vec<IndexKind>,
#[serde(default)]
pub matched: u64,
pub rows_written: u64,
pub rows_deleted: u64,
pub audit_lsn: Lsn,
pub preview: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct ForgetConfirmation {
pub(crate) for_audit_lsn: Lsn,
}
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ForgetOptions {
pub hard: bool,
pub dry_run: bool,
pub confirmation_token: Option<ForgetConfirmation>,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct ForgetRequest {
pub target: ForgetTarget,
pub options: ForgetOptions,
}
impl ForgetTarget {
pub fn hard(self) -> ForgetRequest {
ForgetRequest { target: self, options: ForgetOptions { hard: true, ..Default::default() } }
}
pub fn dry_run(self) -> ForgetRequest {
ForgetRequest {
target: self,
options: ForgetOptions { dry_run: true, ..Default::default() },
}
}
pub fn into_request(self) -> ForgetRequest {
ForgetRequest { target: self, options: ForgetOptions::default() }
}
}
impl ForgetRequest {
pub fn with_token(mut self, token: ForgetConfirmation) -> Self {
self.options.confirmation_token = Some(token);
self
}
pub fn hard_confirmed(target: ForgetTarget, token: ForgetConfirmation) -> Self {
ForgetRequest {
target,
options: ForgetOptions { hard: true, dry_run: false, confirmation_token: Some(token) },
}
}
}
impl From<ForgetTarget> for ForgetRequest {
fn from(target: ForgetTarget) -> Self {
target.into_request()
}
}
impl Lunaris {
#[deprecated(
since = "0.3.0",
note = "use `engine.scoped(scope).forget(request)` — Lunaris::forget routes through Scope::dev() and returns zero matches under any non-_dev_ scope"
)]
pub async fn forget(
&self,
request: impl Into<ForgetRequest>,
) -> Result<ForgetReceipt, LunarisError> {
let request = request.into();
let target_kind = match &request.target {
ForgetTarget::Id(_) => "id",
ForgetTarget::Scope(_) => "scope",
ForgetTarget::Before(_) => "before",
};
let span = tracing::info_span!(
"lunaris.forget",
correlation_id = tracing::field::Empty,
target = %target_kind,
hard = %request.options.hard,
dry_run = %request.options.dry_run,
);
async move {
tracing::warn!(
target: "lunaris::forget",
"Lunaris::forget is hard-coded to Scope::dev() in v0.2.x — \
same-scope forget under any non-`_dev_` scope silently \
returns zero matches. ScopedLunaris::forget(target) lands \
in v0.3. See CHANGELOG.md v0.2.0 Known issues."
);
if request.options.hard && request.options.confirmation_token.is_none() {
return Err(LunarisError::Validate(ValidateError::ConfirmationRequired(
"hard-delete requires confirmation token from prior dry_run".into(),
)));
}
let mut matches =
scan_matches(self.storage.as_ref(), &request.target, self.clock.as_ref()).await?;
drop_already_closed(&mut matches, request.options.hard);
if request.options.dry_run {
let receipt = ForgetReceipt {
target: request.target.clone(),
indices_affected: classify_indices(&request.target),
matched: matches.len() as u64,
rows_written: 0,
rows_deleted: 0,
audit_lsn: Lsn { wall_ms: 0, counter: 0 },
preview: true,
};
let audit_offset = publish_audit_event(
&self.storage,
&lunaris_core::Scope::dev(),
AuditEvent::Forget((&receipt).into()),
)
.await
.unwrap_or(0);
return Ok(ForgetReceipt {
audit_lsn: Lsn { wall_ms: audit_offset, counter: 0 },
..receipt
});
}
let ops: Vec<WriteOp> = if request.options.hard {
matches.iter().map(|m| WriteOp::KvDelete { key: m.key.clone() }).collect()
} else {
matches
.iter()
.map(|m| build_soft_delete_op(m, self.clock.as_ref()))
.collect::<Result<Vec<_>, _>>()?
};
if !ops.is_empty() {
let _lsn = self
.storage
.atomic_write(&lunaris_core::Scope::dev(), &ops)
.await
.map_err(LunarisError::Storage)?;
}
let receipt = ForgetReceipt {
target: request.target.clone(),
indices_affected: classify_indices(&request.target),
matched: matches.len() as u64,
rows_written: if request.options.hard { 0 } else { matches.len() as u64 },
rows_deleted: if request.options.hard { matches.len() as u64 } else { 0 },
audit_lsn: Lsn { wall_ms: 0, counter: 0 },
preview: false,
};
let audit_offset = publish_audit_event(
&self.storage,
&lunaris_core::Scope::dev(),
AuditEvent::Forget((&receipt).into()),
)
.await
.unwrap_or(0);
Ok(ForgetReceipt { audit_lsn: Lsn { wall_ms: audit_offset, counter: 0 }, ..receipt })
}
.instrument(span)
.await
}
pub async fn confirm_hard_forget(
&self,
dry_run_receipt: ForgetReceipt,
) -> Result<ForgetConfirmation, LunarisError> {
if !dry_run_receipt.preview {
return Err(LunarisError::Validate(ValidateError::ConfirmationRequired(
"confirm_hard_forget requires a dry_run receipt (preview=true)".into(),
)));
}
Ok(ForgetConfirmation { for_audit_lsn: dry_run_receipt.audit_lsn })
}
}
#[derive(Clone, Debug)]
struct ForgetMatch {
key: Vec<u8>,
payload: Vec<u8>,
bt: BiTemporal,
}
async fn scan_matches(
storage: &dyn StoragePort,
target: &ForgetTarget,
clock: &HlcClock,
) -> Result<Vec<ForgetMatch>, LunarisError> {
use futures::stream::StreamExt;
if let ForgetTarget::Id(ulid) = target {
let key_str = format!("episode:{ulid}");
let now = clock.tick();
let row = storage
.read_as_of(&lunaris_core::Scope::dev(), key_str.as_bytes(), now)
.await
.map_err(LunarisError::Storage)?;
return Ok(match row {
Some(r) => {
vec![ForgetMatch { key: key_str.into_bytes(), payload: r.value.to_vec(), bt: r.bt }]
}
None => Vec::new(),
});
}
let prefix: &[u8] = b"episode:";
let mut stream = storage
.scan_range(&lunaris_core::Scope::dev(), prefix, None)
.await
.map_err(LunarisError::Storage)?;
let mut out = Vec::new();
while let Some(item) = stream.next().await {
let (k, v) = item.map_err(LunarisError::Storage)?;
if matches_target(&k, &v, target) {
let now = clock.tick();
let row_opt = storage
.read_as_of(&lunaris_core::Scope::dev(), &k, now)
.await
.map_err(LunarisError::Storage)?;
let bt = row_opt
.as_ref()
.map(|r| r.bt)
.unwrap_or_else(|| BiTemporal { valid: (Hlc::ZERO, None), sys: (Hlc::ZERO, None) });
out.push(ForgetMatch { key: k.to_vec(), payload: v.to_vec(), bt });
}
}
Ok(out)
}
fn matches_target(_key: &[u8], value: &[u8], target: &ForgetTarget) -> bool {
match target {
ForgetTarget::Id(_) => true,
ForgetTarget::Scope(spec) => match_scope(value, spec),
ForgetTarget::Before(hlc) => match_before(value, *hlc),
}
}
fn match_scope(value: &[u8], spec: &ScopeSpec) -> bool {
let json: serde_json::Value = match serde_json::from_slice(value) {
Ok(j) => j,
Err(_) => return false,
};
match spec {
ScopeSpec::BySource(prefix) => json
.get("source")
.and_then(|s| s.as_str())
.map(|s| s.starts_with(prefix))
.unwrap_or(false),
ScopeSpec::ByMetadata(key, expected) => json
.get("metadata")
.and_then(|m| m.get(key))
.and_then(|v| v.as_str())
.map(|v| v == expected)
.unwrap_or(false),
ScopeSpec::ByEpisode(ulid) => {
json.get("id").and_then(|v| v.as_str()).map(|s| s == ulid.to_string()).unwrap_or(false)
}
}
}
fn match_before(value: &[u8], hlc: Hlc) -> bool {
let json: serde_json::Value = match serde_json::from_slice(value) {
Ok(j) => j,
Err(_) => return false,
};
json.get("bt")
.and_then(|bt| bt.get("valid"))
.and_then(|valid| valid.get(0))
.and_then(|valid_from_json| {
serde_json::from_value::<Hlc>(valid_from_json.clone()).ok()
})
.map(|valid_from| valid_from < hlc)
.unwrap_or(false)
}
fn build_soft_delete_op(m: &ForgetMatch, clock: &HlcClock) -> Result<WriteOp, LunarisError> {
let now = clock.tick();
let mut bt = m.bt;
bt.invalidate_sys(now);
let mut json: serde_json::Value = serde_json::from_slice(&m.payload).map_err(|e| {
LunarisError::Storage(StorageError::Backend(format!("forget payload parse: {e}")))
})?;
if let Some(bt_obj) = json.get_mut("bt")
&& let Some(sys_arr) = bt_obj.get_mut("sys")
&& let Some(arr) = sys_arr.as_array_mut()
&& arr.len() == 2
{
arr[1] = serde_json::to_value(now).unwrap_or(serde_json::Value::Null);
}
let value = serde_json::to_vec(&json).map_err(|e| {
LunarisError::Storage(StorageError::Backend(format!("forget payload serialize: {e}")))
})?;
Ok(WriteOp::KvPut { key: m.key.clone(), value })
}
fn classify_indices(target: &ForgetTarget) -> Vec<IndexKind> {
match target {
ForgetTarget::Id(_) | ForgetTarget::Scope(_) | ForgetTarget::Before(_) => {
vec![IndexKind::Kv, IndexKind::Vector, IndexKind::Graph]
}
}
}
pub(crate) async fn forget_scoped(
storage: &std::sync::Arc<dyn StoragePort>,
clock: &HlcClock,
scope: &lunaris_core::Scope,
request: ForgetRequest,
) -> Result<ForgetReceipt, LunarisError> {
let target_kind = match &request.target {
ForgetTarget::Id(_) => "id",
ForgetTarget::Scope(_) => "scope",
ForgetTarget::Before(_) => "before",
};
let span = tracing::info_span!(
"lunaris.forget",
correlation_id = tracing::field::Empty,
target = %target_kind,
scope = %scope.as_str(),
hard = %request.options.hard,
dry_run = %request.options.dry_run,
);
async move {
if request.options.hard && request.options.confirmation_token.is_none() {
return Err(LunarisError::Validate(ValidateError::ConfirmationRequired(
"hard-delete requires confirmation token from prior dry_run".into(),
)));
}
let mut matches =
scan_matches_scoped(storage.as_ref(), scope, &request.target, clock).await?;
drop_already_closed(&mut matches, request.options.hard);
let matched = episode_match_count(&matches, scope);
if request.options.dry_run {
let receipt = ForgetReceipt {
target: request.target.clone(),
indices_affected: classify_indices(&request.target),
matched,
rows_written: 0,
rows_deleted: 0,
audit_lsn: Lsn { wall_ms: 0, counter: 0 },
preview: true,
};
let audit_offset =
publish_audit_event(storage, scope, AuditEvent::Forget((&receipt).into()))
.await
.unwrap_or(0);
return Ok(ForgetReceipt {
audit_lsn: Lsn { wall_ms: audit_offset, counter: 0 },
..receipt
});
}
let ops: Vec<WriteOp> = if request.options.hard {
matches.iter().map(|m| WriteOp::KvDelete { key: m.key.clone() }).collect()
} else {
matches.iter().map(|m| build_soft_delete_op(m, clock)).collect::<Result<Vec<_>, _>>()?
};
if !ops.is_empty() {
let _lsn = storage.atomic_write(scope, &ops).await.map_err(LunarisError::Storage)?;
}
let receipt = ForgetReceipt {
target: request.target.clone(),
indices_affected: classify_indices(&request.target),
matched,
rows_written: if request.options.hard { 0 } else { matches.len() as u64 },
rows_deleted: if request.options.hard { matches.len() as u64 } else { 0 },
audit_lsn: Lsn { wall_ms: 0, counter: 0 },
preview: false,
};
let audit_offset =
publish_audit_event(storage, scope, AuditEvent::Forget((&receipt).into()))
.await
.unwrap_or(0);
Ok(ForgetReceipt { audit_lsn: Lsn { wall_ms: audit_offset, counter: 0 }, ..receipt })
}
.instrument(span)
.await
}
async fn scan_matches_scoped(
storage: &dyn StoragePort,
scope: &lunaris_core::Scope,
target: &ForgetTarget,
clock: &HlcClock,
) -> Result<Vec<ForgetMatch>, LunarisError> {
let mut out = scan_episode_matches_scoped(storage, scope, target, clock).await?;
let episode_ids = matched_episode_ids(&out);
if !episode_ids.is_empty() {
out.extend(scan_chunk_matches_scoped(storage, scope, &episode_ids, clock).await?);
}
Ok(out)
}
async fn scan_episode_matches_scoped(
storage: &dyn StoragePort,
scope: &lunaris_core::Scope,
target: &ForgetTarget,
clock: &HlcClock,
) -> Result<Vec<ForgetMatch>, LunarisError> {
use futures::stream::StreamExt;
if let ForgetTarget::Id(ulid) = target {
let key = lunaris_core::keyspace::episode_key(scope, *ulid);
let now = clock.tick();
let row = storage.read_as_of(scope, &key, now).await.map_err(LunarisError::Storage)?;
return Ok(match row {
Some(r) => vec![ForgetMatch { key, payload: r.value.to_vec(), bt: r.bt }],
None => Vec::new(),
});
}
let prefix = format!("{}episode:", lunaris_core::keyspace::scope_prefix(scope)).into_bytes();
let mut stream =
storage.scan_range(scope, &prefix, None).await.map_err(LunarisError::Storage)?;
let mut out = Vec::new();
while let Some(item) = stream.next().await {
let (k, v) = item.map_err(LunarisError::Storage)?;
if matches_target(&k, &v, target) {
let now = clock.tick();
let row_opt =
storage.read_as_of(scope, &k, now).await.map_err(LunarisError::Storage)?;
let bt = row_opt
.as_ref()
.map(|r| r.bt)
.unwrap_or_else(|| BiTemporal { valid: (Hlc::ZERO, None), sys: (Hlc::ZERO, None) });
out.push(ForgetMatch { key: k.to_vec(), payload: v.to_vec(), bt });
}
}
Ok(out)
}
fn is_sys_closed(payload: &[u8]) -> bool {
serde_json::from_slice::<serde_json::Value>(payload)
.ok()
.and_then(|v| {
v.get("bt")
.and_then(|bt| bt.get("sys"))
.and_then(|sys| sys.as_array())
.and_then(|a| a.get(1))
.map(|to| !to.is_null())
})
.unwrap_or(false)
}
fn drop_already_closed(matches: &mut Vec<ForgetMatch>, hard: bool) {
if !hard {
matches.retain(|m| !is_sys_closed(&m.payload));
}
}
fn episode_match_count(matches: &[ForgetMatch], scope: &lunaris_core::Scope) -> u64 {
let prefix = lunaris_core::keyspace::episode_prefix(scope);
matches.iter().filter(|m| m.key.starts_with(&prefix)).count() as u64
}
fn matched_episode_ids(matches: &[ForgetMatch]) -> std::collections::HashSet<Ulid> {
matches
.iter()
.filter_map(|m| serde_json::from_slice::<serde_json::Value>(&m.payload).ok())
.filter_map(|v| v.get("id").and_then(|id| id.as_str()).and_then(|s| s.parse::<Ulid>().ok()))
.collect()
}
async fn scan_chunk_matches_scoped(
storage: &dyn StoragePort,
scope: &lunaris_core::Scope,
episode_ids: &std::collections::HashSet<Ulid>,
clock: &HlcClock,
) -> Result<Vec<ForgetMatch>, LunarisError> {
use futures::stream::StreamExt;
let prefix = lunaris_core::keyspace::chunk_prefix(scope);
let mut stream =
storage.scan_range(scope, &prefix, None).await.map_err(LunarisError::Storage)?;
let mut out = Vec::new();
while let Some(item) = stream.next().await {
let (k, v) = item.map_err(LunarisError::Storage)?;
let Ok(chunk) = serde_json::from_slice::<lunaris_core::Chunk>(&v) else {
continue; };
if !episode_ids.contains(&chunk.episode_id) {
continue;
}
let now = clock.tick();
let row_opt = storage.read_as_of(scope, &k, now).await.map_err(LunarisError::Storage)?;
let bt = row_opt
.as_ref()
.map(|r| r.bt)
.unwrap_or_else(|| BiTemporal { valid: (Hlc::ZERO, None), sys: (Hlc::ZERO, None) });
out.push(ForgetMatch { key: k.to_vec(), payload: v.to_vec(), bt });
}
Ok(out)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn dry_run_builder_sets_options() {
let req = ForgetTarget::Id(Ulid::new()).dry_run();
assert!(req.options.dry_run);
assert!(!req.options.hard);
assert!(req.options.confirmation_token.is_none());
}
#[test]
fn hard_builder_sets_options() {
let req = ForgetTarget::Id(Ulid::new()).hard();
assert!(req.options.hard);
assert!(!req.options.dry_run);
assert!(req.options.confirmation_token.is_none());
}
#[test]
fn with_token_attaches_confirmation() {
let token = ForgetConfirmation { for_audit_lsn: Lsn { wall_ms: 42, counter: 0 } };
let req = ForgetTarget::Id(Ulid::new()).hard().with_token(token.clone());
assert_eq!(req.options.confirmation_token, Some(token));
}
#[test]
fn classify_indices_returns_three_indices() {
let v = classify_indices(&ForgetTarget::Id(Ulid::new()));
assert!(v.contains(&IndexKind::Kv));
assert!(v.contains(&IndexKind::Vector));
assert!(v.contains(&IndexKind::Graph));
}
#[test]
fn match_scope_by_source_prefix() {
let payload = serde_json::to_vec(&serde_json::json!({
"source": "helios:fs/session-42/foo.md",
}))
.unwrap();
assert!(match_scope(&payload, &ScopeSpec::BySource("helios:fs/session-42/".into())));
assert!(!match_scope(&payload, &ScopeSpec::BySource("helios:fs/session-99/".into())));
}
#[test]
fn match_scope_by_metadata_exact() {
let payload = serde_json::to_vec(&serde_json::json!({
"metadata": { "tenant": "acme" },
}))
.unwrap();
assert!(match_scope(&payload, &ScopeSpec::ByMetadata("tenant".into(), "acme".into())));
assert!(!match_scope(&payload, &ScopeSpec::ByMetadata("tenant".into(), "other".into())));
}
#[test]
fn match_before_uses_typed_hlc_compare() {
let earlier = Hlc { wall_ms: 100, counter: 0, node_id: 0 };
let later = Hlc { wall_ms: 200, counter: 0, node_id: 0 };
let cutoff = Hlc { wall_ms: 150, counter: 0, node_id: 0 };
let payload_earlier = serde_json::to_vec(&serde_json::json!({
"bt": { "valid": [earlier, null], "sys": [earlier, null] }
}))
.unwrap();
let payload_later = serde_json::to_vec(&serde_json::json!({
"bt": { "valid": [later, null], "sys": [later, null] }
}))
.unwrap();
assert!(match_before(&payload_earlier, cutoff), "earlier < cutoff");
assert!(!match_before(&payload_later, cutoff), "later >= cutoff");
}
#[test]
fn forget_target_into_request_defaults_to_soft() {
let req: ForgetRequest = ForgetTarget::Id(Ulid::new()).into_request();
assert!(!req.options.hard);
assert!(!req.options.dry_run);
}
#[test]
fn from_impl_allows_bare_target_in_forget_arg() {
let req: ForgetRequest = ForgetTarget::Id(Ulid::new()).into();
assert!(!req.options.hard);
}
#[test]
fn build_soft_delete_op_patches_bt_sys_to() {
let t0 = Hlc { wall_ms: 1, counter: 0, node_id: 0 };
let payload = serde_json::to_vec(&serde_json::json!({
"id": Ulid::new().to_string(),
"source": "test:src",
"bt": { "valid": [t0, null], "sys": [t0, null] },
}))
.unwrap();
let bt = BiTemporal { valid: (t0, None), sys: (t0, None) };
let m = ForgetMatch { key: b"episode:test".to_vec(), payload, bt };
let clock = HlcClock::new(0);
let op = build_soft_delete_op(&m, clock.as_ref()).expect("build_soft_delete_op");
match op {
WriteOp::KvPut { key, value } => {
assert_eq!(key, b"episode:test");
let json: serde_json::Value = serde_json::from_slice(&value).unwrap();
let sys_to = json
.get("bt")
.and_then(|bt| bt.get("sys"))
.and_then(|sys| sys.get(1))
.cloned()
.unwrap_or(serde_json::Value::Null);
assert!(!sys_to.is_null(), "sys[1] (sys_to) MUST be patched non-null");
}
other => panic!("expected KvPut, got {other:?}"),
}
}
}