109 lines
3.8 KiB
Python
109 lines
3.8 KiB
Python
import tempfile
|
|
|
|
from celery import Celery
|
|
from kombu import Queue
|
|
from werkzeug.local import LocalProxy
|
|
from redbeat import RedBeatScheduler
|
|
|
|
celery_app = Celery()
|
|
|
|
|
|
def _create_ssl_cert_file(cert_data: str) -> str:
|
|
"""Create temporary certificate file for Celery SSL"""
|
|
if not cert_data:
|
|
return None
|
|
|
|
with tempfile.NamedTemporaryFile(mode='w', delete=False, suffix='.pem') as cert_file:
|
|
cert_file.write(cert_data)
|
|
return cert_file.name
|
|
|
|
|
|
def init_celery(celery, app, is_beat=False):
|
|
celery_app.main = app.name
|
|
|
|
celery_config = {
|
|
'broker_url': app.config.get('CELERY_BROKER_URL', 'redis://localhost:6379/0'),
|
|
'result_backend': app.config.get('CELERY_RESULT_BACKEND', 'redis://localhost:6379/0'),
|
|
'task_serializer': app.config.get('CELERY_TASK_SERIALIZER', 'json'),
|
|
'result_serializer': app.config.get('CELERY_RESULT_SERIALIZER', 'json'),
|
|
'accept_content': app.config.get('CELERY_ACCEPT_CONTENT', ['json']),
|
|
'timezone': app.config.get('CELERY_TIMEZONE', 'UTC'),
|
|
'enable_utc': app.config.get('CELERY_ENABLE_UTC', True),
|
|
}
|
|
|
|
# Add broker transport options for SSL and connection pooling
|
|
broker_transport_options = {
|
|
'master_name': None,
|
|
'max_connections': 20,
|
|
'retry_on_timeout': True,
|
|
'socket_connect_timeout': 5,
|
|
'socket_timeout': 5,
|
|
}
|
|
|
|
cert_data = app.config.get('REDIS_CERT_DATA')
|
|
if cert_data:
|
|
try:
|
|
ssl_cert_file = _create_ssl_cert_file(cert_data)
|
|
if ssl_cert_file:
|
|
broker_transport_options.update({
|
|
'ssl_cert_reqs': 'required',
|
|
'ssl_ca_certs': ssl_cert_file,
|
|
'ssl_check_hostname': True,
|
|
})
|
|
app.logger.info("SSL configured for Celery Redis connection")
|
|
except Exception as e:
|
|
app.logger.error(f"Failed to configure SSL for Celery: {e}")
|
|
|
|
celery_config['broker_transport_options'] = broker_transport_options
|
|
celery_config['result_backend_transport_options'] = broker_transport_options
|
|
|
|
if is_beat:
|
|
# Add configurations specific to Beat scheduler
|
|
celery_config['beat_scheduler'] = 'redbeat.RedBeatScheduler'
|
|
celery_config['redbeat_lock_key'] = 'redbeat::lock'
|
|
celery_config['beat_max_loop_interval'] = 10 # Adjust as needed
|
|
|
|
celery_app.conf.update(**celery_config)
|
|
|
|
# Task queues for workers only
|
|
if not is_beat:
|
|
celery_app.conf.task_queues = (
|
|
Queue('default', routing_key='task.#'),
|
|
Queue('embeddings', routing_key='embeddings.#', queue_arguments={'x-max-priority': 10}),
|
|
Queue('llm_interactions', routing_key='llm_interactions.#', queue_arguments={'x-max-priority': 5}),
|
|
Queue('entitlements', routing_key='entitlements.#', queue_arguments={'x-max-priority': 10}),
|
|
)
|
|
celery_app.conf.task_routes = {
|
|
'eveai_workers.*': { # All tasks from eveai_workers module
|
|
'queue': 'embeddings',
|
|
'routing_key': 'embeddings.#',
|
|
},
|
|
'eveai_chat_workers.*': { # All tasks from eveai_chat_workers module
|
|
'queue': 'llm_interactions',
|
|
'routing_key': 'llm_interactions.#',
|
|
},
|
|
'eveai_entitlements.*': { # All tasks from eveai_entitlements module
|
|
'queue': 'entitlements',
|
|
'routing_key': 'entitlements.#',
|
|
}
|
|
}
|
|
|
|
# Ensure tasks execute with Flask context
|
|
class ContextTask(celery.Task):
|
|
def __call__(self, *args, **kwargs):
|
|
with app.app_context():
|
|
return self.run(*args, **kwargs)
|
|
|
|
celery.Task = ContextTask
|
|
|
|
|
|
def make_celery(app_name, config):
|
|
return celery_app
|
|
|
|
|
|
def _get_current_celery():
|
|
return celery_app
|
|
|
|
|
|
current_celery = LocalProxy(_get_current_celery)
|