Skip to main content

relay_knowledge/application/worker/
operations.rs

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    /// Returns worker queue and backend readiness status.
53    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    /// Runs one bounded worker task in the foreground.
91    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    /// Lists proposals awaiting or past manual lifecycle decisions.
168    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    /// Shows one proposal and its conflict details.
193    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    /// Accepts a proposal by committing its payload through the graph pipeline.
225    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    /// Rejects or supersedes a proposal without mutating graph facts.
289    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    /// Queries the durable audit sink.
328    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    /// Generates a service-manager plan without executing privileged commands.
353    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    /// Writes the generated service definition into the service directory.
380    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    /// Returns persisted silent-update operator status.
408    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    /// Pauses or resumes the silent-update operator.
429    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}