1#[cfg(unix)]
49use std::future::Future;
50use std::path::{Path, PathBuf};
51#[cfg(unix)]
52use std::time::Duration;
53
54use degenbot_config::EnvVars;
55use serde_json::{json, Map, Value};
56#[cfg(unix)]
57use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
58#[cfg(unix)]
59use tokio::net::UnixStream;
60
61use crate::error::CliError;
62
63pub const SOCKET_ENV: &str = "DEGENBOT_OPERATOR_SOCKET";
66
67pub const SOCKET_DEFAULT: &str = "~/.config/degenbot/operator.sock";
70
71const HOME_ENV: &str = "HOME";
73
74#[cfg(unix)]
77const REQUEST_TIMEOUT: Duration = Duration::from_secs(60);
78
79#[cfg(not(unix))]
84const UDS_UNSUPPORTED: &str = "the operator command channel requires a Unix domain socket, \
85 which this platform does not provide";
86
87pub const FLEET_POSTURE_THRESHOLD_KEYS: [&str; 6] = [
91 "cordon_enter_events",
92 "cordon_enter_window_ms",
93 "cordon_duty_percent",
94 "cordon_duty_window_ms",
95 "cordon_exit_clean_ms",
96 "cordon_sim_intake_floor",
97];
98
99pub const SIM_INTAKE_FLOOR_RESTORE: &str = "null";
104
105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107pub enum PathFamily {
108 V2,
110 V3,
112 V4,
114}
115
116impl PathFamily {
117 #[must_use]
119 pub const fn as_str(self) -> &'static str {
120 match self {
121 Self::V2 => "V2",
122 Self::V3 => "V3",
123 Self::V4 => "V4",
124 }
125 }
126
127 #[must_use]
129 pub fn parse(raw: &str) -> Option<Self> {
130 match raw.to_ascii_uppercase().as_str() {
131 "V2" => Some(Self::V2),
132 "V3" => Some(Self::V3),
133 "V4" => Some(Self::V4),
134 _ => None,
135 }
136 }
137}
138
139#[derive(Debug, Clone, PartialEq, Eq)]
141pub struct PathStep {
142 pub family: PathFamily,
144 pub address: String,
146 pub hash: Option<String>,
149}
150
151#[derive(Debug, Clone, Copy, PartialEq, Eq)]
153pub enum PathDirection {
154 Zfo,
156 Ozf,
158}
159
160impl PathDirection {
161 #[must_use]
163 pub const fn as_str(self) -> &'static str {
164 match self {
165 Self::Zfo => "zfo",
166 Self::Ozf => "ozf",
167 }
168 }
169
170 #[must_use]
172 pub fn parse(raw: &str) -> Option<Self> {
173 match raw {
174 "zfo" => Some(Self::Zfo),
175 "ozf" => Some(Self::Ozf),
176 _ => None,
177 }
178 }
179
180 #[must_use]
182 pub const fn is_zfo(self) -> bool {
183 matches!(self, Self::Zfo)
184 }
185}
186
187pub fn parse_hop_token(hop: &str) -> Result<PathStep, CliError> {
196 let parts: Vec<&str> = hop.split(':').collect();
197 let raw_family = parts.first().copied().unwrap_or_default();
198 let family = PathFamily::parse(raw_family).ok_or_else(|| {
199 CliError::OperatorHygiene(format!("--hop family must be V2|V3|V4, got {raw_family:?}"))
200 })?;
201 let address = parts
202 .get(1)
203 .copied()
204 .filter(|address| !address.is_empty())
205 .ok_or_else(|| CliError::OperatorHygiene(format!("--hop {hop:?} is missing an address")))?;
206 let hash = if family == PathFamily::V4 {
207 parts
208 .get(2)
209 .copied()
210 .filter(|hash| !hash.is_empty())
211 .map(ToString::to_string)
212 } else {
213 None
214 };
215 Ok(PathStep {
216 family,
217 address: address.to_string(),
218 hash,
219 })
220}
221
222#[derive(Debug, Clone, Copy, PartialEq)]
228pub enum PosturePatchValue {
229 Int(u64),
231 Float(f64),
233 Null,
235}
236
237impl PosturePatchValue {
238 #[must_use]
242 pub fn json_value(self) -> Value {
243 match self {
244 Self::Int(value) => Value::from(value),
245 Self::Float(value) => {
246 serde_json::Number::from_f64(value).map_or(Value::Null, Value::Number)
247 }
248 Self::Null => Value::Null,
249 }
250 }
251}
252
253#[derive(Debug, Clone, PartialEq)]
255pub struct PosturePatchEntry {
256 pub key: String,
258 pub value: PosturePatchValue,
260}
261
262impl PosturePatchEntry {
263 #[must_use]
265 pub fn new(key: impl Into<String>, value: PosturePatchValue) -> Self {
266 Self {
267 key: key.into(),
268 value,
269 }
270 }
271
272 #[must_use]
274 pub fn int(key: impl Into<String>, value: u64) -> Self {
275 Self::new(key, PosturePatchValue::Int(value))
276 }
277
278 #[must_use]
280 pub fn float(key: impl Into<String>, value: f64) -> Self {
281 Self::new(key, PosturePatchValue::Float(value))
282 }
283}
284
285pub fn parse_sim_intake_floor(raw: &str) -> Result<PosturePatchValue, CliError> {
294 if raw.eq_ignore_ascii_case(SIM_INTAKE_FLOOR_RESTORE) {
295 return Ok(PosturePatchValue::Null);
296 }
297 raw.trim()
298 .parse::<u64>()
299 .map(PosturePatchValue::Int)
300 .map_err(|_| {
301 CliError::OperatorHygiene(format!(
302 "cordon_sim_intake_floor must be an integer or {SIM_INTAKE_FLOOR_RESTORE:?} to restore half the slot cap, got {raw:?}"
303 ))
304 })
305}
306
307pub fn validate_posture_patch(patch: &[PosturePatchEntry]) -> Result<(), CliError> {
316 let mut unknown: Vec<&str> = patch
317 .iter()
318 .map(|entry| entry.key.as_str())
319 .filter(|key| !FLEET_POSTURE_THRESHOLD_KEYS.contains(key))
320 .collect();
321 unknown.sort_unstable();
322 unknown.dedup();
323 if !unknown.is_empty() {
324 return Err(CliError::OperatorHygiene(format!(
325 "unknown fleet-posture threshold key(s): {}",
326 unknown.join(", ")
327 )));
328 }
329 if patch.is_empty() {
330 return Err(CliError::OperatorHygiene(
331 "set_fleet_posture needs at least one threshold key (empty patch)".to_string(),
332 ));
333 }
334 Ok(())
335}
336
337#[derive(Debug, Clone, PartialEq)]
339pub enum WireRequest {
340 AddPath {
342 steps: Vec<PathStep>,
344 directions: Option<Vec<bool>>,
346 },
347 Discover {
349 bound: Option<u64>,
351 },
352 SetFleetPosture {
354 patch: Vec<PosturePatchEntry>,
356 },
357 GetFleetPosture,
359}
360
361impl WireRequest {
362 #[must_use]
364 pub const fn op_name(&self) -> &'static str {
365 match self {
366 Self::AddPath { .. } => "add_path",
367 Self::Discover { .. } => "discover",
368 Self::SetFleetPosture { .. } => "set_fleet_posture",
369 Self::GetFleetPosture => "get_fleet_posture",
370 }
371 }
372
373 #[must_use]
375 pub fn payload(&self) -> Value {
376 match self {
377 Self::AddPath { steps, directions } => {
378 let steps: Vec<Value> = steps.iter().map(step_json).collect();
379 let directions = directions.as_ref().map_or(Value::Null, |bits| json!(bits));
380 json!({ "steps": steps, "directions": directions })
381 }
382 Self::Discover { bound } => json!({ "bound": bound }),
383 Self::SetFleetPosture { patch } => {
384 let mut payload = Map::new();
385 for entry in patch {
386 payload.insert(entry.key.clone(), entry.value.json_value());
387 }
388 Value::Object(payload)
389 }
390 Self::GetFleetPosture => json!({}),
391 }
392 }
393
394 #[must_use]
396 pub fn encode_line(&self) -> String {
397 let mut envelope = Map::new();
398 envelope.insert("op".to_string(), Value::from(self.op_name()));
399 envelope.insert("payload".to_string(), self.payload());
400 let mut line =
401 serde_json::to_string(&Value::Object(envelope)).unwrap_or_else(|_| "{}".to_string());
402 line.push('\n');
403 line
404 }
405}
406
407fn step_json(step: &PathStep) -> Value {
409 let mut object = Map::new();
410 object.insert("family".to_string(), Value::from(step.family.as_str()));
411 object.insert("address".to_string(), Value::from(step.address.clone()));
412 if let Some(hash) = &step.hash {
413 object.insert("hash".to_string(), Value::from(hash.clone()));
414 }
415 Value::Object(object)
416}
417
418#[derive(Debug, Clone, PartialEq)]
420pub enum WireResponse {
421 Ok {
423 detail: String,
425 effective: Option<Value>,
427 },
428 Err {
430 error: String,
432 },
433}
434
435impl WireResponse {
436 #[must_use]
438 pub const fn is_ok(&self) -> bool {
439 matches!(self, Self::Ok { .. })
440 }
441
442 #[must_use]
444 pub fn detail(&self) -> &str {
445 match self {
446 Self::Ok { detail, .. } => detail,
447 Self::Err { .. } => "",
448 }
449 }
450
451 #[must_use]
453 pub fn effective(&self) -> Option<&Value> {
454 match self {
455 Self::Ok { effective, .. } => effective.as_ref(),
456 Self::Err { .. } => None,
457 }
458 }
459
460 #[must_use]
462 pub fn error(&self) -> Option<&str> {
463 match self {
464 Self::Err { error } => Some(error),
465 Self::Ok { .. } => None,
466 }
467 }
468}
469
470pub fn decode_response(line: &str) -> Result<WireResponse, CliError> {
479 let trimmed = line.trim_end_matches(['\n', '\r']);
480 if trimmed.is_empty() {
481 return Err(CliError::OperatorProtocol(
482 "operator host sent an empty response line".to_string(),
483 ));
484 }
485 let value: Value = serde_json::from_str(trimmed).map_err(|err| {
486 CliError::OperatorProtocol(format!("invalid operator response JSON: {err}"))
487 })?;
488 let Some(object) = value.as_object() else {
489 return Err(CliError::OperatorProtocol(
490 "operator response is not a JSON object".to_string(),
491 ));
492 };
493 match object.get("ok") {
494 Some(Value::Bool(true)) => Ok(WireResponse::Ok {
495 detail: object
496 .get("detail")
497 .and_then(Value::as_str)
498 .unwrap_or_default()
499 .to_string(),
500 effective: object.get("effective").cloned(),
501 }),
502 Some(Value::Bool(false)) => Ok(WireResponse::Err {
503 error: object
504 .get("error")
505 .and_then(Value::as_str)
506 .unwrap_or("operator host reported failure without an error message")
507 .to_string(),
508 }),
509 Some(_) => Err(CliError::OperatorProtocol(
510 "operator response 'ok' is not a boolean".to_string(),
511 )),
512 None => Err(CliError::OperatorProtocol(
513 "operator response is missing 'ok'".to_string(),
514 )),
515 }
516}
517
518#[must_use]
525pub fn render_json_sorted(value: Option<&Value>) -> String {
526 match value {
527 Some(Value::Object(map)) => {
528 let mut entries: Vec<(&String, &Value)> = map.iter().collect();
529 entries.sort_by(|a, b| a.0.cmp(b.0));
530 let inner = entries
531 .iter()
532 .map(|(key, value)| {
533 let key = serde_json::to_string(key).unwrap_or_default();
534 let value = serde_json::to_string(value).unwrap_or_default();
535 format!("{key}: {value}")
536 })
537 .collect::<Vec<_>>()
538 .join(", ");
539 format!("{{{inner}}}")
540 }
541 Some(other) => serde_json::to_string(other).unwrap_or_default(),
542 None => "{}".to_string(),
543 }
544}
545
546#[must_use]
550pub fn resolve_socket(env: &dyn EnvVars, cli_socket: Option<&str>) -> PathBuf {
551 let env_value = env.get(SOCKET_ENV);
552 let raw = match (non_empty(cli_socket), non_empty(env_value.as_deref())) {
553 (Some(cli), _) => cli,
554 (None, Some(value)) => value,
555 (None, None) => SOCKET_DEFAULT,
556 };
557 expand_tilde(env, raw)
558}
559
560fn non_empty(value: Option<&str>) -> Option<&str> {
562 value.filter(|value| !value.is_empty())
563}
564
565fn expand_tilde(env: &dyn EnvVars, raw: &str) -> PathBuf {
568 let home = env.get(HOME_ENV).filter(|home| !home.is_empty());
569 if raw == "~" {
570 if let Some(home) = home {
571 return PathBuf::from(home);
572 }
573 } else if let Some(rest) = raw.strip_prefix("~/") {
574 if let Some(home) = home {
575 return Path::new(&home).join(rest);
576 }
577 }
578 PathBuf::from(raw)
579}
580
581pub fn send_request(socket: &Path, request: &WireRequest) -> Result<WireResponse, CliError> {
592 let line = request.encode_line();
593 exchange_blocking(socket, line)
594}
595
596#[cfg(unix)]
599fn exchange_blocking(socket: &Path, line: String) -> Result<WireResponse, CliError> {
600 block_on_operator(exchange(socket, line))?
601}
602
603#[cfg(not(unix))]
606fn exchange_blocking(socket: &Path, line: String) -> Result<WireResponse, CliError> {
607 let _ = (socket, line);
608 Err(CliError::OperatorProtocol(UDS_UNSUPPORTED.to_string()))
609}
610
611#[cfg(unix)]
616fn block_on_operator<F: Future>(future: F) -> Result<F::Output, CliError> {
617 crate::block::block_on(future).map_err(|err| match err {
618 CliError::RuntimeNested => err,
619 other => CliError::OperatorProtocol(other.message()),
620 })
621}
622
623#[cfg(unix)]
626async fn exchange(socket: &Path, line: String) -> Result<WireResponse, CliError> {
627 let connect = tokio::time::timeout(REQUEST_TIMEOUT, UnixStream::connect(socket))
628 .await
629 .map_err(|_| {
630 CliError::OperatorProtocol(format!(
631 "timed out connecting to operator socket {}",
632 socket.display()
633 ))
634 })?;
635 let mut stream = connect.map_err(|err| {
636 CliError::OperatorProtocol(format!(
637 "cannot reach operator socket {}: {err}",
638 socket.display()
639 ))
640 })?;
641 stream.write_all(line.as_bytes()).await.map_err(|err| {
642 CliError::OperatorProtocol(format!(
643 "failed writing to operator socket {}: {err}",
644 socket.display()
645 ))
646 })?;
647 stream.flush().await.map_err(|err| {
648 CliError::OperatorProtocol(format!(
649 "failed flushing operator socket {}: {err}",
650 socket.display()
651 ))
652 })?;
653 let mut reader = BufReader::new(stream);
654 let mut buffer = Vec::new();
655 let read = tokio::time::timeout(REQUEST_TIMEOUT, reader.read_until(b'\n', &mut buffer))
656 .await
657 .map_err(|_| {
658 CliError::OperatorProtocol(format!(
659 "timed out waiting for a response from operator socket {}",
660 socket.display()
661 ))
662 })?;
663 let read = read.map_err(|err| {
664 CliError::OperatorProtocol(format!(
665 "failed reading from operator socket {}: {err}",
666 socket.display()
667 ))
668 })?;
669 if read == 0 {
670 return Err(CliError::OperatorProtocol(format!(
671 "no response from operator server at {}",
672 socket.display()
673 )));
674 }
675 let text = String::from_utf8(buffer).map_err(|err| {
676 CliError::OperatorProtocol(format!("operator response is not valid UTF-8: {err}"))
677 })?;
678 decode_response(&text)
679}