use std::collections::HashMap;
use anyhow::Context;
use serde_json::json;
use serde_json::Value;
use tokio::fs;
use uuid::Uuid;
use crate::audit::append_protected_audit_event;
use crate::automation_v2::governance::*;
use crate::{now_ms, AppState};
const GOVERNANCE_AUDIT_EVENT_PREFIX: &str = "automation.governance";
fn bind_governance_record_to_tenant(
record: &mut AutomationGovernanceRecord,
tenant_context: &tandem_types::TenantContext,
) -> anyhow::Result<bool> {
let scope = governance_tenant_scope(tenant_context);
let mut changed = false;
match record.tenant_context.as_ref() {
Some(owner) if !governance_tenant_matches(owner, &scope) => {
anyhow::bail!("automation governance tenant ownership mismatch");
}
Some(_) => {}
None => {
record.tenant_context = Some(scope.clone());
changed = true;
}
}
for grant in record
.modify_grants
.iter_mut()
.chain(record.capability_grants.iter_mut())
{
match grant.tenant_context.as_ref() {
Some(owner) if !governance_tenant_matches(owner, &scope) => {
anyhow::bail!("automation governance grant tenant ownership mismatch");
}
Some(_) => {}
None => {
grant.tenant_context = Some(scope.clone());
changed = true;
}
}
}
Ok(changed)
}
fn governance_record_owned_by(
record: &AutomationGovernanceRecord,
tenant_context: &tandem_types::TenantContext,
) -> bool {
let record_matches = record
.tenant_context
.as_ref()
.is_some_and(|owner| governance_tenant_matches(owner, tenant_context))
|| (record.tenant_context.is_none() && tenant_context.is_local_implicit());
record_matches
&& record
.modify_grants
.iter()
.chain(record.capability_grants.iter())
.all(|grant| {
grant
.tenant_context
.as_ref()
.is_some_and(|owner| governance_tenant_matches(owner, tenant_context))
|| (grant.tenant_context.is_none() && tenant_context.is_local_implicit())
})
}
fn governance_agent_scope_key(
tenant_context: &tandem_types::TenantContext,
agent_id: &str,
) -> String {
let agent_id = agent_id.trim();
if tenant_context.is_local_implicit() {
return agent_id.to_string();
}
let deployment_id = tenant_context.deployment_id.as_deref().unwrap_or("");
format!(
"tenant:{}",
serde_json::to_string(&(
tenant_context.org_id.as_str(),
tenant_context.workspace_id.as_str(),
deployment_id,
agent_id
))
.unwrap_or_else(|_| agent_id.to_string())
)
}
fn scoped_agent_id_maps(
tenant_context: &tandem_types::TenantContext,
agent_ids: &[String],
) -> (
Vec<String>,
HashMap<String, String>,
HashMap<String, String>,
) {
let mut scoped_ids = Vec::with_capacity(agent_ids.len());
let mut raw_to_scoped = HashMap::new();
let mut scoped_to_raw = HashMap::new();
for agent_id in agent_ids {
let scoped_id = governance_agent_scope_key(tenant_context, agent_id);
scoped_ids.push(scoped_id.clone());
raw_to_scoped.insert(agent_id.clone(), scoped_id.clone());
scoped_to_raw.insert(scoped_id, agent_id.clone());
}
(scoped_ids, raw_to_scoped, scoped_to_raw)
}
fn scoped_agent_id_for_storage(agent_id: &str, raw_to_scoped: &HashMap<String, String>) -> String {
raw_to_scoped
.get(agent_id)
.cloned()
.unwrap_or_else(|| agent_id.to_string())
}
fn display_agent_id(agent_id: &str, scoped_to_raw: &HashMap<String, String>) -> String {
scoped_to_raw
.get(agent_id)
.cloned()
.unwrap_or_else(|| agent_id.to_string())
}
pub(super) fn approval_receipt_matches_tenant(
approval: &GovernanceApprovalRequest,
caller: &tandem_types::TenantContext,
) -> bool {
let local = tandem_types::TenantContext::local_implicit();
let owner = approval.tenant_context.as_ref().unwrap_or(&local);
owner.org_id == caller.org_id
&& owner.workspace_id == caller.workspace_id
&& owner.deployment_id == caller.deployment_id
}
#[derive(Default)]
pub struct UnavailableGovernanceEngine;
impl GovernancePolicyEngine for UnavailableGovernanceEngine {
fn premium_enabled(&self) -> bool {
false
}
fn authorize_create(
&self,
_snapshot: &GovernanceContextSnapshot,
actor: &GovernanceActorRef,
_provenance: &AutomationProvenanceRecord,
_declared_capabilities: &AutomationDeclaredCapabilities,
_now_ms: u64,
) -> Result<(), GovernanceError> {
if actor.kind == GovernanceActorKind::Human {
return Ok(());
}
Err(GovernanceError::feature_unavailable(
"premium governance is required for agent-authored automation creation",
))
}
fn authorize_capability_escalation(
&self,
_snapshot: &GovernanceContextSnapshot,
actor: &GovernanceActorRef,
_previous: &AutomationDeclaredCapabilities,
_next: &AutomationDeclaredCapabilities,
_now_ms: u64,
) -> Result<(), GovernanceError> {
if actor.kind == GovernanceActorKind::Human {
return Ok(());
}
Err(GovernanceError::feature_unavailable(
"premium governance is required for agent capability escalation",
))
}
fn authorize_mutation(
&self,
record: &AutomationGovernanceRecord,
actor: &GovernanceActorRef,
_destructive: bool,
) -> Result<(), GovernanceError> {
if actor.kind != GovernanceActorKind::Human {
return Err(GovernanceError::feature_unavailable(
"premium governance is required for agent-owned automation mutation",
));
}
let owner = &record.provenance.creator;
if owner.kind == GovernanceActorKind::Human {
if let (Some(owner_id), Some(actor_id)) =
(owner.actor_id.as_deref(), actor.actor_id.as_deref())
{
let owner_id = owner_id.trim();
let actor_id = actor_id.trim();
if !owner_id.is_empty()
&& !actor_id.is_empty()
&& !owner_id.eq_ignore_ascii_case(actor_id)
{
return Err(GovernanceError::forbidden(
"AUTOMATION_V2_NOT_OWNER",
"only the automation owner may mutate it in this build",
));
}
}
}
Ok(())
}
fn create_approval_request(
&self,
_snapshot: &GovernanceContextSnapshot,
_input: GovernanceApprovalDraftInput,
_now_ms: u64,
) -> Result<GovernanceApprovalRequest, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance approval flows are not available in this build",
))
}
fn decide_approval_request(
&self,
_existing: &GovernanceApprovalRequest,
_reviewer: GovernanceActorRef,
_approved: bool,
_notes: Option<String>,
_now_ms: u64,
) -> Result<GovernanceApprovalRequest, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance approval flows are not available in this build",
))
}
fn evaluate_creation_review_progress(
&self,
_snapshot: &GovernanceContextSnapshot,
_agent_id: &str,
_automation_id: &str,
_now_ms: u64,
) -> Result<GovernanceCreationReviewEvaluation, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance review tracking is not available in this build",
))
}
fn evaluate_run_review_progress(
&self,
_snapshot: &GovernanceContextSnapshot,
_automation_id: &str,
_reason: AutomationLifecycleReviewKind,
_run_id: Option<String>,
_detail: Option<String>,
_now_ms: u64,
) -> Result<Option<GovernanceAutomationReviewEvaluation>, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance review tracking is not available in this build",
))
}
fn evaluate_dependency_revocation(
&self,
_snapshot: &GovernanceContextSnapshot,
_input: GovernanceDependencyRevocationInput,
_now_ms: u64,
) -> Result<GovernanceAutomationReviewEvaluation, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance dependency revocation is not available in this build",
))
}
fn health_check_run_window(&self, _limits: &GovernanceLimits) -> usize {
0
}
fn summarize_run_health(
&self,
_runs: &[GovernanceRunHealthObservation],
) -> GovernanceRunHealthSummary {
GovernanceRunHealthSummary::default()
}
fn default_automation_expiry(
&self,
limits: &GovernanceLimits,
provenance: &AutomationProvenanceRecord,
now_ms: u64,
) -> Option<u64> {
if provenance.creator.kind != GovernanceActorKind::Agent {
return None;
}
if limits.default_expires_after_ms == 0 {
return None;
}
Some(now_ms.saturating_add(limits.default_expires_after_ms))
}
fn acknowledge_creation_review(
&self,
existing: Option<AgentCreationReviewSummary>,
agent_id: &str,
notes: Option<String>,
now_ms: u64,
) -> AgentCreationReviewSummary {
let mut summary = existing
.unwrap_or_else(|| AgentCreationReviewSummary::new(agent_id.to_string(), now_ms));
summary.created_since_review = 0;
summary.review_required = false;
summary.review_kind = None;
summary.review_requested_at_ms = None;
summary.review_request_id = None;
summary.last_reviewed_at_ms = Some(now_ms);
summary.last_review_notes = notes;
summary.updated_at_ms = now_ms;
summary
}
fn acknowledge_automation_review(
&self,
record: &AutomationGovernanceRecord,
now_ms: u64,
) -> AutomationGovernanceRecord {
let mut record = record.clone();
record.review_required = false;
record.review_kind = None;
record.review_requested_at_ms = None;
record.review_request_id = None;
record.last_reviewed_at_ms = Some(now_ms);
record.runs_since_review = 0;
record.health_findings.clear();
record.health_last_checked_at_ms = Some(now_ms);
record.updated_at_ms = now_ms;
record
}
fn evaluate_health_check(
&self,
_snapshot: &GovernanceContextSnapshot,
_input: GovernanceHealthCheckInput,
_now_ms: u64,
) -> Result<Option<GovernanceHealthCheckEvaluation>, GovernanceError> {
Ok(None)
}
fn evaluate_retirement(
&self,
_input: GovernanceRetirementInput,
_now_ms: u64,
) -> Result<AutomationGovernanceRecord, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance retirement logic is not available in this build",
))
}
fn evaluate_retirement_extension(
&self,
_input: GovernanceRetirementExtensionInput,
_now_ms: u64,
) -> Result<AutomationGovernanceRecord, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance retirement logic is not available in this build",
))
}
fn evaluate_spend_usage(
&self,
_snapshot: &GovernanceContextSnapshot,
_input: &GovernanceSpendInput,
_now_ms: u64,
) -> Result<GovernanceSpendEvaluation, GovernanceError> {
Err(GovernanceError::feature_unavailable(
"premium governance spend tracking is not available in this build",
))
}
}
fn default_human_provenance(
creator_id: Option<String>,
source: impl Into<String>,
) -> AutomationProvenanceRecord {
AutomationProvenanceRecord::human(creator_id, source)
}
fn declared_capabilities_for_automation(
automation: &crate::AutomationV2Spec,
) -> AutomationDeclaredCapabilities {
AutomationDeclaredCapabilities::from_metadata(automation.metadata.as_ref())
}
fn governance_run_health_status(status: &crate::AutomationRunStatus) -> GovernanceRunHealthStatus {
match status {
crate::AutomationRunStatus::Queued => GovernanceRunHealthStatus::Queued,
crate::AutomationRunStatus::Running => GovernanceRunHealthStatus::Running,
crate::AutomationRunStatus::Pausing => GovernanceRunHealthStatus::Pausing,
crate::AutomationRunStatus::Paused => GovernanceRunHealthStatus::Paused,
crate::AutomationRunStatus::AwaitingApproval => GovernanceRunHealthStatus::AwaitingApproval,
crate::AutomationRunStatus::Completed => GovernanceRunHealthStatus::Completed,
crate::AutomationRunStatus::Blocked => GovernanceRunHealthStatus::Blocked,
crate::AutomationRunStatus::Failed => GovernanceRunHealthStatus::Failed,
crate::AutomationRunStatus::Cancelled => GovernanceRunHealthStatus::Cancelled,
}
}
fn new_governance_record(
automation_id: String,
provenance: AutomationProvenanceRecord,
declared_capabilities: AutomationDeclaredCapabilities,
created_at_ms: u64,
updated_at_ms: u64,
) -> AutomationGovernanceRecord {
AutomationGovernanceRecord {
automation_id,
tenant_context: None,
provenance,
declared_capabilities,
modify_grants: Vec::new(),
capability_grants: Vec::new(),
created_at_ms,
updated_at_ms,
deleted_at_ms: None,
delete_retention_until_ms: None,
published_externally: false,
creation_paused: false,
review_required: false,
review_kind: None,
review_requested_at_ms: None,
review_request_id: None,
last_reviewed_at_ms: None,
runs_since_review: 0,
expires_at_ms: None,
expired_at_ms: None,
retired_at_ms: None,
retire_reason: None,
paused_for_lifecycle: false,
health_last_checked_at_ms: None,
health_findings: Vec::new(),
}
}
impl AppState {
pub fn premium_governance_enabled(&self) -> bool {
self.governance_engine.premium_enabled()
}
fn governance_snapshot(&self, state: &GovernanceState) -> GovernanceContextSnapshot {
state.snapshot()
}
pub async fn load_automation_governance(&self) -> anyhow::Result<()> {
if !self.automation_governance_path.exists() {
return Ok(());
}
let raw = fs::read_to_string(&self.automation_governance_path).await?;
let parsed = serde_json::from_str::<GovernanceState>(&raw).map_err(|error| {
anyhow::anyhow!("failed to parse automation governance state: {error}")
})?;
*self.automation_governance.write().await = parsed;
Ok(())
}
pub async fn persist_automation_governance(&self) -> anyhow::Result<()> {
let guard = self.automation_governance.read().await;
self.persist_automation_governance_snapshot(&guard).await
}
async fn persist_automation_governance_snapshot(
&self,
governance: &GovernanceState,
) -> anyhow::Result<()> {
if let Some(parent) = self.automation_governance_path.parent() {
fs::create_dir_all(parent).await?;
}
let payload = serde_json::to_string_pretty(governance)?;
fs::write(&self.automation_governance_path, payload).await?;
Ok(())
}
async fn persist_automation_governance_locked(&self) -> anyhow::Result<()> {
self.persist_automation_governance().await
}
pub async fn bootstrap_automation_governance(&self) -> anyhow::Result<usize> {
let automations = self.list_automations_v2().await;
let now = now_ms();
let mut changed = 0usize;
let previous = self.automation_governance.read().await.clone();
let mut quarantined_automations = Vec::new();
{
let mut guard = self.automation_governance.write().await;
for mut automation in automations {
let tenant_context = automation.tenant_context();
if let Some(record) = guard.records.get_mut(&automation.automation_id) {
match bind_governance_record_to_tenant(record, &tenant_context) {
Ok(true) => {
record.updated_at_ms = now;
changed += 1;
}
Ok(false) => {}
Err(_) => {
quarantine_governance_record_for_tenant(record, &tenant_context, now);
automation.status = crate::AutomationV2Status::Paused;
quarantined_automations.push(automation);
changed += 1;
}
}
continue;
}
let mut record = new_governance_record(
automation.automation_id.clone(),
default_human_provenance(
Some(automation.creator_id.clone()),
"migration_or_legacy_default",
),
declared_capabilities_for_automation(&automation),
automation.created_at_ms.max(now),
now,
);
bind_governance_record_to_tenant(&mut record, &tenant_context)?;
guard
.records
.insert(automation.automation_id.clone(), record);
changed += 1;
}
if changed > 0 {
guard.updated_at_ms = now;
}
}
for automation in quarantined_automations {
if let Err(error) = self.put_automation_v2(automation).await {
*self.automation_governance.write().await = previous;
return Err(error);
}
}
if changed > 0 {
if let Err(error) = self.persist_automation_governance().await {
*self.automation_governance.write().await = previous;
return Err(error);
}
}
Ok(changed)
}
pub async fn get_automation_governance(
&self,
automation_id: &str,
) -> Option<AutomationGovernanceRecord> {
self.automation_governance
.read()
.await
.records
.get(automation_id)
.cloned()
}
async fn require_active_automation_governance_tenant(
&self,
automation_id: &str,
tenant_context: &tandem_types::TenantContext,
) -> anyhow::Result<crate::AutomationV2Spec> {
let automation = self
.get_automation_v2(automation_id)
.await
.ok_or_else(|| anyhow::anyhow!("automation not found"))?;
if !governance_tenant_matches(&automation.tenant_context(), tenant_context) {
anyhow::bail!("automation not found");
}
let record = self
.get_automation_governance(automation_id)
.await
.ok_or_else(|| anyhow::anyhow!("automation governance record not found"))?;
if !governance_record_owned_by(&record, tenant_context) {
anyhow::bail!("automation governance record not found");
}
Ok(automation)
}
pub async fn get_automation_governance_for_tenant(
&self,
automation_id: &str,
tenant_context: &tandem_types::TenantContext,
) -> Option<AutomationGovernanceRecord> {
self.require_active_automation_governance_tenant(automation_id, tenant_context)
.await
.ok()?;
self.get_automation_governance(automation_id).await
}
pub async fn get_or_bootstrap_automation_governance(
&self,
automation: &crate::AutomationV2Spec,
) -> AutomationGovernanceRecord {
let tenant_context = automation.tenant_context();
if let Some(mut record) = self
.get_automation_governance(&automation.automation_id)
.await
{
match bind_governance_record_to_tenant(&mut record, &tenant_context) {
Ok(true) => {
if self
.upsert_automation_governance(record.clone())
.await
.is_err()
{
record.creation_paused = true;
record.paused_for_lifecycle = true;
record.review_required = true;
}
}
Ok(false) => {}
Err(_) => {
quarantine_governance_record_for_tenant(&mut record, &tenant_context, now_ms());
let _ = self.upsert_automation_governance(record.clone()).await;
let mut quarantined = automation.clone();
quarantined.status = crate::AutomationV2Status::Paused;
let _ = self.put_automation_v2(quarantined).await;
}
}
return record;
}
let mut record = new_governance_record(
automation.automation_id.clone(),
default_human_provenance(Some(automation.creator_id.clone()), "legacy_default"),
declared_capabilities_for_automation(automation),
automation.created_at_ms,
now_ms(),
);
if bind_governance_record_to_tenant(&mut record, &tenant_context).is_err()
|| self
.upsert_automation_governance(record.clone())
.await
.is_err()
{
record.creation_paused = true;
record.paused_for_lifecycle = true;
record.review_required = true;
}
record
}
pub async fn upsert_automation_governance(
&self,
mut record: AutomationGovernanceRecord,
) -> anyhow::Result<AutomationGovernanceRecord> {
if record.automation_id.trim().is_empty() {
anyhow::bail!("automation_id is required");
}
let now = now_ms();
if record.created_at_ms == 0 {
record.created_at_ms = now;
}
record.updated_at_ms = now;
let tenant_context = record
.tenant_context
.clone()
.unwrap_or_else(tandem_types::TenantContext::local_implicit);
bind_governance_record_to_tenant(&mut record, &tenant_context)?;
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.record.updated"),
&tenant_context,
record
.provenance
.creator
.actor_id
.clone()
.or_else(|| record.provenance.creator.source.clone()),
json!({
"automationID": record.automation_id,
"provenance": record.provenance,
"declaredCapabilities": record.declared_capabilities,
"publishedExternally": record.published_externally,
"creationPaused": record.creation_paused,
}),
)
.await?;
let previous = {
let mut guard = self.automation_governance.write().await;
let previous = guard.clone();
guard
.records
.insert(record.automation_id.clone(), record.clone());
guard.updated_at_ms = now;
previous
};
if let Err(error) = self.persist_automation_governance().await {
*self.automation_governance.write().await = previous;
return Err(error);
}
Ok(record)
}
pub async fn set_automation_governance_provenance(
&self,
automation_id: &str,
provenance: AutomationProvenanceRecord,
) -> anyhow::Result<AutomationGovernanceRecord> {
let mut record = self
.get_automation_governance(automation_id)
.await
.unwrap_or_else(|| {
new_governance_record(
automation_id.to_string(),
provenance.clone(),
AutomationDeclaredCapabilities::default(),
now_ms(),
now_ms(),
)
});
record.provenance = provenance;
if record.expires_at_ms.is_none() {
let limits = self.automation_governance.read().await.limits.clone();
record.expires_at_ms = self.governance_engine.default_automation_expiry(
&limits,
&record.provenance,
now_ms(),
);
}
let stored = self.upsert_automation_governance(record).await?;
if let Some(agent_id) = stored
.provenance
.creator
.actor_id
.as_deref()
.filter(|_| stored.provenance.creator.kind == GovernanceActorKind::Agent)
{
let _ = self
.record_agent_creation_review_progress(agent_id, &stored.automation_id)
.await;
}
Ok(stored)
}
pub async fn sync_automation_governance_from_spec(
&self,
automation: &crate::AutomationV2Spec,
provenance: Option<AutomationProvenanceRecord>,
) -> anyhow::Result<AutomationGovernanceRecord> {
let now = now_ms();
let tenant_context = automation.tenant_context();
let mut record = self
.get_automation_governance(&automation.automation_id)
.await
.unwrap_or_else(|| {
new_governance_record(
automation.automation_id.clone(),
provenance.clone().unwrap_or_else(|| {
default_human_provenance(
Some(automation.creator_id.clone()),
"sync_default",
)
}),
declared_capabilities_for_automation(automation),
automation.created_at_ms,
now,
)
});
if let Some(provenance) = provenance {
record.provenance = provenance;
}
bind_governance_record_to_tenant(&mut record, &tenant_context)?;
record.declared_capabilities = declared_capabilities_for_automation(automation);
if record.created_at_ms == 0 {
record.created_at_ms = automation.created_at_ms;
}
record.updated_at_ms = now;
let record = self.upsert_automation_governance(record).await?;
if let Some(agent_id) = record
.provenance
.creator
.actor_id
.as_deref()
.filter(|_| record.provenance.creator.kind == GovernanceActorKind::Agent)
{
let _ = self
.record_agent_creation_review_progress(agent_id, &record.automation_id)
.await;
}
Ok(record)
}
pub async fn pause_automation_creation_for_agent(
&self,
agent_id: &str,
paused: bool,
) -> anyhow::Result<()> {
let mut guard = self.automation_governance.write().await;
if paused {
if !guard.paused_agents.iter().any(|value| value == agent_id) {
guard.paused_agents.push(agent_id.to_string());
}
} else {
guard.paused_agents.retain(|value| value != agent_id);
}
guard.updated_at_ms = now_ms();
drop(guard);
self.persist_automation_governance().await?;
Ok(())
}
pub async fn can_create_automation_for_actor(
&self,
tenant_context: &tandem_types::TenantContext,
actor: &GovernanceActorRef,
provenance: &AutomationProvenanceRecord,
declared_capabilities: &AutomationDeclaredCapabilities,
) -> Result<(), GovernanceError> {
let snapshot = {
let guard = self.automation_governance.read().await;
let mut snapshot = self.governance_snapshot(&guard);
if !tenant_context.is_local_implicit() && actor.kind == GovernanceActorKind::Agent {
if let Some(agent_id) = actor
.actor_id
.as_deref()
.filter(|value| !value.trim().is_empty())
{
let scoped_agent_id = governance_agent_scope_key(tenant_context, agent_id);
if guard.is_agent_spend_paused(&scoped_agent_id)
&& !guard.has_approved_agent_quota_override(&scoped_agent_id)
{
return Err(GovernanceError::too_many_requests(
"AUTOMATION_V2_AGENT_SPEND_CAP_EXCEEDED",
"this agent is paused after reaching its spend cap",
));
}
snapshot
.spend_paused_agents
.retain(|paused_agent_id| paused_agent_id != agent_id);
}
}
snapshot
};
self.governance_engine.authorize_create(
&snapshot,
actor,
provenance,
declared_capabilities,
now_ms(),
)
}
pub async fn can_escalate_declared_capabilities(
&self,
actor: &GovernanceActorRef,
previous: &AutomationDeclaredCapabilities,
next: &AutomationDeclaredCapabilities,
) -> Result<(), GovernanceError> {
let snapshot = {
let guard = self.automation_governance.read().await;
self.governance_snapshot(&guard)
};
self.governance_engine.authorize_capability_escalation(
&snapshot,
actor,
previous,
next,
now_ms(),
)
}
pub async fn record_automation_creation(
&self,
automation: &crate::AutomationV2Spec,
provenance: AutomationProvenanceRecord,
) -> anyhow::Result<AutomationGovernanceRecord> {
let mut record = new_governance_record(
automation.automation_id.clone(),
provenance,
declared_capabilities_for_automation(automation),
automation.created_at_ms,
now_ms(),
);
bind_governance_record_to_tenant(&mut record, &automation.tenant_context())?;
if record.expires_at_ms.is_none() {
let limits = self.automation_governance.read().await.limits.clone();
record.expires_at_ms = self.governance_engine.default_automation_expiry(
&limits,
&record.provenance,
now_ms(),
);
}
let stored = self.upsert_automation_governance(record).await?;
if let Some(agent_id) = stored
.provenance
.creator
.actor_id
.as_deref()
.filter(|_| stored.provenance.creator.kind == GovernanceActorKind::Agent)
{
let _ = self
.record_agent_creation_review_progress(agent_id, &stored.automation_id)
.await;
}
Ok(stored)
}
pub async fn grant_automation_modify_access(
&self,
automation_id: &str,
granted_to: GovernanceActorRef,
granted_by: GovernanceActorRef,
reason: Option<String>,
tenant_context: &tandem_types::TenantContext,
) -> anyhow::Result<AutomationGrantRecord> {
self.require_active_automation_governance_tenant(automation_id, tenant_context)
.await?;
let grant = AutomationGrantRecord {
grant_id: format!("grant-{}", Uuid::new_v4()),
automation_id: automation_id.to_string(),
tenant_context: Some(governance_tenant_scope(tenant_context)),
grant_kind: AutomationGrantKind::Modify,
granted_to,
granted_by,
capability_key: None,
created_at_ms: now_ms(),
revoked_at_ms: None,
revoke_reason: reason,
};
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.grant.created"),
tenant_context,
grant
.granted_by
.actor_id
.clone()
.or_else(|| grant.granted_by.source.clone()),
json!({
"automationID": automation_id,
"grant": grant,
}),
)
.await?;
self.require_active_automation_governance_tenant(automation_id, tenant_context)
.await?;
let previous = {
let mut guard = self.automation_governance.write().await;
let previous = guard.clone();
let Some(record) = guard.records.get_mut(automation_id) else {
anyhow::bail!("automation governance record not found");
};
if !governance_record_owned_by(record, tenant_context) {
anyhow::bail!("automation governance record not found");
}
record.modify_grants.push(grant.clone());
record.updated_at_ms = now_ms();
guard.updated_at_ms = now_ms();
previous
};
if let Err(error) = self.persist_automation_governance().await {
*self.automation_governance.write().await = previous;
return Err(error);
}
Ok(grant)
}
pub async fn revoke_automation_modify_access(
&self,
automation_id: &str,
grant_id: &str,
revoked_by: GovernanceActorRef,
reason: Option<String>,
tenant_context: &tandem_types::TenantContext,
) -> anyhow::Result<Option<AutomationGrantRecord>> {
self.require_active_automation_governance_tenant(automation_id, tenant_context)
.await?;
let Some(mut stored) = self
.automation_governance
.read()
.await
.records
.get(automation_id)
.and_then(|record| {
record
.modify_grants
.iter()
.find(|grant| grant.grant_id == grant_id)
.cloned()
})
else {
return Ok(None);
};
if stored.revoked_at_ms.is_some() {
return Ok(Some(stored));
}
stored.revoked_at_ms = Some(now_ms());
stored.revoke_reason = reason.clone();
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.grant.revoked"),
tenant_context,
revoked_by
.actor_id
.clone()
.or_else(|| revoked_by.source.clone()),
json!({
"automationID": automation_id,
"grantID": grant_id,
"reason": reason,
}),
)
.await?;
self.require_active_automation_governance_tenant(automation_id, tenant_context)
.await?;
let previous = {
let mut guard = self.automation_governance.write().await;
let previous = guard.clone();
let Some(record) = guard.records.get_mut(automation_id) else {
anyhow::bail!("automation governance record not found");
};
if !governance_record_owned_by(record, tenant_context) {
anyhow::bail!("automation governance record not found");
}
let Some(grant) = record
.modify_grants
.iter_mut()
.find(|grant| grant.grant_id == grant_id && grant.revoked_at_ms.is_none())
else {
return Ok(None);
};
*grant = stored.clone();
record.updated_at_ms = now_ms();
guard.updated_at_ms = now_ms();
previous
};
if let Err(error) = self.persist_automation_governance().await {
*self.automation_governance.write().await = previous;
return Err(error);
}
Ok(Some(stored))
}
pub async fn list_approval_requests(
&self,
request_type: Option<GovernanceApprovalRequestType>,
status: Option<GovernanceApprovalStatus>,
) -> Vec<GovernanceApprovalRequest> {
let mut rows = self
.automation_governance
.read()
.await
.approvals
.values()
.filter(|request| {
request_type
.map(|value| request.request_type == value)
.unwrap_or(true)
&& status.map(|value| request.status == value).unwrap_or(true)
})
.cloned()
.collect::<Vec<_>>();
rows.sort_by(|a, b| b.updated_at_ms.cmp(&a.updated_at_ms));
rows
}
pub async fn get_governance_approval_request(
&self,
approval_id: &str,
) -> Option<GovernanceApprovalRequest> {
self.automation_governance
.read()
.await
.approvals
.get(approval_id)
.cloned()
}
pub async fn get_governance_approval_request_for_tenant(
&self,
approval_id: &str,
tenant_context: &tandem_types::TenantContext,
) -> Option<GovernanceApprovalRequest> {
self.get_governance_approval_request(approval_id)
.await
.filter(|request| approval_receipt_matches_tenant(request, tenant_context))
}
pub async fn list_approval_requests_for_tenant(
&self,
request_type: Option<GovernanceApprovalRequestType>,
status: Option<GovernanceApprovalStatus>,
tenant_context: &tandem_types::TenantContext,
) -> Vec<GovernanceApprovalRequest> {
self.list_approval_requests(request_type, status)
.await
.into_iter()
.filter(|request| approval_receipt_matches_tenant(request, tenant_context))
.collect()
}
pub async fn delete_automation_v2_with_governance(
&self,
automation_id: &str,
deleted_by: GovernanceActorRef,
) -> anyhow::Result<Option<crate::AutomationV2Spec>> {
let _guard = self.automations_v2_persistence.lock().await;
let removed = self.automations_v2.write().await.remove(automation_id);
if let Some(automation) = removed.clone() {
let now = now_ms();
{
let mut governance = self.automation_governance.write().await;
let record = governance
.records
.entry(automation_id.to_string())
.or_insert_with(|| {
new_governance_record(
automation_id.to_string(),
default_human_provenance(
Some(automation.creator_id.clone()),
"delete_default",
),
declared_capabilities_for_automation(&automation),
automation.created_at_ms,
now,
)
});
record.deleted_at_ms = Some(now);
record.delete_retention_until_ms =
Some(now.saturating_add(7 * 24 * 60 * 60 * 1000));
record.updated_at_ms = now;
governance.deleted_automations.insert(
automation_id.to_string(),
DeletedAutomationRecord {
automation: automation.clone(),
deleted_at_ms: now,
deleted_by: deleted_by.clone(),
restore_until_ms: now.saturating_add(7 * 24 * 60 * 60 * 1000),
},
);
governance.updated_at_ms = now;
}
self.persist_automation_governance().await?;
self.persist_automations_v2_locked().await?;
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.deleted"),
&tandem_types::TenantContext::local_implicit(),
deleted_by
.actor_id
.clone()
.or_else(|| deleted_by.source.clone()),
json!({
"automationID": automation_id,
"deletedBy": deleted_by,
"deletedAtMs": now,
}),
)
.await?;
}
Ok(removed)
}
pub async fn get_deleted_automation_v2(
&self,
automation_id: &str,
) -> Option<crate::AutomationV2Spec> {
self.automation_governance
.read()
.await
.deleted_automations
.get(automation_id)
.map(|deleted| deleted.automation.clone())
}
pub async fn restore_deleted_automation_v2(
&self,
automation_id: &str,
restored_by: GovernanceActorRef,
approval_id: Option<String>,
tenant_context: &tandem_types::TenantContext,
) -> anyhow::Result<Option<crate::AutomationV2Spec>> {
let Some(candidate) = self
.automation_governance
.read()
.await
.deleted_automations
.get(automation_id)
.map(|deleted| deleted.automation.clone())
else {
return Ok(None);
};
if !governance_tenant_matches(&candidate.tenant_context(), tenant_context) {
return Ok(None);
}
let Some(record) = self.get_automation_governance(automation_id).await else {
return Ok(None);
};
if !governance_record_owned_by(&record, tenant_context) {
return Ok(None);
}
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.restored"),
tenant_context,
restored_by
.actor_id
.clone()
.or_else(|| restored_by.source.clone()),
json!({
"automationID": automation_id,
"restoredBy": restored_by,
"approvalID": approval_id,
}),
)
.await?;
let current = self
.automation_governance
.read()
.await
.deleted_automations
.get(automation_id)
.map(|deleted| deleted.automation.clone());
if current.as_ref().is_none_or(|automation| {
!governance_tenant_matches(&automation.tenant_context(), tenant_context)
}) {
return Ok(None);
}
let (restored, previous_governance) = {
let mut governance = self.automation_governance.write().await;
let previous = governance.clone();
let Some(record) = governance.records.get(automation_id) else {
return Ok(None);
};
if !governance_record_owned_by(record, tenant_context) {
return Ok(None);
}
let Some(deleted) = governance.deleted_automations.remove(automation_id) else {
return Ok(None);
};
let automation = deleted.automation;
if let Some(record) = governance.records.get_mut(automation_id) {
record.deleted_at_ms = None;
record.delete_retention_until_ms = None;
record.updated_at_ms = now_ms();
}
governance.updated_at_ms = now_ms();
(automation, previous)
};
if let Err(error) = self.persist_automation_governance().await {
*self.automation_governance.write().await = previous_governance;
return Err(error);
}
self.automations_v2
.write()
.await
.insert(automation_id.to_string(), restored.clone());
if let Err(error) = self.persist_automations_v2().await {
self.automations_v2.write().await.remove(automation_id);
*self.automation_governance.write().await = previous_governance;
let _ = self.persist_automation_governance().await;
return Err(error);
}
debug_assert_eq!(candidate.automation_id, restored.automation_id);
Ok(Some(restored))
}
pub async fn agent_spend_summary(&self, agent_id: &str) -> Option<AgentSpendSummary> {
self.automation_governance
.read()
.await
.agent_spend_summary(agent_id)
}
pub async fn tenant_agent_spend_summary(
&self,
tenant_context: &tandem_types::TenantContext,
agent_id: &str,
) -> Option<AgentSpendSummary> {
let scoped_agent_id = governance_agent_scope_key(tenant_context, agent_id);
self.automation_governance
.read()
.await
.agent_spend_summary(&scoped_agent_id)
}
pub async fn tenant_agent_has_quota_override(
&self,
tenant_context: &tandem_types::TenantContext,
agent_id: &str,
) -> bool {
let scoped_agent_id = governance_agent_scope_key(tenant_context, agent_id);
self.automation_governance
.read()
.await
.has_approved_agent_quota_override(&scoped_agent_id)
}
pub async fn tenant_agent_spend_paused_without_quota_override(
&self,
tenant_context: &tandem_types::TenantContext,
agent_id: &str,
) -> bool {
let scoped_agent_id = governance_agent_scope_key(tenant_context, agent_id);
let governance = self.automation_governance.read().await;
governance.is_agent_spend_paused(&scoped_agent_id)
&& !governance.has_approved_agent_quota_override(&scoped_agent_id)
}
pub async fn list_agent_spend_summaries(&self) -> Vec<AgentSpendSummary> {
self.automation_governance
.read()
.await
.agent_spend_summaries()
}
pub async fn agent_creation_review_summary(
&self,
agent_id: &str,
) -> Option<AgentCreationReviewSummary> {
self.automation_governance
.read()
.await
.agent_creation_review_summary(agent_id)
}
pub async fn list_agent_creation_review_summaries(&self) -> Vec<AgentCreationReviewSummary> {
self.automation_governance
.read()
.await
.agent_creation_review_summaries()
}
pub async fn record_agent_creation_review_progress(
&self,
agent_id: &str,
automation_id: &str,
) -> anyhow::Result<()> {
let now = now_ms();
let snapshot = {
let guard = self.automation_governance.read().await;
self.governance_snapshot(&guard)
};
let evaluation = self
.governance_engine
.evaluate_creation_review_progress(&snapshot, agent_id, automation_id, now)
.map_err(|error| anyhow::anyhow!(error.message))?;
let approval = evaluation.approval_request.clone();
{
let mut guard = self.automation_governance.write().await;
guard
.agent_creation_reviews
.insert(agent_id.to_string(), evaluation.summary);
if let Some(approval) = approval.clone() {
guard
.approvals
.insert(approval.approval_id.clone(), approval);
}
guard.updated_at_ms = now;
}
self.persist_automation_governance().await?;
if let Some(approval) = approval {
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.approval.requested"),
&tandem_types::TenantContext::local_implicit(),
approval
.requested_by
.actor_id
.clone()
.or_else(|| approval.requested_by.source.clone()),
json!({
"approvalID": approval.approval_id,
"request": approval,
}),
)
.await?;
}
Ok(())
}
pub async fn pause_automation_for_dependency_revocation(
&self,
automation_id: &str,
reason: String,
evidence: Value,
tenant_context: &tandem_types::TenantContext,
) -> anyhow::Result<()> {
let automation = self
.require_active_automation_governance_tenant(automation_id, tenant_context)
.await?;
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.dependency_revocation.requested"),
tenant_context,
Some("automation_dependency_revocation".to_string()),
json!({
"automationID": automation_id,
"reason": reason,
"evidence": evidence,
}),
)
.await?;
self.require_active_automation_governance_tenant(automation_id, tenant_context)
.await?;
let (evaluation, created_review_id, paused_runs, dependency_context) = {
let mut guard = self.automation_governance.write().await;
let Some(current_record) = guard.records.get(automation_id) else {
anyhow::bail!("automation governance record not found");
};
if !governance_record_owned_by(current_record, tenant_context) {
anyhow::bail!("automation governance record not found");
}
let now = now_ms();
let mut matching_dependency_review_id = None;
let mut superseded_review_ids = Vec::new();
for approval in guard.approvals.values() {
if approval.status != GovernanceApprovalStatus::Pending
|| approval.request_type != GovernanceApprovalRequestType::LifecycleReview
|| approval.target_resource.resource_type != "automation"
|| approval.target_resource.id != automation_id
|| !approval_receipt_matches_tenant(approval, tenant_context)
{
continue;
}
let exact_dependency_review = approval.expires_at_ms > now
&& approval.context.get("trigger")
== Some(&Value::String("dependency_revoked".to_string()))
&& approval.context.get("reason") == Some(&Value::String(reason.clone()))
&& approval.context.pointer("/evidence/evidence") == Some(&evidence);
if exact_dependency_review && matching_dependency_review_id.is_none() {
matching_dependency_review_id = Some(approval.approval_id.clone());
} else {
superseded_review_ids.push(approval.approval_id.clone());
}
}
if !superseded_review_ids.is_empty() {
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.approval.superseded"),
tenant_context,
Some("automation_dependency_revocation".to_string()),
json!({
"automationID": automation_id,
"supersededApprovalIDs": superseded_review_ids,
"reason": "dependency_revocation_requires_exact_review",
}),
)
.await?;
let superseded_at_ms = now_ms();
for approval_id in &superseded_review_ids {
if let Some(approval) = guard.approvals.get_mut(approval_id) {
approval.status = GovernanceApprovalStatus::Expired;
approval.updated_at_ms = superseded_at_ms;
}
}
}
let record = guard
.records
.get_mut(automation_id)
.expect("validated governance record remains present");
record.creation_paused = true;
record.paused_for_lifecycle = true;
record.review_required = true;
record.review_kind = Some(AutomationLifecycleReviewKind::DependencyRevoked);
record.review_requested_at_ms = Some(now);
record.review_request_id = matching_dependency_review_id.clone();
record.updated_at_ms = now;
guard.updated_at_ms = now;
if let Err(error) = self.persist_automation_governance_snapshot(&guard).await {
return Err(error.context(
"failed to persist dependency admission closure; admission remains paused in memory",
));
}
let paused_runs = self
.pause_running_automation_v2_runs(
automation_id,
reason.clone(),
crate::AutomationStopKind::GuardrailStopped,
)
.await
.context("failed to persist paused runs; admission remains durably closed")?;
let dependency_context = json!({
"trigger": "dependency_revoked",
"reason": reason.clone(),
"evidence": evidence.clone(),
"pausedRunIDs": paused_runs.clone(),
});
let snapshot = self.governance_snapshot(&guard);
let current_record = guard.records.get(automation_id).cloned();
let evaluation_now = now_ms();
let mut evaluation = self
.governance_engine
.evaluate_dependency_revocation(
&snapshot,
GovernanceDependencyRevocationInput {
automation_id: automation_id.to_string(),
current_record,
default_provenance: default_human_provenance(
Some(automation.creator_id.clone()),
"dependency_revocation_default",
),
declared_capabilities: declared_capabilities_for_automation(&automation),
reason: reason.clone(),
evidence: dependency_context.clone(),
},
evaluation_now,
)
.map_err(|error| anyhow::anyhow!(error.message))?;
if let Some(review_id) = matching_dependency_review_id.as_ref() {
let expired_at_evaluation = guard
.approvals
.get(review_id)
.is_none_or(|approval| approval.expires_at_ms <= evaluation_now);
if expired_at_evaluation {
anyhow::ensure!(
evaluation
.approval_request
.as_ref()
.is_some_and(|approval| {
approval.approval_id.as_str() != review_id.as_str()
}),
"dependency revocation did not replace an expired lifecycle review"
);
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.approval.superseded"),
tenant_context,
Some("automation_dependency_revocation".to_string()),
json!({
"automationID": automation_id,
"supersededApprovalIDs": [review_id],
"reason": "dependency_revocation_exact_review_expired_during_transition",
}),
)
.await?;
let expired_at_ms = now_ms();
if let Some(approval) = guard.approvals.get_mut(review_id) {
if approval.status == GovernanceApprovalStatus::Pending {
approval.status = GovernanceApprovalStatus::Expired;
approval.updated_at_ms = expired_at_ms;
}
}
}
}
bind_governance_record_to_tenant(&mut evaluation.record, tenant_context)?;
evaluation.record.creation_paused = true;
evaluation.record.paused_for_lifecycle = true;
evaluation.record.review_required = true;
if let Some(approval) = evaluation.approval_request.as_mut() {
match approval.tenant_context.as_ref() {
Some(owner) if !governance_tenant_matches(owner, tenant_context) => {
anyhow::bail!("governance approval tenant ownership mismatch");
}
Some(_) => {}
None => {
approval.tenant_context = Some(governance_tenant_scope(tenant_context));
}
}
}
let created_review_id = evaluation
.approval_request
.as_ref()
.map(|approval| approval.approval_id.clone())
.or_else(|| evaluation.record.review_request_id.clone());
guard
.records
.insert(automation_id.to_string(), evaluation.record.clone());
if let Some(approval) = evaluation.approval_request.clone() {
guard
.approvals
.insert(approval.approval_id.clone(), approval);
}
guard.updated_at_ms = evaluation_now;
if let Err(error) = self.persist_automation_governance_snapshot(&guard).await {
return Err(error.context(
"dependency revocation persistence failed; admission remains paused in memory",
));
}
(
evaluation,
created_review_id,
paused_runs,
dependency_context,
)
};
if let Some(approval) = evaluation.approval_request {
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.approval.requested"),
tenant_context,
approval
.requested_by
.actor_id
.clone()
.or_else(|| approval.requested_by.source.clone()),
json!({
"approvalID": approval.approval_id,
"request": approval,
}),
)
.await?;
}
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.dependency_revoked"),
tenant_context,
Some("automation_dependency_revocation".to_string()),
json!({
"automationID": automation_id,
"reason": reason,
"pausedRunIDs": paused_runs,
"evidence": dependency_context.clone(),
"reviewRequestID": created_review_id,
}),
)
.await?;
Ok(())
}
async fn pause_running_automation_v2_runs(
&self,
automation_id: &str,
reason: String,
stop_kind: crate::AutomationStopKind,
) -> anyhow::Result<Vec<String>> {
let (pausing_transitions, cancellations, previously_paused) = {
let mut guard = self.automation_v2_runs.write().await;
let mut transitions = Vec::new();
let mut cancellations = Vec::new();
let mut previously_paused = Vec::new();
for run in guard
.values_mut()
.filter(|run| run.automation_id == automation_id)
{
if matches!(
run.status,
crate::AutomationRunStatus::Queued
| crate::AutomationRunStatus::Running
| crate::AutomationRunStatus::Pausing
) {
let previous_status = run.status.clone();
let previous_gate = run.checkpoint.awaiting_gate.clone();
cancellations.push((
run.run_id.clone(),
run.active_session_ids.clone(),
run.active_instance_ids.clone(),
));
run.status = crate::AutomationRunStatus::Pausing;
run.pause_reason = Some(reason.clone());
run.scheduler = None;
run.execution_claim = None;
run.updated_at_ms = now_ms();
transitions.push((previous_status, previous_gate, run.clone()));
} else if run.status == crate::AutomationRunStatus::Paused
&& run.pause_reason.as_deref() == Some(reason.as_str())
&& run.stop_kind.as_ref() == Some(&stop_kind)
{
previously_paused.push(run.run_id.clone());
}
}
(transitions, cancellations, previously_paused)
};
if cancellations.is_empty() {
if !previously_paused.is_empty() {
self.persist_automation_v2_runs()
.await
.context("failed to retry persistence for governance-paused automation runs")?;
}
return Ok(previously_paused);
}
let _ = self.persist_automation_v2_runs().await;
for (previous_status, previous_gate, run) in pausing_transitions {
self.sync_automation_scheduler_for_run_transition(previous_status.clone(), &run)
.await;
let _ = self.persist_automation_v2_run_status_json(&run).await;
self.project_automation_v2_stateful_boundaries_or_warn(&run)
.await;
self.sync_automation_v2_durable_wait_transition(previous_status, previous_gate, &run)
.await;
}
for (_, session_ids, instance_ids) in &cancellations {
for session_id in session_ids {
let _ = self.cancellations.cancel(session_id).await;
}
for instance_id in instance_ids {
let _ = self
.agent_teams
.cancel_instance(self, instance_id, &reason)
.await;
}
self.forget_automation_v2_sessions(session_ids).await;
}
let paused_transitions = {
let mut guard = self.automation_v2_runs.write().await;
let mut transitions = Vec::new();
for (run_id, _, _) in &cancellations {
let Some(run) = guard.get_mut(run_id) else {
continue;
};
if !matches!(
run.status,
crate::AutomationRunStatus::Queued
| crate::AutomationRunStatus::Running
| crate::AutomationRunStatus::Pausing
) {
continue;
}
let previous_status = run.status.clone();
let previous_gate = run.checkpoint.awaiting_gate.clone();
run.status = crate::AutomationRunStatus::Paused;
run.active_session_ids.clear();
run.active_instance_ids.clear();
run.pause_reason = Some(reason.clone());
run.stop_kind = Some(stop_kind.clone());
run.stop_reason = Some(reason.clone());
run.scheduler = None;
run.execution_claim = None;
run.updated_at_ms = now_ms();
crate::app::state::automation::lifecycle::record_automation_lifecycle_event(
run,
"run_paused_governance",
Some(reason.clone()),
Some(stop_kind.clone()),
);
transitions.push((previous_status, previous_gate, run.clone()));
}
transitions
};
if !paused_transitions.is_empty() {
self.persist_automation_v2_runs().await.context(
"failed to persist governance-paused automation runs; runs remain paused in memory",
)?;
}
for (previous_status, previous_gate, run) in paused_transitions {
self.sync_automation_scheduler_for_run_transition(previous_status.clone(), &run)
.await;
let _ = self.persist_automation_v2_run_status_json(&run).await;
self.project_automation_v2_stateful_boundaries_or_warn(&run)
.await;
self.sync_automation_v2_durable_wait_transition(previous_status, previous_gate, &run)
.await;
}
let mut paused_run_ids = cancellations
.into_iter()
.map(|(run_id, _, _)| run_id)
.collect::<Vec<_>>();
paused_run_ids.extend(previously_paused);
paused_run_ids.sort();
paused_run_ids.dedup();
Ok(paused_run_ids)
}
pub async fn run_automation_governance_health_check(&self) -> anyhow::Result<usize> {
if !self.premium_governance_enabled() {
return Ok(0);
}
let now = now_ms();
let limits = self.automation_governance.read().await.limits.clone();
let automations = self.list_automations_v2().await;
let mut finding_count = 0usize;
let run_window = self.governance_engine.health_check_run_window(&limits);
for automation in automations {
let tenant_context = automation.tenant_context();
let runs = self
.list_automation_v2_runs(Some(&automation.automation_id), run_window)
.await;
let observations = runs
.iter()
.map(|run| GovernanceRunHealthObservation {
run_id: run.run_id.clone(),
status: governance_run_health_status(&run.status),
produced_node_outputs: !run.checkpoint.node_outputs.is_empty(),
guardrail_stopped: run.stop_kind
== Some(crate::AutomationStopKind::GuardrailStopped),
})
.collect::<Vec<_>>();
let run_health = self.governance_engine.summarize_run_health(&observations);
let evaluation = {
let mut guard = self.automation_governance.write().await;
let snapshot = self.governance_snapshot(&guard);
let mut evaluation = self
.governance_engine
.evaluate_health_check(
&snapshot,
GovernanceHealthCheckInput {
automation_id: automation.automation_id.clone(),
current_record: guard.records.get(&automation.automation_id).cloned(),
default_provenance: default_human_provenance(
Some(automation.creator_id.clone()),
"health_check_default",
),
declared_capabilities: declared_capabilities_for_automation(
&automation,
),
terminal_run_count: run_health.terminal_run_count,
failure_count: run_health.failure_count,
empty_output_count: run_health.empty_output_count,
guardrail_stop_count: run_health.guardrail_stop_count,
last_terminal_run_id: run_health.last_terminal_run_id.clone(),
},
now,
)
.map_err(|error| anyhow::anyhow!(error.message))?;
if let Some(evaluation) = evaluation.as_mut() {
bind_governance_record_to_tenant(&mut evaluation.record, &tenant_context)?;
for approval in &mut evaluation.approval_requests {
bind_governance_approval_to_tenant(approval, &tenant_context)?;
}
let previous = guard.clone();
guard
.records
.insert(automation.automation_id.clone(), evaluation.record.clone());
for approval in &evaluation.approval_requests {
guard
.approvals
.insert(approval.approval_id.clone(), approval.clone());
}
guard.updated_at_ms = now;
if let Err(error) = self.persist_automation_governance_snapshot(&guard).await {
*guard = previous;
return Err(error);
}
}
evaluation
};
let Some(evaluation) = evaluation else {
continue;
};
if evaluation.pause_automation && automation.status != crate::AutomationV2Status::Paused
{
let mut paused = automation.clone();
paused.status = crate::AutomationV2Status::Paused;
let _ = self.put_automation_v2(paused).await;
let _ = self
.pause_running_automation_v2_runs(
&automation.automation_id,
format!(
"automation expired after reaching {}ms retention",
limits.default_expires_after_ms
),
crate::AutomationStopKind::GuardrailStopped,
)
.await;
}
for approval in &evaluation.approval_requests {
append_protected_audit_event(
self,
format!("{GOVERNANCE_AUDIT_EVENT_PREFIX}.approval.requested"),
&tenant_context,
approval
.requested_by
.actor_id
.clone()
.or_else(|| approval.requested_by.source.clone()),
json!({
"approvalID": approval.approval_id,
"request": approval,
}),
)
.await?;
}
finding_count += evaluation.record.health_findings.len();
}
Ok(finding_count)
}
}
include!("governance_parts/part01.rs");