Skip to main content

relay_knowledge/application/worker/
operations.rs

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    /// Returns worker queue and backend readiness status.
55    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    /// Runs one bounded worker task in the foreground.
93    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    /// Lists proposals awaiting or past manual lifecycle decisions.
170    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    /// Shows one proposal and its conflict details.
195    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    /// Accepts a proposal by committing its payload through the graph pipeline.
227    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    /// Rejects or supersedes a proposal without mutating graph facts.
291    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    /// Queries the durable audit sink.
330    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    /// Generates a service-manager plan without executing privileged commands.
355    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    /// Writes the generated service definition into the service directory.
382    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    /// Returns persisted silent-update operator status.
410    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    /// Pauses or resumes the silent-update operator.
431    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}