Skip to main content

ironflow_api/routes/
get_run_logs.rs

1//! `GET /api/v1/runs/:id/logs` -- Retrieve persisted log lines for a run.
2
3use 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/// Query parameters for listing run logs.
25#[derive(Debug, Deserialize)]
26#[cfg_attr(feature = "openapi", derive(utoipa::IntoParams, utoipa::ToSchema))]
27pub struct GetRunLogsQuery {
28    /// Filter by step ID.
29    pub step_id: Option<Uuid>,
30    /// Filter by output stream (`stdout`, `stderr`, `system`).
31    pub stream: Option<LogStream>,
32    /// Cursor for pagination (last entry ID from previous page).
33    pub cursor: Option<Uuid>,
34    /// Number of entries to return (default: 100, max: 1000).
35    pub limit: Option<u32>,
36}
37
38/// Cursor-based pagination metadata for log entries.
39#[derive(Debug, Serialize)]
40#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
41pub struct LogCursorMeta {
42    /// Cursor to pass for the next page. `None` when there are no more entries.
43    pub next_cursor: Option<Uuid>,
44    /// Whether more entries exist after this page.
45    pub has_more: bool,
46}
47
48/// Retrieve persisted log lines for a run with cursor-based pagination.
49///
50/// Returns log entries ordered by time (UUID v7 ascending). Use the
51/// `cursor` query parameter with the last entry's `id` to fetch the
52/// next page.
53#[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}