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