dig_node_control_interface/traits.rs
1//! The two contract traits: the client-facing call builder/parser and the node-facing handler.
2//!
3//! * [`ControlCall`] binds a typed params struct to its [`ControlMethod`] and its typed result — so
4//! a caller writes `client.request(&SetCapParams { cap_bytes })` and gets back a `SetCapResult`,
5//! never a stringly-typed `Value`.
6//! * [`ControlClient`] is what a CLIENT depends on: build a JSON-RPC request from a typed call, and
7//! parse a response back into the typed result (or a [`ControlError`]). Pure — no transport; the
8//! consumer carries the bytes over dig-ipc / loopback-mTLS itself.
9//! * [`ControlHandler`] is what a NODE implements to SERVE the surface: one typed method per control
10//! method, plus a provided [`dispatch`](ControlHandler::dispatch) that routes a raw request to the
11//! right method — the single anti-drift seam the conformance KATs exercise.
12
13use async_trait::async_trait;
14use serde::de::DeserializeOwned;
15use serde::Serialize;
16use serde_json::Value;
17
18use crate::envelope::{JsonRpcRequest, JsonRpcResponse, RequestId};
19use crate::error::{ControlError, ControlErrorCode};
20use crate::method::ControlMethod;
21use crate::params;
22use crate::results;
23
24/// A typed control call: a params struct that knows its [`ControlMethod`] and its result type.
25///
26/// Implemented by every struct in [`crate::params`]; this is what makes
27/// [`ControlClient::parse_response`] return the right typed result for each method at compile time.
28pub trait ControlCall: Serialize {
29 /// The wire method this call invokes.
30 const METHOD: ControlMethod;
31 /// The typed result this call returns on success.
32 type Output: DeserializeOwned;
33}
34
35/// Serialize a typed call's params into a JSON object (`{}` for a no-param call, never `null`).
36fn params_value<C: ControlCall>(call: &C) -> Value {
37 match serde_json::to_value(call) {
38 Ok(Value::Null) => Value::Object(Default::default()),
39 Ok(v) => v,
40 Err(_) => Value::Object(Default::default()),
41 }
42}
43
44/// Build the JSON-RPC request envelope for a typed control call. Pure.
45pub fn build_request<C: ControlCall>(id: RequestId, call: &C) -> JsonRpcRequest {
46 JsonRpcRequest::new(id, C::METHOD.name(), params_value(call))
47}
48
49/// Parse a JSON-RPC response into a typed result, or the [`ControlError`] it carried. Pure.
50pub fn parse_response<C: ControlCall>(
51 response: JsonRpcResponse,
52) -> Result<C::Output, ControlError> {
53 let value = response.into_result()?;
54 serde_json::from_value(value).map_err(|e| {
55 ControlError::of(
56 ControlErrorCode::ControlError,
57 format!("failed to parse {} result: {e}", C::METHOD.name()),
58 )
59 })
60}
61
62/// The client-facing half of the contract: turn typed calls into requests and responses back into
63/// typed results.
64///
65/// The default implementations cover every client; a consumer implements this trait only to
66/// customise request construction (e.g. attaching the control token in a bespoke way). The blanket
67/// [`DefaultControlClient`] gives callers the standard behaviour for free.
68pub trait ControlClient {
69 /// Build the request envelope for a typed call with the given request `id`.
70 fn build_request<C: ControlCall>(&self, id: RequestId, call: &C) -> JsonRpcRequest {
71 build_request(id, call)
72 }
73
74 /// Parse a response envelope into the typed result for call type `C`.
75 fn parse_response<C: ControlCall>(
76 &self,
77 response: JsonRpcResponse,
78 ) -> Result<C::Output, ControlError> {
79 parse_response::<C>(response)
80 }
81}
82
83/// The standard, zero-configuration [`ControlClient`] using the default request/response behaviour.
84#[derive(Debug, Clone, Copy, Default)]
85pub struct DefaultControlClient;
86
87impl ControlClient for DefaultControlClient {}
88
89/// The node-facing half of the contract: a running node implements this to SERVE the control
90/// surface. Each method is typed to the catalog's params/results; the provided
91/// [`dispatch`](ControlHandler::dispatch) routes a raw [`JsonRpcRequest`] to the right method so a
92/// server needs only one entry point and can never mis-route.
93///
94/// Open/proxied shapes (the updater beacon status, the pairing list, the peer-pool snapshot) return
95/// [`Value`] rather than a frozen struct, matching the catalog's [`ControlCall::Output`] for those
96/// methods.
97#[async_trait]
98pub trait ControlHandler: Sync {
99 /// `control.status`
100 async fn status(&self) -> Result<results::StatusResult, ControlError>;
101 /// `control.config.get`
102 async fn config_get(&self) -> Result<results::ConfigResult, ControlError>;
103 /// `control.config.setUpstream`
104 async fn config_set_upstream(
105 &self,
106 params: params::SetUpstreamParams,
107 ) -> Result<results::SetUpstreamResult, ControlError>;
108 /// `control.log.setLevel`
109 async fn log_set_level(
110 &self,
111 params: params::SetLevelParams,
112 ) -> Result<results::SetLevelResult, ControlError>;
113 /// `control.cache.get`
114 async fn cache_get(&self) -> Result<results::CacheView, ControlError>;
115 /// `control.cache.setCap`
116 async fn cache_set_cap(
117 &self,
118 params: params::SetCapParams,
119 ) -> Result<results::SetCapResult, ControlError>;
120 /// `control.cache.clear`
121 async fn cache_clear(&self) -> Result<results::CacheClearResult, ControlError>;
122 /// `control.hostedStores.list`
123 async fn hosted_stores_list(&self) -> Result<results::HostedStoresListResult, ControlError>;
124 /// `control.hostedStores.pin`
125 async fn hosted_stores_pin(
126 &self,
127 params: params::PinParams,
128 ) -> Result<results::PinResult, ControlError>;
129 /// `control.hostedStores.unpin`
130 async fn hosted_stores_unpin(
131 &self,
132 params: params::UnpinParams,
133 ) -> Result<results::UnpinResult, ControlError>;
134 /// `control.hostedStores.status`
135 async fn hosted_stores_status(
136 &self,
137 params: params::HostedStoreStatusParams,
138 ) -> Result<results::HostedStoreStatusResult, ControlError>;
139 /// `control.sync.status`
140 async fn sync_status(&self) -> Result<results::SyncStatusResult, ControlError>;
141 /// `control.sync.trigger`
142 async fn sync_trigger(
143 &self,
144 params: params::SyncTriggerParams,
145 ) -> Result<results::SyncTriggerResult, ControlError>;
146 /// `control.updater.status`
147 async fn updater_status(&self) -> Result<Value, ControlError>;
148 /// `control.updater.setChannel`
149 async fn updater_set_channel(
150 &self,
151 params: params::SetChannelParams,
152 ) -> Result<Value, ControlError>;
153 /// `control.updater.pause`
154 async fn updater_pause(&self, params: params::PauseParams) -> Result<Value, ControlError>;
155 /// `control.updater.resume`
156 async fn updater_resume(&self) -> Result<Value, ControlError>;
157 /// `control.updater.checkNow`
158 async fn updater_check_now(&self) -> Result<Value, ControlError>;
159 /// `control.pairing.list`
160 async fn pairing_list(&self) -> Result<Value, ControlError>;
161 /// `control.pairing.approve`
162 async fn pairing_approve(
163 &self,
164 params: params::ApproveParams,
165 ) -> Result<results::PairingApproveResult, ControlError>;
166 /// `control.pairing.revoke`
167 async fn pairing_revoke(
168 &self,
169 params: params::RevokeParams,
170 ) -> Result<results::PairingRevokeResult, ControlError>;
171 /// `control.peerStatus`
172 async fn peer_status(&self) -> Result<Value, ControlError>;
173 /// `control.peers.connect`
174 async fn peers_connect(
175 &self,
176 params: params::PeersConnectParams,
177 ) -> Result<results::PeersConnectResult, ControlError>;
178 /// `control.peers.disconnect`
179 async fn peers_disconnect(
180 &self,
181 params: params::PeersDisconnectParams,
182 ) -> Result<results::PeersDisconnectResult, ControlError>;
183 /// `control.subscribe`
184 async fn subscribe(
185 &self,
186 params: params::SubscribeParams,
187 ) -> Result<results::SubscribeResult, ControlError>;
188 /// `control.unsubscribe`
189 async fn unsubscribe(
190 &self,
191 params: params::UnsubscribeParams,
192 ) -> Result<results::UnsubscribeResult, ControlError>;
193 /// `control.listSubscriptions`
194 async fn list_subscriptions(&self) -> Result<results::ListSubscriptionsResult, ControlError>;
195 /// `control.wallet.balance` (READ-only)
196 async fn wallet_balance(
197 &self,
198 params: params::WalletBalanceParams,
199 ) -> Result<results::WalletBalanceResult, ControlError>;
200 /// `control.wallet.coins` (READ-only, OPEN)
201 ///
202 /// An empty `coins` list MUST mean "a chain was consulted and this address holds nothing".
203 /// A read that could not consult a chain MUST return the matching catalogued error instead.
204 async fn wallet_coins(
205 &self,
206 params: params::WalletCoinsParams,
207 ) -> Result<results::WalletCoinsResult, ControlError>;
208 /// `control.wallet.coinById` (READ-only, OPEN)
209 ///
210 /// `Ok(coin: None)` MUST mean "a chain was consulted and holds no such coin". A read that could
211 /// not consult a chain MUST return the matching catalogued error instead — a caller that cannot
212 /// tell those apart reports a spent mint as pending forever.
213 ///
214 /// The params are validated at DESERIALIZATION (lowercase 64-hex, `0x` stripped), so any path
215 /// that decodes `WalletCoinByIdParams` refuses malformed ids as `INVALID_PARAMS` before this
216 /// method is called.
217 async fn wallet_coin_by_id(
218 &self,
219 params: params::WalletCoinByIdParams,
220 ) -> Result<results::WalletCoinByIdResult, ControlError>;
221 /// `control.wallet.coinSpend` (READ-only, OPEN)
222 ///
223 /// `Ok(spend: None)` MUST mean "a chain was consulted and holds no spend of that coin" — the
224 /// coin is unspent, or unknown. A read that could not consult a chain MUST return the matching
225 /// catalogued error instead: a caller following a singleton forward reads "no spend" as *this is
226 /// the tip* and stops walking, so a failure disguised as absence produces a spend built against
227 /// a superseded singleton.
228 ///
229 /// A returned spend's `puzzle_reveal` MUST tree-hash to the spent coin's own `puzzle_hash`, and
230 /// the implementation MUST fail closed — an error, never an unverified reveal — when it does not
231 /// or when the reveal will not parse. The reveal comes from a peer, and a peer can lie.
232 ///
233 /// The params are validated at DESERIALIZATION (lowercase 64-hex, `0x` stripped), so any path
234 /// that decodes `WalletCoinSpendParams` refuses malformed ids as `INVALID_PARAMS` before this
235 /// method is called.
236 async fn wallet_coin_spend(
237 &self,
238 params: params::WalletCoinSpendParams,
239 ) -> Result<results::WalletCoinSpendResult, ControlError>;
240 /// `control.wallet.coinsByParent` (READ-only, OPEN)
241 ///
242 /// Returns the parent's DIRECT children and nothing further. An implementation MUST NOT recurse:
243 /// a transitive walk over caller-supplied input is unbounded work the caller cannot bound, and a
244 /// partial walk returned as a complete one is a lineage with a silent hole in it.
245 ///
246 /// An empty list MUST mean "a chain was consulted and this parent created no known children".
247 /// A read that could not consult a chain MUST return the matching catalogued error instead.
248 ///
249 /// The answer is ONE PAGE. An implementation MUST return at most
250 /// `params.effective_limit()` records, in ASCENDING `coin_id` order, starting strictly after
251 /// `params.after_coin_id` when one is given; it MUST set `complete` to whether the page carries
252 /// the last child; and it MUST set `cursor` to the last record it actually returned (`None` for
253 /// an empty page). It MUST NOT report `complete: true` on a page it truncated — a caller reads
254 /// that as the end of a lineage branch. The params are validated at DESERIALIZATION, so an
255 /// out-of-range page size is refused as `INVALID_PARAMS` before this method is called.
256 ///
257 /// Every record MUST report `asset: None`: naming a coin by its parent classifies nothing, and
258 /// asserting a class this read never verified is a claim a caller would then spend against.
259 async fn wallet_coins_by_parent(
260 &self,
261 params: params::WalletCoinsByParentParams,
262 ) -> Result<results::WalletCoinsByParentResult, ControlError>;
263 /// `control.wallet.arrivals` (READ-only, TOKEN-GATED)
264 ///
265 /// Gated although it is a read: the caller supplies only a cursor, so the answer names this
266 /// node's OWN watched puzzle hashes and the receive history behind them.
267 ///
268 /// Every returned row MUST be a CONFIRMED arrival that the node itself judged: above its arrival
269 /// baseline, not previously reported, and not the wallet's own change. An implementation MUST NOT
270 /// emit a mempool sighting here, and MUST answer an empty page rather than an error when it has
271 /// no baseline — "nothing arrived" is the honest answer from a wallet that cannot yet tell
272 /// history from news.
273 ///
274 /// `cursor` MUST be the position of the last row actually returned (or the caller's `after_seq`
275 /// for an empty page) and MUST NOT be `latest`; see
276 /// [`WalletArrivalsResult::latest`](results::WalletArrivalsResult::latest).
277 async fn wallet_arrivals(
278 &self,
279 params: params::WalletArrivalsParams,
280 ) -> Result<results::WalletArrivalsResult, ControlError>;
281 /// `control.wallet.peak` (READ-only, OPEN)
282 async fn wallet_peak(&self) -> Result<results::WalletPeakResult, ControlError>;
283 /// `control.peerCounts` (READ-only, OPEN)
284 ///
285 /// `dig_peer_count` MUST be dig-node-core's `connected_peers` — the same figure
286 /// `control.peerStatus` reports — and `chia_peer_count` MUST be the SAME observation
287 /// `wallet_sync_status` reports, served from ONE source so the two answers agree. `None` means
288 /// the count cannot be observed; a network that is not running is UNKNOWN, never `Some(0)`.
289 async fn peer_counts(&self) -> Result<results::PeerCountsResult, ControlError>;
290 /// `control.wallet.syncStatus` (READ-only, OPEN)
291 ///
292 /// `WalletSyncPhase::Synced` MUST require BOTH that the initial catch-up completed and that at
293 /// least one Chia peer connection is live now, which makes it strictly stronger than
294 /// `WalletPeakResult::synced`. `peak_height` MUST be the node's OWN replica's height or `None`,
295 /// never an oracle's, and `chia_peer_count` counts CHIA full-node peers -- never DIG peers.
296 async fn wallet_sync_status(&self) -> Result<results::WalletSyncStatusResult, ControlError>;
297 /// `control.wallet.broadcast` (TOKEN-GATED)
298 ///
299 /// Pushes an ALREADY-SIGNED bundle: the implementation never signs, and never receives anything
300 /// it could sign with (§908). A mempool refusal is `Ok` with `accepted: false`; failing to
301 /// reach a mempool is `Err`.
302 async fn wallet_broadcast(
303 &self,
304 params: params::WalletBroadcastParams,
305 ) -> Result<results::WalletBroadcastResult, ControlError>;
306 /// `pairing.request` (OPEN)
307 async fn pairing_request(
308 &self,
309 params: params::RequestParams,
310 ) -> Result<results::PairingRequestResult, ControlError>;
311 /// `pairing.poll` (OPEN)
312 async fn pairing_poll(
313 &self,
314 params: params::PollParams,
315 ) -> Result<results::PairingPollResult, ControlError>;
316
317 /// Route a raw JSON-RPC request to the right typed method and build the response envelope.
318 ///
319 /// Deserializes the params for methods that take them, calls the handler, and serializes the
320 /// typed result. An unknown method → `METHOD_NOT_FOUND`; malformed params → `INVALID_PARAMS`.
321 /// This is the single seam a server dispatches through — the KATs exercise it end-to-end.
322 async fn dispatch(&self, request: JsonRpcRequest) -> JsonRpcResponse {
323 let id = request.id.clone();
324 let Some(method) = ControlMethod::from_name(&request.method) else {
325 return JsonRpcResponse::error(
326 id,
327 ControlError::of(
328 ControlErrorCode::MethodNotFound,
329 format!("unknown control method: {}", request.method),
330 ),
331 );
332 };
333 match self.dispatch_method(method, request.params).await {
334 Ok(result) => JsonRpcResponse::success(id, result),
335 Err(err) => JsonRpcResponse::error(id, err),
336 }
337 }
338
339 /// Route to the typed method by [`ControlMethod`], returning the result as a [`Value`]. Split
340 /// from [`dispatch`](ControlHandler::dispatch) so the envelope wrapping stays in one place.
341 #[doc(hidden)]
342 async fn dispatch_method(
343 &self,
344 method: ControlMethod,
345 params: Value,
346 ) -> Result<Value, ControlError> {
347 /// Deserialize a method's params, mapping a shape error to `INVALID_PARAMS`.
348 fn decode<T: DeserializeOwned>(params: Value) -> Result<T, ControlError> {
349 serde_json::from_value(params)
350 .map_err(|e| ControlError::of(ControlErrorCode::InvalidParams, e.to_string()))
351 }
352 /// Serialize a typed result to a `Value` (infallible for our derive-Serialize results).
353 fn encode<T: Serialize>(value: T) -> Result<Value, ControlError> {
354 serde_json::to_value(value)
355 .map_err(|e| ControlError::of(ControlErrorCode::ControlError, e.to_string()))
356 }
357 match method {
358 ControlMethod::Status => encode(self.status().await?),
359 ControlMethod::ConfigGet => encode(self.config_get().await?),
360 ControlMethod::ConfigSetUpstream => {
361 encode(self.config_set_upstream(decode(params)?).await?)
362 }
363 ControlMethod::LogSetLevel => encode(self.log_set_level(decode(params)?).await?),
364 ControlMethod::CacheGet => encode(self.cache_get().await?),
365 ControlMethod::CacheSetCap => encode(self.cache_set_cap(decode(params)?).await?),
366 ControlMethod::CacheClear => encode(self.cache_clear().await?),
367 ControlMethod::HostedStoresList => encode(self.hosted_stores_list().await?),
368 ControlMethod::HostedStoresPin => {
369 encode(self.hosted_stores_pin(decode(params)?).await?)
370 }
371 ControlMethod::HostedStoresUnpin => {
372 encode(self.hosted_stores_unpin(decode(params)?).await?)
373 }
374 ControlMethod::HostedStoresStatus => {
375 encode(self.hosted_stores_status(decode(params)?).await?)
376 }
377 ControlMethod::SyncStatus => encode(self.sync_status().await?),
378 ControlMethod::SyncTrigger => encode(self.sync_trigger(decode(params)?).await?),
379 ControlMethod::UpdaterStatus => self.updater_status().await,
380 ControlMethod::UpdaterSetChannel => self.updater_set_channel(decode(params)?).await,
381 ControlMethod::UpdaterPause => self.updater_pause(decode(params)?).await,
382 ControlMethod::UpdaterResume => self.updater_resume().await,
383 ControlMethod::UpdaterCheckNow => self.updater_check_now().await,
384 ControlMethod::PairingList => self.pairing_list().await,
385 ControlMethod::PairingApprove => encode(self.pairing_approve(decode(params)?).await?),
386 ControlMethod::PairingRevoke => encode(self.pairing_revoke(decode(params)?).await?),
387 ControlMethod::PeerStatus => self.peer_status().await,
388 ControlMethod::PeerCounts => encode(self.peer_counts().await?),
389 ControlMethod::PeersConnect => encode(self.peers_connect(decode(params)?).await?),
390 ControlMethod::PeersDisconnect => encode(self.peers_disconnect(decode(params)?).await?),
391 ControlMethod::Subscribe => encode(self.subscribe(decode(params)?).await?),
392 ControlMethod::Unsubscribe => encode(self.unsubscribe(decode(params)?).await?),
393 ControlMethod::ListSubscriptions => encode(self.list_subscriptions().await?),
394 ControlMethod::WalletBalance => encode(self.wallet_balance(decode(params)?).await?),
395 ControlMethod::WalletCoins => encode(self.wallet_coins(decode(params)?).await?),
396 // Re-validated here idempotently; deserialization already enforced the same rule.
397 ControlMethod::WalletCoinById => {
398 let params: params::WalletCoinByIdParams = decode(params)?;
399 encode(self.wallet_coin_by_id(params.validated()?).await?)
400 }
401 // Re-validated here idempotently; deserialization already enforced the same rule.
402 ControlMethod::WalletCoinSpend => {
403 let params: params::WalletCoinSpendParams = decode(params)?;
404 encode(self.wallet_coin_spend(params.validated()?).await?)
405 }
406 ControlMethod::WalletCoinsByParent => {
407 let params: params::WalletCoinsByParentParams = decode(params)?;
408 encode(self.wallet_coins_by_parent(params.validated()?).await?)
409 }
410 ControlMethod::WalletArrivals => encode(self.wallet_arrivals(decode(params)?).await?),
411 ControlMethod::WalletPeak => encode(self.wallet_peak().await?),
412 ControlMethod::WalletSyncStatus => encode(self.wallet_sync_status().await?),
413 ControlMethod::WalletBroadcast => encode(self.wallet_broadcast(decode(params)?).await?),
414 ControlMethod::PairingRequest => encode(self.pairing_request(decode(params)?).await?),
415 ControlMethod::PairingPoll => encode(self.pairing_poll(decode(params)?).await?),
416 }
417 }
418}