diff --git a/backend/.env.example b/backend/.env.example index c54d4b9..54f01dc 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -39,3 +39,12 @@ OPENAI_CONNECT_RETRIES=3 OPENAI_BASE_URL= OPENAI_ORGANIZATION= OPENAI_PROJECT= + +REDIS_URL=redis://localhost:6379/0 +TASKIQ_QUEUE_NAME=inbox +TASKIQ_MAX_RETRIES=3 +TASKIQ_RETRY_DELAY=5 +TASKIQ_MAX_DELAY=120 +TASKIQ_DLQ_STREAM=taskiq:dlq +TASKIQ_IDLE_TIMEOUT_MS=600000 +APP_VERSION=dev diff --git a/backend/Dockerfile b/backend/Dockerfile new file mode 100644 index 0000000..e462a2e --- /dev/null +++ b/backend/Dockerfile @@ -0,0 +1,12 @@ +FROM python:3.12-slim + +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY . . + +# Runs the Taskiq worker against taskiq_management.broker_setup. +# docker-compose overrides this command if needed. +CMD ["taskiq", "worker", "taskiq_management.broker_setup:broker", "inbox.tasks", "taskiq_management.tasks"] diff --git a/backend/inbox/app.py b/backend/inbox/app.py index 2670808..ae0906f 100644 --- a/backend/inbox/app.py +++ b/backend/inbox/app.py @@ -1,4 +1,4 @@ -from fastapi import APIRouter,BackgroundTasks,Depends, Query +from fastapi import APIRouter,Depends, Query from fastapi.responses import JSONResponse from fastapi import HTTPException from db_setup import get_session @@ -12,7 +12,6 @@ router = APIRouter() @router.get("/email/fetch") async def fetch_email( - background_tasks: BackgroundTasks, top:int=Query(100), skip:int=Query(0,ge=0), token=Query(...), @@ -31,7 +30,7 @@ async def fetch_email( items_lst.append({"message_id":message_id,"email_contents":service_per_email}) if service.pending_match_ids: - background_tasks.add_task(Email.run_inbox_matching, list(service.pending_match_ids)) + await service.enqueue_matching(list(service.pending_match_ids),force=False) return JSONResponse(content={"data":items_lst,"status_code":200}) @@ -67,15 +66,14 @@ async def fetch_inbox( @router.post("/inbox/{record_id}/match") async def rematch_inbox( record_id: str, - background_tasks: BackgroundTasks, current_user: dict = Depends(require_permission(PermissionTag.INBOX_EDIT)), session: AsyncSession = Depends(get_session), ): try: service=Email(session=session) queued_id=await service.queue_rematch(record_id) - background_tasks.add_task(Email.run_inbox_matching, [queued_id], force=True) - return JSONResponse(content={"data":{"queued":True},"status_code":200}) + task_ids=await service.enqueue_matching([queued_id], force=True) + return JSONResponse(content={"data":{"queued":True,"task_ids":task_ids},"status_code":200}) except HTTPException: raise except Exception as e: diff --git a/backend/inbox/plugins.py b/backend/inbox/plugins.py index aaac986..3c31eb4 100644 --- a/backend/inbox/plugins.py +++ b/backend/inbox/plugins.py @@ -1,4 +1,4 @@ -"""Inbox helpers — attachment loading and other non-routing checks.""" +"""Inbox helpers — attachment loading and resume text extraction.""" from __future__ import annotations @@ -8,55 +8,58 @@ from pathlib import Path from inbox.models import Inbox_Messages from job.candidate.views import FileRead +_ATTACHMENTS_DIR=Path(__file__).resolve().parent/"decoded_attachments" -def load_message_files(message: Inbox_Messages) -> list[dict]: - """Read files from file_path when they exist on disk.""" + +def resolve_attachment_path(path_str:str) -> Path: + """Prefer stored path; fall back to basename under decoded_attachments.""" + path=Path(path_str.strip()) + if path.is_file(): + return path + fallback=_ATTACHMENTS_DIR/path.name + if fallback.is_file(): + return fallback + return path + + +def load_message_files(message:Inbox_Messages) -> list[dict]: if not message.file_path: return [] - - files: list[dict] = [] + files=[] for path_str in message.file_path.split(","): - path = Path(path_str.strip()) + path=resolve_attachment_path(path_str) if not path.is_file(): continue try: - raw = path.read_bytes() + raw=path.read_bytes() except OSError: continue - files.append( - { - "file_name": path.name, - "content_base64": base64.b64encode(raw).decode("ascii"), - "size": len(raw), - } - ) + files.append({ + "file_name":path.name, + "content_base64":base64.b64encode(raw).decode("ascii"), + "size":len(raw), + }) return files -async def extract_resume_text(file_paths: list[str]) -> tuple[str, str]: - """Extract text from the PDFs among file_paths. Returns (combined_text, error).""" - pdf_paths = [ - Path(p.strip()) - for p in (file_paths or []) - if p and p.strip() and Path(p.strip()).suffix.lower() == ".pdf" - ] - existing = [p for p in pdf_paths if p.is_file()] +async def extract_resume_text(file_paths:list[str]) -> tuple[str,str]: + candidates=[resolve_attachment_path(p) for p in (file_paths or []) if p and p.strip()] + existing=[p for p in candidates if p.is_file() and p.suffix.lower()==".pdf"] if not existing: - return "", "no PDF attachment to extract (.doc/.docx not supported)" + return "","no PDF attachment to extract (.doc/.docx not supported)" - texts: list[str] = [] - errors: list[str] = [] + texts=[] + errors=[] for path in existing: try: - raw = path.read_bytes() - result = await FileRead(session=None, filename=path.name, file=raw).read_file() - text = (result.get("text") or "").strip() + raw=path.read_bytes() + result=await FileRead(session=None,filename=path.name,file=raw).read_file() + text=(result.get("text") or "").strip() if text: texts.append(text) except Exception as exc: errors.append(f"{path.name}: {exc}") if not texts: - return "", "; ".join(errors) if errors else "no text extracted from PDF" - - return "\n\n---\n\n".join(texts), "" + return "","; ".join(errors) if errors else "no text extracted from PDF" + return "\n\n---\n\n".join(texts),"" diff --git a/backend/inbox/tasks.py b/backend/inbox/tasks.py new file mode 100644 index 0000000..c0d65cb --- /dev/null +++ b/backend/inbox/tasks.py @@ -0,0 +1,100 @@ +"""Inbox Taskiq tasks — CV → job-post matching.""" + +from __future__ import annotations + +import logging +from datetime import datetime, timezone + +from db_setup import session_scope +from inbox.models import Inbox_Messages +from inbox.plugins import extract_resume_text +from job.job_post.models import JobPosts +from job.job_post.serializers import serialize_job_post +from taskiq_management.broker_setup import MAX_RETRIES, RETRY_DELAY, broker +from taskiq_management.middleware import PermanentTaskError + +logger=logging.getLogger("inbox.tasks") + +_DONE_STATUSES=frozenset({"matched","skipped","no_text","failed","dlq"}) + + +@broker.task( + task_name="inbox.match_message", + retry_on_error=True, + max_retries=MAX_RETRIES, + delay=RETRY_DELAY, +) +async def match_inbox_message(record_id:str,force:bool=False) -> dict: + if not record_id or not str(record_id).strip(): + raise PermanentTaskError("record_id is required") + + record_id=str(record_id).strip() + + async with session_scope() as session: + row=await Inbox_Messages.get_inbox_message_by_id(session,record_id) + if not row: + raise PermanentTaskError(f"inbox message {record_id} not found") + + if not force and row.match_status in _DONE_STATUSES: + logger.info("skip %s — already %s",record_id,row.match_status) + return {"status":row.match_status,"skipped":True} + + if not row.attachment or not row.file_path: + raise PermanentTaskError("message has no attachment to match") + + row.match_status="processing" + row.match_error=None + row.matched_at=datetime.now(timezone.utc) + session.add(row) + await session.commit() + await session.refresh(row) + + paths=[p.strip() for p in (row.file_path or "").split(",") if p.strip()] + subject=row.message_subject or "" + + text,extract_err=await extract_resume_text(paths) + if not text: + async with session_scope() as session: + await Inbox_Messages.set_match_result( + session, + record_id, + status="no_text", + error=extract_err or "no text extracted", + ) + return {"status":"no_text","error":extract_err} + + from agent.execute_agent import run_agent + + async with session_scope() as session: + posts=await JobPosts.get_active_job_posts(session) + job_posts=[serialize_job_post(p) for p in posts] + + try: + result=await run_agent(subject=subject,resume_text=text,job_posts=job_posts) + except Exception as exc: + logger.exception("agent failed for %s",record_id) + raise RuntimeError(f"agent matching failed: {exc}") from exc + + status=result.get("status") or "failed" + error=result.get("error") or "" + + if status=="failed": + raise RuntimeError(error or "agent returned failed status") + + async with session_scope() as session: + await Inbox_Messages.set_match_result( + session, + record_id, + resume_text=text, + suggested_job_post_ids=result.get("suggested_job_post_ids") or [], + summary=result.get("summary") or "", + reasoning=result.get("reasoning") or "", + status=status, + error=error, + ) + + logger.info("matched inbox %s status=%s ids=%s",record_id,status,result.get("suggested_job_post_ids")) + return { + "status":status, + "suggested_job_post_ids":result.get("suggested_job_post_ids") or [], + } diff --git a/backend/inbox/views.py b/backend/inbox/views.py index 82742e7..9ae3313 100644 --- a/backend/inbox/views.py +++ b/backend/inbox/views.py @@ -2,19 +2,15 @@ import logging import httpx,os from fastapi import HTTPException from inbox.models import Inbox_Messages -from inbox.file_decoder import decode_attachment, AttachmentDecodeError +from inbox.file_decoder import decode_attachment from inbox.serializers import serialize_message -from inbox.plugins import load_message_files, extract_resume_text +from inbox.plugins import load_message_files from dotenv import load_dotenv load_dotenv() from sqlalchemy.ext.asyncio import AsyncSession -from pydantic import BaseModel +from datetime import datetime,timezone -from db_setup import session_scope -from job.job_post.models import JobPosts -from job.job_post.serializers import serialize_job_post - -logger = logging.getLogger("inbox.match") +logger=logging.getLogger("inbox.match") class Email: @@ -22,75 +18,7 @@ class Email: self.session=session self.get_url=os.getenv("EMAIL_URL") self.token=token - self.pending_match_ids: list[str] = [] - - @staticmethod - async def run_inbox_matching(inbox_ids: list[str], *, force: bool = False) -> None: - """Background worker: extract CV text, run the matching agent, persist results. - - Uses its own session_scope per message - never the request session. - """ - if not inbox_ids: - return - - # Local import so inbox routes still load when langgraph/openai are missing. - from agent.execute_agent import run_agent - - cached_posts = None - for record_id in inbox_ids: - try: - async with session_scope() as session: - row = await Inbox_Messages.get_inbox_message_by_id(session, record_id) - if not row: - continue - if not force and row.match_status is not None: - continue - - paths = [p.strip() for p in (row.file_path or "").split(",") if p.strip()] - text, extract_err = await extract_resume_text(paths) - if not text: - await Inbox_Messages.set_match_result( - session, - record_id, - status="no_text", - error=extract_err or "no text extracted", - ) - continue - - if cached_posts is None: - posts = await JobPosts.get_active_job_posts(session) - cached_posts = [serialize_job_post(p) for p in posts] - - result = await run_agent( - subject=row.message_subject or "", - resume_text=text, - job_posts=cached_posts, - ) - await Inbox_Messages.set_match_result( - session, - record_id, - resume_text=text, - suggested_job_post_ids=result.get("suggested_job_post_ids") or [], - summary=result.get("summary") or "", - reasoning=result.get("reasoning") or "", - status=result.get("status") or "failed", - error=result.get("error") or "", - ) - logger.info( - "matched inbox %s status=%s ids=%s", - record_id, - result.get("status"), - result.get("suggested_job_post_ids"), - ) - except Exception as exc: - logger.exception("inbox matching failed for %s", record_id) - try: - async with session_scope() as session: - await Inbox_Messages.set_match_result( - session, record_id, status="failed", error=str(exc) - ) - except Exception: - logger.exception("failed to persist match failure for %s", record_id) + self.pending_match_ids:list[str]=[] async def service_email(self,top,skip): async with httpx.AsyncClient() as client: @@ -157,5 +85,18 @@ class Email: raise HTTPException(status_code=400,detail="Message has no attachment to match") return str(message.id) + async def enqueue_matching(self,inbox_ids,force=False): + from inbox.tasks import match_inbox_message + task_ids=[] + for record_id in inbox_ids or []: + created_at=datetime.now(timezone.utc).isoformat() + task=await match_inbox_message.kicker().with_labels( + created_at=created_at, + correlation_id=str(record_id), + queue="inbox", + ).kiq(str(record_id),force=force) + task_ids.append(task.task_id) + return task_ids + async def count_inbox_messages(self,search=None): return await Inbox_Messages.count_inbox_messages(self.session,search) diff --git a/backend/main.py b/backend/main.py index 4ffbe83..87ec33d 100644 --- a/backend/main.py +++ b/backend/main.py @@ -1,10 +1,8 @@ import logging from contextlib import asynccontextmanager -import fastapi from fastapi.middleware.cors import CORSMiddleware -from pydantic import BaseModel -from fastapi import FastAPI,APIRouter +from fastapi import FastAPI from db_setup import lifespan as db_lifespan from inbox.app import router as inbox_router from users.app import router as users_router @@ -12,31 +10,38 @@ from role.app import router as role_router from forget_password.app import router as forget_password_router from job.app import router as candidate_router from notifications.app import router as confirmation_router -# Without this the db/migration logs have no handler and are swallowed under uvicorn. -logging.basicConfig(level=logging.INFO, format="%(levelname)-8s %(name)s: %(message)s") -logger = logging.getLogger("main") + +logging.basicConfig(level=logging.INFO,format="%(levelname)-8s %(name)s: %(message)s") +logger=logging.getLogger("main") @asynccontextmanager async def lifespan(app): - """Compose DB lifespan with optional LLM/agent warm-up (degrades on failure).""" async with db_lifespan(app): - llm_ready = False - agent_ready = False - close_llm = None - close_agent = None + broker_ready=False + llm_ready=False + agent_ready=False + close_llm=None + close_agent=None + broker=None try: - from llm_setup import init_llm, close_llm as _close_llm - from agent.agent_setup import init_agent, close_agent as _close_agent - - close_llm = _close_llm - close_agent = _close_agent - await init_llm() - llm_ready = True - await init_agent() - agent_ready = True + from taskiq_management.broker_setup import broker as _broker + broker=_broker + await broker.startup() + broker_ready=True except Exception as exc: - logger.warning("llm/agent startup skipped: %s", exc) + logger.warning("taskiq broker startup skipped: %s",exc) + try: + from llm_setup import init_llm,close_llm as _close_llm + from agent.agent_setup import init_agent,close_agent as _close_agent + close_llm=_close_llm + close_agent=_close_agent + await init_llm() + llm_ready=True + await init_agent() + agent_ready=True + except Exception as exc: + logger.warning("llm/agent startup skipped: %s",exc) try: yield finally: @@ -44,11 +49,11 @@ async def lifespan(app): await close_agent() if llm_ready and close_llm is not None: await close_llm() + if broker_ready and broker is not None: + await broker.shutdown() -# lifespan connects to Postgres and brings migrations up to head on startup, -# and disposes of the connection pool on shutdown. -app = FastAPI(lifespan=lifespan) +app=FastAPI(lifespan=lifespan) app.add_middleware( CORSMiddleware, allow_origins=["*"], diff --git a/backend/requirements.txt b/backend/requirements.txt index 5ebc715..63ad39f 100644 --- a/backend/requirements.txt +++ b/backend/requirements.txt @@ -30,6 +30,11 @@ bcrypt==5.0.0 # password hashing in users/plugins.py # --- PDF extraction -------------------------------------------------------- pypdf==5.1.0 +# --- task queue ------------------------------------------------------------ +taskiq>=0.11,<0.12 # broker + worker/scheduler CLI (taskiq_management/) +taskiq-redis>=1.0,<2.0 # RedisStreamBroker / result backend / schedule source +redis>=5.0,<6.0 # DLQ middleware (taskiq_management/middleware.py) async client + # --- LLM ------------------------------------------------------------------- openai==2.53.0 # AsyncOpenAI client in llm_setup.py langgraph==1.2.10 # StateGraph agent framework in agent/agent_setup.py diff --git a/backend/taskiq_management/broker_setup.py b/backend/taskiq_management/broker_setup.py new file mode 100644 index 0000000..6313ab8 --- /dev/null +++ b/backend/taskiq_management/broker_setup.py @@ -0,0 +1,55 @@ +"""Taskiq broker — Redis Streams + smart retry + DLQ. + +Worker: taskiq worker taskiq_management.broker_setup:broker inbox.tasks taskiq_management.tasks +Scheduler: taskiq scheduler taskiq_management.broker_setup:scheduler +""" + +from __future__ import annotations + +import os + +from dotenv import load_dotenv +from taskiq import TaskiqScheduler +from taskiq.middlewares import SmartRetryMiddleware +from taskiq_redis import ( + ListRedisScheduleSource, + RedisAsyncResultBackend, + RedisStreamBroker, +) + +from taskiq_management.middleware import DeadLetterMiddleware + +load_dotenv() + +REDIS_URL=os.getenv("REDIS_URL","redis://localhost:6379/0") +QUEUE_NAME=os.getenv("TASKIQ_QUEUE_NAME","inbox") +# 2 retries after first failure → max_retries=3 +MAX_RETRIES=int(os.getenv("TASKIQ_MAX_RETRIES","3")) +RETRY_DELAY=float(os.getenv("TASKIQ_RETRY_DELAY","5")) + +result_backend=RedisAsyncResultBackend(redis_url=REDIS_URL) +schedule_source=ListRedisScheduleSource(url=REDIS_URL,prefix="taskiq:schedule") + +broker=( + RedisStreamBroker( + url=REDIS_URL, + queue_name=QUEUE_NAME, + consumer_group_name=os.getenv("TASKIQ_CONSUMER_GROUP","taskiq"), + idle_timeout=int(os.getenv("TASKIQ_IDLE_TIMEOUT_MS","600000")), + ) + .with_result_backend(result_backend) + .with_middlewares( + DeadLetterMiddleware(redis_url=REDIS_URL), + SmartRetryMiddleware( + default_retry_count=MAX_RETRIES, + default_retry_label=True, + default_delay=RETRY_DELAY, + use_jitter=True, + use_delay_exponent=True, + max_delay_exponent=float(os.getenv("TASKIQ_MAX_DELAY","120")), + schedule_source=schedule_source, + ), + ) +) + +scheduler=TaskiqScheduler(broker=broker,sources=[schedule_source]) diff --git a/backend/taskiq_management/middleware.py b/backend/taskiq_management/middleware.py new file mode 100644 index 0000000..a031dc0 --- /dev/null +++ b/backend/taskiq_management/middleware.py @@ -0,0 +1,102 @@ +"""PermanentTaskError + Redis Stream DLQ middleware for Taskiq. + +Middleware order: DeadLetterMiddleware before SmartRetryMiddleware so +permanent failures can set retry_on_error=False before SmartRetry runs. + +Pure module: no FastAPI imports. +""" + +from __future__ import annotations + +import json +import logging +from typing import Any + +import redis.asyncio as redis +from taskiq import TaskiqMiddleware +from taskiq.message import TaskiqMessage +from taskiq.result import TaskiqResult + +from taskiq_management.models import DLQ_STREAM +from taskiq_management.serializers import serialize_dlq_payload + +logger=logging.getLogger("taskiq.dlq") + + +class PermanentTaskError(Exception): + """Validation / business failure — DLQ immediately, no retries.""" + + +class DeadLetterMiddleware(TaskiqMiddleware): + def __init__(self,redis_url:str,stream:str=DLQ_STREAM): + super().__init__() + self.redis_url=redis_url + self.stream=stream + self._redis:redis.Redis|None=None + + async def startup(self) -> None: + self._redis=redis.from_url(self.redis_url,decode_responses=True) + + async def shutdown(self) -> None: + if self._redis is not None: + await self._redis.aclose() + self._redis=None + + def _client(self) -> redis.Redis: + if self._redis is None: + self._redis=redis.from_url(self.redis_url,decode_responses=True) + return self._redis + + async def on_error( + self, + message:TaskiqMessage, + result:TaskiqResult[Any], + exception:BaseException, + ) -> None: + retries=int(message.labels.get("_retries",0)) + max_retries=int(message.labels.get("max_retries",2)) + is_permanent=isinstance(exception,PermanentTaskError) + retries_exhausted=(retries+1)>=max_retries + + if is_permanent: + message.labels["retry_on_error"]=False + + if not is_permanent and not retries_exhausted: + return + + queue=getattr(self.broker,"queue_name",None) + payload=serialize_dlq_payload(message,exception,retries=retries+1,queue=queue) + try: + await self._client().xadd(self.stream,{"payload":json.dumps(payload,ensure_ascii=False,default=str)}) + logger.error( + "task %s (%s) sent to DLQ after %s", + message.task_name, + message.task_id, + "permanent failure" if is_permanent else f"{retries+1} attempts", + ) + except Exception: + logger.exception("failed to write DLQ entry for %s",message.task_id) + + await self._mark_inbox_dlq(message,exception) + + async def _mark_inbox_dlq(self,message:TaskiqMessage,exception:BaseException) -> None: + if message.task_name!="inbox.match_message": + return + record_id=(message.kwargs or {}).get("record_id") + if not record_id and message.args: + record_id=message.args[0] + if not record_id: + return + try: + from db_setup import session_scope + from inbox.models import Inbox_Messages + + async with session_scope() as session: + await Inbox_Messages.set_match_result( + session, + record_id, + status="dlq", + error=f"{type(exception).__name__}: {exception}", + ) + except Exception: + logger.exception("failed to mark inbox %s as dlq",record_id) diff --git a/backend/taskiq_management/models.py b/backend/taskiq_management/models.py new file mode 100644 index 0000000..f66669b --- /dev/null +++ b/backend/taskiq_management/models.py @@ -0,0 +1,15 @@ +"""Taskiq constants — DLQ stream + app version defaults. + +Pure module: no FastAPI imports. +""" + +from __future__ import annotations + +import os + +from dotenv import load_dotenv + +load_dotenv() + +DLQ_STREAM=os.getenv("TASKIQ_DLQ_STREAM","taskiq:dlq") +APP_VERSION=os.getenv("APP_VERSION","dev") diff --git a/backend/taskiq_management/serializers.py b/backend/taskiq_management/serializers.py new file mode 100644 index 0000000..497b8be --- /dev/null +++ b/backend/taskiq_management/serializers.py @@ -0,0 +1,44 @@ +"""DLQ payload serializers for Taskiq dead-letter entries. + +Pure module: no FastAPI imports. +""" + +from __future__ import annotations + +import os +import socket +import sys +import traceback +from datetime import datetime, timezone + +from taskiq.message import TaskiqMessage + +from taskiq_management.models import APP_VERSION + + +def serialize_dlq_payload( + message:TaskiqMessage, + exception:BaseException, + *, + retries:int, + queue:str|None=None, +) -> dict: + now=datetime.now(timezone.utc).isoformat() + return { + "task_name":message.task_name, + "task_id":message.task_id, + "kwargs":message.kwargs or {}, + "args":list(message.args or []), + "exception":type(exception).__name__, + "message":str(exception), + "traceback":"".join(traceback.format_exception(type(exception),exception,exception.__traceback__)), + "retry_count":retries, + "worker":os.getenv("TASKIQ_WORKER_NAME") or os.getenv("HOSTNAME") or socket.gethostname(), + "queue":message.labels.get("queue") or queue or "taskiq", + "created_at":message.labels.get("created_at") or now, + "failed_at":now, + "correlation_id":message.labels.get("correlation_id") or message.task_id, + "hostname":socket.gethostname(), + "python_version":sys.version.split()[0], + "app_version":APP_VERSION, + } diff --git a/backend/taskiq_management/tasks.py b/backend/taskiq_management/tasks.py new file mode 100644 index 0000000..ad5702c --- /dev/null +++ b/backend/taskiq_management/tasks.py @@ -0,0 +1,10 @@ +"""Framework smoke tasks for Taskiq — domain tasks stay in their packages.""" + +from __future__ import annotations + +from taskiq_management.broker_setup import broker + + +@broker.task(task_name="ping") +async def ping() -> str: + return "pong" diff --git a/docker-compose.yml b/docker-compose.yml index b2b1c23..3acc1cf 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,61 +1,63 @@ services: - minio: - image: minio/minio:RELEASE.2025-04-22T22-12-26Z - container_name: hrms-minio - command: server /data --console-address ":9001" - environment: - MINIO_ROOT_USER: ${MINIO_ROOT_USER:-minioadmin} - MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD:-minioadmin} + redis: + image: redis:7-alpine + container_name: hrms-redis + command: ["redis-server", "--appendonly", "yes"] ports: - - "9000:9000" # S3 API - - "9001:9001" # web console + - "${REDIS_PORT:-6379}:6379" volumes: - - minio-data:/data + - redis-data:/data healthcheck: - test: ["CMD", "mc", "ready", "local"] + test: ["CMD", "redis-cli", "ping"] interval: 10s timeout: 5s retries: 5 - start_period: 10s restart: unless-stopped - # One-shot: creates the attachments bucket, then exits. - minio-init: - image: minio/mc:RELEASE.2025-04-16T18-13-26Z - container_name: hrms-minio-init - depends_on: - minio: - condition: service_healthy + taskiq-worker: + build: + context: ./backend + container_name: hrms-taskiq-worker + command: + [ + "taskiq", + "worker", + "taskiq_management.broker_setup:broker", + "inbox.tasks", + "taskiq_management.tasks", + "--workers", + "1", + ] + env_file: + - ./backend/.env environment: - MINIO_ROOT_USER: ${MINIO_ROOT_USER:-minioadmin} - MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD:-minioadmin} - MINIO_BUCKET: ${MINIO_BUCKET:-hrms-attachments} - entrypoint: > - /bin/sh -c " - mc alias set local http://minio:9000 \"$$MINIO_ROOT_USER\" \"$$MINIO_ROOT_PASSWORD\" && - mc mb --ignore-existing local/\"$$MINIO_BUCKET\" && - mc version enable local/\"$$MINIO_BUCKET\" && - echo 'bucket ready: '\"$$MINIO_BUCKET\" - " - - postgres: - image: postgres:16-alpine - container_name: hrms-postgres - environment: - POSTGRES_USER: ${DB_USERNAME:-postgres} - POSTGRES_PASSWORD: ${DB_PASSWORD:-postgres} - POSTGRES_DB: ${DB_NAME:-hrms} - ports: - - "${DB_PORT:-5432}:5432" + REDIS_URL: redis://redis:6379/0 + TASKIQ_QUEUE_NAME: inbox + TASKIQ_WORKER_NAME: worker-01 + DB_HOST: ${DB_HOST:-host.docker.internal} + extra_hosts: + - "host.docker.internal:host-gateway" volumes: - - postgres-data:/var/lib/postgresql/data - healthcheck: - test: ["CMD-SHELL", "pg_isready -U ${DB_USERNAME:-postgres} -d ${DB_NAME:-hrms}"] - interval: 10s - timeout: 5s - retries: 5 + - ./backend/inbox/decoded_attachments:/app/inbox/decoded_attachments + depends_on: + redis: + condition: service_healthy + restart: unless-stopped + + taskiq-scheduler: + build: + context: ./backend + container_name: hrms-taskiq-scheduler + command: ["taskiq", "scheduler", "taskiq_management.broker_setup:scheduler"] + env_file: + - ./backend/.env + environment: + REDIS_URL: redis://redis:6379/0 + TASKIQ_QUEUE_NAME: inbox + depends_on: + redis: + condition: service_healthy restart: unless-stopped volumes: - minio-data: - postgres-data: + redis-data: