backend manual uplaod

pull/10/head
ahmed.mujtaba 2026-08-11 20:34:02 +05:00
parent e897ac8d5e
commit 326842ebba
11 changed files with 500 additions and 12 deletions

View File

@ -46,9 +46,11 @@ OPENAI_PROJECT=
REDIS_URL=redis://localhost:6379/0 REDIS_URL=redis://localhost:6379/0
TASKIQ_QUEUE_NAME=inbox TASKIQ_QUEUE_NAME=inbox
TASKIQ_CV_QUEUE_NAME=cv_upload
TASKIQ_MAX_RETRIES=3 TASKIQ_MAX_RETRIES=3
TASKIQ_RETRY_DELAY=5 TASKIQ_RETRY_DELAY=5
TASKIQ_MAX_DELAY=120 TASKIQ_MAX_DELAY=120
TASKIQ_DLQ_STREAM=taskiq:dlq TASKIQ_DLQ_STREAM=taskiq:dlq
TASKIQ_IDLE_TIMEOUT_MS=600000 TASKIQ_IDLE_TIMEOUT_MS=600000
MANUAL_UPLOAD_TO_ADDRESS=manual-cv-upload@hr-ats.local
APP_VERSION=dev APP_VERSION=dev

View File

@ -47,9 +47,13 @@ flowchart TB
API -->|enqueue| REDIS[("Redis Streams")] API -->|enqueue| REDIS[("Redis Streams")]
REDIS --> W["Taskiq worker\ninbox.tasks + inbox.sync_tasks"] REDIS --> W["Taskiq worker\ninbox.tasks + inbox.sync_tasks"]
REDIS --> WCV["Taskiq CV worker\ninbox.cv_tasks"]
SCHED["Taskiq scheduler\ncron"] --> REDIS SCHED["Taskiq scheduler\ncron"] --> REDIS
SCHEDCV["Taskiq CV scheduler\nretries"] --> REDIS
W --> PG W --> PG
WCV --> PG
W --> AGENT["LangGraph agent\nagent/"] W --> AGENT["LangGraph agent\nagent/"]
WCV --> AGENT
AGENT --> OAI["OpenAI"] AGENT --> OAI["OpenAI"]
API -->|GET /emails, /sync/read-status| MAILAPI["Email API (MS Graph proxy)"] API -->|GET /emails, /sync/read-status| MAILAPI["Email API (MS Graph proxy)"]
@ -305,7 +309,7 @@ Base URL: `http://localhost:8000`. Interactive docs at `/docs`.
| GET | `/jobs/alias` | public | Accepted platform shorthands (`fb`, `ig`, `li`, `x`, …) | | GET | `/jobs/alias` | public | Accepted platform shorthands (`fb`, `ig`, `li`, `x`, …) |
| POST | `/job/post-job` | `job_board.create` | Render the ad, create the Buffer post, persist the result | | POST | `/job/post-job` | `job_board.create` | Render the ad, create the Buffer post, persist the result |
| GET | `/job/buffer/channels` | `job_board.view` | Connected Buffer channels across all organizations | | GET | `/job/buffer/channels` | `job_board.view` | Connected Buffer channels across all organizations |
| POST | `/candidate/cv_upload` | `candidates.create` | Upload a PDF, get extracted text back | | POST | `/candidate/cv_upload` | `candidates.create` | Upload a PDF; extract email, persist like an emailed CV, enqueue matching on the CV stream |
| POST | `/candidate/inbox-match?inbox_message_id=` | `candidates.edit` | Queue a forced re-match for a stored message | | POST | `/candidate/inbox-match?inbox_message_id=` | `candidates.edit` | Queue a forced re-match for a stored message |
| GET | `/candidate/fetch?user_id=` | `candidates.view` | Candidate profile via the `inbox` join | | GET | `/candidate/fetch?user_id=` | `candidates.view` | Candidate profile via the `inbox` join |
@ -367,7 +371,8 @@ Taskiq over **Redis Streams**, with a result backend and a Redis-backed schedule
| Task | Trigger | What it does | | Task | Trigger | What it does |
|---|---|---| |---|---|---|
| `inbox.match_message` | enqueued by `/email/fetch`, `/inbox/{id}/match`, `/candidate/inbox-match` | Extract résumé text → run the agent → write match results | | `inbox.match_message` | enqueued by `/email/fetch`, `/inbox/{id}/match`, `/candidate/inbox-match` onto the `inbox` stream | Extract résumé text → run the agent → write match results |
| `inbox.match_message` (CV broker) | enqueued by `/candidate/cv_upload` onto the `cv_upload` stream | Same matcher as above; isolated so uploads never sit behind `/email/fetch` backlog |
| `inbox.sync_read_status` | cron, `EMAIL_SYNC_CRON` (default every minute) | Pull read-status deltas from the Email API and apply them | | `inbox.sync_read_status` | cron, `EMAIL_SYNC_CRON` (default every minute) | Pull read-status deltas from the Email API and apply them |
| `ping` | manual | Framework smoke test | | `ping` | manual | Framework smoke test |
@ -512,6 +517,7 @@ own keys with `os.getenv`.
|---|---| |---|---|
| `REDIS_URL` | `redis://localhost:6379/0` | | `REDIS_URL` | `redis://localhost:6379/0` |
| `TASKIQ_QUEUE_NAME` | `inbox` | | `TASKIQ_QUEUE_NAME` | `inbox` |
| `TASKIQ_CV_QUEUE_NAME` | `cv_upload` |
| `TASKIQ_CONSUMER_GROUP` | `taskiq` | | `TASKIQ_CONSUMER_GROUP` | `taskiq` |
| `TASKIQ_MAX_RETRIES` | `3` | | `TASKIQ_MAX_RETRIES` | `3` |
| `TASKIQ_RETRY_DELAY` | `5` | | `TASKIQ_RETRY_DELAY` | `5` |
@ -519,6 +525,7 @@ own keys with `os.getenv`.
| `TASKIQ_IDLE_TIMEOUT_MS` | `600000` | | `TASKIQ_IDLE_TIMEOUT_MS` | `600000` |
| `TASKIQ_DLQ_STREAM` | `taskiq:dlq` | | `TASKIQ_DLQ_STREAM` | `taskiq:dlq` |
| `TASKIQ_WORKER_NAME` | falls back to `HOSTNAME` | | `TASKIQ_WORKER_NAME` | falls back to `HOSTNAME` |
| `MANUAL_UPLOAD_TO_ADDRESS` | `manual-cv-upload@hr-ats.local` — To address stamped on synthetic inbox rows so source resolves to `Manual CV Upload` |
| `APP_VERSION` | `dev` | | `APP_VERSION` | `dev` |
--- ---
@ -556,12 +563,24 @@ taskiq worker taskiq_management.broker_setup:broker \
inbox.tasks inbox.sync_tasks taskiq_management.tasks inbox.tasks inbox.sync_tasks taskiq_management.tasks
``` ```
**CV-upload worker** (isolated stream for manual uploads):
```bash
taskiq worker taskiq_management.cv_broker_setup:cv_broker inbox.cv_tasks
```
**Scheduler** (cron ticks for `inbox.sync_read_status`): **Scheduler** (cron ticks for `inbox.sync_read_status`):
```bash ```bash
taskiq scheduler taskiq_management.broker_setup:scheduler inbox.sync_tasks taskiq scheduler taskiq_management.broker_setup:scheduler inbox.sync_tasks
``` ```
**CV-upload scheduler** (retries for the CV stream):
```bash
taskiq scheduler taskiq_management.cv_broker_setup:cv_scheduler inbox.cv_tasks
```
Docs: <http://localhost:8000/docs> Docs: <http://localhost:8000/docs>
--- ---
@ -594,13 +613,14 @@ and a module called `alembic.py` would shadow the installed package.
## Docker ## Docker
The repo-root `docker-compose.yml` runs Redis plus the two Taskiq processes; the API itself is The repo-root `docker-compose.yml` runs Redis plus the four Taskiq processes (inbox
expected to run on the host (the compose file points the containers at worker/scheduler and CV-upload worker/scheduler); the API itself is expected to run on the
`host.docker.internal` for the database). host (the compose file points the containers at `host.docker.internal` for the database).
```bash ```bash
docker compose up -d # from the repo root docker compose up -d # from the repo root
docker compose logs -f taskiq-worker docker compose logs -f taskiq-worker
docker compose logs -f taskiq-cv-worker
``` ```
`backend/Dockerfile` builds a `python:3.12-slim` image whose default command is the Taskiq `backend/Dockerfile` builds a `python:3.12-slim` image whose default command is the Taskiq

17
backend/inbox/cv_tasks.py Normal file
View File

@ -0,0 +1,17 @@
"""CV-upload Taskiq tasks — same matcher as inbox.tasks, own broker/stream."""
from __future__ import annotations
from inbox.tasks import match_inbox_message
from taskiq_management.broker_setup import MAX_RETRIES,RETRY_DELAY
from taskiq_management.cv_broker_setup import cv_broker
@cv_broker.task(
task_name="inbox.match_message",
retry_on_error=True,
max_retries=MAX_RETRIES,
delay=RETRY_DELAY,
)
async def match_uploaded_cv(record_id:str,force:bool=False) -> dict:
return await match_inbox_message(record_id,force)

View File

@ -12,7 +12,7 @@ from users.permissions import PermissionTag, require_permission
from job.job_post.views import JobPost,JobPostCreate from job.job_post.views import JobPost,JobPostCreate
import logging import logging
from job.job_post.plugins import PlatformAlias from job.job_post.plugins import PlatformAlias
from fastapi import UploadFile, File from fastapi import UploadFile, File, Form
from dotenv import load_dotenv from dotenv import load_dotenv
from datetime import datetime, time, timezone from datetime import datetime, time, timezone
from pydantic import BaseModel from pydantic import BaseModel
@ -95,10 +95,52 @@ async def get_job_alias():
except Exception as e: except Exception as e:
raise HTTPException(status_code=500,detail=str(e)) raise HTTPException(status_code=500,detail=str(e))
@router.post("/candidate/create/candidate")
async def create_manual_candidate(
file: UploadFile = File(...),
candidate_email: str | None = Form(None),
candidate_name: str | None = Form(None),
candidate_phone: str | None = Form(None),
job_post_id: str | None = Form(None),
current_company: str | None = Form(None),
platform: str | None = Form(None),
experience: str | None = Form(None),
status: str | None = Form(None),
current_user: dict = Depends(require_permission(PermissionTag.CANDIDATES_CREATE)),
session: AsyncSession = Depends(get_session),
):
try:
file_content = await file.read()
logger.info(f"Received file: {file.filename} ({len(file_content)} bytes)")
reader=FileRead(session=session,filename=file.filename,file=file_content)
parsed=await reader.injest_manual_upload()
service=CandidateView(session=session)
data=await service.create_candidate(
candidate_email=candidate_email,
candidate_name=candidate_name,
candidate_phone=candidate_phone,
job_post_id=job_post_id,
current_company=current_company,
platform=platform,
experience=experience,
status=status,
full_text=parsed.get("text") or "",
current_user=current_user.get("id"),
)
return JSONResponse(content={"data":data,"status_code":200})
except HTTPException:
raise
except Exception as e:
raise HTTPException(status_code=500,detail=str(e))
@router.post("/candidate/cv_upload") @router.post("/candidate/cv_upload")
async def cv_upload( async def cv_upload(
file: UploadFile = File(...), file: UploadFile = File(...),
candidate_email: str | None = Form(None),
candidate_name: str | None = Form(None),
candidate_phone: str | None = Form(None),
job_post_id: str | None = Form(None),
current_user: dict = Depends(require_permission(PermissionTag.CANDIDATES_CREATE)), current_user: dict = Depends(require_permission(PermissionTag.CANDIDATES_CREATE)),
session: AsyncSession = Depends(get_session), session: AsyncSession = Depends(get_session),
): ):
@ -106,8 +148,9 @@ async def cv_upload(
file_content = await file.read() file_content = await file.read()
logger.info(f"Received file: {file.filename} ({len(file_content)} bytes)") logger.info(f"Received file: {file.filename} ({len(file_content)} bytes)")
service=FileRead(session=session,filename=file.filename,file=file_content) service=FileRead(session=session,filename=file.filename,file=file_content)
data=await service.read_file() data=await service.ingest_upload(
candidate_email=candidate_email,candidate_name=candidate_name,
)
return JSONResponse(content={"data":data,"status_code":200}) return JSONResponse(content={"data":data,"status_code":200})
except HTTPException: except HTTPException:
raise raise

View File

@ -9,6 +9,7 @@ from sqlmodel import Field, Relationship, SQLModel, select
if TYPE_CHECKING: if TYPE_CHECKING:
from inbox.models import Inbox from inbox.models import Inbox
from users.models import Users from users.models import Users
from job.job_post.models import JobPosts
def _now() -> datetime: def _now() -> datetime:
@ -20,6 +21,73 @@ def _now() -> datetime:
# `datetime` to TIMESTAMP WITHOUT TIME ZONE, and asyncpg refuses to bind an aware # `datetime` to TIMESTAMP WITHOUT TIME ZONE, and asyncpg refuses to bind an aware
# value to one — "can't subtract offset-naive and offset-aware datetimes" — which # value to one — "can't subtract offset-naive and offset-aware datetimes" — which
# turns every insert here into a 500. Same pairing as job/job_post/models.py. # turns every insert here into a 500. Same pairing as job/job_post/models.py.
class Manual_UPLOAD_CANDIDATE(SQLModel, table=True):
__tablename__ = "manual_upload_candidate"
id: uuid.UUID = Field(default_factory=uuid.uuid4, primary_key=True)
candidate_email: str = Field(default="")
candidate_name: str = Field(default="")
candidate_phone: str = Field(default="")
job_post_id: uuid.UUID | None = Field(default=None, foreign_key="job_posts.id")
full_text: str = Field(default="")
current_company: str = Field(default="")
user_id: uuid.UUID | None = Field(default=None, foreign_key="users.id")
platform: str = Field(default="")
created_by: uuid.UUID | None = Field(default=None, foreign_key="users.id")
experience: str = Field(default="")
status: str = Field(default="")
created_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True))
updated_at: datetime = Field(default_factory=_now, sa_type=DateTime(timezone=True))
@staticmethod
def _as_uuid(record_id) -> uuid.UUID | None:
if record_id in (None, ""):
return None
try:
return uuid.UUID(str(record_id))
except ValueError:
return None
@classmethod
async def create_manual_upload_candidate(cls, session: AsyncSession, fields: dict):
import os
from role.models import EnumRoles, Roles
from users.models import Users
from users.plugins import hash_password
email=(fields.get("candidate_email") or "").strip().lower()
name=(fields.get("candidate_name") or "").strip() or email
default_pw=os.getenv("DEFAULT_CANDIDATE_PASSWORD","Utopia!@#")
user=await Users.get_user_by_email(session,email)
if not user:
role=await Roles.get_role_by_name(session,EnumRoles.CANDIDATE.value)
user=await Users.insert_user(session,{
"name":name,
"email":email,
"role_id":role.id if role else 8,
"password":hash_password(default_pw),
"is_active":True,
"is_deleted":False,
})
row=cls(
candidate_email=email,
candidate_name=name,
candidate_phone=(fields.get("candidate_phone") or "").strip(),
job_post_id=cls._as_uuid(fields.get("job_post_id")),
full_text=fields.get("full_text") or "",
current_company=(fields.get("current_company") or "").strip(),
user_id=user.id,
platform=(fields.get("platform") or "").strip(),
created_by=cls._as_uuid(fields.get("created_by")),
experience=(fields.get("experience") or "").strip(),
status=(fields.get("status") or "").strip(),
)
session.add(row)
await session.commit()
await session.refresh(row)
return row
class Interviews(SQLModel, table=True): class Interviews(SQLModel, table=True):

View File

@ -71,3 +71,111 @@ def documents_from_message(file_name: str | None, file_path: str | None) -> list
out.append({"name": path.rsplit("/", 1)[-1].rsplit("\\", 1)[-1], "path": path}) out.append({"name": path.rsplit("/", 1)[-1].rsplit("\\", 1)[-1], "path": path})
return out return out
# Same system prefixes inbox/models._is_linkable_sender rejects — anything we
# accept here must remain linkable when insert_email creates the users row.
_SKIP_SENDER_PREFIXES = (
"noreply", "no-reply", "donotreply", "do-not-reply",
"mailer-daemon", "postmaster", "bounce",
)
_ROLE_LOCAL_PARTS = frozenset({
"info", "hr", "careers", "jobs", "admin", "support", "contact", "sales",
"recruitment", "office", "team", "hello", "enquiry", "inquiry", "recruit",
"talent", "hiring", "apply", "applications", "webmaster", "helpdesk",
})
_EMAIL_RE = re.compile(
r"(?i)\b([a-z0-9][a-z0-9._%+\-]{0,63})@([a-z0-9](?:[a-z0-9\-]{0,61}[a-z0-9])?"
r"(?:\.[a-z0-9](?:[a-z0-9\-]{0,61}[a-z0-9])?)+)\b"
)
_PHONE_RE = re.compile(r"(?:\+?\d[\d\s\-().]{7,}\d)")
_LABEL_RE = re.compile(r"(?i)\b(?:e[\-\s]?mail|mail[\s\-]?id|contact)\b")
_REF_HEADING_RE = re.compile(r"(?i)^\s*(?:references?|referees?)\b")
_REF_MENTION_RE = re.compile(
r"(?i)\b(?:reference|referee|manager|supervisor|contact\s+person)\b"
)
_HEADER_LINE_COUNT = 12
_MIN_ACCEPT_SCORE = 3
def _presumed_name_tokens(lines: list[str]) -> list[str]:
"""First non-empty line with 2+ alpha tokens and no digits/@ — CV name header."""
for line in lines:
stripped = line.strip()
if not stripped:
continue
if any(ch.isdigit() for ch in stripped) or "@" in stripped:
continue
tokens = [_letters_only(t) for t in re.split(r"\s+", stripped) if _letters_only(t)]
if len(tokens) >= 2:
return tokens
return []
def _email_local_ok(local: str) -> bool:
lowered = (local or "").lower()
if lowered in _ROLE_LOCAL_PARTS:
return False
return not lowered.startswith(_SKIP_SENDER_PREFIXES)
def extract_candidate_email(text: str) -> tuple[str | None, list[str]]:
"""Pick the candidate's own email from CV text, or None when ambiguous/absent.
Returns ``(best, all_plausible)``. Ambiguity is intentional a wrong guess
would create a user under a stranger's address and mail them a confirm link.
"""
if not text or not text.strip():
return None, []
lines = text.splitlines()
name_tokens = _presumed_name_tokens(lines)
in_references = False
scored: list[tuple[int, int, str]] = [] # (score, first_line_idx, email)
seen: dict[str, int] = {} # lower email -> index in scored
for idx, line in enumerate(lines):
if _REF_HEADING_RE.search(line):
in_references = True
for match in _EMAIL_RE.finditer(line):
local, domain = match.group(1), match.group(2)
if not _email_local_ok(local):
continue
email = f"{local}@{domain}".lower()
score = 0
if idx < _HEADER_LINE_COUNT:
score += 3
local_letters = _letters_only(local)
if local_letters and any(
tok and (tok in local_letters or local_letters in tok)
for tok in name_tokens
):
score += 3
if _LABEL_RE.search(line) or _PHONE_RE.search(line):
score += 1
if in_references:
score -= 5
if _REF_MENTION_RE.search(line):
score -= 3
if email in seen:
prev_i = seen[email]
prev_score, _, _ = scored[prev_i]
if score > prev_score:
scored[prev_i] = (score, idx, email)
continue
seen[email] = len(scored)
scored.append((score, idx, email))
if not scored:
return None, []
scored.sort(key=lambda t: (-t[0], t[1]))
plausible = [email for _, _, email in scored]
top_score, _, top_email = scored[0]
runner_up = scored[1][0] if len(scored) > 1 else None
if top_score < _MIN_ACCEPT_SCORE:
return None, plausible
if runner_up is not None and top_score <= runner_up:
return None, plausible
return top_email, plausible

View File

@ -7,6 +7,25 @@ from job.activity.serializers import serialize_activity
from job.feedback.serializers import serialize_feedback from job.feedback.serializers import serialize_feedback
def serialize_manual_upload_candidate(row) -> Dict[str,Any]:
return {
"id":str(row.id) if row.id else None,
"candidate_email":row.candidate_email,
"candidate_name":row.candidate_name,
"candidate_phone":row.candidate_phone,
"job_post_id":str(row.job_post_id) if row.job_post_id else None,
"full_text":row.full_text,
"current_company":row.current_company,
"user_id":str(row.user_id) if row.user_id else None,
"platform":row.platform,
"created_by":str(row.created_by) if row.created_by else None,
"experience":row.experience,
"status":row.status,
"created_at":row.created_at.isoformat() if row.created_at else None,
"updated_at":row.updated_at.isoformat() if row.updated_at else None,
}
def serialize_candidate_profile( def serialize_candidate_profile(
link:Inbox|List[Inbox]|Dict[str,Any]|List[Dict[str,Any]], link:Inbox|List[Inbox]|Dict[str,Any]|List[Dict[str,Any]],
*, *,

View File

@ -1,6 +1,7 @@
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
import os,logging,io import base64,io,logging,os,uuid
from datetime import datetime,timezone from datetime import datetime,timezone
from dotenv import load_dotenv
from fastapi import HTTPException from fastapi import HTTPException
from pypdf import PdfReader from pypdf import PdfReader
from sqlalchemy import select from sqlalchemy import select
@ -8,11 +9,18 @@ from sqlalchemy.orm import selectinload
from sqlmodel import true from sqlmodel import true
from job.job_post.models import JobPosts from job.job_post.models import JobPosts
from job.job_post.serializers import serialize_job_post from job.job_post.serializers import serialize_job_post
from job.candidate.serializers import serialize_candidate_profile from job.candidate.serializers import serialize_candidate_profile,serialize_manual_upload_candidate
from job.candidate.models import Notes from job.candidate.models import Notes,Manual_UPLOAD_CANDIDATE
from job.notes.serializers import serialize_note from job.notes.serializers import serialize_note
from inbox.models import Inbox_Messages,Inbox from inbox.models import Inbox_Messages,Inbox
from job.candidate.plugins import normalize_spaced_text from job.candidate.plugins import extract_candidate_email,normalize_spaced_text
load_dotenv()
logger=logging.getLogger("job.candidate.views")
CV_QUEUE_NAME=os.getenv("TASKIQ_CV_QUEUE_NAME","cv_upload")
MANUAL_UPLOAD_TO_ADDRESS=os.getenv(
"MANUAL_UPLOAD_TO_ADDRESS","manual-cv-upload@hr-ats.local"
)
class FileRead: class FileRead:
def __init__(self,session:AsyncSession,filename=None,file=None): def __init__(self,session:AsyncSession,filename=None,file=None):
@ -35,6 +43,102 @@ class FileRead:
raise raise
except Exception as e: except Exception as e:
raise HTTPException(400, str(e)) raise HTTPException(400, str(e))
async def injest_manual_upload(self):
try:
parsed=await self.read_file()
text=(parsed.get("text") or "").strip()
if not text:
raise HTTPException(status_code=400,detail="No usable text could be extracted from the PDF")
return parsed
except HTTPException:
raise
except Exception as e:
raise HTTPException(status_code=400,detail=str(e))
async def ingest_upload(self,candidate_email=None,candidate_name=None):
"""Persist a recruiter-uploaded CV with full email-ingestion parity."""
from inbox.file_decoder import AttachmentDecodeError,decode_attachment
from inbox.cv_tasks import match_uploaded_cv
from inbox.views import Email
parsed=await self.read_file()
text=parsed.get("text") or ""
detected,emails_found=extract_candidate_email(text)
supplied=(candidate_email or "").strip().lower() or None
email=supplied or detected
email_source="recruiter" if supplied else ("cv" if detected else None)
if not email:
raise HTTPException(
status_code=422,
detail={
"error_code":"CANDIDATE_EMAIL_REQUIRED",
"filename":parsed.get("filename"),
"num_pages":parsed.get("num_pages"),
"emails_found":emails_found,
"text":text,
},
)
filename=self.filename or "resume.pdf"
try:
paths=await decode_attachment([{
"name":filename,
"contentBytes":base64.b64encode(self.file).decode("ascii"),
}])
except AttachmentDecodeError as e:
raise HTTPException(status_code=400,detail=str(e))
if not paths:
raise HTTPException(status_code=400,detail="attachment could not be saved")
now=datetime.now(timezone.utc).isoformat()
email_data={
"id":f"manual-cv:{uuid.uuid4()}",
"subject":f"Manual CV upload — {filename}",
"body":{"content":"","contentType":"text"},
"hasAttachments":True,
"attachments":[{"name":filename}],
"from":{"emailAddress":{"address":email,"name":(candidate_name or "").strip()}},
"toRecipients":[{"emailAddress":{"address":MANUAL_UPLOAD_TO_ADDRESS}}],
"ccRecipients":[],
"bccRecipients":[],
"replyTo":[],
"isRead":False,
"sentDateTime":now,
"receivedDateTime":now,
}
row,new_user_email=await Inbox_Messages.insert_email(
self.session,email_data,file_path=paths,
)
created_at=datetime.now(timezone.utc).isoformat()
task=await match_uploaded_cv.kicker().with_labels(
created_at=created_at,
correlation_id=str(row.id),
queue=CV_QUEUE_NAME,
).kiq(str(row.id),force=False)
account_setup=None
if new_user_email:
try:
account_setup=await Email(session=self.session).send_account_setup(
[new_user_email]
)
except Exception as e:
logger.warning("account setup mail failed for %s: %s",new_user_email,e)
account_setup=[{"email":new_user_email,"sent":False}]
return {
"queued":True,
"inbox_message_id":str(row.id),
"task_id":task.task_id,
"filename":parsed.get("filename"),
"num_pages":parsed.get("num_pages"),
"candidate_email":email,
"email_source":email_source,
"account_setup":account_setup,
"text":text,
}
async def match_inbox_cv(self,inbox_message_id): async def match_inbox_cv(self,inbox_message_id):
from inbox.plugins import resolve_attachment_path from inbox.plugins import resolve_attachment_path
@ -78,6 +182,32 @@ class CandidateView:
def __init__(self,session:AsyncSession): def __init__(self,session:AsyncSession):
self.session=session self.session=session
async def create_candidate(self,candidate_email=None,candidate_name=None,candidate_phone=None,job_post_id=None,current_company=None,platform=None,experience=None,status=None,full_text=None,current_user=None):
try:
email=(candidate_email or "").strip().lower()
if not email:
raise HTTPException(status_code=422,detail="candidate_email is required")
if not current_user:
raise HTTPException(status_code=400,detail="created_by is required")
data={
"candidate_email":email,
"candidate_name":(candidate_name or "").strip(),
"candidate_phone":(candidate_phone or "").strip(),
"job_post_id":job_post_id,
"current_company":(current_company or "").strip(),
"platform":(platform or "").strip(),
"experience":(experience or "").strip(),
"status":(status or "").strip(),
"full_text":full_text or "",
"created_by":current_user,
}
row=await Manual_UPLOAD_CANDIDATE.create_manual_upload_candidate(session=self.session,fields=data)
return serialize_manual_upload_candidate(row)
except HTTPException:
raise
except Exception as e:
raise HTTPException(status_code=500,detail=str(e))
async def get_candidate(self,user_id=None,limit=10,offset=0,search=None): async def get_candidate(self,user_id=None,limit=10,offset=0,search=None):
try: try:
detail=bool(user_id) detail=bool(user_id)

View File

@ -19,11 +19,13 @@ logger=logging.getLogger("main")
async def lifespan(app): async def lifespan(app):
async with db_lifespan(app): async with db_lifespan(app):
broker_ready=False broker_ready=False
cv_broker_ready=False
llm_ready=False llm_ready=False
agent_ready=False agent_ready=False
close_llm=None close_llm=None
close_agent=None close_agent=None
broker=None broker=None
cv_broker=None
try: try:
from taskiq_management.broker_setup import broker as _broker from taskiq_management.broker_setup import broker as _broker
broker=_broker broker=_broker
@ -31,6 +33,13 @@ async def lifespan(app):
broker_ready=True broker_ready=True
except Exception as exc: except Exception as exc:
logger.warning("taskiq broker startup skipped: %s",exc) logger.warning("taskiq broker startup skipped: %s",exc)
try:
from taskiq_management.cv_broker_setup import cv_broker as _cv_broker
cv_broker=_cv_broker
await cv_broker.startup()
cv_broker_ready=True
except Exception as exc:
logger.warning("taskiq cv broker startup skipped: %s",exc)
try: try:
from llm_setup import init_llm,close_llm as _close_llm from llm_setup import init_llm,close_llm as _close_llm
from agent.agent_setup import init_agent,close_agent as _close_agent from agent.agent_setup import init_agent,close_agent as _close_agent
@ -49,6 +58,8 @@ async def lifespan(app):
await close_agent() await close_agent()
if llm_ready and close_llm is not None: if llm_ready and close_llm is not None:
await close_llm() await close_llm()
if cv_broker_ready and cv_broker is not None:
await cv_broker.shutdown()
if broker_ready and broker is not None: if broker_ready and broker is not None:
await broker.shutdown() await broker.shutdown()

View File

@ -55,6 +55,15 @@ class Users(SQLModel, table=True):
is_active: bool = Field(default=False) is_active: bool = Field(default=False)
is_deleted: bool = Field(default=False) is_deleted: bool = Field(default=False)
@classmethod
async def get_user_id(cls, session: AsyncSession, user_id: str):
uid = cls._as_uuid(user_id)
if uid is None:
return None
statement = select(cls).where(cls.id == uid)
result = await session.execute(statement)
return result.scalars().first()
@classmethod @classmethod
def _search_filter(cls, search: str): def _search_filter(cls, search: str):
pattern = f"%{search}%" pattern = f"%{search}%"

View File

@ -72,5 +72,66 @@ services:
condition: service_healthy condition: service_healthy
restart: unless-stopped restart: unless-stopped
taskiq-cv-worker:
build:
context: ./backend
container_name: hrms-taskiq-cv-worker
working_dir: /app
command:
[
"taskiq",
"worker",
"taskiq_management.cv_broker_setup:cv_broker",
"inbox.cv_tasks",
"--workers",
"1",
]
env_file:
- ./backend/.env
environment:
PYTHONPATH: /app
REDIS_URL: redis://redis:6379/0
TASKIQ_CV_QUEUE_NAME: cv_upload
TASKIQ_WORKER_NAME: cv-worker-01
DB_HOST: host.docker.internal
EMAIL_URL: http://host.docker.internal:5000
BACKEND_URL: http://host.docker.internal:8000
extra_hosts:
- "host.docker.internal:host-gateway"
volumes:
- ./backend/inbox/decoded_attachments:/app/inbox/decoded_attachments
depends_on:
redis:
condition: service_healthy
restart: unless-stopped
taskiq-cv-scheduler:
build:
context: ./backend
container_name: hrms-taskiq-cv-scheduler
working_dir: /app
command:
[
"taskiq",
"scheduler",
"taskiq_management.cv_broker_setup:cv_scheduler",
"inbox.cv_tasks",
]
env_file:
- ./backend/.env
environment:
PYTHONPATH: /app
REDIS_URL: redis://redis:6379/0
TASKIQ_CV_QUEUE_NAME: cv_upload
DB_HOST: host.docker.internal
EMAIL_URL: http://host.docker.internal:5000
BACKEND_URL: http://host.docker.internal:8000
extra_hosts:
- "host.docker.internal:host-gateway"
depends_on:
redis:
condition: service_healthy
restart: unless-stopped
volumes: volumes:
redis-data: redis-data: