#[cfg(doc)]
use super::VerbRegistry;
use super::{
json_type_name, Arc, AuditEvent, AuditObligationFailure, Event, EventKind, EventOutcome,
EventStore, GateDecision, GateRequest, Namespace, RuntimeError, SubstrateKind, Value,
};
fn target_id_from_args(verb: &str, args: &serde_json::Value) -> Option<uuid::Uuid> {
let alias = args.get("target").filter(|_| verb == "link");
args.get("target_id")
.or(alias)
.and_then(serde_json::Value::as_str)
.and_then(|s| s.parse::<uuid::Uuid>().ok())
}
pub(super) fn masked_audit_event(
gate_req: &GateRequest,
decision: &GateDecision,
gate_impl: &str,
) -> AuditEvent {
let mut audit = AuditEvent::from_check(gate_req, decision, gate_impl)
.with_operation_attribution(
khive_storage::operation_context::current_operation_attribution(),
);
if let Some(reason) = audit.deny_reason.take() {
audit.deny_reason = Some(crate::secret_gate::bounded_masked_log_text(&reason));
}
audit
}
pub(super) fn build_audit_storage_event(
gate_req: &GateRequest,
audit: &AuditEvent,
outcome: EventOutcome,
resource: Option<Value>,
) -> Event {
let mut audit_data = serde_json::to_value(audit).unwrap_or_else(|e| {
tracing::warn!(error = %e, "failed to serialize AuditEvent for EventStore");
serde_json::Value::Null
});
if let Some(resource) = resource {
if let Value::Object(ref mut map) = audit_data {
map.insert("resource".to_string(), resource);
}
}
let mut storage_event = Event::new(
gate_req.namespace.as_str(),
gate_req.verb.as_str(),
EventKind::Audit,
SubstrateKind::Event,
format!("{}:{}", gate_req.actor.kind, gate_req.actor.id),
)
.with_outcome(outcome)
.with_payload(audit_data);
storage_event.op_index = audit.op_index;
storage_event.ref_resolution = audit.ref_resolution;
if let Some(target_id) = target_id_from_args(&gate_req.verb, &gate_req.args) {
storage_event = storage_event.with_target(target_id);
}
storage_event
}
static AUDIT_APPEND_FAILURES: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
pub(crate) fn audit_append_failure_count() -> u64 {
AUDIT_APPEND_FAILURES.load(std::sync::atomic::Ordering::Relaxed)
}
static AUDIT_OBLIGATION_APPEND_FAILURES: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);
pub(crate) fn audit_obligation_append_failure_count() -> u64 {
AUDIT_OBLIGATION_APPEND_FAILURES.load(std::sync::atomic::Ordering::Relaxed)
}
static AUDIT_ADMISSION_REFUSED_OBLIGATIONS: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);
static AUDIT_ADMISSION_REFUSED_OBLIGATIONS_LAST_MS: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);
pub fn audit_admission_refused_obligation_count() -> u64 {
AUDIT_ADMISSION_REFUSED_OBLIGATIONS.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn audit_admission_refused_obligation_last_at_ms() -> Option<u64> {
match AUDIT_ADMISSION_REFUSED_OBLIGATIONS_LAST_MS.load(std::sync::atomic::Ordering::Relaxed) {
0 => None,
at => Some(at),
}
}
static AUDIT_ADMISSION_UNRESOLVED_OBLIGATIONS: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);
static AUDIT_ADMISSION_UNRESOLVED_OBLIGATIONS_LAST_MS: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);
pub fn audit_admission_unresolved_obligation_count() -> u64 {
AUDIT_ADMISSION_UNRESOLVED_OBLIGATIONS.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn audit_admission_unresolved_obligation_last_at_ms() -> Option<u64> {
match AUDIT_ADMISSION_UNRESOLVED_OBLIGATIONS_LAST_MS.load(std::sync::atomic::Ordering::Relaxed)
{
0 => None,
at => Some(at),
}
}
fn mark_admission_obligation_counter(mark: &std::sync::atomic::AtomicU64) {
let now = chrono::Utc::now().timestamp_millis();
let now = u64::try_from(now).unwrap_or(1).max(1);
mark.store(now, std::sync::atomic::Ordering::Relaxed);
}
const GIT_DIGEST_RECEIPT_FAILURE: &str =
"git_digest_receipt_persist_failed: git.digest writes may have committed, but no durable \
success receipt was confirmed; inspect ingest state before retrying";
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum GitDigestReceiptOutcome {
Persisted,
BuildRejected,
PersistenceUnavailable,
}
fn fail_git_digest_receipt(
result: &mut Result<Value, RuntimeError>,
failure: AuditObligationFailure,
) {
let Ok(value) = result else {
return;
};
let domain_result = std::mem::take(value);
*result = Err(RuntimeError::AuditObligation {
failure: Box::new(failure),
domain_result,
});
}
pub(super) async fn persist_git_digest_receipt(
store: Option<&Arc<dyn EventStore>>,
audit_batch: Option<&Arc<crate::audit_batch::AuditBatch>>,
gate_req: &GateRequest,
audit: Option<&AuditEvent>,
result: &mut Result<Value, RuntimeError>,
duration_us: i64,
resource: Option<Value>,
) -> GitDigestReceiptOutcome {
let Ok(report) = result else {
return GitDigestReceiptOutcome::PersistenceUnavailable;
};
let Some(store) = store else {
tracing::error!(
verb = "git.digest",
"durable receipt store is not configured"
);
fail_git_digest_receipt(
result,
AuditObligationFailure::git_digest_receipt("event store is not configured"),
);
return GitDigestReceiptOutcome::PersistenceUnavailable;
};
let Some(audit) = audit else {
tracing::error!(
verb = "git.digest",
"durable receipt cannot be built because the gate produced no audit decision"
);
fail_git_digest_receipt(
result,
AuditObligationFailure::git_digest_receipt("gate audit decision is absent"),
);
return GitDigestReceiptOutcome::PersistenceUnavailable;
};
let Some(report_object) = report.as_object_mut() else {
tracing::error!(
verb = "git.digest",
"digest handler returned a non-object report"
);
fail_git_digest_receipt(
result,
AuditObligationFailure::git_digest_receipt("handler report is not an object"),
);
return GitDigestReceiptOutcome::BuildRejected;
};
let Some(project_id) = report_object
.get("project_id")
.and_then(Value::as_str)
.and_then(|raw| raw.parse::<uuid::Uuid>().ok())
else {
tracing::error!(
verb = "git.digest",
"digest handler report omitted a valid project_id"
);
fail_git_digest_receipt(
result,
AuditObligationFailure::git_digest_receipt("handler report has no valid project_id"),
);
return GitDigestReceiptOutcome::BuildRejected;
};
let mut event = Event::new(
gate_req.namespace.as_str(),
gate_req.verb.as_str(),
EventKind::Audit,
SubstrateKind::Event,
format!("{}:{}", gate_req.actor.kind, gate_req.actor.id),
)
.with_outcome(EventOutcome::Success)
.with_target(project_id)
.with_payload_schema_version(2)
.with_duration_us(duration_us);
let receipt_id = event.id;
report_object.insert(
"receipt_id".to_string(),
Value::String(receipt_id.to_string()),
);
let mut payload = serde_json::to_value(audit).unwrap_or_else(|serialize_err| {
tracing::error!(
verb = "git.digest",
error = %serialize_err,
"failed to serialize gate audit for durable digest receipt"
);
Value::Null
});
let Value::Object(payload_object) = &mut payload else {
tracing::error!(
verb = "git.digest",
"gate audit serialization did not produce an object"
);
fail_git_digest_receipt(
result,
AuditObligationFailure::git_digest_receipt("gate audit payload is not an object"),
);
return GitDigestReceiptOutcome::BuildRejected;
};
if let Some(resource) = resource {
payload_object.insert("resource".to_string(), resource);
}
payload_object.insert("result".to_string(), report.clone());
event.payload = payload;
let submit_result = if let Some(audit_batch) = audit_batch {
audit_batch
.submit_until_resolved(crate::audit_batch::PreparedAuditRow {
event,
producer: crate::audit_batch::AuditProducer::GitDigestReceipt,
})
.await
.map(|_outcome| ())
.map_err(|reason| AuditObligationFailure::new("git.digest", reason))
} else {
store
.append_event(event)
.await
.map_err(|error| AuditObligationFailure::from_store("git.digest", error))
};
if let Err(mut failure) = submit_result {
AUDIT_OBLIGATION_APPEND_FAILURES.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::error!(
verb = "git.digest",
error = %failure,
receipt_id = %receipt_id,
"durable digest receipt append failed"
);
failure.message = format!(
"{GIT_DIGEST_RECEIPT_FAILURE}; audit submission failed ({})",
failure.wire_code()
);
fail_git_digest_receipt(result, failure);
return GitDigestReceiptOutcome::PersistenceUnavailable;
}
GitDigestReceiptOutcome::Persisted
}
pub(super) async fn append_audit_event_best_effort(
audit_batch: Option<&Arc<crate::audit_batch::AuditBatch>>,
store: &Arc<dyn EventStore>,
event: Event,
verb: &str,
producer: crate::audit_batch::AuditProducer,
degrade_allowlisted: bool,
) -> Result<(), AuditObligationFailure> {
use crate::audit_batch::{
classify, AuditBatchControl, AuditProducer, AuditProductionClass, AuditTerminalReason,
};
let is_obligation = classify(producer) == AuditProductionClass::DispatchObligation;
let admission_degrade_eligible =
degrade_allowlisted && producer == AuditProducer::DispatchSucceeded;
let enqueued_row_outlives_deadline = producer == AuditProducer::DispatchSucceeded;
if let Some(audit_batch) = audit_batch {
let row = crate::audit_batch::PreparedAuditRow { event, producer };
let submit_result =
if producer == AuditProducer::DispatchSucceeded && !admission_degrade_eligible {
audit_batch.submit_until_resolved(row).await
} else {
audit_batch.submit(row).await
};
if let Err(reason) = submit_result {
if is_obligation {
if enqueued_row_outlives_deadline
&& reason == AuditTerminalReason::AdmissionDeadlineExpired
{
AUDIT_ADMISSION_UNRESOLVED_OBLIGATIONS
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
mark_admission_obligation_counter(
&AUDIT_ADMISSION_UNRESOLVED_OBLIGATIONS_LAST_MS,
);
tracing::warn!(
verb,
reason = ?reason,
degrade_allowlisted,
"audit obligation row was still enqueued and unresolved when \
the caller's admission wait deadline elapsed; its generation \
commits it independently of this response. Dispatch reports \
its own committed result (non-fatal)"
);
return Ok(());
}
if admission_degrade_eligible
&& reason == AuditTerminalReason::QueueAdmissionExhausted
{
AUDIT_ADMISSION_REFUSED_OBLIGATIONS
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
mark_admission_obligation_counter(&AUDIT_ADMISSION_REFUSED_OBLIGATIONS_LAST_MS);
tracing::warn!(
verb,
reason = ?reason,
"read verb's audit obligation row was refused before \
enqueue under audit-lane admission pressure; dispatch \
still reports its own result (non-fatal)"
);
return Ok(());
}
AUDIT_OBLIGATION_APPEND_FAILURES.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::error!(
verb,
reason = ?reason,
"audit obligation batch submission failed; failing dispatch"
);
return Err(AuditObligationFailure::new(verb, reason));
}
AUDIT_APPEND_FAILURES.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::warn!(
verb,
reason = ?reason,
"audit event batch submission failed (non-fatal)"
);
}
return Ok(());
}
if let Err(store_err) = store.append_event(event).await {
if is_obligation {
AUDIT_OBLIGATION_APPEND_FAILURES.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::error!(
verb,
error = %store_err,
"audit obligation store write failed; failing dispatch"
);
return Err(AuditObligationFailure::from_store(verb, store_err));
}
AUDIT_APPEND_FAILURES.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tracing::warn!(
verb,
error = %store_err,
"audit event store write failed (non-fatal)"
);
}
Ok(())
}
pub(super) fn fold_audit_obligation<T>(
result: Result<T, RuntimeError>,
audit_outcome: Result<(), AuditObligationFailure>,
domain_value: impl FnOnce(T) -> Value,
) -> Result<T, RuntimeError> {
match (result, audit_outcome) {
(Ok(value), Ok(())) => Ok(value),
(Ok(value), Err(failure)) => Err(RuntimeError::AuditObligation {
failure: Box::new(failure),
domain_result: domain_value(value),
}),
(Err(err), _) => Err(err),
}
}
#[derive(Debug, Clone, serde::Serialize)]
struct LinkAuditSuccessV2 {
#[serde(flatten)]
audit: AuditEvent,
edge_id: uuid::Uuid,
source_id: uuid::Uuid,
target_id: uuid::Uuid,
relation: String,
weight: f64,
}
pub(super) fn link_audit_success_from_result(
audit: AuditEvent,
result: &serde_json::Value,
) -> Option<(uuid::Uuid, serde_json::Value)> {
let edge_id = result.get("id")?.as_str()?.parse::<uuid::Uuid>().ok()?;
let source_id = result
.get("source_id")?
.as_str()?
.parse::<uuid::Uuid>()
.ok()?;
let target_id = result
.get("target_id")?
.as_str()?
.parse::<uuid::Uuid>()
.ok()?;
let relation = result.get("relation")?.as_str()?.to_string();
let weight = result.get("weight")?.as_f64()?;
let enriched = LinkAuditSuccessV2 {
audit,
edge_id,
source_id,
target_id,
relation,
weight,
};
let payload = serde_json::to_value(&enriched).ok()?;
Some((edge_id, payload))
}
pub fn resolve_explicit_namespace(
params: &Value,
default_namespace: &str,
) -> Result<Namespace, RuntimeError> {
match params.get("namespace") {
None => Namespace::parse(default_namespace)
.map_err(|e| RuntimeError::InvalidInput(format!("invalid namespace: {e}"))),
Some(Value::String(ns_str)) => Namespace::parse(ns_str)
.map_err(|e| RuntimeError::InvalidInput(format!("invalid namespace {ns_str:?}: {e}"))),
Some(other) => Err(RuntimeError::InvalidInput(format!(
"invalid namespace: expected string when present, got {}",
json_type_name(other),
))),
}
}