/
alexefan136
/
flowstack
Обзор
Документация
Войти
/
alexefan136
/
flowstack
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
core/gateway/src/proxy/tasks.rs
225 строк
8 KB
Alexander Efanov
upd fix
31 июл 2026, 19:17
31 июл 2026, 19:17
d146d86
Код
Авторство
О чём код?
//! Proxy handlers для Tasks API (Gateway → Engine). //! //! Использует generic [`proxy_to_engine`] из [`super::forward`]: //! - Inject `X-User-ID` / `X-Workspace-ID` из AuthContext //! - Circuit breaker + error handling //! - Прозрачное проксирование (включая 204 No Content) //! //! Endpoints: //! - GET /api/v1/tasks — список задач //! - POST /api/v1/tasks — создать задачу //! - GET /api/v1/tasks/stats — статистика //! - GET /api/v1/tasks/overdue — просроченные задачи //! - GET /api/v1/tasks/{task_id} — получить задачу //! - PATCH /api/v1/tasks/{task_id} — обновить задачу //! - DELETE /api/v1/tasks/{task_id} — удалить задачу //! - POST /api/v1/tasks/{task_id}/start — в работу (in_progress) //! - POST /api/v1/tasks/{task_id}/complete — завершить (done) //! - POST /api/v1/tasks/{task_id}/cancel — отменить (cancelled) //! - GET /api/v1/tasks/{task_id}/runs — история выполнений //! - POST /api/v1/tasks/{task_id}/runs — добавить run //! //! ⚠️ Порядок роутов в main.rs: `/stats`, `/overdue` ДО `/{task_id}` //! (статические пути до динамических). use axum::{ body::Bytes, extract::{Path, RawQuery, State}, http::{HeaderMap, Method}, response::Response, }; use tracing::info; use super::forward::{ProxyState, proxy_to_engine}; use crate::auth::OptionalAuth; // ============================================================================ // COLLECTION ENDPOINTS // ============================================================================ /// GET /api/v1/tasks — список задач с фильтрами /// /// Query params (status, priority, assigned_agent_id, limit, offset) /// пробрасываются в engine через [`RawQuery`]. pub async fn proxy_tasks_list( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, RawQuery(query): RawQuery, ) -> Response { let path = match query { Some(q) if !q.is_empty() => format!("/api/v1/tasks?{}", q), _ => "/api/v1/tasks".to_string(), }; proxy_to_engine(&state, Method::GET, &path, auth, headers, None).await } /// POST /api/v1/tasks — создать новую задачу pub async fn proxy_tasks_create( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { info!("proxy_tasks_create"); proxy_to_engine( &state, Method::POST, "/api/v1/tasks", auth, headers, Some(body), ) .await } /// GET /api/v1/tasks/stats — агрегированная статистика по задачам /// /// ⚠️ Должен быть зарегистрирован ДО `/{task_id}` в main.rs. pub async fn proxy_tasks_stats( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { proxy_to_engine( &state, Method::GET, "/api/v1/tasks/stats", auth, headers, None, ) .await } /// GET /api/v1/tasks/overdue — список просроченных задач /// /// ⚠️ Должен быть зарегистрирован ДО `/{task_id}` в main.rs. pub async fn proxy_tasks_overdue( State(state): State<ProxyState>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { proxy_to_engine( &state, Method::GET, "/api/v1/tasks/overdue", auth, headers, None, ) .await } // ============================================================================ // ITEM ENDPOINTS // ============================================================================ /// GET /api/v1/tasks/{task_id} — получить задачу по ID pub async fn proxy_tasks_get( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { let path = format!("/api/v1/tasks/{}", task_id); proxy_to_engine(&state, Method::GET, &path, auth, headers, None).await } /// PATCH /api/v1/tasks/{task_id} — обновить задачу pub async fn proxy_tasks_update( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { info!(task_id = %task_id, "proxy_tasks_update"); let path = format!("/api/v1/tasks/{}", task_id); proxy_to_engine(&state, Method::PATCH, &path, auth, headers, Some(body)).await } /// DELETE /api/v1/tasks/{task_id} — удалить задачу /// /// Engine возвращает 204 No Content — `proxy_to_engine` пробрасывает прозрачно. pub async fn proxy_tasks_delete( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { info!(task_id = %task_id, "proxy_tasks_delete"); let path = format!("/api/v1/tasks/{}", task_id); proxy_to_engine(&state, Method::DELETE, &path, auth, headers, None).await } // ============================================================================ // STATUS TRANSITION ENDPOINTS // ============================================================================ /// POST /api/v1/tasks/{task_id}/start — перевести в работу (in_progress) pub async fn proxy_tasks_start( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { info!(task_id = %task_id, "proxy_tasks_start"); let path = format!("/api/v1/tasks/{}/start", task_id); proxy_to_engine(&state, Method::POST, &path, auth, headers, None).await } /// POST /api/v1/tasks/{task_id}/complete — завершить (done + completed_at) pub async fn proxy_tasks_complete( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { info!(task_id = %task_id, "proxy_tasks_complete"); let path = format!("/api/v1/tasks/{}/complete", task_id); proxy_to_engine(&state, Method::POST, &path, auth, headers, None).await } /// POST /api/v1/tasks/{task_id}/cancel — отменить (cancelled) pub async fn proxy_tasks_cancel( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { info!(task_id = %task_id, "proxy_tasks_cancel"); let path = format!("/api/v1/tasks/{}/cancel", task_id); proxy_to_engine(&state, Method::POST, &path, auth, headers, None).await } // ============================================================================ // TASK RUNS ENDPOINTS // ============================================================================ /// GET /api/v1/tasks/{task_id}/runs — история выполнений задачи pub async fn proxy_tasks_runs_list( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, ) -> Response { let path = format!("/api/v1/tasks/{}/runs", task_id); proxy_to_engine(&state, Method::GET, &path, auth, headers, None).await } /// POST /api/v1/tasks/{task_id}/runs — добавить запись о выполнении /// /// Body: `{ "status": "...", "output": "...", "tokens_used": N, "error_message": "..." }`. /// Engine использует атомарный JSONB concat (защита от race condition). pub async fn proxy_tasks_runs_create( State(state): State<ProxyState>, Path(task_id): Path<String>, OptionalAuth(auth): OptionalAuth, headers: HeaderMap, body: Bytes, ) -> Response { info!(task_id = %task_id, "proxy_tasks_runs_create"); let path = format!("/api/v1/tasks/{}/runs", task_id); proxy_to_engine(&state, Method::POST, &path, auth, headers, Some(body)).await }