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
180 changes: 0 additions & 180 deletions apps/worker/app/core/gevent_worker_shutdown.py

This file was deleted.

91 changes: 42 additions & 49 deletions apps/worker/app/core/worker_bootstrap.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,13 @@
import subprocess
import sys

from celery.signals import worker_init, worker_shutdown, worker_shutting_down
from celery.worker import WorkController
from celery.signals import worker_init, worker_shutdown
from loguru import logger

from app.core.gevent_worker_shutdown import GeventWorkerShutdownController
from shared.core.celery_app import celery_app
from shared.core.logging import setup_logging
from shared.services.worker_health import start_worker_heartbeat, stop_worker_heartbeat

_worker_shutdown_controller: GeventWorkerShutdownController | None = None
_CHILD_PROCESS_TERM_TIMEOUT_SECONDS: float = 5
_CHILD_PROCESS_KILL_TIMEOUT_SECONDS: float = 5

Expand Down Expand Up @@ -44,28 +41,54 @@ def _stop_child_process(


@worker_init.connect
def init_worker(
sender: WorkController | None = None,
**kwargs: object,
) -> None:
def init_worker(**kwargs: object) -> None:
"""Initialize structured logging and sync Redis when worker process starts."""
global _worker_shutdown_controller

setup_logging(service_name="knowhere-worker")
start_worker_heartbeat()

if sender is None:
logger.error("Cannot configure bounded worker shutdown without worker sender")
else:
_worker_shutdown_controller = GeventWorkerShutdownController(
worker=sender,
timeout_seconds=float(celery_app.conf.worker_soft_shutdown_timeout),
)
# Do not cancel gevent tasks from Celery's reconnect or shutdown lifecycle.
# Staging proved that forced greenlet cancellation races with parser child
# processes and still misses the ECS stop deadline. Interrupted reservations
# instead recover through the independent Redis visibility watchdog.
try:
from celery.concurrency.gevent import TaskPool as GeventTaskPool

if not hasattr(GeventTaskPool, "_original_terminate_job"):

def _ignore_gevent_cancellation(
self: GeventTaskPool,
pid: int,
signal: int | None = None,
) -> None:
logger.warning(
f"gevent pool cannot kill greenlet (pid={pid}), "
"relying on visibility recovery and RedisJobLock for "
"redelivery deduplication"
)

setattr(
GeventTaskPool,
"_original_terminate_job",
getattr(
GeventTaskPool,
"terminate_job",
None,
),
)
GeventTaskPool.terminate_job = _ignore_gevent_cancellation
logger.info(
"Patched gevent TaskPool.terminate_job for visibility recovery"
)
except Exception as exc:
logger.warning(f"Could not patch gevent TaskPool: {exc}")

try:
from shared.services.redis.redis_sync_service import SyncRedisServiceFactory
from shared.services.redis.redis_sync_service import (
SyncRedisService,
SyncRedisServiceFactory,
)

service = SyncRedisServiceFactory.get_service()
service: SyncRedisService = SyncRedisServiceFactory.get_service()
if service.ping():
logger.info("Worker sync Redis connection verified")
else:
Expand All @@ -74,39 +97,9 @@ def init_worker(
logger.warning(f"Worker sync Redis init deferred: {exc}")


@worker_shutting_down.connect
def begin_bounded_worker_shutdown(
sender: str | None = None,
sig: str | None = None,
how: str | None = None,
**kwargs: object,
) -> None:
"""Bound ECS SIGTERM without entering Celery's gevent-unsafe cold path."""
if sig != "SIGTERM" or how != "Warm":
return

shutdown_controller: GeventWorkerShutdownController | None = (
_worker_shutdown_controller
)
if shutdown_controller is None:
logger.error("Cannot schedule bounded worker shutdown before worker init")
return

shutdown_controller.schedule()


@worker_shutdown.connect
def shutdown_worker(**kwargs: object) -> None:
"""Clean up shared resources on worker shutdown."""
global _worker_shutdown_controller

shutdown_controller: GeventWorkerShutdownController | None = (
_worker_shutdown_controller
)
_worker_shutdown_controller = None
if shutdown_controller is not None:
shutdown_controller.close()

try:
stop_worker_heartbeat()
logger.info("Worker heartbeat stopped")
Expand Down
Loading
Loading