1use 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 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}