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
19pub trait Transport {
21 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#[derive(Clone)]
31pub struct HttpTransport {
32 agent: ureq::Agent,
33 trace: Option<Arc<Mutex<RpcTraceWriter>>>,
34}
35
36#[cfg(feature = "http")]
37#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
39pub enum RpcTraceMode {
40 #[default]
42 Full,
43 Summary,
45 RequestOnly,
47 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 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 pub fn with_trace_jsonl(path: &Path) -> Result<Self> {
81 Self::with_trace_jsonl_mode(path, RpcTraceMode::Full)
82 }
83
84 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}