Files
Alex adb6963523 fix(rename): address review on the package rename
- CI installs the backend requirements from docsgpt/; the old cd into
  application/ silently installed nothing.
- The root .dockerignore re-admits only application/__init__.py. An upgraded
  checkout may still hold gitignored application/{inputs,indexes,vectors,.env}
  from the old layout, and the directory rule shipped them into the image.
- The compose files keep the host bind mounts on application/{indexes,inputs,
  vectors}, so an upgrade does not start with empty data. The move comes with
  the packaging work, together with an upgrade note.
- The alias loader puts the real docsgpt spec back on the shared module object
  after import (the import machinery stamped the alias spec on it, which made
  importlib.reload rename the module and skip re-execution) and delegates
  get_code/get_source/get_filename to the target loader, so
  python -m application.<name> runs.
- Each legacy application.* task name is registered as its own task object,
  a subclass carrying the old name. Registering the same object under two
  keys made Celery's tracer log every run under whichever name it built last.
- The redbeat key prefix stays redbeat:docsgpt:; the three schedule_syncs
  entries get stable names instead. redbeat tracks its static entries and
  deletes the ones that vanish from beat_schedule at start-up, and rewrites the
  task path of named entries in place, so neither a prefix bump nor a cleanup
  pass is needed (checked against redbeat 2.4.2 with a seeded Redis).
2026-09-07 12:02:07 +01:00

70 lines
2.9 KiB
Python

from kombu import Queue
from docsgpt.core.settings import settings
# Pydantic loads .env into ``settings`` but does not inject values into
# ``os.environ`` — read directly from settings so beat startup (which
# imports this module before any explicit env load) sees a real URL.
broker_url = settings.CELERY_BROKER_URL
result_backend = settings.CELERY_RESULT_BACKEND
task_serializer = 'json'
result_serializer = 'json'
accept_content = ['json']
# Autodiscover tasks
imports = (
'docsgpt.api.user.tasks',
'docsgpt.vectorstore.embeddings_tasks',
)
# Project-scoped queue so a stray sibling worker on the same broker
# (other repo, same default ``celery`` queue) can't grab DocsGPT tasks.
task_default_queue = "docsgpt"
task_default_exchange = "docsgpt"
task_default_routing_key = "docsgpt"
# Route document parsing to a dedicated queue so a parse enqueued from inside a
# Celery worker (headless/scheduled agent) is served by a separate parsing worker
# and never self-deadlocks the awaiting worker. The tool also passes the queue at
# apply_async time, so this routing is the default for any other enqueuer.
# Query embedding gets its own queue for the same reason parsing does: a query
# waiting behind a multi-minute ingest is a query that has timed out. A bare
# worker still consumes it, but its concurrency is shared -- run a separate
# ``-Q embeddings`` worker to actually isolate query latency from ingest.
task_routes = {
"docsgpt.api.user.tasks.parse_document": {"queue": settings.DOCUMENT_PARSE_QUEUE},
"docsgpt.vectorstore.embeddings_tasks.embed_texts": {"queue": settings.EMBEDDINGS_QUEUE},
}
# Declare every queue so a bare ``celery worker`` (no -Q) consumes ALL of them —
# the default worker does the whole job, parsing included. Operators who want
# heavy OCR isolated run one worker with ``-Q docsgpt`` and another with
# ``-Q parsing``. (dict.fromkeys dedupes if DOCUMENT_PARSE_QUEUE == "docsgpt".)
task_queues = tuple(
Queue(name)
for name in dict.fromkeys(
["docsgpt", settings.DOCUMENT_PARSE_QUEUE, settings.EMBEDDINGS_QUEUE]
)
)
beat_scheduler = "redbeat.RedBeatScheduler"
redbeat_redis_url = broker_url
redbeat_key_prefix = "redbeat:docsgpt:"
redbeat_lock_timeout = 90
# Survive worker SIGKILL/OOM without silently dropping in-flight tasks.
task_acks_late = True
task_reject_on_worker_lost = True
worker_prefetch_multiplier = settings.CELERY_WORKER_PREFETCH_MULTIPLIER
broker_transport_options = {"visibility_timeout": settings.CELERY_VISIBILITY_TIMEOUT}
result_expires = 86400 * 7
task_track_started = True
# Recycle the prefork worker child to bound native-heap growth from
# docling/torch parsing. Left unset (Celery's unlimited default) when 0.
if settings.CELERY_WORKER_MAX_MEMORY_PER_CHILD > 0:
worker_max_memory_per_child = settings.CELERY_WORKER_MAX_MEMORY_PER_CHILD
if settings.CELERY_WORKER_MAX_TASKS_PER_CHILD > 0:
worker_max_tasks_per_child = settings.CELERY_WORKER_MAX_TASKS_PER_CHILD