Backend_CODEBASE
parent
dcf0bf5d50
commit
588afd0a28
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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"]
|
||||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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),""
|
||||
|
|
|
|||
|
|
@ -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 [],
|
||||
}
|
||||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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=["*"],
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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])
|
||||
|
|
@ -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)
|
||||
|
|
@ -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")
|
||||
|
|
@ -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,
|
||||
}
|
||||
|
|
@ -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"
|
||||
|
|
@ -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:
|
||||
|
|
|
|||
Loading…
Reference in New Issue