Skip to content

Architecture

Components

flowchart TB
    subgraph API["cuectl serve (FastAPI)"]
        R[Routers] --> S[Services]
    end
    subgraph W["cuectl worker"]
        Q[Job runner] --> S2[Services]
    end
    S -- "rows + jobs, one transaction" --> DB[(PostgreSQL / SQLite)]
    Q -- "claim (SKIP LOCKED)" --> DB
    S2 --> CH[Channel providers]
    S2 -.optional.-> AI[Pydantic AI]
    CH --> EXT[FCM · Twilio · SMTP · Telegram · webhooks]
Package Responsibility
cue_notify.core Pure logic, no I/O: JSON Logic, templating, quiet hours, locales, rollout, content models.
cue_notify.db SQLAlchemy models, portable types, Alembic migrations.
cue_notify.queue Durable job queue on the application database.
cue_notify.services Use cases: ingest, evaluate, compose, deliver, broadcast, engagement.
cue_notify.channels Provider contract, built-in providers, plugin registry.
cue_notify.ai Personalisation contract, guardrails and the Pydantic AI implementation.
cue_notify.api / cue_notify.cli Thin adapters over services.
cue_notify.runtime Long-lived resources (engine, HTTP client, channels, generator) shared by all entry points.

The life of an event

sequenceDiagram
    participant P as Your backend
    participant A as API
    participant D as Database
    participant W as Worker
    participant C as Channel
    P->>A: POST /v1/events
    A->>D: insert event + "process_event" job (1 tx)
    A-->>P: 202 Accepted
    W->>D: claim job
    W->>D: lock recipient, evaluate rules & policy,<br/>insert message + "deliver" job at send_at (1 tx)
    W->>D: claim "deliver" job when due
    W->>W: optional AI personalisation
    W->>C: send to each address
    W->>D: record outcome

Design decisions

The database is the queue. Jobs are rows written in the same transaction as the data they act on (transactional outbox), so there is no window where an event exists without its processing job or a message without its delivery. Delays and quiet hours are a future run_at, which survives restarts without a separate scheduler. Workers claim with SELECT … FOR UPDATE SKIP LOCKED and a conditional update that doubles as a lease; crashed workers' jobs reappear when the lease expires. This removes Redis and a broker from the deployment at the cost of throughput in the low thousands of jobs per second per database — far beyond what notification volumes need.

One composer. Rules, direct sends and broadcasts all build an intent and pass it to services.composer. Policy cannot be bypassed by adding a new way to send.

Policy is serialised per recipient. Composing locks the recipient row, so caps and cooldowns are exact under concurrency. On SQLite, BEGIN IMMEDIATE provides the same guarantee by serialising writers.

Decisions are data. Suppressed and failed messages are stored with reasons; events store per-rule outcomes. Debugging is a query, not a log search.

Render early, personalise late. Templates render when the message is composed (so errors surface immediately and content is auditable); AI personalisation runs at send time, after policy, so tokens are never spent on messages that will not be sent.

Idempotency everywhere. Events dedupe on (source, idempotency_key); messages on dedup_key (event:<id>:<rule>, broadcast:<id>:<recipient>, api:<key>:<key>); jobs on key. Every job handler re-checks state before acting, so at-least-once execution is safe.

UUIDv7 keys. Time-ordered identifiers keep indexes compact and double as stable pagination cursors.

Plugins over configuration flags. Channels are classes discovered through entry points, configured by name. Adding a provider never touches the core.