1use 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
74pub struct TracesAPI {
76 base_url: String,
77 client: Client,
78 session_id: Option<String>,
79}
80
81impl TracesAPI {
82 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 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 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 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 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 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 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 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 pub fn current_session_id(&self) -> Option<&str> {
294 self.session_id.as_deref()
295 }
296}
297
298pub use crate::types::{TraceEvent, TraceEventType, TraceStatus};