from __future__ import annotations from datetime import date, datetime, timezone from sqlalchemy import func, or_, select from sqlalchemy.orm import Session, selectinload from app.modules.clients.models import Client from app.modules.documents.models import EngagementDocument from app.modules.core.iam.models import User from app.modules.services.client_services import quality_required_for_engagement, QUALITY_APPROVED from app.modules.services.models import ( ClientServiceSubscription, ClientServiceTaskInstance, EngagementClosureChecklist, ServiceTaskComment, FirmServiceTaskTemplate, ServiceCatalogue, ) TASK_STATUSES = [ ("pending", "Pending"), ("in_progress", "In Progress"), ("completed", "Completed"), ("blocked", "Blocked"), ("not_applicable", "Not Applicable"), ("cancelled", "Cancelled"), ] OPEN_TASK_STATUSES = {"pending", "in_progress", "blocked"} CLOSED_TASK_STATUSES = {"completed", "not_applicable", "cancelled"} TASK_COMMENT_TYPES = [ ("internal_note", "Internal Note"), ("staff_work_note", "Staff Work Note"), ("client_clarification", "Client Clarification"), ("consultant_clarification", "Consultant Clarification"), ("manager_review_note", "Manager Review Note"), ("partner_review_note", "Partner Review Note"), ("review_partner_review_note", "Review Partner Review Note"), ("rework_note", "Rework Note"), ] AQMM_REVIEW_STATUSES = [ ("not_required", "Not Required"), ("pending", "Pending"), ("reviewed", "Reviewed"), ("rework_required", "Rework Required"), ] AQMM_EVIDENCE_STATUSES = [ ("not_required", "Not Required"), ("pending", "Pending"), ("uploaded", "Uploaded"), ("accepted", "Accepted"), ("rejected", "Rejected / Insufficient"), ] TASK_COMMENT_VISIBILITIES = [ ("internal", "Internal"), ("client", "Client"), ("consultant", "Consultant"), ] TASK_PRIORITIES = [ ("low", "Low"), ("normal", "Normal"), ("high", "High"), ("urgent", "Urgent"), ] def _normalise_status(status: str | None) -> str: allowed = {code for code, _label in TASK_STATUSES} value = (status or "pending").strip().lower() return value if value in allowed else "pending" def _normalise_priority(priority: str | None) -> str: allowed = {code for code, _label in TASK_PRIORITIES} value = (priority or "normal").strip().lower() return value if value in allowed else "normal" def _normalise_comment_type(comment_type: str | None) -> str: allowed = {code for code, _label in TASK_COMMENT_TYPES} value = (comment_type or "internal_note").strip().lower() return value if value in allowed else "internal_note" def _normalise_visibility(visibility: str | None) -> str: allowed = {code for code, _label in TASK_COMMENT_VISIBILITIES} value = (visibility or "internal").strip().lower() return value if value in allowed else "internal" def parse_date_value(value: str | None) -> date | None: text = (value or "").strip() if not text: return None return date.fromisoformat(text) def _default_assignee_for_template(subscription: ClientServiceSubscription, template: FirmServiceTaskTemplate) -> int | None: role = (template.default_role_name or "").strip().lower() if "partner" in role: return subscription.assigned_partner_user_id if "manager" in role: return subscription.assigned_manager_user_id if "staff" in role or "employee" in role: return subscription.assigned_staff_user_id return subscription.assigned_staff_user_id or subscription.assigned_manager_user_id or subscription.assigned_partner_user_id def get_subscription_for_execution(db: Session, *, tenant_id: int, subscription_id: int) -> ClientServiceSubscription | None: return db.execute( select(ClientServiceSubscription) .options( selectinload(ClientServiceSubscription.client), selectinload(ClientServiceSubscription.catalogue), selectinload(ClientServiceSubscription.assigned_partner), selectinload(ClientServiceSubscription.assigned_manager), selectinload(ClientServiceSubscription.assigned_staff), ) .where( ClientServiceSubscription.id == subscription_id, ClientServiceSubscription.tenant_id == tenant_id, ) ).scalar_one_or_none() def list_subscription_execution_payload(db: Session, *, tenant_id: int, branch_id: int | None = None, financial_year: str | None = None, q: str = ""): query = ( select(ClientServiceSubscription) .options( selectinload(ClientServiceSubscription.client), selectinload(ClientServiceSubscription.catalogue), selectinload(ClientServiceSubscription.assigned_partner), selectinload(ClientServiceSubscription.assigned_manager), selectinload(ClientServiceSubscription.assigned_staff), ) .where( ClientServiceSubscription.tenant_id == tenant_id, ClientServiceSubscription.is_active.is_(True), ClientServiceSubscription.status == "active", ClientServiceSubscription.is_locked.is_(False), ) ) if financial_year: query = query.where(ClientServiceSubscription.financial_year == financial_year.strip()) if branch_id: query = query.where(ClientServiceSubscription.branch_id == branch_id) if q.strip(): term = f"%{q.strip()}%" query = ( query.join(Client, Client.id == ClientServiceSubscription.client_id) .join(ServiceCatalogue, ServiceCatalogue.id == ClientServiceSubscription.service_catalogue_id) .where( or_( Client.client_name.ilike(term), Client.client_code.ilike(term), ServiceCatalogue.service_code.ilike(term), ServiceCatalogue.service_name.ilike(term), ) ) ) rows = db.execute(query.order_by(ClientServiceSubscription.id.desc())).scalars().all() payload = [] for sub in rows: total = db.execute( select(func.count(ClientServiceTaskInstance.id)).where( ClientServiceTaskInstance.subscription_id == sub.id, ClientServiceTaskInstance.is_active.is_(True), ) ).scalar_one() completed = db.execute( select(func.count(ClientServiceTaskInstance.id)).where( ClientServiceTaskInstance.subscription_id == sub.id, ClientServiceTaskInstance.is_active.is_(True), ClientServiceTaskInstance.status == "completed", ) ).scalar_one() open_tasks = db.execute( select(func.count(ClientServiceTaskInstance.id)).where( ClientServiceTaskInstance.subscription_id == sub.id, ClientServiceTaskInstance.is_active.is_(True), ClientServiceTaskInstance.status.in_(list(OPEN_TASK_STATUSES)), ) ).scalar_one() payload.append({"subscription": sub, "total_tasks": total, "completed_tasks": completed, "open_tasks": open_tasks}) return payload def _default_internal_target_date(subscription: ClientServiceSubscription, template: FirmServiceTaskTemplate) -> date | None: # Phase 4B keeps task target dates internal. Existing task templates do not yet have # an offset field, so new generated tasks start blank and can be assigned through the tracker. return None def generate_tasks_for_subscription(db: Session, *, subscription: ClientServiceSubscription, user_id: int) -> int: if quality_required_for_engagement(subscription.engagement_type) and getattr(subscription, "quality_acceptance_status", None) != QUALITY_APPROVED: raise ValueError(getattr(subscription, "quality_block_reason", None) or "AQMM acceptance is pending for this assurance engagement. Complete AQMM before generating tasks.") templates = db.execute( select(FirmServiceTaskTemplate) .where( FirmServiceTaskTemplate.tenant_id == subscription.tenant_id, FirmServiceTaskTemplate.service_catalogue_id == subscription.service_catalogue_id, FirmServiceTaskTemplate.is_active.is_(True), ) .order_by(FirmServiceTaskTemplate.sequence_no.asc(), FirmServiceTaskTemplate.id.asc()) ).scalars().all() created = 0 for template in templates: existing = db.execute( select(ClientServiceTaskInstance.id).where( ClientServiceTaskInstance.subscription_id == subscription.id, ClientServiceTaskInstance.firm_task_template_id == template.id, ClientServiceTaskInstance.financial_year == subscription.financial_year, ) ).first() if existing: continue db.add( ClientServiceTaskInstance( tenant_id=subscription.tenant_id, branch_id=subscription.branch_id, subscription_id=subscription.id, client_id=subscription.client_id, service_catalogue_id=subscription.service_catalogue_id, firm_task_template_id=template.id, financial_year=subscription.financial_year, assessment_year=subscription.assessment_year, task_name=template.task_name, description=template.description, sequence_no=template.sequence_no, default_role_name=template.default_role_name, assigned_to_user_id=_default_assignee_for_template(subscription, template), internal_target_date=_default_internal_target_date(subscription, template), status="pending", priority="normal", is_aqmm_task=getattr(template, "is_aqmm_task", False), aqmm_mandatory=getattr(template, "aqmm_mandatory", False), aqmm_evidence_required=getattr(template, "aqmm_evidence_required", False), aqmm_manager_review_required=getattr(template, "aqmm_manager_review_required", False), aqmm_partner_review_required=getattr(template, "aqmm_partner_review_required", False), aqmm_review_partner_required=getattr(template, "aqmm_review_partner_required", False), aqmm_blocks_final_release=getattr(template, "aqmm_blocks_final_release", False), aqmm_reference=getattr(template, "aqmm_reference", None), aqmm_status="pending" if getattr(template, "is_aqmm_task", False) else "not_required", aqmm_review_status="pending_review" if (getattr(template, "aqmm_manager_review_required", False) or getattr(template, "aqmm_partner_review_required", False) or getattr(template, "aqmm_review_partner_required", False)) else "not_required", is_active=True, created_by_user_id=user_id, updated_by_user_id=user_id, ) ) created += 1 return created def _latest_task_documents(db: Session, task_id: int) -> list[EngagementDocument]: return db.execute( select(EngagementDocument).where( EngagementDocument.task_instance_id == int(task_id), EngagementDocument.is_deleted.is_(False), ) ).scalars().all() def _task_documents_for_aqmm(db: Session, task: ClientServiceTaskInstance) -> list[EngagementDocument]: docs = _latest_task_documents(db, task.id) if getattr(task, "is_aqmm_task", False): return [doc for doc in docs if getattr(doc, "evidence_status", "not_required") != "rejected"] return docs def _task_has_evidence(db: Session, task_id: int) -> bool: return bool( db.execute( select(EngagementDocument.id).where( EngagementDocument.task_instance_id == int(task_id), EngagementDocument.is_deleted.is_(False), EngagementDocument.evidence_status != "rejected", ) ).first() ) def _evidence_status_for_task(db: Session, task: ClientServiceTaskInstance) -> str: if not getattr(task, "is_aqmm_task", False) or not getattr(task, "aqmm_evidence_required", False): return "not_required" docs = _latest_task_documents(db, task.id) if not docs: return "pending" if any(getattr(doc, "evidence_status", "uploaded") == "accepted" for doc in docs): return "accepted" if any(getattr(doc, "evidence_status", "uploaded") in {"uploaded", "pending_review"} for doc in docs): return "uploaded" return "pending" def _required_review_issue(task: ClientServiceTaskInstance) -> str | None: if getattr(task, "aqmm_manager_review_required", False) and task.manager_review_status != "reviewed": return "Manager review pending" if task.manager_review_status != "rework_required" else "Manager requested rework" if getattr(task, "aqmm_partner_review_required", False) and task.partner_review_status != "reviewed": return "Partner review pending" if task.partner_review_status != "rework_required" else "Partner requested rework" if getattr(task, "aqmm_review_partner_required", False) and task.review_partner_review_status != "reviewed": return "Review partner review pending" if task.review_partner_review_status != "rework_required" else "Review partner requested rework" return None def recalculate_task_aqmm_status(db: Session, task: ClientServiceTaskInstance) -> None: if not getattr(task, "is_aqmm_task", False): task.evidence_status = "not_required" task.aqmm_status = "not_required" task.aqmm_review_status = "not_required" task.aqmm_completed_at_utc = None return task.evidence_status = _evidence_status_for_task(db, task) review_issue = _required_review_issue(task) if review_issue: task.aqmm_review_status = "rework_required" if "rework" in review_issue.lower() else "pending_review" else: task.aqmm_review_status = "reviewed" if ( task.aqmm_manager_review_required or task.aqmm_partner_review_required or task.aqmm_review_partner_required ) else "not_required" issue = _aqmm_task_issue(db, task) if issue: if "rework" in issue.lower(): task.aqmm_status = "rework_required" elif "evidence" in issue.lower(): task.aqmm_status = "pending_evidence" elif "review" in issue.lower(): task.aqmm_status = "pending_review" else: task.aqmm_status = "pending" task.aqmm_completed_at_utc = None else: task.aqmm_status = "completed" if not task.aqmm_completed_at_utc: task.aqmm_completed_at_utc = datetime.now(timezone.utc) def _aqmm_task_issue(db: Session, task: ClientServiceTaskInstance) -> str | None: """Return blocking reason for one AQMM-tagged task, or None if it passes. The task itself is the AQMM evidence container. For AQMM tasks, task documents, staff notes and manager/partner/review-partner reviews attached to this task are the quality evidence trail. """ if not getattr(task, "is_aqmm_task", False): return None if getattr(task, "rework_status", "none") == "open": return "Rework is open" if getattr(task, "aqmm_mandatory", False) and task.status != "completed": return "Task not completed" if getattr(task, "aqmm_evidence_required", False) and not _task_has_evidence(db, task.id): return "Evidence not uploaded" review_issue = _required_review_issue(task) if review_issue: return review_issue return None def list_aqmm_quality_tasks(db: Session, *, subscription_id: int) -> list[dict]: tasks = db.execute( select(ClientServiceTaskInstance) .where( ClientServiceTaskInstance.subscription_id == int(subscription_id), ClientServiceTaskInstance.is_active.is_(True), ClientServiceTaskInstance.is_aqmm_task.is_(True), ) .order_by(ClientServiceTaskInstance.sequence_no.asc(), ClientServiceTaskInstance.id.asc()) ).scalars().all() rows: list[dict] = [] for task in tasks: has_evidence = _task_has_evidence(db, task.id) issue = _aqmm_task_issue(db, task) recalculate_task_aqmm_status(db, task) rows.append({ "task": task, "has_evidence": has_evidence, "evidence_status": getattr(task, "evidence_status", "not_required"), "review_status": getattr(task, "aqmm_review_status", "not_required"), "issue": issue, "passes": issue is None, }) return rows def aqmm_task_summary_for_subscription(db: Session, *, subscription_id: int) -> dict: rows = list_aqmm_quality_tasks(db, subscription_id=subscription_id) total = len(rows) completed = sum(1 for r in rows if r["task"].status == "completed") mandatory = sum(1 for r in rows if getattr(r["task"], "aqmm_mandatory", False)) mandatory_done = sum(1 for r in rows if getattr(r["task"], "aqmm_mandatory", False) and r["passes"]) evidence_required = sum(1 for r in rows if getattr(r["task"], "aqmm_evidence_required", False)) evidence_missing = sum(1 for r in rows if getattr(r["task"], "aqmm_evidence_required", False) and not r["has_evidence"]) manager_review_pending = sum(1 for r in rows if getattr(r["task"], "aqmm_manager_review_required", False) and r["task"].manager_review_status != "reviewed") partner_review_pending = sum(1 for r in rows if getattr(r["task"], "aqmm_partner_review_required", False) and r["task"].partner_review_status != "reviewed") review_partner_review_pending = sum(1 for r in rows if getattr(r["task"], "aqmm_review_partner_required", False) and r["task"].review_partner_review_status != "reviewed") rework_open = sum(1 for r in rows if getattr(r["task"], "rework_status", "none") == "open") blockers = [r for r in rows if getattr(r["task"], "aqmm_blocks_final_release", False) and not r["passes"]] status = "not_required" if total: status = "completed" if mandatory_done == mandatory and evidence_missing == 0 and manager_review_pending == 0 and partner_review_pending == 0 and review_partner_review_pending == 0 and rework_open == 0 else "in_progress" return { "rows": rows, "total": total, "completed": completed, "mandatory": mandatory, "mandatory_done": mandatory_done, "evidence_required": evidence_required, "evidence_missing": evidence_missing, "manager_review_pending": manager_review_pending, "partner_review_pending": partner_review_pending, "review_partner_review_pending": review_partner_review_pending, "rework_open": rework_open, "blockers": blockers, "status": status, } def assert_aqmm_quality_tasks_complete(db: Session, *, subscription_id: int) -> None: summary = aqmm_task_summary_for_subscription(db, subscription_id=subscription_id) if summary["blockers"]: first = summary["blockers"][0] task = first["task"] raise ValueError(f"AQMM final release blocked: {task.task_name} - {first['issue']}") if summary["mandatory"] and summary["mandatory_done"] < summary["mandatory"]: raise ValueError("AQMM mandatory quality checklist tasks are pending.") if summary["evidence_missing"]: raise ValueError("AQMM evidence upload is pending for one or more quality checklist tasks.") if summary["manager_review_pending"]: raise ValueError("AQMM manager review is pending for one or more quality checklist tasks.") if summary["partner_review_pending"]: raise ValueError("AQMM partner review is pending for one or more quality checklist tasks.") if summary["review_partner_review_pending"]: raise ValueError("AQMM review partner review is pending for one or more quality checklist tasks.") if summary["rework_open"]: raise ValueError("AQMM rework is open for one or more quality checklist tasks.") # ----------------------------------------------------------------------------- # Phase 5 - Engagement closure checklist # ----------------------------------------------------------------------------- CLOSURE_NOT_STARTED = "not_started" CLOSURE_PENDING = "pending" CLOSURE_READY = "ready_for_partner_closure" CLOSURE_CLOSED = "closed" CLOSURE_REOPENED = "reopened" def get_or_create_engagement_closure( db: Session, *, subscription: ClientServiceSubscription, actor_user_id: int | None = None, ) -> EngagementClosureChecklist: row = db.execute( select(EngagementClosureChecklist).where( EngagementClosureChecklist.subscription_id == subscription.id ) ).scalar_one_or_none() if row: return row row = EngagementClosureChecklist( tenant_id=subscription.tenant_id, branch_id=subscription.branch_id, client_id=subscription.client_id, subscription_id=subscription.id, closure_status=CLOSURE_NOT_STARTED, created_by_user_id=actor_user_id, ) db.add(row) return row def _normal_tasks_completed(db: Session, *, subscription_id: int) -> bool: open_task = db.execute( select(ClientServiceTaskInstance.id).where( ClientServiceTaskInstance.subscription_id == subscription_id, ClientServiceTaskInstance.is_active.is_(True), ClientServiceTaskInstance.status.notin_(list(CLOSED_TASK_STATUSES)), ) ).first() return open_task is None def _final_documents_released(db: Session, *, subscription_id: int, require_final_document: bool) -> bool: released = db.execute( select(EngagementDocument.id).where( EngagementDocument.engagement_id == subscription_id, EngagementDocument.is_deleted.is_(False), EngagementDocument.final_release_status == "released", ) ).first() if released: return True if not require_final_document: # Non-assurance closures can proceed even where no final document release # workflow was used for that engagement. Assurance closures need at least # one released final document. return True return False def _udin_completed(db: Session, *, subscription_id: int) -> bool: pending = db.execute( select(EngagementDocument.id).where( EngagementDocument.engagement_id == subscription_id, EngagementDocument.is_deleted.is_(False), EngagementDocument.udin_required.is_(True), or_( EngagementDocument.udin_number.is_(None), EngagementDocument.udin_number == "", EngagementDocument.udin_status.in_(["pending", "required", "not_generated"]), ), ) ).first() return pending is None def closure_readiness_for_subscription(db: Session, *, subscription: ClientServiceSubscription) -> dict: assurance = quality_required_for_engagement(subscription.engagement_type) aqmm_summary = aqmm_task_summary_for_subscription(db, subscription_id=subscription.id) normal_tasks_done = _normal_tasks_completed(db, subscription_id=subscription.id) aqmm_acceptance_completed = (not assurance) or getattr(subscription, "quality_acceptance_status", None) == QUALITY_APPROVED aqmm_tasks_completed = (not assurance) or ( aqmm_summary["mandatory_done"] == aqmm_summary["mandatory"] and aqmm_summary["evidence_missing"] == 0 and aqmm_summary["manager_review_pending"] == 0 and aqmm_summary["partner_review_pending"] == 0 and aqmm_summary["review_partner_review_pending"] == 0 and aqmm_summary["rework_open"] == 0 ) evidence_review_completed = (not assurance) or ( aqmm_summary["evidence_missing"] == 0 and aqmm_summary["manager_review_pending"] == 0 and aqmm_summary["partner_review_pending"] == 0 and aqmm_summary["review_partner_review_pending"] == 0 and aqmm_summary["rework_open"] == 0 ) final_documents_released = _final_documents_released( db, subscription_id=subscription.id, require_final_document=assurance ) udin_completed = _udin_completed(db, subscription_id=subscription.id) return { "assurance": assurance, "aqmm_task_summary": aqmm_summary, "normal_tasks_completed": normal_tasks_done, "aqmm_acceptance_completed": aqmm_acceptance_completed, "aqmm_tasks_completed": aqmm_tasks_completed, "evidence_review_completed": evidence_review_completed, "final_documents_released": final_documents_released, "udin_completed": udin_completed, } def update_engagement_closure_from_sources( db: Session, *, subscription: ClientServiceSubscription, actor_user_id: int | None = None, ) -> tuple[EngagementClosureChecklist, dict]: row = get_or_create_engagement_closure(db, subscription=subscription, actor_user_id=actor_user_id) summary = closure_readiness_for_subscription(db, subscription=subscription) row.aqmm_acceptance_completed = bool(summary["aqmm_acceptance_completed"]) row.aqmm_tasks_completed = bool(summary["aqmm_tasks_completed"]) row.evidence_review_completed = bool(summary["evidence_review_completed"]) row.final_documents_released = bool(summary["final_documents_released"]) row.udin_completed = bool(summary["udin_completed"]) row.normal_tasks_completed = bool(summary["normal_tasks_completed"]) blockers: list[str] = [] if not row.normal_tasks_completed: blockers.append("open work tracker tasks pending") if summary["assurance"]: if not row.aqmm_acceptance_completed: blockers.append("AQMM acceptance not approved") if not row.aqmm_tasks_completed: blockers.append("mandatory AQMM tasks/evidence/reviews pending") if not row.evidence_review_completed: blockers.append("AQMM evidence or review notes pending") if not row.final_documents_released: blockers.append("final document not released") if not row.udin_completed: blockers.append("UDIN pending for UDIN-required document") if not row.deliverables_sent_to_client: blockers.append("deliverables sent to client not confirmed") if not row.billing_reviewed: blockers.append("billing/fee status not reviewed") if not row.open_points_closed: blockers.append("open points not closed") if not row.client_communication_completed: blockers.append("client communication not completed") row.closure_block_reason = "; ".join(blockers) or None if row.closure_status == CLOSURE_CLOSED: pass elif blockers: row.closure_status = CLOSURE_PENDING if row.closure_status != CLOSURE_REOPENED else CLOSURE_REOPENED else: row.closure_status = CLOSURE_READY return row, {**summary, "blockers": blockers} def save_engagement_closure_confirmations( db: Session, *, subscription: ClientServiceSubscription, deliverables_sent_to_client: bool, billing_reviewed: bool, open_points_closed: bool, client_communication_completed: bool, closure_note: str | None, actor_user_id: int, ) -> tuple[EngagementClosureChecklist, dict]: row = get_or_create_engagement_closure(db, subscription=subscription, actor_user_id=actor_user_id) if row.closure_status == CLOSURE_CLOSED: raise ValueError("Closed engagement cannot be edited. Reopen it first.") row.deliverables_sent_to_client = bool(deliverables_sent_to_client) row.billing_reviewed = bool(billing_reviewed) row.open_points_closed = bool(open_points_closed) row.client_communication_completed = bool(client_communication_completed) row.closure_note = (closure_note or "").strip() or None return update_engagement_closure_from_sources(db, subscription=subscription, actor_user_id=actor_user_id) def approve_engagement_closure( db: Session, *, subscription: ClientServiceSubscription, actor_user_id: int, ) -> EngagementClosureChecklist: row, summary = update_engagement_closure_from_sources(db, subscription=subscription, actor_user_id=actor_user_id) if summary["blockers"]: raise ValueError(row.closure_block_reason or "Engagement closure checklist is not complete.") now = datetime.now(timezone.utc) row.closure_status = CLOSURE_CLOSED row.closure_block_reason = None row.closure_approved_by_user_id = actor_user_id row.closure_approved_at_utc = now subscription.status = "completed" subscription.is_locked = True subscription.locked_at_utc = now subscription.locked_by_user_id = actor_user_id subscription.updated_by_user_id = actor_user_id return row def reopen_engagement_closure( db: Session, *, subscription: ClientServiceSubscription, reason: str, actor_user_id: int, ) -> EngagementClosureChecklist: clean_reason = (reason or "").strip() if not clean_reason: raise ValueError("Reopen reason is required.") row = get_or_create_engagement_closure(db, subscription=subscription, actor_user_id=actor_user_id) now = datetime.now(timezone.utc) row.closure_status = CLOSURE_REOPENED row.reopened_by_user_id = actor_user_id row.reopened_at_utc = now row.reopen_reason = clean_reason row.closure_approved_by_user_id = None row.closure_approved_at_utc = None subscription.is_locked = False if subscription.status == "completed": subscription.status = "active" subscription.locked_at_utc = None subscription.locked_by_user_id = None subscription.updated_by_user_id = actor_user_id return row def _decorate_task_for_tracker(task: ClientServiceTaskInstance, *, today: date) -> ClientServiceTaskInstance: target_date = getattr(task, "internal_target_date", None) subscription = getattr(task, "subscription", None) engagement_due_date = getattr(subscription, "current_due_date", None) if subscription else None task.is_task_overdue = bool(target_date and target_date < today and task.status not in CLOSED_TASK_STATUSES) task.is_due_today = bool(target_date and target_date == today and task.status not in CLOSED_TASK_STATUSES) task.is_engagement_due_overdue = bool( engagement_due_date and engagement_due_date < today and task.status not in CLOSED_TASK_STATUSES ) task.tracker_status_label = dict(TASK_STATUSES).get(task.status, task.status) task.priority_label = dict(TASK_PRIORITIES).get(task.priority, task.priority) return task def _apply_partner_visibility_filter(query, partner_user_id: int | None): if not partner_user_id: return query return query.where( or_( ClientServiceTaskInstance.subscription.has( ClientServiceSubscription.assigned_partner_user_id == partner_user_id ), ClientServiceTaskInstance.client.has(Client.partner_id == partner_user_id), ) ) def list_tasks_payload( db: Session, *, tenant_id: int, branch_id: int | None = None, assigned_to_user_id: int | None = None, partner_user_id: int | None = None, status: str = "", q: str = "", include_inactive: bool = False, financial_year: str | None = None, ): today = date.today() special_filter = (status or "").strip().lower() query = ( select(ClientServiceTaskInstance) .options( selectinload(ClientServiceTaskInstance.client), selectinload(ClientServiceTaskInstance.catalogue), selectinload(ClientServiceTaskInstance.assigned_to), selectinload(ClientServiceTaskInstance.subscription), ) .where(ClientServiceTaskInstance.tenant_id == tenant_id) ) if financial_year: query = query.where(ClientServiceTaskInstance.financial_year == financial_year.strip()) if branch_id: query = query.where(ClientServiceTaskInstance.branch_id == branch_id) if assigned_to_user_id: query = query.where(ClientServiceTaskInstance.assigned_to_user_id == assigned_to_user_id) query = _apply_partner_visibility_filter(query, partner_user_id) if special_filter and special_filter not in {"overdue", "due_today", "unassigned"}: query = query.where(ClientServiceTaskInstance.status == special_filter) if special_filter == "overdue": query = query.where( ClientServiceTaskInstance.internal_target_date.is_not(None), ClientServiceTaskInstance.internal_target_date < today, ClientServiceTaskInstance.status.notin_(list(CLOSED_TASK_STATUSES)), ) elif special_filter == "due_today": query = query.where( ClientServiceTaskInstance.internal_target_date == today, ClientServiceTaskInstance.status.notin_(list(CLOSED_TASK_STATUSES)), ) elif special_filter == "unassigned": query = query.where(ClientServiceTaskInstance.assigned_to_user_id.is_(None)) if not include_inactive: query = query.where(ClientServiceTaskInstance.is_active.is_(True)) if q.strip(): term = f"%{q.strip()}%" query = ( query.join(Client, Client.id == ClientServiceTaskInstance.client_id) .join(ServiceCatalogue, ServiceCatalogue.id == ClientServiceTaskInstance.service_catalogue_id) .where( or_( ClientServiceTaskInstance.task_name.ilike(term), Client.client_name.ilike(term), Client.client_code.ilike(term), ServiceCatalogue.service_code.ilike(term), ServiceCatalogue.service_name.ilike(term), ) ) ) rows = db.execute( query.order_by( ClientServiceTaskInstance.internal_target_date.is_(None), ClientServiceTaskInstance.internal_target_date.asc(), ClientServiceTaskInstance.status.asc(), ClientServiceTaskInstance.sequence_no.asc(), ClientServiceTaskInstance.id.desc(), ) ).scalars().all() return [_decorate_task_for_tracker(task, today=today) for task in rows] def get_task( db: Session, *, tenant_id: int, task_id: int, branch_id: int | None = None, assigned_to_user_id: int | None = None, partner_user_id: int | None = None, financial_year: str | None = None, ) -> ClientServiceTaskInstance | None: query = ( select(ClientServiceTaskInstance) .options( selectinload(ClientServiceTaskInstance.client), selectinload(ClientServiceTaskInstance.catalogue), selectinload(ClientServiceTaskInstance.assigned_to), selectinload(ClientServiceTaskInstance.subscription), ) .where( ClientServiceTaskInstance.id == task_id, ClientServiceTaskInstance.tenant_id == tenant_id, ) ) if branch_id: query = query.where(ClientServiceTaskInstance.branch_id == branch_id) if financial_year: query = query.where(ClientServiceTaskInstance.financial_year == financial_year.strip()) if assigned_to_user_id: query = query.where(ClientServiceTaskInstance.assigned_to_user_id == assigned_to_user_id) query = _apply_partner_visibility_filter(query, partner_user_id) task = db.execute(query).scalar_one_or_none() return _decorate_task_for_tracker(task, today=date.today()) if task else None def list_assignees_for_execution(db: Session, *, tenant_id: int, branch_id: int | None = None): query = select(User).where(User.tenant_id == tenant_id, User.is_active.is_(True)) if branch_id: query = query.where((User.branch_id == branch_id) | (User.branch_id.is_(None))) return db.execute(query.order_by(User.full_name.asc(), User.email.asc())).scalars().all() def apply_task_update( task: ClientServiceTaskInstance, *, status: str, priority: str, assigned_to_user_id: int | None, internal_target_date: date | None, remarks: str, is_active: bool, user_id: int, ): if getattr(task, "is_locked", False) or getattr(getattr(task, "subscription", None), "is_locked", False): return previous_status = task.status task.status = _normalise_status(status) task.priority = _normalise_priority(priority) task.assigned_to_user_id = assigned_to_user_id task.internal_target_date = internal_target_date task.remarks = remarks.strip() or None task.is_active = is_active task.updated_by_user_id = user_id now = datetime.now(timezone.utc) if previous_status != "in_progress" and task.status == "in_progress" and not task.started_at_utc: task.started_at_utc = now if task.status == "completed" and not task.completed_at_utc: task.completed_at_utc = now if task.status != "completed": task.completed_at_utc = None def apply_bulk_task_update( tasks: list[ClientServiceTaskInstance], *, status: str | None, assigned_to_user_id: int | None, update_assignee: bool, internal_target_date: date | None, update_internal_target_date: bool, user_id: int, ) -> tuple[int, int]: updated = 0 skipped = 0 for task in tasks: if getattr(task, "is_locked", False) or getattr(getattr(task, "subscription", None), "is_locked", False): skipped += 1 continue previous_status = task.status if status: task.status = _normalise_status(status) if update_assignee: task.assigned_to_user_id = assigned_to_user_id if update_internal_target_date: task.internal_target_date = internal_target_date task.updated_by_user_id = user_id now = datetime.now(timezone.utc) if previous_status != "in_progress" and task.status == "in_progress" and not task.started_at_utc: task.started_at_utc = now if task.status == "completed" and not task.completed_at_utc: task.completed_at_utc = now if task.status != "completed": task.completed_at_utc = None updated += 1 return updated, skipped def get_tasks_for_bulk_update( db: Session, *, tenant_id: int, task_ids: list[int], branch_id: int | None = None, assigned_to_user_id: int | None = None, partner_user_id: int | None = None, financial_year: str | None = None, ) -> list[ClientServiceTaskInstance]: if not task_ids: return [] query = ( select(ClientServiceTaskInstance) .options(selectinload(ClientServiceTaskInstance.subscription)) .where( ClientServiceTaskInstance.tenant_id == tenant_id, ClientServiceTaskInstance.id.in_(task_ids), ClientServiceTaskInstance.is_active.is_(True), ) ) if branch_id: query = query.where(ClientServiceTaskInstance.branch_id == branch_id) if financial_year: query = query.where(ClientServiceTaskInstance.financial_year == financial_year.strip()) if assigned_to_user_id: query = query.where(ClientServiceTaskInstance.assigned_to_user_id == assigned_to_user_id) query = _apply_partner_visibility_filter(query, partner_user_id) return db.execute(query).scalars().all() def _append_system_task_comment( db: Session, *, task: ClientServiceTaskInstance, comment_type: str, message: str, user_id: int, ) -> None: if not (message or "").strip(): return 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=_normalise_comment_type(comment_type), visibility="internal", message=message.strip(), created_by_user_id=user_id, ) ) def submit_task_for_review(db: Session, *, task: ClientServiceTaskInstance, note: str, user_id: int) -> None: if getattr(task, "is_locked", False) or getattr(getattr(task, "subscription", None), "is_locked", False): return if getattr(task, "is_aqmm_task", False) and getattr(task, "aqmm_evidence_required", False) and not _task_has_evidence(db, task.id): raise ValueError("Evidence is required before submitting this AQMM task for review.") now = datetime.now(timezone.utc) task.submitted_for_review_by_user_id = user_id task.submitted_for_review_at_utc = now if getattr(task, "aqmm_manager_review_required", False) and task.manager_review_status != "reviewed": task.manager_review_status = "pending" if getattr(task, "aqmm_partner_review_required", False) and task.partner_review_status != "reviewed": task.partner_review_status = "pending" if getattr(task, "aqmm_review_partner_required", False) and task.review_partner_review_status != "reviewed": task.review_partner_review_status = "pending" if task.status == "pending": task.status = "in_progress" if getattr(task, "rework_status", "none") == "open": task.rework_status = "resolved" task.rework_resolved_at_utc = now task.updated_by_user_id = user_id _append_system_task_comment(db, task=task, comment_type="staff_work_note", message=note or "Submitted for review.", user_id=user_id) recalculate_task_aqmm_status(db, task) def apply_task_review( db: Session, *, task: ClientServiceTaskInstance, review_level: str, decision: str, note: str, user_id: int, ) -> None: if getattr(task, "is_locked", False) or getattr(getattr(task, "subscription", None), "is_locked", False): return clean_decision = (decision or "reviewed").strip().lower() if clean_decision not in {"reviewed", "rework_required"}: clean_decision = "reviewed" now = datetime.now(timezone.utc) clean_note = (note or "").strip() if clean_decision == "rework_required" and not clean_note: raise ValueError("Rework reason is required.") if review_level == "manager": task.manager_review_status = clean_decision task.manager_review_note = clean_note or task.manager_review_note task.manager_reviewed_by_user_id = user_id task.manager_reviewed_at_utc = now comment_type = "manager_review_note" elif review_level == "partner": task.partner_review_status = clean_decision task.partner_review_note = clean_note or task.partner_review_note task.partner_reviewed_by_user_id = user_id task.partner_reviewed_at_utc = now comment_type = "partner_review_note" elif review_level == "review_partner": task.review_partner_review_status = clean_decision task.review_partner_review_note = clean_note or task.review_partner_review_note task.review_partner_reviewed_by_user_id = user_id task.review_partner_reviewed_at_utc = now comment_type = "review_partner_review_note" else: raise ValueError("Invalid review level.") if clean_decision == "rework_required": task.rework_status = "open" task.rework_reason = clean_note task.rework_requested_by_user_id = user_id task.rework_requested_at_utc = now task.rework_resolved_at_utc = None task.status = "blocked" comment_type = "rework_note" else: if getattr(task, "rework_status", "none") == "open": task.rework_status = "resolved" task.rework_resolved_at_utc = now if getattr(task, "is_aqmm_task", False) and _required_review_issue(task) is None and (not getattr(task, "aqmm_evidence_required", False) or _task_has_evidence(db, task.id)): task.status = "completed" if not task.completed_at_utc: task.completed_at_utc = now task.updated_by_user_id = user_id _append_system_task_comment(db, task=task, comment_type=comment_type, message=clean_note or clean_decision.replace("_", " ").title(), user_id=user_id) recalculate_task_aqmm_status(db, task) def list_task_comments(db: Session, *, tenant_id: int, task_id: int) -> list[ServiceTaskComment]: return db.execute( select(ServiceTaskComment) .options(selectinload(ServiceTaskComment.created_by)) .where( ServiceTaskComment.tenant_id == tenant_id, ServiceTaskComment.task_instance_id == task_id, ServiceTaskComment.is_deleted.is_(False), ) .order_by(ServiceTaskComment.created_at_utc.desc(), ServiceTaskComment.id.desc()) ).scalars().all() def add_task_comment( db: Session, *, task: ClientServiceTaskInstance, comment_type: str, visibility: str, message: str, user_id: int, ) -> ServiceTaskComment | None: if getattr(task, "is_locked", False) or getattr(getattr(task, "subscription", None), "is_locked", False): return None clean_message = (message or "").strip() if not clean_message: return None row = ServiceTaskComment( tenant_id=task.tenant_id, branch_id=task.branch_id, subscription_id=task.subscription_id, task_instance_id=task.id, comment_type=_normalise_comment_type(comment_type), visibility=_normalise_visibility(visibility), message=clean_message, created_by_user_id=user_id, ) db.add(row) task.updated_by_user_id = user_id return row def list_client_visible_task_comments( db: Session, *, tenant_id: int, client_id: int, limit: int = 20, ) -> list[ServiceTaskComment]: """Return client-visible task communication for one client dashboard. This is intentionally read-only and scoped by tenant + client. Internal and consultant-only notes are never returned to the client portal. """ safe_limit = max(1, min(int(limit or 20), 100)) rows = db.execute( select(ServiceTaskComment) .options( selectinload(ServiceTaskComment.created_by), selectinload(ServiceTaskComment.task).selectinload(ClientServiceTaskInstance.catalogue), selectinload(ServiceTaskComment.task).selectinload(ClientServiceTaskInstance.subscription), ) .join(ClientServiceTaskInstance, ClientServiceTaskInstance.id == ServiceTaskComment.task_instance_id) .where( ServiceTaskComment.tenant_id == tenant_id, ServiceTaskComment.visibility == "client", ServiceTaskComment.is_deleted.is_(False), ClientServiceTaskInstance.client_id == client_id, ClientServiceTaskInstance.tenant_id == tenant_id, ClientServiceTaskInstance.is_active.is_(True), ) .order_by(ServiceTaskComment.created_at_utc.desc(), ServiceTaskComment.id.desc()) .limit(safe_limit) ).scalars().all() type_labels = dict(TASK_COMMENT_TYPES) visibility_labels = dict(TASK_COMMENT_VISIBILITIES) for row in rows: row.comment_type_label = type_labels.get(row.comment_type, row.comment_type) row.visibility_label = visibility_labels.get(row.visibility, row.visibility) return rows def dashboard_stats( db: Session, *, tenant_id: int, branch_id: int | None = None, assigned_to_user_id: int | None = None, partner_user_id: int | None = None, financial_year: str | None = None, ): today = date.today() base = select(ClientServiceTaskInstance).where( ClientServiceTaskInstance.tenant_id == tenant_id, ClientServiceTaskInstance.is_active.is_(True), ) if branch_id: base = base.where(ClientServiceTaskInstance.branch_id == branch_id) if financial_year: base = base.where(ClientServiceTaskInstance.financial_year == financial_year.strip()) if assigned_to_user_id: base = base.where(ClientServiceTaskInstance.assigned_to_user_id == assigned_to_user_id) base = _apply_partner_visibility_filter(base, partner_user_id) subq = base.subquery() total = db.execute(select(func.count()).select_from(subq)).scalar_one() pending = db.execute(select(func.count()).select_from(subq).where(subq.c.status == "pending")).scalar_one() progress = db.execute(select(func.count()).select_from(subq).where(subq.c.status == "in_progress")).scalar_one() blocked = db.execute(select(func.count()).select_from(subq).where(subq.c.status == "blocked")).scalar_one() completed = db.execute(select(func.count()).select_from(subq).where(subq.c.status == "completed")).scalar_one() not_applicable = db.execute(select(func.count()).select_from(subq).where(subq.c.status == "not_applicable")).scalar_one() unassigned = db.execute(select(func.count()).select_from(subq).where(subq.c.assigned_to_user_id.is_(None))).scalar_one() overdue = db.execute( select(func.count()).select_from(subq).where( subq.c.internal_target_date.is_not(None), subq.c.internal_target_date < today, subq.c.status.notin_(list(CLOSED_TASK_STATUSES)), ) ).scalar_one() due_today = db.execute( select(func.count()).select_from(subq).where( subq.c.internal_target_date == today, subq.c.status.notin_(list(CLOSED_TASK_STATUSES)), ) ).scalar_one() return { "total": total, "pending": pending, "in_progress": progress, "blocked": blocked, "completed": completed, "not_applicable": not_applicable, "unassigned": unassigned, "overdue": overdue, "due_today": due_today, }