from __future__ import annotations import csv import io import zipfile from dataclasses import dataclass from datetime import datetime, timezone from typing import Any from sqlalchemy import and_, func, or_, select from sqlalchemy.orm import joinedload from app.modules.clients.models import Client from app.modules.core.rbac.deps import get_user_roles from app.modules.documents.models import EngagementDocument, EngagementDocumentVersion from app.modules.services.models import ( ClientServiceSubscription, ClientServiceTaskInstance, EngagementClosureChecklist, EngagementKycVerification, EngagementLetter, EngagementQualityDeclaration, ServiceCatalogue, ) ASSURANCE_TYPE = "assurance" APPROVED_STATUSES = {"approved", "completed", "complete"} REVIEWED_STATUSES = {"reviewed", "approved", "completed", "complete", "not_required"} EVIDENCE_OK_STATUSES = {"accepted", "approved", "reviewed", "uploaded", "not_required"} CLOSURE_CLOSED_STATUSES = {"closed", "completed"} @dataclass(frozen=True) class PeerReviewScope: tenant_id: int | None branch_id: int | None def _clean(value: Any) -> str: return str(value or "").strip() def can_view_peer_review_export(db, user) -> bool: """Peer review export is restricted to Firm Admin and Partner roles only.""" roles = set(get_user_roles(db, user.id)) return bool(roles.intersection({"Firm Admin", "Partner"})) def _active_tenant_id(request, user) -> int | None: value = ( request.session.get("active_tenant_id") or request.session.get("selected_tenant_id") or request.session.get("tenant_id") or getattr(user, "tenant_id", None) ) return int(value) if value not in (None, "", 0, "0") else None def _active_branch_id(request, user) -> int | None: value = request.session.get("active_branch_id") or getattr(user, "branch_id", None) return int(value) if value not in (None, "", 0, "0") else None def build_scope(request, user) -> PeerReviewScope: return PeerReviewScope(tenant_id=_active_tenant_id(request, user), branch_id=_active_branch_id(request, user)) def _subscription_filters(scope: PeerReviewScope, financial_year: str = "", q: str = ""): filters = [ClientServiceSubscription.engagement_type == ASSURANCE_TYPE] if scope.tenant_id: filters.append(ClientServiceSubscription.tenant_id == scope.tenant_id) if scope.branch_id: filters.append(ClientServiceSubscription.branch_id == scope.branch_id) if financial_year: filters.append(ClientServiceSubscription.financial_year == financial_year) if q: like = f"%{q.lower()}%" filters.append( or_( func.lower(Client.client_name).like(like), func.lower(Client.client_code).like(like), func.lower(ServiceCatalogue.service_name).like(like), func.lower(ServiceCatalogue.service_code).like(like), ) ) return filters def list_peer_review_engagements(db, request, user, *, financial_year: str = "", q: str = "") -> dict[str, Any]: financial_year = _clean(financial_year) q = _clean(q) scope = build_scope(request, user) rows = db.execute( select(ClientServiceSubscription, EngagementClosureChecklist) .join(Client, Client.id == ClientServiceSubscription.client_id) .join(ServiceCatalogue, ServiceCatalogue.id == ClientServiceSubscription.service_catalogue_id) .outerjoin(EngagementClosureChecklist, EngagementClosureChecklist.subscription_id == ClientServiceSubscription.id) .options( joinedload(ClientServiceSubscription.client), joinedload(ClientServiceSubscription.catalogue), joinedload(ClientServiceSubscription.assigned_partner), joinedload(ClientServiceSubscription.assigned_manager), joinedload(ClientServiceSubscription.assigned_staff), joinedload(ClientServiceSubscription.review_partner), ) .where(and_(*_subscription_filters(scope, financial_year, q))) .order_by(ClientServiceSubscription.financial_year.desc(), ClientServiceSubscription.updated_at_utc.desc(), ClientServiceSubscription.id.desc()) .limit(200) ).all() return {"financial_year": financial_year, "q": q, "scope": scope, "engagements": rows} def get_subscription_for_export(db, request, user, subscription_id: int) -> ClientServiceSubscription | None: scope = build_scope(request, user) filters = [ClientServiceSubscription.id == subscription_id, ClientServiceSubscription.engagement_type == ASSURANCE_TYPE] if scope.tenant_id: filters.append(ClientServiceSubscription.tenant_id == scope.tenant_id) if scope.branch_id: filters.append(ClientServiceSubscription.branch_id == scope.branch_id) return db.execute( select(ClientServiceSubscription) .options( joinedload(ClientServiceSubscription.client), joinedload(ClientServiceSubscription.catalogue), joinedload(ClientServiceSubscription.assigned_partner), joinedload(ClientServiceSubscription.assigned_manager), joinedload(ClientServiceSubscription.assigned_staff), joinedload(ClientServiceSubscription.review_partner), ) .where(and_(*filters)) ).scalar_one_or_none() def _task_blockers(task: ClientServiceTaskInstance) -> list[str]: blockers: list[str] = [] if task.aqmm_evidence_required and (task.evidence_status or "pending") not in EVIDENCE_OK_STATUSES: blockers.append("Evidence pending") if task.aqmm_manager_review_required and (task.manager_review_status or "pending") not in REVIEWED_STATUSES: blockers.append("Manager review pending") if task.aqmm_partner_review_required and (task.partner_review_status or "pending") not in REVIEWED_STATUSES: blockers.append("Partner review pending") if task.aqmm_review_partner_required and (task.review_partner_review_status or "pending") not in REVIEWED_STATUSES: blockers.append("Review partner review pending") if (task.rework_status or "").lower() in {"requested", "pending", "rework_required"}: blockers.append("Rework pending") return blockers def build_peer_review_evidence(db, request, user, subscription_id: int) -> dict[str, Any]: subscription = get_subscription_for_export(db, request, user, subscription_id) if not subscription: raise ValueError("Assurance engagement not found or not accessible") tasks = db.execute( select(ClientServiceTaskInstance) .where(ClientServiceTaskInstance.subscription_id == subscription.id) .order_by(ClientServiceTaskInstance.sequence_no.asc(), ClientServiceTaskInstance.id.asc()) ).scalars().all() documents = db.execute( select(EngagementDocument) .options(joinedload(EngagementDocument.versions)) .where(EngagementDocument.engagement_id == subscription.id, EngagementDocument.is_deleted.is_(False)) .order_by(EngagementDocument.is_aqmm_evidence.desc(), EngagementDocument.document_type.asc(), EngagementDocument.id.asc()) ).scalars().unique().all() declarations = db.execute( select(EngagementQualityDeclaration) .options(joinedload(EngagementQualityDeclaration.requested_user), joinedload(EngagementQualityDeclaration.created_by)) .where(EngagementQualityDeclaration.subscription_id == subscription.id) .order_by(EngagementQualityDeclaration.declaration_type.asc(), EngagementQualityDeclaration.requested_user_id.asc()) ).scalars().all() kyc = db.execute( select(EngagementKycVerification) .options(joinedload(EngagementKycVerification.verified_by), joinedload(EngagementKycVerification.created_by)) .where(EngagementKycVerification.subscription_id == subscription.id) .order_by(EngagementKycVerification.updated_at_utc.desc()) ).scalars().all() letters = db.execute( select(EngagementLetter) .options( joinedload(EngagementLetter.partner_approved_by), joinedload(EngagementLetter.client_accepted_by), joinedload(EngagementLetter.manual_verified_by), joinedload(EngagementLetter.created_by), ) .where(EngagementLetter.subscription_id == subscription.id) .order_by(EngagementLetter.version_no.desc(), EngagementLetter.updated_at_utc.desc()) ).scalars().all() closure = db.execute( select(EngagementClosureChecklist) .options(joinedload(EngagementClosureChecklist.closure_approved_by), joinedload(EngagementClosureChecklist.reopened_by)) .where(EngagementClosureChecklist.subscription_id == subscription.id) ).scalar_one_or_none() aqmm_tasks = [t for t in tasks if getattr(t, "is_aqmm_task", False)] blockers = [] if subscription.quality_workflow_required and (subscription.quality_workflow_status or "pending") not in APPROVED_STATUSES: blockers.append("AQMM acceptance pending") for task in aqmm_tasks: for b in _task_blockers(task): blockers.append(f"{task.task_name}: {b}") if closure and (closure.closure_status or "not_started") not in CLOSURE_CLOSED_STATUSES: blockers.append("Closure not completed") return { "generated_at": datetime.now(timezone.utc), "subscription": subscription, "tasks": tasks, "aqmm_tasks": aqmm_tasks, "documents": documents, "declarations": declarations, "kyc": kyc, "letters": letters, "closure": closure, "blockers": blockers, } def _writerow_dict(writer, row: dict[str, Any], headers: list[str]) -> None: writer.writerow([row.get(h, "") for h in headers]) def _dt(value) -> str: return value.isoformat() if value else "" def _user_label(user) -> str: if not user: return "" return getattr(user, "full_name", None) or getattr(user, "email", None) or str(getattr(user, "id", "")) def _make_csv(headers: list[str], rows: list[dict[str, Any]]) -> bytes: buf = io.StringIO() writer = csv.writer(buf) writer.writerow(headers) for row in rows: _writerow_dict(writer, row, headers) return buf.getvalue().encode("utf-8-sig") def build_peer_review_zip(db, request, user, subscription_id: int) -> bytes: payload = build_peer_review_evidence(db, request, user, subscription_id) sub = payload["subscription"] client_name = sub.client.client_name if sub.client else "" service_name = sub.catalogue.service_name if sub.catalogue else "" summary_headers = ["field", "value"] summary_rows = [ {"field": "Generated At UTC", "value": _dt(payload["generated_at"])}, {"field": "Client", "value": client_name}, {"field": "Service", "value": service_name}, {"field": "Financial Year", "value": sub.financial_year or ""}, {"field": "Engagement Type", "value": sub.engagement_type or ""}, {"field": "Engagement Status", "value": sub.status or ""}, {"field": "AQMM Required", "value": "Yes" if sub.quality_workflow_required else "No"}, {"field": "AQMM Status", "value": sub.quality_workflow_status or ""}, {"field": "Assigned Partner", "value": _user_label(sub.assigned_partner)}, {"field": "Assigned Manager", "value": _user_label(sub.assigned_manager)}, {"field": "Assigned Staff", "value": _user_label(sub.assigned_staff)}, {"field": "Review Partner", "value": _user_label(sub.review_partner)}, {"field": "Blockers", "value": "; ".join(payload["blockers"])}, ] task_headers = [ "sequence_no", "task_name", "status", "is_aqmm_task", "aqmm_mandatory", "evidence_required", "evidence_status", "manager_review_required", "manager_review_status", "manager_review_note", "partner_review_required", "partner_review_status", "partner_review_note", "review_partner_required", "review_partner_review_status", "review_partner_review_note", "rework_status", "rework_reason", "staff_work_note", "aqmm_reference", "blockers", ] task_rows = [] for t in payload["tasks"]: task_rows.append({ "sequence_no": t.sequence_no, "task_name": t.task_name, "status": t.status or "", "is_aqmm_task": "Yes" if t.is_aqmm_task else "No", "aqmm_mandatory": "Yes" if t.aqmm_mandatory else "No", "evidence_required": "Yes" if t.aqmm_evidence_required else "No", "evidence_status": t.evidence_status or "", "manager_review_required": "Yes" if t.aqmm_manager_review_required else "No", "manager_review_status": t.manager_review_status or "", "manager_review_note": t.manager_review_note or "", "partner_review_required": "Yes" if t.aqmm_partner_review_required else "No", "partner_review_status": t.partner_review_status or "", "partner_review_note": t.partner_review_note or "", "review_partner_required": "Yes" if t.aqmm_review_partner_required else "No", "review_partner_review_status": t.review_partner_review_status or "", "review_partner_review_note": t.review_partner_review_note or "", "rework_status": t.rework_status or "", "rework_reason": t.rework_reason or "", "staff_work_note": t.staff_work_note or "", "aqmm_reference": t.aqmm_reference or "", "blockers": "; ".join(_task_blockers(t)), }) document_headers = [ "document_code", "title", "document_type", "task_instance_id", "is_aqmm_evidence", "evidence_status", "evidence_type", "evidence_description", "final_release_status", "udin_required", "udin_status", "udin_number", "version_no", "original_filename", "file_hash_sha256", "local_relative_path", "uploaded_at_utc", ] document_rows = [] for d in payload["documents"]: versions = list(d.versions or []) if not versions: document_rows.append({ "document_code": d.document_code, "title": d.title, "document_type": d.document_type, "task_instance_id": d.task_instance_id or "", "is_aqmm_evidence": "Yes" if d.is_aqmm_evidence else "No", "evidence_status": d.evidence_status or "", "evidence_type": d.evidence_type or "", "evidence_description": d.evidence_description or "", "final_release_status": d.final_release_status or "", "udin_required": "Yes" if d.udin_required else "No", "udin_status": d.udin_status or "", "udin_number": d.udin_number or "", "version_no": "", "original_filename": "", "file_hash_sha256": "", "local_relative_path": "", "uploaded_at_utc": "", }) for v in versions: document_rows.append({ "document_code": d.document_code, "title": d.title, "document_type": d.document_type, "task_instance_id": d.task_instance_id or "", "is_aqmm_evidence": "Yes" if d.is_aqmm_evidence else "No", "evidence_status": d.evidence_status or "", "evidence_type": d.evidence_type or "", "evidence_description": d.evidence_description or "", "final_release_status": d.final_release_status or "", "udin_required": "Yes" if d.udin_required else "No", "udin_status": d.udin_status or "", "udin_number": d.udin_number or "", "version_no": v.version_no, "original_filename": v.original_filename, "file_hash_sha256": v.file_hash_sha256, "local_relative_path": v.local_relative_path, "uploaded_at_utc": _dt(v.uploaded_at_utc), }) declaration_headers = ["type", "requested_user", "requested_role", "status", "response_notes", "responded_at_utc", "ip_address"] declaration_rows = [ { "type": d.declaration_type, "requested_user": _user_label(d.requested_user), "requested_role": d.requested_role or "", "status": d.status or "", "response_notes": d.response_notes or "", "responded_at_utc": _dt(d.responded_at_utc), "ip_address": d.ip_address or "", } for d in payload["declarations"] ] kyc_headers = ["status", "source", "verification_notes", "verified_by", "verified_at_utc"] kyc_rows = [ {"status": k.status or "", "source": k.source or "", "verification_notes": k.verification_notes or "", "verified_by": _user_label(k.verified_by), "verified_at_utc": _dt(k.verified_at_utc)} for k in payload["kyc"] ] letter_headers = ["title", "version_no", "status", "acceptance_mode", "partner_approved_by", "partner_approved_at_utc", "client_accepted_by", "client_accepted_at_utc", "manual_verified_by", "manual_verified_at_utc", "pdf_sha256", "manual_upload_path"] letter_rows = [ { "title": l.title, "version_no": l.version_no, "status": l.status or "", "acceptance_mode": l.acceptance_mode or "", "partner_approved_by": _user_label(l.partner_approved_by), "partner_approved_at_utc": _dt(l.partner_approved_at_utc), "client_accepted_by": _user_label(l.client_accepted_by), "client_accepted_at_utc": _dt(l.client_accepted_at_utc), "manual_verified_by": _user_label(l.manual_verified_by), "manual_verified_at_utc": _dt(l.manual_verified_at_utc), "pdf_sha256": l.pdf_sha256 or "", "manual_upload_path": l.manual_upload_path or "", } for l in payload["letters"] ] closure_headers = ["closure_status", "aqmm_acceptance_completed", "aqmm_tasks_completed", "evidence_review_completed", "final_documents_released", "udin_completed", "normal_tasks_completed", "deliverables_sent_to_client", "billing_reviewed", "open_points_closed", "client_communication_completed", "closure_note", "closure_block_reason", "closure_approved_by", "closure_approved_at_utc"] c = payload["closure"] closure_rows = [] if c: closure_rows.append({ "closure_status": c.closure_status or "", "aqmm_acceptance_completed": "Yes" if c.aqmm_acceptance_completed else "No", "aqmm_tasks_completed": "Yes" if c.aqmm_tasks_completed else "No", "evidence_review_completed": "Yes" if c.evidence_review_completed else "No", "final_documents_released": "Yes" if c.final_documents_released else "No", "udin_completed": "Yes" if c.udin_completed else "No", "normal_tasks_completed": "Yes" if c.normal_tasks_completed else "No", "deliverables_sent_to_client": "Yes" if c.deliverables_sent_to_client else "No", "billing_reviewed": "Yes" if c.billing_reviewed else "No", "open_points_closed": "Yes" if c.open_points_closed else "No", "client_communication_completed": "Yes" if c.client_communication_completed else "No", "closure_note": c.closure_note or "", "closure_block_reason": c.closure_block_reason or "", "closure_approved_by": _user_label(c.closure_approved_by), "closure_approved_at_utc": _dt(c.closure_approved_at_utc), }) readme = f"""Peer Review Evidence Export\nGenerated At UTC: {_dt(payload['generated_at'])}\nClient: {client_name}\nService: {service_name}\nFinancial Year: {sub.financial_year or ''}\n\nFiles included:\n- 01_engagement_summary.csv\n- 02_aqmm_and_task_review_evidence.csv\n- 03_document_and_udin_evidence.csv\n- 04_independence_conflict_declarations.csv\n- 05_kyc_verification.csv\n- 06_engagement_letter_evidence.csv\n- 07_closure_checklist.csv\n\nNote: This export contains the peer review evidence index, statuses, notes, hashes and storage paths. Physical documents remain in the ERP document vault and can be opened from the engagement/task documents screens.\n""" bio = io.BytesIO() with zipfile.ZipFile(bio, "w", compression=zipfile.ZIP_DEFLATED) as zf: zf.writestr("README.txt", readme) zf.writestr("01_engagement_summary.csv", _make_csv(summary_headers, summary_rows)) zf.writestr("02_aqmm_and_task_review_evidence.csv", _make_csv(task_headers, task_rows)) zf.writestr("03_document_and_udin_evidence.csv", _make_csv(document_headers, document_rows)) zf.writestr("04_independence_conflict_declarations.csv", _make_csv(declaration_headers, declaration_rows)) zf.writestr("05_kyc_verification.csv", _make_csv(kyc_headers, kyc_rows)) zf.writestr("06_engagement_letter_evidence.csv", _make_csv(letter_headers, letter_rows)) zf.writestr("07_closure_checklist.csv", _make_csv(closure_headers, closure_rows)) return bio.getvalue()