diff --git a/litellm-proxy-extras/litellm_proxy_extras/migrations/20260429000000_add_scheduled_tasks/migration.sql b/litellm-proxy-extras/litellm_proxy_extras/migrations/20260429000000_add_scheduled_tasks/migration.sql index 489814e0d7a..f8fe979aeb5 100644 --- a/litellm-proxy-extras/litellm_proxy_extras/migrations/20260429000000_add_scheduled_tasks/migration.sql +++ b/litellm-proxy-extras/litellm_proxy_extras/migrations/20260429000000_add_scheduled_tasks/migration.sql @@ -19,6 +19,8 @@ CREATE TABLE "LiteLLM_ScheduledTaskTable" ( "fire_once" BOOLEAN NOT NULL DEFAULT true, "status" TEXT NOT NULL DEFAULT 'pending', "last_fired_at" TIMESTAMP(3), + "consecutive_errors" INTEGER NOT NULL DEFAULT 0, + "last_error" TEXT, "created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, "updated_at" TIMESTAMP(3) NOT NULL, @@ -54,7 +56,7 @@ ALTER TABLE "LiteLLM_ScheduledTaskTable" CHECK (schedule_kind IN ('interval','cron','once')); ALTER TABLE "LiteLLM_ScheduledTaskTable" ADD CONSTRAINT "LiteLLM_ScheduledTaskTable_status_check" - CHECK (status IN ('pending','fired','expired','cancelled')); + CHECK (status IN ('pending','fired','expired','cancelled','failed')); ALTER TABLE "LiteLLM_ScheduledTaskTable" ADD CONSTRAINT "LiteLLM_ScheduledTaskTable_action_prompt_check" CHECK (action <> 'check' OR check_prompt IS NOT NULL); diff --git a/litellm-proxy-extras/litellm_proxy_extras/schema.prisma b/litellm-proxy-extras/litellm_proxy_extras/schema.prisma index 8137c616b4b..62c6c16c3e4 100644 --- a/litellm-proxy-extras/litellm_proxy_extras/schema.prisma +++ b/litellm-proxy-extras/litellm_proxy_extras/schema.prisma @@ -1321,6 +1321,8 @@ model LiteLLM_ScheduledTaskTable { status String @default("pending") last_fired_at DateTime? + consecutive_errors Int @default(0) + last_error String? created_at DateTime @default(now()) updated_at DateTime @updatedAt diff --git a/litellm/proxy/scheduled_tasks/endpoints.py b/litellm/proxy/scheduled_tasks/endpoints.py index 5c2e17ae14f..628bd8c7f0a 100644 --- a/litellm/proxy/scheduled_tasks/endpoints.py +++ b/litellm/proxy/scheduled_tasks/endpoints.py @@ -25,6 +25,7 @@ from litellm.proxy.scheduled_tasks.types import ( DueTaskResponse, DueTasksResponse, ListScheduledTasksResponse, + ReportTaskResultRequest, ScheduledTaskResponse, UpdateScheduledTaskRequest, ) @@ -66,6 +67,8 @@ def _row_to_response(row) -> ScheduledTaskResponse: fire_once=row.fire_once, status=row.status, last_fired_at=row.last_fired_at, + consecutive_errors=getattr(row, "consecutive_errors", 0) or 0, + last_error=getattr(row, "last_error", None), created_at=row.created_at, updated_at=row.updated_at, ) @@ -318,3 +321,36 @@ async def cancel_scheduled_task( if row is None: raise HTTPException(status_code=404, detail="task not found") return _row_to_response(row) + + +@router.post( + "/v1/tasks/{task_id}/report", + tags=["scheduled tasks"], + response_model=ScheduledTaskResponse, +) +async def report_task_result( + task_id: str, + data: ReportTaskResultRequest, + user_api_key_dict: UserAPIKeyAuth = Depends(user_api_key_auth), +): + """ + Agent reports outcome of one dispatch attempt. + + success → resets the consecutive-error counter. + error → increments it. On the Nth consecutive error + (store.MAX_CONSECUTIVE_ERRORS), status flips to 'failed' and + /due stops returning the task. Caller's next list() shows it + as failed with last_error set. + """ + owner_token = _require_token(user_api_key_dict) + prisma_client = _get_prisma_client() + row = await store.report_task_result( + prisma_client, + task_id=task_id, + owner_token=owner_token, + result=data.result, + reason=data.reason, + ) + if row is None: + raise HTTPException(status_code=404, detail="task not found") + return _row_to_response(row) diff --git a/litellm/proxy/scheduled_tasks/store.py b/litellm/proxy/scheduled_tasks/store.py index 45314195ce8..4c770441d0f 100644 --- a/litellm/proxy/scheduled_tasks/store.py +++ b/litellm/proxy/scheduled_tasks/store.py @@ -49,7 +49,8 @@ UPDATABLE_FIELDS = frozenset( JSON_FIELDS = frozenset({"action_args", "metadata"}) MAX_ACTIVE_TASKS_PER_KEY = 10 -TERMINAL_STATUSES = ("fired", "expired", "cancelled") +MAX_CONSECUTIVE_ERRORS = 3 +TERMINAL_STATUSES = ("fired", "expired", "cancelled", "failed") async def count_active_for_owner(prisma_client: Any, owner_token: str) -> int: @@ -101,12 +102,32 @@ async def create_task( return await prisma_client.db.litellm_scheduledtasktable.create(data=data) +async def sweep_expired_for_owner(prisma_client: Any, owner_token: str) -> int: + """ + Lazy expiry: flip pending rows past expires_at to 'expired' before any + read returns them. Without this, a task that never fires after its + expires_at sits in 'pending' until the next claim that happens to scan + it — which may be never if no other rows are ever due. + + Returns count of rows flipped. + """ + return await prisma_client.db.litellm_scheduledtasktable.update_many( + where={ + "owner_token": owner_token, + "status": "pending", + "expires_at": {"lte": datetime.now(timezone.utc)}, + }, + data={"status": "expired"}, + ) + + async def list_tasks_for_owner( prisma_client: Any, *, owner_token: str, include_terminal: bool, ) -> List[Any]: + await sweep_expired_for_owner(prisma_client, owner_token) where: Dict[str, Any] = {"owner_token": owner_token} if not include_terminal: where["status"] = "pending" @@ -122,11 +143,55 @@ async def get_task_for_owner( task_id: str, owner_token: str, ) -> Optional[Any]: + await sweep_expired_for_owner(prisma_client, owner_token) return await prisma_client.db.litellm_scheduledtasktable.find_first( where={"task_id": task_id, "owner_token": owner_token}, ) +async def report_task_result( + prisma_client: Any, + *, + task_id: str, + owner_token: str, + result: str, + reason: Optional[str], +) -> Optional[Any]: + """ + Agent reports outcome of one dispatch attempt. + + success → reset consecutive_errors, clear last_error. + error → bump consecutive_errors. If >= MAX_CONSECUTIVE_ERRORS, flip + status to 'failed' so /due stops re-emitting it. + + Scoped by (task_id, owner_token). Returns updated row, or None if not + found / not owned. + """ + existing = await prisma_client.db.litellm_scheduledtasktable.find_first( + where={"task_id": task_id, "owner_token": owner_token}, + ) + if existing is None: + return None + + if result == "success": + return await prisma_client.db.litellm_scheduledtasktable.update( + where={"task_id": task_id}, + data={"consecutive_errors": 0, "last_error": None}, + ) + + new_count = (existing.consecutive_errors or 0) + 1 + update: Dict[str, Any] = { + "consecutive_errors": new_count, + "last_error": reason, + } + if new_count >= MAX_CONSECUTIVE_ERRORS and existing.status == "pending": + update["status"] = "failed" + return await prisma_client.db.litellm_scheduledtasktable.update( + where={"task_id": task_id}, + data=update, + ) + + async def update_task_for_owner( prisma_client: Any, *, diff --git a/litellm/proxy/scheduled_tasks/types.py b/litellm/proxy/scheduled_tasks/types.py index e415ccc17d6..0209a8749f1 100644 --- a/litellm/proxy/scheduled_tasks/types.py +++ b/litellm/proxy/scheduled_tasks/types.py @@ -4,7 +4,8 @@ from typing import Any, Dict, List, Literal, Optional from pydantic import BaseModel, Field ScheduleKind = Literal["interval", "cron", "once"] -TaskStatus = Literal["pending", "fired", "expired", "cancelled"] +TaskStatus = Literal["pending", "fired", "expired", "cancelled", "failed"] +ReportResult = Literal["success", "error"] class CreateScheduledTaskRequest(BaseModel): @@ -62,11 +63,20 @@ class ScheduledTaskResponse(BaseModel): status: str last_fired_at: Optional[datetime] + consecutive_errors: int = 0 + last_error: Optional[str] = None created_at: datetime updated_at: datetime +class ReportTaskResultRequest(BaseModel): + """Agent reports the outcome of one dispatch attempt.""" + + result: ReportResult + reason: Optional[str] = None + + class DueTaskResponse(BaseModel): """Trimmed shape returned by /due — only fields the agent needs to dispatch.""" diff --git a/litellm/proxy/schema.prisma b/litellm/proxy/schema.prisma index 8137c616b4b..62c6c16c3e4 100644 --- a/litellm/proxy/schema.prisma +++ b/litellm/proxy/schema.prisma @@ -1321,6 +1321,8 @@ model LiteLLM_ScheduledTaskTable { status String @default("pending") last_fired_at DateTime? + consecutive_errors Int @default(0) + last_error String? created_at DateTime @default(now()) updated_at DateTime @updatedAt diff --git a/schema.prisma b/schema.prisma index 8137c616b4b..62c6c16c3e4 100644 --- a/schema.prisma +++ b/schema.prisma @@ -1321,6 +1321,8 @@ model LiteLLM_ScheduledTaskTable { status String @default("pending") last_fired_at DateTime? + consecutive_errors Int @default(0) + last_error String? created_at DateTime @default(now()) updated_at DateTime @updatedAt diff --git a/tests/test_litellm/proxy/scheduled_tasks/test_endpoints.py b/tests/test_litellm/proxy/scheduled_tasks/test_endpoints.py index 31161b3ecb0..5efd0e683cf 100644 --- a/tests/test_litellm/proxy/scheduled_tasks/test_endpoints.py +++ b/tests/test_litellm/proxy/scheduled_tasks/test_endpoints.py @@ -53,6 +53,8 @@ def _make_row(**kwargs) -> MagicMock: "fire_once": True, "status": "pending", "last_fired_at": None, + "consecutive_errors": 0, + "last_error": None, "created_at": now, "updated_at": now, } @@ -120,6 +122,29 @@ class _FakeScheduledTaskTable: _ = order return self._filter(where) + async def update_many(self, where: Dict[str, Any], data: Dict[str, Any]) -> int: + # Minimal Prisma-style filter: support {"lte":