khive-storage 0.9.0

Storage capability contracts and backend-neutral request context.
Documentation
//! Entity storage capability — graph node CRUD.

use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use uuid::Uuid;

use crate::attachment::Attachment;
use crate::types::{
    BatchWriteSummary, DeleteMode, Page, PageRequest, SeekCursor, SeekPage, StorageResult,
};

/// Storage-level entity record. Flat SQL-friendly representation.
/// Maps to the `entities` substrate table.
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Entity {
    pub id: Uuid,
    pub namespace: String,
    pub kind: String,
    /// Pack-governed subtype token. Maps to `entities.entity_type` column.
    pub entity_type: Option<String>,
    pub name: String,
    pub description: Option<String>,
    pub properties: Option<Value>,
    pub tags: Vec<String>,
    pub created_at: i64,
    pub updated_at: i64,
    /// Persisted row revision: inserts start at one; each committed update advances once.
    /// Input to a snapshot replacement is the revision read before normalization.
    #[serde(default = "initial_entity_version")]
    pub version: i64,
    pub deleted_at: Option<i64>,
    /// When this entity was tombstoned by a merge, the `into` entity's ID.
    pub merged_into: Option<Uuid>,
    /// Opaque event ID for the merge that tombstoned this entity.
    pub merge_event_id: Option<Uuid>,
    /// Read-only compatibility projection of attachment role `"content"`.
    ///
    /// Entity writes ignore this field. Callers publish content through the
    /// attachment substrate; reads populate it so existing response payloads
    /// keep their `content_ref` field during the coordinated cutover.
    pub content_ref: Option<String>,
}

fn initial_entity_version() -> i64 {
    1
}

impl Entity {
    /// Create a new entity with a generated UUID and current timestamp.
    pub fn new(
        namespace: impl Into<String>,
        kind: impl Into<String>,
        name: impl Into<String>,
    ) -> Self {
        let now = chrono::Utc::now().timestamp_micros();
        Self {
            id: Uuid::new_v4(),
            namespace: namespace.into(),
            kind: kind.into(),
            entity_type: None,
            name: name.into(),
            description: None,
            properties: None,
            tags: Vec::new(),
            created_at: now,
            updated_at: now,
            version: 1,
            deleted_at: None,
            merged_into: None,
            merge_event_id: None,
            content_ref: None,
        }
    }

    /// Set the pack-governed entity subtype token.
    pub fn with_entity_type(mut self, t: Option<impl Into<String>>) -> Self {
        self.entity_type = t.map(Into::into);
        self
    }

    /// Set the entity description.
    pub fn with_description(mut self, d: impl Into<String>) -> Self {
        self.description = Some(d.into());
        self
    }

    /// Set the entity properties JSON blob.
    pub fn with_properties(mut self, p: Value) -> Self {
        self.properties = Some(p);
        self
    }

    /// Set the entity tags.
    pub fn with_tags(mut self, t: Vec<String>) -> Self {
        self.tags = t;
        self
    }
}

/// Entity filter for query operations.
#[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct EntityFilter {
    pub ids: Vec<Uuid>,
    pub kinds: Vec<String>,
    /// Filter by exact `entity_type` value. Multiple values are ORed.
    pub entity_types: Vec<String>,
    /// Kind-qualified accepted subtype spellings. Groups are ORed; each group
    /// requires its kind and one of its values. An empty value group matches nothing.
    #[serde(default)]
    pub entity_types_by_kind: std::collections::BTreeMap<String, Vec<String>>,
    /// For entity listing, fall back to a string `properties.type` only when
    /// `entity_type` is null. Does not change the returned entity or apply when
    /// both type filters are empty. Other query callers retain exact-column filtering.
    #[serde(default)]
    pub legacy_entity_type_fallback: bool,
    pub name_prefix: Option<String>,
    /// Deterministic, case-sensitive equality on `entities.name` (binary
    /// comparison — SQLite's default collation for `=` on a `TEXT` column
    /// without an explicit `COLLATE NOCASE`). Distinct from `name_prefix`:
    /// that stage's `LIKE` is inherently prefix-shaped and, with SQLite's
    /// default `NOCASE`-free `LIKE` on ASCII, still ranks a page by
    /// `created_at DESC` — a match that is exact but not the newest can be
    /// paged out. `name_exact` skips paging risk entirely by filtering to
    /// only rows that equal `name` at the SQL layer.
    pub name_exact: Option<String>,
    pub tags_any: Vec<String>,
    /// When non-empty, restricts results to any of these namespaces using
    /// `namespace IN (...)`. Takes precedence over the `namespace` string
    /// parameter passed to `query_entities` / `count_entities`. When empty the
    /// caller-supplied `namespace` parameter is used (single-namespace path,
    /// backward-compatible default).
    #[serde(default)]
    pub namespaces: Vec<String>,
    /// ASCII-case-insensitive batched exact-name match (ADR-104 Stage C).
    /// Compares a caller-bounded set of raw and ASCII-lowercased candidate
    /// strings to `LOWER(name)`. Cased non-ASCII characters require exact form.
    /// Distinct from single-value, case-sensitive `name_exact`. Results contain
    /// at most one representative row per folded candidate before page limits
    /// and offsets are applied.
    /// Implementations may omit the page total to keep this lookup page-limited
    /// instead of issuing a separate count.
    #[serde(default)]
    pub names_ci: Vec<String>,
}

/// Entity CRUD operations over the entities substrate table.
#[async_trait]
pub trait EntityStore: Send + Sync + 'static {
    /// Insert at version one or update the current row and advance its version once.
    /// Incoming version values do not override the stored counter. Raw SQL replacement
    /// of an existing entity is forbidden: use a conflict UPDATE or the typed store.
    async fn upsert_entity(&self, entity: Entity) -> StorageResult<()>;
    /// Insert an entity only when no row with its id or another conflicting
    /// key exists. Returns `true` when this call inserted the row and `false`
    /// when an existing row won the race. The existing row is never updated.
    ///
    /// The default returns `Unsupported` rather than falling back to
    /// [`EntityStore::upsert_entity`], because an upsert would overwrite the
    /// winning row and violate this method's conditional-insert contract.
    async fn insert_entity_if_absent(&self, _entity: Entity) -> StorageResult<bool> {
        Err(crate::StorageError::Unsupported {
            capability: crate::StorageCapability::Entities,
            operation: "insert_entity_if_absent".into(),
            message: "this backend does not implement conditional entity insert".into(),
        })
    }
    /// Atomically insert/update an entity and all supplied attachment roles.
    async fn upsert_entity_with_attachments(
        &self,
        _entity: Entity,
        _attachments: Vec<Attachment>,
    ) -> StorageResult<()> {
        Err(crate::StorageError::Unsupported {
            capability: crate::StorageCapability::Attachments,
            operation: "upsert_entity_with_attachments".into(),
            message: "this backend does not implement atomic entity attachment publication".into(),
        })
    }
    /// Insert or update a batch of entities.
    async fn upsert_entities(&self, entities: Vec<Entity>) -> StorageResult<BatchWriteSummary>;
    /// Replace an entity only when the persisted row still matches the
    /// caller's read snapshot.
    ///
    /// `entity.version` is the persisted revision read with the snapshot;
    /// a successful replacement increments the stored version once.
    /// `expected_updated_at` is an additional snapshot timestamp guard and
    /// `expected_deleted_at` closes the soft-delete race. The replacement
    /// entity's `updated_at` must be strictly greater than that persisted
    /// revision. Returns `false` when the row disappeared, changed, or was
    /// supplied a non-advancing replacement revision. This is the full-entity
    /// compare-and-swap seam used when a caller derives coupled fields from
    /// that snapshot before persistence — mirrors
    /// [`crate::NoteStore::replace_note_if_unchanged`]. The default returns
    /// `Unsupported` rather than falling back to an unguarded upsert and
    /// reintroducing the stale-snapshot race.
    async fn replace_entity_if_unchanged(
        &self,
        _entity: Entity,
        _expected_updated_at: i64,
        _expected_deleted_at: Option<i64>,
    ) -> StorageResult<bool> {
        Err(crate::StorageError::Unsupported {
            capability: crate::StorageCapability::Entities,
            operation: "replace_entity_if_unchanged".into(),
            message: "this backend does not implement guarded entity replacement".into(),
        })
    }
    /// Fetch an entity by UUID, returning `None` if absent.
    async fn get_entity(&self, id: Uuid) -> StorageResult<Option<Entity>>;
    /// Delete an entity by UUID using the specified delete mode.
    async fn delete_entity(&self, id: Uuid, mode: DeleteMode) -> StorageResult<bool>;
    /// Query entities by namespace with filter and pagination.
    async fn query_entities(
        &self,
        namespace: &str,
        filter: EntityFilter,
        page: PageRequest,
    ) -> StorageResult<Page<Entity>>;
    /// Resolve an entity id to its immutable insertion sequence.
    async fn entity_sequence(&self, _id: Uuid) -> StorageResult<Option<i64>> {
        Err(crate::StorageError::Unsupported {
            capability: crate::StorageCapability::Entities,
            operation: "entity_sequence".into(),
            message: "this backend does not implement entity insertion sequences".into(),
        })
    }
    /// Query an immutable insertion-sequence keyset page.
    ///
    /// Backends that do not implement seek pagination may retain the default
    /// unsupported result; callers must not silently fall back to offset
    /// paging because that would weaken the no-gap/no-duplicate contract.
    async fn query_entities_after(
        &self,
        _namespace: &str,
        _filter: EntityFilter,
        _after: Option<SeekCursor>,
        _limit: u32,
    ) -> StorageResult<SeekPage<Entity>> {
        Err(crate::StorageError::Unsupported {
            capability: crate::StorageCapability::Entities,
            operation: "query_entities_after".into(),
            message: "this backend does not implement entity seek pagination".into(),
        })
    }
    /// Count entities in a namespace matching the given filter.
    async fn count_entities(&self, namespace: &str, filter: EntityFilter) -> StorageResult<u64>;
    /// Fetch an entity by UUID regardless of soft-deletion state.
    ///
    /// Returns the entity row even when `deleted_at` is set. Callers use this
    /// to distinguish "soft-deleted" from "never existed".
    async fn get_entity_including_deleted(&self, id: Uuid) -> StorageResult<Option<Entity>>;
}