Skip to main content

relay_knowledge/application/worker/
operations.rs

1use std::{
2    path::PathBuf,
3    sync::Arc,
4    time::{SystemTime, UNIX_EPOCH},
5};
6
7use serde_json::json;
8
9use crate::{
10    api::{
11        ApiError, ApiMetadata, AuditQueryApiRequest, AuditQueryResponse, IngestRequest,
12        InterfaceKind, ProposalDecisionApiRequest, ProposalDecisionResponse,
13        ProposalListApiRequest, ProposalListResponse, ProposalShowResponse, RequestContext,
14        ServiceDefinitionWriteResponse, ServiceOperatorResponse, ServicePlanRequest,
15        ServicePlanResponse, WorkerRunRequest, WorkerRunResponse, WorkerStatusRequest,
16        WorkerStatusResponse,
17    },
18    domain::{
19        AuditStatus, EvidenceModality, EvidenceRecord, GraphVersion, ProposalState,
20        ServiceDefinitionPlan, ServiceManagerAction, ServiceOperatorState, WorkerBackendState,
21        WorkerKind, WorkerStatus, WorkerTaskRecord, normalize_actor,
22    },
23    project::PROJECT_NAME,
24    storage::{
25        AuditQueryRequest, KnowledgeStore, NewAuditEvent, NewProposal, ProposalDecision,
26        ProposalListRequest, ServiceOperatorUpdate, StorageError, WorkerTaskClaimRequest,
27        WorkerTaskCompletion, WorkerTaskSeed,
28    },
29};
30
31use super::proposals::{fallback_proposal, proposal_from_worker_response, worker_request_payload};
32use crate::application::{
33    knowledge::{
34        index_refresh::{metadata_for_indexes, refresh_index_kinds},
35        ingest::mutation_batch_from_request,
36    },
37    service::RelayKnowledgeService,
38};
39
40const WORKER_LEASE_MS: u64 = 30_000;
41const WORKER_MAX_ATTEMPTS: u32 = 3;
42const PROPOSAL_LIST_LIMIT: usize = 50;
43const AUDIT_QUERY_LIMIT: usize = 100;
44
45struct AuditRecordInput<'a> {
46    operation: &'static str,
47    context: &'a RequestContext,
48    status: AuditStatus,
49    actor: Option<String>,
50    source_scope: Option<String>,
51    graph_version: GraphVersion,
52    detail: serde_json::Value,
53}
54
55impl RelayKnowledgeService {
56    /// Returns worker queue and backend readiness status.
57    pub async fn worker_status(
58        &self,
59        request: WorkerStatusRequest,
60        context: RequestContext,
61    ) -> Result<WorkerStatusResponse, ApiError> {
62        let store = self.storage.get().await.map_err(storage_api_error)?;
63        let graph_version = store
64            .current_graph_version()
65            .await
66            .map_err(storage_api_error)?;
67        let mut workers = overlay_worker_runtime(
68            store.worker_statuses().await.map_err(storage_api_error)?,
69            &self.runtime.workers,
70        );
71        if let Some(kind) = request.kind {
72            workers.retain(|status| status.kind == kind);
73        }
74        self.record_audit(
75            &store,
76            AuditRecordInput {
77                operation: "worker.status",
78                context: &context,
79                status: AuditStatus::Completed,
80                actor: None,
81                source_scope: None,
82                graph_version,
83                detail: json!({"worker_count": workers.len()}),
84            },
85        )
86        .await;
87
88        Ok(WorkerStatusResponse {
89            metadata: ApiMetadata::graph_only(&context, graph_version),
90            workers,
91        })
92    }
93
94    /// Runs one bounded worker task in the foreground.
95    pub async fn run_worker_once(
96        &self,
97        request: WorkerRunRequest,
98        context: RequestContext,
99    ) -> Result<WorkerRunResponse, ApiError> {
100        let store = self.storage.get().await.map_err(storage_api_error)?;
101        let graph_version = store
102            .current_graph_version()
103            .await
104            .map_err(storage_api_error)?;
105        let lease_owner = format!("worker-run-once-{}", std::process::id());
106        let Some(task) = store
107            .claim_worker_task(WorkerTaskClaimRequest {
108                kind: request.kind,
109                lease_owner: lease_owner.clone(),
110                lease_duration_ms: WORKER_LEASE_MS,
111                max_attempts: WORKER_MAX_ATTEMPTS,
112                now_ms: now_millis(),
113            })
114            .await
115            .map_err(storage_api_error)?
116        else {
117            let workers = overlay_worker_runtime(
118                store.worker_statuses().await.map_err(storage_api_error)?,
119                &self.runtime.workers,
120            );
121            return Ok(WorkerRunResponse {
122                metadata: ApiMetadata::graph_only(&context, graph_version),
123                task: None,
124                proposals: Vec::new(),
125                workers,
126                degraded_reason: None,
127            });
128        };
129
130        let (proposal, degraded_reason) = self.proposal_from_worker_task(&task).await?;
131        let proposal = store
132            .insert_proposal(proposal)
133            .await
134            .map_err(storage_api_error)?;
135        let completed = store
136            .complete_worker_task(WorkerTaskCompletion {
137                task_id: task.task_id.clone(),
138                lease_owner,
139                attempt_count: task.attempt_count,
140                now_ms: now_millis(),
141            })
142            .await
143            .map_err(storage_api_error)?;
144        let workers = overlay_worker_runtime(
145            store.worker_statuses().await.map_err(storage_api_error)?,
146            &self.runtime.workers,
147        );
148        self.record_audit(
149            &store,
150            AuditRecordInput {
151                operation: "worker.run_once",
152                context: &context,
153                status: AuditStatus::Completed,
154                actor: None,
155                source_scope: Some(task.source_scope.clone()),
156                graph_version,
157                detail: json!({"task_id": task.task_id, "proposal_id": proposal.proposal_id}),
158            },
159        )
160        .await;
161
162        Ok(WorkerRunResponse {
163            metadata: ApiMetadata::graph_only(&context, graph_version),
164            task: Some(completed),
165            proposals: vec![proposal],
166            workers,
167            degraded_reason,
168        })
169    }
170
171    /// Lists proposals awaiting or past manual lifecycle decisions.
172    pub async fn list_proposals(
173        &self,
174        request: ProposalListApiRequest,
175        context: RequestContext,
176    ) -> Result<ProposalListResponse, ApiError> {
177        let store = self.storage.get().await.map_err(storage_api_error)?;
178        let graph_version = store
179            .current_graph_version()
180            .await
181            .map_err(storage_api_error)?;
182        let proposals = store
183            .list_proposals(ProposalListRequest {
184                state: request.state,
185                limit: bounded_limit(request.limit, PROPOSAL_LIST_LIMIT),
186            })
187            .await
188            .map_err(storage_api_error)?;
189
190        Ok(ProposalListResponse {
191            metadata: ApiMetadata::graph_only(&context, graph_version),
192            proposals,
193        })
194    }
195
196    /// Shows one proposal and its conflict details.
197    pub async fn show_proposal(
198        &self,
199        proposal_id: String,
200        context: RequestContext,
201    ) -> Result<ProposalShowResponse, ApiError> {
202        let store = self.storage.get().await.map_err(storage_api_error)?;
203        let graph_version = store
204            .current_graph_version()
205            .await
206            .map_err(storage_api_error)?;
207        let proposal = store
208            .proposal_by_id(proposal_id.clone())
209            .await
210            .map_err(storage_api_error)?
211            .ok_or_else(|| {
212                ApiError::invalid_argument(format!("proposal '{proposal_id}' not found"))
213            })?;
214        let conflicts = store
215            .proposal_conflicts(proposal_id)
216            .await
217            .map_err(storage_api_error)?;
218        let payload = proposal.payload_value();
219
220        Ok(ProposalShowResponse {
221            metadata: ApiMetadata::graph_only(&context, graph_version),
222            proposal,
223            conflicts,
224            payload,
225        })
226    }
227
228    /// Accepts a proposal by committing its payload through the graph pipeline.
229    pub async fn accept_proposal(
230        &self,
231        proposal_id: String,
232        request: ProposalDecisionApiRequest,
233        context: RequestContext,
234    ) -> Result<ProposalDecisionResponse, ApiError> {
235        let actor = normalize_actor(request.actor)
236            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
237        let store = self.storage.get().await.map_err(storage_api_error)?;
238        let proposal = store
239            .proposal_by_id(proposal_id.clone())
240            .await
241            .map_err(storage_api_error)?
242            .ok_or_else(|| {
243                ApiError::invalid_argument(format!("proposal '{proposal_id}' not found"))
244            })?;
245        let ingest = serde_json::from_str::<IngestRequest>(&proposal.payload_json)
246            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
247        let batch = mutation_batch_from_request(ingest)
248            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
249        let evidence = batch.evidence.clone();
250        let receipt = store
251            .commit_mutation_batch(batch)
252            .await
253            .map_err(storage_api_error)?;
254        self.queue_worker_tasks_for_evidence(&store, &evidence, receipt.graph_version)
255            .await?;
256        let decided = store
257            .decide_proposal(ProposalDecision {
258                proposal_id,
259                next_state: ProposalState::Accepted,
260                actor,
261                reason: request.reason,
262                now_ms: now_millis(),
263            })
264            .await
265            .map_err(storage_api_error)?;
266        let (metadata, index_refresh_error) = match refresh_index_kinds(
267            &store,
268            crate::domain::IndexKind::ALL,
269            receipt.graph_version,
270            &self.runtime.retrieval,
271        )
272        .await
273        {
274            Ok(outcome) => (
275                metadata_for_indexes(&context, receipt.graph_version, &outcome.indexes),
276                None,
277            ),
278            Err(error) => (
279                ApiMetadata::indexed(&context, receipt.graph_version, None, None, true),
280                Some(error.message),
281            ),
282        };
283
284        Ok(ProposalDecisionResponse {
285            metadata,
286            proposal: decided,
287            receipt: Some(receipt),
288            index_refresh_error,
289        })
290    }
291
292    /// Rejects or supersedes a proposal without mutating graph facts.
293    pub async fn decide_proposal_without_commit(
294        &self,
295        proposal_id: String,
296        next_state: ProposalState,
297        request: ProposalDecisionApiRequest,
298        context: RequestContext,
299    ) -> Result<ProposalDecisionResponse, ApiError> {
300        if next_state == ProposalState::Accepted {
301            return Err(ApiError::invalid_argument(
302                "accept decisions must use accept_proposal".to_owned(),
303            ));
304        }
305        let actor = normalize_actor(request.actor)
306            .map_err(|error| ApiError::invalid_argument(error.to_string()))?;
307        let store = self.storage.get().await.map_err(storage_api_error)?;
308        let graph_version = store
309            .current_graph_version()
310            .await
311            .map_err(storage_api_error)?;
312        let proposal = store
313            .decide_proposal(ProposalDecision {
314                proposal_id,
315                next_state,
316                actor,
317                reason: request.reason,
318                now_ms: now_millis(),
319            })
320            .await
321            .map_err(storage_api_error)?;
322
323        Ok(ProposalDecisionResponse {
324            metadata: ApiMetadata::graph_only(&context, graph_version),
325            proposal,
326            receipt: None,
327            index_refresh_error: None,
328        })
329    }
330
331    /// Queries the durable audit sink.
332    pub async fn query_audit(
333        &self,
334        request: AuditQueryApiRequest,
335        context: RequestContext,
336    ) -> Result<AuditQueryResponse, ApiError> {
337        let store = self.storage.get().await.map_err(storage_api_error)?;
338        let graph_version = store
339            .current_graph_version()
340            .await
341            .map_err(storage_api_error)?;
342        let events = store
343            .query_audit_events(AuditQueryRequest {
344                operation: request.operation,
345                limit: bounded_limit(request.limit, AUDIT_QUERY_LIMIT),
346            })
347            .await
348            .map_err(storage_api_error)?;
349
350        Ok(AuditQueryResponse {
351            metadata: ApiMetadata::graph_only(&context, graph_version),
352            events,
353        })
354    }
355
356    /// Generates a service-manager plan without executing privileged commands.
357    pub async fn service_plan(
358        &self,
359        request: ServicePlanRequest,
360        context: RequestContext,
361    ) -> Result<ServicePlanResponse, ApiError> {
362        let store = self.storage.get().await.map_err(storage_api_error)?;
363        let graph_version = store
364            .current_graph_version()
365            .await
366            .map_err(storage_api_error)?;
367        let plan = self.render_service_plan(request.action);
368
369        Ok(ServicePlanResponse {
370            metadata: ApiMetadata::graph_only(&context, graph_version),
371            plan,
372        })
373    }
374
375    /// Writes the generated service definition into the service directory.
376    pub async fn write_service_definition(
377        &self,
378        context: RequestContext,
379    ) -> Result<ServiceDefinitionWriteResponse, ApiError> {
380        let store = self.storage.get().await.map_err(storage_api_error)?;
381        let graph_version = store
382            .current_graph_version()
383            .await
384            .map_err(storage_api_error)?;
385        let plan = self.render_service_plan(ServiceManagerAction::Install);
386        let path = PathBuf::from(&plan.definition_path);
387        let contents = plan.definition.clone();
388        tokio::task::spawn_blocking(move || {
389            if let Some(parent) = path.parent() {
390                std::fs::create_dir_all(parent)?;
391            }
392            std::fs::write(path, contents)
393        })
394        .await
395        .map_err(|error| ApiError::storage_unavailable(error.to_string()))?
396        .map_err(|error| ApiError::storage_unavailable(error.to_string()))?;
397
398        Ok(ServiceDefinitionWriteResponse {
399            metadata: ApiMetadata::graph_only(&context, graph_version),
400            plan,
401            written: true,
402        })
403    }
404
405    /// Returns persisted silent-update operator status.
406    pub async fn service_operator_status(
407        &self,
408        context: RequestContext,
409    ) -> Result<ServiceOperatorResponse, ApiError> {
410        let store = self.storage.get().await.map_err(storage_api_error)?;
411        let graph_version = store
412            .current_graph_version()
413            .await
414            .map_err(storage_api_error)?;
415        let operator = store
416            .service_operator_status()
417            .await
418            .map_err(storage_api_error)?;
419
420        Ok(ServiceOperatorResponse {
421            metadata: ApiMetadata::graph_only(&context, graph_version),
422            operator,
423        })
424    }
425
426    /// Pauses or resumes the silent-update operator.
427    pub async fn set_service_operator_state(
428        &self,
429        state: ServiceOperatorState,
430        context: RequestContext,
431    ) -> Result<ServiceOperatorResponse, ApiError> {
432        let store = self.storage.get().await.map_err(storage_api_error)?;
433        let graph_version = store
434            .current_graph_version()
435            .await
436            .map_err(storage_api_error)?;
437        let current = store
438            .service_operator_status()
439            .await
440            .map_err(storage_api_error)?;
441        let operator = store
442            .update_service_operator(ServiceOperatorUpdate {
443                state,
444                silent_updates_enabled: self.runtime.workers.silent_updates_enabled,
445                allowed_scopes: current.allowed_scopes,
446                last_error: current.last_error,
447                now_ms: now_millis(),
448            })
449            .await
450            .map_err(storage_api_error)?;
451
452        Ok(ServiceOperatorResponse {
453            metadata: ApiMetadata::graph_only(&context, graph_version),
454            operator,
455        })
456    }
457
458    pub(in crate::application) async fn queue_worker_tasks_for_evidence(
459        &self,
460        store: &Arc<dyn KnowledgeStore>,
461        evidence: &[EvidenceRecord],
462        graph_version: GraphVersion,
463    ) -> Result<(), ApiError> {
464        let now_ms = now_millis();
465        let seeds = evidence
466            .iter()
467            .flat_map(|record| worker_task_seeds(record, graph_version, now_ms))
468            .collect::<Vec<_>>();
469        if seeds.is_empty() {
470            return Ok(());
471        }
472        store
473            .queue_worker_tasks(seeds)
474            .await
475            .map(|_| ())
476            .map_err(storage_api_error)
477    }
478
479    async fn proposal_from_worker_task(
480        &self,
481        task: &WorkerTaskRecord,
482    ) -> Result<(NewProposal, Option<String>), ApiError> {
483        let fallback = fallback_proposal(task, WORKER_LEASE_MS, WORKER_MAX_ATTEMPTS)
484            .map_err(ApiError::invalid_argument)?;
485        let Some(endpoint) = self.runtime.workers.endpoint_for(task.kind) else {
486            return Ok((fallback, None));
487        };
488        let network = self.runtime.network.current();
489        let timeout_ms =
490            u64::try_from(network.http.request_timeout.as_millis()).unwrap_or(u64::MAX);
491        let payload = worker_request_payload(
492            task,
493            timeout_ms,
494            WORKER_LEASE_MS,
495            WORKER_MAX_ATTEMPTS,
496            self.runtime.workers.max_in_flight,
497        );
498        let response = crate::net::http::post_json(&network.http, endpoint, &payload).await;
499        match response {
500            Ok(value) => match proposal_from_worker_response(
501                task,
502                value,
503                WORKER_LEASE_MS,
504                WORKER_MAX_ATTEMPTS,
505            ) {
506                Ok(proposal) => Ok((proposal, None)),
507                Err(_) => Ok((
508                    fallback,
509                    Some(
510                        "worker response did not match proposal contract; deterministic fallback used"
511                            .to_owned(),
512                    ),
513                )),
514            },
515            Err(error) => Ok((
516                fallback,
517                Some(format!("external worker unavailable: {error}; deterministic fallback used")),
518            )),
519        }
520    }
521
522    async fn record_audit(&self, store: &Arc<dyn KnowledgeStore>, input: AuditRecordInput<'_>) {
523        let detail_json = serde_json::to_string(&input.detail).unwrap_or_else(|_| "{}".to_owned());
524        let _ = store
525            .insert_audit_event(NewAuditEvent {
526                operation: input.operation.to_owned(),
527                interface: interface_label(input.context.interface).to_owned(),
528                request_id: input.context.request_id.clone(),
529                trace_id: input.context.trace_id.clone(),
530                status: input.status,
531                actor: input.actor,
532                source_scope: input.source_scope,
533                graph_version: input.graph_version.get(),
534                detail_json,
535                message: None,
536                now_ms: now_millis(),
537            })
538            .await;
539    }
540
541    fn render_service_plan(&self, action: ServiceManagerAction) -> ServiceDefinitionPlan {
542        let definition_path = self
543            .runtime
544            .paths
545            .service_dir
546            .join(service_definition_filename())
547            .display()
548            .to_string();
549        let executable = std::env::current_exe()
550            .ok()
551            .map(|path| path.display().to_string())
552            .unwrap_or_else(|| "relay-knowledge".to_owned());
553        let platform = if cfg!(target_os = "windows") {
554            "windows"
555        } else if cfg!(target_os = "macos") {
556            "macos"
557        } else {
558            "linux"
559        }
560        .to_owned();
561        let definition = service_definition(
562            &platform,
563            &executable,
564            &self.runtime.paths.data_dir.display().to_string(),
565        );
566        let checksum = format!("{:016x}", stable_hash64(definition.as_bytes()));
567        let mut runtime_state_paths = vec![
568            self.runtime.paths.database_file().display().to_string(),
569            self.runtime.paths.config_dir.display().to_string(),
570            self.runtime.paths.state_dir.display().to_string(),
571            self.runtime.paths.log_dir.display().to_string(),
572            self.runtime.paths.cache_dir.display().to_string(),
573        ];
574        let mut warnings = vec![
575            "service manager commands are a generated plan; privileged install, uninstall, backup, migration, and rollback remain caller-executed".to_owned(),
576        ];
577        if self.runtime.storage.topology == crate::storage::StorageTopology::PartitionedSqlite {
578            runtime_state_paths.push(
579                self.runtime
580                    .paths
581                    .repository_shards_dir()
582                    .display()
583                    .to_string(),
584            );
585            warnings.push(
586                "partitioned_sqlite backup, migration, rollback, and uninstall confirmation must include both the control database and repository shard directory"
587                    .to_owned(),
588            );
589        }
590
591        ServiceDefinitionPlan {
592            action,
593            platform: platform.clone(),
594            service_name: PROJECT_NAME.to_owned(),
595            definition_path,
596            install_command: install_command(&platform),
597            uninstall_command: uninstall_command(&platform),
598            start_command: start_command(&platform),
599            stop_command: stop_command(&platform),
600            runtime_state_paths,
601            warnings,
602            definition,
603            checksum,
604        }
605    }
606}
607
608fn worker_task_seeds(
609    record: &EvidenceRecord,
610    graph_version: GraphVersion,
611    now_ms: u64,
612) -> Vec<WorkerTaskSeed> {
613    let kinds = match record.extraction.modality {
614        EvidenceModality::TextSpan | EvidenceModality::OcrText | EvidenceModality::Caption => {
615            vec![WorkerKind::Embedding, WorkerKind::Extractor]
616        }
617        EvidenceModality::ImageAsset => {
618            vec![WorkerKind::Ocr, WorkerKind::Vision, WorkerKind::Embedding]
619        }
620        EvidenceModality::ImageEmbedding => Vec::new(),
621        EvidenceModality::Table | EvidenceModality::LayoutRegion => vec![WorkerKind::Extractor],
622    };
623
624    kinds
625        .into_iter()
626        .map(|kind| {
627            let payload = json!({
628                "evidence_id": record.id,
629                "source_scope": record.source_scope.as_str(),
630                "modality": record.extraction.modality.as_str(),
631                "source_path": record.source_path.as_deref(),
632                "source_uri": record.extraction.source_uri.as_deref(),
633                "source_hash": record.extraction.source_hash.as_deref(),
634                "media_hash": record.extraction.media_hash.as_deref(),
635            });
636            WorkerTaskSeed {
637                kind,
638                source_scope: record.source_scope.as_str().to_owned(),
639                evidence_id: Some(record.id.clone()),
640                target_graph_version: graph_version,
641                input_fingerprint: format!(
642                    "{}:{}:{}",
643                    kind.as_str(),
644                    record.id,
645                    graph_version.get()
646                ),
647                payload_json: serde_json::to_string(&payload).unwrap_or_else(|_| "{}".to_owned()),
648                now_ms,
649            }
650        })
651        .collect()
652}
653
654pub(in crate::application) fn overlay_worker_runtime(
655    mut statuses: Vec<WorkerStatus>,
656    runtime: &crate::application::WorkerRuntimeConfig,
657) -> Vec<WorkerStatus> {
658    if statuses.is_empty() {
659        statuses = WorkerKind::ALL
660            .into_iter()
661            .map(|kind| WorkerStatus {
662                kind,
663                backend_state: WorkerBackendState::Fallback,
664                endpoint_configured: false,
665                queue_depth: 0,
666                running_count: 0,
667                retrying_count: 0,
668                dead_letter_count: 0,
669                last_error: None,
670            })
671            .collect();
672    }
673    for status in &mut statuses {
674        status.endpoint_configured = runtime.endpoint_for(status.kind).is_some();
675        status.backend_state = if status.dead_letter_count > 0 {
676            WorkerBackendState::Degraded
677        } else if status.endpoint_configured {
678            WorkerBackendState::Configured
679        } else {
680            WorkerBackendState::Fallback
681        };
682    }
683
684    statuses
685}
686
687fn bounded_limit(value: usize, default_limit: usize) -> usize {
688    if value == 0 {
689        default_limit
690    } else {
691        value.min(default_limit)
692    }
693}
694
695fn storage_api_error(error: StorageError) -> ApiError {
696    ApiError::storage_unavailable(error.to_string())
697}
698
699fn interface_label(interface: InterfaceKind) -> &'static str {
700    match interface {
701        InterfaceKind::Cli => "cli",
702        InterfaceKind::Web => "web",
703        InterfaceKind::Api => "api",
704        InterfaceKind::Mcp => "mcp",
705        InterfaceKind::Acp => "acp",
706    }
707}
708
709fn service_definition_filename() -> &'static str {
710    if cfg!(target_os = "windows") {
711        "relay-knowledge-service.xml"
712    } else if cfg!(target_os = "macos") {
713        "com.coolplayagent.relay-knowledge.plist"
714    } else {
715        "relay-knowledge.service"
716    }
717}
718
719fn service_definition(platform: &str, executable: &str, data_dir: &str) -> String {
720    match platform {
721        "windows" => format!(
722            "<service><id>relay-knowledge</id><name>relay-knowledge</name><executable>{executable}</executable><arguments>service run --web --mcp streamable-http</arguments><env name=\"RELAY_KNOWLEDGE_DATA_DIR\" value=\"{data_dir}\"/></service>\n"
723        ),
724        "macos" => format!(
725            "<?xml version=\"1.0\" encoding=\"UTF-8\"?><plist version=\"1.0\"><dict><key>Label</key><string>com.coolplayagent.relay-knowledge</string><key>ProgramArguments</key><array><string>{executable}</string><string>service</string><string>run</string><string>--web</string><string>--mcp</string><string>streamable-http</string></array><key>RunAtLoad</key><true/></dict></plist>\n"
726        ),
727        _ => format!(
728            "[Unit]\nDescription=relay-knowledge background service\nAfter=network-online.target\n\n[Service]\nType=simple\nExecStart={executable} service run --web --mcp streamable-http\nEnvironment=RELAY_KNOWLEDGE_DATA_DIR={data_dir}\nRestart=on-failure\n\n[Install]\nWantedBy=default.target\n"
729        ),
730    }
731}
732
733fn install_command(platform: &str) -> Vec<String> {
734    match platform {
735        "windows" => vec![
736            "powershell".to_owned(),
737            "New-Service".to_owned(),
738            "relay-knowledge".to_owned(),
739        ],
740        "macos" => vec![
741            "launchctl".to_owned(),
742            "load".to_owned(),
743            "<definition_path>".to_owned(),
744        ],
745        _ => vec![
746            "systemctl".to_owned(),
747            "--user".to_owned(),
748            "enable".to_owned(),
749            "--now".to_owned(),
750            "relay-knowledge.service".to_owned(),
751        ],
752    }
753}
754
755fn uninstall_command(platform: &str) -> Vec<String> {
756    match platform {
757        "windows" => vec![
758            "powershell".to_owned(),
759            "Remove-Service".to_owned(),
760            "relay-knowledge".to_owned(),
761        ],
762        "macos" => vec![
763            "launchctl".to_owned(),
764            "unload".to_owned(),
765            "<definition_path>".to_owned(),
766        ],
767        _ => vec![
768            "systemctl".to_owned(),
769            "--user".to_owned(),
770            "disable".to_owned(),
771            "--now".to_owned(),
772            "relay-knowledge.service".to_owned(),
773        ],
774    }
775}
776
777fn start_command(platform: &str) -> Vec<String> {
778    match platform {
779        "windows" => vec![
780            "powershell".to_owned(),
781            "Start-Service".to_owned(),
782            "relay-knowledge".to_owned(),
783        ],
784        "macos" => vec![
785            "launchctl".to_owned(),
786            "start".to_owned(),
787            "com.coolplayagent.relay-knowledge".to_owned(),
788        ],
789        _ => vec![
790            "systemctl".to_owned(),
791            "--user".to_owned(),
792            "start".to_owned(),
793            "relay-knowledge.service".to_owned(),
794        ],
795    }
796}
797
798fn stop_command(platform: &str) -> Vec<String> {
799    match platform {
800        "windows" => vec![
801            "powershell".to_owned(),
802            "Stop-Service".to_owned(),
803            "relay-knowledge".to_owned(),
804        ],
805        "macos" => vec![
806            "launchctl".to_owned(),
807            "stop".to_owned(),
808            "com.coolplayagent.relay-knowledge".to_owned(),
809        ],
810        _ => vec![
811            "systemctl".to_owned(),
812            "--user".to_owned(),
813            "stop".to_owned(),
814            "relay-knowledge.service".to_owned(),
815        ],
816    }
817}
818
819fn now_millis() -> u64 {
820    SystemTime::now()
821        .duration_since(UNIX_EPOCH)
822        .map_or(0, |duration| {
823            u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
824        })
825}
826
827fn stable_hash64(bytes: &[u8]) -> u64 {
828    const FNV_OFFSET_BASIS: u64 = 0xcbf29ce484222325;
829    const FNV_PRIME: u64 = 0x100000001b3;
830
831    let mut hash = FNV_OFFSET_BASIS;
832    for byte in bytes {
833        hash ^= u64::from(*byte);
834        hash = hash.wrapping_mul(FNV_PRIME);
835    }
836
837    hash
838}