From 4f4e7513f445b30cc02fc4b8a659b9ce764f8679 Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 15:01:27 +0800 Subject: [PATCH 1/8] feat: async resume extraction and crm update flow --- .env.example | 7 +- DEVELOPMENT.md | 1 + README.md | 9 +- .../src/five08/discord_bot/cogs/crm.py | 556 +++++++++++++++--- .../src/five08/discord_bot/config.py | 1 + apps/worker/src/five08/worker/actors.py | 18 +- apps/worker/src/five08/worker/api.py | 161 ++++- apps/worker/src/five08/worker/config.py | 1 + .../worker/crm/resume_profile_processor.py | 462 +++++++++++++++ .../src/five08/worker/crm/skills_extractor.py | 2 +- apps/worker/src/five08/worker/jobs.py | 37 ++ apps/worker/src/five08/worker/models.py | 54 ++ packages/shared/src/five08/queue.py | 19 +- packages/shared/src/five08/settings.py | 6 +- tests/unit/test_resume_profile_processor.py | 79 +++ tests/unit/test_skills_extractor.py | 14 + tests/unit/test_worker_api.py | 97 ++- 17 files changed, 1413 insertions(+), 111 deletions(-) create mode 100644 apps/worker/src/five08/worker/crm/resume_profile_processor.py create mode 100644 tests/unit/test_resume_profile_processor.py create mode 100644 tests/unit/test_skills_extractor.py diff --git a/.env.example b/.env.example index 5284fe9a..bdc21b44 100644 --- a/.env.example +++ b/.env.example @@ -32,21 +32,23 @@ MINIO_HOST_BIND=127.0.0.1 MINIO_API_PORT=9000 MINIO_CONSOLE_PORT=9001 -# API (optional defaults; WEBHOOK_SHARED_SECRET required to accept ingest calls) +# API (optional defaults; API_SHARED_SECRET required to accept ingest calls) WEBHOOK_INGEST_HOST=0.0.0.0 WEBHOOK_INGEST_PORT=8090 # Required: ingest requests are rejected when unset -WEBHOOK_SHARED_SECRET= +API_SHARED_SECRET= # Worker / consumer (optional defaults) WORKER_NAME=integrations-worker WORKER_QUEUE_NAMES=jobs.default WORKER_BURST=false +CRM_LINKEDIN_FIELD=cLinkedInUrl MAX_ATTACHMENTS_PER_CONTACT=3 MAX_FILE_SIZE_MB=10 ALLOWED_FILE_TYPES=pdf,doc,docx,txt RESUME_KEYWORDS=resume,cv,curriculum OPENAI_API_KEY= +# For OpenRouter, set OPENAI_BASE_URL=https://openrouter.ai/api/v1 OPENAI_BASE_URL= OPENAI_MODEL=gpt-4o-mini @@ -55,6 +57,7 @@ DISCORD_BOT_TOKEN=your_bot_token_here HEALTHCHECK_PORT=3000 DISCORD_SENDMSG_CHARACTER_LIMIT=2000 CHECK_EMAIL_WAIT=2 +WORKER_API_BASE_URL=http://worker-api:8090 # Required for Discord bot commands that post to channels CHANNEL_ID=1391742724666822798 diff --git a/DEVELOPMENT.md b/DEVELOPMENT.md index c005b454..88126f4d 100644 --- a/DEVELOPMENT.md +++ b/DEVELOPMENT.md @@ -118,6 +118,7 @@ Use `.env.example` as source of truth. Key categories: - Bot credentials/integrations: Discord, email, Espo, Kimai - Worker controls: `WORKER_NAME`, `WORKER_QUEUE_NAMES`, `WORKER_BURST` - Worker CRM processing: `MAX_ATTACHMENTS_PER_CONTACT`, `MAX_FILE_SIZE_MB`, `ALLOWED_FILE_TYPES`, `RESUME_KEYWORDS`, `OPENAI_API_KEY`, `OPENAI_BASE_URL`, `OPENAI_MODEL` +- Resume upload UX wiring: `WORKER_API_BASE_URL` on bot, `CRM_LINKEDIN_FIELD` on worker. ## CI Notes diff --git a/README.md b/README.md index e7c137c3..454be3fa 100644 --- a/README.md +++ b/README.md @@ -45,6 +45,9 @@ Migrations: ### Worker API Endpoints - `GET /health`: Redis/Postgres/worker health check. +- `GET /jobs/{job_id}`: Fetch queued job status/result payload. +- `POST /jobs/resume-extract`: Enqueue resume profile extraction. +- `POST /jobs/resume-apply`: Enqueue confirmed CRM field apply. - `POST /webhooks/{source}`: Generic webhook enqueue endpoint. - `POST /webhooks/espocrm`: EspoCRM webhook endpoint (expects array payload). - `POST /process-contact/{contact_id}`: Manually enqueue one contact skills job. @@ -96,7 +99,7 @@ docker compose up --build - `ESPO_API_KEY` (required by both bot and worker) - `JOB_TIMEOUT_SECONDS` (default: `600`) - `JOB_RESULT_TTL_SECONDS` (default: `3600`) -- `WEBHOOK_SHARED_SECRET` (required; requests are rejected when unset) +- `API_SHARED_SECRET` (required; requests are rejected when unset) - `POSTGRES_URL` (default: `postgresql://postgres:postgres@postgres:5432/workflows`) - `POSTGRES_DB` (default: `workflows`) - `POSTGRES_USER` (default: `postgres`) @@ -120,6 +123,7 @@ docker compose up --build - `DISCORD_BOT_TOKEN` - `CHANNEL_ID` +- `WORKER_API_BASE_URL` (default: `http://worker-api:8090`) - `EMAIL_USERNAME` - `EMAIL_PASSWORD` - `IMAP_SERVER` @@ -137,8 +141,9 @@ docker compose up --build - `MAX_FILE_SIZE_MB` (default: `10`) - `ALLOWED_FILE_TYPES` (default: `pdf,doc,docx,txt`) - `RESUME_KEYWORDS` (default: `resume,cv,curriculum`) +- `CRM_LINKEDIN_FIELD` (default: `cLinkedInUrl`) - `OPENAI_API_KEY` (optional; if unset, heuristic extraction is used) -- `OPENAI_BASE_URL` (optional) +- `OPENAI_BASE_URL` (optional; set `https://openrouter.ai/api/v1` for OpenRouter) - `OPENAI_MODEL` (default: `gpt-4o-mini`) ## Commands diff --git a/apps/discord_bot/src/five08/discord_bot/cogs/crm.py b/apps/discord_bot/src/five08/discord_bot/cogs/crm.py index b321f9a4..42e6174f 100644 --- a/apps/discord_bot/src/five08/discord_bot/cogs/crm.py +++ b/apps/discord_bot/src/five08/discord_bot/cogs/crm.py @@ -5,12 +5,15 @@ It allows team members to quickly access CRM data without leaving Discord. """ -import logging +import asyncio import io +import logging from typing import Any -from discord.ext import commands -from discord import app_commands + +import aiohttp import discord +from discord import app_commands +from discord.ext import commands from five08.discord_bot.config import settings from five08.clients import espo @@ -321,6 +324,133 @@ async def cancel_upload( pass +class ResumeUpdateConfirmationView(discord.ui.View): + """Confirm extracted profile updates before writing to CRM.""" + + def __init__( + self, + *, + crm_cog: "CRMCog", + requester_id: int, + contact_id: str, + contact_name: str, + proposed_updates: dict[str, str], + link_discord: dict[str, str] | None = None, + ) -> None: + super().__init__(timeout=300) + self.crm_cog = crm_cog + self.requester_id = requester_id + self.contact_id = contact_id + self.contact_name = contact_name + self.proposed_updates = proposed_updates + self.link_discord = link_discord + + async def interaction_check(self, interaction: discord.Interaction) -> bool: + """Allow only the original requester to confirm/cancel.""" + if interaction.user.id != self.requester_id: + await interaction.response.send_message( + "โŒ Only the command requester can confirm these updates.", + ephemeral=True, + ) + return False + return True + + @discord.ui.button(label="Confirm Updates", style=discord.ButtonStyle.primary) + async def confirm_updates( + self, + interaction: discord.Interaction, + button: discord.ui.Button["ResumeUpdateConfirmationView"], + ) -> None: + """Apply confirmed updates through the worker.""" + await interaction.response.defer(ephemeral=True) + + try: + apply_job_id = await self.crm_cog._enqueue_resume_apply_job( + contact_id=self.contact_id, + updates=self.proposed_updates, + link_discord=self.link_discord, + ) + except Exception as exc: + logger.error("Failed to enqueue resume apply job: %s", exc) + await interaction.followup.send( + "โŒ Failed to enqueue CRM apply job. Please try again." + ) + return + + await interaction.followup.send( + "๐Ÿ› ๏ธ Applying confirmed updates to CRM...", + ephemeral=True, + ) + apply_result = await self.crm_cog._wait_for_worker_job_result(apply_job_id) + + if not apply_result: + await interaction.followup.send( + "โš ๏ธ Timed out waiting for apply job. Please check again shortly." + ) + return + + status = str(apply_result.get("status", "unknown")) + if status != "succeeded": + await interaction.followup.send( + f"โŒ Apply job failed (status: {status}). " + f"Error: {apply_result.get('last_error') or 'Unknown error'}" + ) + return + + result = apply_result.get("result") + updated_fields: list[str] = [] + if isinstance(result, dict): + raw_fields = result.get("updated_fields") + if isinstance(raw_fields, list): + updated_fields = [str(field) for field in raw_fields] + + embed = discord.Embed( + title="โœ… CRM Updated", + description=f"Applied updates for **{self.contact_name}**.", + color=0x00FF00, + ) + embed.add_field( + name="Updated Fields", + value=", ".join(updated_fields) if updated_fields else "No field changes", + inline=False, + ) + profile_url = f"{self.crm_cog.base_url}/#Contact/view/{self.contact_id}" + embed.add_field(name="๐Ÿ”— CRM Profile", value=f"[View in CRM]({profile_url})") + await interaction.followup.send(embed=embed) + + for item in self.children: + if isinstance(item, discord.ui.Button): + item.disabled = True + if interaction.message: + try: + await interaction.message.edit(view=self) + except discord.NotFound: + pass + except discord.HTTPException as exc: + logger.warning("Failed to update confirmation view: %s", exc) + + @discord.ui.button(label="Cancel", style=discord.ButtonStyle.secondary) + async def cancel_updates( + self, + interaction: discord.Interaction, + button: discord.ui.Button["ResumeUpdateConfirmationView"], + ) -> None: + """Cancel CRM updates after preview.""" + await interaction.response.send_message( + "No CRM profile updates were applied.", ephemeral=True + ) + for item in self.children: + if isinstance(item, discord.ui.Button): + item.disabled = True + if interaction.message: + try: + await interaction.message.edit(view=self) + except discord.NotFound: + pass + except discord.HTTPException as exc: + logger.warning("Failed to update confirmation view: %s", exc) + + class CRMCog(commands.Cog): """CRM integration cog for EspoCRM operations.""" @@ -332,6 +462,291 @@ def __init__(self, bot: commands.Bot) -> None: # Store base URL for profile links self.base_url = settings.espo_base_url.rstrip("/") + def _worker_headers(self) -> dict[str, str]: + """Build auth headers for internal worker API calls.""" + if not settings.api_shared_secret: + raise ValueError("API_SHARED_SECRET is required for worker API requests.") + return { + "X-API-Secret": settings.api_shared_secret, + "Content-Type": "application/json", + } + + def _worker_url(self, path: str) -> str: + return f"{settings.worker_api_base_url.rstrip('/')}{path}" + + async def _enqueue_resume_extract_job( + self, *, contact_id: str, attachment_id: str, filename: str + ) -> str: + payload = { + "contact_id": contact_id, + "attachment_id": attachment_id, + "filename": filename, + } + async with aiohttp.ClientSession() as session: + async with session.post( + self._worker_url("/jobs/resume-extract"), + headers=self._worker_headers(), + json=payload, + timeout=aiohttp.ClientTimeout(total=30), + ) as response: + data = await response.json() + if response.status != 202: + raise ValueError(f"Worker extract enqueue failed: {data}") + job_id = data.get("job_id") + if not isinstance(job_id, str) or not job_id: + raise ValueError("Missing worker extract job_id in response.") + return job_id + + async def _enqueue_resume_apply_job( + self, + *, + contact_id: str, + updates: dict[str, str], + link_discord: dict[str, str] | None = None, + ) -> str: + payload = { + "contact_id": contact_id, + "updates": updates, + "link_discord": link_discord, + } + async with aiohttp.ClientSession() as session: + async with session.post( + self._worker_url("/jobs/resume-apply"), + headers=self._worker_headers(), + json=payload, + timeout=aiohttp.ClientTimeout(total=30), + ) as response: + data = await response.json() + if response.status != 202: + raise ValueError(f"Worker apply enqueue failed: {data}") + job_id = data.get("job_id") + if not isinstance(job_id, str) or not job_id: + raise ValueError("Missing worker apply job_id in response.") + return job_id + + async def _get_worker_job_status(self, job_id: str) -> dict[str, Any]: + async with aiohttp.ClientSession() as session: + async with session.get( + self._worker_url(f"/jobs/{job_id}"), + headers=self._worker_headers(), + timeout=aiohttp.ClientTimeout(total=30), + ) as response: + data = await response.json() + if response.status != 200: + raise ValueError(f"Worker job status failed: {data}") + if not isinstance(data, dict): + raise ValueError("Worker job status response must be an object.") + return data + + async def _wait_for_worker_job_result( + self, job_id: str, *, timeout_seconds: int = 180, poll_seconds: float = 2.0 + ) -> dict[str, Any] | None: + """Poll worker job status until terminal or timeout.""" + terminal = {"succeeded", "dead", "canceled"} + max_attempts = max(1, int(timeout_seconds / poll_seconds)) + + for _ in range(max_attempts): + job = await self._get_worker_job_status(job_id) + status = str(job.get("status", "")) + if status in terminal: + return job + await asyncio.sleep(poll_seconds) + + return None + + def _build_resume_preview_embed( + self, + *, + contact_id: str, + contact_name: str, + result: dict[str, Any], + link_member: discord.Member | None, + ) -> tuple[discord.Embed, dict[str, str]]: + """Render worker extraction result as a Discord preview embed.""" + proposed_updates_raw = result.get("proposed_updates") + proposed_updates: dict[str, str] = {} + if isinstance(proposed_updates_raw, dict): + proposed_updates = { + str(field): str(value) + for field, value in proposed_updates_raw.items() + if value is not None and str(value).strip() + } + + changes = result.get("proposed_changes") + new_skills = result.get("new_skills") + skipped = result.get("skipped") + extracted_profile = result.get("extracted_profile") + + embed = discord.Embed( + title="๐Ÿงพ Resume Parsed", + description=f"Review extracted updates for **{contact_name}**.", + color=0x0099FF, + ) + + if isinstance(changes, list) and changes: + lines: list[str] = [] + for change in changes[:8]: + if not isinstance(change, dict): + continue + label = str(change.get("label", change.get("field", "Field"))) + current = str(change.get("current", "None")) + proposed = str(change.get("proposed", "")) + lines.append(f"**{label}**: `{current}` โ†’ `{proposed}`") + embed.add_field( + name="Proposed Changes", + value="\n".join(lines) if lines else "No changes", + inline=False, + ) + else: + embed.add_field( + name="Proposed Changes", + value="No CRM field updates were extracted.", + inline=False, + ) + + if isinstance(new_skills, list) and new_skills: + formatted_skills = ", ".join(str(skill) for skill in new_skills[:25]) + embed.add_field( + name="New Skills", + value=formatted_skills, + inline=False, + ) + + if isinstance(skipped, list) and skipped: + skip_lines: list[str] = [] + for item in skipped[:4]: + if not isinstance(item, dict): + continue + field = str(item.get("field", "field")) + reason = str(item.get("reason", "Skipped")) + value = str(item.get("value", "")) + skip_lines.append(f"`{field}`: `{value}` ({reason})") + if skip_lines: + embed.add_field( + name="Skipped", + value="\n".join(skip_lines), + inline=False, + ) + + if isinstance(extracted_profile, dict): + confidence = extracted_profile.get("confidence") + source = extracted_profile.get("source") + if confidence is not None or source: + embed.add_field( + name="Extraction", + value=f"Source: `{source or 'unknown'}` | Confidence: `{confidence}`", + inline=False, + ) + + if link_member: + embed.add_field( + name="Discord Link", + value=f"Will link contact to {link_member.mention}", + inline=False, + ) + + profile_url = f"{self.base_url}/#Contact/view/{contact_id}" + embed.add_field(name="๐Ÿ”— CRM Profile", value=f"[View in CRM]({profile_url})") + return embed, proposed_updates + + async def _run_resume_extract_and_preview( + self, + *, + interaction: discord.Interaction, + contact_id: str, + contact_name: str, + attachment_id: str, + filename: str, + link_member: discord.Member | None, + ) -> None: + """Kick off worker extraction and show confirmation preview.""" + try: + job_id = await self._enqueue_resume_extract_job( + contact_id=contact_id, + attachment_id=attachment_id, + filename=filename, + ) + except Exception as exc: + logger.error("Failed to enqueue resume extract job: %s", exc) + await interaction.followup.send( + "โš ๏ธ Resume uploaded, but extraction job could not be enqueued.", + ephemeral=True, + ) + return + + await interaction.followup.send( + "๐Ÿ“ฅ Resume uploaded. Extracting profile fields now...", + ephemeral=True, + ) + + try: + job = await self._wait_for_worker_job_result(job_id) + except Exception as exc: + logger.error("Worker polling failed for job_id=%s error=%s", job_id, exc) + await interaction.followup.send( + "โš ๏ธ Resume uploaded, but extraction polling failed.", + ephemeral=True, + ) + return + if not job: + await interaction.followup.send( + "โš ๏ธ Timed out waiting for extraction result. Try again in a moment.", + ephemeral=True, + ) + return + + status = str(job.get("status", "unknown")) + if status != "succeeded": + await interaction.followup.send( + f"โŒ Extraction job failed (status: {status}). " + f"Error: {job.get('last_error') or 'Unknown error'}", + ephemeral=True, + ) + return + + result = job.get("result") + if not isinstance(result, dict): + await interaction.followup.send( + "โŒ Extraction result was empty or malformed.", + ephemeral=True, + ) + return + + if not result.get("success", False): + await interaction.followup.send( + f"โŒ Resume extraction failed: {result.get('error') or 'Unknown error'}", + ephemeral=True, + ) + return + + embed, proposed_updates = self._build_resume_preview_embed( + contact_id=contact_id, + contact_name=contact_name, + result=result, + link_member=link_member, + ) + + if not proposed_updates and not link_member: + await interaction.followup.send(embed=embed, ephemeral=True) + return + + link_discord_payload: dict[str, str] | None = None + if link_member: + link_discord_payload = { + "user_id": str(link_member.id), + "username": str(link_member), + } + + view = ResumeUpdateConfirmationView( + crm_cog=self, + requester_id=interaction.user.id, + contact_id=contact_id, + contact_name=contact_name, + proposed_updates=proposed_updates, + link_discord=link_discord_payload, + ) + await interaction.followup.send(embed=embed, view=view, ephemeral=True) + async def _download_and_send_resume( self, interaction: discord.Interaction, contact_name: str, resume_id: str ) -> None: @@ -1244,12 +1659,13 @@ async def _update_contact_resume( @app_commands.command( name="upload-resume", - description="Upload a resume file to CRM (your own or for someone else if Steering Committee+)", + description="Upload resume, extract profile fields, and preview CRM updates", ) @app_commands.describe( - file="Resume file to upload (PDF, DOC, DOCX)", - search_term="Email, name, or contact ID to find contact (optional - if not provided, uploads to your own profile)", - overwrite="Whether to replace all existing resumes instead of adding (default: False)", + file="Resume file to upload (PDF, DOC, DOCX, TXT)", + search_term="Email, name, or contact ID (optional - defaults to your linked profile)", + overwrite="Replace existing resumes instead of appending", + link_user="Discord user to link to this CRM contact (optional, Steering Committee+ for others)", ) async def upload_resume( self, @@ -1257,17 +1673,20 @@ async def upload_resume( file: discord.Attachment, search_term: str | None = None, overwrite: bool = False, + link_user: discord.Member | None = None, ) -> None: - """Upload a resume file to the CRM. - - By default, new resumes are added alongside existing ones. - Use overwrite=True to replace all existing resumes with this one. - """ + """Upload resume and run worker extraction to preview CRM updates.""" try: await interaction.response.defer(ephemeral=True) + if not settings.api_shared_secret: + await interaction.followup.send( + "โŒ API_SHARED_SECRET is not configured for worker API access." + ) + return + # Validate file type - valid_extensions = {".pdf", ".doc", ".docx"} + valid_extensions = {".pdf", ".doc", ".docx", ".txt"} file_extension = ( "." + file.filename.split(".")[-1].lower() if "." in file.filename @@ -1276,7 +1695,7 @@ async def upload_resume( if file_extension not in valid_extensions: await interaction.followup.send( - f"โŒ Invalid file type. Please upload a PDF, DOC, or DOCX file.\nYou uploaded: `{file.filename}`" + f"โŒ Invalid file type. Please upload a PDF, DOC, DOCX, or TXT file.\nYou uploaded: `{file.filename}`" ) return @@ -1288,22 +1707,24 @@ async def upload_resume( ) return + requires_steering = bool(search_term) or ( + link_user is not None and link_user.id != interaction.user.id + ) + if requires_steering and ( + not hasattr(interaction.user, "roles") + or not check_user_roles_with_hierarchy( + interaction.user.roles, ["Steering Committee"] + ) + ): + await interaction.followup.send( + "โŒ You must have Steering Committee role or higher for this upload." + ) + return + # Determine target contact target_contact = None if search_term: - # Uploading for someone else - requires Steering Committee+ role - if not hasattr( - interaction.user, "roles" - ) or not check_user_roles_with_hierarchy( - interaction.user.roles, ["Steering Committee"] - ): - await interaction.followup.send( - "โŒ You must have Steering Committee role or higher to upload resumes for other people." - ) - return - - # Search for target contact contacts = await self._search_contact_for_linking(search_term) if not contacts: await interaction.followup.send( @@ -1336,40 +1757,6 @@ async def upload_resume( await interaction.followup.send("โŒ Contact ID not found.") return - # Check for existing resume with same name and size - has_duplicate, existing_resume_id = await self._check_existing_resume( - contact_id, file.filename, file.size - ) - - if has_duplicate: - # Show confirmation dialog - embed = discord.Embed( - title="โš ๏ธ Duplicate Resume Detected", - description=f"A resume with the same name and file size already exists for **{contact_name}**.", - color=0xFFA500, - ) - embed.add_field(name="๐Ÿ“„ File", value=file.filename, inline=True) - embed.add_field( - name="๐Ÿ“ Size", value=f"{file.size / 1024:.1f} KB", inline=True - ) - embed.add_field( - name="โ“ Question", - value="Do you want to upload this resume anyway? This will add it as a new resume without replacing the existing one.", - inline=False, - ) - - view = ResumeConfirmationView( - self, - interaction, - file, - contact_id, - contact_name, - existing_resume_id or "", - overwrite, - ) - await interaction.followup.send(embed=embed, view=view) - return - # Download file content from Discord file_content = await file.read() @@ -1388,40 +1775,29 @@ async def upload_resume( await interaction.followup.send("โŒ Failed to upload file to CRM.") return - # Update contact's resume field - if await self._update_contact_resume( + if not await self._update_contact_resume( contact_id, attachment_id, overwrite ): - # Create success embed - embed = discord.Embed( - title="โœ… Resume Uploaded Successfully", - description="Resume has been uploaded and linked to the contact.", - color=0x00FF00, - ) - embed.add_field(name="๐Ÿ‘ค Contact", value=contact_name, inline=True) - embed.add_field(name="๐Ÿ“„ File", value=file.filename, inline=True) - embed.add_field( - name="๐Ÿ“ Size", value=f"{file.size / 1024:.1f} KB", inline=True - ) - - # Add CRM link - profile_url = f"{self.base_url}/#Contact/view/{contact_id}" - embed.add_field( - name="๐Ÿ”— CRM Profile", - value=f"[View in CRM]({profile_url})", - inline=False, - ) - - await interaction.followup.send(embed=embed) - - logger.info( - f"Resume uploaded for {contact_name} (ID: {contact_id}) " - f"by {interaction.user.name}: {file.filename}" - ) - else: await interaction.followup.send( - "โš ๏ธ File uploaded but failed to link to contact. Please check CRM manually." + "โš ๏ธ File uploaded, but failed to link in contact resume field." ) + return + + logger.info( + "Resume uploaded for %s (contact_id=%s, attachment_id=%s) by %s", + contact_name, + contact_id, + attachment_id, + interaction.user.name, + ) + await self._run_resume_extract_and_preview( + interaction=interaction, + contact_id=contact_id, + contact_name=contact_name, + attachment_id=attachment_id, + filename=file.filename, + link_member=link_user, + ) except EspoAPIError as e: logger.error(f"Failed to upload file to EspoCRM: {e}") diff --git a/apps/discord_bot/src/five08/discord_bot/config.py b/apps/discord_bot/src/five08/discord_bot/config.py index 43a748d4..97d5ca98 100644 --- a/apps/discord_bot/src/five08/discord_bot/config.py +++ b/apps/discord_bot/src/five08/discord_bot/config.py @@ -34,6 +34,7 @@ class Settings(SharedSettings): # CRM/EspoCRM settings espo_api_key: str espo_base_url: str + worker_api_base_url: str = "http://worker-api:8090" # Kimai time tracking settings kimai_base_url: str diff --git a/apps/worker/src/five08/worker/actors.py b/apps/worker/src/five08/worker/actors.py index 6aa029a5..ceb2adc6 100644 --- a/apps/worker/src/five08/worker/actors.py +++ b/apps/worker/src/five08/worker/actors.py @@ -21,7 +21,12 @@ ) from five08.queue import parse_queue_names from five08.worker.config import settings -from five08.worker.jobs import process_contact_skills_job, process_webhook_event +from five08.worker.jobs import ( + apply_resume_profile_job, + extract_resume_profile_job, + process_contact_skills_job, + process_webhook_event, +) from five08.logging import configure_logging @@ -37,6 +42,8 @@ _HANDLERS: dict[str, Any] = { process_webhook_event.__name__: process_webhook_event, process_contact_skills_job.__name__: process_contact_skills_job, + extract_resume_profile_job.__name__: extract_resume_profile_job, + apply_resume_profile_job.__name__: apply_resume_profile_job, } @@ -97,8 +104,13 @@ def _run_job(job_id: str) -> None: try: args, kwargs = _extract_call_args(job) - handler(*args, **kwargs) - mark_job_succeeded(settings, job_id) + result = handler(*args, **kwargs) + mark_job_succeeded( + settings, + job_id, + result=result, + base_payload=job.payload, + ) logger.info("Completed job_id=%s type=%s", job_id, job.type) except Exception as exc: next_attempt = job.attempts + 1 diff --git a/apps/worker/src/five08/worker/api.py b/apps/worker/src/five08/worker/api.py index 1541a9d7..74833b7d 100644 --- a/apps/worker/src/five08/worker/api.py +++ b/apps/worker/src/five08/worker/api.py @@ -4,9 +4,10 @@ import logging import secrets from datetime import datetime, timezone +from typing import Any from aiohttp import web -from pydantic import ValidationError +from pydantic import BaseModel, ValidationError from redis import Redis from five08.logging import configure_logging @@ -14,13 +15,19 @@ EnqueuedJob, QueueClient, enqueue_job, + get_job, get_redis_connection, is_postgres_healthy, ) from five08.worker.config import settings from five08.worker.db_migrations import run_job_migrations from five08.worker.dispatcher import build_queue_client -from five08.worker.jobs import process_contact_skills_job, process_webhook_event +from five08.worker.jobs import ( + apply_resume_profile_job, + extract_resume_profile_job, + process_contact_skills_job, + process_webhook_event, +) from five08.worker.models import EspoCRMWebhookPayload logger = logging.getLogger(__name__) @@ -28,16 +35,30 @@ QUEUE_KEY = web.AppKey("queue", QueueClient) +class ResumeExtractRequest(BaseModel): + """Request schema for queued resume extraction.""" + + contact_id: str + attachment_id: str + filename: str + + +class ResumeApplyRequest(BaseModel): + """Request schema for queued resume apply updates.""" + + contact_id: str + updates: dict[str, str] + link_discord: dict[str, str] | None = None + + def _is_authorized(request: web.Request) -> bool: - """Validate webhook secret.""" - if not settings.webhook_shared_secret: - logger.error( - "Rejecting webhook request: WEBHOOK_SHARED_SECRET is not configured" - ) + """Validate shared API secret.""" + if not settings.api_shared_secret: + logger.error("Rejecting request: API_SHARED_SECRET is not configured") return False - provided_secret = request.headers.get("X-Webhook-Secret", "") - return secrets.compare_digest(provided_secret, settings.webhook_shared_secret) + provided_secret = request.headers.get("X-API-Secret", "") + return secrets.compare_digest(provided_secret, settings.api_shared_secret) def _extract_idempotency_key(value: object) -> str | None: @@ -190,6 +211,125 @@ async def process_contact_handler(request: web.Request) -> web.Response: ) +async def resume_extract_handler(request: web.Request) -> web.Response: + """Enqueue resume extraction job for one uploaded attachment.""" + if not _is_authorized(request): + return web.json_response({"error": "unauthorized"}, status=401) + + try: + payload_data = await request.json() + except Exception: + return web.json_response({"error": "invalid_json"}, status=400) + + try: + payload = ResumeExtractRequest.model_validate(payload_data) + except ValidationError as exc: + return web.json_response( + {"error": "invalid_resume_extract_payload", "detail": str(exc)}, status=400 + ) + + queue = request.app[QUEUE_KEY] + job = await asyncio.to_thread( + enqueue_job, + queue=queue, + fn=extract_resume_profile_job, + args=(payload.contact_id, payload.attachment_id, payload.filename), + settings=settings, + idempotency_key=f"resume-extract:{payload.contact_id}:{payload.attachment_id}", + ) + logger.info( + "Enqueued resume extract job contact_id=%s attachment_id=%s job_id=%s created=%s", + payload.contact_id, + payload.attachment_id, + job.id, + job.created, + ) + return web.json_response( + { + "status": "queued", + "job_id": job.id, + "contact_id": payload.contact_id, + "attachment_id": payload.attachment_id, + "created": job.created, + }, + status=202, + ) + + +async def resume_apply_handler(request: web.Request) -> web.Response: + """Enqueue CRM apply job after user confirmation in Discord.""" + if not _is_authorized(request): + return web.json_response({"error": "unauthorized"}, status=401) + + try: + payload_data = await request.json() + except Exception: + return web.json_response({"error": "invalid_json"}, status=400) + + try: + payload = ResumeApplyRequest.model_validate(payload_data) + except ValidationError as exc: + return web.json_response( + {"error": "invalid_resume_apply_payload", "detail": str(exc)}, status=400 + ) + + queue = request.app[QUEUE_KEY] + manual_nonce = datetime.now(tz=timezone.utc).isoformat() + job = await asyncio.to_thread( + enqueue_job, + queue=queue, + fn=apply_resume_profile_job, + args=(payload.contact_id, payload.updates, payload.link_discord), + settings=settings, + idempotency_key=f"resume-apply:{payload.contact_id}:{manual_nonce}", + ) + logger.info( + "Enqueued resume apply job contact_id=%s job_id=%s created=%s", + payload.contact_id, + job.id, + job.created, + ) + return web.json_response( + { + "status": "queued", + "job_id": job.id, + "contact_id": payload.contact_id, + }, + status=202, + ) + + +async def job_status_handler(request: web.Request) -> web.Response: + """Return persisted status and worker result payload for one job.""" + if not _is_authorized(request): + return web.json_response({"error": "unauthorized"}, status=401) + + job_id = request.match_info.get("job_id", "").strip() + if not job_id: + return web.json_response({"error": "job_id_required"}, status=400) + + job = await asyncio.to_thread(get_job, settings, job_id) + if job is None: + return web.json_response({"error": "job_not_found"}, status=404) + + result: Any = None + payload = job.payload if isinstance(job.payload, dict) else {} + if "result" in payload: + result = payload["result"] + + return web.json_response( + { + "job_id": job.id, + "type": job.type, + "status": job.status.value, + "attempts": job.attempts, + "max_attempts": job.max_attempts, + "last_error": job.last_error, + "result": result, + } + ) + + async def on_startup(app: web.Application) -> None: """Initialize queue dependencies.""" await asyncio.to_thread(run_job_migrations) @@ -211,6 +351,9 @@ def create_app() -> web.Application: app.on_cleanup.append(on_cleanup) app.router.add_get("/", health_handler) app.router.add_get("/health", health_handler) + app.router.add_get("/jobs/{job_id}", job_status_handler) + app.router.add_post("/jobs/resume-extract", resume_extract_handler) + app.router.add_post("/jobs/resume-apply", resume_apply_handler) app.router.add_post("/webhooks/espocrm", espocrm_webhook_handler) app.router.add_post("/webhooks/{source}", ingest_handler) app.router.add_post("/process-contact/{contact_id}", process_contact_handler) diff --git a/apps/worker/src/five08/worker/config.py b/apps/worker/src/five08/worker/config.py index 5e8d7b00..07784f54 100644 --- a/apps/worker/src/five08/worker/config.py +++ b/apps/worker/src/five08/worker/config.py @@ -12,6 +12,7 @@ class WorkerSettings(SharedSettings): espo_base_url: str espo_api_key: str + crm_linkedin_field: str = "cLinkedInUrl" openai_api_key: str | None = None openai_base_url: str | None = None diff --git a/apps/worker/src/five08/worker/crm/resume_profile_processor.py b/apps/worker/src/five08/worker/crm/resume_profile_processor.py new file mode 100644 index 00000000..f58b6964 --- /dev/null +++ b/apps/worker/src/five08/worker/crm/resume_profile_processor.py @@ -0,0 +1,462 @@ +"""Resume extraction + CRM profile update workflow.""" + +from __future__ import annotations + +import json +import logging +import re +from collections.abc import Callable +from typing import Any + +from five08.clients.espo import EspoAPI, EspoAPIError +from five08.worker.config import settings +from five08.worker.crm.document_processor import DocumentProcessor +from five08.worker.crm.skills_extractor import SkillsExtractor +from five08.worker.models import ( + ResumeApplyResult, + ResumeExtractedProfile, + ResumeExtractionResult, + ResumeFieldChange, + ResumeSkipReason, +) + +logger = logging.getLogger(__name__) + +try: # pragma: no cover - import success depends on environment + from openai import OpenAI as OpenAIClient +except Exception: # pragma: no cover + OpenAIClient = None # type: ignore[misc,assignment] + + +class ResumeEspoClient: + """Minimal EspoCRM client wrapper for resume profile flows.""" + + def __init__(self) -> None: + api_url = settings.espo_base_url.rstrip("/") + "/api/v1" + self.api = EspoAPI(api_url, settings.espo_api_key) + + def get_contact(self, contact_id: str) -> dict[str, Any]: + return self.api.request("GET", f"Contact/{contact_id}") + + def download_attachment(self, attachment_id: str) -> bytes: + return self.api.download_file(f"Attachment/file/{attachment_id}") + + def update_contact(self, contact_id: str, updates: dict[str, Any]) -> None: + self.api.request("PUT", f"Contact/{contact_id}", updates) + + +class ResumeProfileExtractor: + """Extract candidate profile fields from resume text.""" + + def __init__(self) -> None: + self.model = settings.openai_model + self.client: Any = None + + if settings.openai_api_key and OpenAIClient is not None: + self.client = OpenAIClient( + api_key=settings.openai_api_key, + base_url=settings.openai_base_url, + ) + + def extract(self, resume_text: str) -> ResumeExtractedProfile: + """Return extracted fields from resume text.""" + if self.client is None: + return self._heuristic_extract(resume_text) + + try: + response = self.client.chat.completions.create( + model=self.model, + messages=[ + { + "role": "system", + "content": ( + "Extract contact fields from resume text. Return JSON only." + ), + }, + { + "role": "user", + "content": self._build_prompt(resume_text), + }, + ], + temperature=0.1, + max_tokens=800, + ) + raw_content = response.choices[0].message.content + if not raw_content: + raise ValueError("LLM returned empty content") + + parsed = self._parse_json(raw_content) + return ResumeExtractedProfile( + email=self._normalize_email(parsed.get("email")), + github_username=self._normalize_github(parsed.get("github_username")), + linkedin_url=self._normalize_linkedin(parsed.get("linkedin_url")), + phone=self._normalize_phone(parsed.get("phone")), + confidence=self._bounded_confidence(parsed.get("confidence", 0.75)), + source=self.model, + ) + except Exception as exc: + logger.warning("LLM resume extraction failed, using fallback: %s", exc) + return self._heuristic_extract(resume_text) + + def _heuristic_extract(self, resume_text: str) -> ResumeExtractedProfile: + email_match = re.search( + r"\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,}\b", resume_text + ) + github_match = re.search( + r"(?:https?://)?(?:www\.)?github\.com/([A-Za-z0-9-]{1,39})", + resume_text, + flags=re.IGNORECASE, + ) + linkedin_match = re.search( + r"(?:https?://)?(?:[\w.-]+\.)?linkedin\.com/in/[A-Za-z0-9\-_%]+/?", + resume_text, + flags=re.IGNORECASE, + ) + phone_match = re.search( + r"(?:\+?\d[\d\s().-]{7,}\d)", + resume_text, + ) + + github_value: str | None = None + if github_match: + github_value = github_match.group(1) + + linkedin_value: str | None = None + if linkedin_match: + linkedin_value = linkedin_match.group(0) + + phone_value: str | None = None + if phone_match: + phone_value = phone_match.group(0) + + return ResumeExtractedProfile( + email=self._normalize_email(email_match.group(0) if email_match else None), + github_username=self._normalize_github(github_value), + linkedin_url=self._normalize_linkedin(linkedin_value), + phone=self._normalize_phone(phone_value), + confidence=0.45, + source="heuristic", + ) + + def _build_prompt(self, resume_text: str) -> str: + snippet = resume_text[:12000] + return ( + "Extract contact fields from this resume.\\n" + "Return JSON with exact keys:\\n" + '{"email": string|null, "github_username": string|null, ' + '"linkedin_url": string|null, "phone": string|null, ' + '"confidence": number}\\n\\n' + f"Resume:\\n{snippet}" + ) + + def _parse_json(self, content: str) -> dict[str, Any]: + raw = content.strip() + if raw.startswith("```"): + lines = [line for line in raw.splitlines() if not line.startswith("```")] + raw = "\n".join(lines).strip() + + parsed = json.loads(raw) + if not isinstance(parsed, dict): + raise ValueError("Model output was not a JSON object") + return parsed + + def _bounded_confidence(self, value: Any) -> float: + try: + numeric = float(value) + except Exception: + numeric = 0.0 + return max(0.0, min(1.0, numeric)) + + def _normalize_email(self, value: Any) -> str | None: + if not isinstance(value, str): + return None + normalized = value.strip().lower() + return normalized or None + + def _normalize_github(self, value: Any) -> str | None: + if not isinstance(value, str): + return None + candidate = value.strip() + if not candidate: + return None + + github_match = re.search( + r"(?:https?://)?(?:www\.)?github\.com/([A-Za-z0-9-]{1,39})", + candidate, + flags=re.IGNORECASE, + ) + if github_match: + candidate = github_match.group(1) + + candidate = candidate.lstrip("@").strip().strip("/") + return candidate or None + + def _normalize_linkedin(self, value: Any) -> str | None: + if not isinstance(value, str): + return None + candidate = value.strip() + if not candidate: + return None + + if "linkedin.com" not in candidate.lower(): + return None + if not candidate.lower().startswith(("http://", "https://")): + candidate = f"https://{candidate}" + return candidate.rstrip("/") + + def _normalize_phone(self, value: Any) -> str | None: + if not isinstance(value, str): + return None + candidate = value.strip() + if not candidate: + return None + + digits = re.sub(r"\D", "", candidate) + if len(digits) < 7: + return None + + if candidate.startswith("+"): + return "+" + digits + return digits + + +class ResumeProfileProcessor: + """End-to-end extraction and apply operations for uploaded resumes.""" + + def __init__(self) -> None: + self.crm = ResumeEspoClient() + self.extractor = ResumeProfileExtractor() + self.skills_extractor = SkillsExtractor() + self.document_processor = DocumentProcessor() + + def extract_profile_proposal( + self, + *, + contact_id: str, + attachment_id: str, + filename: str, + ) -> ResumeExtractionResult: + """Build preview proposal from an uploaded resume attachment.""" + try: + contact = self.crm.get_contact(contact_id) + content = self.crm.download_attachment(attachment_id) + text = self.document_processor.extract_text(content, filename) + extracted = self.extractor.extract(text) + extracted_skills_result = self.skills_extractor.extract_skills(text) + extracted_skills = extracted_skills_result.skills + existing_skills = self._parse_existing_skills(contact.get("skills")) + existing_lower = {item.casefold() for item in existing_skills} + new_skills = [ + skill + for skill in extracted_skills + if skill.casefold() not in existing_lower + ] + + proposed_updates: dict[str, str] = {} + proposed_changes: list[ResumeFieldChange] = [] + skipped: list[ResumeSkipReason] = [] + + self._collect_change( + crm_field="emailAddress", + label="Email", + current=contact.get("emailAddress"), + proposed=extracted.email, + proposed_updates=proposed_updates, + proposed_changes=proposed_changes, + skipped=skipped, + blocked_reason="Skipped because @508.dev emails are managed separately", + is_blocked=lambda value: value.lower().endswith("@508.dev"), + ) + self._collect_change( + crm_field="cGitHubUsername", + label="GitHub", + current=contact.get("cGitHubUsername"), + proposed=extracted.github_username, + proposed_updates=proposed_updates, + proposed_changes=proposed_changes, + skipped=skipped, + ) + self._collect_change( + crm_field=settings.crm_linkedin_field, + label="LinkedIn", + current=contact.get(settings.crm_linkedin_field), + proposed=extracted.linkedin_url, + proposed_updates=proposed_updates, + proposed_changes=proposed_changes, + skipped=skipped, + ) + self._collect_change( + crm_field="phoneNumber", + label="Phone", + current=contact.get("phoneNumber"), + proposed=extracted.phone, + proposed_updates=proposed_updates, + proposed_changes=proposed_changes, + skipped=skipped, + ) + if new_skills: + merged_skills = existing_skills + new_skills + proposed_updates["skills"] = ", ".join(merged_skills) + proposed_changes.append( + ResumeFieldChange( + field="skills", + label="Skills", + current=", ".join(existing_skills) if existing_skills else None, + proposed=", ".join(merged_skills), + reason=f"Added {len(new_skills)} skills from resume extraction", + ) + ) + + return ResumeExtractionResult( + contact_id=contact_id, + attachment_id=attachment_id, + proposed_updates=proposed_updates, + proposed_changes=proposed_changes, + skipped=skipped, + extracted_profile=extracted, + extracted_skills=extracted_skills, + new_skills=new_skills, + success=True, + ) + except Exception as exc: + logger.error( + "Resume extraction proposal failed contact_id=%s attachment_id=%s error=%s", + contact_id, + attachment_id, + exc, + ) + return ResumeExtractionResult( + contact_id=contact_id, + attachment_id=attachment_id, + proposed_updates={}, + proposed_changes=[], + skipped=[], + extracted_profile=ResumeExtractedProfile( + email=None, + github_username=None, + linkedin_url=None, + phone=None, + confidence=0.0, + source="error", + ), + extracted_skills=[], + new_skills=[], + success=False, + error=str(exc), + ) + + def apply_profile_updates( + self, + *, + contact_id: str, + updates: dict[str, str], + link_discord: dict[str, str] | None = None, + ) -> ResumeApplyResult: + """Apply confirmed updates to contact in CRM.""" + try: + allowed_fields = { + "emailAddress", + "cGitHubUsername", + settings.crm_linkedin_field, + "phoneNumber", + "skills", + } + sanitized_updates: dict[str, str] = { + field: value + for field, value in updates.items() + if field in allowed_fields and value + } + + email_value = sanitized_updates.get("emailAddress") + if email_value and email_value.lower().endswith("@508.dev"): + sanitized_updates.pop("emailAddress") + + link_applied = False + if link_discord: + discord_user_id = str(link_discord.get("user_id", "")).strip() + discord_username = str(link_discord.get("username", "")).strip() + if discord_user_id and discord_username: + sanitized_updates["cDiscordUserID"] = discord_user_id + sanitized_updates["cDiscordUsername"] = ( + f"{discord_username} (ID: {discord_user_id})" + ) + link_applied = True + + if sanitized_updates: + self.crm.update_contact(contact_id, sanitized_updates) + + return ResumeApplyResult( + contact_id=contact_id, + updated_fields=sorted(sanitized_updates.keys()), + link_discord_applied=link_applied, + success=True, + ) + except EspoAPIError as exc: + logger.error("EspoCRM apply failed contact_id=%s error=%s", contact_id, exc) + return ResumeApplyResult( + contact_id=contact_id, + updated_fields=[], + success=False, + error=str(exc), + ) + except Exception as exc: + logger.error( + "Unexpected apply error contact_id=%s error=%s", contact_id, exc + ) + return ResumeApplyResult( + contact_id=contact_id, + updated_fields=[], + success=False, + error=str(exc), + ) + + def _collect_change( + self, + *, + crm_field: str, + label: str, + current: Any, + proposed: str | None, + proposed_updates: dict[str, str], + proposed_changes: list[ResumeFieldChange], + skipped: list[ResumeSkipReason], + blocked_reason: str | None = None, + is_blocked: Callable[[str], bool] | None = None, + ) -> None: + if not proposed: + return + + if callable(is_blocked) and is_blocked(proposed): + skipped.append( + ResumeSkipReason( + field=crm_field, + value=proposed, + reason=blocked_reason or "Update blocked by policy", + ) + ) + return + + current_value = str(current).strip() if current is not None else None + if current_value and current_value == proposed: + return + + proposed_updates[crm_field] = proposed + proposed_changes.append( + ResumeFieldChange( + field=crm_field, + label=label, + current=current_value, + proposed=proposed, + reason="Extracted from uploaded resume", + ) + ) + + def _parse_existing_skills(self, value: Any) -> list[str]: + if value is None: + return [] + if isinstance(value, list): + normalized = [str(item).strip() for item in value if str(item).strip()] + return normalized + + parsed = [item.strip() for item in str(value).split(",")] + return [item for item in parsed if item] diff --git a/apps/worker/src/five08/worker/crm/skills_extractor.py b/apps/worker/src/five08/worker/crm/skills_extractor.py index 840834f4..3cf6e670 100644 --- a/apps/worker/src/five08/worker/crm/skills_extractor.py +++ b/apps/worker/src/five08/worker/crm/skills_extractor.py @@ -98,7 +98,7 @@ def extract_skills(self, resume_text: str) -> ExtractedSkills: def _extract_skills_heuristic(self, resume_text: str) -> ExtractedSkills: """Simple keyword and token-based extraction fallback.""" lowered = resume_text.lower() - token_matches = re.findall(r"\b[a-z][a-z0-9+#\-.]{2,24}\b", lowered) + token_matches = re.findall(r"\b[a-z][a-z0-9+#\-.]{1,24}\b", lowered) detected: set[str] = set() for token in token_matches: if token in COMMON_SKILLS: diff --git a/apps/worker/src/five08/worker/jobs.py b/apps/worker/src/five08/worker/jobs.py index 532cf72c..8c6719d5 100644 --- a/apps/worker/src/five08/worker/jobs.py +++ b/apps/worker/src/five08/worker/jobs.py @@ -5,6 +5,7 @@ from typing import Any from five08.worker.crm.processor import ContactSkillsProcessor +from five08.worker.crm.resume_profile_processor import ResumeProfileProcessor logger = logging.getLogger(__name__) @@ -28,3 +29,39 @@ def process_webhook_event(source: str, payload: dict[str, Any]) -> dict[str, Any "received_at": received_at, "payload_keys": sorted(payload.keys()), } + + +def extract_resume_profile_job( + contact_id: str, + attachment_id: str, + filename: str, +) -> dict[str, Any]: + """Extract profile updates from an uploaded resume attachment.""" + logger.info( + "Processing resume extract job contact_id=%s attachment_id=%s", + contact_id, + attachment_id, + ) + processor = ResumeProfileProcessor() + result = processor.extract_profile_proposal( + contact_id=contact_id, + attachment_id=attachment_id, + filename=filename, + ) + return result.model_dump() + + +def apply_resume_profile_job( + contact_id: str, + updates: dict[str, str], + link_discord: dict[str, str] | None = None, +) -> dict[str, Any]: + """Apply confirmed CRM profile updates after bot-side confirmation.""" + logger.info("Processing resume apply job contact_id=%s", contact_id) + processor = ResumeProfileProcessor() + result = processor.apply_profile_updates( + contact_id=contact_id, + updates=updates, + link_discord=link_discord, + ) + return result.model_dump() diff --git a/apps/worker/src/five08/worker/models.py b/apps/worker/src/five08/worker/models.py index 93d3442e..40415be3 100644 --- a/apps/worker/src/five08/worker/models.py +++ b/apps/worker/src/five08/worker/models.py @@ -53,3 +53,57 @@ class SkillsExtractionResult(BaseModel): updated_skills: list[str] success: bool error: str | None = None + + +class ResumeExtractedProfile(BaseModel): + """Normalized profile fields extracted from resume text.""" + + email: str | None = None + github_username: str | None = None + linkedin_url: str | None = None + phone: str | None = None + confidence: float = Field(..., ge=0.0, le=1.0) + source: str + + +class ResumeFieldChange(BaseModel): + """Single proposed CRM field update.""" + + field: str + label: str + current: str | None = None + proposed: str + reason: str + + +class ResumeSkipReason(BaseModel): + """Field extraction skip explanation for preview UX.""" + + field: str + value: str + reason: str + + +class ResumeExtractionResult(BaseModel): + """Worker output used by bot preview/confirmation flow.""" + + contact_id: str + attachment_id: str + proposed_updates: dict[str, str] + proposed_changes: list[ResumeFieldChange] + skipped: list[ResumeSkipReason] + extracted_profile: ResumeExtractedProfile + extracted_skills: list[str] = Field(default_factory=list) + new_skills: list[str] = Field(default_factory=list) + success: bool + error: str | None = None + + +class ResumeApplyResult(BaseModel): + """CRM apply-phase result.""" + + contact_id: str + updated_fields: list[str] + link_discord_applied: bool = False + success: bool + error: str | None = None diff --git a/packages/shared/src/five08/queue.py b/packages/shared/src/five08/queue.py index 81fe4c1c..ee147b5a 100644 --- a/packages/shared/src/five08/queue.py +++ b/packages/shared/src/five08/queue.py @@ -209,6 +209,7 @@ def _mark_job( *, status: JobStatus | None = None, attempts: int | None = None, + payload: Any = _UNSET, locked_at: Any = _UNSET, locked_by: Any = _UNSET, run_after: Any = _UNSET, @@ -223,6 +224,9 @@ def _mark_job( if attempts is not None: updates.append("attempts = %s") params.append(attempts) + if payload is not _UNSET: + updates.append("payload = %s") + params.append(Jsonb(payload)) if locked_at is not _UNSET: updates.append("locked_at = %s") params.append(locked_at) @@ -266,12 +270,25 @@ def mark_job_running( ) -def mark_job_succeeded(settings: SharedSettings, job_id: str) -> None: +def mark_job_succeeded( + settings: SharedSettings, + job_id: str, + *, + result: Any | None = None, + base_payload: dict[str, Any] | None = None, +) -> None: """Mark successful completion.""" + payload: Any = _UNSET + if result is not None: + merged_payload = dict(base_payload or {}) + merged_payload["result"] = result + payload = merged_payload + _mark_job( settings, job_id, status=JobStatus.SUCCEEDED, + payload=payload, locked_at=None, locked_by=None, run_after=None, diff --git a/packages/shared/src/five08/settings.py b/packages/shared/src/five08/settings.py index 721e0031..9b47ce5c 100644 --- a/packages/shared/src/five08/settings.py +++ b/packages/shared/src/five08/settings.py @@ -1,5 +1,6 @@ """Shared configuration settings across services.""" +from pydantic import AliasChoices, Field from pydantic_settings import BaseSettings, SettingsConfigDict @@ -33,7 +34,10 @@ class SharedSettings(BaseSettings): webhook_ingest_host: str = "0.0.0.0" webhook_ingest_port: int = 8090 - webhook_shared_secret: str | None = None + api_shared_secret: str | None = Field( + default=None, + validation_alias=AliasChoices("API_SHARED_SECRET", "WEBHOOK_SHARED_SECRET"), + ) model_config = SettingsConfigDict(env_file=".env", extra="ignore") diff --git a/tests/unit/test_resume_profile_processor.py b/tests/unit/test_resume_profile_processor.py new file mode 100644 index 00000000..75c191f3 --- /dev/null +++ b/tests/unit/test_resume_profile_processor.py @@ -0,0 +1,79 @@ +"""Unit tests for resume profile worker processor.""" + +from unittest.mock import Mock + +from five08.worker.crm.resume_profile_processor import ResumeProfileProcessor +from five08.worker.models import ExtractedSkills, ResumeExtractedProfile + + +def test_extract_profile_proposal_filters_508_email() -> None: + """Extract proposal should skip @508.dev email updates by policy.""" + processor = ResumeProfileProcessor() + processor.crm = Mock() + processor.extractor = Mock() + processor.skills_extractor = Mock() + processor.document_processor = Mock() + + processor.crm.get_contact.return_value = { + "emailAddress": "member@example.com", + "cGitHubUsername": "old-gh", + "cLinkedInUrl": "https://linkedin.com/in/old", + "phoneNumber": "1234567890", + } + processor.crm.download_attachment.return_value = b"resume-bytes" + processor.document_processor.extract_text.return_value = "resume text" + processor.extractor.extract.return_value = ResumeExtractedProfile( + email="new@508.dev", + github_username="new-gh", + linkedin_url="https://linkedin.com/in/new", + phone="14155551234", + confidence=0.9, + source="gpt-4o-mini", + ) + processor.skills_extractor.extract_skills.return_value = ExtractedSkills( + skills=["Python", "FastAPI"], + confidence=0.8, + source="gpt-4o-mini", + ) + + result = processor.extract_profile_proposal( + contact_id="contact-1", + attachment_id="att-1", + filename="resume.pdf", + ) + + assert result.success is True + assert "emailAddress" not in result.proposed_updates + assert result.proposed_updates["cGitHubUsername"] == "new-gh" + assert result.proposed_updates["cLinkedInUrl"] == "https://linkedin.com/in/new" + assert result.proposed_updates["phoneNumber"] == "14155551234" + assert result.proposed_updates["skills"] == "Python, FastAPI" + assert result.new_skills == ["Python", "FastAPI"] + assert any(item.field == "emailAddress" for item in result.skipped) + + +def test_apply_profile_updates_adds_discord_and_filters_email() -> None: + """Apply should include Discord link values and prevent @508.dev email writes.""" + processor = ResumeProfileProcessor() + processor.crm = Mock() + + result = processor.apply_profile_updates( + contact_id="contact-1", + updates={ + "emailAddress": "member@508.dev", + "cGitHubUsername": "new-gh", + "phoneNumber": "14155551234", + "skills": "Python, FastAPI", + }, + link_discord={"user_id": "123", "username": "member#0001"}, + ) + + assert result.success is True + processor.crm.update_contact.assert_called_once() + update_payload = processor.crm.update_contact.call_args[0][1] + assert "emailAddress" not in update_payload + assert update_payload["cGitHubUsername"] == "new-gh" + assert update_payload["phoneNumber"] == "14155551234" + assert update_payload["skills"] == "Python, FastAPI" + assert update_payload["cDiscordUserID"] == "123" + assert update_payload["cDiscordUsername"] == "member#0001 (ID: 123)" diff --git a/tests/unit/test_skills_extractor.py b/tests/unit/test_skills_extractor.py new file mode 100644 index 00000000..1dada5cb --- /dev/null +++ b/tests/unit/test_skills_extractor.py @@ -0,0 +1,14 @@ +"""Unit tests for skills extraction fallback heuristics.""" + +from five08.worker.crm.skills_extractor import SkillsExtractor + + +def test_heuristic_extract_includes_two_letter_skill_go() -> None: + """Heuristic fallback should detect 2-letter skills in COMMON_SKILLS.""" + extractor = SkillsExtractor() + extractor.client = None + + result = extractor.extract_skills("Built distributed services in Go and Docker") + + assert "go" in result.skills + assert "docker" in result.skills diff --git a/tests/unit/test_worker_api.py b/tests/unit/test_worker_api.py index 2daa7f23..ae818355 100644 --- a/tests/unit/test_worker_api.py +++ b/tests/unit/test_worker_api.py @@ -23,8 +23,8 @@ def ping(self) -> bool: @pytest.fixture def auth_headers(monkeypatch: pytest.MonkeyPatch) -> dict[str, str]: """Configure webhook secret and return matching auth headers.""" - monkeypatch.setattr(api.settings, "webhook_shared_secret", "test-secret") - return {"X-Webhook-Secret": "test-secret"} + monkeypatch.setattr(api.settings, "api_shared_secret", "test-secret") + return {"X-API-Secret": "test-secret"} @pytest.mark.asyncio @@ -167,3 +167,96 @@ async def test_process_contact_handler_enqueues_single_contact( assert response.status == 202 assert payload["contact_id"] == "c-123" assert payload["job_id"] == "job-123" + + +@pytest.mark.asyncio +async def test_resume_extract_handler_enqueues_job( + auth_headers: dict[str, str], +) -> None: + """Resume extract endpoint should enqueue extraction job.""" + app_obj = web.Application() + app_obj[api.QUEUE_KEY] = Mock() + request = make_mocked_request( + "POST", "/jobs/resume-extract", app=app_obj, headers=auth_headers + ) + request.json = AsyncMock( + return_value={ + "contact_id": "c-1", + "attachment_id": "a-1", + "filename": "resume.pdf", + } + ) # type: ignore[method-assign] + + with patch("five08.worker.api.enqueue_job") as mock_enqueue: + mock_enqueue.return_value = Mock(id="job-extract", created=True) + response = await api.resume_extract_handler(request) + + payload = json.loads(response.text) + assert response.status == 202 + assert payload["job_id"] == "job-extract" + assert payload["contact_id"] == "c-1" + assert payload["attachment_id"] == "a-1" + + +@pytest.mark.asyncio +async def test_resume_apply_handler_enqueues_job( + auth_headers: dict[str, str], +) -> None: + """Resume apply endpoint should enqueue apply job.""" + app_obj = web.Application() + app_obj[api.QUEUE_KEY] = Mock() + request = make_mocked_request( + "POST", "/jobs/resume-apply", app=app_obj, headers=auth_headers + ) + request.json = AsyncMock( + return_value={ + "contact_id": "c-1", + "updates": {"emailAddress": "dev@example.com"}, + "link_discord": {"user_id": "123", "username": "dev#1111"}, + } + ) # type: ignore[method-assign] + + with patch("five08.worker.api.enqueue_job") as mock_enqueue: + mock_enqueue.return_value = Mock(id="job-apply", created=True) + response = await api.resume_apply_handler(request) + + payload = json.loads(response.text) + assert response.status == 202 + assert payload["job_id"] == "job-apply" + assert payload["contact_id"] == "c-1" + + +@pytest.mark.asyncio +async def test_job_status_handler_returns_result( + auth_headers: dict[str, str], +) -> None: + """Job status endpoint should expose persisted result payload.""" + app_obj = web.Application() + request = make_mocked_request( + "GET", + "/jobs/job-123", + app=app_obj, + headers=auth_headers, + match_info={"job_id": "job-123"}, + ) + + mock_status = Mock() + mock_status.value = "succeeded" + mock_job = Mock( + id="job-123", + type="extract_resume_profile_job", + status=mock_status, + attempts=1, + max_attempts=8, + last_error=None, + payload={"result": {"success": True}}, + ) + + with patch("five08.worker.api.get_job", return_value=mock_job): + response = await api.job_status_handler(request) + + payload = json.loads(response.text) + assert response.status == 200 + assert payload["job_id"] == "job-123" + assert payload["status"] == "succeeded" + assert payload["result"] == {"success": True} From 5e52043a8f59fafb70c8eccd4f3987d4550c256c Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 15:09:10 +0800 Subject: [PATCH 2/8] feat: track resume LLM processing timestamp --- .../worker/crm/resume_profile_processor.py | 18 ++++++++++++++++++ tests/unit/test_resume_profile_processor.py | 4 ++++ 2 files changed, 22 insertions(+) diff --git a/apps/worker/src/five08/worker/crm/resume_profile_processor.py b/apps/worker/src/five08/worker/crm/resume_profile_processor.py index f58b6964..c09fde56 100644 --- a/apps/worker/src/five08/worker/crm/resume_profile_processor.py +++ b/apps/worker/src/five08/worker/crm/resume_profile_processor.py @@ -6,6 +6,7 @@ import logging import re from collections.abc import Callable +from datetime import datetime, timezone from typing import Any from five08.clients.espo import EspoAPI, EspoAPIError @@ -307,6 +308,9 @@ def extract_profile_proposal( ) ) + # Track extraction completion before user confirmation/apply step. + self._mark_resume_processed(contact_id) + return ResumeExtractionResult( contact_id=contact_id, attachment_id=attachment_id, @@ -460,3 +464,17 @@ def _parse_existing_skills(self, value: Any) -> list[str]: parsed = [item.strip() for item in str(value).split(",")] return [item for item in parsed if item] + + def _mark_resume_processed(self, contact_id: str) -> None: + """Best-effort update for extraction completion tracking.""" + processed_at = datetime.now(tz=timezone.utc).isoformat() + try: + self.crm.update_contact( + contact_id, {"cResumeLastProcessed": processed_at} + ) + except Exception as exc: + logger.warning( + "Failed to update cResumeLastProcessed contact_id=%s error=%s", + contact_id, + exc, + ) diff --git a/tests/unit/test_resume_profile_processor.py b/tests/unit/test_resume_profile_processor.py index 75c191f3..5cce3bac 100644 --- a/tests/unit/test_resume_profile_processor.py +++ b/tests/unit/test_resume_profile_processor.py @@ -50,6 +50,10 @@ def test_extract_profile_proposal_filters_508_email() -> None: assert result.proposed_updates["skills"] == "Python, FastAPI" assert result.new_skills == ["Python", "FastAPI"] assert any(item.field == "emailAddress" for item in result.skipped) + processor.crm.update_contact.assert_called_once() + update_contact_payload = processor.crm.update_contact.call_args.args[1] + assert "cResumeLastProcessed" in update_contact_payload + assert isinstance(update_contact_payload["cResumeLastProcessed"], str) def test_apply_profile_updates_adds_discord_and_filters_email() -> None: From 2b6c18e25d12f2c83cc604ebe7f94cb3f9a705ac Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 15:17:40 +0800 Subject: [PATCH 3/8] feat: add resume processing ledger table and keys --- .env.example | 1 + DEVELOPMENT.md | 2 +- README.md | 1 + apps/worker/src/five08/worker/api.py | 13 ++- apps/worker/src/five08/worker/config.py | 1 + .../worker/crm/resume_profile_processor.py | 87 +++++++++++++++++- ...200_create_resume_processing_runs_table.py | 91 +++++++++++++++++++ tests/unit/test_resume_profile_processor.py | 36 ++++++++ tests/unit/test_worker_api.py | 17 ++++ 9 files changed, 244 insertions(+), 5 deletions(-) create mode 100644 apps/worker/src/five08/worker/migrations/versions/20260221_0200_create_resume_processing_runs_table.py diff --git a/.env.example b/.env.example index bdc21b44..bc933bd8 100644 --- a/.env.example +++ b/.env.example @@ -51,6 +51,7 @@ OPENAI_API_KEY= # For OpenRouter, set OPENAI_BASE_URL=https://openrouter.ai/api/v1 OPENAI_BASE_URL= OPENAI_MODEL=gpt-4o-mini +RESUME_EXTRACTOR_VERSION=v1 # Discord bot (required for bot runtime) DISCORD_BOT_TOKEN=your_bot_token_here diff --git a/DEVELOPMENT.md b/DEVELOPMENT.md index 88126f4d..4994d401 100644 --- a/DEVELOPMENT.md +++ b/DEVELOPMENT.md @@ -117,7 +117,7 @@ Use `.env.example` as source of truth. Key categories: - Shared queue/runtime: `REDIS_URL`, `REDIS_QUEUE_NAME`, `POSTGRES_URL`, `JOB_MAX_ATTEMPTS`, `JOB_RETRY_BASE_SECONDS`, `JOB_RETRY_MAX_SECONDS`, `LOG_LEVEL`, webhook settings - Bot credentials/integrations: Discord, email, Espo, Kimai - Worker controls: `WORKER_NAME`, `WORKER_QUEUE_NAMES`, `WORKER_BURST` -- Worker CRM processing: `MAX_ATTACHMENTS_PER_CONTACT`, `MAX_FILE_SIZE_MB`, `ALLOWED_FILE_TYPES`, `RESUME_KEYWORDS`, `OPENAI_API_KEY`, `OPENAI_BASE_URL`, `OPENAI_MODEL` +- Worker CRM processing: `MAX_ATTACHMENTS_PER_CONTACT`, `MAX_FILE_SIZE_MB`, `ALLOWED_FILE_TYPES`, `RESUME_KEYWORDS`, `OPENAI_API_KEY`, `OPENAI_BASE_URL`, `OPENAI_MODEL`, `RESUME_EXTRACTOR_VERSION` - Resume upload UX wiring: `WORKER_API_BASE_URL` on bot, `CRM_LINKEDIN_FIELD` on worker. ## CI Notes diff --git a/README.md b/README.md index 454be3fa..0892bfe5 100644 --- a/README.md +++ b/README.md @@ -145,6 +145,7 @@ docker compose up --build - `OPENAI_API_KEY` (optional; if unset, heuristic extraction is used) - `OPENAI_BASE_URL` (optional; set `https://openrouter.ai/api/v1` for OpenRouter) - `OPENAI_MODEL` (default: `gpt-4o-mini`) +- `RESUME_EXTRACTOR_VERSION` (default: `v1`; used in resume processing idempotency/ledger keys) ## Commands diff --git a/apps/worker/src/five08/worker/api.py b/apps/worker/src/five08/worker/api.py index 74833b7d..29321ff6 100644 --- a/apps/worker/src/five08/worker/api.py +++ b/apps/worker/src/five08/worker/api.py @@ -67,6 +67,12 @@ def _extract_idempotency_key(value: object) -> str | None: return None +def _resume_extract_model_name() -> str: + if settings.openai_api_key and settings.openai_model: + return settings.openai_model + return "heuristic" + + async def health_handler(request: web.Request) -> web.Response: """Simple health endpoint.""" redis_conn = request.app[REDIS_CONN_KEY] @@ -229,13 +235,18 @@ async def resume_extract_handler(request: web.Request) -> web.Response: ) queue = request.app[QUEUE_KEY] + model_name = _resume_extract_model_name() + idempotency_key = ( + f"resume-extract:{payload.contact_id}:{payload.attachment_id}:" + f"{settings.resume_extractor_version}:{model_name}" + ) job = await asyncio.to_thread( enqueue_job, queue=queue, fn=extract_resume_profile_job, args=(payload.contact_id, payload.attachment_id, payload.filename), settings=settings, - idempotency_key=f"resume-extract:{payload.contact_id}:{payload.attachment_id}", + idempotency_key=idempotency_key, ) logger.info( "Enqueued resume extract job contact_id=%s attachment_id=%s job_id=%s created=%s", diff --git a/apps/worker/src/five08/worker/config.py b/apps/worker/src/five08/worker/config.py index 07784f54..84083c56 100644 --- a/apps/worker/src/five08/worker/config.py +++ b/apps/worker/src/five08/worker/config.py @@ -17,6 +17,7 @@ class WorkerSettings(SharedSettings): openai_api_key: str | None = None openai_base_url: str | None = None openai_model: str = "gpt-4o-mini" + resume_extractor_version: str = "v1" max_file_size_mb: int = 10 allowed_file_types: str = "pdf,doc,docx,txt" diff --git a/apps/worker/src/five08/worker/crm/resume_profile_processor.py b/apps/worker/src/five08/worker/crm/resume_profile_processor.py index c09fde56..50ef55e9 100644 --- a/apps/worker/src/five08/worker/crm/resume_profile_processor.py +++ b/apps/worker/src/five08/worker/crm/resume_profile_processor.py @@ -10,6 +10,7 @@ from typing import Any from five08.clients.espo import EspoAPI, EspoAPIError +from five08.queue import get_postgres_connection from five08.worker.config import settings from five08.worker.crm.document_processor import DocumentProcessor from five08.worker.crm.skills_extractor import SkillsExtractor @@ -238,11 +239,15 @@ def extract_profile_proposal( filename: str, ) -> ResumeExtractionResult: """Build preview proposal from an uploaded resume attachment.""" + content_hash: str | None = None + model_name = self._configured_model_name() try: contact = self.crm.get_contact(contact_id) content = self.crm.download_attachment(attachment_id) + content_hash = self.document_processor.get_content_hash(content) text = self.document_processor.extract_text(content, filename) extracted = self.extractor.extract(text) + model_name = extracted.source extracted_skills_result = self.skills_extractor.extract_skills(text) extracted_skills = extracted_skills_result.skills existing_skills = self._parse_existing_skills(contact.get("skills")) @@ -310,6 +315,13 @@ def extract_profile_proposal( # Track extraction completion before user confirmation/apply step. self._mark_resume_processed(contact_id) + self._record_processing_run( + contact_id=contact_id, + attachment_id=attachment_id, + content_hash=content_hash, + model_name=model_name, + status="succeeded", + ) return ResumeExtractionResult( contact_id=contact_id, @@ -329,6 +341,14 @@ def extract_profile_proposal( attachment_id, exc, ) + self._record_processing_run( + contact_id=contact_id, + attachment_id=attachment_id, + content_hash=content_hash, + model_name=model_name, + status="failed", + last_error=str(exc), + ) return ResumeExtractionResult( contact_id=contact_id, attachment_id=attachment_id, @@ -469,12 +489,73 @@ def _mark_resume_processed(self, contact_id: str) -> None: """Best-effort update for extraction completion tracking.""" processed_at = datetime.now(tz=timezone.utc).isoformat() try: - self.crm.update_contact( - contact_id, {"cResumeLastProcessed": processed_at} - ) + self.crm.update_contact(contact_id, {"cResumeLastProcessed": processed_at}) except Exception as exc: logger.warning( "Failed to update cResumeLastProcessed contact_id=%s error=%s", contact_id, exc, ) + + def _configured_model_name(self) -> str: + """Model identity used for idempotency/ledger keys.""" + if settings.openai_api_key and settings.openai_model: + return settings.openai_model + return "heuristic" + + def _record_processing_run( + self, + *, + contact_id: str, + attachment_id: str, + content_hash: str | None, + model_name: str, + status: str, + last_error: str | None = None, + ) -> None: + """Persist one processing result keyed by contact+attachment+version+model.""" + query = """ + INSERT INTO resume_processing_runs ( + contact_id, + attachment_id, + content_hash, + extractor_version, + model_name, + status, + last_error, + processed_at + ) + VALUES (%s, %s, %s, %s, %s, %s, %s, NOW()) + ON CONFLICT (contact_id, attachment_id, extractor_version, model_name) + DO UPDATE SET + content_hash = EXCLUDED.content_hash, + status = EXCLUDED.status, + last_error = EXCLUDED.last_error, + processed_at = NOW(); + """ + try: + with get_postgres_connection(settings) as conn: + with conn.cursor() as cursor: + cursor.execute( + query, + ( + contact_id, + attachment_id, + content_hash, + settings.resume_extractor_version, + model_name, + status, + last_error, + ), + ) + except Exception as exc: + logger.warning( + "Failed to persist resume processing run contact_id=%s attachment_id=%s " + "version=%s model=%s status=%s error=%s", + contact_id, + attachment_id, + settings.resume_extractor_version, + model_name, + status, + exc, + ) diff --git a/apps/worker/src/five08/worker/migrations/versions/20260221_0200_create_resume_processing_runs_table.py b/apps/worker/src/five08/worker/migrations/versions/20260221_0200_create_resume_processing_runs_table.py new file mode 100644 index 00000000..5b27a68f --- /dev/null +++ b/apps/worker/src/five08/worker/migrations/versions/20260221_0200_create_resume_processing_runs_table.py @@ -0,0 +1,91 @@ +"""Create resume processing runs table.""" + +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op + +revision = "202602210200" +down_revision = "202602210100" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + """Create per-resume processing ledger table.""" + op.create_table( + "resume_processing_runs", + sa.Column("contact_id", sa.Text(), nullable=False), + sa.Column("attachment_id", sa.Text(), nullable=False), + sa.Column("content_hash", sa.Text(), nullable=True), + sa.Column("extractor_version", sa.Text(), nullable=False), + sa.Column("model_name", sa.Text(), nullable=False), + sa.Column("status", sa.Text(), nullable=False), + sa.Column("last_error", sa.Text(), nullable=True), + sa.Column( + "processed_at", + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.text("NOW()"), + ), + sa.Column( + "created_at", + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.text("NOW()"), + ), + sa.Column( + "updated_at", + sa.DateTime(timezone=True), + nullable=False, + server_default=sa.text("NOW()"), + ), + sa.PrimaryKeyConstraint( + "contact_id", + "attachment_id", + "extractor_version", + "model_name", + name="pk_resume_processing_runs", + ), + sa.CheckConstraint( + "status IN ('succeeded', 'failed')", + name="ck_resume_processing_runs_status", + ), + ) + op.create_index( + "idx_resume_processing_runs_processed_at", + "resume_processing_runs", + ["processed_at"], + ) + op.execute( + """ + CREATE FUNCTION resume_processing_runs_set_updated_at_fn() + RETURNS TRIGGER AS $$ + BEGIN + NEW.updated_at = NOW(); + RETURN NEW; + END; + $$ LANGUAGE plpgsql; + """ + ) + op.execute( + """ + CREATE TRIGGER resume_processing_runs_set_updated_at_tr + BEFORE UPDATE ON resume_processing_runs + FOR EACH ROW + EXECUTE FUNCTION resume_processing_runs_set_updated_at_fn(); + """ + ) + + +def downgrade() -> None: + """Drop resume processing ledger table.""" + op.execute( + "DROP TRIGGER IF EXISTS resume_processing_runs_set_updated_at_tr ON resume_processing_runs" + ) + op.execute("DROP FUNCTION IF EXISTS resume_processing_runs_set_updated_at_fn()") + op.drop_index( + "idx_resume_processing_runs_processed_at", + table_name="resume_processing_runs", + ) + op.drop_table("resume_processing_runs") diff --git a/tests/unit/test_resume_profile_processor.py b/tests/unit/test_resume_profile_processor.py index 5cce3bac..51599d77 100644 --- a/tests/unit/test_resume_profile_processor.py +++ b/tests/unit/test_resume_profile_processor.py @@ -13,6 +13,7 @@ def test_extract_profile_proposal_filters_508_email() -> None: processor.extractor = Mock() processor.skills_extractor = Mock() processor.document_processor = Mock() + processor._record_processing_run = Mock() processor.crm.get_contact.return_value = { "emailAddress": "member@example.com", @@ -22,6 +23,7 @@ def test_extract_profile_proposal_filters_508_email() -> None: } processor.crm.download_attachment.return_value = b"resume-bytes" processor.document_processor.extract_text.return_value = "resume text" + processor.document_processor.get_content_hash.return_value = "hash-1" processor.extractor.extract.return_value = ResumeExtractedProfile( email="new@508.dev", github_username="new-gh", @@ -54,6 +56,11 @@ def test_extract_profile_proposal_filters_508_email() -> None: update_contact_payload = processor.crm.update_contact.call_args.args[1] assert "cResumeLastProcessed" in update_contact_payload assert isinstance(update_contact_payload["cResumeLastProcessed"], str) + processor._record_processing_run.assert_called_once() + record_kwargs = processor._record_processing_run.call_args.kwargs + assert record_kwargs["status"] == "succeeded" + assert record_kwargs["contact_id"] == "contact-1" + assert record_kwargs["attachment_id"] == "att-1" def test_apply_profile_updates_adds_discord_and_filters_email() -> None: @@ -81,3 +88,32 @@ def test_apply_profile_updates_adds_discord_and_filters_email() -> None: assert update_payload["skills"] == "Python, FastAPI" assert update_payload["cDiscordUserID"] == "123" assert update_payload["cDiscordUsername"] == "member#0001 (ID: 123)" + + +def test_extract_profile_proposal_records_failed_run() -> None: + """Failed extraction should still be written to the processing ledger.""" + processor = ResumeProfileProcessor() + processor.crm = Mock() + processor.extractor = Mock() + processor.skills_extractor = Mock() + processor.document_processor = Mock() + processor._record_processing_run = Mock() + + processor.crm.get_contact.return_value = {"emailAddress": "member@example.com"} + processor.crm.download_attachment.return_value = b"resume-bytes" + processor.document_processor.get_content_hash.return_value = "hash-2" + processor.document_processor.extract_text.side_effect = ValueError("parse failed") + + result = processor.extract_profile_proposal( + contact_id="contact-2", + attachment_id="att-2", + filename="broken.pdf", + ) + + assert result.success is False + processor._record_processing_run.assert_called_once() + record_kwargs = processor._record_processing_run.call_args.kwargs + assert record_kwargs["status"] == "failed" + assert record_kwargs["contact_id"] == "contact-2" + assert record_kwargs["attachment_id"] == "att-2" + assert record_kwargs["content_hash"] == "hash-2" diff --git a/tests/unit/test_worker_api.py b/tests/unit/test_worker_api.py index ae818355..f3126cc5 100644 --- a/tests/unit/test_worker_api.py +++ b/tests/unit/test_worker_api.py @@ -171,9 +171,14 @@ async def test_process_contact_handler_enqueues_single_contact( @pytest.mark.asyncio async def test_resume_extract_handler_enqueues_job( + monkeypatch: pytest.MonkeyPatch, auth_headers: dict[str, str], ) -> None: """Resume extract endpoint should enqueue extraction job.""" + monkeypatch.setattr(api.settings, "resume_extractor_version", "v7") + monkeypatch.setattr(api.settings, "openai_api_key", "key") + monkeypatch.setattr(api.settings, "openai_model", "gpt-test") + app_obj = web.Application() app_obj[api.QUEUE_KEY] = Mock() request = make_mocked_request( @@ -196,6 +201,8 @@ async def test_resume_extract_handler_enqueues_job( assert payload["job_id"] == "job-extract" assert payload["contact_id"] == "c-1" assert payload["attachment_id"] == "a-1" + call_kwargs = mock_enqueue.call_args.kwargs + assert call_kwargs["idempotency_key"] == "resume-extract:c-1:a-1:v7:gpt-test" @pytest.mark.asyncio @@ -260,3 +267,13 @@ async def test_job_status_handler_returns_result( assert payload["job_id"] == "job-123" assert payload["status"] == "succeeded" assert payload["result"] == {"success": True} + + +def test_resume_extract_model_name_uses_heuristic_without_api_key( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Model identity should be heuristic when OpenAI key is absent.""" + monkeypatch.setattr(api.settings, "openai_api_key", None) + monkeypatch.setattr(api.settings, "openai_model", "gpt-test") + + assert api._resume_extract_model_name() == "heuristic" From d0cc5d4daecde99dec6f32c128492d4b5f2537cb Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 15:58:50 +0800 Subject: [PATCH 4/8] feat: improve resume extraction, skills normalization, and CRM sync --- .env.example | 2 + README.md | 3 +- .../src/five08/discord_bot/cogs/crm.py | 232 +++++++++++++++++- apps/worker/src/five08/worker/api.py | 4 +- apps/worker/src/five08/worker/config.py | 26 ++ .../worker/src/five08/worker/crm/processor.py | 15 +- .../worker/crm/resume_profile_processor.py | 149 ++++++++++- .../src/five08/worker/crm/skills_extractor.py | 137 +++++++++-- apps/worker/src/five08/worker/models.py | 7 + packages/shared/src/five08/skills.py | 71 ++++++ tests/unit/test_resume_profile_processor.py | 70 ++++++ tests/unit/test_shared_skills.py | 19 ++ tests/unit/test_skills_extractor.py | 33 +++ tests/unit/test_worker_api.py | 16 +- 14 files changed, 740 insertions(+), 44 deletions(-) create mode 100644 packages/shared/src/five08/skills.py create mode 100644 tests/unit/test_shared_skills.py diff --git a/.env.example b/.env.example index 0459ae07..549cf6f3 100644 --- a/.env.example +++ b/.env.example @@ -55,6 +55,8 @@ OPENAI_API_KEY= # For OpenRouter, set OPENAI_BASE_URL=https://openrouter.ai/api/v1 OPENAI_BASE_URL= OPENAI_MODEL=gpt-4o-mini +# Resume model name without provider prefix; OpenRouter is auto-prefixed to openai/ +RESUME_AI_MODEL=gpt-4o-mini RESUME_EXTRACTOR_VERSION=v1 CRM_SYNC_ENABLED=true CRM_SYNC_INTERVAL_SECONDS=900 diff --git a/README.md b/README.md index 4fc9ab1a..7e8fd3fd 100644 --- a/README.md +++ b/README.md @@ -159,7 +159,8 @@ Use `.env.example` as the source of truth for defaults. - `Optional`: `RESUME_KEYWORDS` (default: `resume,cv,curriculum`) - `Optional`: `OPENAI_API_KEY` (if unset, heuristic extraction is used) - `Optional`: `OPENAI_BASE_URL` (set `https://openrouter.ai/api/v1` for OpenRouter) -- `Optional`: `OPENAI_MODEL` (default: `gpt-4o-mini`) +- `Optional`: `RESUME_AI_MODEL` (default: `gpt-4o-mini`; use plain names like `gpt-4o-mini`, OpenRouter gets auto-prefixed to `openai/`) +- `Optional`: `OPENAI_MODEL` (default: `gpt-4o-mini`; fallback/legacy model setting) - `Optional`: `RESUME_EXTRACTOR_VERSION` (default: `v1`; used in resume processing idempotency/ledger keys) ### Discord Bot Core diff --git a/apps/discord_bot/src/five08/discord_bot/cogs/crm.py b/apps/discord_bot/src/five08/discord_bot/cogs/crm.py index 362548a4..bd8dee2e 100644 --- a/apps/discord_bot/src/five08/discord_bot/cogs/crm.py +++ b/apps/discord_bot/src/five08/discord_bot/cogs/crm.py @@ -17,6 +17,7 @@ from five08.discord_bot.config import settings from five08.clients import espo +from five08.skills import normalize_skill_list from five08.discord_bot.utils.audit import DiscordAuditLogger from five08.discord_bot.utils.role_decorators import ( require_role, @@ -398,6 +399,19 @@ async def confirm_updates( ) except Exception as exc: logger.error("Failed to enqueue resume apply job: %s", exc) + self.crm_cog._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "stage": "apply_enqueue", + "error": str(exc), + "proposed_updates_count": len(self.proposed_updates), + "link_member_requested": bool(self.link_discord), + }, + resource_type="crm_contact", + resource_id=self.contact_id, + ) await interaction.followup.send( "โŒ Failed to enqueue CRM apply job. Please try again." ) @@ -410,6 +424,17 @@ async def confirm_updates( apply_result = await self.crm_cog._wait_for_worker_job_result(apply_job_id) if not apply_result: + self.crm_cog._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "stage": "apply_timeout", + "job_id": apply_job_id, + }, + resource_type="crm_contact", + resource_id=self.contact_id, + ) await interaction.followup.send( "โš ๏ธ Timed out waiting for apply job. Please check again shortly." ) @@ -417,6 +442,19 @@ async def confirm_updates( status = str(apply_result.get("status", "unknown")) if status != "succeeded": + self.crm_cog._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "stage": "apply_failed", + "job_id": apply_job_id, + "job_status": status, + "last_error": str(apply_result.get("last_error", "")), + }, + resource_type="crm_contact", + resource_id=self.contact_id, + ) await interaction.followup.send( f"โŒ Apply job failed (status: {status}). " f"Error: {apply_result.get('last_error') or 'Unknown error'}" @@ -442,6 +480,20 @@ async def confirm_updates( ) profile_url = f"{self.crm_cog.base_url}/#Contact/view/{self.contact_id}" embed.add_field(name="๐Ÿ”— CRM Profile", value=f"[View in CRM]({profile_url})") + self.crm_cog._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="success", + metadata={ + "stage": "apply_succeeded", + "job_id": apply_job_id, + "updated_fields": updated_fields, + "proposed_updates_count": len(self.proposed_updates), + "link_member_requested": bool(self.link_discord), + }, + resource_type="crm_contact", + resource_id=self.contact_id, + ) await interaction.followup.send(embed=embed) for item in self.children: @@ -462,6 +514,18 @@ async def cancel_updates( button: discord.ui.Button["ResumeUpdateConfirmationView"], ) -> None: """Cancel CRM updates after preview.""" + self.crm_cog._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="denied", + metadata={ + "stage": "apply_cancelled", + "proposed_updates_count": len(self.proposed_updates), + "link_member_requested": bool(self.link_discord), + }, + resource_type="crm_contact", + resource_id=self.contact_id, + ) await interaction.response.send_message( "No CRM profile updates were applied.", ephemeral=True ) @@ -719,6 +783,19 @@ async def _run_resume_extract_and_preview( ) except Exception as exc: logger.error("Failed to enqueue resume extract job: %s", exc) + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "stage": "extract_enqueue", + "error": str(exc), + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send( "โš ๏ธ Resume uploaded, but extraction job could not be enqueued.", ephemeral=True, @@ -734,12 +811,39 @@ async def _run_resume_extract_and_preview( job = await self._wait_for_worker_job_result(job_id) except Exception as exc: logger.error("Worker polling failed for job_id=%s error=%s", job_id, exc) + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "job_id": job_id, + "stage": "extract_polling", + "error": str(exc), + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send( "โš ๏ธ Resume uploaded, but extraction polling failed.", ephemeral=True, ) return if not job: + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "job_id": job_id, + "stage": "extract_timeout", + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send( "โš ๏ธ Timed out waiting for extraction result. Try again in a moment.", ephemeral=True, @@ -748,6 +852,21 @@ async def _run_resume_extract_and_preview( status = str(job.get("status", "unknown")) if status != "succeeded": + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "job_id": job_id, + "stage": "extract_failed", + "job_status": status, + "last_error": str(job.get("last_error", "")), + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send( f"โŒ Extraction job failed (status: {status}). " f"Error: {job.get('last_error') or 'Unknown error'}", @@ -757,6 +876,19 @@ async def _run_resume_extract_and_preview( result = job.get("result") if not isinstance(result, dict): + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "job_id": job_id, + "stage": "extract_malformed_result", + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send( "โŒ Extraction result was empty or malformed.", ephemeral=True, @@ -764,6 +896,20 @@ async def _run_resume_extract_and_preview( return if not result.get("success", False): + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "job_id": job_id, + "stage": "extract_unsuccessful", + "error": str(result.get("error", "")), + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send( f"โŒ Resume extraction failed: {result.get('error') or 'Unknown error'}", ephemeral=True, @@ -778,6 +924,19 @@ async def _run_resume_extract_and_preview( ) if not proposed_updates and not link_member: + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="success", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "job_id": job_id, + "stage": "preview_no_changes", + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send(embed=embed, ephemeral=True) return @@ -796,6 +955,21 @@ async def _run_resume_extract_and_preview( proposed_updates=proposed_updates, link_discord=link_discord_payload, ) + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="success", + metadata={ + "filename": filename, + "attachment_id": attachment_id, + "job_id": job_id, + "stage": "preview_ready", + "proposed_updates_count": len(proposed_updates), + "link_member_requested": bool(link_member), + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send(embed=embed, view=view, ephemeral=True) async def _download_and_send_resume( @@ -849,11 +1023,12 @@ async def search_members( await interaction.response.defer(ephemeral=True) query_value = (query or "").strip() - skills_list = ( + raw_skills_list = ( [skill.strip() for skill in skills.split(",") if skill.strip()] if skills else [] ) + skills_list = normalize_skill_list(raw_skills_list) if not query_value and not skills_list: self._audit_command( @@ -2038,6 +2213,15 @@ async def upload_resume( await interaction.response.defer(ephemeral=True) if not settings.api_shared_secret: + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "filename": file.filename, + "reason": "api_shared_secret_missing", + }, + ) await interaction.followup.send( "โŒ API_SHARED_SECRET is not configured for worker API access." ) @@ -2093,6 +2277,17 @@ async def upload_resume( interaction.user.roles, ["Steering Committee"] ) ): + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="denied", + metadata={ + "search_term": search_term, + "filename": file.filename, + "target_scope": "other" if search_term else "self", + "reason": "missing_required_role", + }, + ) await interaction.followup.send( "โŒ You must have Steering Committee role or higher for this upload." ) @@ -2107,7 +2302,7 @@ async def upload_resume( self._audit_command( interaction=interaction, action="crm.upload_resume", - result="success", + result="denied", metadata={ "search_term": search_term, "filename": file.filename, @@ -2123,7 +2318,7 @@ async def upload_resume( self._audit_command( interaction=interaction, action="crm.upload_resume", - result="success", + result="denied", metadata={ "search_term": search_term, "filename": file.filename, @@ -2210,6 +2405,21 @@ async def upload_resume( if not await self._update_contact_resume( contact_id, attachment_id, overwrite ): + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="error", + metadata={ + "search_term": search_term, + "filename": file.filename, + "attachment_id": attachment_id, + "overwrite": overwrite, + "target_scope": "other" if search_term else "self", + "reason": "resume_link_update_failed", + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await interaction.followup.send( "โš ๏ธ File uploaded, but failed to link in contact resume field." ) @@ -2222,6 +2432,22 @@ async def upload_resume( attachment_id, interaction.user.name, ) + self._audit_command( + interaction=interaction, + action="crm.upload_resume", + result="success", + metadata={ + "search_term": search_term, + "filename": file.filename, + "size_bytes": file.size, + "overwrite": overwrite, + "target_scope": "other" if search_term else "self", + "attachment_id": attachment_id, + "stage": "uploaded_and_linked", + }, + resource_type="crm_contact", + resource_id=str(contact_id), + ) await self._run_resume_extract_and_preview( interaction=interaction, contact_id=contact_id, diff --git a/apps/worker/src/five08/worker/api.py b/apps/worker/src/five08/worker/api.py index 3391a3ef..2eddb023 100644 --- a/apps/worker/src/five08/worker/api.py +++ b/apps/worker/src/five08/worker/api.py @@ -84,8 +84,8 @@ def _extract_idempotency_key(value: object) -> str | None: def _resume_extract_model_name() -> str: - if settings.openai_api_key and settings.openai_model: - return settings.openai_model + if settings.openai_api_key: + return settings.resolved_resume_ai_model return "heuristic" diff --git a/apps/worker/src/five08/worker/config.py b/apps/worker/src/five08/worker/config.py index bec5b292..17284d0e 100644 --- a/apps/worker/src/five08/worker/config.py +++ b/apps/worker/src/five08/worker/config.py @@ -1,5 +1,7 @@ """Configuration for webhook ingest and worker services.""" +from urllib.parse import urlparse + from five08.settings import SharedSettings @@ -17,6 +19,7 @@ class WorkerSettings(SharedSettings): openai_api_key: str | None = None openai_base_url: str | None = None openai_model: str = "gpt-4o-mini" + resume_ai_model: str = "gpt-4o-mini" resume_extractor_version: str = "v1" max_file_size_mb: int = 10 @@ -41,5 +44,28 @@ def parsed_resume_keywords(self) -> set[str]: if keyword.strip() } + @property + def resolved_resume_ai_model(self) -> str: + """Resolve provider-specific resume model name (e.g. OpenRouter prefixes).""" + candidate = self.resume_ai_model.strip() + if not candidate: + candidate = self.openai_model.strip() + if not candidate: + return "gpt-4o-mini" + + # Keep explicit provider prefixes intact. + if "/" in candidate: + return candidate + + base_url = (self.openai_base_url or "").strip() + if not base_url: + return candidate + + parsed = urlparse(base_url) + host = (parsed.netloc or parsed.path).split("/")[0].split(":")[0].lower() + if host.endswith("openrouter.ai"): + return f"openai/{candidate}" + return candidate + settings = WorkerSettings() # type: ignore[call-arg] diff --git a/apps/worker/src/five08/worker/crm/processor.py b/apps/worker/src/five08/worker/crm/processor.py index e14e245a..f7b3cc9b 100644 --- a/apps/worker/src/five08/worker/crm/processor.py +++ b/apps/worker/src/five08/worker/crm/processor.py @@ -187,8 +187,19 @@ def _extract_from_attachments( def _parse_existing_skills(self, skills_text: str | None) -> list[str]: if not skills_text: return [] - parsed = [skill.strip() for skill in skills_text.split(",")] - return [skill for skill in parsed if skill] + parsed = [skill.strip() for skill in skills_text.split(",") if skill.strip()] + normalized: list[str] = [] + seen: set[str] = set() + for skill in parsed: + canonical = self.skills_extractor.canonicalize_skill(skill) + if not canonical: + continue + key = canonical.casefold() + if key in seen: + continue + seen.add(key) + normalized.append(canonical) + return normalized def _filter_resume_attachments( self, attachments: list[dict[str, Any]] diff --git a/apps/worker/src/five08/worker/crm/resume_profile_processor.py b/apps/worker/src/five08/worker/crm/resume_profile_processor.py index 50ef55e9..c845022f 100644 --- a/apps/worker/src/five08/worker/crm/resume_profile_processor.py +++ b/apps/worker/src/five08/worker/crm/resume_profile_processor.py @@ -20,6 +20,7 @@ ResumeExtractionResult, ResumeFieldChange, ResumeSkipReason, + SkillAttributes, ) logger = logging.getLogger(__name__) @@ -51,7 +52,7 @@ class ResumeProfileExtractor: """Extract candidate profile fields from resume text.""" def __init__(self) -> None: - self.model = settings.openai_model + self.model = settings.resolved_resume_ai_model self.client: Any = None if settings.openai_api_key and OpenAIClient is not None: @@ -72,7 +73,9 @@ def extract(self, resume_text: str) -> ResumeExtractedProfile: { "role": "system", "content": ( - "Extract contact fields from resume text. Return JSON only." + "You extract candidate contact fields from resumes for a CRM. " + "Return JSON only with no commentary. Be conservative: when unsure, use null. " + "Prefer candidate-owned contact info and ignore references or company contact details." ), }, { @@ -143,11 +146,18 @@ def _heuristic_extract(self, resume_text: str) -> ResumeExtractedProfile: def _build_prompt(self, resume_text: str) -> str: snippet = resume_text[:12000] return ( - "Extract contact fields from this resume.\\n" - "Return JSON with exact keys:\\n" + "Extract candidate contact fields from this resume.\\n" + "Return JSON with exact keys and no extras:\\n" '{"email": string|null, "github_username": string|null, ' '"linkedin_url": string|null, "phone": string|null, ' - '"confidence": number}\\n\\n' + '"confidence": number}\\n' + "Rules:\\n" + "- prefer explicit values from header/contact sections\\n" + "- for github_username return username only (no URL, no @)\\n" + "- for linkedin_url return full linkedin profile URL when available\\n" + "- for phone return digits with optional leading +\\n" + "- use null for unknown/ambiguous fields\\n" + "- confidence is 0-1 for overall extraction reliability\\n\\n" f"Resume:\\n{snippet}" ) @@ -251,12 +261,19 @@ def extract_profile_proposal( extracted_skills_result = self.skills_extractor.extract_skills(text) extracted_skills = extracted_skills_result.skills existing_skills = self._parse_existing_skills(contact.get("skills")) + existing_skill_attrs = self._parse_skill_attrs(contact.get("cSkillAttrs")) existing_lower = {item.casefold() for item in existing_skills} new_skills = [ skill for skill in extracted_skills if skill.casefold() not in existing_lower ] + merged_skills = existing_skills + new_skills + merged_skill_attrs = self._merge_skill_attrs( + existing_attrs=existing_skill_attrs, + extracted_attrs=extracted_skills_result.skill_attrs, + merged_skills=merged_skills, + ) proposed_updates: dict[str, str] = {} proposed_changes: list[ResumeFieldChange] = [] @@ -301,7 +318,6 @@ def extract_profile_proposal( skipped=skipped, ) if new_skills: - merged_skills = existing_skills + new_skills proposed_updates["skills"] = ", ".join(merged_skills) proposed_changes.append( ResumeFieldChange( @@ -313,6 +329,24 @@ def extract_profile_proposal( ) ) + if merged_skill_attrs and merged_skill_attrs != existing_skill_attrs: + proposed_updates["cSkillAttrs"] = self._serialize_skill_attrs( + merged_skill_attrs + ) + proposed_changes.append( + ResumeFieldChange( + field="cSkillAttrs", + label="Skill Attributes", + current=( + f"{len(existing_skill_attrs)} skills rated" + if existing_skill_attrs + else None + ), + proposed=f"{len(merged_skill_attrs)} skills rated (strength 1-5)", + reason="Updated structured skill strengths from resume extraction", + ) + ) + # Track extraction completion before user confirmation/apply step. self._mark_resume_processed(contact_id) self._record_processing_run( @@ -384,8 +418,9 @@ def apply_profile_updates( settings.crm_linkedin_field, "phoneNumber", "skills", + "cSkillAttrs", } - sanitized_updates: dict[str, str] = { + sanitized_updates: dict[str, Any] = { field: value for field, value in updates.items() if field in allowed_fields and value @@ -395,6 +430,16 @@ def apply_profile_updates( if email_value and email_value.lower().endswith("@508.dev"): sanitized_updates.pop("emailAddress") + if "cSkillAttrs" in sanitized_updates: + parsed_attrs = self._parse_skill_attrs(sanitized_updates["cSkillAttrs"]) + # Be forgiving: if value is malformed, overwrite with an empty object. + if parsed_attrs: + sanitized_updates["cSkillAttrs"] = json.loads( + self._serialize_skill_attrs(parsed_attrs) + ) + else: + sanitized_updates["cSkillAttrs"] = {} + link_applied = False if link_discord: discord_user_id = str(link_discord.get("user_id", "")).strip() @@ -478,12 +523,90 @@ def _collect_change( def _parse_existing_skills(self, value: Any) -> list[str]: if value is None: return [] + if isinstance(value, list): - normalized = [str(item).strip() for item in value if str(item).strip()] - return normalized + raw_skills = [str(item).strip() for item in value if str(item).strip()] + else: + raw_skills = [ + item.strip() for item in str(value).split(",") if item.strip() + ] - parsed = [item.strip() for item in str(value).split(",")] - return [item for item in parsed if item] + normalized: list[str] = [] + seen: set[str] = set() + for skill in raw_skills: + canonical = self.skills_extractor.canonicalize_skill(skill) + if not canonical: + continue + key = canonical.casefold() + if key in seen: + continue + seen.add(key) + normalized.append(canonical) + return normalized + + def _parse_skill_attrs(self, value: Any) -> dict[str, int]: + if value is None: + return {} + + candidate = value + if isinstance(value, str): + raw = value.strip() + if not raw: + return {} + try: + candidate = json.loads(raw) + except Exception: + return {} + + if not isinstance(candidate, dict): + return {} + + parsed: dict[str, int] = {} + for raw_skill, raw_payload in candidate.items(): + skill = self.skills_extractor.canonicalize_skill(str(raw_skill)).casefold() + if not skill: + continue + + strength_value = raw_payload + if isinstance(raw_payload, dict): + strength_value = raw_payload.get("strength") + + try: + strength = int(float(strength_value)) + except Exception: + strength = 0 + parsed[skill] = max(1, min(5, strength)) if strength else 0 + + return {skill: strength for skill, strength in parsed.items() if strength > 0} + + def _merge_skill_attrs( + self, + *, + existing_attrs: dict[str, int], + extracted_attrs: dict[str, SkillAttributes], + merged_skills: list[str], + ) -> dict[str, int]: + merged: dict[str, int] = dict(existing_attrs) + + for skill in merged_skills: + key = str(skill).strip().casefold() + if key and key not in merged: + merged[key] = 3 + + for skill, attrs in extracted_attrs.items(): + key = str(skill).strip().casefold() + if key: + merged[key] = max(1, min(5, int(attrs.strength))) + + return merged + + def _serialize_skill_attrs(self, attrs: dict[str, int]) -> str: + payload = { + skill: {"strength": max(1, min(5, int(strength)))} + for skill, strength in sorted(attrs.items()) + if skill + } + return json.dumps(payload, separators=(",", ":"), sort_keys=True) def _mark_resume_processed(self, contact_id: str) -> None: """Best-effort update for extraction completion tracking.""" @@ -499,8 +622,8 @@ def _mark_resume_processed(self, contact_id: str) -> None: def _configured_model_name(self) -> str: """Model identity used for idempotency/ledger keys.""" - if settings.openai_api_key and settings.openai_model: - return settings.openai_model + if settings.openai_api_key: + return settings.resolved_resume_ai_model return "heuristic" def _record_processing_run( diff --git a/apps/worker/src/five08/worker/crm/skills_extractor.py b/apps/worker/src/five08/worker/crm/skills_extractor.py index 3cf6e670..094f7145 100644 --- a/apps/worker/src/five08/worker/crm/skills_extractor.py +++ b/apps/worker/src/five08/worker/crm/skills_extractor.py @@ -5,8 +5,9 @@ import re from typing import Any +from five08.skills import normalize_skill from five08.worker.config import settings -from five08.worker.models import ExtractedSkills +from five08.worker.models import ExtractedSkills, SkillAttributes logger = logging.getLogger(__name__) @@ -22,10 +23,11 @@ "java", "go", "rust", + "node", "docker", "kubernetes", - "aws", - "gcp", + "amazon web services", + "google cloud", "azure", "postgresql", "mysql", @@ -36,14 +38,25 @@ "fastapi", "git", "linux", + "product management", + "go to market", + "ab testing", + "search engine optimization", + "search engine marketing", + "customer relationship management", + "google analytics", + "product marketing", + "content marketing", } +DEFAULT_SKILL_STRENGTH = 3 + class SkillsExtractor: """Extract skills with LLM when configured, fallback heuristics otherwise.""" def __init__(self) -> None: - self.model = settings.openai_model + self.model = settings.resolved_resume_ai_model self.client: Any = None if settings.openai_api_key and OpenAIClient is not None: @@ -65,8 +78,12 @@ def extract_skills(self, resume_text: str) -> ExtractedSkills: { "role": "system", "content": ( - "Extract professional and technical skills from resume text. " - "Return only JSON." + "You extract professional skills from resumes for a CRM. " + "Focus on white-collar skills for product development orgs: " + "engineering, product, data, design, growth, and marketing. " + "Return JSON only, no prose. " + "Normalize skills to concise canonical names, lowercase. " + "Provide a strength from 1-5 for each skill, where 5 is strongest." ), }, {"role": "user", "content": prompt}, @@ -78,17 +95,12 @@ def extract_skills(self, resume_text: str) -> ExtractedSkills: if not content: raise ValueError("LLM returned empty content") - parsed = json.loads(content) - skills = parsed.get("skills", []) + parsed = self._parse_llm_json(content) confidence = float(parsed.get("confidence", 0.7)) - - if not isinstance(skills, list): - raise ValueError("skills must be a list") - - normalized = [skill.strip() for skill in skills if str(skill).strip()] - return ExtractedSkills( - skills=sorted(set(normalized)), - confidence=max(0.0, min(1.0, confidence)), + return self._normalize_extracted_payload( + skills_value=parsed.get("skills", []), + skill_attrs_value=parsed.get("skill_attrs", {}), + confidence=confidence, source=self.model, ) except Exception as exc: @@ -101,12 +113,18 @@ def _extract_skills_heuristic(self, resume_text: str) -> ExtractedSkills: token_matches = re.findall(r"\b[a-z][a-z0-9+#\-.]{1,24}\b", lowered) detected: set[str] = set() for token in token_matches: - if token in COMMON_SKILLS: - detected.add(token) + canonical = self._normalize_skill_name(token) + if canonical in COMMON_SKILLS: + detected.add(canonical) + sorted_skills = sorted(detected) return ExtractedSkills( - skills=sorted(detected), - confidence=0.45 if detected else 0.2, + skills=sorted_skills, + skill_attrs={ + skill: SkillAttributes(strength=DEFAULT_SKILL_STRENGTH) + for skill in sorted_skills + }, + confidence=0.45 if sorted_skills else 0.2, source="heuristic", ) @@ -115,7 +133,84 @@ def _create_prompt(self, resume_text: str) -> str: snippet = resume_text[:8000] return ( "Analyze the resume and extract a concise skill list.\n" + "Use white-collar/product-development relevance only: engineering, product, " + "data, design, growth, and marketing.\n" + "Exclude personal traits and vague soft skills unless role-critical.\n" "Return JSON with this exact schema:\n" - '{"skills": ["skill1", "skill2"], "confidence": 0.8}\n\n' + '{"skills": ["skill1", "skill2"], ' + '"skill_attrs": {"skill1": {"strength": 4}}, ' + '"confidence": 0.8}\n' + "Rules:\n" + "- skills must be lowercase canonical names with minimal punctuation\n" + '- prefer forms like "nodejs", "ab testing", "go to market"\n' + "- skill_attrs keys must match skills\n" + "- strength is integer 1-5 (5 strongest)\n" + "- no extra keys\n\n" f"Resume:\n{snippet}" ) + + def _parse_llm_json(self, content: str) -> dict[str, Any]: + raw = content.strip() + if raw.startswith("```"): + lines = [line for line in raw.splitlines() if not line.startswith("```")] + raw = "\n".join(lines).strip() + + parsed = json.loads(raw) + if not isinstance(parsed, dict): + raise ValueError("skills extraction output was not a JSON object") + return parsed + + def _normalize_extracted_payload( + self, + *, + skills_value: Any, + skill_attrs_value: Any, + confidence: float, + source: str, + ) -> ExtractedSkills: + raw_skills = skills_value if isinstance(skills_value, list) else [] + normalized_skills: list[str] = [] + for skill in raw_skills: + canonical = self._normalize_skill_name(str(skill)) + if canonical: + normalized_skills.append(canonical) + + attrs_map: dict[str, SkillAttributes] = {} + if isinstance(skill_attrs_value, dict): + for raw_name, raw_attr in skill_attrs_value.items(): + canonical = self._normalize_skill_name(str(raw_name)) + if not canonical: + continue + attrs_map[canonical] = SkillAttributes( + strength=self._parse_strength(raw_attr) + ) + + # Ensure attrs exists for every skill and include attr-only entries in skill list. + deduped_skills = sorted(set(normalized_skills) | set(attrs_map.keys())) + for skill in deduped_skills: + if skill not in attrs_map: + attrs_map[skill] = SkillAttributes(strength=DEFAULT_SKILL_STRENGTH) + + return ExtractedSkills( + skills=deduped_skills, + skill_attrs=attrs_map, + confidence=max(0.0, min(1.0, confidence)), + source=source, + ) + + def _parse_strength(self, value: Any) -> int: + raw: Any = value + if isinstance(value, dict): + raw = value.get("strength") + try: + numeric = int(float(raw)) + except Exception: + numeric = DEFAULT_SKILL_STRENGTH + return max(1, min(5, numeric)) + + def _normalize_skill_name(self, value: str) -> str: + return normalize_skill(value) + + def canonicalize_skill(self, value: str) -> str: + """Public helper for consistent skill normalization across processors.""" + return self._normalize_skill_name(value) diff --git a/apps/worker/src/five08/worker/models.py b/apps/worker/src/five08/worker/models.py index e629e398..e5c4e84d 100644 --- a/apps/worker/src/five08/worker/models.py +++ b/apps/worker/src/five08/worker/models.py @@ -41,10 +41,17 @@ class ExtractedSkills(BaseModel): """Skills extraction response.""" skills: list[str] + skill_attrs: dict[str, "SkillAttributes"] = Field(default_factory=dict) confidence: float = Field(..., ge=0.0, le=1.0) source: str +class SkillAttributes(BaseModel): + """Structured per-skill metadata for CRM persistence.""" + + strength: int = Field(..., ge=1, le=5) + + class SkillsExtractionResult(BaseModel): """End-to-end processing result.""" diff --git a/packages/shared/src/five08/skills.py b/packages/shared/src/five08/skills.py new file mode 100644 index 00000000..ac88c17d --- /dev/null +++ b/packages/shared/src/five08/skills.py @@ -0,0 +1,71 @@ +"""Skill normalization helpers shared across bot and worker services.""" + +from __future__ import annotations + +import re + +# Canonicalization map tuned for Discord-friendly search terms and CRM consistency. +SKILL_ALIASES: dict[str, str] = { + "js": "javascript", + "ts": "typescript", + "node.js": "node", + "nodejs": "node", + "node js": "node", + "golang": "go", + "py": "python", + "postgres": "postgresql", + "k8s": "kubernetes", + "gcp": "google cloud", + "google cloud platform": "google cloud", + "aws": "amazon web services", + "g suite": "google workspace", + "ab testing": "ab testing", + "a/b testing": "ab testing", + "a b testing": "ab testing", + "experimentation": "ab testing", + "product mgmt": "product management", + "product manager": "product management", + "pm": "product management", + "gtm": "go to market", + "go-to-market": "go to market", + "go to market": "go to market", + "seo": "search engine optimization", + "sem": "search engine marketing", + "crm": "customer relationship management", + "ga4": "google analytics", + "google analytics 4": "google analytics", +} + + +def normalize_skill(value: str) -> str: + """Normalize one skill string into a canonical, punctuation-light form.""" + normalized = value.strip().lower() + if not normalized: + return "" + + normalized = re.sub(r"\s+", " ", normalized).strip(" .,_-:/") + if normalized in SKILL_ALIASES: + return SKILL_ALIASES[normalized] + + punctuation_light = re.sub(r"[./_-]+", " ", normalized) + punctuation_light = re.sub(r"\s+", " ", punctuation_light).strip() + if punctuation_light in SKILL_ALIASES: + return SKILL_ALIASES[punctuation_light] + + return punctuation_light or normalized + + +def normalize_skill_list(values: list[str]) -> list[str]: + """Normalize and de-duplicate skills while preserving first-seen order.""" + normalized: list[str] = [] + seen: set[str] = set() + for raw in values: + skill = normalize_skill(raw) + if not skill: + continue + key = skill.casefold() + if key in seen: + continue + seen.add(key) + normalized.append(skill) + return normalized diff --git a/tests/unit/test_resume_profile_processor.py b/tests/unit/test_resume_profile_processor.py index 51599d77..c58fca44 100644 --- a/tests/unit/test_resume_profile_processor.py +++ b/tests/unit/test_resume_profile_processor.py @@ -1,5 +1,6 @@ """Unit tests for resume profile worker processor.""" +import json from unittest.mock import Mock from five08.worker.crm.resume_profile_processor import ResumeProfileProcessor @@ -34,6 +35,10 @@ def test_extract_profile_proposal_filters_508_email() -> None: ) processor.skills_extractor.extract_skills.return_value = ExtractedSkills( skills=["Python", "FastAPI"], + skill_attrs={ + "python": {"strength": 5}, + "fastapi": {"strength": 4}, + }, confidence=0.8, source="gpt-4o-mini", ) @@ -50,6 +55,10 @@ def test_extract_profile_proposal_filters_508_email() -> None: assert result.proposed_updates["cLinkedInUrl"] == "https://linkedin.com/in/new" assert result.proposed_updates["phoneNumber"] == "14155551234" assert result.proposed_updates["skills"] == "Python, FastAPI" + assert "cSkillAttrs" in result.proposed_updates + attrs_payload = json.loads(result.proposed_updates["cSkillAttrs"]) + assert attrs_payload["python"]["strength"] == 5 + assert attrs_payload["fastapi"]["strength"] == 4 assert result.new_skills == ["Python", "FastAPI"] assert any(item.field == "emailAddress" for item in result.skipped) processor.crm.update_contact.assert_called_once() @@ -63,6 +72,62 @@ def test_extract_profile_proposal_filters_508_email() -> None: assert record_kwargs["attachment_id"] == "att-1" +def test_extract_profile_proposal_normalizes_existing_skill_punctuation() -> None: + """Existing punctuation-heavy skills should normalize to search-friendly canonical forms.""" + processor = ResumeProfileProcessor() + processor.crm = Mock() + processor.extractor = Mock() + processor.skills_extractor = Mock() + processor.document_processor = Mock() + processor._record_processing_run = Mock() + + processor.crm.get_contact.return_value = { + "emailAddress": "member@example.com", + "skills": "Node.js, A/B Testing", + "cSkillAttrs": '{"node.js":{"strength":4},"a/b testing":{"strength":2}}', + } + processor.crm.download_attachment.return_value = b"resume-bytes" + processor.document_processor.extract_text.return_value = "resume text" + processor.document_processor.get_content_hash.return_value = "hash-10" + processor.extractor.extract.return_value = ResumeExtractedProfile( + email=None, + github_username=None, + linkedin_url=None, + phone=None, + confidence=0.9, + source="gpt-4o-mini", + ) + processor.skills_extractor.extract_skills.return_value = ExtractedSkills( + skills=["node", "ab testing", "product management"], + skill_attrs={ + "node": {"strength": 5}, + "ab testing": {"strength": 3}, + "product management": {"strength": 4}, + }, + confidence=0.8, + source="gpt-4o-mini", + ) + processor.skills_extractor.canonicalize_skill.side_effect = lambda v: { + "node.js": "node", + "a/b testing": "ab testing", + "node": "node", + "ab testing": "ab testing", + }.get(str(v).strip().lower(), str(v).strip().lower()) + + result = processor.extract_profile_proposal( + contact_id="contact-10", + attachment_id="att-10", + filename="resume.pdf", + ) + + assert result.success is True + assert result.proposed_updates["skills"] == "node, ab testing, product management" + attrs_payload = json.loads(result.proposed_updates["cSkillAttrs"]) + assert attrs_payload["node"]["strength"] == 5 + assert attrs_payload["ab testing"]["strength"] == 3 + assert attrs_payload["product management"]["strength"] == 4 + + def test_apply_profile_updates_adds_discord_and_filters_email() -> None: """Apply should include Discord link values and prevent @508.dev email writes.""" processor = ResumeProfileProcessor() @@ -75,6 +140,7 @@ def test_apply_profile_updates_adds_discord_and_filters_email() -> None: "cGitHubUsername": "new-gh", "phoneNumber": "14155551234", "skills": "Python, FastAPI", + "cSkillAttrs": '{"python":{"strength":4},"fastapi":{"strength":3}}', }, link_discord={"user_id": "123", "username": "member#0001"}, ) @@ -86,6 +152,10 @@ def test_apply_profile_updates_adds_discord_and_filters_email() -> None: assert update_payload["cGitHubUsername"] == "new-gh" assert update_payload["phoneNumber"] == "14155551234" assert update_payload["skills"] == "Python, FastAPI" + assert update_payload["cSkillAttrs"] == { + "fastapi": {"strength": 3}, + "python": {"strength": 4}, + } assert update_payload["cDiscordUserID"] == "123" assert update_payload["cDiscordUsername"] == "member#0001 (ID: 123)" diff --git a/tests/unit/test_shared_skills.py b/tests/unit/test_shared_skills.py new file mode 100644 index 00000000..3a4c85da --- /dev/null +++ b/tests/unit/test_shared_skills.py @@ -0,0 +1,19 @@ +"""Unit tests for shared skill normalization helpers.""" + +from five08.skills import normalize_skill, normalize_skill_list + + +def test_normalize_skill_prefers_discord_friendly_canonical_forms() -> None: + """Canonical outputs should avoid punctuation-heavy variants and initials.""" + assert normalize_skill("Node.js") == "node" + assert normalize_skill("A/B Testing") == "ab testing" + assert normalize_skill("GTM") == "go to market" + assert normalize_skill("CRM") == "customer relationship management" + assert normalize_skill("SEO") == "search engine optimization" + + +def test_normalize_skill_list_dedupes_after_aliasing() -> None: + """List normalization should dedupe across equivalent aliases.""" + normalized = normalize_skill_list(["node.js", "node", "A/B Testing", "ab testing"]) + + assert normalized == ["node", "ab testing"] diff --git a/tests/unit/test_skills_extractor.py b/tests/unit/test_skills_extractor.py index 79eb3bce..c1bc75ae 100644 --- a/tests/unit/test_skills_extractor.py +++ b/tests/unit/test_skills_extractor.py @@ -12,6 +12,8 @@ def test_heuristic_extract_includes_two_letter_skill_go() -> None: assert "go" in result.skills assert "docker" in result.skills + assert result.skill_attrs["go"].strength == 3 + assert result.skill_attrs["docker"].strength == 3 def test_heuristic_extractor_includes_two_letter_go_skill() -> None: @@ -21,3 +23,34 @@ def test_heuristic_extractor_includes_two_letter_go_skill() -> None: assert "go" in result.skills assert "python" in result.skills + + +def test_normalize_extracted_payload_canonicalizes_and_clamps_strength() -> None: + """LLM payload normalization should map aliases and clamp strengths to 1-5.""" + extractor = SkillsExtractor() + + result = extractor._normalize_extracted_payload( + skills_value=["JS", " PM ", "A/B Testing", "Node.js", "Go-To-Market"], + skill_attrs_value={ + "javascript": {"strength": 9}, + "product management": {"strength": 4}, + "a/b testing": {"strength": 0}, + "node.js": {"strength": 2}, + "go-to-market": {"strength": 4}, + }, + confidence=0.8, + source="model", + ) + + assert result.skills == [ + "ab testing", + "go to market", + "javascript", + "node", + "product management", + ] + assert result.skill_attrs["javascript"].strength == 5 + assert result.skill_attrs["product management"].strength == 4 + assert result.skill_attrs["ab testing"].strength == 1 + assert result.skill_attrs["node"].strength == 2 + assert result.skill_attrs["go to market"].strength == 4 diff --git a/tests/unit/test_worker_api.py b/tests/unit/test_worker_api.py index c5a2992a..6a7f68c1 100644 --- a/tests/unit/test_worker_api.py +++ b/tests/unit/test_worker_api.py @@ -173,7 +173,8 @@ async def test_resume_extract_handler_enqueues_job( """Resume extract endpoint should enqueue extraction job.""" monkeypatch.setattr(api.settings, "resume_extractor_version", "v7") monkeypatch.setattr(api.settings, "openai_api_key", "key") - monkeypatch.setattr(api.settings, "openai_model", "gpt-test") + monkeypatch.setattr(api.settings, "openai_base_url", None) + monkeypatch.setattr(api.settings, "resume_ai_model", "gpt-test") app_obj = web.Application() app_obj[api.QUEUE_KEY] = Mock() @@ -270,11 +271,22 @@ def test_resume_extract_model_name_uses_heuristic_without_api_key( ) -> None: """Model identity should be heuristic when OpenAI key is absent.""" monkeypatch.setattr(api.settings, "openai_api_key", None) - monkeypatch.setattr(api.settings, "openai_model", "gpt-test") + monkeypatch.setattr(api.settings, "resume_ai_model", "gpt-test") assert api._resume_extract_model_name() == "heuristic" +def test_resume_extract_model_name_prefixes_openrouter_model( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """OpenRouter base URL should map plain resume model to openai/.""" + monkeypatch.setattr(api.settings, "openai_api_key", "key") + monkeypatch.setattr(api.settings, "openai_base_url", "https://openrouter.ai/api/v1") + monkeypatch.setattr(api.settings, "resume_ai_model", "gpt-4o-mini") + + assert api._resume_extract_model_name() == "openai/gpt-4o-mini" + + @pytest.mark.asyncio async def test_sync_people_handler_enqueues_full_sync( auth_headers: dict[str, str], From 531473bb84f99231349d1763aa8bcae2a20e17ee Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 16:01:23 +0800 Subject: [PATCH 5/8] chore: run mypy hook via repo script --- .pre-commit-config.yaml | 17 +++++------------ 1 file changed, 5 insertions(+), 12 deletions(-) diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml index fd0466d8..cb70429c 100644 --- a/.pre-commit-config.yaml +++ b/.pre-commit-config.yaml @@ -8,18 +8,11 @@ repos: # Run the formatter on changed files only - id: ruff-format - - repo: https://github.com/pre-commit/mirrors-mypy - rev: v1.14.0 + - repo: local hooks: - id: mypy - args: [--strict] + name: mypy + entry: ./scripts/mypy.sh + language: system + pass_filenames: false files: ^(apps/discord_bot/src/five08/|apps/worker/src/five08/|packages/shared/src/five08/) - additional_dependencies: - - aiohttp>=3.13.1 - - discord.py~=2.6.0 - - pydantic~=2.10 - - pydantic-settings~=2.8 - - redis>=6.4.0 - - requests~=2.31 - - rq>=2.6.0 - - types-requests>=2.32.4.20250913 From c1a1c86a4eb3637b5bef3eab2e2f425a2c19698a Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 16:03:22 +0800 Subject: [PATCH 6/8] test: align worker processor mock with canonical skills --- tests/unit/test_worker_processor.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/tests/unit/test_worker_processor.py b/tests/unit/test_worker_processor.py index 65537523..47e5ce56 100644 --- a/tests/unit/test_worker_processor.py +++ b/tests/unit/test_worker_processor.py @@ -21,14 +21,17 @@ def test_process_contact_skills_merges_and_updates() -> None: processor.espocrm_client.download_attachment.return_value = b"file-content" processor.document_processor.extract_text.return_value = "Python FastAPI Docker" processor.skills_extractor.extract_skills.return_value = ExtractedSkills( - skills=["Python", "FastAPI", "Docker"], + skills=["python", "fastapi", "docker"], confidence=0.9, source="heuristic", ) + processor.skills_extractor.canonicalize_skill.side_effect = ( + lambda value: str(value).strip().lower() + ) processor.espocrm_client.update_contact_skills.return_value = True result = processor.process_contact_skills("contact-1") assert result.success is True - assert sorted(result.new_skills) == ["Docker", "FastAPI"] - assert set(result.updated_skills) == {"Python", "Redis", "FastAPI", "Docker"} + assert sorted(result.new_skills) == ["docker", "fastapi"] + assert set(result.updated_skills) == {"python", "redis", "fastapi", "docker"} From 97a9a6e73c30f94401582639652d713da9479898 Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 16:05:54 +0800 Subject: [PATCH 7/8] fix: address resume apply review feedback --- AGENTS.md | 2 +- DEVELOPMENT.md | 2 +- .../src/five08/discord_bot/cogs/crm.py | 125 ++++++++++++------ .../worker/crm/resume_profile_processor.py | 22 +-- packages/shared/src/five08/settings.py | 7 +- 5 files changed, 99 insertions(+), 59 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index a0172dd8..08b1fbab 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -93,7 +93,7 @@ async def setup(bot: commands.Bot) -> None: - Add shared env/config in `packages/shared/src/five08/settings.py`. - Add service-specific settings in local service `config.py` by subclassing shared settings. - Keep secrets in env vars, not code. -- For Discord CRM audit writes, use `AUDIT_API_BASE_URL` and shared `WEBHOOK_SHARED_SECRET`. +- For Discord CRM audit writes, use `AUDIT_API_BASE_URL` and shared `API_SHARED_SECRET`. ## Agent Guidelines diff --git a/DEVELOPMENT.md b/DEVELOPMENT.md index 76a50196..fa012fed 100644 --- a/DEVELOPMENT.md +++ b/DEVELOPMENT.md @@ -123,7 +123,7 @@ Use `.env.example` as source of truth. Key categories: - Shared queue/runtime: `REDIS_URL`, `REDIS_QUEUE_NAME`, `POSTGRES_URL`, `JOB_MAX_ATTEMPTS`, `JOB_RETRY_BASE_SECONDS`, `JOB_RETRY_MAX_SECONDS`, `LOG_LEVEL`, webhook settings - Bot credentials/integrations: Discord, email, Espo, Kimai -- Discord CRM audit writer: `AUDIT_API_BASE_URL`, `AUDIT_API_TIMEOUT_SECONDS` (plus shared `WEBHOOK_SHARED_SECRET`) +- Discord CRM audit writer: `AUDIT_API_BASE_URL`, `AUDIT_API_TIMEOUT_SECONDS` (plus shared `API_SHARED_SECRET`) - Worker controls: `WORKER_NAME`, `WORKER_QUEUE_NAMES`, `WORKER_BURST` - Worker CRM processing: `MAX_ATTACHMENTS_PER_CONTACT`, `MAX_FILE_SIZE_MB`, `ALLOWED_FILE_TYPES`, `RESUME_KEYWORDS`, `OPENAI_API_KEY`, `OPENAI_BASE_URL`, `OPENAI_MODEL`, `RESUME_EXTRACTOR_VERSION` - Resume upload UX wiring: `WORKER_API_BASE_URL` on bot, `CRM_LINKEDIN_FIELD` on worker. diff --git a/apps/discord_bot/src/five08/discord_bot/cogs/crm.py b/apps/discord_bot/src/five08/discord_bot/cogs/crm.py index bd8dee2e..5484fc35 100644 --- a/apps/discord_bot/src/five08/discord_bot/cogs/crm.py +++ b/apps/discord_bot/src/five08/discord_bot/cogs/crm.py @@ -389,7 +389,24 @@ async def confirm_updates( button: discord.ui.Button["ResumeUpdateConfirmationView"], ) -> None: """Apply confirmed updates through the worker.""" - await interaction.response.defer(ephemeral=True) + await interaction.response.defer(thinking=True, ephemeral=True) + + def _audit_apply_event(result: str, metadata: dict[str, Any]) -> None: + try: + self.crm_cog._audit_command( + interaction=interaction, + action="crm.upload_resume.apply", + result=result, + metadata=metadata, + resource_type="crm_contact", + resource_id=self.contact_id, + ) + except Exception as exc: + logger.warning( + "Failed to write resume apply audit event for contact_id=%s: %s", + self.contact_id, + exc, + ) try: apply_job_id = await self.crm_cog._enqueue_resume_apply_job( @@ -399,21 +416,21 @@ async def confirm_updates( ) except Exception as exc: logger.error("Failed to enqueue resume apply job: %s", exc) - self.crm_cog._audit_command( - interaction=interaction, - action="crm.upload_resume", - result="error", - metadata={ + _audit_apply_event( + "error", + { + "contact_id": self.contact_id, "stage": "apply_enqueue", "error": str(exc), + "updated_fields": [], "proposed_updates_count": len(self.proposed_updates), "link_member_requested": bool(self.link_discord), + "link_discord_applied": None, }, - resource_type="crm_contact", - resource_id=self.contact_id, ) await interaction.followup.send( - "โŒ Failed to enqueue CRM apply job. Please try again." + "โŒ Failed to enqueue CRM apply job. Please try again.", + ephemeral=True, ) return @@ -421,53 +438,81 @@ async def confirm_updates( "๐Ÿ› ๏ธ Applying confirmed updates to CRM...", ephemeral=True, ) - apply_result = await self.crm_cog._wait_for_worker_job_result(apply_job_id) + try: + apply_result = await self.crm_cog._wait_for_worker_job_result(apply_job_id) + except Exception as exc: + logger.error( + "Worker polling failed for apply_job_id=%s contact_id=%s error=%s", + apply_job_id, + self.contact_id, + exc, + ) + _audit_apply_event( + "error", + { + "contact_id": self.contact_id, + "stage": "apply_polling_failed", + "job_id": apply_job_id, + "error": str(exc), + "updated_fields": [], + "link_discord_applied": None, + }, + ) + await interaction.followup.send( + "โš ๏ธ Resume apply polling failed. Please retry or check CRM manually.", + ephemeral=True, + ) + return if not apply_result: - self.crm_cog._audit_command( - interaction=interaction, - action="crm.upload_resume", - result="error", - metadata={ + _audit_apply_event( + "error", + { + "contact_id": self.contact_id, "stage": "apply_timeout", "job_id": apply_job_id, + "updated_fields": [], + "link_discord_applied": None, }, - resource_type="crm_contact", - resource_id=self.contact_id, ) await interaction.followup.send( - "โš ๏ธ Timed out waiting for apply job. Please check again shortly." + "โš ๏ธ Timed out waiting for apply job. Please check again shortly.", + ephemeral=True, ) return status = str(apply_result.get("status", "unknown")) + result = apply_result.get("result") + updated_fields: list[str] = [] + link_discord_applied: bool | None = None + if isinstance(result, dict): + raw_fields = result.get("updated_fields") + if isinstance(raw_fields, list): + updated_fields = [str(field) for field in raw_fields] + raw_link_applied = result.get("link_discord_applied") + if isinstance(raw_link_applied, bool): + link_discord_applied = raw_link_applied + if status != "succeeded": - self.crm_cog._audit_command( - interaction=interaction, - action="crm.upload_resume", - result="error", - metadata={ + _audit_apply_event( + "error", + { + "contact_id": self.contact_id, "stage": "apply_failed", "job_id": apply_job_id, "job_status": status, "last_error": str(apply_result.get("last_error", "")), + "updated_fields": updated_fields, + "link_discord_applied": link_discord_applied, }, - resource_type="crm_contact", - resource_id=self.contact_id, ) await interaction.followup.send( f"โŒ Apply job failed (status: {status}). " - f"Error: {apply_result.get('last_error') or 'Unknown error'}" + f"Error: {apply_result.get('last_error') or 'Unknown error'}", + ephemeral=True, ) return - result = apply_result.get("result") - updated_fields: list[str] = [] - if isinstance(result, dict): - raw_fields = result.get("updated_fields") - if isinstance(raw_fields, list): - updated_fields = [str(field) for field in raw_fields] - embed = discord.Embed( title="โœ… CRM Updated", description=f"Applied updates for **{self.contact_name}**.", @@ -480,21 +525,19 @@ async def confirm_updates( ) profile_url = f"{self.crm_cog.base_url}/#Contact/view/{self.contact_id}" embed.add_field(name="๐Ÿ”— CRM Profile", value=f"[View in CRM]({profile_url})") - self.crm_cog._audit_command( - interaction=interaction, - action="crm.upload_resume", - result="success", - metadata={ + _audit_apply_event( + "success", + { + "contact_id": self.contact_id, "stage": "apply_succeeded", "job_id": apply_job_id, "updated_fields": updated_fields, "proposed_updates_count": len(self.proposed_updates), "link_member_requested": bool(self.link_discord), + "link_discord_applied": link_discord_applied, }, - resource_type="crm_contact", - resource_id=self.contact_id, ) - await interaction.followup.send(embed=embed) + await interaction.followup.send(embed=embed, ephemeral=True) for item in self.children: if isinstance(item, discord.ui.Button): diff --git a/apps/worker/src/five08/worker/crm/resume_profile_processor.py b/apps/worker/src/five08/worker/crm/resume_profile_processor.py index c845022f..9697c4bf 100644 --- a/apps/worker/src/five08/worker/crm/resume_profile_processor.py +++ b/apps/worker/src/five08/worker/crm/resume_profile_processor.py @@ -146,19 +146,19 @@ def _heuristic_extract(self, resume_text: str) -> ResumeExtractedProfile: def _build_prompt(self, resume_text: str) -> str: snippet = resume_text[:12000] return ( - "Extract candidate contact fields from this resume.\\n" - "Return JSON with exact keys and no extras:\\n" + "Extract candidate contact fields from this resume.\n" + "Return JSON with exact keys and no extras:\n" '{"email": string|null, "github_username": string|null, ' '"linkedin_url": string|null, "phone": string|null, ' - '"confidence": number}\\n' - "Rules:\\n" - "- prefer explicit values from header/contact sections\\n" - "- for github_username return username only (no URL, no @)\\n" - "- for linkedin_url return full linkedin profile URL when available\\n" - "- for phone return digits with optional leading +\\n" - "- use null for unknown/ambiguous fields\\n" - "- confidence is 0-1 for overall extraction reliability\\n\\n" - f"Resume:\\n{snippet}" + '"confidence": number}\n' + "Rules:\n" + "- prefer explicit values from header/contact sections\n" + "- for github_username return username only (no URL, no @)\n" + "- for linkedin_url return full linkedin profile URL when available\n" + "- for phone return digits with optional leading +\n" + "- use null for unknown/ambiguous fields\n" + "- confidence is 0-1 for overall extraction reliability\n\n" + f"Resume:\n{snippet}" ) def _parse_json(self, content: str) -> dict[str, Any]: diff --git a/packages/shared/src/five08/settings.py b/packages/shared/src/five08/settings.py index feb11561..55f3d1ae 100644 --- a/packages/shared/src/five08/settings.py +++ b/packages/shared/src/five08/settings.py @@ -1,6 +1,6 @@ """Shared configuration settings across services.""" -from pydantic import AliasChoices, Field, model_validator +from pydantic import model_validator from pydantic_settings import BaseSettings, SettingsConfigDict @@ -35,10 +35,7 @@ class SharedSettings(BaseSettings): webhook_ingest_host: str = "0.0.0.0" webhook_ingest_port: int = 8090 - api_shared_secret: str | None = Field( - default=None, - validation_alias=AliasChoices("API_SHARED_SECRET", "WEBHOOK_SHARED_SECRET"), - ) + api_shared_secret: str | None = None model_config = SettingsConfigDict(env_file=".env", extra="ignore") From 47afdabca5ad0aa63123c02692befadb014e17c6 Mon Sep 17 00:00:00 2001 From: Michael Wu Date: Sat, 21 Feb 2026 16:08:10 +0800 Subject: [PATCH 8/8] style: format worker processor test --- tests/unit/test_worker_processor.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/unit/test_worker_processor.py b/tests/unit/test_worker_processor.py index 47e5ce56..6966489d 100644 --- a/tests/unit/test_worker_processor.py +++ b/tests/unit/test_worker_processor.py @@ -25,8 +25,8 @@ def test_process_contact_skills_merges_and_updates() -> None: confidence=0.9, source="heuristic", ) - processor.skills_extractor.canonicalize_skill.side_effect = ( - lambda value: str(value).strip().lower() + processor.skills_extractor.canonicalize_skill.side_effect = lambda value: ( + str(value).strip().lower() ) processor.espocrm_client.update_contact_skills.return_value = True