Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -77,3 +77,21 @@ MINIO_CONSOLE_PORT=9001
STORAGE_BACKEND=s3
LOCAL_STORAGE_ROOT=./data
LOCAL_STORAGE_HOST_PATH=./local_storage_data

# Operational services (see doc 30)
# Staging→canonical sync: `mc mirror` runs on the worker (Celery task + Beat cron).
# STAGING_LOCATION accepts a local path or an s3:// URI.
STAGING_LOCATION=./data/staging
STAGING_SYNC_INTERVAL_SECONDS=300
STAGING_SYNC_SOFT_TIME_LIMIT_SECONDS=1800
# Per-task time-limit contract (doc 30): every Celery task must declare its own
# ceiling explicitly (enforced at class-definition time, see base_task.py).
INGEST_BATCH_SOFT_TIME_LIMIT_SECONDS=7200
MOSAIC_DRIZZLE_SOFT_TIME_LIMIT_SECONDS=3600
CLI_TOOL_SOFT_TIME_LIMIT_SECONDS=3600
DB_BACKUP_SOFT_TIME_LIMIT_SECONDS=3600
DLQ_DUMP_SOFT_TIME_LIMIT_SECONDS=30
RECONCILE_STUCK_JOBS_SOFT_TIME_LIMIT_SECONDS=120
ENABLE_STAGING_SYNC_CRON=true
# Stuck-job watchdog: seconds an IN_PROCESS job may age before it is failed.
JOB_STALENESS_TIMEOUT_SECONDS=3600
23 changes: 23 additions & 0 deletions docker-compose.prod.yml
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,29 @@ services:
networks:
- diffpype_net

beat:
image: ghcr.io/davecoulter/diffpype_claude-worker:${IMAGE_TAG:-main}
# Single database-backed Celery Beat process (see docker-compose.yml for the
# rationale). Consumes no task queues; publishes scheduled tasks to Redis.
command: celery -A src.worker.celery_app beat -S sqlalchemy_celery_beat.schedulers:DatabaseScheduler --loglevel=info
environment:
DATABASE_URL: postgresql+psycopg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
REDIS_URL: redis://redis:6379/0
LOG_LEVEL: ${LOG_LEVEL:-INFO}
OTEL_EXPORTER_OTLP_ENDPOINT: ${OTEL_EXPORTER_OTLP_ENDPOINT}
OTEL_SERVICE_NAME: diffpype-beat
STAGING_SYNC_INTERVAL_SECONDS: ${STAGING_SYNC_INTERVAL_SECONDS:-300}
ENABLE_STAGING_SYNC_CRON: ${ENABLE_STAGING_SYNC_CRON:-true}
JOB_STALENESS_TIMEOUT_SECONDS: ${JOB_STALENESS_TIMEOUT_SECONDS:-3600}
depends_on:
db:
condition: service_healthy
redis:
condition: service_healthy
restart: unless-stopped
networks:
- diffpype_net

networks:
diffpype_net:
name: diffpype_net
Expand Down
29 changes: 29 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,35 @@ services:
networks:
- diffpype_net

beat:
build:
context: .
dockerfile: docker/worker.Dockerfile
# Dedicated Celery Beat process (exactly one, ever) using the database-backed
# scheduler so schedules edited at runtime via SQLAdmin take effect without a
# restart. Not a --beat flag on a worker: that would double-fire schedules if
# the worker ever scaled past one replica. Requires the scheduler tables to
# exist (run `alembic upgrade head` first); it consumes no task queues.
command: celery -A src.worker.celery_app beat -S sqlalchemy_celery_beat.schedulers:DatabaseScheduler --loglevel=info
environment:
DATABASE_URL: postgresql+psycopg://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB}
REDIS_URL: redis://redis:6379/0
LOG_LEVEL: ${LOG_LEVEL:-INFO}
OTEL_EXPORTER_OTLP_ENDPOINT: ${OTEL_EXPORTER_OTLP_ENDPOINT:-http://jaeger:4317}
OTEL_SERVICE_NAME: diffpype-beat
STAGING_SYNC_INTERVAL_SECONDS: ${STAGING_SYNC_INTERVAL_SECONDS:-300}
ENABLE_STAGING_SYNC_CRON: ${ENABLE_STAGING_SYNC_CRON:-true}
JOB_STALENESS_TIMEOUT_SECONDS: ${JOB_STALENESS_TIMEOUT_SECONDS:-3600}
volumes:
- ./src:/app/src
depends_on:
db:
condition: service_healthy
redis:
condition: service_healthy
networks:
- diffpype_net

flower:
image: mher/flower:2.0.1
environment:
Expand Down
11 changes: 11 additions & 0 deletions docker/worker.Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,17 @@ FROM python:3.12-slim

COPY --from=ghcr.io/astral-sh/uv:latest /uv /bin/uv

# MinIO client (`mc`), used by run_staging_sync's `mc mirror` staging->canonical
# sync (doc 30 §1). Arch-aware so the image builds on both CI's linux/amd64 and
# an Apple-Silicon linux/arm64 build. Only the worker image gets `mc` — the api
# image never runs the sync (it is dispatched to the worker).
RUN apt-get update \
&& apt-get install -y --no-install-recommends curl ca-certificates \
&& curl -fsSL "https://dl.min.io/client/mc/release/linux-$(dpkg --print-architecture)/mc" \
-o /usr/local/bin/mc \
&& chmod +x /usr/local/bin/mc \
&& rm -rf /var/lib/apt/lists/*

WORKDIR /app

# Layer 1: install deps only (cached unless pyproject.toml/uv.lock change).
Expand Down
158 changes: 158 additions & 0 deletions docs/architecture/30_operational_services_and_watchdog.md

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions docs/architecture/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,4 +36,5 @@ work for that stage.
27_schema_storage_spatial_types
28_domain_graph_population
29_psycopg3_healpix_dummy_cleanup
30_operational_services_and_watchdog
```
43 changes: 39 additions & 4 deletions docs/cli_guide.rst
Original file line number Diff line number Diff line change
Expand Up @@ -29,10 +29,18 @@ Overview
.. note::

This guide covers the foundational database-management commands. The domain
commands added in later stages (``create-project``, ``ingest``,
``tessellate-tiles``, ``create-mosaic``, ``populate-demo-project``, and their
``*-status`` pollers) share the same Service Layer and follow the same
API/CLI-parity contract; run ``diffpype-manage --help`` for the full list.
and operational commands added in later stages (``create-project``,
``ingest``, ``tessellate-tiles``, ``create-mosaic``, ``sync-staging``,
``reconcile-stuck-jobs``, ``populate-demo-project``, and their ``*-status``
pollers) share the same Service Layer and follow the same API/CLI-parity
contract; run ``diffpype-manage --help`` for the full list.

``tessellate-tiles``/``create-tiles`` take a ``--region-source``
(``cone`` | ``project_footprint`` | ``bounding_box``) with the fields that
mode needs (e.g. ``--ra/--decl/--radius-deg`` for ``cone``,
``--min-ra/--max-ra/--min-decl/--max-decl`` for ``bounding_box``), plus
``--overlap-only/--no-overlap-only`` to trim the grid to the region or
materialize it fully.

``seed-db``
-----------
Expand Down Expand Up @@ -62,3 +70,30 @@ usable. Intended for local development only.
Schema reset complete. Auto-seeding foundational records...
Seeding database: inserting foundational sysadmin + reference records...
Done.

``sync-staging``
----------------

Dispatches a staging→canonical storage sync to the worker (which runs
``mc mirror`` in a streamed, restart-safe Celery task). ``--staging-prefix``
accepts a local path or an ``s3://`` URI; ``--canonical-prefix`` defaults to the
bucket root.

.. code-block:: console

$ docker compose run --rm api diffpype-manage sync-staging --staging-prefix ./data/staging --canonical-prefix raw
Dispatched staging sync. job_id=<celery-task-id>

``reconcile-stuck-jobs``
------------------------

Fails any job left in ``IN_PROCESS`` past the staleness threshold (an
uncatchable worker crash or OOM kill can't run a task's own failure handler).
``--threshold-seconds`` overrides ``JOB_STALENESS_TIMEOUT_SECONDS`` for this
sweep; it also runs automatically on a Celery Beat schedule.

.. code-block:: console

$ docker compose run --rm api diffpype-manage reconcile-stuck-jobs --threshold-seconds 3600
Reconciled 1 stuck job(s).
IngestBatch id=9 (age=7200s) -> FAILED
1 change: 1 addition & 0 deletions docs/conf.py
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@
"slugify",
"sqladmin",
"sqlalchemy",
"sqlalchemy_celery_beat",
"src.api.admin",
"starlette",
"starlette_exporter",
Expand Down
14 changes: 12 additions & 2 deletions docs/diagrams/infrastructure_topology.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,10 @@ flowchart TB
w_app --> w_core
end

subgraph beat_c["beat (Celery Beat)"]
beat_app["Celery Beat<br/>DatabaseScheduler"]:::workerLayer
end

subgraph ui_c["ui (React / Vite)"]
direction TB
ui_view["Dashboard<br/>DashboardPage.tsx"]:::uiLayer
Expand Down Expand Up @@ -112,6 +116,11 @@ flowchart TB
w_tasks -->|store/fetch FITS| minio_node
portainer_node -.-> minio_c

%% Beat scheduler (declared last so existing link indices stay valid)
beat_app -->|read/write schedules| db
beat_app -->|publish scheduled tasks| redis
portainer_node -.-> beat_c

classDef apiLayer fill:#3B6EA5,stroke:#1F4066,color:#fff
classDef workerLayer fill:#C97A3D,stroke:#8A4F24,color:#fff
classDef repo fill:#4C8C6B,stroke:#2E5842,color:#fff
Expand All @@ -123,15 +132,15 @@ flowchart TB
classDef spacer fill:transparent,stroke:transparent,color:transparent
classDef networkBg fill:transparent,stroke:#5A6C77,color:#fff

class db_c,redis_c,minio_c,api_c,worker_c,ui_c,jaeger_c,flower_c,portainer_c containerBg
class db_c,redis_c,minio_c,api_c,worker_c,beat_c,ui_c,jaeger_c,flower_c,portainer_c containerBg
class network networkBg

%% Tier 1 (default): solid amber — configured application data flow between components
linkStyle default stroke:#D9A64A,stroke-width:4px
%% Invisible spacer links used only to center nodes within db_c, worker_c, portainer_c — must override the amber default or they render as dangling solid lines
linkStyle 0,1,11,17,18 stroke:none,stroke-width:0
%% Tier 3: dashed cool steel gray — Portainer manages containers at the Docker daemon level, independent of what's running inside them
linkStyle 25,26,27,28,29,30,31,38 stroke:#8F97A0,stroke-width:4px
linkStyle 25,26,27,28,29,30,31,38,41 stroke:#8F97A0,stroke-width:4px
%% Tier 2: dashed rose — Jaeger/Flower/DBeaver observe a specific process's exposed port/protocol (OTLP, Celery/Redis state, Postgres wire protocol)
linkStyle 32,33,34,35 stroke:#C0546A,stroke-width:4px
```
Expand Down Expand Up @@ -167,6 +176,7 @@ flowchart TB
| `minio` | Local S3-compatible object store (mock) for FITS payloads; the api and workers read/write files here via `S3StorageService` |
| `api` | Serves the HTTP API, admin panel, and CLI; validates input and dispatches work |
| `worker` (×2: `worker_light`, `worker_heavy`) | Same image/codebase, deployed as two instances with different queue subscriptions and resource limits — `light` for fast I/O-bound tasks, `heavy_memory` for memory/compute-intensive ones |
| `beat` | Single Celery Beat process (same worker image) using the database-backed scheduler; reads/writes runtime-editable schedules in Postgres (`celery_schema`) and publishes due tasks to Redis. Consumes no task queues |
| `ui` | Frontend for dispatching jobs and viewing status |
| `jaeger` | Collects and visualizes distributed traces via OTLP |
| `flower` | Real-time Celery task/worker monitoring dashboard |
Expand Down
8 changes: 8 additions & 0 deletions docs/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,14 @@ API (FastAPI)
:members:
:undoc-members:

.. automodule:: src.api.routes.storage
:members:
:undoc-members:

.. automodule:: src.api.routes.jobs
:members:
:undoc-members:

.. automodule:: src.api.routes.tiles
:members:
:undoc-members:
Expand Down
11 changes: 10 additions & 1 deletion migrations/env.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,12 @@

target_metadata = Base.metadata

# Note: sqlalchemy-celery-beat's tables (doc 30 §4) live in a dedicated
# `celery_schema` Postgres schema, created by migration 0014. Autogenerate runs
# with the default include_schemas=False, so it only ever reflects the public
# schema and never sees — or tries to drop — those package-owned tables. No
# include_object filter is needed.


def get_url() -> str:
# Allow callers (e.g. the integration-test fixture) to override the URL
Expand All @@ -38,7 +44,10 @@ def run_migrations_online() -> None:
cfg["sqlalchemy.url"] = get_url()
connectable = engine_from_config(cfg, prefix="sqlalchemy.", poolclass=pool.NullPool)
with connectable.connect() as connection:
context.configure(connection=connection, target_metadata=target_metadata)
context.configure(
connection=connection,
target_metadata=target_metadata,
)
with context.begin_transaction():
context.run_migrations()

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
"""add JobConfiguration.task_name provenance column

Revision ID: 0013
Revises: 0012
Create Date: 2026-07-29

Doc 30 §3. Records which Celery task (or CLI tool) a JobConfiguration row's
kwargs/command correspond to, e.g. "src.worker.tasks.run_ingest_batch". Nullable,
so existing provenance rows need no backfill.
"""

import sqlalchemy as sa
from alembic import op

revision = "0013"
down_revision = "0012"
branch_labels = None
depends_on = None


def upgrade() -> None:
op.add_column(
"job_configurations",
sa.Column("task_name", sa.String(), nullable=True),
)


def downgrade() -> None:
op.drop_column("job_configurations", "task_name")
38 changes: 38 additions & 0 deletions migrations/versions/20260729_0014_add_beat_scheduler_tables.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
"""add sqlalchemy-celery-beat scheduler schema and tables

Revision ID: 0014
Revises: 0013
Create Date: 2026-07-29

Doc 30 §4. Creates the dedicated ``celery_schema`` Postgres schema and the
relational tables `sqlalchemy-celery-beat` uses for its database-backed Celery
Beat scheduler (DatabaseScheduler), so schedules can be created/edited/paused at
runtime via SQLAdmin.

Table (and Enum type) creation is delegated to the package's own
``ModelBase.metadata`` rather than transcribing six tables + two schema-qualified
enums by hand: that keeps this migration faithful to the package's real schema and
lets a future package upgrade re-run the same delegation instead of drifting from
a hand-copied definition. The tables live in their own schema, so autogenerate
(which runs with the default ``include_schemas=False``) never sees them and this
migration is the sole authority over their lifecycle. Downgrade drops the whole
schema (tables + enum types) in one cascade.
"""

from alembic import op

revision = "0014"
down_revision = "0013"
branch_labels = None
depends_on = None


def upgrade() -> None:
from sqlalchemy_celery_beat.models import ModelBase as BeatBase

op.execute("CREATE SCHEMA IF NOT EXISTS celery_schema")
BeatBase.metadata.create_all(op.get_bind())


def downgrade() -> None:
op.execute("DROP SCHEMA IF EXISTS celery_schema CASCADE")
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ dependencies = [
"scipy==1.18.0",
"scikit-learn==1.9.0",
"pandas==3.0.5",
"sqlalchemy-celery-beat==0.8.4",
]

[project.scripts]
Expand Down
Loading
Loading