Skip to main content

surreal_client/engines/
debug.rs

1//! Debug wrapper for Engine trait that logs all RPC operations
2
3use async_trait::async_trait;
4use ciborium::Value as CborValue;
5use serde_json::Value;
6use tokio::sync::mpsc;
7
8use crate::live::Notification;
9use crate::{Engine, Result};
10
11/// Wrapper around an Engine that logs all RPC operations
12pub struct DebugEngine {
13    inner: Box<dyn Engine>,
14}
15
16impl DebugEngine {
17    /// Create a new DebugEngine wrapping an existing engine
18    pub fn wrap(engine: Box<dyn Engine>) -> Box<dyn Engine> {
19        Box::new(Self { inner: engine })
20    }
21
22    /// Log an RPC method call
23    fn log_request(&self, method: &str, params: &Value) {
24        let params_str = serde_json::to_string(params).unwrap_or_default();
25        println!("🔍 Surreal RPC: {} {}", method, params_str);
26    }
27
28    /// Log an RPC response
29    fn log_response(&self, response: &Value) {
30        // Check if response contains error
31        let icon = if let Value::Array(results) = response {
32            if results
33                .iter()
34                .any(|r| r.get("status").and_then(|s| s.as_str()) == Some("ERR"))
35            {
36                "❌"
37            } else {
38                "✅"
39            }
40        } else if response.get("error").is_some() {
41            "❌"
42        } else {
43            "✅"
44        };
45
46        let response_str = serde_json::to_string(response).unwrap_or_default();
47        println!("{} {}", icon, response_str);
48    }
49}
50
51#[async_trait]
52impl Engine for DebugEngine {
53    async fn send_message(&mut self, method: &str, params: Value) -> Result<Value> {
54        self.log_request(method, &params);
55        let response = self.inner.send_message(method, params).await?;
56        self.log_response(&response);
57        Ok(response)
58    }
59
60    async fn send_message_cbor(&mut self, method: &str, params: CborValue) -> Result<CborValue> {
61        println!("🔍 Surreal CBOR RPC: {} {:?}", method, params);
62        let response = self.inner.send_message_cbor(method, params).await?;
63        println!("✅ CBOR Response: {:?}", response);
64        Ok(response)
65    }
66
67    async fn register_live(
68        &mut self,
69        query_id: &str,
70    ) -> Result<mpsc::UnboundedReceiver<Notification>> {
71        self.inner.register_live(query_id).await
72    }
73
74    async fn unregister_live(&mut self, query_id: &str) {
75        self.inner.unregister_live(query_id).await
76    }
77}
78
79#[cfg(test)]
80mod tests {
81    use super::*;
82
83    struct MockEngine;
84
85    #[async_trait]
86    impl Engine for MockEngine {
87        async fn send_message(&mut self, _method: &str, _params: Value) -> Result<Value> {
88            Ok(serde_json::json!({
89                "status": "OK",
90                "result": []
91            }))
92        }
93
94        async fn send_message_cbor(
95            &mut self,
96            _method: &str,
97            _params: CborValue,
98        ) -> Result<CborValue> {
99            Ok(CborValue::Map(vec![
100                (
101                    CborValue::Text("status".to_string()),
102                    CborValue::Text("OK".to_string()),
103                ),
104                (
105                    CborValue::Text("result".to_string()),
106                    CborValue::Array(vec![]),
107                ),
108            ]))
109        }
110    }
111
112    #[tokio::test]
113    async fn test_debug_engine_wraps_correctly() {
114        let mock = Box::new(MockEngine);
115        let mut debug_engine = DebugEngine::wrap(mock);
116
117        let result = debug_engine
118            .send_message("test", serde_json::json!(["param1"]))
119            .await;
120
121        assert!(result.is_ok());
122    }
123}