1use std::collections::HashMap;
4
5use axum::Json;
6use axum::extract::{Path, Query, State};
7use axum::response::IntoResponse;
8use serde::{Deserialize, Serialize};
9use serde_json::Value;
10use uuid::Uuid;
11
12use ironflow_auth::extractor::Authenticated;
13#[cfg(feature = "openapi")]
14use ironflow_store::entities::LogEntry;
15use ironflow_store::entities::{LogFilter, LogStream};
16use ironflow_types::{ApiMeta, ApiResponse};
17
18use crate::error::ApiError;
19use crate::state::AppState;
20
21const MAX_LIMIT: u32 = 1000;
22const DEFAULT_LIMIT: u32 = 100;
23
24#[derive(Debug, Deserialize)]
26#[cfg_attr(feature = "openapi", derive(utoipa::IntoParams, utoipa::ToSchema))]
27pub struct GetRunLogsQuery {
28 pub step_id: Option<Uuid>,
30 pub stream: Option<LogStream>,
32 pub cursor: Option<Uuid>,
34 pub limit: Option<u32>,
36}
37
38#[derive(Debug, Serialize)]
40#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
41pub struct LogCursorMeta {
42 pub next_cursor: Option<Uuid>,
44 pub has_more: bool,
46}
47
48#[cfg_attr(
54 feature = "openapi",
55 utoipa::path(
56 get,
57 path = "/api/v1/runs/{id}/logs",
58 tags = ["runs"],
59 params(
60 ("id" = Uuid, Path, description = "Run ID"),
61 GetRunLogsQuery,
62 ),
63 responses(
64 (status = 200, description = "Log entries with cursor-based pagination", body = Vec<LogEntry>),
65 (status = 401, description = "Unauthorized"),
66 (status = 404, description = "Run not found")
67 ),
68 security(("Bearer" = []))
69 )
70)]
71pub async fn get_run_logs(
72 _auth: Authenticated,
73 State(state): State<AppState>,
74 Path(run_id): Path<Uuid>,
75 Query(params): Query<GetRunLogsQuery>,
76) -> Result<impl IntoResponse, ApiError> {
77 state.get_run_or_404(run_id).await?;
78
79 let limit = match params.limit {
80 Some(0) | None => DEFAULT_LIMIT,
81 Some(l) => l.min(MAX_LIMIT),
82 };
83
84 let filter = LogFilter {
85 step_id: params.step_id,
86 stream: params.stream,
87 };
88
89 let entries = state
90 .store
91 .get_logs(run_id, filter, params.cursor, limit + 1)
92 .await?;
93
94 let has_more = entries.len() > limit as usize;
95 let entries: Vec<_> = entries.into_iter().take(limit as usize).collect();
96 let next_cursor = if has_more {
97 entries.last().map(|e| e.id)
98 } else {
99 None
100 };
101
102 let cursor_meta = LogCursorMeta {
103 next_cursor,
104 has_more,
105 };
106 let extra = match serde_json::to_value(cursor_meta) {
107 Ok(Value::Object(map)) => map.into_iter().collect(),
108 _ => HashMap::new(),
109 };
110
111 Ok(Json(ApiResponse {
112 data: entries,
113 meta: Some(ApiMeta {
114 page: None,
115 per_page: None,
116 total: None,
117 extra,
118 }),
119 }))
120}
121
122#[cfg(test)]
123mod tests {
124 use std::collections::HashMap;
125 use std::sync::Arc;
126
127 use axum::Router;
128 use axum::body::Body;
129 use axum::http::{Request, StatusCode};
130 use axum::routing::get;
131 use http_body_util::BodyExt;
132 use serde_json::{Value as JsonValue, from_slice, json};
133 use tokio::sync::broadcast;
134 use tower::ServiceExt;
135 use uuid::Uuid;
136
137 use ironflow_auth::jwt::{AccessToken, JwtConfig};
138 use ironflow_core::providers::claude::ClaudeCodeProvider;
139 use ironflow_engine::engine::Engine;
140 use ironflow_engine::notify::Event;
141 use ironflow_store::entities::{LogStream, NewLogEntries, NewRun, TriggerKind};
142 use ironflow_store::memory::InMemoryStore;
143
144 use super::*;
145
146 fn test_state() -> AppState {
147 let store = Arc::new(InMemoryStore::new());
148 let provider = Arc::new(ClaudeCodeProvider::new());
149 let engine = Arc::new(Engine::new(store.clone(), provider));
150 let jwt_config = Arc::new(JwtConfig {
151 secret: "test-secret".to_string(),
152 access_token_ttl_secs: 900,
153 refresh_token_ttl_secs: 604800,
154 cookie_domain: None,
155 cookie_secure: false,
156 });
157 let (event_sender, _) = broadcast::channel::<Event>(1);
158 AppState::new(
159 store,
160 engine,
161 jwt_config,
162 "test-worker-token".to_string(),
163 event_sender,
164 )
165 }
166
167 fn make_auth_header(state: &AppState) -> String {
168 let user_id = Uuid::now_v7();
169 let token = AccessToken::for_user(user_id, "testuser", false, &state.jwt_config).unwrap();
170 format!("Bearer {}", token.0)
171 }
172
173 async fn create_run(state: &AppState) -> Uuid {
174 state
175 .store
176 .create_run(NewRun {
177 created_by: None,
178 workflow_name: "test".to_string(),
179 trigger: TriggerKind::Manual,
180 payload: json!({}),
181 max_retries: 0,
182 handler_version: None,
183 labels: HashMap::new(),
184 scheduled_at: None,
185 idempotency_key: None,
186 max_cost_usd: None,
187 })
188 .await
189 .unwrap()
190 .into_run()
191 .id
192 }
193
194 async fn push_logs(state: &AppState, run_id: Uuid, step_id: Uuid, stream: LogStream, n: usize) {
195 state
196 .store
197 .append_logs(NewLogEntries {
198 ids: (0..n).map(|_| Uuid::now_v7()).collect(),
199 run_id,
200 step_id,
201 step_name: "build".to_string(),
202 stream,
203 lines: (0..n).map(|i| format!("line {i}")).collect(),
204 })
205 .await
206 .unwrap();
207 }
208
209 #[tokio::test]
210 async fn returns_persisted_logs() {
211 let state = test_state();
212 let auth_header = make_auth_header(&state);
213 let run_id = create_run(&state).await;
214 let step_id = Uuid::now_v7();
215
216 push_logs(&state, run_id, step_id, LogStream::Stdout, 3).await;
217
218 let app = Router::new()
219 .route("/runs/{id}/logs", get(get_run_logs))
220 .with_state(state);
221
222 let req = Request::builder()
223 .uri(format!("/runs/{run_id}/logs"))
224 .header("authorization", auth_header)
225 .body(Body::empty())
226 .unwrap();
227
228 let resp = app.oneshot(req).await.unwrap();
229 assert_eq!(resp.status(), StatusCode::OK);
230
231 let body = resp.into_body().collect().await.unwrap().to_bytes();
232 let json_val: JsonValue = from_slice(&body).unwrap();
233 assert_eq!(json_val["data"].as_array().unwrap().len(), 3);
234 assert_eq!(json_val["data"][0]["line"], "line 0");
235 assert_eq!(json_val["data"][0]["stream"], "stdout");
236 assert_eq!(json_val["meta"]["has_more"], false);
237 assert!(json_val["meta"]["next_cursor"].is_null());
238 }
239
240 #[tokio::test]
241 async fn filters_by_step_id() {
242 let state = test_state();
243 let auth_header = make_auth_header(&state);
244 let run_id = create_run(&state).await;
245 let step_a = Uuid::now_v7();
246 let step_b = Uuid::now_v7();
247
248 push_logs(&state, run_id, step_a, LogStream::Stdout, 2).await;
249 push_logs(&state, run_id, step_b, LogStream::Stdout, 3).await;
250
251 let app = Router::new()
252 .route("/runs/{id}/logs", get(get_run_logs))
253 .with_state(state);
254
255 let req = Request::builder()
256 .uri(format!("/runs/{run_id}/logs?step_id={step_a}"))
257 .header("authorization", auth_header)
258 .body(Body::empty())
259 .unwrap();
260
261 let resp = app.oneshot(req).await.unwrap();
262 assert_eq!(resp.status(), StatusCode::OK);
263
264 let body = resp.into_body().collect().await.unwrap().to_bytes();
265 let json_val: JsonValue = from_slice(&body).unwrap();
266 assert_eq!(json_val["data"].as_array().unwrap().len(), 2);
267 }
268
269 #[tokio::test]
270 async fn filters_by_stream() {
271 let state = test_state();
272 let auth_header = make_auth_header(&state);
273 let run_id = create_run(&state).await;
274 let step_id = Uuid::now_v7();
275
276 push_logs(&state, run_id, step_id, LogStream::Stdout, 2).await;
277 push_logs(&state, run_id, step_id, LogStream::Stderr, 1).await;
278
279 let app = Router::new()
280 .route("/runs/{id}/logs", get(get_run_logs))
281 .with_state(state);
282
283 let req = Request::builder()
284 .uri(format!("/runs/{run_id}/logs?stream=stderr"))
285 .header("authorization", auth_header)
286 .body(Body::empty())
287 .unwrap();
288
289 let resp = app.oneshot(req).await.unwrap();
290 let body = resp.into_body().collect().await.unwrap().to_bytes();
291 let json_val: JsonValue = from_slice(&body).unwrap();
292 assert_eq!(json_val["data"].as_array().unwrap().len(), 1);
293 assert_eq!(json_val["data"][0]["stream"], "stderr");
294 }
295
296 #[tokio::test]
297 async fn cursor_based_pagination() {
298 let state = test_state();
299 let auth_header = make_auth_header(&state);
300 let run_id = create_run(&state).await;
301 let step_id = Uuid::now_v7();
302
303 push_logs(&state, run_id, step_id, LogStream::Stdout, 5).await;
304
305 let app = Router::new()
306 .route("/runs/{id}/logs", get(get_run_logs))
307 .with_state(state);
308
309 let req = Request::builder()
310 .uri(format!("/runs/{run_id}/logs?limit=2"))
311 .header("authorization", &auth_header)
312 .body(Body::empty())
313 .unwrap();
314
315 let resp = app.clone().oneshot(req).await.unwrap();
316 let body = resp.into_body().collect().await.unwrap().to_bytes();
317 let page1: JsonValue = from_slice(&body).unwrap();
318 assert_eq!(page1["data"].as_array().unwrap().len(), 2);
319 assert_eq!(page1["meta"]["has_more"], true);
320
321 let cursor = page1["meta"]["next_cursor"].as_str().unwrap();
322
323 let req = Request::builder()
324 .uri(format!("/runs/{run_id}/logs?limit=2&cursor={cursor}"))
325 .header("authorization", &auth_header)
326 .body(Body::empty())
327 .unwrap();
328
329 let resp = app.clone().oneshot(req).await.unwrap();
330 let body = resp.into_body().collect().await.unwrap().to_bytes();
331 let page2: JsonValue = from_slice(&body).unwrap();
332 assert_eq!(page2["data"].as_array().unwrap().len(), 2);
333 assert_eq!(page2["data"][0]["line"], "line 2");
334 assert_eq!(page2["meta"]["has_more"], true);
335
336 let cursor = page2["meta"]["next_cursor"].as_str().unwrap();
337
338 let req = Request::builder()
339 .uri(format!("/runs/{run_id}/logs?limit=2&cursor={cursor}"))
340 .header("authorization", &auth_header)
341 .body(Body::empty())
342 .unwrap();
343
344 let resp = app.oneshot(req).await.unwrap();
345 let body = resp.into_body().collect().await.unwrap().to_bytes();
346 let page3: JsonValue = from_slice(&body).unwrap();
347 assert_eq!(page3["data"].as_array().unwrap().len(), 1);
348 assert_eq!(page3["meta"]["has_more"], false);
349 assert!(page3["meta"]["next_cursor"].is_null());
350 }
351
352 #[tokio::test]
353 async fn run_not_found_returns_404() {
354 let state = test_state();
355 let auth_header = make_auth_header(&state);
356 let app = Router::new()
357 .route("/runs/{id}/logs", get(get_run_logs))
358 .with_state(state);
359
360 let req = Request::builder()
361 .uri(format!("/runs/{}/logs", Uuid::now_v7()))
362 .header("authorization", auth_header)
363 .body(Body::empty())
364 .unwrap();
365
366 let resp = app.oneshot(req).await.unwrap();
367 assert_eq!(resp.status(), StatusCode::NOT_FOUND);
368 }
369
370 #[tokio::test]
371 async fn unauthenticated_returns_401() {
372 let state = test_state();
373 let run_id = create_run(&state).await;
374 let app = Router::new()
375 .route("/runs/{id}/logs", get(get_run_logs))
376 .with_state(state);
377
378 let req = Request::builder()
379 .uri(format!("/runs/{run_id}/logs"))
380 .body(Body::empty())
381 .unwrap();
382
383 let resp = app.oneshot(req).await.unwrap();
384 assert_eq!(resp.status(), StatusCode::UNAUTHORIZED);
385 }
386
387 #[tokio::test]
388 async fn limit_capped_at_1000() {
389 let state = test_state();
390 let auth_header = make_auth_header(&state);
391 let run_id = create_run(&state).await;
392 let step_id = Uuid::now_v7();
393
394 push_logs(&state, run_id, step_id, LogStream::Stdout, 2).await;
395
396 let app = Router::new()
397 .route("/runs/{id}/logs", get(get_run_logs))
398 .with_state(state);
399
400 let req = Request::builder()
401 .uri(format!("/runs/{run_id}/logs?limit=5000"))
402 .header("authorization", auth_header)
403 .body(Body::empty())
404 .unwrap();
405
406 let resp = app.oneshot(req).await.unwrap();
407 assert_eq!(resp.status(), StatusCode::OK);
408
409 let body = resp.into_body().collect().await.unwrap().to_bytes();
410 let json_val: JsonValue = from_slice(&body).unwrap();
411 assert_eq!(json_val["data"].as_array().unwrap().len(), 2);
412 }
413
414 #[tokio::test]
415 async fn empty_logs_returns_empty_array() {
416 let state = test_state();
417 let auth_header = make_auth_header(&state);
418 let run_id = create_run(&state).await;
419
420 let app = Router::new()
421 .route("/runs/{id}/logs", get(get_run_logs))
422 .with_state(state);
423
424 let req = Request::builder()
425 .uri(format!("/runs/{run_id}/logs"))
426 .header("authorization", auth_header)
427 .body(Body::empty())
428 .unwrap();
429
430 let resp = app.oneshot(req).await.unwrap();
431 assert_eq!(resp.status(), StatusCode::OK);
432
433 let body = resp.into_body().collect().await.unwrap().to_bytes();
434 let json_val: JsonValue = from_slice(&body).unwrap();
435 assert_eq!(json_val["data"].as_array().unwrap().len(), 0);
436 assert_eq!(json_val["meta"]["has_more"], false);
437 }
438}