Skip to content

forge.jobs

Background jobs and scheduling module with queue-backed execution, retry, concurrency control, and cron-like scheduled tasks.

forge.jobs

Background jobs and scheduling module.

Provides queue-backed async job execution with automatic retry, concurrency control, dead-letter queue, progress tracking, and cron-like scheduled tasks.

Usage::

from forge.jobs import job, schedule, JobsModule

@job(queue="emails", retry=3)
async def send_welcome_email(user_id: int):
    ...

await send_welcome_email.enqueue(user_id=123)

@schedule(cron="0 9 * * *")
async def daily_report():
    ...

Classes

CronExpression

Standard five-field cron expression parser.

Fields: minute hour day-of-month month day-of-week

Supported syntax: - * — any value - 5 — exact value - 1,15 — list of values - 1-5 — range (inclusive) - */15 — step (every N) - 0 9 * * * — daily at 09:00 - */5 * * * 1-5 — every 5 minutes on weekdays

Source code in src/forge/jobs/scheduler.py
class CronExpression:
    """
    Standard five-field cron expression parser.

    Fields: minute hour day-of-month month day-of-week

    Supported syntax:
    - ``*`` — any value
    - ``5`` — exact value
    - ``1,15`` — list of values
    - ``1-5`` — range (inclusive)
    - ``*/15`` — step (every N)
    - ``0 9 * * *`` — daily at 09:00
    - ``*/5 * * * 1-5`` — every 5 minutes on weekdays
    """

    MONTH_NAMES: ClassVar[dict[str, int]] = {
        "JAN": 1,
        "FEB": 2,
        "MAR": 3,
        "APR": 4,
        "MAY": 5,
        "JUN": 6,
        "JUL": 7,
        "AUG": 8,
        "SEP": 9,
        "OCT": 10,
        "NOV": 11,
        "DEC": 12,
    }
    DOW_NAMES: ClassVar[dict[str, int]] = {
        "SUN": 0,
        "MON": 1,
        "TUE": 2,
        "WED": 3,
        "THU": 4,
        "FRI": 5,
        "SAT": 6,
    }

    _CRON_FIELD_COUNT = 5
    _MONTH_MAX = 12
    _DOW_MAX = 6

    def __init__(self, expression: str) -> None:
        self._raw = expression
        parts = expression.strip().split()
        if len(parts) != self._CRON_FIELD_COUNT:
            raise ValueError(
                f"Cron expression must have {self._CRON_FIELD_COUNT} fields, got {len(parts)}: '{expression}'"
            )
        self._minute = self._parse_field(parts[0], 0, 59)
        self._hour = self._parse_field(parts[1], 0, 23)
        self._dom = self._parse_field(parts[2], 1, 31)
        self._month = self._parse_field(parts[3], 1, 12)
        self._dow = self._parse_field(parts[4], 0, 6)

    def matches(self, dt: datetime) -> bool:
        dow = (dt.weekday() + 1) % 7
        return (
            dt.minute in self._minute
            and dt.hour in self._hour
            and dt.day in self._dom
            and dt.month in self._month
            and dow in self._dow
        )

    def next_after(self, dt: datetime) -> datetime:
        candidate = dt.replace(second=0, microsecond=0) + timedelta(minutes=1)
        for _ in range(525600):
            if self.matches(candidate):
                return candidate
            candidate += timedelta(minutes=1)
        raise RuntimeError(f"No matching time found for cron '{self._raw}' within 1 year")

    def _parse_field(self, field: str, min_val: int, max_val: int) -> set[int]:
        result: set[int] = set()
        for part in field.split(","):
            result |= self._parse_part(part.strip(), min_val, max_val)
        return result

    def _parse_part(self, part: str, min_val: int, max_val: int) -> set[int]:
        base_min, base_max = min_val, max_val

        part = part.upper()
        if min_val == 1 and max_val == self._MONTH_MAX:
            for name, val in self.MONTH_NAMES.items():
                part = part.replace(name, str(val))
        if min_val == 0 and max_val == self._DOW_MAX:
            for name, val in self.DOW_NAMES.items():
                part = part.replace(name, str(val))
            part = part.replace("7", "0")

        step = 1
        if "/" in part:
            range_part, step_str = part.split("/", 1)
            step = int(step_str)
            part = range_part

        if part == "*":
            return set(range(base_min, base_max + 1, step))

        if "-" in part:
            start_str, end_str = part.split("-", 1)
            start = int(start_str)
            end = int(end_str)
            return set(range(start, end + 1, step))

        value = int(part)
        if step > 1:
            return set(range(value, max_val + 1, step))
        return {value}

    def __repr__(self) -> str:
        return f"CronExpression({self._raw!r})"

Job

Represents a single enqueued background job.

Tracks identity, payload, retry state, progress, and result.

Source code in src/forge/jobs/queue.py
class Job:
    """
    Represents a single enqueued background job.

    Tracks identity, payload, retry state, progress, and result.
    """

    def __init__(
        self,
        job_id: str,
        queue: str,
        func_name: str,
        args: tuple[Any, ...],
        kwargs: dict[str, Any],
        max_retries: int,
    ) -> None:
        self.job_id = job_id
        self.queue = queue
        self.func_name = func_name
        self.args = args
        self.kwargs = kwargs
        self.max_retries = max_retries
        self.retry_count: int = 0
        self.status: str = JobStatus.PENDING
        self.result: JobResult | None = None
        self.progress: float = 0.0
        self.created_at: float = time.monotonic()
        self.started_at: float | None = None
        self.finished_at: float | None = None

    def to_dict(self) -> dict[str, Any]:
        return {
            "job_id": self.job_id,
            "queue": self.queue,
            "func_name": self.func_name,
            "retry_count": self.retry_count,
            "max_retries": self.max_retries,
            "status": self.status,
            "progress": self.progress,
            "created_at": self.created_at,
            "started_at": self.started_at,
            "finished_at": self.finished_at,
            "error": str(self.result.error) if self.result and self.result.error else None,
        }

JobDefinition

Wrapper returned by the @job decorator.

Holds the original async function and metadata. Calling .enqueue(...) pushes a job onto the queue.

Source code in src/forge/jobs/module.py
class JobDefinition:
    """
    Wrapper returned by the ``@job`` decorator.

    Holds the original async function and metadata. Calling
    ``.enqueue(...)`` pushes a job onto the queue.
    """

    def __init__(
        self,
        func: Callable[..., Any],
        queue: str,
        max_retries: int,
    ) -> None:
        self._func = func
        self.queue = queue
        self.max_retries = max_retries
        functools.update_wrapper(self, func)

    async def __call__(self, *args: Any, **kwargs: Any) -> Any:
        return await self._func(*args, **kwargs)

    async def enqueue(self, *args: Any, **kwargs: Any) -> Any:
        from forge.jobs._state import get_job_queue

        jq = get_job_queue()
        if jq is None:
            raise RuntimeError(
                "Jobs module is not initialized. "
                "Ensure the JobsModule is registered with the runtime and "
                "await runtime.init() has been called."
            )
        job = await jq.enqueue(
            queue=self.queue,
            func_name=self._func.__qualname__,
            args=args,
            kwargs=kwargs,
            max_retries=self.max_retries,
        )
        return job

JobQueue

High-level job queue manager.

Coordinates enqueueing, processing, retries, concurrency, dead-letter handling, and progress tracking.

Source code in src/forge/jobs/queue.py
class JobQueue:
    """
    High-level job queue manager.

    Coordinates enqueueing, processing, retries, concurrency,
    dead-letter handling, and progress tracking.
    """

    def __init__(
        self,
        backend: QueueBackend,
        default_retry: int = 3,
        concurrency: int = 10,
        retry_backoff_base: float = 1.0,
    ) -> None:
        self._backend = backend
        self._default_retry = default_retry
        self._concurrency = concurrency
        self._retry_backoff_base = retry_backoff_base
        self._registry: dict[str, Any] = {}
        self._semaphore: asyncio.Semaphore | None = None
        self._running: dict[str, asyncio.Task[None]] = {}
        self._started = False

    @property
    def backend(self) -> QueueBackend:
        return self._backend

    def register(self, func_name: str, func: Any) -> None:
        self._registry[func_name] = func

    async def start(self) -> None:
        self._semaphore = asyncio.Semaphore(self._concurrency)
        self._started = True

    async def stop(self) -> None:
        self._started = False
        for task in self._running.values():
            task.cancel()
        if self._running:
            await asyncio.gather(*self._running.values(), return_exceptions=True)
        self._running.clear()
        await self._backend.close()

    async def enqueue(
        self,
        queue: str,
        func_name: str,
        args: tuple[Any, ...] = (),
        kwargs: dict[str, Any] | None = None,
        max_retries: int | None = None,
    ) -> Job:
        job = Job(
            job_id=uuid.uuid4().hex,
            queue=queue,
            func_name=func_name,
            args=args,
            kwargs=kwargs or {},
            max_retries=max_retries if max_retries is not None else self._default_retry,
        )
        await self._backend.enqueue(job)
        _logger.info("Enqueued job %s on queue '%s'", job.job_id, queue)
        return job

    async def process(self, queue: str) -> None:
        if not self._started or self._semaphore is None:
            return
        while self._started:
            job = await self._backend.dequeue(queue)
            if job is None:
                await asyncio.sleep(0.05)
                continue
            task = asyncio.create_task(self._execute_job(job))
            self._running[job.job_id] = task
            task.add_done_callback(lambda _t, jid=job.job_id: self._running.pop(jid, None))  # type: ignore[misc]

    async def _execute_job(self, job: Job) -> None:
        if self._semaphore is None:
            return
        async with self._semaphore:
            job.status = JobStatus.RUNNING
            job.started_at = time.monotonic()
            await self._backend.update_job(job)

            func = self._registry.get(job.func_name)
            if func is None:
                job.status = JobStatus.FAILED
                job.result = JobResult(
                    error=RuntimeError(f"Job function '{job.func_name}' not registered")
                )
                job.finished_at = time.monotonic()
                await self._backend.enqueue_dead(job)
                _logger.error(
                    "Job %s: function '%s' not registered, sent to DLQ", job.job_id, job.func_name
                )
                return

            try:
                result = await func(*job.args, **job.kwargs)
                job.status = JobStatus.SUCCESS
                job.result = JobResult(value=result)
                job.finished_at = time.monotonic()
                await self._backend.update_job(job)
                _logger.info("Job %s completed successfully", job.job_id)
            except Exception as exc:
                job.retry_count += 1
                _logger.warning(
                    "Job %s failed (attempt %s/%s): %s",
                    job.job_id,
                    job.retry_count,
                    job.max_retries,
                    exc,
                )
                if job.retry_count < job.max_retries:
                    job.result = JobResult(error=exc)
                    job.status = JobStatus.PENDING
                    job.started_at = None
                    job.finished_at = None
                    await self._backend.update_job(job)
                    delay = self._retry_backoff_base * (2 ** (job.retry_count - 1))
                    await asyncio.sleep(delay)
                    await self._backend.enqueue(job)
                else:
                    job.status = JobStatus.FAILED
                    job.result = JobResult(error=exc)
                    job.finished_at = time.monotonic()
                    await self._backend.enqueue_dead(job)
                    _logger.exception(
                        "Job %s exhausted %s retries, sent to DLQ",
                        job.job_id,
                        job.max_retries,
                    )

    async def get_job_status(self, job_id: str) -> dict[str, Any] | None:
        job = await self._backend.get_job(job_id)
        if job is None:
            return None
        return job.to_dict()

    async def get_queue_size(self, queue: str) -> int:
        return await self._backend.size(queue)

    async def get_dead_letter_jobs(self) -> list[Job]:
        return await self._backend.dead_letter_jobs()

    async def get_dead_letter_size(self) -> int:
        return await self._backend.dead_letter_size()

    async def requeue_dead_letter(self, job_id: str) -> Job | None:
        return await self._backend.requeue_dead(job_id)

    async def update_progress(self, job_id: str, progress: float) -> None:
        job = await self._backend.get_job(job_id)
        if job is not None:
            job.progress = max(0.0, min(1.0, progress))
            await self._backend.update_job(job)

    @property
    def active_count(self) -> int:
        return len(self._running)

JobResult

Stores the result or error of a completed job.

Source code in src/forge/jobs/queue.py
class JobResult:
    """Stores the result or error of a completed job."""

    __slots__ = ("error", "value")

    def __init__(self, value: Any = None, error: BaseException | None = None) -> None:
        self.value = value
        self.error = error

JobStatus

Enumeration of possible job lifecycle states.

Source code in src/forge/jobs/queue.py
class JobStatus:
    """Enumeration of possible job lifecycle states."""

    PENDING = "pending"
    RUNNING = "running"
    SUCCESS = "success"
    FAILED = "failed"
    DEAD = "dead"

JobsModule

Bases: ForgeModule

Background jobs and scheduling module.

Provides queue-backed async job execution with automatic retry, concurrency control, dead-letter queue, progress tracking, and cron-like scheduled tasks.

Source code in src/forge/jobs/module.py
class JobsModule(ForgeModule):
    """
    Background jobs and scheduling module.

    Provides queue-backed async job execution with automatic retry,
    concurrency control, dead-letter queue, progress tracking, and
    cron-like scheduled tasks.
    """

    name = "jobs"
    dependencies: ClassVar[list[str]] = ["config"]

    def __init__(
        self,
        *,
        backend: str = "memory",
        default_retry: int = 3,
        concurrency: int = 10,
        retry_backoff_base: float = 1.0,
        redis_url: str | None = None,
        redis_key_prefix: str = "forge:jobs:",
        redis_max_connections: int = 10,
    ) -> None:
        super().__init__()
        self._backend_type = backend
        self._default_retry = default_retry
        self._concurrency = concurrency
        self._retry_backoff_base = retry_backoff_base
        self._redis_url = redis_url
        self._redis_key_prefix = redis_key_prefix
        self._redis_max_connections = redis_max_connections
        self._queue: JobQueue | None = None
        self._scheduler: Scheduler | None = None
        self._queue_backend: QueueBackend | None = None
        self._job_defs: list[JobDefinition] = []
        self._schedule_defs: list[ScheduleDefinition] = []
        self._worker_task: asyncio.Task[None] | None = None
        self._scheduler_task: asyncio.Task[None] | None = None

    async def setup(self, runtime: ForgeRuntime) -> None:
        from forge.config.module import ConfigModule
        from forge.jobs._state import set_job_queue

        config_module: ConfigModule = runtime.get(ConfigModule)  # type: ignore[assignment]
        jobs_cfg = getattr(config_module.config, "jobs", None)

        backend = self._backend_type
        default_retry = self._default_retry
        concurrency = self._concurrency
        retry_backoff_base = self._retry_backoff_base
        redis_url = self._redis_url
        redis_key_prefix = self._redis_key_prefix
        redis_max_connections = self._redis_max_connections

        if jobs_cfg is not None:
            backend = getattr(jobs_cfg, "backend", backend)
            default_retry = getattr(jobs_cfg, "default_retry", default_retry)
            concurrency = getattr(jobs_cfg, "concurrency", concurrency)
            retry_backoff_base = getattr(jobs_cfg, "retry_backoff_base", retry_backoff_base)
            redis_cfg = getattr(jobs_cfg, "redis", None)
            if redis_cfg is not None:
                redis_url = getattr(redis_cfg, "url", redis_url)
                redis_key_prefix = getattr(redis_cfg, "key_prefix", redis_key_prefix)
                redis_max_connections = getattr(redis_cfg, "max_connections", redis_max_connections)

        if backend == "redis":
            b = RedisBackend(
                redis_url=redis_url or "redis://localhost:6379/0",
                key_prefix=redis_key_prefix,
                max_connections=redis_max_connections,
            )
            await b.connect()
            self._queue_backend = b
        else:
            self._queue_backend = MemoryBackend()

        self._queue = JobQueue(
            backend=self._queue_backend,
            default_retry=default_retry,
            concurrency=concurrency,
            retry_backoff_base=retry_backoff_base,
        )

        for jd in self._job_defs:
            self._queue.register(jd._func.__qualname__, jd._func)

        await self._queue.start()
        set_job_queue(self._queue)

        self._scheduler = Scheduler()
        for sd in self._schedule_defs:
            self._scheduler.register(
                name=sd._func.__qualname__,
                cron_expression=sd._cron,
                func=sd._func,
                queue=sd._queue,
            )

        self._worker_task = asyncio.create_task(self._queue.process("default"))
        self._scheduler_task = asyncio.create_task(self._scheduler.start())

        _logger.info("Jobs module initialized (backend=%s, concurrency=%d)", backend, concurrency)

    async def teardown(self) -> None:
        from forge.jobs._state import set_job_queue

        if self._scheduler_task is not None:
            self._scheduler_task.cancel()
            with contextlib.suppress(asyncio.CancelledError):
                await self._scheduler_task
            self._scheduler_task = None

        if self._scheduler is not None:
            await self._scheduler.stop()
            self._scheduler = None

        if self._worker_task is not None:
            self._worker_task.cancel()
            with contextlib.suppress(asyncio.CancelledError):
                await self._worker_task
            self._worker_task = None

        if self._queue is not None:
            await self._queue.stop()
            self._queue = None

        set_job_queue(None)

    def health_check(self) -> HealthResult:
        if self._queue is None:
            return HealthResult.error("Jobs module not initialized")
        return HealthResult.ok()

    def register_job(self, job_def: JobDefinition) -> None:
        self._job_defs.append(job_def)

    def register_schedule(self, schedule_def: ScheduleDefinition) -> None:
        self._schedule_defs.append(schedule_def)

    async def enqueue(
        self,
        func_name: str,
        queue: str = "default",
        args: tuple[Any, ...] = (),
        kwargs: dict[str, Any] | None = None,
        max_retries: int | None = None,
    ) -> Any:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        return await self._queue.enqueue(
            queue=queue,
            func_name=func_name,
            args=args,
            kwargs=kwargs or {},
            max_retries=max_retries,
        )

    async def get_job_status(self, job_id: str) -> dict[str, Any] | None:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        return await self._queue.get_job_status(job_id)

    async def get_queue_size(self, queue: str = "default") -> int:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        return await self._queue.get_queue_size(queue)

    async def get_dead_letter_jobs(self) -> list[Any]:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        return await self._queue.get_dead_letter_jobs()

    async def get_dead_letter_size(self) -> int:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        return await self._queue.get_dead_letter_size()

    async def requeue_dead_letter(self, job_id: str) -> Any:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        return await self._queue.requeue_dead_letter(job_id)

    async def update_progress(self, job_id: str, progress: float) -> None:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        await self._queue.update_progress(job_id, progress)

    @property
    def queue(self) -> JobQueue:
        if self._queue is None:
            raise RuntimeError("Jobs module not initialized")
        return self._queue

    @property
    def scheduler_instance(self) -> Scheduler:
        if self._scheduler is None:
            raise RuntimeError("Jobs module not initialized")
        return self._scheduler

MemoryBackend

Bases: QueueBackend

In-memory queue backend for development.

Uses OrderedDict per-queue for FIFO ordering. Stores both active queues and a dead-letter queue.

Source code in src/forge/jobs/queue.py
class MemoryBackend(QueueBackend):
    """
    In-memory queue backend for development.

    Uses OrderedDict per-queue for FIFO ordering.
    Stores both active queues and a dead-letter queue.
    """

    def __init__(self, max_dead_letter: int = 1000) -> None:
        self._queues: dict[str, OrderedDict[str, Job]] = {}
        self._jobs: dict[str, Job] = {}
        self._dead_letter: OrderedDict[str, Job] = OrderedDict()
        self._max_dead_letter = max_dead_letter

    async def enqueue(self, job: Job) -> None:
        self._jobs[job.job_id] = job
        q = self._queues.setdefault(job.queue, OrderedDict())
        q[job.job_id] = job

    async def dequeue(self, queue: str) -> Job | None:
        q = self._queues.get(queue)
        if not q:
            return None
        while q:
            job_id, job = next(iter(q.items()))
            if job.status == JobStatus.PENDING:
                del q[job_id]
                return job
            del q[job_id]
        return None

    async def size(self, queue: str) -> int:
        q = self._queues.get(queue)
        if not q:
            return 0
        return sum(1 for j in q.values() if j.status == JobStatus.PENDING)

    async def get_job(self, job_id: str) -> Job | None:
        return self._jobs.get(job_id)

    async def update_job(self, job: Job) -> None:
        self._jobs[job.job_id] = job

    async def enqueue_dead(self, job: Job) -> None:
        job.status = JobStatus.DEAD
        self._dead_letter[job.job_id] = job
        if len(self._dead_letter) > self._max_dead_letter:
            self._dead_letter.popitem(last=False)

    async def dead_letter_size(self) -> int:
        return len(self._dead_letter)

    async def dead_letter_jobs(self) -> list[Job]:
        return list(self._dead_letter.values())

    async def requeue_dead(self, job_id: str) -> Job | None:
        job = self._dead_letter.get(job_id)
        if job is None:
            return None
        del self._dead_letter[job_id]
        job.status = JobStatus.PENDING
        job.retry_count = 0
        job.result = None
        job.started_at = None
        job.finished_at = None
        await self.enqueue(job)
        return job

    async def close(self) -> None:
        self._queues.clear()
        self._jobs.clear()
        self._dead_letter.clear()

QueueBackend

Bases: ABC

Abstract base class for job queue storage backends.

Source code in src/forge/jobs/queue.py
class QueueBackend(ABC):
    """Abstract base class for job queue storage backends."""

    @abstractmethod
    async def enqueue(self, job: Job) -> None: ...

    @abstractmethod
    async def dequeue(self, queue: str) -> Job | None: ...

    @abstractmethod
    async def size(self, queue: str) -> int: ...

    @abstractmethod
    async def get_job(self, job_id: str) -> Job | None: ...

    @abstractmethod
    async def update_job(self, job: Job) -> None: ...

    @abstractmethod
    async def enqueue_dead(self, job: Job) -> None: ...

    @abstractmethod
    async def dead_letter_size(self) -> int: ...

    @abstractmethod
    async def dead_letter_jobs(self) -> list[Job]: ...

    @abstractmethod
    async def requeue_dead(self, job_id: str) -> Job | None: ...

    @abstractmethod
    async def close(self) -> None: ...

RedisBackend

Bases: QueueBackend

Optional Redis-backed queue for production.

Uses redis.asyncio with connection pooling and JSON serialization. Requires the redis optional extra.

Source code in src/forge/jobs/queue.py
class RedisBackend(QueueBackend):
    """
    Optional Redis-backed queue for production.

    Uses redis.asyncio with connection pooling and JSON serialization.
    Requires the ``redis`` optional extra.
    """

    def __init__(
        self,
        redis_url: str = "redis://localhost:6379/0",
        key_prefix: str = "forge:jobs:",
        max_connections: int = 10,
        max_dead_letter: int = 1000,
    ) -> None:
        self._url = redis_url
        self._prefix = key_prefix
        self._max_connections = max_connections
        self._max_dead_letter = max_dead_letter
        self._pool: Any = None
        self._client: Any = None

    async def connect(self) -> None:
        import redis.asyncio as aioredis

        self._pool = aioredis.ConnectionPool.from_url(
            self._url,
            max_connections=self._max_connections,
            decode_responses=True,
        )
        self._client = aioredis.Redis(connection_pool=self._pool)

    async def _ensure_connected(self) -> Any:
        if self._client is None:
            await self.connect()
        return self._client

    async def enqueue(self, job: Job) -> None:
        client = await self._ensure_connected()
        import json

        data = json.dumps(job.to_dict())
        await client.rpush(f"{self._prefix}queue:{job.queue}", job.job_id)
        await client.set(f"{self._prefix}job:{job.job_id}", data)

    async def dequeue(self, queue: str) -> Job | None:
        client = await self._ensure_connected()
        import json

        job_id = await client.lpop(f"{self._prefix}queue:{queue}")
        if not job_id:
            return None
        raw = await client.get(f"{self._prefix}job:{job_id}")
        if not raw:
            return None
        data = json.loads(raw)
        job = self._reconstruct_job(data)
        return job

    async def size(self, queue: str) -> int:
        client = await self._ensure_connected()
        return await client.llen(f"{self._prefix}queue:{queue}")  # type: ignore[no-any-return]

    async def get_job(self, job_id: str) -> Job | None:
        client = await self._ensure_connected()
        import json

        raw = await client.get(f"{self._prefix}job:{job_id}")
        if not raw:
            return None
        data = json.loads(raw)
        return self._reconstruct_job(data)

    async def update_job(self, job: Job) -> None:
        client = await self._ensure_connected()
        import json

        data = json.dumps(job.to_dict())
        await client.set(f"{self._prefix}job:{job.job_id}", data)

    async def enqueue_dead(self, job: Job) -> None:
        client = await self._ensure_connected()
        import json

        job.status = JobStatus.DEAD
        data = json.dumps(job.to_dict())
        key = f"{self._prefix}dead:{job.job_id}"
        await client.set(key, data)
        await client.rpush(f"{self._prefix}dead_letter", job.job_id)
        count = await client.llen(f"{self._prefix}dead_letter")
        if count > self._max_dead_letter:
            old_id = await client.lpop(f"{self._prefix}dead_letter")
            if old_id:
                await client.delete(f"{self._prefix}dead:{old_id}")

    async def dead_letter_size(self) -> int:
        client = await self._ensure_connected()
        return await client.llen(f"{self._prefix}dead_letter")  # type: ignore[no-any-return]

    async def dead_letter_jobs(self) -> list[Job]:
        client = await self._ensure_connected()
        import json

        ids = await client.lrange(f"{self._prefix}dead_letter", 0, -1)
        jobs: list[Job] = []
        for jid in ids:
            raw = await client.get(f"{self._prefix}dead:{jid}")
            if raw:
                data = json.loads(raw)
                jobs.append(self._reconstruct_job(data))
        return jobs

    async def requeue_dead(self, job_id: str) -> Job | None:
        client = await self._ensure_connected()
        import json

        key = f"{self._prefix}dead:{job_id}"
        raw = await client.get(key)
        if not raw:
            return None
        data = json.loads(raw)
        job = self._reconstruct_job(data)
        job.status = JobStatus.PENDING
        job.retry_count = 0
        job.result = None
        job.started_at = None
        job.finished_at = None
        await client.delete(key)
        await client.lrem(f"{self._prefix}dead_letter", 0, job_id)
        await self.enqueue(job)
        return job

    async def close(self) -> None:
        if self._client:
            await self._client.aclose()
        if self._pool:
            await self._pool.disconnect()
        self._client = None
        self._pool = None

    @staticmethod
    def _reconstruct_job(data: dict[str, Any]) -> Job:
        job = Job(
            job_id=data["job_id"],
            queue=data["queue"],
            func_name=data["func_name"],
            args=(),
            kwargs={},
            max_retries=data.get("max_retries", 0),
        )
        job.retry_count = data.get("retry_count", 0)
        job.status = data.get("status", JobStatus.PENDING)
        job.progress = data.get("progress", 0.0)
        job.created_at = data.get("created_at", 0.0)
        job.started_at = data.get("started_at")
        job.finished_at = data.get("finished_at")
        return job

ScheduleDefinition

Wrapper returned by the @schedule decorator.

Registers the function with the Scheduler on runtime init. The decorated function is still callable directly.

Source code in src/forge/jobs/module.py
class ScheduleDefinition:
    """
    Wrapper returned by the ``@schedule`` decorator.

    Registers the function with the Scheduler on runtime init.
    The decorated function is still callable directly.
    """

    def __init__(
        self,
        func: Callable[..., Any],
        cron: str,
        queue: str,
    ) -> None:
        self._func = func
        self._cron = cron
        self._queue = queue
        self._cron_expr: CronExpression | None = None
        self._scheduled_job: ScheduledJob | None = None
        functools.update_wrapper(self, func)

    async def __call__(self, *args: Any, **kwargs: Any) -> Any:
        result = self._func(*args, **kwargs)
        if asyncio.iscoroutine(result):
            return await result
        return result

ScheduledJob

A job that runs on a cron schedule.

Source code in src/forge/jobs/scheduler.py
class ScheduledJob:
    """A job that runs on a cron schedule."""

    def __init__(
        self,
        name: str,
        cron: CronExpression,
        func: Callable[..., Any],
        queue: str = "scheduled",
    ) -> None:
        self.name = name
        self.cron = cron
        self.func = func
        self.queue = queue
        self.last_run: datetime | None = None
        self.next_run: datetime | None = None
        self.run_count: int = 0
        self.last_error: str | None = None

Scheduler

Cron-like scheduler that evaluates scheduled jobs at minute boundaries.

Runs a loop that checks every ~1 second whether the current minute matches any registered schedule and executes matching jobs.

Source code in src/forge/jobs/scheduler.py
class Scheduler:
    """
    Cron-like scheduler that evaluates scheduled jobs at minute boundaries.

    Runs a loop that checks every ~1 second whether the current minute
    matches any registered schedule and executes matching jobs.
    """

    def __init__(self) -> None:
        self._jobs: dict[str, ScheduledJob] = {}
        self._running: dict[str, asyncio.Task[None]] = {}
        self._started = False

    def register(
        self,
        name: str,
        cron_expression: str,
        func: Callable[..., Any],
        queue: str = "scheduled",
    ) -> ScheduledJob:
        cron = CronExpression(cron_expression)
        job = ScheduledJob(name=name, cron=cron, func=func, queue=queue)
        self._jobs[name] = job
        return job

    async def start(self) -> None:
        self._started = True
        _logger.info("Scheduler started with %d job(s)", len(self._jobs))
        _last_minute: int = -1
        try:
            while self._started:
                now = datetime.now(tz=UTC)
                current_minute = now.minute + now.hour * 60 + now.day * 60 * 24
                if current_minute != _last_minute:
                    _last_minute = current_minute
                    await self._check_and_fire(now)
                await asyncio.sleep(0.5)
        except asyncio.CancelledError:
            pass

    async def stop(self) -> None:
        self._started = False
        for task in self._running.values():
            task.cancel()
        if self._running:
            await asyncio.gather(*self._running.values(), return_exceptions=True)
        self._running.clear()
        _logger.info("Scheduler stopped")

    async def _check_and_fire(self, now: datetime) -> None:
        _min_dedup_seconds = 60
        for job in self._jobs.values():
            if not job.cron.matches(now):
                continue
            if job.last_run is not None:
                since_last = (now - job.last_run).total_seconds()
                if since_last < _min_dedup_seconds:
                    continue
            task = asyncio.create_task(self._run_scheduled(job))
            self._running[job.name] = task
            task.add_done_callback(lambda _t, n=job.name: self._running.pop(n, None))  # type: ignore[misc]

    async def _run_scheduled(self, job: ScheduledJob) -> None:
        now = datetime.now(tz=UTC)
        job.last_run = now
        job.next_run = job.cron.next_after(now)
        try:
            result = job.func()
            if asyncio.iscoroutine(result):
                await result
            job.run_count += 1
            job.last_error = None
            _logger.info("Scheduled job '%s' completed (run #%d)", job.name, job.run_count)
        except Exception as exc:
            job.last_error = str(exc)
            _logger.exception("Scheduled job '%s' failed", job.name)

    @property
    def jobs(self) -> dict[str, ScheduledJob]:
        return dict(self._jobs)

    def get_status(self) -> list[dict[str, Any]]:
        return [
            {
                "name": job.name,
                "cron": job.cron._raw,
                "last_run": job.last_run.isoformat() if job.last_run else None,
                "next_run": job.next_run.isoformat() if job.next_run else None,
                "run_count": job.run_count,
                "last_error": job.last_error,
            }
            for job in self._jobs.values()
        ]

Functions:

job

job(_func: Callable[..., Any] | None = None, *, queue: str = 'default', retry: int = 3) -> Any

Decorator for defining an async background job.

Parameters

queue: Name of the queue this job belongs to. retry: Maximum number of retry attempts on failure.

Usage::

@job(queue="emails", retry=3)
async def send_welcome_email(user_id: int):
    ...

await send_welcome_email.enqueue(user_id=123)
Source code in src/forge/jobs/module.py
def job(
    _func: Callable[..., Any] | None = None,
    *,
    queue: str = "default",
    retry: int = 3,
) -> Any:
    """
    Decorator for defining an async background job.

    Parameters
    ----------
    queue:
        Name of the queue this job belongs to.
    retry:
        Maximum number of retry attempts on failure.

    Usage::

        @job(queue="emails", retry=3)
        async def send_welcome_email(user_id: int):
            ...

        await send_welcome_email.enqueue(user_id=123)
    """

    def decorator(fn: Callable[..., Any]) -> JobDefinition:
        return JobDefinition(func=fn, queue=queue, max_retries=retry)

    if _func is not None:
        return decorator(_func)
    return decorator

schedule

schedule(_func: Callable[..., Any] | None = None, *, cron: str = '* * * * *', queue: str = 'scheduled') -> Any

Decorator for cron-like scheduled tasks.

Parameters

cron: Standard five-field cron expression. queue: Queue name for tracking.

Usage::

@schedule(cron="0 9 * * *")
async def daily_report():
    ...
Source code in src/forge/jobs/module.py
def schedule(
    _func: Callable[..., Any] | None = None,
    *,
    cron: str = "* * * * *",
    queue: str = "scheduled",
) -> Any:
    """
    Decorator for cron-like scheduled tasks.

    Parameters
    ----------
    cron:
        Standard five-field cron expression.
    queue:
        Queue name for tracking.

    Usage::

        @schedule(cron="0 9 * * *")
        async def daily_report():
            ...
    """

    def decorator(fn: Callable[..., Any]) -> ScheduleDefinition:
        return ScheduleDefinition(func=fn, cron=cron, queue=queue)

    if _func is not None:
        return decorator(_func)
    return decorator