Skip to main content

KhiveRuntime

Struct KhiveRuntime 

Source
pub struct KhiveRuntime { /* private fields */ }
Expand description

for each storage capability, plus a lazily-loaded embedder.

Implementations§

Source§

impl KhiveRuntime

Source

pub async fn create_sender_transport( &self, token: &NamespaceToken, envelope: SenderEnvelope, ) -> RuntimeResult<SenderRecord>

Persist immutable bytes before first submission. An exact retry is a no-op; a different envelope with the same identity is refused. The namespace is caller attribution and is always taken from the token. Sender assurance is claimed until authenticated verification is available.

Source

pub async fn reencrypt_sender_transport_after_confirmed_key_change( &self, token: &NamespaceToken, envelope: SenderEnvelope, ) -> RuntimeResult<SenderRecord>

Re-encrypt only after the caller confirms the recipient key change. The old record must be on a key-change hold, and the new epoch must increase. Cryptographic construction and directory confirmation belong to the caller. Caller-supplied assurance is normalized to claimed as on initial creation.

Source

pub async fn sender_transport( &self, key: EnvelopeKey, ) -> RuntimeResult<Option<SenderRecord>>

Read by durable identity. Absence means unknown; unknown is not persisted.

Source

pub async fn pending_sender_transports( &self, token: &NamespaceToken, kind: &str, slug: &str, now: i64, limit: u32, ) -> RuntimeResult<Vec<SenderRecord>>

Due unheld submissions for the exact configured kind and slug. Times use Unix microseconds; choosing the backoff and jitter is the caller’s job.

Source

pub async fn record_sender_transport_admission( &self, key: EnvelopeKey, admitted_at: DateTime<Utc>, ) -> RuntimeResult<()>

Record service admission time and schedule resubmission 600 seconds later. The timestamp is stored as Unix microseconds and survives restarts.

Source

pub async fn record_sender_transport_failure( &self, key: EnvelopeKey, class: FailureClass, next_retry_at: Option<i64>, ) -> RuntimeResult<()>

Record a failed attempt. Authentication errors remain pending and do not advance the attempt count. This never modifies the envelope bytes.

Source

pub async fn hold_sender_transport( &self, key: EnvelopeKey, reason: Option<HoldReason>, ) -> RuntimeResult<()>

Set/release credit or policy holds, or set a key-change hold. A policy hold carries the evaluated mode and revision. Key-change holds can only be superseded through explicit confirmed re-encryption.

Source

pub async fn accept_verified_recipient_receipt( &self, key: EnvelopeKey, state: TransportState, receipt: VerifiedRecipientReceipt, ) -> RuntimeResult<()>

Accept an already signature-verified recipient receipt, including after a local failure. Every known binding field and target disposition must match durable data. This method performs no signature verification itself.

Source§

impl KhiveRuntime

Source

pub async fn update_entity( &self, token: &NamespaceToken, id: Uuid, patch: EntityPatch, ) -> RuntimeResult<Entity>

Source

pub async fn update_entity_with_embedding_report( &self, token: &NamespaceToken, id: Uuid, patch: EntityPatch, ) -> RuntimeResult<(Entity, EmbeddingTruncationReport)>

Source

pub async fn update_entity_with_expected_version_and_embedding_report( &self, token: &NamespaceToken, id: Uuid, patch: EntityPatch, expected_version: Option<i64>, ) -> RuntimeResult<(Entity, EmbeddingTruncationReport)>

Entity update with an optional caller revision, checked inside the writer transaction.

Source

pub async fn update_entity_if_unchanged( &self, token: &NamespaceToken, expected: &Entity, patch: EntityPatch, remove_properties: &[&str], ) -> RuntimeResult<Entity>

Apply an admin patch only if the entity still matches the full read snapshot. Property removals apply after the normal merge and preserve all other keys. Missing keys alone are a no-op; reserved runtime-owned keys cannot be removed. A changed, deleted, or missing entity returns a conflict without writing.

Source

pub async fn merge_entity( &self, token: &NamespaceToken, into_id: Uuid, from_id: Uuid, strategy: EntityDedupMergePolicy, content_strategy: ContentMergeStrategy, dry_run: bool, ) -> RuntimeResult<MergeSummary>

Merge from_id into into_id.

All edges incident to from_id are rewired to into_id. Self-loops that would result from the rewire are dropped. Properties and tags are merged per strategy. from_id is tombstoned with merge provenance and removed from indexes. Returns a summary.

If dry_run is true, computes and returns the planned summary without mutating any rows.

Atomic: all SQL (entity reads/writes, edge rewires, FTS updates, vec-index delete, merge event with destructive edge preimages) runs on one pool connection inside one BEGIN IMMEDIATE transaction via merge_entity_sql. If embedding vectors are configured, the vector re-insert for into_id is performed after the transaction (requires async embedding computation).

Source

pub async fn merge_entity_with_reason( &self, token: &NamespaceToken, into_id: Uuid, from_id: Uuid, strategy: EntityDedupMergePolicy, content_strategy: ContentMergeStrategy, dry_run: bool, reason: Option<String>, ) -> RuntimeResult<MergeSummary>

Merge from_id into into_id and include an optional reason in the audit event.

Source

pub async fn merge_entity_with_reason_and_force( &self, token: &NamespaceToken, into_id: Uuid, from_id: Uuid, strategy: EntityDedupMergePolicy, content_strategy: ContentMergeStrategy, dry_run: bool, reason: Option<String>, force: bool, ) -> RuntimeResult<MergeSummary>

Merge two entities with an explicit override for the entity safety floor.

Non-forced calls enforce entity kind, name similarity, and project compatibility against the rows reread inside the merge transaction. Legacy merge methods retain their historical same-kind-only policy. A non-dry-run override is recorded as force: true in the merge event.

Source

pub async fn update_note( &self, token: &NamespaceToken, id: Uuid, patch: NotePatch, ) -> RuntimeResult<Note>

Patch-style note update.

Source

pub async fn update_note_with_embedding_report( &self, token: &NamespaceToken, id: Uuid, patch: NotePatch, ) -> RuntimeResult<(Note, EmbeddingTruncationReport)>

Source

pub async fn update_note_from_snapshot_with_embedding_report( &self, token: &NamespaceToken, snapshot: Note, patch: NotePatch, ) -> RuntimeResult<(Note, EmbeddingTruncationReport)>

Patch and persist one note from a caller-owned read snapshot.

This is the canonical seam for kind hooks that normalize coupled fields from the current note. The same snapshot feeds normalization, patch application, and the compare-and-swap write; a concurrent note change therefore refuses the write instead of persisting derivations computed from stale state.

Source

pub async fn update_note_from_snapshot_with_kind_effects( &self, token: &NamespaceToken, snapshot: Note, args: &Value, policy: NoteUpdatePolicy, registry: &VerbRegistry, ) -> RuntimeResult<(Note, EmbeddingTruncationReport)>

Commit a normalized and validated kind-owned update, including its typed graph companions. Callers must first run prepare_note_update_policy against this exact snapshot and pass the policy it returned; the shared atomic prepare seam checks all patch fields before asking the kind hook to derive any graph effects.

Source

pub async fn claim_outbound_message_external_id( &self, token: &NamespaceToken, id: Uuid, external_id: String, ) -> RuntimeResult<Note>

Claim external_id on an outbound message note through the ADR-124-sanctioned store-level owner path, bypassing the caller-facing owner-established-property refusal in Self::update_note (and its crate-internal prepare path). This is deliberately NOT exposed through any registered verb (ADR-124’s stated bound): it is reachable only from pack/runtime code that owns outbox bookkeeping for the message note kind.

Refuses (returns Err, never writes) unless the live row is a message note, properties.direction == "outbound", and properties.external_id is currently absent or empty. The claim is committed against that exact snapshot and advances its timestamp so competing claims and delivery-outcome CAS writes cannot overwrite it.

Source

pub async fn list_undelivered_outbound_messages( &self, token: &NamespaceToken, to_prefix: Option<&str>, limit: u32, ) -> RuntimeResult<Vec<Note>>

Non-wire outbox scan for the channel delivery loops.

Fetches live message notes matching the SQL-side pending predicate newest-first (created_at DESC, id ASC), bounded by an internal scan cap, and returns those that are still due, capped at limit. Direction, delivered_at, terminal delivery state and the optional to_actor channel prefix are all filtered by SQLite; only a valid next_attempt_at remains a Rust check. Pending means delivered_at is absent or null, properties.delivery carries no terminal state ("delivered" / "failed"), and a valid next_attempt_at is absent or due (ADR-122 §1). Malformed legacy deadlines fail open so a bad property cannot strand mail forever.

The channel prefix has to be in the statement, not applied to the fetched page: every actor-to-actor outbound row matches the pending predicate forever (nothing marks those delivered), so that population outgrows any scan cap and a page-then-filter scan never reaches a channel’s rows once enough other rows sort ahead of them. The prefix renders as an index range served by idx_comm_message_outbound_recipient, and the newest-first order means a future predicate miss still surfaces new rows first. This lives on the runtime rather than going through the wire registry for the same reason as Self::claim_outbound_message_external_id: the delivery loop must scan the backend that actually holds comm’s notes, and under a [packs.comm] backend assignment that is not the backend serving the generic kg verbs.

Source

pub async fn mark_outbound_message_transient_failure( &self, token: &NamespaceToken, id: Uuid, attempted_at: DateTime<Utc>, last_error: String, base_delay: Duration, max_delay: Duration, ) -> RuntimeResult<Note>

Record a transient transport failure while leaving the message pending. The retry deadline is derived from the incremented persisted attempt count, using base_delay * 2^(attempt - 1) capped at max_delay.

Source

pub async fn mark_outbound_message_claim_transient_failure( &self, token: &NamespaceToken, id: Uuid, attempted_at: DateTime<Utc>, last_error: String, base_delay: Duration, max_delay: Duration, ) -> RuntimeResult<Note>

Schedule an external-id claim retry only for a still-unclaimed outbound snapshot. Existing claims and terminal outcomes are unchanged.

Source

pub async fn mark_outbound_message_delivered( &self, token: &NamespaceToken, id: Uuid, delivered_at: String, transport_message_id: Option<String>, ) -> RuntimeResult<Note>

Mark an outbound message note delivered by merging the ADR-122 §1 terminal-outcome properties (delivery = "delivered", delivered_at, and transport_message_id when the transport minted one), and clearing delivery_attempts / next_attempt_at. The compare-and-swap protects unrelated properties from a concurrent full-row overwrite; delivered_at remains deliberately caller-patchable, pinned by generic_update_can_still_patch_delivered_at_on_message_note. Refuses (InvalidInput) unless id names a live outbound message note. Non-wire companion to Self::list_undelivered_outbound_messages so the delivery loop writes the backend that holds the note.

Source

pub async fn mark_outbound_message_failed( &self, token: &NamespaceToken, id: Uuid, failed_at: String, last_error: String, ) -> RuntimeResult<Note>

Record a permanent delivery failure on an outbound message note: delivery = "failed", failed_at, last_error (ADR-122 §2 — an allowlist rejection must be recorded, not skipped, or the row stays pending forever while the caller saw ok: true). Any retry counter and deadline are cleared because the outcome is terminal. Refuses (InvalidInput) unless id names a live outbound message note.

Source

pub async fn mark_outbound_message_claim_failed( &self, token: &NamespaceToken, id: Uuid, failed_at: String, last_error: String, ) -> RuntimeResult<Note>

Park a deterministic external-id claim refusal only while the exact current outbound snapshot remains unclaimed. An existing claim or a terminal delivery outcome is returned unchanged. Non-wire owner API: a second worker must not turn another worker’s successful claim into a permanent delivery failure.

Source

pub async fn merge_note( &self, token: &NamespaceToken, into_id: Uuid, from_id: Uuid, strategy: EntityDedupMergePolicy, content_strategy: ContentMergeStrategy, dry_run: bool, ) -> RuntimeResult<MergeSummary>

Merge from_id note into into_id note.

Both notes must exist in the namespace and have the same kind. Content is merged per content_strategy. Properties are merged per strategy. from_id is tombstoned (status=‘deleted’, deleted_at set). Returns a summary.

If dry_run is true, computes and returns the planned summary without mutating any rows, edges, or indexes. The NoteMerged event, including destructive edge preimages, commits in the same SQL transaction as the note and edge changes.

Source

pub async fn merge_note_with_reason( &self, token: &NamespaceToken, into_id: Uuid, from_id: Uuid, strategy: EntityDedupMergePolicy, content_strategy: ContentMergeStrategy, dry_run: bool, reason: Option<String>, ) -> RuntimeResult<MergeSummary>

Merge from_id note into into_id note and include an optional audit reason.

Source§

impl KhiveRuntime

Source

pub async fn hybrid_search_with_strategy( &self, token: &NamespaceToken, query_text: &str, query_vector: Option<Vec<f32>>, strategy: FusionStrategy, limit: u32, ) -> RuntimeResult<Vec<SearchHit>>

Hybrid search with a caller-supplied fusion strategy.

FusionStrategy::Custom { name, .. } is resolved against this runtime’s registered executors (see register_fusion_strategy); an unregistered name fails closed with RuntimeError::UnknownFusionStrategy.

Source§

impl KhiveRuntime

Source

pub async fn bfs_traverse( &self, token: &NamespaceToken, start: Uuid, options: TraversalOptions, ) -> RuntimeResult<Vec<PathNode>>

BFS traversal from start, returning nodes in level order.

The first element is always the start node (via_edge = None, depth = 0). Nodes already visited are skipped so the result set is deduplicated.

DB round-trips are O(max_depth): one batch_neighbors call and one get_edges call per BFS level, rather than one call per node/edge.

Source

pub async fn shortest_path( &self, token: &NamespaceToken, from: Uuid, to: Uuid, max_depth: usize, ) -> RuntimeResult<Option<Vec<PathNode>>>

Bidirectional BFS shortest path from from to to.

Returns Some(path) where path[0] is from and path.last() is to, or None if no path exists within max_depth hops. For from == to returns Some with a single-node path immediately.

DB round-trips are O(max_depth): one batch_neighbors per frontier expansion level, plus one get_edges call during path reconstruction.

Source§

impl KhiveRuntime

Source

pub fn authorize_mailbox_view( &self, token: &NamespaceToken, verb: &str, selector: Option<&str>, args: &Value, ) -> RuntimeResult<MailboxView>

Authorize a read-only mailbox view without changing token identity.

The original arguments are retained as gate input. Direct handler calls consult both the existing policy and the separate mailbox capability; a permissive or unconfigured gate never grants cross-actor reads.

Source§

impl KhiveRuntime

Source

pub async fn list_notes_filtered( &self, token: &NamespaceToken, filter: NoteFilter, limit: u32, offset: u32, ) -> RuntimeResult<Vec<Note>>

Source

pub async fn list_notes_filtered_after( &self, token: &NamespaceToken, filter: NoteFilter, after: Option<Uuid>, limit: u32, ) -> RuntimeResult<(Vec<Note>, Option<Uuid>)>

Source§

impl KhiveRuntime

Source

pub async fn get_note_by_key( &self, token: &NamespaceToken, key: &str, kind: Option<&str>, after_key: bool, ) -> RuntimeResult<Note>

Source

pub async fn create_note_with_options( &self, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, embedding_content: Option<&str>, salience: Option<f64>, decay_factor: Option<f64>, properties: Option<Value>, annotates: Vec<Uuid>, embedding_model: Option<&str>, options: NoteWriteOptions, ) -> RuntimeResult<(Note, EmbeddingTruncationReport)>

Source§

impl KhiveRuntime

Source

pub async fn create_entity( &self, token: &NamespaceToken, kind: &str, entity_type: Option<&str>, name: &str, description: Option<&str>, properties: Option<Value>, tags: Vec<String>, ) -> RuntimeResult<Entity>

Create and persist a new entity.

Indexing failures trigger compensation across the entity row, FTS document, and any vector models touched by this call. If compensation also fails, the returned structured internal error identifies possible partial persistence, includes both failure classes, and carries the entity ID as a reconciliation handle.

Source

pub async fn create_entity_with_attachments( &self, token: &NamespaceToken, kind: &str, entity_type: Option<&str>, name: &str, description: Option<&str>, properties: Option<Value>, tags: Vec<String>, attachments: Vec<NewAttachment>, ) -> RuntimeResult<Entity>

Create an entity with role-keyed bytes already published to BlobStore.

Every NewAttachment carries a typed content reference, so malformed references cannot enter through this consumer seam. Blob existence is checked before the database write. The entity row and all attachment rows then commit in one storage transaction; the FTS/vector compensation path hard-deletes the entity and its attachments together if a later indexing step fails. Published bytes remain recoverable by the BlobStore grace-period orphan policy when any post-publication step fails.

Source

pub async fn create_entity_with_embedding_report( &self, token: &NamespaceToken, kind: &str, entity_type: Option<&str>, name: &str, description: Option<&str>, properties: Option<Value>, tags: Vec<String>, ) -> RuntimeResult<(Entity, EmbeddingTruncationReport)>

Source

pub async fn get_entity( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Entity>

Retrieve an entity by ID.

UUID v4 is globally unique: no namespace filter on by-ID ops.

Interim identifier-continuity disclosure (precedes the full transitive redirect chase): a miss is probed once against the tombstone row. If the id was consumed by merge(into_id, from_id) — merged_into set — the NotFound message names the kept id so the caller can requery it directly. Single-level only: it does not chase a chain of merges and does not return the kept entity in place of the miss. The probe only runs after the live-row lookup misses, so the happy path pays no extra query.

Source

pub async fn get_entity_including_deleted( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Entity>>

Retrieve an entity by ID including soft-deleted rows.

UUID v4 is globally unique: no namespace filter on by-ID ops.

Source

pub async fn get_note_including_deleted( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Note>>

Retrieve a note by ID including soft-deleted rows.

UUID v4 is globally unique: no namespace filter on by-ID ops.

Source

pub async fn get_entities_by_ids( &self, token: &NamespaceToken, ids: &[Uuid], ) -> RuntimeResult<Vec<Entity>>

Fetch multiple entities by ID, returning only those that exist in the caller’s namespace. Missing or namespace-mismatched IDs are silently omitted so that batch lookups don’t abort on a single stale reference.

Source

pub async fn list_entities( &self, token: &NamespaceToken, kind: Option<&str>, entity_type: Option<&str>, limit: u32, offset: u32, ) -> RuntimeResult<Vec<Entity>>

List entities visible to the token, optionally filtered by kind and entity_type. A null entity_type falls back to a string properties.type for filtering only.

When the token carries a multi-namespace visible set, entities from all visible namespaces are returned. When the visible set is [primary] (the default) this behaves identically to the pre-visibility behaviour.

Source

pub async fn list_entities_filtered( &self, token: &NamespaceToken, filter: EntityFilter, limit: u32, offset: u32, ) -> RuntimeResult<Vec<Entity>>

Apply a composed entity predicate before offset pagination. Namespace visibility is supplied by the token, just as for the scalar list API.

Source

pub async fn list_entities_after( &self, token: &NamespaceToken, kind: Option<&str>, entity_type: Option<&str>, tags_any: &[String], after: Option<Uuid>, limit: u32, ) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)>

List an immutable insertion-sequence page of visible entities.

The public cursor remains the UUID of the last returned entity. We resolve its immutable database-assigned sequence before querying so callers do not need to serialize storage details. A missing or out-of-scope cursor fails explicitly instead of silently resuming from the wrong boundary.

Source

pub async fn list_entities_after_filtered( &self, token: &NamespaceToken, filter: EntityFilter, after: Option<Uuid>, limit: u32, ) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)>

Apply a composed entity predicate before insertion-sequence pagination, preserving the scalar API’s cursor validation and token visibility.

Source

pub async fn list_entities_tagged( &self, token: &NamespaceToken, kind: Option<&str>, domain_tag: Option<&str>, limit: u32, offset: u32, ) -> RuntimeResult<Vec<Entity>>

List entities filtered by kind, optional domain tag, limit, and offset.

When domain_tag is Some, the query is restricted at the storage layer via EntityFilter::tags_any so the page result already reflects the domain constraint. This avoids the silent truncation that occurs when filtering post-page (K-3). Multi-namespace visibility from the token is applied.

Source

pub async fn count_entities_tagged( &self, token: &NamespaceToken, kind: Option<&str>, domain_tag: Option<&str>, ) -> RuntimeResult<u64>

Count entities filtered by kind and optional domain tag.

Used to report a meaningful total alongside a paginated listing (K-6). Multi-namespace visibility from the token is applied.

Source

pub async fn list_events( &self, token: &NamespaceToken, filter: EventFilter, page: PageRequest, ) -> RuntimeResult<Page<Event>>

List events in the namespace proven by the caller token.

Public delegator for cross-backend link validation.

Exposes validate_edge_relation_endpoints for the SubstrateCoordinator so it can validate the relation before writing the edge on the source backend.

Validate an edge relation using pre-fetched endpoint records.

For cross-backend links the source and target live on different backends — the source runtime cannot resolve the target. The coordinator fetches each endpoint from its own backend, then calls this method to enforce the kind-pairing rules without a second DB round-trip.

src and tgt are the resolve_edge_endpoint results from each backend. The token supplies the pack edge rules installed on this (source) runtime; no DB access is performed.

Source

pub fn validate_annotates_endpoint_kinds( &self, source_id: Uuid, target_id: Uuid, source: Option<EdgeEndpointKind>, target: Option<EdgeEndpointKind>, ) -> RuntimeResult<()>

Validate an annotates edge relation using pre-located endpoint kinds.

Sibling of Self::validate_link_endpoints_by_resolved for callers that only have an EdgeEndpointKind (entity/note/event/edge) rather than a full Resolved record — the SubstrateCoordinator’s cross-backend locate_endpoint resolves edge-substrate UUIDs too (matching get’s by-ID resolution order), but edges have no Resolved variant, so validate_link_endpoints_by_resolved cannot express them.

annotates is the only relation this covers: source must be a note, target may be any substrate (entity, note, event, or edge).

Create a directed edge between two substrates.

Enforces the three-case relation contract via validate_edge_relation_endpoints. See that method for the full contract.

For symmetric relations (competes_with, composed_with) the endpoint pair is canonicalised to source_uuid < target_uuid so that A→B and B→A deduplicate to one row.

metadata is validated against governed keys; dependency_kind is inferred for depends_on edges when absent.

target_backend is always None for locally-routed edges written through this path. Both endpoints must exist in the local namespace, so setting target_backend = None is the only valid choice.

Endpoint existence is a by-ID check and namespace-agnostic: a record that exists in a different namespace than the caller still resolves, exactly as get() would.

Observable form of Self::link. Live natural-key conflicts retain the accepted replace semantics, while tombstones require the explicit resurrect opt-in. The returned preimage and disposition are derived inside the graph writer transaction and drive the lifecycle event.

Write an edge with an explicit target_backend stamp (ADR-029 D3).

Called by the SubstrateCoordinator when source and target are on different backends. The coordinator validates endpoints before calling this method via Self::validate_link_endpoints, so endpoint validation is skipped here. The edge is written on the source backend only.

Policy-aware cross-backend form of Self::link_observed. Endpoint validation remains the coordinator’s responsibility; mutation classification and tombstone handling stay inside the source store.

Source

pub async fn latest_annotating_note( &self, token: &NamespaceToken, node_id: Uuid, kind: &str, tag: &str, ) -> RuntimeResult<Option<Uuid>>

Find the newest live annotation note with an exact kind and tag across the token’s visible edge namespaces, on this runtime’s bound backend. Each store selects one eligible candidate before returning; note bodies and the complete annotation history are never hydrated here.

Source

pub async fn neighbors( &self, token: &NamespaceToken, node_id: Uuid, direction: Direction, limit: Option<u32>, relations: Option<Vec<EdgeRelation>>, ) -> RuntimeResult<Vec<NeighborHit>>

Get immediate neighbors of a node, optionally filtered by relation type.

Pass relations: Some(vec![EdgeRelation::Annotates]) to retrieve only annotation edges, enabling cross-substrate navigation.

Symmetric relations (competes_with, composed_with) are stored with the canonical source as the lower UUID. Direction normalization is applied in neighbors_with_query so both callers see correct results.

Source

pub async fn neighbors_with_query( &self, token: &NamespaceToken, node_id: Uuid, query: NeighborQuery, ) -> RuntimeResult<Vec<NeighborHit>>

Get neighbors with full query control (includes min_weight).

Applies symmetric-relation direction normalization: if the relations filter contains only symmetric relations the direction is overridden to Both so edges stored in canonical order are always found.

Soft-deleted entity nodes are excluded from results unless the caller explicitly requested them (future: include_deleted flag; currently always false).

Source

pub async fn neighbors_with_query_page( &self, token: &NamespaceToken, node_id: Uuid, query: NeighborQuery, after: Option<NeighborCursor>, neighbor_kinds: Option<Vec<String>>, enrich: bool, ) -> RuntimeResult<Vec<NeighborHit>>

Get a deterministic neighbor page, optionally applying a continuation cursor and filtering entity/note kinds before the storage limit. enrich is false for the lightweight edge projection.

Source

pub async fn annotation_neighbors_by_target_id( &self, target_id: Uuid, ) -> RuntimeResult<Vec<NeighborHit>>

Find live annotates edges targeting one record without applying a namespace predicate.

This is the graph counterpart to the namespace-agnostic by-ID get contract (ADR-007 Rev 6). Multi-record neighbor traversal remains visibility-scoped; callers should use this only after resolving a live target through a by-ID operation.

Source

pub async fn neighbors_with_query_directed( &self, token: &NamespaceToken, node_id: Uuid, query: NeighborQuery, ) -> RuntimeResult<Vec<(NeighborHit, Direction)>>

Get both-direction neighbors, each tagged with the direction (Out/ In) it was found in, via a single storage query per visible namespace instead of two separate direction-scoped neighbors_with_query calls: halving the neighbor SELECT count for context(direction="both") expansion. query.direction is ignored: always both.

Mirrors neighbors_with_query’s dedup/enrich/soft-delete-filter/order pipeline exactly, carrying the per-hit direction tag through unchanged.

Source

pub async fn traverse( &self, token: &NamespaceToken, request: TraversalRequest, ) -> RuntimeResult<Vec<GraphPath>>

Traverse the graph from a set of root nodes.

Full-UUID roots use the by-ID contract; expansion and returned edges remain scoped to the caller’s visible namespaces. Missing roots refuse. Soft-deleted entity nodes are excluded from results.

Source

pub async fn create_note( &self, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, salience: Option<f64>, properties: Option<Value>, annotates: Vec<Uuid>, ) -> RuntimeResult<Note>

Create and persist a note, optionally with properties and annotation targets.

After creating the note:

  • Always indexes into FTS5 at the notes_<namespace> key.
  • If an embedding model is configured, indexes into the vector store with SubstrateKind::Note.
  • For each UUID in annotates, creates an EdgeRelation::Annotates edge from the note to that target.
Source

pub async fn create_note_with_embedding_content( &self, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, embedding_content: Option<&str>, salience: Option<f64>, properties: Option<Value>, annotates: Vec<Uuid>, ) -> RuntimeResult<Note>

Like Self::create_note, but lets the caller supply a smaller text to send to the vector embedder while the note’s stored/FTS-indexed content remains the full text.

embedding_content, when Some, must be non-empty and a proper prefix of content — anything else is rejected with InvalidInput before any write. None behaves exactly like Self::create_note. Use this when content may exceed an embedder’s input cap (e.g. a very long commit message) and only a capped head prefix should be embedded, while the full text is still stored and searchable via FTS.

Source

pub async fn create_note_with_embedding_content_and_report( &self, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, embedding_content: Option<&str>, salience: Option<f64>, properties: Option<Value>, annotates: Vec<Uuid>, ) -> RuntimeResult<(Note, EmbeddingTruncationReport)>

Source

pub async fn create_note_with_decay( &self, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, salience: Option<f64>, decay_factor: f64, properties: Option<Value>, annotates: Vec<Uuid>, ) -> RuntimeResult<Note>

Like Self::create_note but also sets a non-zero decay factor on the note.

Source

pub async fn create_note_with_decay_for_embedding_model( &self, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, salience: Option<f64>, decay_factor: f64, properties: Option<Value>, annotates: Vec<Uuid>, embedding_model: Option<&str>, ) -> RuntimeResult<Note>

Like Self::create_note_with_decay but targets a specific embedding model.

Source

pub async fn try_create_note( &self, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, properties: Option<Value>, ) -> RuntimeResult<Option<Note>>

Insert a note using INSERT OR IGNORE semantics for atomic deduplication.

Returns Ok(Some(note)) when the note was newly written. Returns Ok(None) when a unique constraint (e.g. the external_id partial index on comm message notes) was already satisfied by an existing row, making this call a no-op. FTS indexing and vector embedding are attempted on success but treated as best-effort: failures are logged and do not abort the write.

This method is intentionally narrower than create_note: it skips salience/decay, annotates edges, and embedding-model selection, which are not needed for channel-ingest paths.

Rejects quarantined / channel_kind / channel_slug on a message note: those three properties are transport-owned evidence that comm.health trusts at face value, and this fast path (unlike the generic create verb funnel) is not covered by the pack-installed note-write validator. Only the trusted channel-ingest path may establish them — see Self::try_create_note_as_trusted_ingest.

Source

pub async fn try_create_note_as_trusted_ingest( &self, _capability: &ChannelIngestCapability, token: &NamespaceToken, kind: &str, name: Option<&str>, content: &str, properties: Option<Value>, ) -> RuntimeResult<Option<Note>>

Like Self::try_create_note but permits the caller to establish the transport-owned message properties (quarantined, channel_kind, channel_slug).

This is a deliberately named, separate entry point rather than a flag on try_create_note so the trust decision is visible at every call site: comm.ingest (khive-pack-comm/src/handlers.rs) is the sole legitimate caller, because it is the only code that has just derived quarantine disposition and channel provenance from the inbound transport itself. The caller set is bounded by possession, not documentation: the required crate::ChannelIngestCapability is constructible only inside this crate and granted at pack registration exclusively to channel-transport packs. Every other write path uses try_create_note, which rejects those three properties unconditionally.

Source

pub async fn list_notes( &self, token: &NamespaceToken, kind: Option<&str>, limit: u32, offset: u32, ) -> RuntimeResult<Vec<Note>>

List notes visible to the token, optionally filtered by kind.

When the token carries a multi-namespace visible set, notes from all visible namespaces are returned. When the visible set is [primary] (the default) this behaves identically to the pre-visibility behaviour.

Source

pub async fn list_notes_after( &self, token: &NamespaceToken, kind: Option<&str>, after: Option<Uuid>, limit: u32, ) -> RuntimeResult<(Vec<Note>, Option<Uuid>)>

List an immutable insertion-sequence page of visible notes.

Soft-deleting the prior page’s last note does not invalidate the cursor because the boundary is resolved including tombstones. A hard deletion makes the cursor unresolvable and returns an explicit error.

Source

pub async fn count_notes( &self, token: &NamespaceToken, kind: Option<&str>, ) -> RuntimeResult<u64>

Count notes matching kind across the caller’s visible namespaces.

Source

pub async fn search_notes( &self, token: &NamespaceToken, query_text: &str, query_vector: Option<Vec<f32>>, limit: u32, note_kind: Option<&str>, include_superseded: bool, tags_any: &[String], properties_filter: Option<&Value>, ) -> RuntimeResult<Vec<NoteSearchHit>>

Search notes using a hybrid FTS5 + vector pipeline with salience weighting.

Pipeline:

  1. FTS5 query against notes_<namespace>.
  2. If embedding model is configured: vector search filtered to kind="note".
  3. RRF fusion (k=60).
  4. Salience-weighted rerank: score *= (0.5 + 0.5 * note.salience).
  5. Filter soft-deleted notes, apply optional kind / tag / properties predicates. Tags and properties are pushed into the per-note fetch loop BEFORE truncation so that matching notes ranked beyond limit in the raw fusion are not silently dropped.
  6. Truncate to limit.

tags_any: when non-empty, only notes that have at least one of these tags (stored in properties["tags"], case-insensitive match) are retained. The check happens inside the alive-note loop, before hits.truncate(limit).

properties_filter: when Some, only notes whose properties JSON object is a superset of the given filter object are retained. Also applied before truncation.

Source

pub async fn search_notes_with_text_mode( &self, token: &NamespaceToken, query_text: &str, query_vector: Option<Vec<f32>>, limit: u32, note_kind: Option<&str>, include_superseded: bool, tags_any: &[String], properties_filter: Option<&Value>, text_mode: TextQueryMode, ) -> RuntimeResult<Vec<NoteSearchHit>>

Note search with an explicit lexical mode for the text arm.

Source

pub async fn search_notes_outcome( &self, token: &NamespaceToken, query_text: &str, limit: u32, note_kind: Option<&str>, include_superseded: bool, tags_any: &[String], properties_filter: Option<&Value>, ) -> RuntimeResult<NoteSearchOutcome>

Coordinator fan-out variant of Self::search_notes: the text arm still fails loud, but a vector-arm failure after a successful text leg is captured instead of discarding the text hits — mirrors Self::hybrid_search_outcome’s contract for the entity substrate. Reserved for SubstrateCoordinator::fan_out_search_with_visibility; every other caller keeps the fail-loud Self::search_notes contract.

Source

pub async fn search_notes_outcome_with_text_mode( &self, token: &NamespaceToken, query_text: &str, limit: u32, note_kind: Option<&str>, include_superseded: bool, tags_any: &[String], properties_filter: Option<&Value>, text_mode: TextQueryMode, ) -> RuntimeResult<NoteSearchOutcome>

Coordinator note-search variant with an explicit lexical mode.

Source

pub async fn resolve_prefix( &self, token: &NamespaceToken, prefix: &str, ) -> RuntimeResult<Option<Uuid>>

Resolve a short UUID prefix (8+ hex chars) to a full UUID.

Searches entities, notes, and edges tables for a UUID starting with the given prefix, scoped to the caller’s primary namespace only. Returns Ok(Some(uuid)) if exactly one match is found, Ok(None) if no matches, or an error if ambiguous (multiple matches).

Source

pub async fn resolve_prefix_including_deleted( &self, token: &NamespaceToken, prefix: &str, ) -> RuntimeResult<Option<Uuid>>

Source

pub async fn resolve_prefix_unfiltered( &self, prefix: &str, ) -> RuntimeResult<Option<Uuid>>

Resolve a short UUID prefix (8+ hex chars) to a full UUID with NO namespace filter at all: mirrors resolve_by_id’s by-ID contract: by-ID resolution is namespace-agnostic, since the Gate (not storage-layer filtering) is the authz seam. Used by the four by-ID CRUD verbs (get/update/delete/merge) so their prefix path matches their already-unfiltered full-UUID path. No token param: unlike resolve_prefix, there is no namespace to derive from one.

Source

pub async fn resolve_prefix_unfiltered_including_deleted( &self, prefix: &str, ) -> RuntimeResult<Option<Uuid>>

resolve_prefix_unfiltered, including soft-deleted rows — used by the hard-delete by-ID path.

Source

pub async fn resolve_by_id( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Resolved>>

Resolve a UUID to its substrate kind with NO namespace filter.

By-ID contract: UUID v4 is globally unique: by-ID substrate inference must return the record regardless of caller namespace. Used by the public update and delete verb handlers when no explicit kind is supplied.

Does NOT consult the visible set or the primary-namespace check. The token is still required to route to the correct backend pool but its namespace value is not used as a filter.

Source

pub async fn resolve_by_id_including_deleted( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Resolved>>

Resolve a UUID to its substrate kind with NO namespace filter, including soft-deleted rows.

Used by the hard-delete path when no explicit kind is supplied, so already-soft-deleted records can still be located by UUID alone.

Source

pub async fn resolve( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Resolved>>

Resolve a UUID to its substrate kind by trying entity, then note, then event stores.

Returns None if the UUID is not found in any substrate. Cost: at most 3 store lookups per call (cheap for v0.1).

Source

pub async fn resolve_edge_endpoint( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Resolved>>

Resolve a UUID to its substrate kind with NO namespace filter, for edge endpoint validation.

link and create’s annotates targets consume by-ID endpoints, so their existence check must follow the same by-ID contract as get(): by-ID ops are namespace-agnostic: the Gate, not storage-layer filtering, is the authz seam. Mirrors resolve_by_id (entity + note, unfiltered) and additionally resolves events, unfiltered, so edge endpoint validation resolves exactly what get() resolves regardless of the caller’s namespace.

Source

pub async fn resolve_primary( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Resolved>>

Resolve a UUID to its substrate kind using primary-namespace-only enforcement.

Unlike resolve, never consults the visible set. Use from GTD dependency validation paths where strict primary ownership is required.

Source

pub async fn resolve_including_deleted( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<Resolved>>

Resolve a UUID to its substrate kind, including soft-deleted rows.

Used exclusively by the hard-delete path to locate records that have already been soft-deleted. Namespace isolation is still enforced.

Source

pub async fn restore_entity( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<(Entity, bool)>>

Restore an entity tombstone owned by the caller’s primary namespace.

The restore is guarded by both the tombstone’s id/namespace and the current uniqueness state, so a caller cannot resurrect over a newer live record. Indexes are rebuilt only after the row restore commits.

Source

pub async fn restore_note( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<(Note, bool)>>

Restore a note tombstone owned by the caller’s primary namespace.

A live note holding the tombstone’s (namespace, kind, key) refuses the operation before any row changes. The same condition is repeated in the guarded restore statement for the concurrent race.

Source

pub async fn restore_edge( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<Option<(Edge, bool)>>

Restore an edge tombstone owned by the caller’s primary namespace.

Source

pub async fn delete_note( &self, token: &NamespaceToken, id: Uuid, hard: bool, ) -> RuntimeResult<bool>

Soft-delete or hard-delete a note by ID.

On hard delete, cascades to remove all incident edges (both inbound and outbound) and cleans up FTS and vector indexes, preventing dangling references for annotates edges that target this note. Soft delete also cleans FTS and vector indexes; edges are left in place.

UUID v4 is globally unique: no namespace filter on by-ID ops. Cascade and index cleanup target the RECORD’s stored namespace, not the caller token’s. Returns Ok(false) if the note does not exist.

Source

pub async fn delete_note_row_first_for_compensation( &self, token: &NamespaceToken, id: Uuid, ) -> RuntimeResult<()>

Row-first compensating delete for rolling back a partially-written note (e.g. dual_write_message rollback after a later delivery step fails). Unlike KhiveRuntime::delete_note, which cleans up graph/FTS/ vector indexes before removing the row, this removes the row first so that a cleanup failure afterward cannot leave the compensated note live.

Returns Ok(()) once the row is gone (whether or not cleanup fully succeeded). Returns Err(RuntimeError::Internal) naming the failed cleanup legs when row removal succeeded but cleanup did not — the message is gone, but stale index entries may remain and should be surfaced to the caller rather than silently discarded.

Returns Ok(()) immediately, with no cleanup attempted, if the note does not exist (nothing to compensate).

Not a general-purpose replacement for delete_note(..., hard=true): normal hard delete still needs cleanup-first semantics (no dangling references) since a caller-visible error there should not remove the row.

Source§

impl KhiveRuntime

Source

pub const EDGE_LIST_MAX_LIMIT: u32 = 1000

Maximum rows returned by a single Self::list_edges / Self::list_edges_after page. A lower bound the docs promise callers can rely on; kept as a named constant so tests can exercise pagination (page tiling, out-of-range offsets) without needing >1000 real rows.

Source

pub async fn query( &self, token: &NamespaceToken, query: &str, ) -> RuntimeResult<Vec<SqlRow>>

Execute a GQL or SPARQL query string, returning raw SQL rows.

The query is compiled to SQL with the namespace scope applied. GQL syntax: MATCH (a:concept)-[e:extends]->(b) RETURN a, b LIMIT 10 SPARQL syntax: SELECT ?a WHERE { ?a :kind "concept" . }

Source

pub async fn query_with_metadata( &self, token: &NamespaceToken, query: &str, opts: CompileOptions, ) -> RuntimeResult<QueryResult>

Execute a GQL/SPARQL query, returning rows and any validation warnings.

Source

pub async fn delete_entity( &self, token: &NamespaceToken, id: Uuid, hard: bool, ) -> RuntimeResult<bool>

Soft-delete or hard-delete an entity by ID (soft delete by default).

On hard delete, cascades to remove all incident edges (both inbound and outbound) to prevent dangling references. Soft delete also cleans FTS and vector indexes; edges are left in place. Routed attachment cleanup is performed by the registry after ownership resolution.

UUID v4 is globally unique: no namespace filter on by-ID ops.

Source

pub async fn count_entities( &self, token: &NamespaceToken, kind: Option<&str>, ) -> RuntimeResult<u64>

Count entities in a namespace, optionally filtered.

Source

pub async fn get_edge( &self, _token: &NamespaceToken, edge_id: Uuid, ) -> RuntimeResult<Option<Edge>>

Fetch a single edge by id.

UUID v4 is globally unique: returns the edge regardless of which namespace the token carries. Ok(None) means the edge does not exist at all.

Source

pub async fn get_edge_visible( &self, token: &NamespaceToken, edge_id: Uuid, ) -> RuntimeResult<Option<Edge>>

Fetch a single edge by id.

Delegates to get_edge: no visible-set check. By-ID ops are namespace-agnostic; UUID v4 is globally unique.

Source

pub async fn get_edge_including_deleted( &self, _token: &NamespaceToken, edge_id: Uuid, ) -> RuntimeResult<Option<Edge>>

Fetch an edge by UUID including soft-deleted rows.

Returns the edge regardless of which namespace the token carries: UUID v4 is globally unique. Used by the hard-delete path so that a soft-deleted edge can still be purged via its edge ID.

Source

pub async fn get_edge_by_natural_key_including_deleted( &self, token: &NamespaceToken, namespace: &str, source_id: Uuid, target_id: Uuid, relation: EdgeRelation, ) -> RuntimeResult<Option<Edge>>

Fetch an edge by natural key (namespace, canonical source/target, relation) including soft-deleted rows. Unlike Self::list_edges/Self::list_edges_after, which always filter deleted_at IS NULL, this can render a tombstoned symmetric-edge survivor (ADR-039 DO NOTHING conflict absorption) — used by the atomic-apply post-commit result renderer, which otherwise reports “not found” for a committed update whose surviving row happens to be soft-deleted.

token selects the store instance; namespace is the natural key’s own namespace and is bound into the query explicitly. These can legitimately differ: the record namespace is fixed at prepare time (EdgeNaturalKey::namespace) and by-ID edge updates are namespace-agnostic (ADR-007 Rev 6), so the caller’s ambient token namespace is never a safe substitute for the record’s own — the prior implicit self.namespace scoping is exactly the bug this parameter closes (khive#1213/#1214).

Source

pub async fn list_edges( &self, token: &NamespaceToken, filter: EdgeListFilter, limit: u32, offset: u32, ) -> RuntimeResult<Vec<Edge>>

List edges matching filter, paging by offset. limit is capped at Self::EDGE_LIST_MAX_LIMIT; defaults to 100. For O(1)-at-depth walks over large edge populations, prefer Self::list_edges_after instead of paging offset deep.

Source

pub async fn list_edges_after( &self, token: &NamespaceToken, filter: EdgeListFilter, after: Option<Uuid>, limit: u32, ) -> RuntimeResult<(Vec<Edge>, Option<Uuid>)>

Keyset (seek) page of edges matching filter, ordered by immutable database-assigned insertion sequence. after is the last edge id from the previous page (exclusive); omit to start from the beginning. Returns (items, next_after) — next_after is Some when more rows remain past this page.

Unlike Self::list_edges, this is O(log n + limit) at any depth and genuinely new inserts are appended after already-issued boundaries. The cursor row is resolved including tombstones; a hard-deleted or out-of-scope cursor fails explicitly rather than hiding an incomplete traversal.

Source

pub async fn count_edges_by_relation( &self, token: &NamespaceToken, ) -> RuntimeResult<HashMap<String, u64>>

Count edges by relation, ignoring soft-deleted rows. Used by stats() to report the true per-relation population so full-graph audits know what they’re sampling from before they walk it.

Source

pub async fn count_edges_by_endpoint_base( &self, token: &NamespaceToken, ) -> RuntimeResult<EdgeEndpointBaseCounts>

Count edges by the base each endpoint resolves against. Used by stats() so a caller can name the denominator of a density figure instead of inheriting the flat edge total, which on a real store is mostly provenance.

The per-namespace fallback sums the same buckets, so the aggregate and the fallback are checkable against each other and against count_edges by the invariant that the buckets sum to the total.

Source

pub async fn update_edge( &self, token: &NamespaceToken, edge_id: Uuid, patch: EdgePatch, ) -> RuntimeResult<Edge>

Patch-style edge update. Only Some(_) fields are applied.

When relation is Some(new_rel), validates that the edge’s existing endpoints are legal for new_rel before persisting. Weight-only updates (relation = None) skip validation. Returns InvalidInput if the new relation would violate the three-case endpoint contract; the edge is NOT mutated on error.

For symmetric relations (competes_with, composed_with), endpoint order is canonicalised to source_uuid < target_uuid after validation. If a canonical row already exists at the target triple, the non-canonical edge is deleted and the existing canonical row is preserved unchanged (ADR-039 ON CONFLICT DO NOTHING, mirroring merge_entity_sql) — its attributes, including a soft-deleted deleted_at, are never overwritten by the discarded edge’s patch.

Source

pub async fn delete_edge( &self, token: &NamespaceToken, edge_id: Uuid, hard: bool, ) -> RuntimeResult<bool>

Hard-delete an edge by id.

Cascades to remove any annotates edges whose target is the deleted edge (annotates is note → anything; deleting an edge target leaves annotation edges dangling if not cleaned up). Returns true if the primary edge was removed.

If edge_id does not refer to an edge (e.g. the caller passes an entity or note UUID by mistake), this method returns Ok(false) immediately with no side effects — it does not cascade inbound edges of the non-edge record.

Source

pub async fn count_edges( &self, token: &NamespaceToken, filter: EdgeListFilter, ) -> RuntimeResult<u64>

Count edges matching filter across the caller’s visible namespaces.

Source

pub async fn build_edge( &self, token: &NamespaceToken, spec: &LinkSpec, ) -> RuntimeResult<Edge>

Validate and construct an edge from a LinkSpec without writing to storage.

Applies the full edge contract (endpoint validation, symmetric canonicalization, dependency_kind inference and metadata validation). Returns the constructed Edge on success; the caller is responsible for persisting it (e.g. via upsert_edge or link_many).

The token must be a pre-authorized namespace token from the dispatch layer. If spec.namespace is set it must match token.namespace(); a mismatch returns RuntimeError::InvalidInput.

Validate and atomically upsert a batch of edges.

All edges are validated and constructed with build_edge before any write. If validation fails for any entry the entire batch is rejected (no writes occur). On success, all edges are persisted in a single atomic transaction via upsert_edges.

After the bulk upsert, each edge is read back by its natural key (namespace, source_id, target_id, relation) so that the returned IDs are always the persisted row IDs, not the locally-generated UUIDs that may have been displaced by an ON CONFLICT DO UPDATE. This mirrors the same read-back applied to singleton link() and prevents phantom-ID exposure when callers upsert overlapping triples with verbose=true.

All specs must share the same namespace; the namespace is taken from token (or validated against it if spec.namespace is set).

Observed all-or-nothing bulk link upsert. Every row carries its own create/update/resurrection disposition, and every tombstone policy is preflighted inside the same writer transaction before any mutation.

Source

pub async fn create_many( &self, token: &NamespaceToken, specs: Vec<EntityCreateSpec>, ) -> RuntimeResult<Vec<Entity>>

Create a batch of entities atomically.

All specs are validated before any write. If ANY spec fails validation (unknown kind, empty name, secret-gate violation), the method returns that error and no entities are written.

Entity rows and their FTS documents are written in one SQLite transaction. Any statement failure rolls back the entire batch across both surfaces. Embedding is intentionally skipped: bulk structural ingest is the expected use-case, and dense vectors are backfilled later via a reindex call.

Source

pub async fn prepare_bulk_entity_plan( &self, token: &NamespaceToken, spec: EntityCreateSpec, ) -> RuntimeResult<(Entity, AtomicOpPlan)>

Validate and prepare one entity item for a bulk create(items=[...]) write: the same admission and row/FTS plan as Self::create_many, with no scheduled reindex, so the vector is deferred to a later reindex exactly as for create_many. The bulk create handler uses this for every entity item so entity and note plans can join one run_atomic_unit call.

Source

pub async fn prepare_bulk_note_plan( &self, token: &NamespaceToken, spec: NoteCreateSpec, ) -> RuntimeResult<(Note, AtomicOpPlan)>

Validate and prepare one note item for a bulk create(items=[...]) write. The note goes through the preparation a singleton note create uses (validate_head, then prepare_atomic_notes: kind validation, owned-identity derivation, secret gate, salience range, row and FTS statements) with embedding switched off, so the plan writes the row and its FTS document and no vector. A later reindex backfills the vector, as it does for bulk entities. The caller commits the plan, alone or joined with its siblings in one run_atomic_unit call.

Source§

impl KhiveRuntime

Source

pub async fn export_kg( &self, token: &NamespaceToken, ) -> RuntimeResult<KgArchive>

Export all entities and edges in a namespace to a portable JSON archive.

Edge collection: all entity IDs in the namespace are gathered first; query_edges is then called with those IDs as source_ids. This captures every edge whose source entity belongs to the namespace.

Source

pub async fn export_kg_json( &self, token: &NamespaceToken, ) -> RuntimeResult<String>

Export to a JSON string (convenience wrapper around export_kg).

Source

pub async fn import_kg( &self, archive: &KgArchive, token: &NamespaceToken, ) -> RuntimeResult<ImportSummary>

Import an archive into target_namespace.

If target_namespace is None, the archive’s own namespace is used.

  • Entities: upserted by ID; existing records are overwritten.
  • Edges: upserted; existing records are overwritten.
  • Validation: format != "khive-kg" or unsupported version → InvalidInput. Invalid edge relations are caught at JSON deserialization time.
Source

pub async fn import_kg_json( &self, json: &str, token: &NamespaceToken, ) -> RuntimeResult<ImportSummary>

Import from a JSON string (convenience wrapper around import_kg).

Source§

impl KhiveRuntime

Source

pub async fn embed(&self, text: &str) -> RuntimeResult<Vec<f32>>

Generate an embedding vector for text using the configured default model.

First call lazily loads model weights (cold start cost). Subsequent calls reuse them. Returns Unconfigured("embedding_model") if no model is configured.

Source

pub async fn embed_with_model( &self, model_name: &str, text: &str, ) -> RuntimeResult<Vec<f32>>

Generate an embedding vector for text using the named model.

Accepts both built-in lattice model names/aliases and custom provider names registered via KhiveRuntime::register_embedder. For lattice models the resolved EmbeddingModel enum is forwarded to embed_one so the service can select the correct model variant. For custom providers, embed_one is called with EmbeddingModel::default() because custom services are expected to ignore the enum argument (they own a single model implicitly).

Applies no instruction prefix (generic role). Use Self::embed_document_with_model / Self::embed_query_with_model for instruction-tuned models where the asymmetric prefix matters.

Returns UnknownModel if model_name is not in the embedder registry.

Source

pub async fn embed_document_with_model( &self, model_name: &str, text: &str, ) -> RuntimeResult<Vec<f32>>

Embed a document/passage for indexing using the named model.

Applies EmbeddingService::embed_passage, which prepends the model’s document_instruction() prefix when defined (e.g. "passage: " for multilingual-e5). For models with no document prefix (MiniLM, BGE) this is identical to Self::embed_with_model.

Use this for all index/store/backfill paths so that instruction-tuned models produce passage-side vectors.

Reindex caveat: switching from an unprefixed model (or a model with no document_instruction) to an instruction-tuned model changes the vector representation. Vectors stored under the old scheme are not comparable to newly prefixed vectors. Operators must trigger a full reindex (knowledge.index(rebuild_ann=true) / kkernel reindex) after changing the embedding model config.

Returns UnknownModel if model_name is not registered.

Source

pub async fn embed_document_with_model_outcome( &self, model_name: &str, text: &str, ) -> RuntimeResult<DocumentEmbeddingOutcome>

Source

pub async fn embed_query_with_model( &self, model_name: &str, text: &str, ) -> RuntimeResult<Vec<f32>>

Embed a query string for retrieval using the named model.

Applies the pre-0.6 query-role behavior. E5 and Qwen models retain their query instructions, while BGE and custom providers receive the original unprefixed query text.

Use this for all search/recall/suggest query embedding paths so that instruction-tuned models land in the correct side of their retrieval space.

Returns UnknownModel if model_name is not registered.

Source

pub async fn embed_document(&self, text: &str) -> RuntimeResult<Vec<f32>>

Embed a document for indexing using the configured default model.

Delegates to Self::embed_document_with_model. Use for entity/note create and reindex paths.

Returns Unconfigured("embedding_model") if no model is configured.

Source

pub async fn embed_query(&self, text: &str) -> RuntimeResult<Vec<f32>>

Embed a query for retrieval using the configured default model.

Delegates to Self::embed_query_with_model. Use for vector search and hybrid search query paths.

Returns Unconfigured("embedding_model") if no model is configured.

Source

pub async fn embed_batch( &self, texts: &[String], ) -> RuntimeResult<Vec<Vec<f32>>>

Generate embeddings for multiple texts in one call using the configured default model.

Delegates to the cached EmbeddingService::embed, so repeated texts within and across calls benefit from the runtime-level LRU cache.

Returns an empty vec for empty input without hitting the embedding service. Returns Unconfigured("embedding_model") if no model is configured.

Source

pub async fn embed_batch_with_model( &self, model_name: &str, texts: &[String], ) -> RuntimeResult<Vec<Vec<f32>>>

Generate embeddings for multiple texts using the named model.

Accepts lattice model names/aliases and custom provider names. Returns UnknownModel if model_name is not in the embedder registry.

Source

pub async fn embed_document_batch_with_model( &self, model_name: &str, texts: &[String], ) -> RuntimeResult<Vec<Vec<f32>>>

Embed a batch of documents for indexing using the named model.

Applies EmbeddingService::embed_passage. Use for all bulk index/backfill/reindex operations to apply the passage-side prefix. A mixed batch is bounded into one ordered owned batch so truncation never fragments one provider batch into sequential singleton inference calls.

Reindex caveat: see Self::embed_document_with_model — the same incomparability applies to batch-indexed vectors when switching models.

Returns UnknownModel if model_name is not registered.

Source

pub async fn embed_document_batch_with_model_outcomes( &self, model_name: &str, texts: &[String], ) -> RuntimeResult<Vec<DocumentEmbeddingOutcome>>

Source

pub async fn embed_document_batch( &self, texts: &[String], ) -> RuntimeResult<Vec<Vec<f32>>>

Embed a batch of documents for indexing using the configured default model.

Convenience delegate to Self::embed_document_batch_with_model. Use for bulk knowledge-atom and section indexing paths.

Returns Unconfigured("embedding_model") if no model is configured.

Source

pub async fn embed_document_batch_outcomes( &self, texts: &[String], ) -> RuntimeResult<Vec<DocumentEmbeddingOutcome>>

Embed documents with the default model and retain actual input-bounding metadata.

Source

pub async fn embed_query_batch_with_model( &self, model_name: &str, texts: &[String], ) -> RuntimeResult<Vec<Vec<f32>>>

Embed a batch of queries for retrieval using the named model.

Applies the same pre-0.6 query-role compatibility behavior as Self::embed_query_with_model.

Returns UnknownModel if model_name is not registered.

Search vectors using either a caller-provided embedding or query text.

Existing callers pass query_embedding: Some(vec) to avoid re-embedding. Text callers pass query_embedding: None, query_text: Some(...) and the runtime embeds internally.

Hybrid search: text (FTS5) + vector retrieval fused via Reciprocal Rank Fusion.

  • Always performs text search over query_text.
  • If query_vector is Some, also performs vector search and fuses both lists.
  • If None, returns text-only results — no vector store needed.
  • If entity_kind is Some, the alive-set query filters to that kind. The text/vector candidate pools are unfiltered up front; the kind filter applies at the alive-check stage where we already fetch each candidate to confirm it isn’t soft-deleted.
  • tags_any: when non-empty, only entities that have at least one of these tags (case-insensitive) survive the alive-set intersection. Applied BEFORE truncation so matches ranked beyond limit in the raw fusion are not lost.
  • properties_filter: when Some, only entities whose properties are a superset of the given JSON object survive. Applied BEFORE truncation.

limit caps the final returned list; internally pulls limit * 4 candidates per path.

§Cross-namespace visibility (entity search — primary namespace only; deferred)

Both the FTS leg and the vector/ANN leg of entity search (hybrid_search) are restricted to the primary namespace only.

Rationale: each namespace owns a separate FTS table (fts_entities_{ns}) and a separate ANN index instance. Cross-namespace entity-search fanout requires iterating over every visible namespace’s store, issuing parallel search requests, and fusing the results: this is deferred.

Note: this is distinct from memory.recall’s cross-namespace fanout, which already iterates visible_namespaces across both the FTS and vector legs. Entity search fanout is the remaining deferred piece; memory recall fanout is not deferred.

The visible_ns list is forwarded in the TextFilter.namespaces field, which limits results to those namespaces within the primary store. Because entities from extra namespaces live in their own FTS tables, this filter has no cross-namespace effect today.

Callers with a multi-namespace visible set can READ cross-namespace entities directly via get_entity / resolve, but hybrid_search returns only primary-namespace hits until entity-search cross-namespace fanout ships.

Source

pub async fn hybrid_search_with_text_mode( &self, token: &NamespaceToken, query_text: &str, query_vector: Option<Vec<f32>>, limit: u32, entity_kind: Option<&str>, entity_type: Option<&str>, tags_any: &[String], properties_filter: Option<&Value>, text_mode: TextQueryMode, ) -> RuntimeResult<Vec<SearchHit>>

Hybrid search with an explicit lexical mode for the text arm.

Source

pub async fn hybrid_search_outcome( &self, token: &NamespaceToken, query_text: &str, limit: u32, entity_kind: Option<&str>, entity_type: Option<&str>, tags_any: &[String], properties_filter: Option<&Value>, ) -> RuntimeResult<HybridSearchOutcome>

Coordinator fan-out variant of Self::hybrid_search: the text arm still fails loud (propagated through hybrid_search_inner’s tolerate_vector_error=false semantics for that leg), but a vector-arm failure after a successful text leg is captured instead of discarding the text hits. Reserved for SubstrateCoordinator::fan_out_search_with_visibility — every other caller keeps the fail-loud Self::hybrid_search contract, so a single backend’s vector-store outage does not misreport that backend’s text arm as failed too.

Source

pub async fn hybrid_search_outcome_with_text_mode( &self, token: &NamespaceToken, query_text: &str, limit: u32, entity_kind: Option<&str>, entity_type: Option<&str>, tags_any: &[String], properties_filter: Option<&Value>, text_mode: TextQueryMode, ) -> RuntimeResult<HybridSearchOutcome>

Coordinator variant with an explicit lexical mode for the text arm.

Source

pub async fn knn( &self, token: &NamespaceToken, query_vector: Vec<f32>, top_k: u32, ) -> RuntimeResult<Vec<VectorSearchHit>>

Exact KNN over the full namespace’s vector store.

sqlite-vec uses brute-force cosine — results are exact, not approximate. Cost is O(N · D) per query. For small-to-medium namespaces (~hundreds of thousands of vectors) this is well within latency budgets.

Source

pub async fn rerank( &self, token: &NamespaceToken, query_vector: &[f32], candidate_ids: &[Uuid], top_k: u32, ) -> RuntimeResult<Vec<VectorSearchHit>>

Exact KNN restricted to a candidate set.

Useful for reranking the top-N results from hybrid_search (or any other retrieval path) with exact cosine similarity against a query vector. Returns hits sorted by similarity (highest first), truncated to top_k.

Source

pub async fn backfill_missing_embeddings( &self, token: &NamespaceToken, ) -> RuntimeResult<u64>

Backfill vector and FTS index entries for entities and notes that are missing them.

Intended to run once at startup as a background task (warm-up sequence steps 2–4). Queries the SQL substrate for entity descriptions and note contents that have no corresponding entry in the vector store for any registered embedding model, then embeds and inserts them. FTS entries missing for notes are also repopulated.

The operation is best-effort: individual embed/insert failures are logged and skipped rather than aborting the whole backfill. If no embedding models are registered, returns immediately with 0.

Returns the total number of records backfilled across all models.

Source

pub async fn sweep_orphan_vectors( &self, token: &NamespaceToken, max_delete_per_model: u32, dry_run: bool, ) -> RuntimeResult<u64>

Sweep orphaned vector entries for all registered embedding models.

A vector entry is orphaned when its subject_id no longer exists as a live row in the entity, note, or knowledge-atom tables (i.e. either the row is absent or has deleted_at IS NOT NULL). Orphaned entries accumulate after hard-deletes because the vector store and SQL substrate are decoupled.

Iterates over every registered embedding model and calls khive_storage::VectorStore::orphan_sweep for the token’s namespace. Models whose backend returns khive_storage::StorageError::Unsupported are skipped without error — this preserves forward-compat when a newly registered model does not yet implement sweep.

Returns the total number of vector rows deleted across all models.

Source§

impl KhiveRuntime

Source

pub fn new(config: RuntimeConfig) -> RuntimeResult<Self>

Create a new runtime with the given config.

The config’s db_path is used to open or create the SQLite backend. This direct constructor is intended for fresh/current single-backend databases and tests. Production and multi-backend hosts must use the async khive-mcp/kkernel builders so secondary inventory and any application-assisted V21 cutover complete before serving. The from_backend seam is likewise only for an already-prepared backend.

Source

pub fn new_readonly(config: RuntimeConfig) -> RuntimeResult<Self>

Open a runtime for read-only inspection (no model registration, no DB creation).

File-backed databases are opened with SQLite read-only/query-only flags and must already be at this build’s current schema version. No migrations or configured-model registration writes are attempted. A None path retains the historical ephemeral in-memory behavior for tests.

Source

pub fn from_backend(backend: Arc<StorageBackend>, config: RuntimeConfig) -> Self

Construct a runtime from an already-opened backend.

This is a low-level, infallible assembly seam for already-prepared multi-backend deployments. It does not inspect or migrate the V21 attachment-cutover state. Production hosts must first run the async kkernel/khive-mcp coordinator and must not expose a server over a pending or incomplete backend. Prefer Self::from_prepared_backend when constructing one fallible host runtime.

The returned runtime has db_path = None and embedding_model = None; all storage access is through the provided backend. Set backend_id and default_namespace via the config builder pattern if non-defaults are needed.

Source

pub fn from_prepared_backend( backend: Arc<StorageBackend>, config: RuntimeConfig, ) -> RuntimeResult<Self>

Construct a single-backend runtime after a host boot coordinator has completed schema preparation and any application-assisted cutover.

Unlike Self::from_backend, configured embedding-model registration is fallible here, preserving Self::new’s single-backend startup semantics. This method never runs migrations itself.

Source

pub fn with_core_backend(self, core: Arc<StorageBackend>) -> Self

Wire this runtime as a secondary-backend runtime pointing at core.

After this call, self.core() returns a handle to core rather than cloning self. The caller (the boot path, not pack code) is responsible for passing the correct main backend.

Panics in debug builds if self.config.backend_id == BackendId::MAIN, because the main runtime does not need a core pointer.

Source

pub fn with_core_embedders_from(self, main: &KhiveRuntime) -> Self

Carry the main runtime’s embedder wiring for core()-routed writes.

Boot-path companion to with_core_backend. Without it, core() shares this pack runtime’s own embedder registry — which under [packs.<name>] no_embed = true is empty, so core-routed concept writes would silently skip embedding on the shared graph.

Source

pub fn core(&self) -> KhiveRuntime

Return a runtime handle bound to the main (shared-graph) backend.

When self is already the main runtime (core_backend is None), this returns a clone of self — no new backend reference is acquired.

When self is a secondary-backend runtime (core_backend is Some), this returns a new KhiveRuntime backed by the main Arc<StorageBackend> and sharing all registry state (embedder_registry, edge_rules, valid_entity_kinds, valid_note_kinds, entity_type_validator, note_mutation_hook, entity_kind_hooks) with self. No database I/O occurs; no embedding models are reloaded.

Use core() for notes and entities that must reside in the shared graph so that memory.recall, cross-pack search, and annotates edges work. Use self (or self.sql()) for pack-auxiliary bulk tables.

Handlers that call core() more than once per request or loop should bind let core = self.core(); once and reuse it, since each call clones RuntimeConfig (a heap-allocated struct containing Vec<String> fields).

Source

pub fn memory() -> RuntimeResult<Self>

Create an in-memory runtime (for tests and ephemeral use).

Source

pub fn backend_id(&self) -> &BackendId

Return the BackendId for this runtime’s backend.

Used by SubstrateCoordinator in kkernel to identify which backend owns a given node, and to detect cross-backend merges.

Source

pub fn visible_namespaces(&self) -> &[Namespace]

Return the extra-visible namespaces assembled at config load.

OSS dispatch uses this set to widen the default multi-record read scope to ['local'] ∪ visible_namespaces. Writes are unchanged: always pinned to 'local'. This set is also available as gate/cloud policy input.

Source

pub fn config(&self) -> &RuntimeConfig

Return a reference to the runtime config.

Source

pub fn vector_arm_selected(&self) -> bool

Whether this runtime selects the vector arm for a hybrid search — true exactly when a default embedding model is configured. Single source of truth for the policy every fan-out and single-backend dispatch path uses to report arm_participation/vector_selected.

Source

pub fn ann_fresh_tail_enabled(&self) -> bool

Return the immutable ADR-118 fresh-tail serving policy captured when this runtime was constructed.

Source

pub fn with_ann_fresh_tail_enabled(self, enabled: bool) -> Self

Override ADR-118’s fresh-tail serving policy for this runtime instance.

This is primarily useful for embedded runtimes and deterministic tests: it avoids mutating process-global environment state. Clones and core() handles preserve the chosen value.

Source

pub fn backend(&self) -> &StorageBackend

Return a reference to the underlying storage backend.

This is an embedder/infrastructure surface (connection pools, schema plans, diagnostics). Stores obtained from it are NOT wrapped by the message-evidence policy that Self::notes enforces: an embedder holding the backend already holds root-equivalent access to the database file, so the policy boundary sits at the typed accessors pack code uses, not here. Pack code must not take note stores from this surface.

Source

pub fn is_read_only(&self) -> bool

Whether this runtime’s bound backend is explicitly or filesystem-mode detected read-only.

Source

pub fn backend_data_dir(&self) -> Option<PathBuf>

Return the directory containing the backend’s database file, or None for an in-memory backend.

Source

pub fn backend_ann_root(&self) -> Option<PathBuf>

Root directory for this database’s ANN segment tree (<db-file>.ann/ beside the file), or None for an in-memory backend. Scoped to the database file itself so two databases sharing a parent directory can never adopt each other’s segments.

Source

pub async fn db_diagnostics(&self) -> RuntimeResult<DbDiagnostics>

Writer-contention, graph-edge integrity, and WAL/checkpoint diagnostics (ADR-091/ADR-135 operator surface): pooled writer and audit-failure counters, build identity, duplicate edge-ID and list-ledger counts, checkpoint counters, a PASSIVE checkpoint probe, WAL file size, and explicitly qualified WAL-pin census. Not write-free: the PASSIVE probe may backfill WAL frames into the database (normal checkpoint I/O). It never changes logical state, escalates to TRUNCATE, creates a missing database file, or deletes sidecar evidence — see khive_db::diagnostics for the narrowings that make those claims hold.

Always targets the main backend via Self::core, regardless of which backend this runtime handle is bound to, so a report never describes a database this handle is not the canonical owner of.

Source

pub async fn db_diagnostics_with_audit_metrics( &self, runtime_audit_batch_metrics: Option<RuntimeAuditBatchMetrics>, ) -> RuntimeResult<DbDiagnostics>

As Self::db_diagnostics, but with the caller supplying the ADR-133 audit-batch health counters from whichever VerbRegistry owns the seam over this runtime’s EventStore (typically VerbRegistry::audit_batch_metrics()). None behaves identically to Self::db_diagnostics.

Source

pub fn entities( &self, token: &NamespaceToken, ) -> RuntimeResult<Arc<dyn EntityStore>>

Get an EntityStore scoped to the token’s namespace.

Source

pub fn graph( &self, token: &NamespaceToken, ) -> RuntimeResult<Arc<dyn GraphStore>>

Get a GraphStore scoped to the token’s namespace.

Source

pub fn notes(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn NoteStore>>

Get a NoteStore scoped to the token’s namespace.

Wrapped in note_store_guard::PolicyEnforcingNoteStore, which refuses any insert/upsert of a kind = "message" note carrying quarantined / channel_kind / channel_slug — the transport-owned evidence comm.health trusts at face value — and refuses patching those keys through the property-mutation seams on any note kind, so the guard cannot be sidestepped by inserting a clean message note and patching the evidence onto it afterward. Full-row writes also preserve existing channel-health coordinates while allowing heartbeat metadata to change. The trusted channel-ingest path does not go through this accessor; see Self::raw_notes and Self::try_create_note_as_trusted_ingest.

Source

pub fn attachments(&self) -> RuntimeResult<Arc<dyn AttachmentStore>>

Return the role-keyed attachment substrate on the canonical main backend.

Attachment rows are the process-shared BlobStore’s sole SQL liveness authority. A runtime bound directly to a secondary pack backend must call Self::core first; accepting a secondary mutation here would create a reference that the main-database GC sweep cannot see or fence.

Source

pub fn events( &self, token: &NamespaceToken, ) -> RuntimeResult<Arc<dyn EventStore>>

Get an EventStore scoped to the token’s namespace.

When the events-daemon split (ADR-170) is configured, the store routes by append class: the ADR-133 idempotent audit-batch lane — the measured bulk of event write volume — persists to the events database (forwarded over the events daemon socket in daemon deployments, or opened directly in embedded/one-shot contexts), while plain appends stay on this runtime’s backend, keeping every raw-SQL consumer of the legacy events table (schedule provenance, kg projection guards, GraphQuery’s substrate union) correct by construction. Reads merge both stores. Unconfigured runtimes (tests, in-memory) keep the legacy main-store behavior. Every returned store is decorated at this typed accessor boundary so append callers cannot override the namespace or actor resolved into the sealed authorization token.

Source

pub fn sql(&self) -> Arc<dyn SqlAccess> ⓘ

Get the raw SQL access capability (for ad-hoc queries).

Source

pub fn events_sidecar_sql_read_only( &self, ) -> RuntimeResult<Option<Arc<dyn SqlAccess>>>

SQL access to the events-split sidecar database for read purposes, when the split (ADR-170) is configured and the sidecar exists on disk. None means every event row lives in the legacy events table, so a raw-SQL consumer needs no second lookup. Consumers that resolve an event by id or hex prefix against self.sql() must also consult this store on a miss: the audit-batch lane’s rows live only in the sidecar.

A writable runtime opens the sidecar through the ordinary writable binding: WAL supports one-writer-many-readers, and the read-only binding’s frozen-snapshot guard refuses any sidecar with a live writer’s -shm beside it — exactly the live-deployment case these reads exist for. A read-only runtime keeps the read-only open (it must neither create nor schema-initialize a sidecar), which serves genuinely frozen snapshots and refuses live ones, matching the read-only arm of events(). Never creates a sidecar as a side effect of a read.

Source

pub fn vectors( &self, token: &NamespaceToken, ) -> RuntimeResult<Arc<dyn VectorStore>>

Get a VectorStore for the configured embedding model, scoped to the token’s namespace.

Returns Unconfigured("embedding_model") if no model is set.

Source

pub fn vectors_for_model( &self, token: &NamespaceToken, model_name: &str, ) -> RuntimeResult<Arc<dyn VectorStore>>

Get a VectorStore for a specific named embedding model, scoped to the token’s namespace.

Accepts both built-in lattice model names/aliases and custom provider names registered via register_embedder. Lattice names are routed through the enum-backed path; custom provider names use the provider’s declared dimensions() directly so that the vector store key is consistent with how vectors were written during remember/recall.

Source

pub async fn vectors_for_named_identity( &self, token: &NamespaceToken, identity: &NamedVectorIdentity, ) -> RuntimeResult<Arc<dyn VectorStore>>

Get a namespace-scoped vector store for a pack-owned immutable identity.

The table key is syntactically validated by NamedVectorIdentity. This accessor additionally verifies the table’s actual sqlite-vec dimension declaration and every persisted embedding_model value before returning the store, so reusing one key for incompatible descriptor geometry or semantics fails before a caller can replace rows.

Source

pub fn embedder_dimensions(&self, model_name: &str) -> Option<usize>

Output dimensions for a named embedding model, resolved from the embedder registry alone — no storage access. Mirrors vectors_for_model’s resolution order: lattice aliases route through the enum when registered, otherwise the custom provider’s declared dimensions(). None when no such model is registered.

Source

pub fn text(&self, token: &NamespaceToken) -> RuntimeResult<Arc<dyn TextSearch>>

Get a TextSearch index for the entity corpus (single shared table).

Source

pub fn text_for_notes( &self, token: &NamespaceToken, ) -> RuntimeResult<Arc<dyn TextSearch>>

Get a TextSearch index for the notes corpus (single shared table).

Source

pub fn authorize(&self, ns: Namespace) -> RuntimeResult<NamespaceToken>

Mint an authorization token for the given namespace.

Consults the configured crate::Gate before minting. With the default AllowAllGate this always succeeds. When a real policy-backed gate is installed, this method enforces it and returns PermissionDenied on denial.

The returned token’s read visibility set defaults to [ns] — identical to the pre-visibility-set behaviour. Use Self::authorize_with_visibility to mint a token that can read additional namespaces.

When actor_id is configured in RuntimeConfig, the token carries that actor label so that comm.inbox filters by to_actor. When unconfigured, the token carries ActorRef::anonymous() and inbox falls back to party-line behavior.

Source

pub fn authorize_with_visibility( &self, primary: Namespace, extra_visible: Vec<Namespace>, ) -> RuntimeResult<NamespaceToken>

Mint an authorization token with an explicit read-visibility set.

primary is the write namespace — all records created via the returned token land there. extra_visible lists additional namespaces the token may read. The primary is always included in the visible set regardless of extra_visible.

Usage (lambda:leo reading both leo and khive namespaces):

ⓘ
let tok = rt.authorize_with_visibility(
    Namespace::parse("lambda:leo").unwrap(),
    vec![Namespace::parse("lambda:khive").unwrap()],
)?;
Source

pub fn install_edge_rules(&self, rules: Vec<EdgeEndpointRule>)

Install the pack-aggregated edge endpoint rules.

Called by the transport layer after the VerbRegistry is built so that runtime-layer edge validation can consult pack rules. Idempotent: later calls overwrite the previous rule set.

Source

pub fn install_blob_hydrator( &self, hydrator: Arc<BlobHydrator>, ) -> RuntimeResult<()>

Install an already-paired blob hydrator into this runtime.

Reinstalling the exact same Arc is idempotent. A different pair is rejected: replacing it would split or reset the aggregate admission budget while requests may still hold leases.

Source

pub fn install_shared_blob_hydrator( &self, hydrator: Arc<BlobHydrator>, ) -> RuntimeResult<()>

Install a boot-shared hydrator whose mutability is governed by the blob runtime’s own mode, not this handle’s domain-store mode.

ADR-160 D3 installs one hydrator Arc on every runtime handle a boot produces, and the documented multi-backend matrix includes a writable blob secondary beside a read-only main: there the shared hydrator is legitimately writable on a read-only domain handle. The mode decision must therefore already be encoded in the hydrator, and it must have been DERIVED, not declared: only hydrators built through crate::BlobHydrator::resolve_for_governing_backend — whose mode comes from the governing backend’s own access mode — are accepted here. A hand-paired hydrator (BlobHydrator::new / for_mode) is refused so a safe downstream caller cannot use this seam to put a writable store on a read-only runtime; such callers use Self::install_blob_hydrator, which holds hydrator mode against this runtime’s own.

This gate is a wrong-wiring guard, not an in-process sandbox: which backend governs is the boot host’s topology assertion, and a caller who deliberately selects an unrelated writable backend as governing is outside what any runtime seam can enforce (see the trust-model note on crate::BlobHydrator::resolve_for_governing_backend).

Source

pub fn install_blob_store(&self, store: Arc<dyn BlobStore>) -> RuntimeResult<()>

Pair and install a store using this runtime’s resolved hydration budget.

Boot paths that own multiple runtimes should instead construct one crate::BlobHydrator and call Self::install_blob_hydrator with the same Arc on every handle.

Source

pub fn blob_hydrator(&self) -> Option<Arc<BlobHydrator>>

Return the installed shared blob hydrator, if boot configured one.

Source

pub fn blob_store(&self) -> Option<Arc<dyn BlobStore>>

Return the installed BlobStore, if the boot path resolved and installed one. None when no [storage.blob] selection was ever installed — e.g. a bare/test runtime constructed without going through the khive-mcp boot path.

Source

pub fn install_kind_registry( &self, entity_kinds: Vec<String>, note_kinds: Vec<String>, )

Install the pack-aggregated valid entity and note kinds.

Called by the transport layer after the VerbRegistry is built so that runtime-layer entity/note creation and import validate kind strings against the merged pack vocabulary. Idempotent: later calls overwrite previous sets.

When no kinds are installed (empty lists), kind validation is skipped at the runtime layer. The pack handler layer remains the primary enforcement point; this provides defense-in-depth for direct Rust callers and import.

Source

pub fn install_pack_owned_note_kinds(&self, kinds: Vec<String>)

Install the pack-owned note kinds aggregated from the pack registry.

Called by the transport after the VerbRegistry is built, same timing as install_kind_registry.

Source

pub fn is_pack_owned_note_kind(&self, kind: &str) -> bool

Whether kind is a note kind owned by a pack (see install_pack_owned_note_kinds).

Always false before the transport installs the list — a bare runtime has no packs, so no kind is pack-owned there.

Source

pub fn install_entity_type_validator(&self, f: EntityTypeValidatorFn)

Install a pack-supplied entity-type validator.

Called by the KgPack during registration so that create_many can validate entity_type values at the runtime layer, closing the hole where direct Rust callers bypass the handler-layer validate_entity_type check.

The callback receives (kind, entity_type) and returns the normalised type string, or RuntimeError::InvalidInput if the type is not registered for that kind. Passing entity_type = None must return Ok(None).

Source

pub fn install_note_mutation_hook(&self, f: NoteMutationHookFn)

Install a pack-owned note-mutation hook.

Overwrites any previously-installed hook, same single-slot semantics as install_entity_type_validator. In practice only one pack (khive-pack-memory) installs one today; if a second pack ever needs this, the slot should be widened to a Vec at that point rather than silently overwritten.

Source

pub fn install_entity_kind_hooks(&self, hooks: EntityKindHooks)

Install the pack-aggregated entity-kind update hooks (issue #2943).

Called by the transport after the VerbRegistry is built, same timing as install_kind_registry — pass registry.entity_kind_hooks(). Idempotent: a later call replaces the set.

Source

pub fn register_fusion_strategy( &self, name: impl Into<String>, executor: Arc<dyn FusionExecutor>, )

Install a pack-owned note-write validator.

Called during pack registration (PackRuntime::register_note_write_validator) so that the covered note-write sites carrying caller-supplied properties derive the owning pack’s identity properties from the authorization token, closing the gap where a direct Rust caller, the generic create verb, or the proposal-apply path (which dispatches no pack hooks) writes them unchecked. Single-slot semantics, same as install_note_mutation_hook: a second installing pack overwrites the first, so a validator must return kinds it does not own unchanged.

Covered sites — each calls derive_note_write_properties before the write: create_note_inner (operations.rs, the generic create verb funnel and every other public create_note* variant), atomic_prepare::prepare_add_note (the proposal-apply add-note path), and atomic_message::create_notes_atomic_with_report (the atomic multi-note writer).

NOT covered by this validator: try_create_note (operations.rs). try_create_note is deliberately excluded — its only caller path is comm.ingest, where properties.from_actor is the external transport sender named by the from parameter, not the authenticated caller, and where transport-owned quarantine/channel properties are legitimately established. Running the generic validator there would stamp every inbound message as the ingesting daemon and reject the evidence the trusted ingest handler just derived. try_create_note instead runs its own narrower reserved-transport-property check inline (operations.rs’s try_create_note_impl), which allows the three message-kind transport properties only when called through Self::try_create_note_as_trusted_ingest with a crate::pack::ChannelIngestCapability.

The NoteStore returned by notes is covered by a different, narrower mechanism: it is wrapped in note_store_guard::PolicyEnforcingNoteStore, which refuses upsert_note / upsert_notes / try_insert_note / replace_note_if_unchanged calls that would write a kind = "message" note carrying quarantined / channel_kind / channel_slug, and refuses set_note_property / try_patch_note_property / patch_note_property_atomic / update_note_properties calls that would patch any of those keys onto any note — unconditionally, since that public accessor has no way to see a trust decision. try_create_note_impl itself reaches storage through Self::raw_notes, the unwrapped accessor, so its own inline check (which can legitimately allow those properties for trusted ingest) is not double-enforced or contradicted by the wrapper. Register a pack-defined custom fusion strategy under name (ADR-012).

Unlike install_entity_type_validator/install_note_mutation_hook, this slot is keyed rather than single-occupancy: multiple packs each register their own named strategy, and a second registration under an already-used name replaces the first. Looked up by FusionStrategy::Custom { name, .. } at the hybrid-search dispatch boundary in crate::fusion; an unregistered name fails closed with RuntimeError::UnknownFusionStrategy rather than silently falling back to RRF.

Source

pub fn install_note_write_validator(&self, f: NoteWriteValidatorFn)

Source

pub fn has_note_write_validator(&self) -> bool

Whether a note-write validator is installed on this runtime.

Exists so a transport’s own tests can assert, per boot path, that the documented startup sequence actually filled the slot. A missing install fails open and silently — an empty slot passes caller-supplied properties straight through, which no write site can distinguish from a validator that approved them — so occupancy is asserted, never assumed.

Source

pub fn pack_edge_rules(&self) -> Vec<EdgeEndpointRule>

Snapshot of currently-installed pack edge rules.

This is the same composed rule set validate_edge_relation_endpoints consults via pack_rule_allows when accepting/rejecting an edge. Public so pack-layer error-hint code (e.g. khive-pack-kg’s valid_relations_for_entity_pair) can derive hints from the exact source the validator uses, rather than maintaining a separate hand-authored table that can drift out of sync.

Source

pub fn default_embedder_name(&self) -> &str

Return the name of the default embedding model (empty string if none configured).

Source

pub fn resolve_embedding_model( &self, name: Option<&str>, ) -> RuntimeResult<EmbeddingModel>

Resolve a model name (or None for the default) to an EmbeddingModel.

Returns UnknownModel if the name is not in the registry, or Unconfigured if None is passed and no default model is set.

Source

pub fn registered_embedding_model_names(&self) -> Vec<String>

Names of all registered embedding models in this runtime.

Includes both built-in lattice models and any custom embedders registered by packs via register_embedder. Useful for operations that must touch every model’s storage (e.g., scoped vector deletion on note delete). The default model is included.

Source

pub async fn embedder( &self, name: &str, ) -> RuntimeResult<Arc<dyn EmbeddingService>>

Get the lazily-initialized embedding service for the named model.

Accepts both built-in lattice model names (e.g. "all-minilm-l6-v2", "paraphrase") and custom provider names registered via register_embedder.

For lattice model names, aliases (e.g. "paraphrase") are resolved to their canonical key before looking up the registry. For custom providers the name must match exactly as supplied during registration.

First call for any name loads the underlying service (cold start cost); subsequent calls are cheap (registry caches the Arc).

Source

pub fn register_embedder(&self, provider: impl EmbedderProvider + 'static)

Register a custom embedding provider with this runtime.

The provider is added to the shared EmbedderRegistry so all clones of this runtime see the new provider immediately. If a provider with the same name already exists it is replaced (last-writer wins — see crate::EmbedderRegistry::register for the rationale).

Packs should call this from crate::PackRuntime::register_embedders (the hook is invoked by the transport during pack initialisation, before the first verb dispatch).

Source

pub async fn list_embedding_models( &self, engine_filter: Option<&str>, ) -> RuntimeResult<Vec<EmbeddingModelRegistryRecord>>

List registered embedding models via SqlAccess, routing through the existing connection pool rather than opening a fresh Connection per call.

Optionally filter by engine_name. Returns an empty vec when the _embedding_models table does not yet exist (e.g. no migrations have run or no models have been registered). All other SQL errors are propagated.

Source§

impl KhiveRuntime

Source

pub async fn stream_append( &self, token: &NamespaceToken, stream: &str, record: &Value, expected_seq: Option<i64>, note_kind: &str, tags: Option<Vec<String>>, fence: Option<NoteFences>, embed: Option<bool>, embedding_model: Option<String>, registry: &VerbRegistry, ) -> RuntimeResult<Value>

Append a JSON value as an immutable note. The sequence precondition and every note/index/ledger statement share the same writer transaction.

Source

pub async fn stream_append_with_outcome( &self, token: &NamespaceToken, stream: &str, record: &Value, expected_seq: Option<i64>, note_kind: &str, tags: Option<Vec<String>>, fence: Option<NoteFences>, embed: Option<bool>, embedding_model: Option<String>, registry: &VerbRegistry, ) -> Result<Value, StreamAppendFailure>

Append once while retaining evidence about the failed append’s own committedness. The original error remains available for stop semantics or the shared structured error projection; this method never retries.

Source

pub async fn stream_batch_atomic( &self, token: &NamespaceToken, members: Vec<StreamBatchMember>, fence: Option<NoteFence>, observed: Vec<StreamObservation>, registry: &VerbRegistry, ) -> RuntimeResult<Result<Vec<Value>, StreamBatchRefusal>>

All members share one transaction. Batch predicates precede member DML; effects become executable only after the complete unit commits.

Source

pub async fn stream_batch_per_member( &self, token: &NamespaceToken, members: Vec<StreamBatchMember>, registry: &VerbRegistry, ) -> RuntimeResult<Vec<Value>>

One writer transaction per member, in list order. Every member is prepared before the first write; a member’s refusal is returned as its value and its siblings stand, so numbers on one stream increase with list position but another writer’s append may fall between them.

Source

pub async fn stream_read( &self, token: &NamespaceToken, stream: &str, after: i64, limit: i64, ) -> RuntimeResult<Value>

Read one ordered page and its head from one SQL snapshot.

Source

pub async fn stream_stat( &self, token: &NamespaceToken, stream: &str, ) -> RuntimeResult<Value>

Independently count entries and read the head in the same statement.

Trait Implementations§

Source§

impl Clone for KhiveRuntime

Source§

fn clone(&self) -> Self

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more