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),
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 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}