Skip to main content

macula_rust/station_link/
dht.rs

1//! The DHT, as macula 12's facade reaches it: station procedures on the zero
2//! realm, targeting the connected station, carrying record wire bytes. Every
3//! record found is verified here before it is handed on, and one that does
4//! not verify is dropped and counted.
5
6use crate::cbor::Value;
7use crate::record::{self, RecordType, Verified};
8
9use super::{now_ms, Call, Link, LinkError};
10
11impl Link {
12    /// Stores a signed record, as its wire bytes, in the station's DHT.
13    pub async fn put_record(&self, wire: &[u8]) -> Result<(), LinkError> {
14        let result = self
15            .dht_call("_dht.put_record", Value::Bytes(wire.to_vec()))
16            .await?;
17        match result {
18            Value::Text(t) if t == "ok" => Ok(()),
19            other => Err(LinkError::UnexpectedReply(format!(
20                "put_record answered {other:?}"
21            ))),
22        }
23    }
24
25    /// The record stored under `key`, verified.
26    pub async fn find_record(&self, key: &[u8; 32]) -> Result<Verified, LinkError> {
27        match self.dht_call("_dht.find_record", key_payload(key)).await? {
28            Value::Text(t) if t == "not_found" => Err(LinkError::RecordNotFound),
29            Value::Bytes(wire) => Ok(record::verify(&wire, self.profile(), now_ms())?),
30            other => Err(LinkError::UnexpectedReply(format!(
31                "find_record answered {other:?}"
32            ))),
33        }
34    }
35
36    /// Every record stored under `key` that verifies, and how many the
37    /// station returned that did not.
38    pub async fn find_records(&self, key: &[u8; 32]) -> Result<(Vec<Verified>, usize), LinkError> {
39        self.verified_list("_dht.find_records", key_payload(key))
40            .await
41    }
42
43    /// Every record of type `t` the station holds that verifies, and how many
44    /// it returned that did not.
45    pub async fn find_records_by_type(
46        &self,
47        t: RecordType,
48    ) -> Result<(Vec<Verified>, usize), LinkError> {
49        let payload = Value::Map(vec![(Value::text("type"), Value::Int(i128::from(t.0)))]);
50        self.verified_list("_dht.find_records_by_type", payload)
51            .await
52    }
53
54    async fn verified_list(
55        &self,
56        procedure: &str,
57        payload: Value,
58    ) -> Result<(Vec<Verified>, usize), LinkError> {
59        let Value::List(items) = self.dht_call(procedure, payload).await? else {
60            return Err(LinkError::UnexpectedReply(format!(
61                "{procedure} answered no list"
62            )));
63        };
64        let now = now_ms();
65        let mut verified = Vec::with_capacity(items.len());
66        let mut dropped = 0;
67        for item in items {
68            match item {
69                Value::Bytes(wire) => match record::verify(&wire, self.profile(), now) {
70                    Ok(r) => verified.push(r),
71                    Err(_) => dropped += 1,
72                },
73                _ => dropped += 1,
74            }
75        }
76        Ok((verified, dropped))
77    }
78
79    async fn dht_call(&self, procedure: &str, payload: Value) -> Result<Value, LinkError> {
80        self.call(Call {
81            procedure: procedure.to_string(),
82            payload,
83            ..Call::default()
84        })
85        .await
86    }
87}
88
89fn key_payload(key: &[u8; 32]) -> Value {
90    Value::Map(vec![(Value::text("key"), Value::Bytes(key.to_vec()))])
91}