Skip to main content

octra_sqlite/client/
transport.rs

1use super::error::Result;
2#[cfg(feature = "http")]
3use super::error::{Error, ErrorKind};
4#[cfg(feature = "http")]
5use crate::private_file::create_new;
6use serde_json::Value;
7#[cfg(feature = "http")]
8use serde_json::json;
9#[cfg(feature = "http")]
10use sha2::{Digest, Sha256};
11#[cfg(feature = "http")]
12use std::{
13    io::Write,
14    path::Path,
15    sync::{Arc, Mutex},
16    time::{Duration, SystemTime, UNIX_EPOCH},
17};
18
19/// Synchronous Octra JSON-RPC transport used by [`crate::Client`].
20pub trait Transport {
21    /// Call one JSON-RPC method and return its decoded result.
22    fn call(&self, rpc: &str, method: &str, params: Value) -> Result<Value>;
23}
24
25#[cfg(feature = "http")]
26const MAX_RPC_ATTEMPTS: usize = 4;
27
28#[cfg(feature = "http")]
29/// Default blocking HTTP transport with bounded read retries.
30#[derive(Clone)]
31pub struct HttpTransport {
32    agent: ureq::Agent,
33    trace: Option<Arc<Mutex<RpcTraceWriter>>>,
34}
35
36#[cfg(feature = "http")]
37/// Disclosure level for JSONL RPC traces.
38#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
39pub enum RpcTraceMode {
40    /// Record complete request and response bodies.
41    #[default]
42    Full,
43    /// Record method, hashes, sizes, status, and timing without bodies.
44    Summary,
45    /// Record the complete request but only response metadata.
46    RequestOnly,
47    /// Record request and response metadata without either body.
48    ResponseMeta,
49}
50
51#[cfg(feature = "http")]
52struct RpcTraceWriter {
53    file: std::fs::File,
54    sequence: u64,
55    mode: RpcTraceMode,
56}
57
58#[cfg(feature = "http")]
59impl Default for HttpTransport {
60    fn default() -> Self {
61        Self::new()
62    }
63}
64
65#[cfg(feature = "http")]
66impl HttpTransport {
67    /// Construct the default blocking HTTP transport.
68    pub fn new() -> Self {
69        let config = ureq::Agent::config_builder()
70            .timeout_global(Some(Duration::from_secs(30)))
71            .http_status_as_error(false)
72            .build();
73        Self {
74            agent: ureq::Agent::new_with_config(config),
75            trace: None,
76        }
77    }
78
79    /// Construct a transport that writes full JSONL RPC traces to a new file.
80    pub fn with_trace_jsonl(path: &Path) -> Result<Self> {
81        Self::with_trace_jsonl_mode(path, RpcTraceMode::Full)
82    }
83
84    /// Construct a transport that writes the selected JSONL trace mode.
85    pub fn with_trace_jsonl_mode(path: &Path, mode: RpcTraceMode) -> Result<Self> {
86        let mut transport = Self::new();
87        if let Some(parent) = path
88            .parent()
89            .filter(|parent| !parent.as_os_str().is_empty())
90        {
91            std::fs::create_dir_all(parent).map_err(|error| {
92                Error::with_kind(
93                    ErrorKind::Io,
94                    format!("creating RPC trace directory {}: {error}", parent.display()),
95                )
96            })?;
97        }
98        let file = create_new(path).map_err(|error| {
99            Error::with_kind(
100                ErrorKind::Io,
101                format!("opening RPC trace {}: {error}", path.display()),
102            )
103        })?;
104        transport.trace = Some(Arc::new(Mutex::new(RpcTraceWriter {
105            file,
106            sequence: 0,
107            mode,
108        })));
109        Ok(transport)
110    }
111
112    fn trace_rpc(
113        &self,
114        rpc: &str,
115        method: &str,
116        request: &Value,
117        response: Option<&Value>,
118        http_status: Option<u16>,
119        error: Option<&str>,
120    ) {
121        let Some(trace) = &self.trace else {
122            return;
123        };
124        let Ok(mut trace) = trace.lock() else {
125            return;
126        };
127        trace.sequence += 1;
128        let mut event = json!({
129            "schema": "octra-sqlite.rpc-trace.v1",
130            "mode": trace_mode_name(trace.mode),
131            "sequence": trace.sequence,
132            "timestamp_ms": unix_timestamp_ms(),
133            "rpc": rpc,
134            "method": method,
135            "http_status": http_status,
136            "ok": error.is_none(),
137            "error": error,
138        });
139        match trace.mode {
140            RpcTraceMode::Full => {
141                event["request"] = request.clone();
142                event["response"] = response.cloned().unwrap_or(Value::Null);
143                event["request_meta"] = trace_value_meta(request);
144                event["response_meta"] = trace_optional_value_meta(response);
145            }
146            RpcTraceMode::Summary => {
147                let request_meta = trace_value_meta(request);
148                let response_meta = trace_optional_value_meta(response);
149                event["request_bytes"] = request_meta["bytes"].clone();
150                event["request_sha256"] = request_meta["sha256"].clone();
151                event["response_bytes"] = response_meta["bytes"].clone();
152                event["response_sha256"] = response_meta["sha256"].clone();
153            }
154            RpcTraceMode::RequestOnly => {
155                event["request"] = request.clone();
156                event["response_meta"] = trace_optional_value_meta(response);
157            }
158            RpcTraceMode::ResponseMeta => {
159                event["request_meta"] = trace_value_meta(request);
160                event["response_meta"] = trace_optional_value_meta(response);
161            }
162        }
163        if serde_json::to_writer(&mut trace.file, &event).is_err() {
164            return;
165        }
166        let _ = trace.file.write_all(b"\n");
167    }
168}
169
170#[cfg(feature = "http")]
171impl Transport for HttpTransport {
172    fn call(&self, rpc: &str, method: &str, params: Value) -> Result<Value> {
173        let body = json!({
174            "jsonrpc": "2.0",
175            "id": 1,
176            "method": method,
177            "params": params,
178        });
179        let retryable = method != "octra_submit";
180        for attempt in 1..=MAX_RPC_ATTEMPTS {
181            let mut response = match self.agent.post(rpc).send_json(&body) {
182                Ok(response) => response,
183                Err(error) => {
184                    let message = format!("calling {method}: {error}");
185                    self.trace_rpc(rpc, method, &body, None, None, Some(&message));
186                    if should_retry_transport(attempt, retryable) {
187                        sleep_before_retry(attempt, None);
188                        continue;
189                    }
190                    return Err(Error::with_kind(ErrorKind::Transport, message));
191                }
192            };
193            let status = response.status();
194            let status_code = status.as_u16();
195            let retry_after = response
196                .headers()
197                .get("retry-after")
198                .and_then(|value| value.to_str().ok())
199                .and_then(parse_retry_after);
200            let text = match response.body_mut().read_to_string() {
201                Ok(text) => text,
202                Err(error) => {
203                    let message = format!("reading {method} response body: {error}");
204                    self.trace_rpc(rpc, method, &body, None, Some(status_code), Some(&message));
205                    if should_retry_http(attempt, retryable, status_code) {
206                        sleep_before_retry(attempt, retry_after);
207                        continue;
208                    }
209                    return Err(Error::with_kind(ErrorKind::Transport, message));
210                }
211            };
212            let payload: Value = match serde_json::from_str(&text) {
213                Ok(payload) => payload,
214                Err(error) => {
215                    let message = format!(
216                        "decoding {method} non-JSON response from HTTP {status}: {error}; body: {}",
217                        preview_text(&text, 512)
218                    );
219                    self.trace_rpc(rpc, method, &body, None, Some(status_code), Some(&message));
220                    if should_retry_http(attempt, retryable, status_code)
221                        || should_retry_non_json(attempt, retryable, &text)
222                    {
223                        sleep_before_retry(attempt, retry_after);
224                        continue;
225                    }
226                    let code = if status_code == 429 {
227                        "rpc_rate_limited"
228                    } else {
229                        "rpc_non_json"
230                    };
231                    return Err(Error::with_code(ErrorKind::Decode, code, message));
232                }
233            };
234            if !status.is_success() {
235                let message = format!(
236                    "{method} failed with HTTP {status}: {}",
237                    preview_value(&payload, 1024)
238                );
239                self.trace_rpc(
240                    rpc,
241                    method,
242                    &body,
243                    Some(&payload),
244                    Some(status_code),
245                    Some(&message),
246                );
247                if should_retry_http(attempt, retryable, status_code) {
248                    sleep_before_retry(attempt, retry_after);
249                    continue;
250                }
251                if status_code == 429 {
252                    return Err(Error::with_code(
253                        ErrorKind::Transport,
254                        "rpc_rate_limited",
255                        message,
256                    ));
257                }
258                return Err(Error::with_kind(ErrorKind::Transport, message));
259            }
260            if let Some(error) = payload.get("error") {
261                let message = format!("{method} failed: {}", preview_value(error, 1024));
262                self.trace_rpc(
263                    rpc,
264                    method,
265                    &body,
266                    Some(&payload),
267                    Some(status_code),
268                    Some(&message),
269                );
270                if should_retry_rpc_error(attempt, retryable, error) {
271                    sleep_before_retry(attempt, retry_after);
272                    continue;
273                }
274                if rpc_error_is_rate_limited(error) {
275                    return Err(Error::with_code(
276                        ErrorKind::Rpc,
277                        "rpc_rate_limited",
278                        message,
279                    ));
280                }
281                if let Some(code) = rpc_error_source_code(error) {
282                    return Err(Error::with_code(ErrorKind::Rpc, code, message));
283                }
284                return Err(Error::with_kind(ErrorKind::Rpc, message));
285            }
286            self.trace_rpc(rpc, method, &body, Some(&payload), Some(status_code), None);
287            return Ok(payload.get("result").cloned().unwrap_or(Value::Null));
288        }
289        Err(Error::with_kind(
290            ErrorKind::Transport,
291            format!("{method} failed after retry attempts"),
292        ))
293    }
294}
295
296#[cfg(feature = "http")]
297fn should_retry_transport(attempt: usize, retryable: bool) -> bool {
298    retryable && attempt < MAX_RPC_ATTEMPTS
299}
300
301#[cfg(feature = "http")]
302fn should_retry_http(attempt: usize, retryable: bool, status_code: u16) -> bool {
303    should_retry_transport(attempt, retryable)
304        && matches!(status_code, 408 | 425 | 429 | 500 | 502 | 503 | 504)
305}
306
307#[cfg(feature = "http")]
308fn should_retry_non_json(attempt: usize, retryable: bool, text: &str) -> bool {
309    let first = text.trim_start().chars().next();
310    should_retry_transport(attempt, retryable) && matches!(first, Some('<' | '\u{feff}'))
311}
312
313#[cfg(feature = "http")]
314fn should_retry_rpc_error(attempt: usize, retryable: bool, error: &Value) -> bool {
315    should_retry_transport(attempt, retryable) && rpc_error_is_rate_limited(error)
316}
317
318#[cfg(feature = "http")]
319fn rpc_error_is_rate_limited(error: &Value) -> bool {
320    let code = error.get("code").and_then(Value::as_i64);
321    let message = error
322        .get("message")
323        .and_then(Value::as_str)
324        .unwrap_or_default()
325        .to_ascii_lowercase();
326    matches!(code, Some(429) | Some(-32029))
327        || message.contains("too many requests")
328        || message.contains("rate limit")
329        || message.contains("temporarily unavailable")
330}
331
332#[cfg(feature = "http")]
333fn rpc_error_source_code(error: &Value) -> Option<&'static str> {
334    match error.get("code").and_then(Value::as_str) {
335        Some("storage_uninitialized") => return Some("storage_uninitialized"),
336        Some("auth_uninitialized") => return Some("auth_uninitialized"),
337        _ => {}
338    }
339    let message = error
340        .as_str()
341        .or_else(|| error.get("message").and_then(Value::as_str))?
342        .to_ascii_lowercase();
343    if message.contains("missing storage cache") || message.contains("storage_uninitialized") {
344        Some("storage_uninitialized")
345    } else if message.contains("auth_uninitialized") {
346        Some("auth_uninitialized")
347    } else {
348        None
349    }
350}
351
352#[cfg(feature = "http")]
353fn parse_retry_after(value: &str) -> Option<Duration> {
354    let seconds = value.trim().parse::<u64>().ok()?;
355    Some(Duration::from_secs(seconds.min(30)))
356}
357
358#[cfg(feature = "http")]
359fn sleep_before_retry(attempt: usize, retry_after: Option<Duration>) {
360    std::thread::sleep(retry_after.unwrap_or_else(|| {
361        let millis = 500_u64.saturating_mul(2_u64.saturating_pow((attempt - 1) as u32));
362        Duration::from_millis(millis.min(5_000))
363    }));
364}
365
366#[cfg(feature = "http")]
367fn preview_value(value: &Value, limit: usize) -> String {
368    preview_text(
369        &serde_json::to_string(value).unwrap_or_else(|_| value.to_string()),
370        limit,
371    )
372}
373
374#[cfg(feature = "http")]
375fn preview_text(text: &str, limit: usize) -> String {
376    let compact = text.split_whitespace().collect::<Vec<_>>().join(" ");
377    if compact.len() <= limit {
378        compact
379    } else {
380        format!(
381            "{}...<truncated {} bytes>",
382            truncate_to_char_boundary(&compact, limit),
383            compact.len()
384        )
385    }
386}
387
388#[cfg(feature = "http")]
389fn truncate_to_char_boundary(text: &str, limit: usize) -> &str {
390    if text.len() <= limit {
391        return text;
392    }
393    let mut end = 0usize;
394    for (index, _) in text.char_indices() {
395        if index > limit {
396            break;
397        }
398        end = index;
399    }
400    &text[..end]
401}
402
403#[cfg(feature = "http")]
404fn unix_timestamp_ms() -> u128 {
405    SystemTime::now()
406        .duration_since(UNIX_EPOCH)
407        .unwrap_or_default()
408        .as_millis()
409}
410
411#[cfg(feature = "http")]
412fn trace_mode_name(mode: RpcTraceMode) -> &'static str {
413    match mode {
414        RpcTraceMode::Full => "full",
415        RpcTraceMode::Summary => "summary",
416        RpcTraceMode::RequestOnly => "request_only",
417        RpcTraceMode::ResponseMeta => "response_meta",
418    }
419}
420
421#[cfg(feature = "http")]
422fn trace_optional_value_meta(value: Option<&Value>) -> Value {
423    match value {
424        Some(value) => trace_value_meta(value),
425        None => json!({
426            "bytes": null,
427            "sha256": null,
428        }),
429    }
430}
431
432#[cfg(feature = "http")]
433fn trace_value_meta(value: &Value) -> Value {
434    let bytes = serde_json::to_vec(value).unwrap_or_default();
435    json!({
436        "bytes": bytes.len(),
437        "sha256": hex::encode(Sha256::digest(&bytes)),
438    })
439}
440
441#[cfg(all(test, feature = "http"))]
442mod tests {
443    use super::*;
444    use serde_json::json;
445
446    #[test]
447    fn rpc_error_source_codes_are_parsed_at_the_transport_boundary() {
448        assert_eq!(
449            rpc_error_source_code(&json!({"message": "missing storage cache: octABC:0000"})),
450            Some("storage_uninitialized")
451        );
452        assert_eq!(
453            rpc_error_source_code(&json!("auth_uninitialized: auth_info is unavailable")),
454            Some("auth_uninitialized")
455        );
456        assert_eq!(
457            rpc_error_source_code(&json!({"code": "storage_uninitialized"})),
458            Some("storage_uninitialized")
459        );
460        assert_eq!(rpc_error_source_code(&json!({"message": "other"})), None);
461    }
462
463    #[test]
464    fn trace_writer_records_json_rpc_request_and_response() {
465        let path = std::env::temp_dir().join(format!(
466            "octra-sqlite-rpc-trace-{}.jsonl",
467            std::process::id()
468        ));
469        let _ = std::fs::remove_file(&path);
470        let transport = HttpTransport::with_trace_jsonl(&path).unwrap();
471        let request = json!({
472            "jsonrpc": "2.0",
473            "id": 1,
474            "method": "octra_circleViewAuth",
475            "params": ["octCircle"]
476        });
477        let response = json!({
478            "jsonrpc": "2.0",
479            "id": 1,
480            "result": "ok"
481        });
482
483        transport.trace_rpc(
484            "https://devnet.octrascan.io/rpc",
485            "octra_circleViewAuth",
486            &request,
487            Some(&response),
488            Some(200),
489            None,
490        );
491
492        let text = std::fs::read_to_string(&path).unwrap();
493        let event: Value = serde_json::from_str(text.trim()).unwrap();
494        assert_eq!(event["schema"], "octra-sqlite.rpc-trace.v1");
495        assert_eq!(event["mode"], "full");
496        assert_eq!(event["sequence"], 1);
497        assert_eq!(event["method"], "octra_circleViewAuth");
498        assert_eq!(event["request"], request);
499        assert_eq!(event["response"], response);
500        assert_eq!(
501            event["request_meta"]["bytes"],
502            serde_json::to_vec(&request).unwrap().len()
503        );
504        assert_eq!(event["request_meta"]["sha256"].as_str().unwrap().len(), 64);
505        assert!(event["error"].is_null());
506        let _ = std::fs::remove_file(path);
507    }
508
509    #[test]
510    fn trace_writer_summary_omits_rpc_bodies() {
511        let path = std::env::temp_dir().join(format!(
512            "octra-sqlite-rpc-trace-summary-{}.jsonl",
513            std::process::id()
514        ));
515        let _ = std::fs::remove_file(&path);
516        let transport = HttpTransport::with_trace_jsonl_mode(&path, RpcTraceMode::Summary).unwrap();
517        transport.trace_rpc(
518            "https://devnet.octrascan.io/rpc",
519            "octra_circleViewAuth",
520            &json!({"params": ["select * from artist;"]}),
521            Some(&json!({"result": {"rows": [[1]]}})),
522            Some(200),
523            None,
524        );
525
526        let text = std::fs::read_to_string(&path).unwrap();
527        let event: Value = serde_json::from_str(text.trim()).unwrap();
528        assert_eq!(event["mode"], "summary");
529        assert!(event.get("request").is_none());
530        assert!(event.get("response").is_none());
531        assert!(event["request_bytes"].as_u64().unwrap() > 0);
532        assert_eq!(event["request_sha256"].as_str().unwrap().len(), 64);
533        assert!(event["response_bytes"].as_u64().unwrap() > 0);
534        assert_eq!(event["response_sha256"].as_str().unwrap().len(), 64);
535        let _ = std::fs::remove_file(path);
536    }
537
538    #[test]
539    fn trace_writer_is_best_effort_after_open() {
540        let path = std::env::temp_dir().join(format!(
541            "octra-sqlite-rpc-trace-poison-{}.jsonl",
542            std::process::id()
543        ));
544        let _ = std::fs::remove_file(&path);
545        let transport = HttpTransport::with_trace_jsonl(&path).unwrap();
546        let trace = transport.trace.as_ref().unwrap().clone();
547        let _ = std::thread::spawn(move || {
548            let _guard = trace.lock().unwrap();
549            panic!("poison trace lock");
550        })
551        .join();
552
553        transport.trace_rpc(
554            "https://devnet.octrascan.io/rpc",
555            "octra_circleViewAuth",
556            &json!({"params": []}),
557            None,
558            None,
559            Some("real rpc error"),
560        );
561
562        let _ = std::fs::remove_file(path);
563    }
564
565    #[test]
566    fn trace_writer_refuses_to_replace_an_existing_file() {
567        let path = std::env::temp_dir().join(format!(
568            "octra-sqlite-rpc-trace-existing-{}.jsonl",
569            std::process::id()
570        ));
571        let _ = std::fs::remove_file(&path);
572        std::fs::write(&path, "keep\n").unwrap();
573        assert!(HttpTransport::with_trace_jsonl(&path).is_err());
574        assert_eq!(std::fs::read_to_string(&path).unwrap(), "keep\n");
575        let _ = std::fs::remove_file(path);
576    }
577
578    #[cfg(unix)]
579    #[test]
580    fn trace_writer_creates_private_files() {
581        use std::os::unix::fs::PermissionsExt;
582        let path = std::env::temp_dir().join(format!(
583            "octra-sqlite-rpc-trace-private-{}.jsonl",
584            std::process::id()
585        ));
586        let _ = std::fs::remove_file(&path);
587        let _transport = HttpTransport::with_trace_jsonl(&path).unwrap();
588        assert_eq!(
589            std::fs::metadata(&path).unwrap().permissions().mode() & 0o777,
590            0o600
591        );
592        let _ = std::fs::remove_file(path);
593    }
594
595    #[test]
596    fn retry_policy_is_mainnet_safe() {
597        assert!(should_retry_http(1, true, 429));
598        assert!(should_retry_http(1, true, 503));
599        assert!(!should_retry_http(MAX_RPC_ATTEMPTS, true, 429));
600        assert!(!should_retry_http(1, false, 429));
601        assert!(should_retry_rpc_error(
602            1,
603            true,
604            &json!({"code":429,"message":"Too Many Requests"})
605        ));
606        assert!(should_retry_rpc_error(
607            1,
608            true,
609            &json!({"message":"temporarily unavailable"})
610        ));
611        assert!(!should_retry_rpc_error(
612            1,
613            true,
614            &json!({"code":-32000,"message":"wasm export returned 1"})
615        ));
616    }
617
618    #[test]
619    fn response_previews_are_compact_and_utf8_safe() {
620        let text = "αβγδε ζηθικ λμνξο";
621        let preview = preview_text(text, 7);
622        assert!(preview.contains("<truncated"));
623        assert!(preview.is_char_boundary(preview.len()));
624        assert!(preview_text("<html> too many requests </html>", 512).contains("<html>"));
625        assert!(should_retry_non_json(1, true, "  <html>429</html>"));
626    }
627}