from __future__ import annotations from collections import OrderedDict from datetime import date, datetime, timedelta, timezone from typing import Any from sqlalchemy import or_, select from sqlalchemy.orm import Session, selectinload from app.modules.clients.models import Client from app.modules.consultants.models import ClientConsultantLink, ConsultantProfile, ConsultantServiceRequest from app.modules.documents.models import EngagementDocument, PermanentClientDocument from app.modules.services.models import ( ClientServiceSubscription, ClientServiceTaskInstance, ServiceCatalogue, ServiceTaskComment, ) CONSULTANT_BOARD_COLUMNS = OrderedDict( [ ("assigned", "Assigned"), ("awaiting_documents", "Awaiting Documents"), ("in_progress", "In Progress"), ("submitted", "Submitted"), ("accepted", "Accepted"), ("closed", "Closed"), ] ) def _allowed_client_ids(db: Session, *, consultant: ConsultantProfile, require_communications: bool = False) -> list[int]: query = select(ClientConsultantLink.client_id).where( ClientConsultantLink.tenant_id == consultant.tenant_id, ClientConsultantLink.consultant_id == consultant.id, ClientConsultantLink.is_active.is_(True), ) if require_communications: query = query.where(ClientConsultantLink.can_view_communications.is_(True)) return [int(x) for x in db.execute(query).scalars().all()] def _board_key_for_status(status: str | None) -> str: value = (status or "pending").strip().lower() if value in {"blocked", "awaiting_documents", "document_pending", "clarification_required"}: return "awaiting_documents" if value in {"in_progress", "under_process", "processing", "started"}: return "in_progress" if value in {"pending_review", "ready_for_review", "submitted", "completed"}: return "submitted" if value in {"approved", "accepted"}: return "accepted" if value in {"closed", "locked", "cancelled", "inactive"}: return "closed" return "assigned" def _matches_search(*values: Any, q: str = "") -> bool: term = (q or "").strip().lower() if not term: return True return any(term in str(v or "").lower() for v in values) def get_consultant_work_board(db: Session, *, consultant: ConsultantProfile, q: str = "", status: str = "") -> dict: links = db.execute(select(ClientConsultantLink).where(ClientConsultantLink.tenant_id == consultant.tenant_id, ClientConsultantLink.consultant_id == consultant.id, ClientConsultantLink.is_active.is_(True), ClientConsultantLink.can_view_engagements.is_(True), ClientConsultantLink.can_view_task_status.is_(True))).scalars().all() links_by_client: dict[int, list[ClientConsultantLink]] = {} for link in links: links_by_client.setdefault(int(link.client_id), []).append(link) columns = {key: {"label": label, "items": []} for key, label in CONSULTANT_BOARD_COLUMNS.items()} items=[] if links_by_client: rows=db.execute(select(ClientServiceTaskInstance, ClientServiceSubscription, Client, ServiceCatalogue).join(ClientServiceSubscription, ClientServiceSubscription.id==ClientServiceTaskInstance.subscription_id).join(Client, Client.id==ClientServiceTaskInstance.client_id).join(ServiceCatalogue, ServiceCatalogue.id==ClientServiceTaskInstance.service_catalogue_id).where(ClientServiceTaskInstance.tenant_id==consultant.tenant_id, ClientServiceTaskInstance.client_id.in_(list(links_by_client)), ClientServiceTaskInstance.is_active.is_(True)).order_by(ClientServiceTaskInstance.internal_target_date.asc(), ClientServiceTaskInstance.id.desc()).limit(500)).all() for task, subscription, client, catalogue in rows: allowed=next((ln for ln in links_by_client[int(client.id)] if ln.service_catalogue_id in (None, task.service_catalogue_id)),None) if not allowed or not _matches_search(client.client_name, getattr(client,"client_code",""), catalogue.service_name, task.task_name, q=q): continue key=_board_key_for_status(task.status) if status and key!=status: continue latest=db.execute(select(ServiceTaskComment).where(ServiceTaskComment.task_instance_id==task.id, ServiceTaskComment.visibility=="consultant", ServiceTaskComment.is_deleted.is_(False)).order_by(ServiceTaskComment.created_at_utc.desc()).limit(1)).scalars().first() item={"comment":latest,"task":task,"subscription":subscription,"client":client,"catalogue":catalogue,"board_key":key,"link":allowed,"is_overdue":bool(task.internal_target_date and task.internal_target_date dict | None: client_ids = _allowed_client_ids(db, consultant=consultant, require_communications=False) if not client_ids: return None row = db.execute( select(ClientServiceTaskInstance, ClientServiceSubscription, Client, ServiceCatalogue) .join(ClientServiceSubscription, ClientServiceSubscription.id == ClientServiceTaskInstance.subscription_id) .join(Client, Client.id == ClientServiceTaskInstance.client_id) .join(ServiceCatalogue, ServiceCatalogue.id == ClientServiceTaskInstance.service_catalogue_id) .where( ClientServiceTaskInstance.id == task_id, ClientServiceTaskInstance.tenant_id == consultant.tenant_id, ClientServiceTaskInstance.client_id.in_(client_ids), ClientServiceTaskInstance.is_active.is_(True), ) ).first() if not row: return None task, subscription, client, catalogue = row link = db.execute(select(ClientConsultantLink).where(ClientConsultantLink.tenant_id == consultant.tenant_id, ClientConsultantLink.client_id == client.id, ClientConsultantLink.consultant_id == consultant.id, ClientConsultantLink.is_active.is_(True), ClientConsultantLink.can_view_engagements.is_(True), or_(ClientConsultantLink.service_catalogue_id.is_(None), ClientConsultantLink.service_catalogue_id == task.service_catalogue_id))).scalars().first() if not link: return None timeline = db.execute( select(ServiceTaskComment) .options(selectinload(ServiceTaskComment.created_by)) .where( ServiceTaskComment.tenant_id == consultant.tenant_id, ServiceTaskComment.task_instance_id == task.id, ServiceTaskComment.visibility == "consultant", ServiceTaskComment.is_deleted.is_(False), ) .order_by(ServiceTaskComment.created_at_utc.asc(), ServiceTaskComment.id.asc()) ).scalars().all() engagement_documents = db.execute( select(EngagementDocument) .options(selectinload(EngagementDocument.versions)) .where( EngagementDocument.tenant_id == consultant.tenant_id, EngagementDocument.client_id == client.id, EngagementDocument.engagement_id == subscription.id, EngagementDocument.is_deleted.is_(False), ) .order_by(EngagementDocument.document_type.asc(), EngagementDocument.title.asc()) .limit(50) ).scalars().all() if not (link.can_view_final_documents or link.can_upload_documents): engagement_documents = [] permanent_documents = db.execute( select(PermanentClientDocument) .options(selectinload(PermanentClientDocument.versions)) .where( PermanentClientDocument.tenant_id == consultant.tenant_id, PermanentClientDocument.client_id == client.id, PermanentClientDocument.is_deleted.is_(False), ) .order_by(PermanentClientDocument.category.asc(), PermanentClientDocument.title.asc()) .limit(50) ).scalars().all() if not link.can_view_permanent_documents: permanent_documents = [] return { "link": link, "task": task, "subscription": subscription, "client": client, "catalogue": catalogue, "timeline": timeline, "engagement_documents": engagement_documents, "permanent_documents": permanent_documents, } def get_consultant_document_centre(db: Session, *, consultant: ConsultantProfile, q: str = "") -> dict: links=db.execute(select(ClientConsultantLink).where(ClientConsultantLink.tenant_id==consultant.tenant_id,ClientConsultantLink.consultant_id==consultant.id,ClientConsultantLink.is_active.is_(True))).scalars().all() engagement_ids={int(x.client_id) for x in links if x.can_view_final_documents or x.can_upload_documents} permanent_ids={int(x.client_id) for x in links if x.can_view_permanent_documents} engagement_documents=[]; permanent_documents=[] if engagement_ids: stmt=select(EngagementDocument).options(selectinload(EngagementDocument.client),selectinload(EngagementDocument.engagement),selectinload(EngagementDocument.versions)).where(EngagementDocument.tenant_id==consultant.tenant_id,EngagementDocument.client_id.in_(engagement_ids),EngagementDocument.is_deleted.is_(False)) if q.strip(): term=f"%{q.strip()}%"; stmt=stmt.where(or_(EngagementDocument.title.ilike(term),EngagementDocument.document_type.ilike(term),EngagementDocument.document_code.ilike(term))) engagement_documents=db.execute(stmt.order_by(EngagementDocument.created_at_utc.desc()).limit(200)).scalars().all() if permanent_ids: stmt=select(PermanentClientDocument).options(selectinload(PermanentClientDocument.client),selectinload(PermanentClientDocument.versions)).where(PermanentClientDocument.tenant_id==consultant.tenant_id,PermanentClientDocument.client_id.in_(permanent_ids),PermanentClientDocument.is_deleted.is_(False)) if q.strip(): term=f"%{q.strip()}%"; stmt=stmt.where(or_(PermanentClientDocument.title.ilike(term),PermanentClientDocument.category.ilike(term),PermanentClientDocument.document_code.ilike(term))) permanent_documents=db.execute(stmt.order_by(PermanentClientDocument.created_at_utc.desc()).limit(200)).scalars().all() return {"engagement_documents":engagement_documents,"permanent_documents":permanent_documents,"total":len(engagement_documents)+len(permanent_documents)} def consultant_can_execute_task(*, consultant: ConsultantProfile, task: ClientServiceTaskInstance) -> bool: return bool(task.execution_mode == "consultant" and task.assigned_consultant_id == consultant.id and task.is_active and not task.is_locked) def update_consultant_assignment(db: Session, *, consultant: ConsultantProfile, task: ClientServiceTaskInstance, action: str, note: str, user_id: int) -> tuple[bool, str | None]: if not consultant_can_execute_task(consultant=consultant, task=task): return False, "This task is not assigned to your consultant profile." action = (action or "").strip().lower() note = (note or "").strip() current = (task.consultant_assignment_status or "not_applicable").lower() now = datetime.now(timezone.utc) message = None if action == "accept" and current == "offered": task.consultant_assignment_status = "accepted" task.consultant_accepted_at_utc = now task.consultant_declined_at_utc = None task.consultant_decline_reason = None task.status = "pending" message = "Consultant accepted the assignment." elif action == "decline" and current == "offered": if not note: return False, "Decline reason is required." task.consultant_assignment_status = "declined" task.consultant_declined_at_utc = now task.consultant_decline_reason = note task.status = "pending" message = f"Consultant declined the assignment: {note}" elif action == "start" and current in {"accepted", "rework_required"}: task.consultant_assignment_status = "in_progress" task.consultant_started_at_utc = task.consultant_started_at_utc or now task.status = "in_progress" if current == "rework_required": task.rework_status = "in_progress" message = "Consultant started work on the assignment." elif action == "submit" and current in {"accepted", "in_progress", "rework_required"}: if not note: return False, "Submission note is required." task.consultant_assignment_status = "submitted" task.consultant_submitted_at_utc = now task.consultant_submission_note = note task.status = "pending_review" if task.rework_status in {"requested", "in_progress"}: task.rework_status = "resolved" task.rework_resolved_at_utc = now message = f"Consultant submitted the assignment for firm review: {note}" else: return False, "This action is not allowed for the current assignment status." task.updated_by_user_id = user_id db.add(ServiceTaskComment(tenant_id=task.tenant_id, branch_id=task.branch_id, subscription_id=task.subscription_id, task_instance_id=task.id, comment_type="consultant_execution", visibility="consultant", message=message, created_by_user_id=user_id)) return True, None