Watch
1
0
Fork
You've already forked Seyyed_arc
0
forked from hesabix/arc
Seyyed_arc/hesabixAPI/app/core/queue.py

349 lines
13 KiB
Python
Executable file
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""
Queue Service با استفاده از RQ (Redis Queue)
برای مدیریت background jobs
"""
from __future__ import annotations
import json
import logging
from typing import Any, Dict, Optional, Callable
from datetime import timedelta
import rq
from rq import Queue, Worker
from rq.job import Job
from rq.registry import FailedJobRegistry
from redis import Redis
from app.core.cache import get_redis_client
from app.core.settings import get_settings
logger = logging.getLogger(__name__)
# Queue names
QUEUE_DEFAULT = "default"
QUEUE_HIGH_PRIORITY = "high"
QUEUE_LOW_PRIORITY = "low"
QUEUE_EMAIL = "email"
QUEUE_REPORTS = "reports"
QUEUE_EXPORTS = "exports"
QUEUE_TAX = "tax" # Queue مخصوص عملیات مالیاتی
def load_redis_settings_from_configuration() -> tuple[bool, str, int, int, Optional[str]]:
"""
خواندن تنظیمات Redis از دیتابیس (تنظیمات کل سیستم) یا متغیرهای محیطی.
اگر مدیر Redis را غیرفعال کرده باشد، enabled=False برمی‌گردد.
"""
settings = get_settings()
redis_enabled = getattr(settings, "redis_enabled", False)
redis_host = getattr(settings, "redis_host", "localhost")
redis_port = getattr(settings, "redis_port", 6379)
redis_db = getattr(settings, "redis_db", 0)
redis_password = getattr(settings, "redis_password", None)
try:
from adapters.db.session import get_db
from app.services.system_settings_service import get_redis_configuration
db_gen = get_db()
db = next(db_gen)
try:
redis_config = get_redis_configuration(db)
redis_enabled = redis_config.get("enabled", False)
redis_host = redis_config.get("host", "localhost")
redis_port = redis_config.get("port", 6379)
redis_db = redis_config.get("db", 0)
redis_password = redis_config.get("password")
finally:
db.close()
except Exception:
pass
return bool(redis_enabled), str(redis_host), int(redis_port), int(redis_db), redis_password
def is_redis_enabled_in_configuration() -> bool:
"""آیا Redis در تنظیمات سیستم روشن است (بدون تست اتصال شبکه)."""
return load_redis_settings_from_configuration()[0]
def get_redis_connection() -> Optional[Redis]:
"""دریافت Redis connection برای RQ"""
client = get_redis_client()
if client is None:
return None
redis_enabled, redis_host, redis_port, redis_db, redis_password = load_redis_settings_from_configuration()
if not redis_enabled:
return None
try:
# RQ نیاز به Redis بدون decode_responses دارد
redis_conn = Redis(
host=redis_host,
port=redis_port,
db=redis_db,
password=redis_password,
decode_responses=False, # RQ نیاز به bytes دارد
socket_connect_timeout=2,
socket_timeout=2,
retry_on_timeout=True,
health_check_interval=30,
)
redis_conn.ping()
return redis_conn
except Exception as e:
logger.warning(f"Failed to create Redis connection for RQ: {e}")
return None
def get_queue(name: str = QUEUE_DEFAULT) -> Optional[Queue]:
"""دریافت queue با نام مشخص"""
redis_conn = get_redis_connection()
if redis_conn is None:
return None
return Queue(name, connection=redis_conn)
class QueueService:
"""سرویس مدیریت Queue"""
def __init__(self):
self.redis_conn = get_redis_connection()
self.enabled = self.redis_conn is not None
def enqueue(
self,
func: Callable,
*args,
queue_name: str = QUEUE_DEFAULT,
timeout: int = 300,
result_ttl: int = 3600,
job_id: Optional[str] = None,
**kwargs
) -> Optional[Job]:
"""
اضافه کردن job به queue
Args:
func: تابعی که باید اجرا شود
*args: آرگومان‌های positional
queue_name: نام queue
timeout: حداکثر زمان اجرا (ثانیه)
result_ttl: زمان نگهداری نتیجه (ثانیه)
job_id: شناسه job (اختیاری)
**kwargs: آرگومان‌های keyword
Returns:
Job object یا None در صورت خطا
"""
if not self.enabled:
logger.warning("Queue service is disabled (Redis not available)")
return None
try:
queue = get_queue(queue_name)
if queue is None:
return None
job = queue.enqueue(
func,
*args,
job_id=job_id,
timeout=timeout,
result_ttl=result_ttl,
**kwargs
)
logger.info(f"Job {job.id} enqueued to {queue_name} queue")
return job
except Exception as e:
logger.error(f"Error enqueueing job: {e}", exc_info=True)
return None
def get_job(self, job_id: str) -> Optional[Job]:
"""دریافت job با شناسه"""
if not self.enabled:
logger.debug(f"Queue service disabled, cannot fetch job {job_id}")
return None
try:
job = Job.fetch(job_id, connection=self.redis_conn)
logger.debug(f"Job {job_id} found in Redis")
return job
except rq.exceptions.NoSuchJobError:
# Job وجود ندارد - این یک وضعیت عادی است
logger.debug(f"Job {job_id} not found in Redis (NoSuchJobError)")
return None
except Exception as e:
logger.warning(f"Error fetching job {job_id}: {e}", exc_info=True)
return None
def get_job_status(self, job_id: str) -> Optional[Dict[str, Any]]:
"""دریافت وضعیت job به صورت dictionary"""
logger.info(f"[QUEUE_SERVICE] get_job_status called for job_id: {job_id}")
job = self.get_job(job_id)
if job is None:
logger.info(f"[QUEUE_SERVICE] Job {job_id} not found in main queue, checking failed registry")
# اگر job در queue اصلی پیدا نشد، در failed registry بررسی کن
failed_status = self._get_failed_job_status(job_id)
if failed_status:
logger.info(f"[QUEUE_SERVICE] Job {job_id} found in failed registry")
else:
logger.info(f"[QUEUE_SERVICE] Job {job_id} not found in failed registry either")
return failed_status
try:
status = {
"id": job.id,
"state": job.get_status(),
"created_at": job.created_at.isoformat() if job.created_at else None,
"started_at": job.started_at.isoformat() if job.started_at else None,
"ended_at": job.ended_at.isoformat() if job.ended_at else None,
"result": None,
"error": None,
"updated_at": (job.ended_at or job.started_at or job.created_at).isoformat() if (job.ended_at or job.started_at or job.created_at) else None,
}
# اضافه کردن نتیجه یا خطا
if job.is_finished:
try:
status["result"] = job.result
except Exception:
pass
elif job.is_failed:
try:
if job.exc_info:
status["error"] = str(job.exc_info)
else:
# اگر exc_info نباشد، سعی می‌کنیم از exception message استفاده کنیم
try:
status["error"] = str(job.exc_info) if hasattr(job, 'exc_info') and job.exc_info else "Job failed without error details"
except Exception:
status["error"] = "Job failed without error details"
except Exception as e:
logger.warning(f"Error getting job error info for {job_id}: {e}")
status["error"] = "Job failed (error details unavailable)"
# اضافه کردن metadata اگر وجود دارد
if job.meta:
status["meta"] = job.meta
return status
except Exception as e:
logger.error(f"Error getting job status {job_id}: {e}", exc_info=True)
return None
def cancel_job(self, job_id: str) -> bool:
"""لغو job"""
job = self.get_job(job_id)
if job is None:
return False
try:
job.cancel()
logger.info(f"Job {job_id} cancelled")
return True
except Exception as e:
logger.error(f"Error cancelling job {job_id}: {e}", exc_info=True)
return False
def delete_job(self, job_id: str) -> bool:
"""حذف job"""
job = self.get_job(job_id)
if job is None:
return False
try:
job.delete()
logger.info(f"Job {job_id} deleted")
return True
except Exception as e:
logger.error(f"Error deleting job {job_id}: {e}", exc_info=True)
return False
def get_queue_length(self, queue_name: str = QUEUE_DEFAULT) -> int:
"""دریافت تعداد jobs در queue"""
queue = get_queue(queue_name)
if queue is None:
return 0
return len(queue)
def _get_failed_job_status(self, job_id: str) -> Optional[Dict[str, Any]]:
"""بررسی job در failed registry"""
if not self.enabled:
return None
try:
# بررسی در تمام queue ها
for queue_name in [QUEUE_DEFAULT, QUEUE_REPORTS, QUEUE_HIGH_PRIORITY,
QUEUE_LOW_PRIORITY, QUEUE_EMAIL, QUEUE_EXPORTS, QUEUE_TAX]:
try:
queue = Queue(queue_name, connection=self.redis_conn)
failed_registry = FailedJobRegistry(queue=queue)
if job_id in failed_registry:
# job در failed registry است
try:
job = Job.fetch(job_id, connection=self.redis_conn)
if job:
error_msg = "Job failed"
try:
if job.exc_info:
error_msg = str(job.exc_info)
elif hasattr(job, 'exc_info') and job.exc_info:
error_msg = str(job.exc_info)
except Exception:
error_msg = "Job failed (error details unavailable)"
return {
"id": job.id,
"state": "failed",
"created_at": job.created_at.isoformat() if job.created_at else None,
"started_at": job.started_at.isoformat() if job.started_at else None,
"ended_at": job.ended_at.isoformat() if job.ended_at else None,
"result": None,
"error": error_msg,
"updated_at": (job.ended_at or job.started_at or job.created_at).isoformat() if (job.ended_at or job.started_at or job.created_at) else None,
"meta": job.meta if job.meta else None,
}
except rq.exceptions.NoSuchJobError:
# Job در failed registry است اما از Redis حذف شده
logger.debug(f"Job {job_id} is in failed registry but not in Redis")
continue
except Exception as e:
logger.debug(f"Error checking queue {queue_name} for failed job {job_id}: {e}")
continue
except Exception as e:
logger.debug(f"Error checking failed registry for job {job_id}: {e}")
return None
def get_failed_jobs(self, limit: int = 10) -> list[Dict[str, Any]]:
"""دریافت لیست jobs ناموفق"""
if not self.enabled:
return []
try:
failed_registry = FailedJobRegistry(queue=Queue(connection=self.redis_conn))
job_ids = failed_registry.get_job_ids(0, limit - 1)
jobs = []
for job_id in job_ids:
status = self.get_job_status(job_id)
if status:
jobs.append(status)
return jobs
except Exception as e:
logger.error(f"Error getting failed jobs: {e}", exc_info=True)
return []
_queue_service: Optional[QueueService] = None
def get_queue_service() -> QueueService:
"""دریافت instance از QueueService"""
global _queue_service
if _queue_service is None:
_queue_service = QueueService()
return _queue_service