use super::*;
pub fn compute_decision_basis_hash(gotchas: &[(&str, &serde_json::Value)]) -> String {
let mut ordered: Vec<&(&str, &serde_json::Value)> = gotchas.iter().collect();
ordered.sort_by(|a, b| a.0.cmp(b.0));
let mut hasher = Sha256::new();
for (key, record_json) in ordered {
hasher.update(key.as_bytes());
let rule = record_json
.pointer("/value")
.and_then(|v| v.as_str())
.unwrap_or("");
hasher.update(rule.as_bytes());
let conf = record_json
.pointer("/confidence/value")
.and_then(|v| v.as_f64())
.unwrap_or(0.0);
hasher.update(format!("{conf}").as_bytes());
}
format!("{:x}", hasher.finalize())
}
static ENFORCEMENT_WRITERS: OnceLock<
StdMutex<HashMap<PathBuf, Arc<tokio::sync::Mutex<EnforcementEventWriter>>>>,
> = OnceLock::new();
pub(crate) async fn shared_writer(
store: &Store,
) -> Result<Arc<tokio::sync::Mutex<EnforcementEventWriter>>> {
let registry = ENFORCEMENT_WRITERS.get_or_init(|| StdMutex::new(HashMap::new()));
if let Some(writer) = registry
.lock()
.expect("enforcement writer registry poisoned")
.get(&store.root)
.cloned()
{
return Ok(writer);
}
let writer = Arc::new(tokio::sync::Mutex::new(
EnforcementEventWriter::new(store).await?,
));
Ok(registry
.lock()
.expect("enforcement writer registry poisoned")
.entry(store.root.clone())
.or_insert(writer)
.clone())
}
#[allow(clippy::too_many_arguments)]
pub async fn record_event(
store: &Store,
event_type: EnforcementEventType,
subject_kind: SubjectKind,
subject_key: String,
agent_type: String,
receipt_id: Option<String>,
decision_reason_code: String,
decision_basis_hash: Option<String>,
) -> Result<Option<EnforcementEvent>> {
record_event_with_session(
store,
event_type,
subject_kind,
subject_key,
agent_type,
receipt_id,
decision_reason_code,
decision_basis_hash,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn record_event_with_session(
store: &Store,
event_type: EnforcementEventType,
subject_kind: SubjectKind,
subject_key: String,
agent_type: String,
receipt_id: Option<String>,
decision_reason_code: String,
decision_basis_hash: Option<String>,
agent_session: Option<String>,
) -> Result<Option<EnforcementEvent>> {
record_event_with_lineage(
store,
event_type,
subject_kind,
subject_key,
agent_type,
receipt_id,
decision_reason_code,
decision_basis_hash,
agent_session,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn record_event_with_lineage(
store: &Store,
event_type: EnforcementEventType,
subject_kind: SubjectKind,
subject_key: String,
agent_type: String,
receipt_id: Option<String>,
decision_reason_code: String,
decision_basis_hash: Option<String>,
agent_session: Option<String>,
agent_id: Option<String>,
) -> Result<Option<EnforcementEvent>> {
record_event_with_nested_lineage(
store,
event_type,
subject_kind,
subject_key,
agent_type,
receipt_id,
decision_reason_code,
decision_basis_hash,
agent_session,
agent_id,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn record_event_with_nested_lineage(
store: &Store,
event_type: EnforcementEventType,
subject_kind: SubjectKind,
subject_key: String,
agent_type: String,
receipt_id: Option<String>,
decision_reason_code: String,
decision_basis_hash: Option<String>,
agent_session: Option<String>,
agent_id: Option<String>,
parent_agent_id: Option<String>,
) -> Result<Option<EnforcementEvent>> {
let mode = get_enforcement_mode(store).await;
let result = async {
let writer = shared_writer(store).await?;
let mut writer = writer.lock().await;
writer.agent_session = agent_session;
writer.agent_id = agent_id;
writer.parent_agent_id = parent_agent_id;
writer
.write(
store,
event_type,
subject_kind,
subject_key,
agent_type,
receipt_id,
decision_reason_code,
decision_basis_hash,
)
.await
}
.await;
match result {
Ok(event) => Ok(Some(event)),
Err(e) => match mode {
EnforcementMode::Advisory => {
tracing::warn!("enforcement event write failed (advisory mode): {e}");
Ok(None)
}
EnforcementMode::Strict => Err(e),
},
}
}