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 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 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 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 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 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 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 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 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 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 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 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}