1use std::{
2 sync::Arc,
3 time::{SystemTime, UNIX_EPOCH},
4};
5
6use serde_json::json;
7
8use crate::{
9 api::{
10 ApiError, ApiMetadata, AuditQueryApiRequest, AuditQueryResponse, IngestRequest,
11 InterfaceKind, ProposalDecisionApiRequest, ProposalDecisionResponse,
12 ProposalListApiRequest, ProposalListResponse, ProposalShowResponse, RequestContext,
13 ServiceDefinitionWriteResponse, ServiceOperatorResponse, ServicePlanRequest,
14 ServicePlanResponse, WorkerRunRequest, WorkerRunResponse, WorkerStatusRequest,
15 WorkerStatusResponse,
16 },
17 domain::{
18 AuditStatus, EvidenceModality, EvidenceRecord, GraphVersion, ProposalState,
19 ServiceManagerAction, ServiceOperatorState, WorkerBackendState, WorkerKind, WorkerStatus,
20 WorkerTaskRecord, normalize_actor,
21 },
22 storage::{
23 AuditQueryRequest, KnowledgeStore, NewAuditEvent, NewProposal, ProposalDecision,
24 ProposalListRequest, ServiceOperatorUpdate, StorageError, WorkerTaskClaimRequest,
25 WorkerTaskCompletion, WorkerTaskSeed,
26 },
27};
28
29use super::proposals::{fallback_proposal, proposal_from_worker_response, worker_request_payload};
30use crate::application::{
31 knowledge::{
32 index_refresh::{metadata_for_indexes, refresh_index_kinds},
33 ingest::mutation_batch_from_request,
34 },
35 service::RelayKnowledgeService,
36};
37
38const WORKER_LEASE_MS: u64 = 30_000;
39const WORKER_MAX_ATTEMPTS: u32 = 3;
40const PROPOSAL_LIST_LIMIT: usize = 50;
41const AUDIT_QUERY_LIMIT: usize = 100;
42
43struct AuditRecordInput<'a> {
44 operation: &'static str,
45 context: &'a RequestContext,
46 status: AuditStatus,
47 actor: Option<String>,
48 source_scope: Option<String>,
49 graph_version: GraphVersion,
50 detail: serde_json::Value,
51}
52
53impl RelayKnowledgeService {
54 pub async fn worker_status(
56 &self,
57 request: WorkerStatusRequest,
58 context: RequestContext,
59 ) -> Result<WorkerStatusResponse, ApiError> {
60 let store = self.storage.get().await.map_err(storage_api_error)?;
61 let graph_version = store
62 .current_graph_version()
63 .await
64 .map_err(storage_api_error)?;
65 let mut workers = overlay_worker_runtime(
66 store.worker_statuses().await.map_err(storage_api_error)?,
67 &self.runtime.workers,
68 );
69 if let Some(kind) = request.kind {
70 workers.retain(|status| status.kind == kind);
71 }
72 self.record_audit(
73 &store,
74 AuditRecordInput {
75 operation: "worker.status",
76 context: &context,
77 status: AuditStatus::Completed,
78 actor: None,
79 source_scope: None,
80 graph_version,
81 detail: json!({"worker_count": workers.len()}),
82 },
83 )
84 .await;
85
86 Ok(WorkerStatusResponse {
87 metadata: ApiMetadata::graph_only(&context, graph_version),
88 workers,
89 })
90 }
91
92 pub async fn run_worker_once(
94 &self,
95 request: WorkerRunRequest,
96 context: RequestContext,
97 ) -> Result<WorkerRunResponse, ApiError> {
98 let store = self.storage.get().await.map_err(storage_api_error)?;
99 let graph_version = store
100 .current_graph_version()
101 .await
102 .map_err(storage_api_error)?;
103 let lease_owner = format!("worker-run-once-{}", std::process::id());
104 let Some(task) = store
105 .claim_worker_task(WorkerTaskClaimRequest {
106 kind: request.kind,
107 lease_owner: lease_owner.clone(),
108 lease_duration_ms: WORKER_LEASE_MS,
109 max_attempts: WORKER_MAX_ATTEMPTS,
110 now_ms: now_millis(),
111 })
112 .await
113 .map_err(storage_api_error)?
114 else {
115 let workers = overlay_worker_runtime(
116 store.worker_statuses().await.map_err(storage_api_error)?,
117 &self.runtime.workers,
118 );
119 return Ok(WorkerRunResponse {
120 metadata: ApiMetadata::graph_only(&context, graph_version),
121 task: None,
122 proposals: Vec::new(),
123 workers,
124 degraded_reason: None,
125 });
126 };
127
128 let (proposal, degraded_reason) = self.proposal_from_worker_task(&task).await?;
129 let proposal = store
130 .insert_proposal(proposal)
131 .await
132 .map_err(storage_api_error)?;
133 let completed = store
134 .complete_worker_task(WorkerTaskCompletion {
135 task_id: task.task_id.clone(),
136 lease_owner,
137 attempt_count: task.attempt_count,
138 now_ms: now_millis(),
139 })
140 .await
141 .map_err(storage_api_error)?;
142 let workers = overlay_worker_runtime(
143 store.worker_statuses().await.map_err(storage_api_error)?,
144 &self.runtime.workers,
145 );
146 self.record_audit(
147 &store,
148 AuditRecordInput {
149 operation: "worker.run_once",
150 context: &context,
151 status: AuditStatus::Completed,
152 actor: None,
153 source_scope: Some(task.source_scope.clone()),
154 graph_version,
155 detail: json!({"task_id": task.task_id, "proposal_id": proposal.proposal_id}),
156 },
157 )
158 .await;
159
160 Ok(WorkerRunResponse {
161 metadata: ApiMetadata::graph_only(&context, graph_version),
162 task: Some(completed),
163 proposals: vec![proposal],
164 workers,
165 degraded_reason,
166 })
167 }
168
169 pub async fn list_proposals(
171 &self,
172 request: ProposalListApiRequest,
173 context: RequestContext,
174 ) -> Result<ProposalListResponse, ApiError> {
175 let store = self.storage.get().await.map_err(storage_api_error)?;
176 let graph_version = store
177 .current_graph_version()
178 .await
179 .map_err(storage_api_error)?;
180 let proposals = store
181 .list_proposals(ProposalListRequest {
182 state: request.state,
183 limit: bounded_limit(request.limit, PROPOSAL_LIST_LIMIT),
184 })
185 .await
186 .map_err(storage_api_error)?;
187
188 Ok(ProposalListResponse {
189 metadata: ApiMetadata::graph_only(&context, graph_version),
190 proposals,
191 })
192 }
193
194 pub async fn show_proposal(
196 &self,
197 proposal_id: String,
198 context: RequestContext,
199 ) -> Result<ProposalShowResponse, ApiError> {
200 let store = self.storage.get().await.map_err(storage_api_error)?;
201 let graph_version = store
202 .current_graph_version()
203 .await
204 .map_err(storage_api_error)?;
205 let proposal = store
206 .proposal_by_id(proposal_id.clone())
207 .await
208 .map_err(storage_api_error)?
209 .ok_or_else(|| {
210 ApiError::invalid_argument(format!("proposal '{proposal_id}' not found"))
211 })?;
212 let conflicts = store
213 .proposal_conflicts(proposal_id)
214 .await
215 .map_err(storage_api_error)?;
216 let payload = proposal.payload_value();
217
218 Ok(ProposalShowResponse {
219 metadata: ApiMetadata::graph_only(&context, graph_version),
220 proposal,
221 conflicts,
222 payload,
223 })
224 }
225
226 pub async fn accept_proposal(
228 &self,
229 proposal_id: String,
230 request: ProposalDecisionApiRequest,
231 context: RequestContext,
232 ) -> Result<ProposalDecisionResponse, ApiError> {
233 let actor = normalize_actor(request.actor)
234 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
235 let store = self.storage.get().await.map_err(storage_api_error)?;
236 let proposal = store
237 .proposal_by_id(proposal_id.clone())
238 .await
239 .map_err(storage_api_error)?
240 .ok_or_else(|| {
241 ApiError::invalid_argument(format!("proposal '{proposal_id}' not found"))
242 })?;
243 let ingest = serde_json::from_str::<IngestRequest>(&proposal.payload_json)
244 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
245 let batch = mutation_batch_from_request(ingest)
246 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
247 let evidence = batch.evidence.clone();
248 let receipt = store
249 .commit_mutation_batch(batch)
250 .await
251 .map_err(storage_api_error)?;
252 self.queue_worker_tasks_for_evidence(&store, &evidence, receipt.graph_version)
253 .await?;
254 let decided = store
255 .decide_proposal(ProposalDecision {
256 proposal_id,
257 next_state: ProposalState::Accepted,
258 actor,
259 reason: request.reason,
260 now_ms: now_millis(),
261 })
262 .await
263 .map_err(storage_api_error)?;
264 let (metadata, index_refresh_error) = match refresh_index_kinds(
265 &store,
266 crate::domain::IndexKind::ALL,
267 receipt.graph_version,
268 &self.runtime.retrieval,
269 )
270 .await
271 {
272 Ok(outcome) => (
273 metadata_for_indexes(&context, receipt.graph_version, &outcome.indexes),
274 None,
275 ),
276 Err(error) => (
277 ApiMetadata::indexed(&context, receipt.graph_version, None, None, true),
278 Some(error.message),
279 ),
280 };
281
282 Ok(ProposalDecisionResponse {
283 metadata,
284 proposal: decided,
285 receipt: Some(receipt),
286 index_refresh_error,
287 })
288 }
289
290 pub async fn decide_proposal_without_commit(
292 &self,
293 proposal_id: String,
294 next_state: ProposalState,
295 request: ProposalDecisionApiRequest,
296 context: RequestContext,
297 ) -> Result<ProposalDecisionResponse, ApiError> {
298 if next_state == ProposalState::Accepted {
299 return Err(ApiError::invalid_argument(
300 "accept decisions must use accept_proposal".to_owned(),
301 ));
302 }
303 let actor = normalize_actor(request.actor)
304 .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
305 let store = self.storage.get().await.map_err(storage_api_error)?;
306 let graph_version = store
307 .current_graph_version()
308 .await
309 .map_err(storage_api_error)?;
310 let proposal = store
311 .decide_proposal(ProposalDecision {
312 proposal_id,
313 next_state,
314 actor,
315 reason: request.reason,
316 now_ms: now_millis(),
317 })
318 .await
319 .map_err(storage_api_error)?;
320
321 Ok(ProposalDecisionResponse {
322 metadata: ApiMetadata::graph_only(&context, graph_version),
323 proposal,
324 receipt: None,
325 index_refresh_error: None,
326 })
327 }
328
329 pub async fn query_audit(
331 &self,
332 request: AuditQueryApiRequest,
333 context: RequestContext,
334 ) -> Result<AuditQueryResponse, ApiError> {
335 let store = self.storage.get().await.map_err(storage_api_error)?;
336 let graph_version = store
337 .current_graph_version()
338 .await
339 .map_err(storage_api_error)?;
340 let events = store
341 .query_audit_events(AuditQueryRequest {
342 operation: request.operation,
343 limit: bounded_limit(request.limit, AUDIT_QUERY_LIMIT),
344 })
345 .await
346 .map_err(storage_api_error)?;
347
348 Ok(AuditQueryResponse {
349 metadata: ApiMetadata::graph_only(&context, graph_version),
350 events,
351 })
352 }
353
354 pub async fn service_plan(
356 &self,
357 request: ServicePlanRequest,
358 context: RequestContext,
359 ) -> Result<ServicePlanResponse, ApiError> {
360 let store = self.storage.get().await.map_err(storage_api_error)?;
361 let graph_version = store
362 .current_graph_version()
363 .await
364 .map_err(storage_api_error)?;
365 let plan = self
366 .render_service_plan_for_request(&request)
367 .map_err(ApiError::invalid_argument)?;
368 let execution = if request.execute {
369 Some(self.execute_service_plan(&plan).await?)
370 } else {
371 None
372 };
373
374 Ok(ServicePlanResponse {
375 metadata: ApiMetadata::graph_only(&context, graph_version),
376 plan,
377 execution,
378 })
379 }
380
381 pub async fn write_service_definition(
383 &self,
384 context: RequestContext,
385 ) -> Result<ServiceDefinitionWriteResponse, ApiError> {
386 let store = self.storage.get().await.map_err(storage_api_error)?;
387 let graph_version = store
388 .current_graph_version()
389 .await
390 .map_err(storage_api_error)?;
391 let plan = self
392 .render_service_plan_for_request(&ServicePlanRequest {
393 action: ServiceManagerAction::Install,
394 dry_run: true,
395 execute: false,
396 target_version: None,
397 install_dir: None,
398 })
399 .map_err(ApiError::invalid_argument)?;
400 self.write_service_definition_from_plan(&plan).await?;
401
402 Ok(ServiceDefinitionWriteResponse {
403 metadata: ApiMetadata::graph_only(&context, graph_version),
404 plan,
405 written: true,
406 })
407 }
408
409 pub async fn service_operator_status(
411 &self,
412 context: RequestContext,
413 ) -> Result<ServiceOperatorResponse, ApiError> {
414 let store = self.storage.get().await.map_err(storage_api_error)?;
415 let graph_version = store
416 .current_graph_version()
417 .await
418 .map_err(storage_api_error)?;
419 let operator = store
420 .service_operator_status()
421 .await
422 .map_err(storage_api_error)?;
423
424 Ok(ServiceOperatorResponse {
425 metadata: ApiMetadata::graph_only(&context, graph_version),
426 operator,
427 })
428 }
429
430 pub async fn set_service_operator_state(
432 &self,
433 state: ServiceOperatorState,
434 context: RequestContext,
435 ) -> Result<ServiceOperatorResponse, ApiError> {
436 let store = self.storage.get().await.map_err(storage_api_error)?;
437 let graph_version = store
438 .current_graph_version()
439 .await
440 .map_err(storage_api_error)?;
441 let current = store
442 .service_operator_status()
443 .await
444 .map_err(storage_api_error)?;
445 let operator = store
446 .update_service_operator(ServiceOperatorUpdate {
447 state,
448 silent_updates_enabled: self.runtime.workers.silent_updates_enabled,
449 allowed_scopes: current.allowed_scopes,
450 last_error: current.last_error,
451 now_ms: now_millis(),
452 })
453 .await
454 .map_err(storage_api_error)?;
455
456 Ok(ServiceOperatorResponse {
457 metadata: ApiMetadata::graph_only(&context, graph_version),
458 operator,
459 })
460 }
461
462 pub(in crate::application) async fn queue_worker_tasks_for_evidence(
463 &self,
464 store: &Arc<dyn KnowledgeStore>,
465 evidence: &[EvidenceRecord],
466 graph_version: GraphVersion,
467 ) -> Result<(), ApiError> {
468 let now_ms = now_millis();
469 let seeds = evidence
470 .iter()
471 .flat_map(|record| worker_task_seeds(record, graph_version, now_ms))
472 .collect::<Vec<_>>();
473 if seeds.is_empty() {
474 return Ok(());
475 }
476 store
477 .queue_worker_tasks(seeds)
478 .await
479 .map(|_| ())
480 .map_err(storage_api_error)
481 }
482
483 async fn proposal_from_worker_task(
484 &self,
485 task: &WorkerTaskRecord,
486 ) -> Result<(NewProposal, Option<String>), ApiError> {
487 let fallback = fallback_proposal(task, WORKER_LEASE_MS, WORKER_MAX_ATTEMPTS)
488 .map_err(ApiError::invalid_argument)?;
489 let Some(endpoint) = self.runtime.workers.endpoint_for(task.kind) else {
490 return Ok((fallback, None));
491 };
492 let network = self.runtime.network.current();
493 let timeout_ms =
494 u64::try_from(network.http.request_timeout.as_millis()).unwrap_or(u64::MAX);
495 let payload = worker_request_payload(
496 task,
497 timeout_ms,
498 WORKER_LEASE_MS,
499 WORKER_MAX_ATTEMPTS,
500 self.runtime.workers.max_in_flight,
501 );
502 let response = crate::net::http::post_json_with_qos(
503 &network.http,
504 &self.runtime.network.qos_runtime(),
505 &network.qos,
506 endpoint,
507 &payload,
508 )
509 .await;
510 match response {
511 Ok(value) => match proposal_from_worker_response(
512 task,
513 value,
514 WORKER_LEASE_MS,
515 WORKER_MAX_ATTEMPTS,
516 ) {
517 Ok(proposal) => Ok((proposal, None)),
518 Err(_) => Ok((
519 fallback,
520 Some(
521 "worker response did not match proposal contract; deterministic fallback used"
522 .to_owned(),
523 ),
524 )),
525 },
526 Err(error) => Ok((
527 fallback,
528 Some(format!("external worker unavailable: {error}; deterministic fallback used")),
529 )),
530 }
531 }
532
533 async fn record_audit(&self, store: &Arc<dyn KnowledgeStore>, input: AuditRecordInput<'_>) {
534 let detail_json = serde_json::to_string(&input.detail).unwrap_or_else(|_| "{}".to_owned());
535 let _ = store
536 .insert_audit_event(NewAuditEvent {
537 operation: input.operation.to_owned(),
538 interface: interface_label(input.context.interface).to_owned(),
539 request_id: input.context.request_id.clone(),
540 trace_id: input.context.trace_id.clone(),
541 status: input.status,
542 actor: input.actor,
543 source_scope: input.source_scope,
544 graph_version: input.graph_version.get(),
545 detail_json,
546 message: None,
547 now_ms: now_millis(),
548 })
549 .await;
550 }
551}
552
553fn worker_task_seeds(
554 record: &EvidenceRecord,
555 graph_version: GraphVersion,
556 now_ms: u64,
557) -> Vec<WorkerTaskSeed> {
558 let kinds = match record.extraction.modality {
559 EvidenceModality::TextSpan | EvidenceModality::OcrText | EvidenceModality::Caption => {
560 vec![WorkerKind::Embedding, WorkerKind::Extractor]
561 }
562 EvidenceModality::ImageAsset => {
563 vec![WorkerKind::Ocr, WorkerKind::Vision, WorkerKind::Embedding]
564 }
565 EvidenceModality::ImageEmbedding => Vec::new(),
566 EvidenceModality::Table | EvidenceModality::LayoutRegion => vec![WorkerKind::Extractor],
567 };
568
569 kinds
570 .into_iter()
571 .map(|kind| {
572 let payload = json!({
573 "evidence_id": record.id,
574 "source_scope": record.source_scope.as_str(),
575 "modality": record.extraction.modality.as_str(),
576 "source_path": record.source_path.as_deref(),
577 "source_uri": record.extraction.source_uri.as_deref(),
578 "source_hash": record.extraction.source_hash.as_deref(),
579 "media_hash": record.extraction.media_hash.as_deref(),
580 });
581 WorkerTaskSeed {
582 kind,
583 source_scope: record.source_scope.as_str().to_owned(),
584 evidence_id: Some(record.id.clone()),
585 target_graph_version: graph_version,
586 input_fingerprint: format!(
587 "{}:{}:{}",
588 kind.as_str(),
589 record.id,
590 graph_version.get()
591 ),
592 payload_json: serde_json::to_string(&payload).unwrap_or_else(|_| "{}".to_owned()),
593 now_ms,
594 }
595 })
596 .collect()
597}
598
599pub(in crate::application) fn overlay_worker_runtime(
600 mut statuses: Vec<WorkerStatus>,
601 runtime: &crate::application::WorkerRuntimeConfig,
602) -> Vec<WorkerStatus> {
603 if statuses.is_empty() {
604 statuses = WorkerKind::ALL
605 .into_iter()
606 .map(|kind| WorkerStatus {
607 kind,
608 backend_state: WorkerBackendState::Fallback,
609 endpoint_configured: false,
610 queue_depth: 0,
611 running_count: 0,
612 retrying_count: 0,
613 dead_letter_count: 0,
614 last_error: None,
615 })
616 .collect();
617 }
618 for status in &mut statuses {
619 status.endpoint_configured = runtime.endpoint_for(status.kind).is_some();
620 status.backend_state = if status.dead_letter_count > 0 {
621 WorkerBackendState::Degraded
622 } else if status.endpoint_configured {
623 WorkerBackendState::Configured
624 } else {
625 WorkerBackendState::Fallback
626 };
627 }
628
629 statuses
630}
631
632fn bounded_limit(value: usize, default_limit: usize) -> usize {
633 if value == 0 {
634 default_limit
635 } else {
636 value.min(default_limit)
637 }
638}
639
640fn storage_api_error(error: StorageError) -> ApiError {
641 ApiError::storage_unavailable(error.to_string())
642}
643
644fn interface_label(interface: InterfaceKind) -> &'static str {
645 match interface {
646 InterfaceKind::Cli => "cli",
647 InterfaceKind::Web => "web",
648 InterfaceKind::Api => "api",
649 InterfaceKind::Mcp => "mcp",
650 InterfaceKind::Acp => "acp",
651 }
652}
653
654fn now_millis() -> u64 {
655 SystemTime::now()
656 .duration_since(UNIX_EPOCH)
657 .map_or(0, |duration| {
658 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
659 })
660}