Skip to content

Scheduler

The Scheduler module provides background job processing with named tasks, cron expressions, priority queues, retries, and pluggable backends (several can run at once).

Enable:

app = create_kirak_app(models_path="...", modules=["scheduler"])

The scheduler starts automatically with the application and stops gracefully on shutdown.


kirak.json lists the queue backends the app runs as provider instances and names the default. Each instance has a name (the key) and a type (which implementation to use). Several can be active at once, for example the database plus Redis.

{
"scheduler": {
"default_provider": "db",
"providers": {
"db": { "type": "database" },
"redis": { "type": "redis" }
},
"poll_interval": 5,
"concurrency": 4
}
}
Built-in type Extra install Notes
database None Uses your existing DB. Zero extra infra for simple apps. Polling. No settings; the job table is always scheduler_jobs. At most one database provider is allowed.
redis pip install "kirak[scheduler-redis]" Polling. Setting: redis_ssl (default false).
rabbitmq pip install "kirak[scheduler-rabbitmq]" Push-based, not polled – see Worker below.
  • Default: a call without provider uses default_provider.
  • Per call: pass "provider": "<instance name>" to enqueue (or in the JSON body of POST /scheduler/enqueue). The value is the instance name (redis), not the type. @kirak.scheduler.cron(..., provider="redis") fires a cron on a named backend.
  • Allow-list: only instances listed in kirak.json can be used. Anything else fails with PROVIDER_NOT_CONFIGURED (400).
  • Every backend runs. Each configured backend is connected at startup and gets its own worker loop with its own concurrency limit, so total in-flight jobs can reach concurrency times the number of backends. If any backend cannot connect, startup fails instead of leaving jobs piling up on a backend with no worker.
  • Connection URLs are secrets, never read from kirak.json. Each redis and rabbitmq instance reads KIRAK_SCHEDULER_<INSTANCE>_URL, e.g. KIRAK_SCHEDULER_REDIS_URL for an instance named redis. When unset they default to redis://localhost:6379 and amqp://guest:guest@localhost/.
  • Database-only features – the jobs routes (/scheduler/jobs) and dynamic schedules (/scheduler/schedules) – need a provider of type database, which does not have to be the default. Dynamic schedules always run on that database provider, whichever provider is the default, because each run is linked back to its schedule by the job’s integer row id. Jobs on Redis or RabbitMQ are not queryable through the API.
KIRAK_SCHEDULER_REDIS_URL=redis://localhost:6379

All tasks must be registered in on_kirak_ready – before the worker starts.

def on_kirak_ready(kirak):
@kirak.scheduler.task("send_welcome_email")
async def send_welcome_email(payload: dict) -> None:
user_id = payload["user_id"]
# A job has no request and no user: without this, Kirak calls run as guest.
token = set_user_context({"role": "system", "user_id": None, "token": None})
try:
user = await kirak.fetch("users", {"id": user_id})
finally:
reset_user_context(token)
email = user["data"][0]["email"]
await kirak.notifications.send_email({
"to": email,
"template_name": "welcome",
})

(set_user_context and reset_user_context come from kirak.core.context.)

The task decorator registers the function under a string name. The function receives payload: dict and should return None. Any exception is caught by the worker and triggers a retry.

The worker sets no identity for a job, so Kirak operations called from it run as guest. Set one as above: the user the job acts for, so their access rules apply, or system for work no user owns – system needs a {"role": "system"} rule in the model’s access block for each operation it uses (the built-in users model has one for fetch). See Who a direct call runs as.

def on_kirak_ready(kirak):
@kirak.scheduler.cron("0 9 * * 1-5") # 09:00 UTC, Mon-Fri
async def morning_digest(payload: dict) -> None:
# Send digest email to all users...
pass
@kirak.scheduler.cron("*/15 * * * *") # every 15 minutes
async def sync_inventory(payload: dict) -> None:
pass
@kirak.scheduler.cron("0 0 1 * *", name="monthly_report")
async def monthly_report_fn(payload: dict) -> None:
pass

5-field cron syntax: minute hour day-of-month month day-of-week

Pattern Meaning
* * * * * Every minute
0 * * * * Every hour
0 9 * * 1-5 09:00 UTC, weekdays
*/15 * * * * Every 15 minutes
0 0 * * 0 Midnight every Sunday
0 0 1 * * Midnight on the 1st of every month

The cron expression is checked at registration time and raises ValueError immediately when:

  • it does not have 5 fields, or a token is not a number, range, list or step (names such as MON / JAN and shortcuts such as @daily are not supported);
  • a value is outside its field’s range: minute 0-59, hour 0-23, day 1-31, month 1-12, weekday 0-7 (0 and 7 are Sunday); a step is below 1; or a range such as 5-3 matches nothing;
  • it would not fire within the next 2 years (0 0 31 2 *).

An expression must fire at least once every 2 years: after a run, the next one is searched for only 2 years ahead. 0 0 29 2 * (leap day) registers, but after it fires it does not fire again.


# From a hook:
@kirak.on("users").hook("after_create")
async def queue_welcome(result):
user_id = result["data"]["id"]
await kirak.scheduler.enqueue({"task": "send_welcome_email", "payload": {"user_id": user_id}})
return result
# From a route or anywhere with kirak access:
result = await kirak.scheduler.enqueue({
"task": "send_welcome_email",
"payload": {"user_id": 42},
"queue": "high", # optional -- named queue (default: "default")
"priority": 10, # optional -- higher = picked up first (default: 0)
"run_at": int(time.time()) + 300, # optional -- Unix timestamp for deferred execution
"max_retries": 5, # optional -- override default max retries
})
job_id = result["data"]["job_id"]
print(f"Job enqueued: {job_id}")

Parameters:

Parameter Default Description
task – Required. Must match a registered @kirak.scheduler.task(name). A missing task raises KeyError, not a KirakException.
payload {} Dict passed to the task function.
provider default_provider Backend instance name (see Providers).
queue "default" Named queue (default from default_queue in kirak.json). Runs in this app only if it is default_queue or listed in queues (see Multiple Queues).
priority 0 Higher value = picked up sooner. See Priority – the effect depends on the backend.
run_at now Unix timestamp. Deferred if in the future.
max_retries 3 Maximum number of runs, including the first (overrides max_retries in kirak.json). 3 means one run and at most two retries.

The result is {"data": {"job_id": ..., "provider": ..., "task": ..., "payload": ..., "queue": ..., "priority": ..., "run_at": ..., "max_retries": ...}} – the values the job was stored with, after any before_enqueue hook. Keep the provider with the job id if you need to know where it went. enqueue raises a KirakException:

Code Status When
SCHEDULER_NOT_STARTED 503 The scheduler has not been started
PROVIDER_NOT_CONFIGURED 400 provider is not listed in kirak.json
TASK_NOT_REGISTERED 400 The task name is not registered
INVALID_PAYLOAD 400 The payload contains a secret-looking key (passwords, tokens, …)
QUEUE_BACKEND_ERROR 503 The backend failed (details["provider"] names it)
Backend How priority is used
database Due jobs are picked ORDER BY priority DESC, run_at ASC: a higher priority always goes first.
redis Jobs are ordered by run_at - priority * 0.001, so priority only reorders jobs due within about the same second. It does not jump a job ahead of one that became due earlier.
rabbitmq Ignored. Messages are delivered in queue order.

The worker is started automatically when "scheduler" is in modules. It runs one loop per configured backend, each acking and retrying jobs on the backend that delivered them. On the database and Redis backends, a loop:

  1. Polls the backend queue every poll_interval seconds (kirak.json, default 5).
  2. Runs up to concurrency jobs simultaneously (kirak.json, default 4).
  3. Stops a job that runs longer than job_timeout seconds (kirak.json, default 3600) and treats it as a failure, so it is retried like any other.
  4. Retries a failed job, or marks it failed when it has no runs left (see Retry Behavior).

On the RabbitMQ backend, the worker is push-based instead: it calls the backend’s start_consuming() and jobs are delivered as they arrive, rather than polled – poll_interval does not apply. Retries, delayed jobs and permanent failures use two extra queues that Kirak declares beside each named queue q:

Queue Holds
q.delay Jobs waiting for their time: a retry after retry_backoff, or a job enqueued with a future run_at. Each is published with an expiry, and RabbitMQ moves it back to q when the time is up. Nothing consumes this queue.
q.failed Jobs that used up all their attempts, with the last error in the message body. Nothing consumes this queue: inspect the messages in the broker, or publish them back to q to replay them.

A delayed job is never run early, but it can run late: RabbitMQ only expires the message at the front of a queue, so a job can wait behind one with a longer delay. This is a difference from the database and Redis backends, which order by due time. Delays longer than 30 days are served in 30-day steps. A job that is running when the app shuts down is given back to the queue and runs again on another worker. If the process stops between publishing a retry and acknowledging the original, the job may run twice; it is never lost.

Every minute the worker also puts back jobs that a crashed worker left in running (see Job Timeout and Crash Recovery).

The worker also checks for due DB-defined cron schedules (see Dynamic Schedules below) every 60 seconds, independent of poll_interval.


If a worker process dies while running a job, the job stays running with nobody working on it. Kirak recovers these jobs without disturbing jobs that other app instances are still running:

  • Every job has a time limit, job_timeout in kirak.json (seconds, default 3600). The worker stops a job that exceeds it and treats that as a failure, so the job is retried with retry_backoff like any other failure. The error is Task '<name>' exceeded job_timeout (<n>s).
  • A job still running more than 60 seconds after job_timeout is considered orphaned and is put back to pending. This is checked when the worker starts and then every minute, on the database and Redis backends. RabbitMQ needs no check: the broker requeues messages that were never acknowledged.
  • A restarting instance no longer touches a job that another instance started recently.

So after a crash an orphaned job is picked up within about job_timeout plus a few minutes, not at the next restart. Set job_timeout above the longest job you run; a job that needs longer than the limit is stopped and retried until max_retries is used up, then marked failed. Recovery counts as an attempt, so an orphaned job still respects max_retries.

Stopping a job cancels its coroutine. A task that blocks the event loop with synchronous code cannot be cancelled, so keep long blocking work in a thread or a separate process.


max_retries is the maximum number of runs, including the first. Each run increments the job’s attempts when it starts. After a failure the job is retried only while attempts < max_retries, after waiting retry_backoff x attempts seconds. With the defaults (max_retries 3, retry_backoff 60):

Run that failed What happens
1 Retry after retry_backoff x 1 seconds (60 s)
2 Retry after retry_backoff x 2 seconds (120 s)
3 (= max_retries) Job status -> failed

The on_job_failure hook fires on every failure, not only the last one. retry_backoff is a kirak.json manifest setting, not an environment variable.


@kirak.scheduler.hook("before_enqueue")
async def before_enqueue(payload):
# payload is an envelope; the job is in payload["data"]:
# task, payload, queue, priority, run_at, max_retries, provider
# Changes to payload, queue, priority, run_at and max_retries are what gets queued.
# task and provider are fixed. Raise to refuse the job.
payload["data"]["queue"] = "high" if payload["data"]["priority"] > 5 else payload["data"]["queue"]
return payload
@kirak.scheduler.hook("after_enqueue")
async def after_enqueue(payload):
# payload["data"] has the same fields as before_enqueue, as queued, plus job_id
return payload
@kirak.scheduler.hook("before_job")
async def before_job(payload):
# payload is the job payload dict -- mutations are forwarded to the task function.
return payload
@kirak.scheduler.hook("after_job")
async def after_job(payload):
# payload is the (possibly mutated) payload that was passed to the task function.
return payload
@kirak.scheduler.hook("on_job_failure")
async def on_job_failure(payload):
job = payload.get("job", {}) # dict with full job fields
error = payload.get("error", "")
# Fires on every failure, including retryable ones.
# Check job["attempts"] >= job["max_attempts"] to act only on permanent failure.
return payload

before_enqueue must return the envelope it received (changed or not), {"data": {...}} with only the fields it changes, or None. Any other return value raises TypeError and nothing is queued.


When a database provider is configured, Kirak uses a scheduler_jobs table (the name is fixed; a table setting used to exist and is now rejected). It holds two kinds of rows: jobs, and dynamic schedule definitions (rows with cron_expression set).

scheduler_jobs
+-- id BIGINT AUTO_INCREMENT / BIGSERIAL primary key
+-- queue VARCHAR(100)
+-- task VARCHAR(500) -- registered task name (api_call for schedules)
+-- payload TEXT / JSONB -- job payload
+-- status VARCHAR(20) -- jobs: pending, running, done, failed
| -- schedule definitions: scheduled, paused
+-- priority INT
+-- run_at BIGINT -- Unix epoch seconds (not a DATETIME -- timezone-agnostic)
+-- attempts INT -- runs started so far (the first run counts)
+-- max_attempts INT -- max_retries: the most runs allowed
+-- error TEXT -- full traceback of the last failure
+-- started_at DATETIME / TIMESTAMP
+-- completed_at DATETIME / TIMESTAMP
+-- cron_expression VARCHAR(100) -- set only on schedule definitions
+-- last_run BIGINT -- schedule definitions: Unix seconds it last fired (or was created)
+-- schedule_id BIGINT -- on a job: the schedule definition that spawned it
+-- created_at DATETIME / TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP / NOW()

A second table, scheduler_cron_locks (name, locked_until), is the database backend’s cron lock: when several app instances run, each cron firing on that backend (code-registered or dynamic) is enqueued by only one of them.


Method Path Description
POST /scheduler/enqueue Enqueue a job via HTTP
GET /scheduler/jobs List recent jobs and their status
GET /scheduler/jobs/{job_id} Get status of a specific job
DELETE /scheduler/jobs/{job_id} Cancel a pending job (sets status to failed; 409 if it is no longer pending)
GET /scheduler/registered List all registered tasks, cron expressions and configured providers
GET /scheduler/schedules List DB-defined dynamic cron schedules (see below)
POST /scheduler/schedules Create a dynamic cron schedule
GET /scheduler/schedules/{id} Get a schedule by id
PUT /scheduler/schedules/{id} Update a schedule
DELETE /scheduler/schedules/{id} Delete a schedule
POST /scheduler/schedules/{id}/run Trigger a schedule immediately, outside its cron timing

All scheduler HTTP endpoints require an admin, system, or superadmin role. Every response uses the standard envelope: {"statusCode": 200, "status": "success", "message": ..., "data": ...} on success and the standard error envelope on failure.

Route Request data in the response
POST /scheduler/enqueue JSON body: task (required, else 422), payload, queue, priority, provider, run_after (delay in seconds, default 0). There is no run_at or max_retries here: the job gets max_retries from kirak.json. {job_id, provider, run_at}
GET /scheduler/jobs Query: status, queue, schedule_id, limit (1-500, default 50), offset. Newest first. Schedule definition rows are not included. list of job rows
GET /scheduler/jobs/{job_id} – the row
DELETE /scheduler/jobs/{job_id} – {job_id}
GET /scheduler/registered – {tasks, crons, providers, default_provider}
GET /scheduler/schedules Query: limit (1-500, default 50), offset list of schedule rows
POST /scheduler/schedules JSON body: cron_expression (required, checked as in Cron Task; 422 if invalid), path (required), method (default GET), headers, body, name, queue the schedule row
GET /scheduler/schedules/{id} – the schedule row
PUT /scheduler/schedules/{id} Any of the POST fields, plus status: scheduled or paused. Fields you leave out keep their value. the schedule row
DELETE /scheduler/schedules/{id} – {id}
POST /scheduler/schedules/{id}/run – {job_id, run_at}

A schedule row has its payload decoded to {name, method, path, headers, body} and a computed next_run_at (Unix seconds, null when paused).


Alongside code-defined @kirak.scheduler.cron(...) tasks (registered at startup), Kirak also supports cron schedules created at runtime through the /scheduler/schedules endpoints above – useful when the schedule itself needs to be admin-configurable without a redeploy.

A dynamic schedule is a row in scheduler_jobs itself, with cron_expression set and status scheduled or paused. When it is due, the worker enqueues a job for the built-in api_call task, which makes an in-process HTTP call back into your own FastAPI app (via httpx.ASGITransport – no real network egress, so it can’t be pointed at an external URL). The worker checks for due dynamic schedules every 60 seconds, independent of the regular polling interval. A new schedule first fires at the first time its expression matches after it was created.

Terminal window
POST /scheduler/schedules
{
"cron_expression": "0 3 * * *",
"path": "/reports/nightly",
"method": "POST"
}
  • Who the call runs as. The call carries only the schedule’s own headers. With none, the target route sees a guest. To call a protected route, put a credential in headers (for example an X-API-Key for a dedicated key) – it is stored in plain text in scheduler_jobs.payload and returned by the schedule routes, so use a key scoped to what the schedule needs. The payload secret-key check only looks at top-level payload keys, so it does not catch headers.Authorization.
  • Failures. Any response with status 400 or higher fails the job, which is then retried like any other job (see Retry Behavior).
  • api_call is an ordinary task. It is listed by GET /scheduler/registered and an admin can enqueue it directly with a payload of {method, path, headers, body}.

Use this when you want a schedule that an admin can create, edit, or trigger on demand via the API, rather than one baked into application code at deploy time.


Use named queues to separate workloads:

await kirak.scheduler.enqueue({"task": "send_email", "payload": payload, "queue": "emails"})
await kirak.scheduler.enqueue({"task": "sync_data", "payload": payload, "queue": "sync"})
await kirak.scheduler.enqueue({"task": "send_report", "payload": payload, "queue": "low-priority"})

List the extra queues in kirak.json and the worker consumes them:

{
"scheduler": {
"default_queue": "default",
"queues": ["emails", "sync", "low-priority"]
}
}

The worker consumes default_queue plus every queue in queues, with one loop per queue on each provider. Listing default_queue again changes nothing. Each provider’s concurrency limit is shared by all of its queues, so a busy queue can use up the slots the others need; raise concurrency if you split work this way.

A queue name is 1 to 100 characters. Enqueueing to a queue that is not listed still works, because another process may consume it, but the job waits until a worker for that queue runs, and Kirak logs a warning once per queue naming the queues this app does consume.


The shared contract, the credential check (check()), entry points and testing are covered in Adding a Provider.

Subclass QueueBackend (kirak/scheduler/backends/base.py), register the type, and list an instance in kirak.json.

from kirak.scheduler.backends.base import Job, QueueBackend
class SqsBackend(QueueBackend):
TYPE_NAME = "sqs" # the "type" in kirak.json
SECRET_FIELDS = ("queue_url",) # read from KIRAK_SCHEDULER_<INSTANCE>_QUEUE_URL
def __init__(self, config, scheduler=None):
super().__init__(config, scheduler)
self._require("region") # fail early on missing config
async def connect(self): ...
async def disconnect(self): ...
async def enqueue(self, task, payload, *, queue="default", priority=0, run_at, max_attempts=3) -> str: ...
async def dequeue(self, queue) -> "Job | None": ...
async def ack(self, job_id): ...
async def nack(self, job_id, error, retry_at): ...
async def reset_stale(self, queue="default") -> int: ...
async def start_consuming(self, queue, handler): ...

Call self._validate_payload(payload) at the top of enqueue so secrets never reach the queue. reset_stale runs at startup and every minute while other app instances may be running jobs on the same queue, so a backend shared between processes must record when a job started and reset only jobs that have been running for more than self._stale_after() seconds. Module-wide settings such as poll_interval and concurrency are on self.scheduler.config. Override try_acquire_cron_lock if the backend can dedupe cron firings across app instances, and is_push_based if it delivers messages instead of being polled.

async def on_kirak_ready(kirak):
kirak.scheduler.register_provider("sqs", SqsBackend)
"scheduler": {
"default_provider": "queue",
"providers": { "queue": { "type": "sqs", "region": "us-east-1" } }
}

register_provider also works as a decorator, raises ValueError if the type name is taken (including by a built-in) and TypeError if the class does not subclass QueueBackend. A backend can instead be published as a pip package with an entry point in the group kirak.scheduler_providers, so users need only pip install plus a type in kirak.json:

[project.entry-points."kirak.scheduler_providers"]
sqs = "your_package.sqs:SqsBackend"

To contribute a backend to Kirak core, add the class under kirak/scheduler/backends/, add one line to BUILTIN_BACKENDS in kirak/scheduler/backends/registry.py, and add tests.