Skip to main content

zeal_sdk/
traces.rs

1//! Traces API for workflow execution tracing
2
3use crate::errors::{Result, ZealError};
4use crate::types::*;
5use reqwest::Client;
6use serde::{Deserialize, Serialize};
7
8#[derive(Debug, Clone, Serialize, Deserialize)]
9pub struct SubmitEventsResponse {
10    pub success: bool,
11    #[serde(rename = "eventsProcessed")]
12    pub events_processed: usize,
13}
14
15#[derive(Debug, Clone, Serialize, Deserialize)]
16pub struct CompleteSessionRequest {
17    pub status: SessionCompletionStatus,
18    pub summary: Option<SessionSummary>,
19    pub error: Option<SessionError>,
20}
21
22#[derive(Debug, Clone, Serialize, Deserialize)]
23#[serde(rename_all = "lowercase")]
24pub enum SessionCompletionStatus {
25    Success,
26    Error,
27    Cancelled,
28}
29
30#[derive(Debug, Clone, Serialize, Deserialize)]
31pub struct SessionSummary {
32    #[serde(rename = "totalNodes")]
33    pub total_nodes: u32,
34    #[serde(rename = "successfulNodes")]
35    pub successful_nodes: u32,
36    #[serde(rename = "failedNodes")]
37    pub failed_nodes: u32,
38    #[serde(rename = "totalDuration")]
39    pub total_duration: u64,
40    #[serde(rename = "totalDataProcessed")]
41    pub total_data_processed: u64,
42}
43
44#[derive(Debug, Clone, Serialize, Deserialize)]
45pub struct SessionError {
46    pub message: String,
47    #[serde(rename = "nodeId")]
48    pub node_id: Option<String>,
49    pub stack: Option<String>,
50}
51
52#[derive(Debug, Clone, Serialize, Deserialize)]
53pub struct CompleteSessionResponse {
54    pub success: bool,
55    #[serde(rename = "sessionId")]
56    pub session_id: String,
57    pub status: String,
58}
59
60#[derive(Debug, Clone, Serialize, Deserialize)]
61pub struct BatchTraceRequest {
62    #[serde(rename = "sessionId")]
63    pub session_id: String,
64    pub events: Vec<TraceEvent>,
65    #[serde(rename = "isComplete")]
66    pub is_complete: Option<bool>,
67}
68
69#[derive(Debug, Clone, Serialize, Deserialize)]
70pub struct BatchTraceResponse {
71    pub success: bool,
72}
73
74/// Traces API for managing execution traces
75pub struct TracesAPI {
76    base_url: String,
77    client: Client,
78    session_id: Option<String>,
79}
80
81impl TracesAPI {
82    /// Create a new Traces API instance
83    pub fn new(base_url: &str) -> Self {
84        Self {
85            base_url: base_url.to_string(),
86            client: Client::new(),
87            session_id: None,
88        }
89    }
90
91    /// Create a new Traces API instance with custom HTTP client
92    pub fn with_client(base_url: &str, client: Client) -> Self {
93        Self {
94            base_url: base_url.to_string(),
95            client,
96            session_id: None,
97        }
98    }
99
100    /// Create a new trace session
101    pub async fn create_session(
102        &mut self,
103        request: CreateTraceSessionRequest,
104    ) -> Result<CreateTraceSessionResponse> {
105        let url = format!(
106            "{}/api/zip/traces/sessions",
107            self.base_url.trim_end_matches('/')
108        );
109
110        let response = self
111            .client
112            .post(&url)
113            .header("Content-Type", "application/json")
114            .json(&request)
115            .send()
116            .await?;
117
118        let status = response.status();
119        if !status.is_success() {
120            let error_text = response
121                .text()
122                .await
123                .unwrap_or_else(|_| "Unknown error".to_string());
124            return Err(ZealError::api_error(
125                status.as_u16(),
126                format!("Failed to create trace session: {}", status),
127                Some(error_text),
128            ));
129        }
130
131        let session_response = response.json::<CreateTraceSessionResponse>().await?;
132        self.session_id = Some(session_response.session_id.clone());
133        Ok(session_response)
134    }
135
136    /// Submit trace events
137    pub async fn submit_events(
138        &self,
139        session_id: &str,
140        events: Vec<TraceEvent>,
141    ) -> Result<SubmitEventsResponse> {
142        let url = format!(
143            "{}/api/zip/traces/{}/events",
144            self.base_url.trim_end_matches('/'),
145            session_id
146        );
147
148        let request_body = serde_json::json!({
149            "events": events
150        });
151
152        let response = self
153            .client
154            .post(&url)
155            .header("Content-Type", "application/json")
156            .json(&request_body)
157            .send()
158            .await?;
159
160        let status = response.status();
161        if !status.is_success() {
162            let error_text = response
163                .text()
164                .await
165                .unwrap_or_else(|_| "Unknown error".to_string());
166            return Err(ZealError::api_error(
167                status.as_u16(),
168                format!("Failed to submit trace events: {}", status),
169                Some(error_text),
170            ));
171        }
172
173        let submit_response = response.json::<SubmitEventsResponse>().await?;
174        Ok(submit_response)
175    }
176
177    /// Submit a single trace event
178    pub async fn submit_event(
179        &self,
180        session_id: &str,
181        event: TraceEvent,
182    ) -> Result<SubmitEventsResponse> {
183        self.submit_events(session_id, vec![event]).await
184    }
185
186    /// Complete a trace session
187    pub async fn complete_session(
188        &mut self,
189        session_id: &str,
190        request: CompleteSessionRequest,
191    ) -> Result<CompleteSessionResponse> {
192        let url = format!(
193            "{}/api/zip/traces/{}/complete",
194            self.base_url.trim_end_matches('/'),
195            session_id
196        );
197
198        let response = self
199            .client
200            .post(&url)
201            .header("Content-Type", "application/json")
202            .json(&request)
203            .send()
204            .await?;
205
206        let status = response.status();
207        if !status.is_success() {
208            let error_text = response
209                .text()
210                .await
211                .unwrap_or_else(|_| "Unknown error".to_string());
212            return Err(ZealError::api_error(
213                status.as_u16(),
214                format!("Failed to complete trace session: {}", status),
215                Some(error_text),
216            ));
217        }
218
219        let complete_response = response.json::<CompleteSessionResponse>().await?;
220
221        if self.session_id.as_deref() == Some(session_id) {
222            self.session_id = None;
223        }
224
225        Ok(complete_response)
226    }
227
228    /// Helper method to trace node execution
229    pub async fn trace_node_execution(
230        &self,
231        session_id: &str,
232        node_id: &str,
233        event_type: TraceEventType,
234        data: serde_json::Value,
235        duration: Option<std::time::Duration>,
236    ) -> Result<()> {
237        let data_str = serde_json::to_string(&data)?;
238        let trace_data = TraceData {
239            size: data_str.len(),
240            data_type: "application/json".to_string(),
241            preview: Some(data.clone()),
242            full_data: Some(data),
243        };
244
245        let event = TraceEvent {
246            timestamp: chrono::Utc::now().timestamp_millis(),
247            node_id: node_id.to_string(),
248            port_id: None,
249            event_type,
250            data: trace_data,
251            duration,
252            metadata: None,
253            error: None,
254        };
255
256        self.submit_event(session_id, event).await?;
257        Ok(())
258    }
259
260    /// Batch trace submission
261    pub async fn submit_batch(&self, request: BatchTraceRequest) -> Result<BatchTraceResponse> {
262        let url = format!(
263            "{}/api/zip/traces/batch",
264            self.base_url.trim_end_matches('/')
265        );
266
267        let response = self
268            .client
269            .post(&url)
270            .header("Content-Type", "application/json")
271            .json(&request)
272            .send()
273            .await?;
274
275        let status = response.status();
276        if !status.is_success() {
277            let error_text = response
278                .text()
279                .await
280                .unwrap_or_else(|_| "Unknown error".to_string());
281            return Err(ZealError::api_error(
282                status.as_u16(),
283                format!("Failed to submit batch trace: {}", status),
284                Some(error_text),
285            ));
286        }
287
288        let batch_response = response.json::<BatchTraceResponse>().await?;
289        Ok(batch_response)
290    }
291
292    /// Get the current session ID
293    pub fn current_session_id(&self) -> Option<&str> {
294        self.session_id.as_deref()
295    }
296}
297
298/// Re-export trace types from types.rs for convenience
299pub use crate::types::{TraceEvent, TraceEventType, TraceStatus};