Skip to main content

relay_knowledge/application/
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::{
32    index_refresh::{metadata_for_indexes, refresh_index_kinds},
33    ingest::mutation_batch_from_request,
34    service::RelayKnowledgeService,
35    worker_proposals::{fallback_proposal, proposal_from_worker_response, worker_request_payload},
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.render_service_plan(request.action);
366
367        Ok(ServicePlanResponse {
368            metadata: ApiMetadata::graph_only(&context, graph_version),
369            plan,
370        })
371    }
372
373    /// Writes the generated service definition into the service directory.
374    pub async fn write_service_definition(
375        &self,
376        context: RequestContext,
377    ) -> Result<ServiceDefinitionWriteResponse, ApiError> {
378        let store = self.storage.get().await.map_err(storage_api_error)?;
379        let graph_version = store
380            .current_graph_version()
381            .await
382            .map_err(storage_api_error)?;
383        let plan = self.render_service_plan(ServiceManagerAction::Install);
384        let path = PathBuf::from(&plan.definition_path);
385        let contents = plan.definition.clone();
386        tokio::task::spawn_blocking(move || {
387            if let Some(parent) = path.parent() {
388                std::fs::create_dir_all(parent)?;
389            }
390            std::fs::write(path, contents)
391        })
392        .await
393        .map_err(|error| ApiError::storage_unavailable(error.to_string()))?
394        .map_err(|error| ApiError::storage_unavailable(error.to_string()))?;
395
396        Ok(ServiceDefinitionWriteResponse {
397            metadata: ApiMetadata::graph_only(&context, graph_version),
398            plan,
399            written: true,
400        })
401    }
402
403    /// Returns persisted silent-update operator status.
404    pub async fn service_operator_status(
405        &self,
406        context: RequestContext,
407    ) -> Result<ServiceOperatorResponse, ApiError> {
408        let store = self.storage.get().await.map_err(storage_api_error)?;
409        let graph_version = store
410            .current_graph_version()
411            .await
412            .map_err(storage_api_error)?;
413        let operator = store
414            .service_operator_status()
415            .await
416            .map_err(storage_api_error)?;
417
418        Ok(ServiceOperatorResponse {
419            metadata: ApiMetadata::graph_only(&context, graph_version),
420            operator,
421        })
422    }
423
424    /// Pauses or resumes the silent-update operator.
425    pub async fn set_service_operator_state(
426        &self,
427        state: ServiceOperatorState,
428        context: RequestContext,
429    ) -> Result<ServiceOperatorResponse, ApiError> {
430        let store = self.storage.get().await.map_err(storage_api_error)?;
431        let graph_version = store
432            .current_graph_version()
433            .await
434            .map_err(storage_api_error)?;
435        let current = store
436            .service_operator_status()
437            .await
438            .map_err(storage_api_error)?;
439        let operator = store
440            .update_service_operator(ServiceOperatorUpdate {
441                state,
442                silent_updates_enabled: self.runtime.workers.silent_updates_enabled,
443                allowed_scopes: current.allowed_scopes,
444                last_error: current.last_error,
445                now_ms: now_millis(),
446            })
447            .await
448            .map_err(storage_api_error)?;
449
450        Ok(ServiceOperatorResponse {
451            metadata: ApiMetadata::graph_only(&context, graph_version),
452            operator,
453        })
454    }
455
456    pub(super) async fn queue_worker_tasks_for_evidence(
457        &self,
458        store: &Arc<dyn KnowledgeStore>,
459        evidence: &[EvidenceRecord],
460        graph_version: GraphVersion,
461    ) -> Result<(), ApiError> {
462        let now_ms = now_millis();
463        let seeds = evidence
464            .iter()
465            .flat_map(|record| worker_task_seeds(record, graph_version, now_ms))
466            .collect::<Vec<_>>();
467        if seeds.is_empty() {
468            return Ok(());
469        }
470        store
471            .queue_worker_tasks(seeds)
472            .await
473            .map(|_| ())
474            .map_err(storage_api_error)
475    }
476
477    async fn proposal_from_worker_task(
478        &self,
479        task: &WorkerTaskRecord,
480    ) -> Result<(NewProposal, Option<String>), ApiError> {
481        let fallback = fallback_proposal(task, WORKER_LEASE_MS, WORKER_MAX_ATTEMPTS)
482            .map_err(ApiError::invalid_argument)?;
483        let Some(endpoint) = self.runtime.workers.endpoint_for(task.kind) else {
484            return Ok((fallback, None));
485        };
486        let network = self.runtime.network.current();
487        let timeout_ms =
488            u64::try_from(network.http.request_timeout.as_millis()).unwrap_or(u64::MAX);
489        let payload = worker_request_payload(
490            task,
491            timeout_ms,
492            WORKER_LEASE_MS,
493            WORKER_MAX_ATTEMPTS,
494            self.runtime.workers.max_in_flight,
495        );
496        let response = crate::net::http::post_json(&network.http, endpoint, &payload).await;
497        match response {
498            Ok(value) => match proposal_from_worker_response(
499                task,
500                value,
501                WORKER_LEASE_MS,
502                WORKER_MAX_ATTEMPTS,
503            ) {
504                Ok(proposal) => Ok((proposal, None)),
505                Err(_) => Ok((
506                    fallback,
507                    Some(
508                        "worker response did not match proposal contract; deterministic fallback used"
509                            .to_owned(),
510                    ),
511                )),
512            },
513            Err(error) => Ok((
514                fallback,
515                Some(format!("external worker unavailable: {error}; deterministic fallback used")),
516            )),
517        }
518    }
519
520    async fn record_audit(&self, store: &Arc<dyn KnowledgeStore>, input: AuditRecordInput<'_>) {
521        let detail_json = serde_json::to_string(&input.detail).unwrap_or_else(|_| "{}".to_owned());
522        let _ = store
523            .insert_audit_event(NewAuditEvent {
524                operation: input.operation.to_owned(),
525                interface: interface_label(input.context.interface).to_owned(),
526                request_id: input.context.request_id.clone(),
527                trace_id: input.context.trace_id.clone(),
528                status: input.status,
529                actor: input.actor,
530                source_scope: input.source_scope,
531                graph_version: input.graph_version.get(),
532                detail_json,
533                message: None,
534                now_ms: now_millis(),
535            })
536            .await;
537    }
538
539    fn render_service_plan(&self, action: ServiceManagerAction) -> ServiceDefinitionPlan {
540        let definition_path = self
541            .runtime
542            .paths
543            .service_dir
544            .join(service_definition_filename())
545            .display()
546            .to_string();
547        let executable = std::env::current_exe()
548            .ok()
549            .map(|path| path.display().to_string())
550            .unwrap_or_else(|| "relay-knowledge".to_owned());
551        let platform = if cfg!(target_os = "windows") {
552            "windows"
553        } else if cfg!(target_os = "macos") {
554            "macos"
555        } else {
556            "linux"
557        }
558        .to_owned();
559        let definition = service_definition(
560            &platform,
561            &executable,
562            &self.runtime.paths.data_dir.display().to_string(),
563        );
564        let checksum = format!("{:016x}", stable_hash64(definition.as_bytes()));
565
566        ServiceDefinitionPlan {
567            action,
568            platform: platform.clone(),
569            service_name: PROJECT_NAME.to_owned(),
570            definition_path,
571            install_command: install_command(&platform),
572            uninstall_command: uninstall_command(&platform),
573            start_command: start_command(&platform),
574            stop_command: stop_command(&platform),
575            definition,
576            checksum,
577        }
578    }
579}
580
581fn worker_task_seeds(
582    record: &EvidenceRecord,
583    graph_version: GraphVersion,
584    now_ms: u64,
585) -> Vec<WorkerTaskSeed> {
586    let kinds = match record.extraction.modality {
587        EvidenceModality::TextSpan | EvidenceModality::OcrText | EvidenceModality::Caption => {
588            vec![WorkerKind::Embedding, WorkerKind::Extractor]
589        }
590        EvidenceModality::ImageAsset => {
591            vec![WorkerKind::Ocr, WorkerKind::Vision, WorkerKind::Embedding]
592        }
593        EvidenceModality::ImageEmbedding => Vec::new(),
594        EvidenceModality::Table | EvidenceModality::LayoutRegion => vec![WorkerKind::Extractor],
595    };
596
597    kinds
598        .into_iter()
599        .map(|kind| {
600            let payload = json!({
601                "evidence_id": record.id,
602                "source_scope": record.source_scope.as_str(),
603                "modality": record.extraction.modality.as_str(),
604                "source_path": record.source_path.as_deref(),
605                "source_uri": record.extraction.source_uri.as_deref(),
606                "source_hash": record.extraction.source_hash.as_deref(),
607                "media_hash": record.extraction.media_hash.as_deref(),
608            });
609            WorkerTaskSeed {
610                kind,
611                source_scope: record.source_scope.as_str().to_owned(),
612                evidence_id: Some(record.id.clone()),
613                target_graph_version: graph_version,
614                input_fingerprint: format!(
615                    "{}:{}:{}",
616                    kind.as_str(),
617                    record.id,
618                    graph_version.get()
619                ),
620                payload_json: serde_json::to_string(&payload).unwrap_or_else(|_| "{}".to_owned()),
621                now_ms,
622            }
623        })
624        .collect()
625}
626
627pub(super) fn overlay_worker_runtime(
628    mut statuses: Vec<WorkerStatus>,
629    runtime: &super::WorkerRuntimeConfig,
630) -> Vec<WorkerStatus> {
631    if statuses.is_empty() {
632        statuses = WorkerKind::ALL
633            .into_iter()
634            .map(|kind| WorkerStatus {
635                kind,
636                backend_state: WorkerBackendState::Fallback,
637                endpoint_configured: false,
638                queue_depth: 0,
639                running_count: 0,
640                retrying_count: 0,
641                dead_letter_count: 0,
642                last_error: None,
643            })
644            .collect();
645    }
646    for status in &mut statuses {
647        status.endpoint_configured = runtime.endpoint_for(status.kind).is_some();
648        status.backend_state = if status.dead_letter_count > 0 {
649            WorkerBackendState::Degraded
650        } else if status.endpoint_configured {
651            WorkerBackendState::Configured
652        } else {
653            WorkerBackendState::Fallback
654        };
655    }
656
657    statuses
658}
659
660fn bounded_limit(value: usize, default_limit: usize) -> usize {
661    if value == 0 {
662        default_limit
663    } else {
664        value.min(default_limit)
665    }
666}
667
668fn storage_api_error(error: StorageError) -> ApiError {
669    ApiError::storage_unavailable(error.to_string())
670}
671
672fn interface_label(interface: InterfaceKind) -> &'static str {
673    match interface {
674        InterfaceKind::Cli => "cli",
675        InterfaceKind::Web => "web",
676        InterfaceKind::Api => "api",
677        InterfaceKind::Mcp => "mcp",
678        InterfaceKind::Acp => "acp",
679    }
680}
681
682fn service_definition_filename() -> &'static str {
683    if cfg!(target_os = "windows") {
684        "relay-knowledge-service.xml"
685    } else if cfg!(target_os = "macos") {
686        "com.coolplayagent.relay-knowledge.plist"
687    } else {
688        "relay-knowledge.service"
689    }
690}
691
692fn service_definition(platform: &str, executable: &str, data_dir: &str) -> String {
693    match platform {
694        "windows" => format!(
695            "<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"
696        ),
697        "macos" => format!(
698            "<?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"
699        ),
700        _ => format!(
701            "[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"
702        ),
703    }
704}
705
706fn install_command(platform: &str) -> Vec<String> {
707    match platform {
708        "windows" => vec![
709            "powershell".to_owned(),
710            "New-Service".to_owned(),
711            "relay-knowledge".to_owned(),
712        ],
713        "macos" => vec![
714            "launchctl".to_owned(),
715            "load".to_owned(),
716            "<definition_path>".to_owned(),
717        ],
718        _ => vec![
719            "systemctl".to_owned(),
720            "--user".to_owned(),
721            "enable".to_owned(),
722            "--now".to_owned(),
723            "relay-knowledge.service".to_owned(),
724        ],
725    }
726}
727
728fn uninstall_command(platform: &str) -> Vec<String> {
729    match platform {
730        "windows" => vec![
731            "powershell".to_owned(),
732            "Remove-Service".to_owned(),
733            "relay-knowledge".to_owned(),
734        ],
735        "macos" => vec![
736            "launchctl".to_owned(),
737            "unload".to_owned(),
738            "<definition_path>".to_owned(),
739        ],
740        _ => vec![
741            "systemctl".to_owned(),
742            "--user".to_owned(),
743            "disable".to_owned(),
744            "--now".to_owned(),
745            "relay-knowledge.service".to_owned(),
746        ],
747    }
748}
749
750fn start_command(platform: &str) -> Vec<String> {
751    match platform {
752        "windows" => vec![
753            "powershell".to_owned(),
754            "Start-Service".to_owned(),
755            "relay-knowledge".to_owned(),
756        ],
757        "macos" => vec![
758            "launchctl".to_owned(),
759            "start".to_owned(),
760            "com.coolplayagent.relay-knowledge".to_owned(),
761        ],
762        _ => vec![
763            "systemctl".to_owned(),
764            "--user".to_owned(),
765            "start".to_owned(),
766            "relay-knowledge.service".to_owned(),
767        ],
768    }
769}
770
771fn stop_command(platform: &str) -> Vec<String> {
772    match platform {
773        "windows" => vec![
774            "powershell".to_owned(),
775            "Stop-Service".to_owned(),
776            "relay-knowledge".to_owned(),
777        ],
778        "macos" => vec![
779            "launchctl".to_owned(),
780            "stop".to_owned(),
781            "com.coolplayagent.relay-knowledge".to_owned(),
782        ],
783        _ => vec![
784            "systemctl".to_owned(),
785            "--user".to_owned(),
786            "stop".to_owned(),
787            "relay-knowledge.service".to_owned(),
788        ],
789    }
790}
791
792fn now_millis() -> u64 {
793    SystemTime::now()
794        .duration_since(UNIX_EPOCH)
795        .map_or(0, |duration| {
796            u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
797        })
798}
799
800fn stable_hash64(bytes: &[u8]) -> u64 {
801    const FNV_OFFSET_BASIS: u64 = 0xcbf29ce484222325;
802    const FNV_PRIME: u64 = 0x100000001b3;
803
804    let mut hash = FNV_OFFSET_BASIS;
805    for byte in bytes {
806        hash ^= u64::from(*byte);
807        hash = hash.wrapping_mul(FNV_PRIME);
808    }
809
810    hash
811}