chore: frontend and schedule module
This commit is contained in:
@@ -567,12 +567,11 @@ export async function getLatestScriptVersion(
|
|||||||
workspaceId: string,
|
workspaceId: string,
|
||||||
scriptId: string,
|
scriptId: string,
|
||||||
): Promise<LatestVersion | null> {
|
): Promise<LatestVersion | null> {
|
||||||
const resp = await apiRequest<{ data: LatestVersion | null }>(
|
return apiRequest<LatestVersion | null>(
|
||||||
`/api/v1/scripts/${scriptId}/latest-version`,
|
`/api/v1/scripts/${scriptId}/latest-version`,
|
||||||
{},
|
{},
|
||||||
workspaceId,
|
workspaceId,
|
||||||
);
|
);
|
||||||
return resp.data;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function listScriptVersions(
|
export async function listScriptVersions(
|
||||||
|
|||||||
@@ -13,10 +13,10 @@ post-back (see ``schedule.service.trigger_schedule``).
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
|
||||||
from datetime import timedelta
|
from datetime import timedelta
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
from sqlalchemy import select
|
from sqlalchemy import select
|
||||||
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
|
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
|
||||||
|
|
||||||
@@ -36,8 +36,6 @@ from schedule.context import (
|
|||||||
TERMINAL_RUN_STATES,
|
TERMINAL_RUN_STATES,
|
||||||
)
|
)
|
||||||
|
|
||||||
LOGGER = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
|
|
||||||
class DispatchOrchestrator:
|
class DispatchOrchestrator:
|
||||||
"""Polls Outbox + advances DAG schedule runs.
|
"""Polls Outbox + advances DAG schedule runs.
|
||||||
@@ -128,7 +126,7 @@ class DispatchOrchestrator:
|
|||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
raise
|
raise
|
||||||
except Exception:
|
except Exception:
|
||||||
LOGGER.exception("database event loop failed")
|
logger.exception("database event loop failed")
|
||||||
await asyncio.sleep(1)
|
await asyncio.sleep(1)
|
||||||
|
|
||||||
async def _execution_loop(self) -> None:
|
async def _execution_loop(self) -> None:
|
||||||
@@ -140,7 +138,7 @@ class DispatchOrchestrator:
|
|||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
raise
|
raise
|
||||||
except Exception:
|
except Exception:
|
||||||
LOGGER.exception("node execute loop failed")
|
logger.exception("node execute loop failed")
|
||||||
await asyncio.sleep(1)
|
await asyncio.sleep(1)
|
||||||
|
|
||||||
async def _claim_execution_events(
|
async def _claim_execution_events(
|
||||||
@@ -216,7 +214,7 @@ class DispatchOrchestrator:
|
|||||||
.with_for_update()
|
.with_for_update()
|
||||||
)
|
)
|
||||||
if item is None:
|
if item is None:
|
||||||
LOGGER.warning("execution event %s disappeared", event_id)
|
logger.warning("execution event {} disappeared", event_id)
|
||||||
return
|
return
|
||||||
if exc is None:
|
if exc is None:
|
||||||
item.event_status = "published"
|
item.event_status = "published"
|
||||||
|
|||||||
@@ -15,11 +15,12 @@ restarts.
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
|
||||||
from datetime import UTC
|
from datetime import UTC
|
||||||
from typing import Awaitable, Callable
|
from typing import Awaitable, Callable
|
||||||
from zoneinfo import ZoneInfo
|
from zoneinfo import ZoneInfo
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
from apscheduler.schedulers.asyncio import AsyncIOScheduler
|
from apscheduler.schedulers.asyncio import AsyncIOScheduler
|
||||||
from apscheduler.triggers.cron import CronTrigger
|
from apscheduler.triggers.cron import CronTrigger
|
||||||
from sqlalchemy import select
|
from sqlalchemy import select
|
||||||
@@ -31,8 +32,6 @@ from common.scheduler import build_sqlalchemy_jobstore
|
|||||||
|
|
||||||
from schedule.context import naive_utc
|
from schedule.context import naive_utc
|
||||||
|
|
||||||
LOGGER = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
# APScheduler's persistent SQLAlchemy job store pickles each job. A bound
|
# APScheduler's persistent SQLAlchemy job store pickles each job. A bound
|
||||||
# ``CronScheduler`` method captures this instance (including SQLAlchemy engine
|
# ``CronScheduler`` method captures this instance (including SQLAlchemy engine
|
||||||
# state) and therefore cannot be pickled. Keep the persisted callable at
|
# state) and therefore cannot be pickled. Keep the persisted callable at
|
||||||
@@ -111,7 +110,7 @@ class CronScheduler:
|
|||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
raise
|
raise
|
||||||
except Exception:
|
except Exception:
|
||||||
LOGGER.exception("cron job synchronization failed")
|
logger.exception("cron job synchronization failed")
|
||||||
await asyncio.sleep(5)
|
await asyncio.sleep(5)
|
||||||
|
|
||||||
async def _sync_once(self) -> None:
|
async def _sync_once(self) -> None:
|
||||||
|
|||||||
@@ -20,11 +20,11 @@ service shares the same MySQL via the Docker network).
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
|
||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
|
from loguru import logger
|
||||||
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
|
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
|
||||||
|
|
||||||
from common.config import settings
|
from common.config import settings
|
||||||
@@ -39,8 +39,6 @@ from schedule.orchestrator import DispatchOrchestrator
|
|||||||
from schedule.scheduler import CronScheduler
|
from schedule.scheduler import CronScheduler
|
||||||
from schedule.worker import NodeExecutor
|
from schedule.worker import NodeExecutor
|
||||||
|
|
||||||
LOGGER = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
|
|
||||||
class SchedulerService:
|
class SchedulerService:
|
||||||
"""Composes cron / orchestrator / worker into one bootable service.
|
"""Composes cron / orchestrator / worker into one bootable service.
|
||||||
@@ -153,8 +151,8 @@ class SchedulerService:
|
|||||||
except TriggerError as exc:
|
except TriggerError as exc:
|
||||||
# Most likely: idempotency_key collision from a previous
|
# Most likely: idempotency_key collision from a previous
|
||||||
# tick in the same minute — silently no-op.
|
# tick in the same minute — silently no-op.
|
||||||
LOGGER.info(
|
logger.info(
|
||||||
"cron trigger no-op for schedule %s: %s", schedule_id, exc,
|
"cron trigger no-op for schedule {}: {}", schedule_id, exc,
|
||||||
)
|
)
|
||||||
|
|
||||||
async def process_pending_events(
|
async def process_pending_events(
|
||||||
|
|||||||
@@ -16,11 +16,11 @@ from __future__ import annotations
|
|||||||
import asyncio
|
import asyncio
|
||||||
import hashlib
|
import hashlib
|
||||||
import json
|
import json
|
||||||
import logging
|
|
||||||
import traceback
|
import traceback
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
from sqlalchemy import select
|
from sqlalchemy import select
|
||||||
|
|
||||||
from common.db import session_scope
|
from common.db import session_scope
|
||||||
@@ -38,8 +38,6 @@ from common.eventing import add_outbox_event, event_time, utcnow
|
|||||||
from schedule.context import TERMINAL_NODE_STATES
|
from schedule.context import TERMINAL_NODE_STATES
|
||||||
from schedule.execution import ExecutionResult, execute_artifact
|
from schedule.execution import ExecutionResult, execute_artifact
|
||||||
|
|
||||||
LOGGER = logging.getLogger(__name__)
|
|
||||||
|
|
||||||
|
|
||||||
class NodeExecutor:
|
class NodeExecutor:
|
||||||
"""Owns the actual execution of one schedule node (notebook / python)."""
|
"""Owns the actual execution of one schedule node (notebook / python)."""
|
||||||
@@ -368,7 +366,7 @@ class NodeExecutor:
|
|||||||
result_id = result_object["storage_object_id"]
|
result_id = result_object["storage_object_id"]
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
upload_error = f"result upload failed: {exc}"[:2000]
|
upload_error = f"result upload failed: {exc}"[:2000]
|
||||||
LOGGER.exception("failed to upload node execution artifacts")
|
logger.exception("failed to upload node execution artifacts")
|
||||||
return log_id, result_id, upload_error
|
return log_id, result_id, upload_error
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
|
|||||||
Reference in New Issue
Block a user