Skip to main content

meerkat_mobkit/runtime/
gating.rs

1//! Gating subsystem — policy evaluation, audit logging, and module-backed decisions.
2
3use super::module_boundary::{
4    CORE_MODULE_MCP_TIMEOUT, MEMORY_CONFLICT_READ_MCP_TOOL, call_module_mcp_tool_json,
5    mcp_required_error, module_uses_mcp,
6};
7use super::*;
8
9impl MobkitRuntimeHandle {
10    fn next_gating_sequence(&mut self) -> u64 {
11        Self::next_sequence(&mut self.gating_sequence)
12    }
13    fn append_gating_audit(&mut self, mut entry: GatingAuditEntry) {
14        let audit_sequence = self.next_gating_sequence();
15        entry.audit_id = format!("gate-audit-{audit_sequence:06}");
16        entry.timestamp_ms = current_time_ms();
17        self.gating_audit.push(entry);
18        while self.gating_audit.len() > GATING_AUDIT_MAX_RETAINED {
19            self.gating_audit.remove(0);
20        }
21    }
22    fn refresh_gating_timeouts(&mut self) {
23        let now_ms = current_time_ms();
24        let expired = self
25            .gating_pending
26            .iter()
27            .filter(|(_, entry)| now_ms >= entry.deadline_at_ms)
28            .map(|(pending_id, _)| pending_id.clone())
29            .collect::<Vec<_>>();
30        for pending_id in expired {
31            if let Some(expired_entry) = self.gating_pending.remove(&pending_id) {
32                self.gating_pending_order
33                    .retain(|candidate| candidate != &pending_id);
34                self.append_gating_audit(GatingAuditEntry {
35                    audit_id: String::new(),
36                    timestamp_ms: 0,
37                    event_type: "timeout_fallback".to_string(),
38                    action_id: expired_entry.action_id.clone(),
39                    pending_id: Some(pending_id),
40                    actor_id: expired_entry.actor_id,
41                    risk_tier: expired_entry.risk_tier,
42                    outcome: GatingOutcome::SafeDraft,
43                    detail: serde_json::json!({
44                        "fallback": "safe_draft",
45                        "reason": "approval_timeout"
46                    }),
47                });
48            }
49        }
50    }
51    fn upsert_gating_pending_entry(&mut self, entry: GatingPendingEntry) {
52        let pending_id = entry.pending_id.clone();
53        self.gating_pending.insert(pending_id.clone(), entry);
54        self.gating_pending_order
55            .retain(|candidate| candidate != &pending_id);
56        self.gating_pending_order.push(pending_id);
57        while self.gating_pending_order.len() > GATING_PENDING_MAX_RETAINED {
58            let oldest = self.gating_pending_order.remove(0);
59            self.gating_pending.remove(&oldest);
60        }
61    }
62
63    fn parse_memory_conflict_mcp_response(
64        response: Value,
65    ) -> Result<Option<MemoryConflictSignal>, RuntimeBoundaryError> {
66        let candidate = response
67            .as_object()
68            .and_then(|payload| payload.get("conflict"))
69            .cloned()
70            .unwrap_or(response);
71        if candidate.is_null() {
72            return Ok(None);
73        }
74        serde_json::from_value::<MemoryConflictSignal>(candidate)
75            .map(Some)
76            .map_err(|error| {
77                RuntimeBoundaryError::Mcp(McpBoundaryError::InvalidToolPayload {
78                    module_id: "memory".to_string(),
79                    tool: MEMORY_CONFLICT_READ_MCP_TOOL.to_string(),
80                    reason: error.to_string(),
81                })
82            })
83    }
84
85    fn gating_memory_conflict_for_reference(
86        &self,
87        entity: Option<&str>,
88        topic: Option<&str>,
89    ) -> Result<Option<MemoryConflictSignal>, RuntimeBoundaryError> {
90        if !self.is_module_loaded("memory") {
91            return Ok(self.memory_conflict_for_reference(entity, topic));
92        }
93
94        let Some((memory_module, pre_spawn)) = self.module_and_prespawn("memory") else {
95            return Err(mcp_required_error("memory", MEMORY_CONFLICT_READ_MCP_TOOL));
96        };
97        if !module_uses_mcp(memory_module, pre_spawn) {
98            return Err(mcp_required_error("memory", MEMORY_CONFLICT_READ_MCP_TOOL));
99        }
100
101        let response = call_module_mcp_tool_json(
102            memory_module,
103            pre_spawn,
104            MEMORY_CONFLICT_READ_MCP_TOOL,
105            &serde_json::json!({
106                "entity": entity,
107                "topic": topic,
108            }),
109            CORE_MODULE_MCP_TIMEOUT,
110        )?;
111        Self::parse_memory_conflict_mcp_response(response)
112    }
113
114    pub fn evaluate_gating_action(
115        &mut self,
116        request: GatingEvaluateRequest,
117    ) -> GatingEvaluateResult {
118        self.refresh_gating_timeouts();
119        let action = request.action.trim().to_string();
120        let actor_id = request.actor_id.trim().to_string();
121        let requested_approver = request
122            .requested_approver
123            .as_deref()
124            .map(str::trim)
125            .filter(|value| !value.is_empty())
126            .map(ToString::to_string);
127        let approval_recipient = request
128            .approval_recipient
129            .as_deref()
130            .map(str::trim)
131            .filter(|value| !value.is_empty())
132            .map(ToString::to_string);
133        let approval_channel = request
134            .approval_channel
135            .as_deref()
136            .map(str::trim)
137            .filter(|value| !value.is_empty())
138            .map(ToString::to_string);
139        let entity = request
140            .entity
141            .as_deref()
142            .map(str::trim)
143            .filter(|value| !value.is_empty())
144            .map(ToString::to_string);
145        let topic = request
146            .topic
147            .as_deref()
148            .map(str::trim)
149            .filter(|value| !value.is_empty())
150            .map(ToString::to_string);
151        let action_sequence = self.next_gating_sequence();
152        let action_id = format!("gate-action-{action_sequence:06}");
153        let risk_tier = request.risk_tier.clone();
154
155        if matches!(request.risk_tier, GatingRiskTier::R2 | GatingRiskTier::R3) {
156            if !self.memory_conflicts.is_empty() && (entity.is_none() || topic.is_none()) {
157                self.append_gating_audit(GatingAuditEntry {
158                    audit_id: String::new(),
159                    timestamp_ms: 0,
160                    event_type: "conflict_blocked".to_string(),
161                    action_id: action_id.clone(),
162                    pending_id: None,
163                    actor_id: actor_id.clone(),
164                    risk_tier: risk_tier.clone(),
165                    outcome: GatingOutcome::SafeDraft,
166                    detail: serde_json::json!({
167                        "policy": "memory_conflict_context_required_v0_1",
168                        "reason": "memory_conflict_context_missing",
169                        "action": action,
170                        "reference": {
171                            "entity": entity,
172                            "topic": topic,
173                        },
174                        "missing_context": {
175                            "entity": entity.is_none(),
176                            "topic": topic.is_none(),
177                        },
178                        "conflict_count": self.memory_conflicts.len(),
179                    }),
180                });
181                return GatingEvaluateResult {
182                    action_id,
183                    action,
184                    actor_id,
185                    risk_tier,
186                    outcome: GatingOutcome::SafeDraft,
187                    pending_id: None,
188                    fallback_reason: Some("memory_conflict_context_missing".to_string()),
189                };
190            }
191            let conflict = match self
192                .gating_memory_conflict_for_reference(entity.as_deref(), topic.as_deref())
193            {
194                Ok(conflict) => conflict,
195                Err(error) => {
196                    self.append_gating_audit(GatingAuditEntry {
197                        audit_id: String::new(),
198                        timestamp_ms: 0,
199                        event_type: "memory_conflict_lookup_failed".to_string(),
200                        action_id: action_id.clone(),
201                        pending_id: None,
202                        actor_id: actor_id.clone(),
203                        risk_tier: risk_tier.clone(),
204                        outcome: GatingOutcome::SafeDraft,
205                        detail: serde_json::json!({
206                            "policy": "memory_conflict_lookup_via_core_mcp",
207                            "reason": "memory_conflict_lookup_failed",
208                            "error": format!("{error:?}"),
209                            "reference": {
210                                "entity": entity,
211                                "topic": topic,
212                            },
213                        }),
214                    });
215                    return GatingEvaluateResult {
216                        action_id,
217                        action,
218                        actor_id,
219                        risk_tier,
220                        outcome: GatingOutcome::SafeDraft,
221                        pending_id: None,
222                        fallback_reason: Some("memory_conflict_lookup_failed".to_string()),
223                    };
224                }
225            };
226            if let Some(conflict) = conflict {
227                self.append_gating_audit(GatingAuditEntry {
228                    audit_id: String::new(),
229                    timestamp_ms: 0,
230                    event_type: "conflict_blocked".to_string(),
231                    action_id: action_id.clone(),
232                    pending_id: None,
233                    actor_id: actor_id.clone(),
234                    risk_tier: risk_tier.clone(),
235                    outcome: GatingOutcome::SafeDraft,
236                    detail: serde_json::json!({
237                        "policy": "memory_conflict_block_v0_1",
238                        "reason": "memory_conflict",
239                        "action": action,
240                        "reference": {
241                            "entity": entity,
242                            "topic": topic,
243                        },
244                        "conflict": conflict,
245                    }),
246                });
247                return GatingEvaluateResult {
248                    action_id,
249                    action,
250                    actor_id,
251                    risk_tier,
252                    outcome: GatingOutcome::SafeDraft,
253                    pending_id: None,
254                    fallback_reason: Some("memory_conflict".to_string()),
255                };
256            }
257        }
258
259        match request.risk_tier {
260            GatingRiskTier::R0 | GatingRiskTier::R1 => {
261                self.append_gating_audit(GatingAuditEntry {
262                    audit_id: String::new(),
263                    timestamp_ms: 0,
264                    event_type: "evaluated".to_string(),
265                    action_id: action_id.clone(),
266                    pending_id: None,
267                    actor_id: actor_id.clone(),
268                    risk_tier: risk_tier.clone(),
269                    outcome: GatingOutcome::Allowed,
270                    detail: serde_json::json!({
271                        "policy": "allow_immediate",
272                        "rationale": request.rationale,
273                        "action": action,
274                    }),
275                });
276                GatingEvaluateResult {
277                    action_id,
278                    action,
279                    actor_id,
280                    risk_tier,
281                    outcome: GatingOutcome::Allowed,
282                    pending_id: None,
283                    fallback_reason: None,
284                }
285            }
286            GatingRiskTier::R2 => {
287                self.append_gating_audit(GatingAuditEntry {
288                    audit_id: String::new(),
289                    timestamp_ms: 0,
290                    event_type: "evaluated".to_string(),
291                    action_id: action_id.clone(),
292                    pending_id: None,
293                    actor_id: actor_id.clone(),
294                    risk_tier: risk_tier.clone(),
295                    outcome: GatingOutcome::AllowedWithAudit,
296                    detail: serde_json::json!({
297                        "policy": "consequence_mode_allow_with_audit_v0_1",
298                        "rationale": request.rationale,
299                        "action": action,
300                    }),
301                });
302                GatingEvaluateResult {
303                    action_id,
304                    action,
305                    actor_id,
306                    risk_tier,
307                    outcome: GatingOutcome::AllowedWithAudit,
308                    pending_id: None,
309                    fallback_reason: None,
310                }
311            }
312            GatingRiskTier::R3 => {
313                let pending_sequence = self.next_gating_sequence();
314                let pending_id = format!("gate-pending-{pending_sequence:06}");
315                let created_at_ms = current_time_ms();
316                // Clamp both ends. The upper bound stops a deadline that
317                // saturates past `u64::MAX` from never expiring (the
318                // timeout-fallback guarantee must stay reachable). The lower
319                // bound stops a tiny/zero `approval_timeout_ms` from minting a
320                // pending entry that `refresh_gating_timeouts` treats as already
321                // expired on the next RPC, which would make R3 approval
322                // impossible to complete. Absent (`None`) still uses the
323                // default, which already sits inside the clamp.
324                let timeout_ms = request
325                    .approval_timeout_ms
326                    .unwrap_or(GATING_APPROVAL_TIMEOUT_DEFAULT_MS)
327                    .clamp(
328                        GATING_APPROVAL_TIMEOUT_MIN_MS,
329                        GATING_APPROVAL_TIMEOUT_MAX_MS,
330                    );
331                let mut approval_route_id = None;
332                let mut approval_delivery_id = None;
333                let mut approval_notification_error = None;
334
335                if let (Some(recipient), Some(channel)) =
336                    (approval_recipient.as_ref(), approval_channel.as_ref())
337                {
338                    if self.is_module_loaded("router") && self.is_module_loaded("delivery") {
339                        match self.resolve_routing(RoutingResolveRequest {
340                            recipient: recipient.clone(),
341                            channel: Some(channel.clone()),
342                            retry_max: None,
343                            backoff_ms: None,
344                            rate_limit_per_minute: None,
345                        }) {
346                            Ok(resolution) => {
347                                approval_route_id = Some(resolution.route_id.clone());
348                                match self.send_delivery(DeliverySendRequest {
349                                    resolution,
350                                    payload: serde_json::json!({
351                                        "kind": "gating_approval_request",
352                                        "pending_id": pending_id,
353                                        "action_id": action_id,
354                                        "action": action,
355                                        "actor_id": actor_id,
356                                        "risk_tier": risk_tier,
357                                        "requested_approver": requested_approver,
358                                        "deadline_at_ms": created_at_ms.saturating_add(timeout_ms),
359                                    }),
360                                    idempotency_key: Some(format!("gating-approval-{pending_id}")),
361                                }) {
362                                    Ok(record) => {
363                                        if record.status == "sent" {
364                                            approval_delivery_id = Some(record.delivery_id);
365                                        } else {
366                                            approval_notification_error = Some(format!(
367                                                "delivery_status:{}:{}",
368                                                record.status, record.delivery_id
369                                            ));
370                                        }
371                                    }
372                                    Err(err) => {
373                                        approval_notification_error =
374                                            Some(format!("delivery:{err:?}"));
375                                    }
376                                }
377                            }
378                            Err(err) => {
379                                approval_notification_error = Some(format!("routing:{err:?}"));
380                            }
381                        }
382                    } else {
383                        let mut missing_modules = Vec::new();
384                        if !self.is_module_loaded("router") {
385                            missing_modules.push("router");
386                        }
387                        if !self.is_module_loaded("delivery") {
388                            missing_modules.push("delivery");
389                        }
390                        approval_notification_error = Some(format!(
391                            "notification_modules_unavailable:{}",
392                            missing_modules.join(",")
393                        ));
394                    }
395                }
396                let pending_entry = GatingPendingEntry {
397                    pending_id: pending_id.clone(),
398                    action_id: action_id.clone(),
399                    action: action.clone(),
400                    actor_id: actor_id.clone(),
401                    risk_tier: risk_tier.clone(),
402                    requested_approver,
403                    approval_recipient,
404                    approval_channel,
405                    approval_route_id,
406                    approval_delivery_id,
407                    created_at_ms,
408                    deadline_at_ms: created_at_ms.saturating_add(timeout_ms),
409                };
410                self.upsert_gating_pending_entry(pending_entry.clone());
411                self.append_gating_audit(GatingAuditEntry {
412                    audit_id: String::new(),
413                    timestamp_ms: 0,
414                    event_type: "pending_created".to_string(),
415                    action_id: action_id.clone(),
416                    pending_id: Some(pending_id.clone()),
417                    actor_id: actor_id.clone(),
418                    risk_tier: risk_tier.clone(),
419                    outcome: GatingOutcome::PendingApproval,
420                    detail: serde_json::json!({
421                        "requested_approver": pending_entry.requested_approver,
422                        "approval_recipient": pending_entry.approval_recipient,
423                        "approval_channel": pending_entry.approval_channel,
424                        "approval_route_id": pending_entry.approval_route_id,
425                        "approval_delivery_id": pending_entry.approval_delivery_id,
426                        "approval_notification_error": approval_notification_error,
427                        "deadline_at_ms": pending_entry.deadline_at_ms,
428                        "action": action,
429                    }),
430                });
431                GatingEvaluateResult {
432                    action_id,
433                    action,
434                    actor_id,
435                    risk_tier,
436                    outcome: GatingOutcome::PendingApproval,
437                    pending_id: Some(pending_id),
438                    fallback_reason: None,
439                }
440            }
441        }
442    }
443
444    pub fn list_gating_pending(&mut self) -> Vec<GatingPendingEntry> {
445        self.refresh_gating_timeouts();
446        self.gating_pending_order
447            .iter()
448            .filter_map(|pending_id| self.gating_pending.get(pending_id).cloned())
449            .collect()
450    }
451
452    pub fn decide_gating_action(
453        &mut self,
454        request: GatingDecideRequest,
455    ) -> Result<GatingDecisionResult, GatingDecideError> {
456        self.refresh_gating_timeouts();
457        let decision = request.decision.clone();
458        let reason = request.reason.clone();
459        let pending_id = request.pending_id.trim().to_string();
460        let approver_id = request.approver_id.trim().to_string();
461        let pending_entry = self
462            .gating_pending
463            .remove(&pending_id)
464            .ok_or_else(|| GatingDecideError::UnknownPendingId(pending_id.clone()))?;
465        self.gating_pending_order
466            .retain(|candidate| candidate != &pending_id);
467
468        if matches!(decision, GatingDecision::Approve) && approver_id == pending_entry.actor_id {
469            self.upsert_gating_pending_entry(pending_entry);
470            return Err(GatingDecideError::SelfApprovalForbidden);
471        }
472        if let Some(expected_approver) = pending_entry.requested_approver.as_deref()
473            && expected_approver != approver_id
474        {
475            let expected = expected_approver.to_string();
476            self.upsert_gating_pending_entry(pending_entry);
477            return Err(GatingDecideError::ApproverMismatch {
478                expected,
479                provided: approver_id,
480            });
481        }
482
483        let mut next_pending_id = None;
484        let (outcome, event_type) = match decision {
485            GatingDecision::Approve => (GatingOutcome::Allowed, "approval_decided"),
486            GatingDecision::Reject => (GatingOutcome::SafeDraft, "rejection_decided"),
487            GatingDecision::Escalate => {
488                let successor_sequence = self.next_gating_sequence();
489                let successor_pending_id = format!("gate-pending-{successor_sequence:06}");
490                let successor_entry = GatingPendingEntry {
491                    pending_id: successor_pending_id.clone(),
492                    action_id: pending_entry.action_id.clone(),
493                    action: pending_entry.action.clone(),
494                    actor_id: pending_entry.actor_id.clone(),
495                    risk_tier: pending_entry.risk_tier.clone(),
496                    requested_approver: None,
497                    approval_recipient: pending_entry.approval_recipient.clone(),
498                    approval_channel: pending_entry.approval_channel.clone(),
499                    approval_route_id: None,
500                    approval_delivery_id: None,
501                    created_at_ms: current_time_ms(),
502                    deadline_at_ms: pending_entry.deadline_at_ms,
503                };
504                self.upsert_gating_pending_entry(successor_entry.clone());
505                next_pending_id = Some(successor_pending_id.clone());
506                self.append_gating_audit(GatingAuditEntry {
507                    audit_id: String::new(),
508                    timestamp_ms: 0,
509                    event_type: "pending_created".to_string(),
510                    action_id: successor_entry.action_id.clone(),
511                    pending_id: Some(successor_pending_id),
512                    actor_id: successor_entry.actor_id.clone(),
513                    risk_tier: successor_entry.risk_tier.clone(),
514                    outcome: GatingOutcome::PendingApproval,
515                    detail: serde_json::json!({
516                        "escalated_from_pending_id": pending_id.clone(),
517                        "requested_approver": successor_entry.requested_approver,
518                        "approval_recipient": successor_entry.approval_recipient,
519                        "approval_channel": successor_entry.approval_channel,
520                        "approval_route_id": successor_entry.approval_route_id,
521                        "approval_delivery_id": successor_entry.approval_delivery_id,
522                        "deadline_at_ms": successor_entry.deadline_at_ms,
523                        "action": successor_entry.action,
524                    }),
525                });
526                (GatingOutcome::PendingApproval, "escalation_decided")
527            }
528        };
529        let decided_at_ms = current_time_ms();
530        self.append_gating_audit(GatingAuditEntry {
531            audit_id: String::new(),
532            timestamp_ms: 0,
533            event_type: event_type.to_string(),
534            action_id: pending_entry.action_id.clone(),
535            pending_id: Some(pending_id.clone()),
536            actor_id: pending_entry.actor_id.clone(),
537            risk_tier: pending_entry.risk_tier.clone(),
538            outcome: outcome.clone(),
539            detail: serde_json::json!({
540                "approver_id": approver_id,
541                "decision": decision,
542                "reason": reason,
543                "approval_route_id": pending_entry.approval_route_id,
544                "approval_delivery_id": pending_entry.approval_delivery_id,
545                "next_pending_id": next_pending_id,
546            }),
547        });
548        Ok(GatingDecisionResult {
549            pending_id,
550            action_id: pending_entry.action_id,
551            approver_id,
552            decision,
553            outcome,
554            decided_at_ms,
555            reason,
556            next_pending_id,
557        })
558    }
559
560    pub fn gating_audit_entries(&mut self, limit: usize) -> Vec<GatingAuditEntry> {
561        self.refresh_gating_timeouts();
562        self.gating_audit
563            .iter()
564            .rev()
565            .take(limit)
566            .cloned()
567            .collect::<Vec<_>>()
568            .into_iter()
569            .rev()
570            .collect()
571    }
572}