Skip to main content

onlyne_client/ops/local_cli/
local.rs

1//! Local CLI handlers for the socket verb surface.
2
3use crate::runtime::intent::{IntentMachine, stamp_op_id};
4use anyhow::{Result, anyhow};
5use onlyne_proto::{
6    AdapterMsg, ClientOp, ControlArgs, Envelope, HistoryArgs, LedgerQuery, PluginOp,
7    QueryFaultsArgs, QueryRolesArgs, QuerySessionsArgs, ResBody, Subscribe,
8};
9
10/// Local CLI handlers share the intent queue with plugin-originated sends.
11#[derive(Clone)]
12pub struct LocalCli {
13    pub intents: IntentMachine,
14    pub role: String,
15}
16
17impl LocalCli {
18    pub fn new(intents: IntentMachine) -> Self {
19        let role = intents
20            .store
21            .config("role")
22            .ok()
23            .flatten()
24            .unwrap_or_default();
25        Self { intents, role }
26    }
27
28    pub fn with_role(intents: IntentMachine, role: impl Into<String>) -> Self {
29        Self {
30            intents,
31            role: role.into(),
32        }
33    }
34
35    /// Queue one outbound send and answer the stamped envelope.
36    ///
37    /// The queue keys its row with an `op_id`; a note arrives without one, so
38    /// the stamp lands before the envelope becomes the `ClientOp` the caller
39    /// writes, and the frame on the wire carries the id the row stored.
40    pub fn send(&self, mut envelope: Envelope) -> Result<ClientOp> {
41        stamp_op_id(&mut envelope);
42        self.intents.enqueue(&envelope)?;
43        Ok(ClientOp::Send(Box::new(envelope)))
44    }
45
46    pub fn reply(&self, envelope: Envelope) -> Result<ClientOp> {
47        self.send(envelope)
48    }
49    pub fn complete(&self, envelope: Envelope) -> Result<ClientOp> {
50        self.send(envelope)
51    }
52    pub fn handoff(&self, envelope: Envelope) -> Result<ClientOp> {
53        self.send(envelope)
54    }
55    pub fn control(&self, args: ControlArgs) -> ClientOp {
56        ClientOp::Control(args)
57    }
58    pub fn query_sessions(&self, args: QuerySessionsArgs) -> ClientOp {
59        ClientOp::QuerySessions(args)
60    }
61
62    /// Build the server query shape for callers that need a wire request.
63    pub fn query_roles(&self, args: QueryRolesArgs) -> ClientOp {
64        ClientOp::QueryRoles(args)
65    }
66
67    /// Answer the role query from durable local prose cache. `QueryRolesArgs`
68    /// is the only role-query type in onlyne-proto and has the field `role`.
69    pub fn query_roles_local(&self, args: &QueryRolesArgs) -> Result<ResBody> {
70        let role = args.role.as_deref().unwrap_or(&self.role);
71        let Some((prose, spec_hash)) = self.intents.store.prose(role)? else {
72            return Ok(ResBody::ok(serde_json::json!({"roles": []})));
73        };
74        Ok(ResBody::ok(serde_json::json!({
75            "roles": [{"name": role, "role": role, "prose": prose, "spec_hash": spec_hash}]
76        })))
77    }
78
79    /// Client-surface export used by `cluster export-prose`.
80    pub fn export_prose(&self) -> Result<ResBody> {
81        let query = QueryRolesArgs {
82            role: Some(self.role.clone()),
83        };
84        self.query_roles_local(&query)
85    }
86
87    pub fn query_ledger(&self, args: LedgerQuery) -> ClientOp {
88        ClientOp::QueryLedger(args)
89    }
90    pub fn subscribe(&self, args: Subscribe) -> ClientOp {
91        ClientOp::Subscribe(args)
92    }
93    pub fn history(&self, args: HistoryArgs) -> Result<ClientOp> {
94        if args.limit == 0 {
95            return Err(anyhow!("history limit must be positive"));
96        }
97        Ok(ClientOp::QueryFaults(QueryFaultsArgs {
98            role: None,
99            task_id: args.task_id,
100            kind: args.kind,
101            open_only: false,
102            limit: args.limit,
103        }))
104    }
105
106    /// Queue one plugin `send` and answer the key its intent row carries.
107    ///
108    /// A note arrives without an `op_id`, so the stamp happens before the
109    /// reply is built: the plugin reads the id the queue stored rather than
110    /// `null`, and the stored envelope is the one that replays.
111    pub fn offline_send(&self, envelope: &Envelope) -> Result<ResBody> {
112        let mut stamped = envelope.clone();
113        let op_id = stamp_op_id(&mut stamped);
114        self.intents.enqueue(&stamped)?;
115        Ok(ResBody::ok(
116            serde_json::json!({"queued": true, "op_id": op_id}),
117        ))
118    }
119
120    pub async fn handle(&self, message: AdapterMsg) -> Result<ResBody> {
121        match message {
122            AdapterMsg::Plugin(PluginOp::Send(envelope)) => self.offline_send(&envelope),
123            AdapterMsg::Plugin(PluginOp::Detach(_)) => Ok(ResBody::ok(serde_json::Value::Null)),
124            AdapterMsg::Plugin(PluginOp::Hello(_)) => Ok(ResBody::err(
125                onlyne_proto::ErrorCode::Invalid,
126                "hello already completed",
127                Some("op".into()),
128            )),
129            _ => Ok(ResBody::err(
130                onlyne_proto::ErrorCode::UnknownOp,
131                "unsupported local cli operation",
132                Some("op".into()),
133            )),
134        }
135    }
136}
137
138pub fn map_send(envelope: Envelope) -> Result<ClientOp> {
139    envelope.validate().map_err(|e| anyhow!(e.to_string()))?;
140    Ok(ClientOp::Send(Box::new(envelope)))
141}
142pub fn map_reply(envelope: Envelope) -> Result<ClientOp> {
143    map_send(envelope)
144}
145pub fn map_complete(envelope: Envelope) -> Result<ClientOp> {
146    map_send(envelope)
147}
148pub fn map_handoff(envelope: Envelope) -> Result<ClientOp> {
149    map_send(envelope)
150}
151pub fn map_control(args: ControlArgs) -> ClientOp {
152    ClientOp::Control(args)
153}
154pub fn map_query_sessions(args: QuerySessionsArgs) -> ClientOp {
155    ClientOp::QuerySessions(args)
156}
157pub fn map_query_roles(args: QueryRolesArgs) -> ClientOp {
158    ClientOp::QueryRoles(args)
159}
160pub fn map_query_ledger(args: LedgerQuery) -> ClientOp {
161    ClientOp::QueryLedger(args)
162}
163pub fn map_subscribe(args: Subscribe) -> ClientOp {
164    ClientOp::Subscribe(args)
165}
166pub fn map_history(args: HistoryArgs) -> ClientOp {
167    ClientOp::QueryFaults(QueryFaultsArgs {
168        role: None,
169        task_id: args.task_id,
170        kind: args.kind,
171        open_only: false,
172        limit: args.limit,
173    })
174}