onlyne_client/ops/local_cli/
local.rs1use 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#[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 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 pub fn query_roles(&self, args: QueryRolesArgs) -> ClientOp {
64 ClientOp::QueryRoles(args)
65 }
66
67 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 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 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}