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