use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use khive_storage::{
BatchWriteSummary, BoundedCount, DeleteMode, Note, NoteFilter, NoteStore, NoteVisibility, Page,
PageRequest, SeekCursor, SeekPage, StorageCapability, StorageError, StorageResult,
};
use serde_json::Value;
use uuid::Uuid;
use crate::curation::kind_owned_properties;
fn transport_owned_message_property_named_in(
properties: &serde_json::Map<String, Value>,
) -> Option<&'static str> {
kind_owned_properties("message")
.iter()
.copied()
.find(|key| properties.contains_key(*key))
}
fn reject_if_forged_message_note(note: &Note, operation: &'static str) -> StorageResult<()> {
if note.kind != "message" {
return Ok(());
}
let Some(key) = note
.properties
.as_ref()
.and_then(Value::as_object)
.and_then(transport_owned_message_property_named_in)
else {
return Ok(());
};
Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: format!(
"`{key}` is transport-owned on a `message` note and cannot be written through the \
public NoteStore accessor; only the trusted channel-ingest path may establish \
quarantine disposition and channel provenance"
),
})
}
fn has_web_receipt_provenance(note: &Note) -> bool {
note.properties.as_ref().is_some_and(|properties| {
properties
.as_object()
.is_some_and(|map| map.contains_key(crate::secret_gate::RESERVED_WEB_RECEIPT_KEY))
})
}
fn web_receipt_write_refused(operation: &'static str) -> StorageError {
StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: "web receipt provenance and its receipt body are web-pack-owned; generic note writes cannot create or alter them".into(),
}
}
fn reject_forged_web_receipt_note(note: &Note, operation: &'static str) -> StorageResult<()> {
if has_web_receipt_provenance(note) {
return Err(web_receipt_write_refused(operation));
}
Ok(())
}
fn reject_reserved_patch_target(target: &str, operation: &'static str) -> StorageResult<()> {
let body = target.strip_prefix('$').unwrap_or(target);
let body = body.strip_prefix('.').unwrap_or(body);
let first_segment = &body[..body.find(['.', '[']).unwrap_or(body.len())];
let is_bare_identifier = !first_segment.is_empty()
&& !first_segment.starts_with(|c: char| c.is_ascii_digit())
&& first_segment
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_');
if !is_bare_identifier {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: format!(
"property target `{target}` is not a bare top-level identifier; the public \
NoteStore accessor refuses target spellings it cannot prove distinct from \
the transport-owned message properties"
),
});
}
if !kind_owned_properties("message").contains(&first_segment) {
return Ok(());
}
Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: format!(
"`{first_segment}` is transport-owned and cannot be patched through the public \
NoteStore accessor; only the trusted channel-ingest path may establish quarantine \
disposition and channel provenance"
),
})
}
fn reject_reserved_replacement_properties(
properties: Option<&Value>,
operation: &'static str,
) -> StorageResult<()> {
if properties
.and_then(Value::as_object)
.is_some_and(|map| map.contains_key(crate::secret_gate::RESERVED_WEB_RECEIPT_KEY))
{
return Err(web_receipt_write_refused(operation));
}
let Some(key) = properties
.and_then(Value::as_object)
.and_then(transport_owned_message_property_named_in)
else {
return Ok(());
};
Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: format!(
"`{key}` is transport-owned and cannot be written through the public NoteStore \
accessor; only the trusted channel-ingest path may establish quarantine disposition \
and channel provenance"
),
})
}
fn reject_reserved_note_properties(
properties: Option<&Value>,
operation: &'static str,
) -> StorageResult<()> {
crate::secret_gate::reject_reserved_secret_gate_property(properties).map_err(|error| {
StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: error.to_string(),
}
})
}
fn reject_existing_secret_gate_property(
existing: Option<&Note>,
operation: &'static str,
) -> StorageResult<()> {
if let Some(existing) = existing {
reject_reserved_note_properties(existing.properties.as_ref(), operation)?;
}
Ok(())
}
fn reject_changed_identity_properties(
existing: &Note,
properties: Option<&Value>,
operation: &'static str,
) -> StorageResult<()> {
if existing.kind == "message" {
return Ok(());
}
for key in kind_owned_properties(&existing.kind) {
let before = existing
.properties
.as_ref()
.and_then(|value| value.get(*key));
let after = properties.and_then(|value| value.get(*key));
if before != after {
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: format!(
"`{key}` is kind-owned on a `{}` note and cannot be changed through the \
public NoteStore accessor",
existing.kind
),
});
}
}
Ok(())
}
fn reject_changed_identity_note(
existing: &Note,
note: &Note,
operation: &'static str,
) -> StorageResult<()> {
if has_web_receipt_provenance(existing) {
return Err(web_receipt_write_refused(operation));
}
if existing.kind != "message"
&& !kind_owned_properties(&existing.kind).is_empty()
&& existing.kind != note.kind
{
return Err(StorageError::InvalidInput {
capability: StorageCapability::Notes,
operation: operation.into(),
message: format!(
"the kind of a `{}` note with kind-owned identity cannot be changed through \
the public NoteStore accessor",
existing.kind
),
});
}
reject_changed_identity_properties(existing, note.properties.as_ref(), operation)
}
pub(crate) struct PolicyEnforcingNoteStore {
inner: Arc<dyn NoteStore>,
}
impl PolicyEnforcingNoteStore {
pub(crate) fn wrap(inner: Arc<dyn NoteStore>) -> Arc<dyn NoteStore> {
Arc::new(Self { inner })
}
async fn reject_identity_change(
&self,
note: &Note,
operation: &'static str,
) -> StorageResult<()> {
if let Some(existing) = self.inner.get_note_including_deleted(note.id).await? {
reject_reserved_note_properties(existing.properties.as_ref(), operation)?;
reject_changed_identity_note(&existing, note, operation)?;
}
Ok(())
}
async fn reject_web_receipt_mutation(
&self,
id: Uuid,
operation: &'static str,
) -> StorageResult<()> {
if self
.inner
.get_note_including_deleted(id)
.await?
.as_ref()
.is_some_and(has_web_receipt_provenance)
{
return Err(web_receipt_write_refused(operation));
}
Ok(())
}
}
#[async_trait]
impl NoteStore for PolicyEnforcingNoteStore {
async fn get_live_notes_by_key(
&self,
namespace: &str,
key: &str,
kind: Option<&str>,
) -> StorageResult<Vec<Note>> {
self.inner.get_live_notes_by_key(namespace, key, kind).await
}
async fn query_keyed_notes(
&self,
namespace: &str,
filter: &NoteFilter,
prefix: &str,
after: Option<&khive_storage::note::NoteKeyCursor>,
page: PageRequest,
) -> StorageResult<(Vec<Note>, Option<khive_storage::note::NoteKeyCursor>)> {
self.inner
.query_keyed_notes(namespace, filter, prefix, after, page)
.await
}
async fn upsert_note(&self, note: Note) -> StorageResult<()> {
reject_reserved_note_properties(note.properties.as_ref(), "upsert_note")?;
reject_if_forged_message_note(¬e, "upsert_note")?;
reject_forged_web_receipt_note(¬e, "upsert_note")?;
self.reject_identity_change(¬e, "upsert_note").await?;
self.inner.upsert_note(note).await
}
async fn insert_note_if_absent(&self, note: Note) -> StorageResult<bool> {
reject_reserved_note_properties(note.properties.as_ref(), "insert_note_if_absent")?;
reject_if_forged_message_note(¬e, "insert_note_if_absent")?;
reject_forged_web_receipt_note(¬e, "insert_note_if_absent")?;
self.inner.insert_note_if_absent(note).await
}
async fn replace_note_if_unchanged(
&self,
note: Note,
expected_updated_at: i64,
expected_deleted_at: Option<i64>,
) -> StorageResult<bool> {
reject_reserved_note_properties(note.properties.as_ref(), "replace_note_if_unchanged")?;
reject_if_forged_message_note(¬e, "replace_note_if_unchanged")?;
if !has_web_receipt_provenance(¬e) {
self.reject_web_receipt_mutation(note.id, "replace_note_if_unchanged")
.await?;
}
reject_forged_web_receipt_note(¬e, "replace_note_if_unchanged")?;
self.reject_identity_change(¬e, "replace_note_if_unchanged")
.await?;
self.inner
.replace_note_if_unchanged(note, expected_updated_at, expected_deleted_at)
.await
}
async fn upsert_notes(&self, notes: Vec<Note>) -> StorageResult<BatchWriteSummary> {
for note in ¬es {
reject_reserved_note_properties(note.properties.as_ref(), "upsert_notes")?;
reject_if_forged_message_note(note, "upsert_notes")?;
reject_forged_web_receipt_note(note, "upsert_notes")?;
}
{
let mut preceding: HashMap<Uuid, &Note> = HashMap::new();
for note in ¬es {
if let Some(existing) = preceding.get(¬e.id) {
reject_changed_identity_note(existing, note, "upsert_notes")?;
} else {
self.reject_identity_change(note, "upsert_notes").await?;
}
preceding.insert(note.id, note);
}
}
self.inner.upsert_notes(notes).await
}
async fn get_note(&self, id: Uuid) -> StorageResult<Option<Note>> {
self.inner.get_note(id).await
}
async fn get_note_including_deleted(&self, id: Uuid) -> StorageResult<Option<Note>> {
self.inner.get_note_including_deleted(id).await
}
async fn delete_note(&self, id: Uuid, mode: DeleteMode) -> StorageResult<bool> {
self.inner.delete_note(id, mode).await
}
async fn update_note_properties(
&self,
id: Uuid,
properties: Option<Value>,
updated_at: i64,
) -> StorageResult<bool> {
reject_reserved_note_properties(properties.as_ref(), "update_note_properties")?;
reject_reserved_replacement_properties(properties.as_ref(), "update_note_properties")?;
self.reject_web_receipt_mutation(id, "update_note_properties")
.await?;
if let Some(existing) = self.inner.get_note_including_deleted(id).await? {
reject_reserved_note_properties(
existing.properties.as_ref(),
"update_note_properties",
)?;
reject_changed_identity_properties(
&existing,
properties.as_ref(),
"update_note_properties",
)?;
}
self.inner
.update_note_properties(id, properties, updated_at)
.await
}
async fn set_note_property(
&self,
id: Uuid,
key: &str,
value: Value,
updated_at: i64,
) -> StorageResult<bool> {
reject_reserved_patch_target(key, "set_note_property")?;
let existing = self.inner.get_note_including_deleted(id).await?;
if existing.as_ref().is_some_and(has_web_receipt_provenance) {
return Err(web_receipt_write_refused("set_note_property"));
}
reject_existing_secret_gate_property(existing.as_ref(), "set_note_property")?;
self.inner
.set_note_property(id, key, value, updated_at)
.await
}
async fn try_patch_note_property(
&self,
id: Uuid,
namespace: &str,
filter: &NoteFilter,
json_path: &str,
value: Value,
updated_at: i64,
) -> StorageResult<bool> {
reject_reserved_patch_target(json_path, "try_patch_note_property")?;
let existing = self.inner.get_note_including_deleted(id).await?;
if existing.as_ref().is_some_and(has_web_receipt_provenance) {
return Err(web_receipt_write_refused("try_patch_note_property"));
}
reject_existing_secret_gate_property(existing.as_ref(), "try_patch_note_property")?;
self.inner
.try_patch_note_property(id, namespace, filter, json_path, value, updated_at)
.await
}
async fn patch_note_property_atomic(
&self,
ids: Vec<Uuid>,
namespace: &str,
filter: &NoteFilter,
json_path: &str,
value: Value,
updated_at: i64,
) -> StorageResult<()> {
reject_reserved_patch_target(json_path, "patch_note_property_atomic")?;
let mut secret_gate_refusal = None;
for window in ids.chunks(128) {
match self.inner.get_notes_batch_including_deleted(window).await {
Ok(notes) => {
let notes: HashMap<Uuid, Note> =
notes.into_iter().map(|note| (note.id, note)).collect();
for id in window {
let existing = notes.get(id);
if existing.is_some_and(has_web_receipt_provenance) {
return Err(web_receipt_write_refused("patch_note_property_atomic"));
}
if secret_gate_refusal.is_none() {
secret_gate_refusal = reject_existing_secret_gate_property(
existing,
"patch_note_property_atomic",
)
.err();
}
}
}
Err(_) => {
for id in window {
let existing = self.inner.get_note_including_deleted(*id).await?;
let existing = existing.as_ref();
if existing.is_some_and(has_web_receipt_provenance) {
return Err(web_receipt_write_refused("patch_note_property_atomic"));
}
if secret_gate_refusal.is_none() {
secret_gate_refusal = reject_existing_secret_gate_property(
existing,
"patch_note_property_atomic",
)
.err();
}
}
}
}
}
if let Some(refusal) = secret_gate_refusal {
return Err(refusal);
}
self.inner
.patch_note_property_atomic(ids, namespace, filter, json_path, value, updated_at)
.await
}
async fn query_notes(
&self,
namespace: &str,
kind: Option<&str>,
page: PageRequest,
) -> StorageResult<Page<Note>> {
self.inner.query_notes(namespace, kind, page).await
}
async fn query_notes_count_free(
&self,
namespace: &str,
kind: Option<&str>,
page: PageRequest,
) -> StorageResult<Page<Note>> {
self.inner
.query_notes_count_free(namespace, kind, page)
.await
}
async fn query_notes_filtered(
&self,
namespace: &str,
filter: &NoteFilter,
page: PageRequest,
) -> StorageResult<Page<Note>> {
self.inner
.query_notes_filtered(namespace, filter, page)
.await
}
async fn query_notes_filtered_count_free(
&self,
namespace: &str,
filter: &NoteFilter,
page: PageRequest,
) -> StorageResult<Page<Note>> {
self.inner
.query_notes_filtered_count_free(namespace, filter, page)
.await
}
async fn count_notes_filtered_in_snapshot(
&self,
namespace: &str,
filters: &[NoteFilter],
) -> StorageResult<Vec<u64>> {
self.inner
.count_notes_filtered_in_snapshot(namespace, filters)
.await
}
async fn count_notes_filtered_bounded_in_snapshot(
&self,
namespace: &str,
filters: &[NoteFilter],
cap: u32,
) -> StorageResult<Vec<BoundedCount>> {
self.inner
.count_notes_filtered_bounded_in_snapshot(namespace, filters, cap)
.await
}
async fn note_sequence(&self, id: Uuid) -> StorageResult<Option<i64>> {
self.inner.note_sequence(id).await
}
async fn query_notes_filtered_after(
&self,
namespace: &str,
filter: &NoteFilter,
after: Option<SeekCursor>,
limit: u32,
) -> StorageResult<SeekPage<Note>> {
self.inner
.query_notes_filtered_after(namespace, filter, after, limit)
.await
}
async fn query_notes_filtered_bounded(
&self,
namespace: &str,
filter: &NoteFilter,
max_rows: u32,
) -> StorageResult<Vec<Note>> {
self.inner
.query_notes_filtered_bounded(namespace, filter, max_rows)
.await
}
async fn count_notes(&self, namespace: &str, kind: Option<&str>) -> StorageResult<u64> {
self.inner.count_notes(namespace, kind).await
}
async fn count_notes_in_namespaces(
&self,
namespaces: &[String],
kind: Option<&str>,
) -> StorageResult<u64> {
self.inner.count_notes_in_namespaces(namespaces, kind).await
}
async fn try_insert_note(&self, note: Note) -> StorageResult<bool> {
reject_reserved_note_properties(note.properties.as_ref(), "try_insert_note")?;
reject_if_forged_message_note(¬e, "try_insert_note")?;
reject_forged_web_receipt_note(¬e, "try_insert_note")?;
self.inner.try_insert_note(note).await
}
async fn get_notes_batch(&self, ids: &[Uuid]) -> StorageResult<Vec<Note>> {
self.inner.get_notes_batch(ids).await
}
async fn get_notes_batch_including_deleted(&self, ids: &[Uuid]) -> StorageResult<Vec<Note>> {
self.inner.get_notes_batch_including_deleted(ids).await
}
async fn get_note_visibility_batch(&self, ids: &[Uuid]) -> StorageResult<Vec<NoteVisibility>> {
self.inner.get_note_visibility_batch(ids).await
}
}
#[cfg(test)]
mod tests;