Merge branch 'refactor/remove-redis' into develop
# Conflicts: # .env.example # CLAUDE.md # backend/Dockerfile # backend/pyproject.toml # backend/src/backend/main.py # backend/src/backend/schedule_runs.py # common/src/common/db/models.py # common/src/common/eventing.py # contracts/README.md # contracts/demo-core-v1.md # contracts/events/README.md # contracts/events/event-envelope-v1.json # contracts/locks/README.md # contracts/locks/file-edit-lock-v1.md # contracts/schedules/schedule-definition-api-v1.md # docker-compose.yml # frontend/README.md # migrations/README.md # migrations/versions/20260724_0001_v1_schema_baseline.py # runtime/Dockerfile # runtime/README.md # runtime/pyproject.toml # runtime/src/runtime/main.py # schedule/Dockerfile # schedule/README.md # schedule/pyproject.toml # schedule/src/schedule/main.py # schedule/src/schedule/service.py
This commit is contained in:
+3
-4
@@ -1,12 +1,11 @@
|
||||
FROM python:3.12-slim-bookworm
|
||||
|
||||
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
|
||||
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1 PYTHONPATH=/app
|
||||
WORKDIR /app
|
||||
COPY --from=ghcr.io/astral-sh/uv:latest /uv /bin/uv
|
||||
COPY pyproject.toml uv.lock ./
|
||||
COPY common ./common
|
||||
COPY schedule ./schedule
|
||||
RUN uv sync --frozen --no-dev --no-editable --package schedule
|
||||
RUN uv pip install --system ./common ./schedule
|
||||
|
||||
EXPOSE 8000
|
||||
CMD ["uv", "run", "--frozen", "--package", "schedule", "uvicorn", "schedule.main:app", "--host", "0.0.0.0", "--port", "8000"]
|
||||
CMD ["uvicorn", "schedule.main:app", "--host", "0.0.0.0", "--port", "8000"]
|
||||
|
||||
+8
-3
@@ -1,4 +1,9 @@
|
||||
# Schedule
|
||||
# Schedule Executor
|
||||
|
||||
独立调度执行服务。负责 Outbox/Redis Streams 轮询、DAG 节点派发、稳定版本
|
||||
执行、重试、状态推进以及日志和结果回写。
|
||||
独立调度执行服务,内置 APScheduler。
|
||||
|
||||
- Cron 任务持久化到 MySQL 的 `apscheduler_jobs` 表;
|
||||
- FastAPI Backend 创建运行记录和 Outbox 事件后,通过 HTTP 尝试立即推送;
|
||||
- HTTP 推送失败时,Executor 继续轮询 MySQL `outbox_events`,保证任务不会丢失;
|
||||
- Executor 负责 DAG 节点派发、稳定版本执行、重试、状态推进和结果回写;
|
||||
- 不依赖 Redis,MySQL 是调度状态与幂等状态的唯一权威。
|
||||
|
||||
@@ -7,14 +7,15 @@ dependencies = [
|
||||
"fastapi==0.116.1",
|
||||
"uvicorn[standard]==0.35.0",
|
||||
"httpx==0.28.1",
|
||||
"redis==5.2.1",
|
||||
"apscheduler==3.11.3",
|
||||
"pymysql==1.2.0",
|
||||
"nbclient==0.10.2",
|
||||
"nbformat==5.10.4",
|
||||
"ipykernel==6.29.5",
|
||||
]
|
||||
|
||||
[tool.uv.sources]
|
||||
common = { workspace = true }
|
||||
common = { path = "../common" }
|
||||
|
||||
[build-system]
|
||||
requires = ["hatchling"]
|
||||
|
||||
@@ -1,38 +1,52 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import secrets
|
||||
from contextlib import asynccontextmanager
|
||||
from pathlib import Path
|
||||
from typing import Any, AsyncIterator
|
||||
|
||||
from fastapi import Depends, Header, HTTPException, Request, status
|
||||
|
||||
from common.db import create_database_engine, create_session_factory
|
||||
from common.service_app import create_service_app
|
||||
from schedule.service import (
|
||||
SchedulerService,
|
||||
build_object_store,
|
||||
build_redis_client,
|
||||
build_storage_http_client,
|
||||
)
|
||||
from schedule.storage_client import SchedulerStorageClient
|
||||
|
||||
|
||||
def verify_internal_service(
|
||||
x_service_token: str = Header(alias="X-Service-Token"),
|
||||
) -> None:
|
||||
expected = os.environ.get("INTERNAL_SERVICE_TOKEN", "")
|
||||
if not expected or not secrets.compare_digest(expected, x_service_token):
|
||||
raise HTTPException(
|
||||
status.HTTP_401_UNAUTHORIZED,
|
||||
"invalid internal service identity",
|
||||
)
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: Any) -> AsyncIterator[None]:
|
||||
engine = create_database_engine(os.environ["DATABASE_URL"])
|
||||
database_url = os.environ["DATABASE_URL"]
|
||||
engine = create_database_engine(database_url)
|
||||
session_factory = create_session_factory(engine)
|
||||
redis = build_redis_client()
|
||||
storage_http_client = build_storage_http_client()
|
||||
backend_http_client = build_storage_http_client()
|
||||
service = SchedulerService(
|
||||
session_factory=session_factory,
|
||||
redis=redis,
|
||||
object_store=build_object_store(),
|
||||
storage_client=SchedulerStorageClient(
|
||||
storage_http_client,
|
||||
backend_http_client,
|
||||
os.environ["INTERNAL_SERVICE_TOKEN"],
|
||||
),
|
||||
backend_http_client=backend_http_client,
|
||||
workspace_root=Path(
|
||||
os.getenv("WORKSPACE_ROOT", "/workspace/workspaces")
|
||||
),
|
||||
database_url=database_url,
|
||||
)
|
||||
app.state.scheduler_service = service
|
||||
await service.start()
|
||||
@@ -40,12 +54,24 @@ async def lifespan(app: Any) -> AsyncIterator[None]:
|
||||
yield
|
||||
finally:
|
||||
await service.close()
|
||||
await storage_http_client.aclose()
|
||||
await redis.aclose()
|
||||
await backend_http_client.aclose()
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
app = create_service_app(
|
||||
os.getenv("SERVICE_NAME", "scheduler-worker"),
|
||||
os.getenv("SERVICE_NAME", "schedule-executor"),
|
||||
lifespan=lifespan,
|
||||
)
|
||||
|
||||
|
||||
@app.post(
|
||||
"/internal/v1/runs/{run_id}/dispatch",
|
||||
dependencies=[Depends(verify_internal_service)],
|
||||
)
|
||||
async def dispatch_run(run_id: str, request: Request) -> dict[str, Any]:
|
||||
processed = await request.app.state.scheduler_service.dispatch_run(run_id)
|
||||
return {
|
||||
"status": "accepted",
|
||||
"run_id": run_id,
|
||||
"processed_events": processed,
|
||||
}
|
||||
|
||||
+190
-171
@@ -5,18 +5,18 @@ import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import socket
|
||||
import traceback
|
||||
from collections.abc import Awaitable, Callable
|
||||
from contextlib import suppress
|
||||
from datetime import timedelta
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from zoneinfo import ZoneInfo
|
||||
|
||||
import boto3
|
||||
import httpx
|
||||
from redis.asyncio import Redis
|
||||
from redis.exceptions import ResponseError
|
||||
from apscheduler.jobstores.sqlalchemy import SQLAlchemyJobStore
|
||||
from apscheduler.schedulers.asyncio import AsyncIOScheduler
|
||||
from apscheduler.triggers.cron import CronTrigger
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
|
||||
|
||||
@@ -30,12 +30,7 @@ from common.db.models import (
|
||||
Versions,
|
||||
Workspaces,
|
||||
)
|
||||
from common.eventing import (
|
||||
STREAM_BY_EVENT_TYPE,
|
||||
add_outbox_event,
|
||||
event_time,
|
||||
utcnow,
|
||||
)
|
||||
from common.eventing import add_outbox_event, event_time, utcnow
|
||||
from common.ids import new_ulid
|
||||
from common.db.session import session_scope
|
||||
from schedule.execution import ExecutionResult, execute_artifact
|
||||
@@ -58,203 +53,236 @@ TERMINAL_RUN_STATES = {
|
||||
"timed_out",
|
||||
}
|
||||
|
||||
_ACTIVE_SERVICE: "SchedulerService | None" = None
|
||||
|
||||
|
||||
def _sync_database_url(value: str) -> str:
|
||||
return value.replace("mysql+asyncmy://", "mysql+pymysql://", 1)
|
||||
|
||||
|
||||
def _naive_utc(value: datetime | None) -> datetime | None:
|
||||
if value is None:
|
||||
return None
|
||||
if value.tzinfo is None:
|
||||
value = value.replace(tzinfo=UTC)
|
||||
return value.astimezone(UTC).replace(tzinfo=None)
|
||||
|
||||
|
||||
async def run_scheduled_job(schedule_id: str) -> None:
|
||||
service = _ACTIVE_SERVICE
|
||||
if service is None:
|
||||
LOGGER.warning("scheduler job skipped because service is not ready")
|
||||
return
|
||||
await service.trigger_schedule(schedule_id)
|
||||
|
||||
|
||||
class SchedulerService:
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
session_factory: async_sessionmaker[AsyncSession],
|
||||
redis: Redis,
|
||||
object_store: Any,
|
||||
storage_client: SchedulerStorageClient,
|
||||
backend_http_client: httpx.AsyncClient,
|
||||
workspace_root: Path,
|
||||
database_url: str,
|
||||
) -> None:
|
||||
self.session_factory = session_factory
|
||||
self.redis = redis
|
||||
self.object_store = object_store
|
||||
self.storage_client = storage_client
|
||||
self.backend_http_client = backend_http_client
|
||||
self.workspace_root = workspace_root
|
||||
self.consumer_name = (
|
||||
os.getenv("SCHEDULER_CONSUMER_NAME")
|
||||
or f"{socket.gethostname()}-{os.getpid()}"
|
||||
)
|
||||
self.tasks: list[asyncio.Task[Any]] = []
|
||||
self.dispatch_lock = asyncio.Lock()
|
||||
self.scheduler = AsyncIOScheduler(
|
||||
jobstores={
|
||||
"default": SQLAlchemyJobStore(
|
||||
url=_sync_database_url(database_url),
|
||||
tablename="apscheduler_jobs",
|
||||
)
|
||||
},
|
||||
timezone=UTC,
|
||||
)
|
||||
|
||||
async def start(self) -> None:
|
||||
await self._ensure_group(
|
||||
"stream:scheduler:commands",
|
||||
"schedule-orchestrator",
|
||||
)
|
||||
await self._ensure_group("stream:jobs:execute", "job-workers")
|
||||
await self._ensure_group("stream:jobs:results", "schedule-results")
|
||||
global _ACTIVE_SERVICE
|
||||
_ACTIVE_SERVICE = self
|
||||
self.scheduler.start()
|
||||
await self._sync_cron_jobs()
|
||||
self.tasks = [
|
||||
asyncio.create_task(
|
||||
self._outbox_loop(),
|
||||
name="scheduler-outbox-publisher",
|
||||
self._database_event_loop(),
|
||||
name="scheduler-database-events",
|
||||
),
|
||||
asyncio.create_task(
|
||||
self._consumer_loop(
|
||||
"stream:scheduler:commands",
|
||||
"schedule-orchestrator",
|
||||
self._handle_run_requested,
|
||||
),
|
||||
name="schedule-orchestrator",
|
||||
),
|
||||
asyncio.create_task(
|
||||
self._consumer_loop(
|
||||
"stream:jobs:execute",
|
||||
"job-workers",
|
||||
self._handle_node_execute,
|
||||
),
|
||||
name="job-worker",
|
||||
),
|
||||
asyncio.create_task(
|
||||
self._consumer_loop(
|
||||
"stream:jobs:results",
|
||||
"schedule-results",
|
||||
self._handle_node_finished,
|
||||
),
|
||||
name="schedule-results",
|
||||
self._schedule_sync_loop(),
|
||||
name="scheduler-cron-sync",
|
||||
),
|
||||
]
|
||||
|
||||
async def close(self) -> None:
|
||||
global _ACTIVE_SERVICE
|
||||
for task in self.tasks:
|
||||
task.cancel()
|
||||
for task in self.tasks:
|
||||
with suppress(asyncio.CancelledError):
|
||||
await task
|
||||
self.tasks.clear()
|
||||
if self.scheduler.running:
|
||||
self.scheduler.shutdown(wait=False)
|
||||
_ACTIVE_SERVICE = None
|
||||
|
||||
async def _ensure_group(self, stream: str, group: str) -> None:
|
||||
try:
|
||||
await self.redis.xgroup_create(
|
||||
stream,
|
||||
group,
|
||||
id="0-0",
|
||||
mkstream=True,
|
||||
)
|
||||
except ResponseError as exc:
|
||||
if "BUSYGROUP" not in str(exc):
|
||||
raise
|
||||
|
||||
async def _outbox_loop(self) -> None:
|
||||
async def _database_event_loop(self) -> None:
|
||||
while True:
|
||||
try:
|
||||
published = await self._publish_outbox_batch()
|
||||
if not published:
|
||||
await asyncio.sleep(0.35)
|
||||
processed = await self.process_pending_events(limit=20)
|
||||
if not processed:
|
||||
await asyncio.sleep(0.25)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception:
|
||||
LOGGER.exception("outbox publisher iteration failed")
|
||||
LOGGER.exception("database event loop failed")
|
||||
await asyncio.sleep(1)
|
||||
|
||||
async def _publish_outbox_batch(self) -> int:
|
||||
now = utcnow()
|
||||
async def process_pending_events(
|
||||
self,
|
||||
*,
|
||||
limit: int = 20,
|
||||
aggregate_id: str | None = None,
|
||||
) -> int:
|
||||
async with self.dispatch_lock:
|
||||
async with session_scope(self.session_factory) as session:
|
||||
statement = (
|
||||
select(OutboxEvents)
|
||||
.where(
|
||||
OutboxEvents.event_status == "pending",
|
||||
OutboxEvents.available_at <= utcnow(),
|
||||
)
|
||||
.order_by(OutboxEvents.created_at)
|
||||
.limit(limit)
|
||||
)
|
||||
if aggregate_id:
|
||||
statement = statement.where(
|
||||
OutboxEvents.aggregate_id == aggregate_id
|
||||
)
|
||||
events = list((await session.scalars(statement)).all())
|
||||
for item in events:
|
||||
try:
|
||||
await self._process_outbox_event(item)
|
||||
item.event_status = "published"
|
||||
item.published_at = utcnow()
|
||||
item.last_error = None
|
||||
except Exception as exc:
|
||||
item.retry_count += 1
|
||||
item.last_error = str(exc)[:2000]
|
||||
if item.retry_count >= 5:
|
||||
item.event_status = "failed"
|
||||
else:
|
||||
item.available_at = utcnow() + timedelta(
|
||||
seconds=min(30, 2 ** item.retry_count)
|
||||
)
|
||||
LOGGER.exception(
|
||||
"failed to process database event %s",
|
||||
item.event_id,
|
||||
)
|
||||
return len(events)
|
||||
|
||||
async def _process_outbox_event(self, item: OutboxEvents) -> None:
|
||||
handlers = {
|
||||
"schedule.run.requested": self._handle_run_requested,
|
||||
"job.node.execute": self._handle_node_execute,
|
||||
"job.node.finished": self._handle_node_finished,
|
||||
}
|
||||
handler = handlers.get(item.event_type)
|
||||
if handler is None:
|
||||
raise ValueError(f"unsupported event type: {item.event_type}")
|
||||
await handler(item.payload_json, f"mysql:{item.event_id}")
|
||||
|
||||
async def dispatch_run(self, run_id: str) -> int:
|
||||
return await self.process_pending_events(
|
||||
limit=50,
|
||||
aggregate_id=run_id,
|
||||
)
|
||||
|
||||
async def _schedule_sync_loop(self) -> None:
|
||||
while True:
|
||||
try:
|
||||
await self._sync_cron_jobs()
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception:
|
||||
LOGGER.exception("cron job synchronization failed")
|
||||
await asyncio.sleep(5)
|
||||
|
||||
async def _sync_cron_jobs(self) -> None:
|
||||
async with session_scope(self.session_factory) as session:
|
||||
events = list(
|
||||
schedules = list(
|
||||
(
|
||||
await session.scalars(
|
||||
select(OutboxEvents)
|
||||
.where(
|
||||
OutboxEvents.event_status == "pending",
|
||||
OutboxEvents.available_at <= now,
|
||||
select(Schedules).where(
|
||||
Schedules.deleted_at.is_(None),
|
||||
Schedules.enabled == 1,
|
||||
Schedules.trigger_type == "cron",
|
||||
Schedules.cron_expression.is_not(None),
|
||||
)
|
||||
.order_by(OutboxEvents.created_at)
|
||||
.limit(20)
|
||||
.with_for_update(skip_locked=True)
|
||||
)
|
||||
).all()
|
||||
)
|
||||
for item in events:
|
||||
stream = STREAM_BY_EVENT_TYPE.get(item.event_type)
|
||||
if stream is None:
|
||||
item.event_status = "failed"
|
||||
item.last_error = f"unsupported event type: {item.event_type}"
|
||||
continue
|
||||
try:
|
||||
await self.redis.xadd(
|
||||
stream,
|
||||
{
|
||||
"event": json.dumps(
|
||||
item.payload_json,
|
||||
ensure_ascii=False,
|
||||
separators=(",", ":"),
|
||||
)
|
||||
},
|
||||
)
|
||||
item.event_status = "published"
|
||||
item.published_at = utcnow()
|
||||
item.last_error = None
|
||||
except Exception as exc:
|
||||
item.retry_count += 1
|
||||
item.last_error = str(exc)[:2000]
|
||||
raise
|
||||
return len(events)
|
||||
|
||||
async def _consumer_loop(
|
||||
self,
|
||||
stream: str,
|
||||
group: str,
|
||||
handler: Callable[[dict[str, Any], str], Awaitable[None]],
|
||||
) -> None:
|
||||
while True:
|
||||
try:
|
||||
messages = await self.redis.xreadgroup(
|
||||
group,
|
||||
self.consumer_name,
|
||||
{stream: ">"},
|
||||
count=5,
|
||||
block=1000,
|
||||
active_job_ids: set[str] = set()
|
||||
for item in schedules:
|
||||
job_id = f"schedule:{item.schedule_id}"
|
||||
active_job_ids.add(job_id)
|
||||
expression = (item.cron_expression or "").strip()
|
||||
trigger = CronTrigger.from_crontab(
|
||||
expression,
|
||||
timezone=ZoneInfo(item.timezone),
|
||||
)
|
||||
entries: list[tuple[str, dict[str, str]]] = []
|
||||
for _, stream_messages in messages:
|
||||
entries.extend(stream_messages)
|
||||
if not entries:
|
||||
claimed = await self.redis.xautoclaim(
|
||||
stream,
|
||||
group,
|
||||
self.consumer_name,
|
||||
min_idle_time=10_000,
|
||||
start_id="0-0",
|
||||
count=5,
|
||||
)
|
||||
if len(claimed) >= 2:
|
||||
entries.extend(claimed[1])
|
||||
for message_id, fields in entries:
|
||||
try:
|
||||
raw = fields.get("event")
|
||||
if not raw:
|
||||
raise ValueError("stream message has no event field")
|
||||
event = json.loads(raw)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except (ValueError, TypeError, json.JSONDecodeError):
|
||||
LOGGER.exception(
|
||||
"discarding malformed message %s from %s",
|
||||
message_id,
|
||||
stream,
|
||||
)
|
||||
await self.redis.xack(stream, group, message_id)
|
||||
continue
|
||||
try:
|
||||
await handler(event, message_id)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception:
|
||||
LOGGER.exception(
|
||||
"consumer %s failed for message %s",
|
||||
group,
|
||||
message_id,
|
||||
)
|
||||
continue
|
||||
await self.redis.xack(stream, group, message_id)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception:
|
||||
LOGGER.exception("consumer loop %s failed", group)
|
||||
await asyncio.sleep(1)
|
||||
job = self.scheduler.add_job(
|
||||
run_scheduled_job,
|
||||
trigger=trigger,
|
||||
args=[item.schedule_id],
|
||||
id=job_id,
|
||||
replace_existing=True,
|
||||
coalesce=True,
|
||||
max_instances=max(1, item.max_concurrency),
|
||||
misfire_grace_time=60,
|
||||
)
|
||||
item.next_run_at = _naive_utc(job.next_run_time)
|
||||
for job in self.scheduler.get_jobs():
|
||||
if job.id.startswith("schedule:") and job.id not in active_job_ids:
|
||||
self.scheduler.remove_job(job.id)
|
||||
|
||||
async def trigger_schedule(self, schedule_id: str) -> None:
|
||||
async with self.session_factory() as session:
|
||||
item = await session.get(Schedules, schedule_id)
|
||||
if (
|
||||
item is None
|
||||
or item.deleted_at is not None
|
||||
or not item.enabled
|
||||
or item.trigger_type != "cron"
|
||||
):
|
||||
return
|
||||
user_id = item.created_by
|
||||
workspace_id = item.workspace_id
|
||||
now = datetime.now(UTC)
|
||||
idempotency_key = (
|
||||
f"cron:{schedule_id}:{now.strftime('%Y%m%d%H%M')}"
|
||||
)
|
||||
response = await self.backend_http_client.post(
|
||||
f"/api/v1/schedules/{schedule_id}/run",
|
||||
headers={
|
||||
"X-User-ID": user_id,
|
||||
"X-Workspace-ID": workspace_id,
|
||||
"X-Request-ID": new_ulid(),
|
||||
"Idempotency-Key": idempotency_key,
|
||||
},
|
||||
json={"reason": "cron"},
|
||||
)
|
||||
if response.is_error:
|
||||
raise RuntimeError(
|
||||
f"backend rejected cron run: {response.status_code} "
|
||||
f"{response.text[:500]}"
|
||||
)
|
||||
|
||||
async def _start_inbox(
|
||||
self,
|
||||
@@ -874,15 +902,6 @@ class SchedulerService:
|
||||
self._finish_inbox(inbox)
|
||||
|
||||
|
||||
def build_redis_client() -> Redis:
|
||||
return Redis(
|
||||
host=os.getenv("REDIS_HOST", "redis"),
|
||||
port=int(os.getenv("REDIS_PORT", "6379")),
|
||||
password=os.getenv("REDIS_PASSWORD") or None,
|
||||
decode_responses=True,
|
||||
)
|
||||
|
||||
|
||||
def build_object_store() -> Any:
|
||||
return boto3.client(
|
||||
"s3",
|
||||
@@ -898,6 +917,6 @@ def build_object_store() -> Any:
|
||||
|
||||
def build_storage_http_client() -> httpx.AsyncClient:
|
||||
return httpx.AsyncClient(
|
||||
base_url=os.getenv("STORAGE_API_URL", "http://storage_api:8000"),
|
||||
base_url=os.getenv("BACKEND_API_URL", "http://backend:8000"),
|
||||
timeout=httpx.Timeout(60.0),
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user