From a362cf438cd2b5d157517d5d27cd4ad7d0dfccd8 Mon Sep 17 00:00:00 2001 From: A R R R Associates Date: Mon, 31 Aug 2026 15:24:09 +0530 Subject: [PATCH] Add task tool result bridge and evidence attachment workflow --- .../20260831_task_tool_result_bridge.py | 61 +++++++ app/modules/services/models.py | 36 +++- app/modules/services/task_tool_results.py | 145 ++++++++++++++++ app/modules/services/task_tools.py | 23 +++ .../services/work_tracker/task_form.html | 78 ++++++++- app/modules/services/work_tracker_ui.py | 160 +++++++++++++++++- 6 files changed, 497 insertions(+), 6 deletions(-) create mode 100644 alembic/versions/20260831_task_tool_result_bridge.py create mode 100644 app/modules/services/task_tool_results.py diff --git a/alembic/versions/20260831_task_tool_result_bridge.py b/alembic/versions/20260831_task_tool_result_bridge.py new file mode 100644 index 0000000..bb23732 --- /dev/null +++ b/alembic/versions/20260831_task_tool_result_bridge.py @@ -0,0 +1,61 @@ +"""task tool result bridge + +Revision ID: 20260831_task_tool_result_bridge +Revises: 20260830_task_admin_tools +""" +from alembic import op +import sqlalchemy as sa + +revision = "20260831_task_tool_result_bridge" +down_revision = "20260830_task_admin_tools" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "service_task_tool_runs", + sa.Column("id", sa.Integer(), autoincrement=True, nullable=False), + sa.Column("tenant_id", sa.Integer(), nullable=False), + sa.Column("task_instance_id", sa.Integer(), nullable=False), + sa.Column("client_id", sa.Integer(), nullable=False), + sa.Column("tool_code", sa.String(length=80), nullable=False), + sa.Column("status", sa.String(length=30), server_default="started", nullable=False), + sa.Column("summary", sa.Text(), nullable=True), + sa.Column("result_json", sa.Text(), nullable=True), + sa.Column("exception_count", sa.Integer(), nullable=True), + sa.Column("clean_pass", sa.Boolean(), nullable=True), + sa.Column("output_filename", sa.String(length=255), nullable=True), + sa.Column("output_content_type", sa.String(length=120), nullable=True), + sa.Column("output_bytes", sa.LargeBinary(), nullable=True), + sa.Column("started_by_user_id", sa.Integer(), nullable=True), + sa.Column("completed_by_user_id", sa.Integer(), nullable=True), + sa.Column("started_at_utc", sa.DateTime(timezone=True), nullable=False), + sa.Column("completed_at_utc", sa.DateTime(timezone=True), nullable=True), + sa.Column("created_at_utc", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at_utc", sa.DateTime(timezone=True), nullable=False), + sa.ForeignKeyConstraint(["tenant_id"], ["tenants.id"], ondelete="CASCADE"), + sa.ForeignKeyConstraint(["task_instance_id"], ["client_service_task_instances.id"], ondelete="CASCADE"), + sa.ForeignKeyConstraint(["client_id"], ["clients.id"], ondelete="CASCADE"), + sa.ForeignKeyConstraint(["started_by_user_id"], ["users.id"], ondelete="SET NULL"), + sa.ForeignKeyConstraint(["completed_by_user_id"], ["users.id"], ondelete="SET NULL"), + sa.PrimaryKeyConstraint("id"), + ) + op.create_index("ix_service_task_tool_runs_tenant_id", "service_task_tool_runs", ["tenant_id"], unique=False) + op.create_index("ix_service_task_tool_runs_task_instance_id", "service_task_tool_runs", ["task_instance_id"], unique=False) + op.create_index("ix_service_task_tool_runs_client_id", "service_task_tool_runs", ["client_id"], unique=False) + op.create_index("ix_service_task_tool_runs_tool_code", "service_task_tool_runs", ["tool_code"], unique=False) + op.create_index("ix_service_task_tool_runs_status", "service_task_tool_runs", ["status"], unique=False) + op.create_index("ix_service_task_tool_runs_started_by_user_id", "service_task_tool_runs", ["started_by_user_id"], unique=False) + op.create_index("ix_service_task_tool_runs_completed_by_user_id", "service_task_tool_runs", ["completed_by_user_id"], unique=False) + + +def downgrade(): + op.drop_index("ix_service_task_tool_runs_completed_by_user_id", table_name="service_task_tool_runs") + op.drop_index("ix_service_task_tool_runs_started_by_user_id", table_name="service_task_tool_runs") + op.drop_index("ix_service_task_tool_runs_status", table_name="service_task_tool_runs") + op.drop_index("ix_service_task_tool_runs_tool_code", table_name="service_task_tool_runs") + op.drop_index("ix_service_task_tool_runs_client_id", table_name="service_task_tool_runs") + op.drop_index("ix_service_task_tool_runs_task_instance_id", table_name="service_task_tool_runs") + op.drop_index("ix_service_task_tool_runs_tenant_id", table_name="service_task_tool_runs") + op.drop_table("service_task_tool_runs") diff --git a/app/modules/services/models.py b/app/modules/services/models.py index b4a4e84..88a6a90 100644 --- a/app/modules/services/models.py +++ b/app/modules/services/models.py @@ -2,7 +2,7 @@ from __future__ import annotations from datetime import date, datetime, timezone -from sqlalchemy import Boolean, Date, DateTime, ForeignKey, Integer, Numeric, String, Text, UniqueConstraint +from sqlalchemy import Boolean, Date, DateTime, ForeignKey, Integer, LargeBinary, Numeric, String, Text, UniqueConstraint from sqlalchemy.orm import Mapped, mapped_column, relationship from app.core.db.common import CommonBase @@ -940,6 +940,40 @@ class ServiceTaskDocumentRequest(CommonBase): verified_by = relationship("User", foreign_keys=[verified_by_user_id]) + + +class ServiceTaskToolRun(CommonBase): + """Persisted execution/result record for an ERP tool launched from a service task. + + The result attachment is stored with the run so a tool can return its generated + workbook/PDF/JSON to the originating task without changing the existing document + module or historical EngagementDocument records. + """ + + __tablename__ = "service_task_tool_runs" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) + tenant_id: Mapped[int] = mapped_column(ForeignKey("tenants.id", ondelete="CASCADE"), nullable=False, index=True) + task_instance_id: Mapped[int] = mapped_column(ForeignKey("client_service_task_instances.id", ondelete="CASCADE"), nullable=False, index=True) + client_id: Mapped[int] = mapped_column(ForeignKey("clients.id", ondelete="CASCADE"), nullable=False, index=True) + tool_code: Mapped[str] = mapped_column(String(80), nullable=False, index=True) + status: Mapped[str] = mapped_column(String(30), nullable=False, default="started", index=True) + summary: Mapped[str | None] = mapped_column(Text, nullable=True) + result_json: Mapped[str | None] = mapped_column(Text, nullable=True) + exception_count: Mapped[int | None] = mapped_column(Integer, nullable=True) + clean_pass: Mapped[bool | None] = mapped_column(Boolean, nullable=True) + output_filename: Mapped[str | None] = mapped_column(String(255), nullable=True) + output_content_type: Mapped[str | None] = mapped_column(String(120), nullable=True) + output_bytes: Mapped[bytes | None] = mapped_column(LargeBinary, nullable=True) + started_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True, index=True) + completed_by_user_id: Mapped[int | None] = mapped_column(ForeignKey("users.id", ondelete="SET NULL"), nullable=True, index=True) + started_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False) + completed_at_utc: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + created_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), nullable=False) + updated_at_utc: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=lambda: datetime.now(timezone.utc), onupdate=lambda: datetime.now(timezone.utc), nullable=False) + + task = relationship("ClientServiceTaskInstance", foreign_keys=[task_instance_id]) + class ServiceTaskComment(CommonBase): """Communication timeline entry linked to a service task instance.""" diff --git a/app/modules/services/task_tool_results.py b/app/modules/services/task_tool_results.py new file mode 100644 index 0000000..92389ba --- /dev/null +++ b/app/modules/services/task_tool_results.py @@ -0,0 +1,145 @@ +from __future__ import annotations + +import json +from datetime import datetime, timezone +from typing import Any + +from sqlalchemy import select +from sqlalchemy.orm import Session + +from app.modules.services.models import ClientServiceTaskInstance, ServiceTaskToolRun +from app.modules.services.task_tools import get_task_tool, normalize_task_tool_code + +MAX_TOOL_OUTPUT_BYTES = 25 * 1024 * 1024 +VALID_RESULT_STATUSES = {"completed", "completed_with_exceptions", "failed", "cancelled"} + + +def start_task_tool_run( + db: Session, + *, + task: ClientServiceTaskInstance, + user_id: int | None, +) -> ServiceTaskToolRun: + tool_code = normalize_task_tool_code(getattr(task, "task_tool_code", "NONE")) + if tool_code == "NONE": + raise ValueError("This task does not have an ERP tool linked.") + get_task_tool(tool_code) + run = ServiceTaskToolRun( + tenant_id=int(task.tenant_id), + task_instance_id=int(task.id), + client_id=int(task.client_id), + tool_code=tool_code, + status="started", + started_by_user_id=user_id, + ) + db.add(run) + db.flush() + db.commit() + db.refresh(run) + return run + + +def list_task_tool_runs(db: Session, *, tenant_id: int, task_id: int) -> list[ServiceTaskToolRun]: + return list( + db.scalars( + select(ServiceTaskToolRun) + .where( + ServiceTaskToolRun.tenant_id == int(tenant_id), + ServiceTaskToolRun.task_instance_id == int(task_id), + ) + .order_by(ServiceTaskToolRun.started_at_utc.desc(), ServiceTaskToolRun.id.desc()) + ).all() + ) + + +def get_task_tool_run(db: Session, *, tenant_id: int, run_id: int) -> ServiceTaskToolRun | None: + return db.scalar( + select(ServiceTaskToolRun).where( + ServiceTaskToolRun.id == int(run_id), + ServiceTaskToolRun.tenant_id == int(tenant_id), + ) + ) + + +def _json_text(payload: Any) -> str | None: + if payload is None or payload == "": + return None + if isinstance(payload, str): + text = payload.strip() + if not text: + return None + try: + parsed = json.loads(text) + return json.dumps(parsed, ensure_ascii=False, separators=(",", ":")) + except Exception: + return json.dumps({"value": text}, ensure_ascii=False, separators=(",", ":")) + return json.dumps(payload, ensure_ascii=False, default=str, separators=(",", ":")) + + +def record_task_tool_result( + db: Session, + *, + run: ServiceTaskToolRun, + status: str, + summary: str | None = None, + result_payload: Any = None, + clean_pass: bool | None = None, + exception_count: int | None = None, + output_filename: str | None = None, + output_content_type: str | None = None, + output_bytes: bytes | None = None, + completed_by_user_id: int | None = None, +) -> ServiceTaskToolRun: + normalized_status = (status or "completed").strip().lower() + if normalized_status not in VALID_RESULT_STATUSES: + raise ValueError(f"Unsupported tool result status: {normalized_status}") + if output_bytes is not None and len(output_bytes) > MAX_TOOL_OUTPUT_BYTES: + raise ValueError("Tool output exceeds the 25 MB task-result attachment limit.") + + run.status = normalized_status + run.summary = (summary or "").strip() or None + run.result_json = _json_text(result_payload) + run.clean_pass = clean_pass + run.exception_count = exception_count + run.output_filename = (output_filename or "").strip() or None + run.output_content_type = (output_content_type or "").strip() or None + run.output_bytes = output_bytes + run.completed_by_user_id = completed_by_user_id + run.completed_at_utc = datetime.now(timezone.utc) + db.add(run) + db.commit() + db.refresh(run) + return run + + +def record_task_tool_result_by_id( + db: Session, + *, + tenant_id: int, + run_id: int, + status: str, + summary: str | None = None, + result_payload: Any = None, + clean_pass: bool | None = None, + exception_count: int | None = None, + output_filename: str | None = None, + output_content_type: str | None = None, + output_bytes: bytes | None = None, + completed_by_user_id: int | None = None, +) -> ServiceTaskToolRun: + run = get_task_tool_run(db, tenant_id=tenant_id, run_id=run_id) + if run is None: + raise ValueError("Task tool run not found.") + return record_task_tool_result( + db, + run=run, + status=status, + summary=summary, + result_payload=result_payload, + clean_pass=clean_pass, + exception_count=exception_count, + output_filename=output_filename, + output_content_type=output_content_type, + output_bytes=output_bytes, + completed_by_user_id=completed_by_user_id, + ) diff --git a/app/modules/services/task_tools.py b/app/modules/services/task_tools.py index b649835..7973d8e 100644 --- a/app/modules/services/task_tools.py +++ b/app/modules/services/task_tools.py @@ -1,6 +1,7 @@ from __future__ import annotations from dataclasses import dataclass +from urllib.parse import urlencode @dataclass(frozen=True) @@ -46,3 +47,25 @@ def normalize_task_tool_code(code: str | None) -> str: if value not in valid: raise ValueError(f"Unsupported task tool: {value}") return value + + +def build_task_tool_launch_url( + tool: TaskToolDefinition, + *, + client_id: int, + task_id: int, + tool_run_id: int, +) -> str: + if not tool.launch_url: + raise ValueError("This ERP task tool does not expose a launch URL.") + query = urlencode( + { + "client_id": int(client_id), + "task_id": int(task_id), + "task_tool_run_id": int(tool_run_id), + "task_return_url": f"/services/work-tracker/tasks/{int(task_id)}/edit", + "task_result_url": f"/services/work-tracker/tool-runs/{int(tool_run_id)}/result", + } + ) + separator = "&" if "?" in tool.launch_url else "?" + return f"{tool.launch_url}{separator}{query}" diff --git a/app/modules/services/templates/services/work_tracker/task_form.html b/app/modules/services/templates/services/work_tracker/task_form.html index 6c89b3e..10df635 100644 --- a/app/modules/services/templates/services/work_tracker/task_form.html +++ b/app/modules/services/templates/services/work_tracker/task_form.html @@ -88,8 +88,16 @@ {% if task_tool and task_tool.code != 'NONE' %}
-
ERP Task Tool · {{ task_tool.name }}

{{ task_tool.description }}

{% if task_tool.launch_url %}Open Tool{% endif %}
-

Opening the tool does not automatically mark this task complete. Save the task response/evidence after completing the verification.

+
+
+
ERP Task Tool · {{ task_tool.name }}
+

{{ task_tool.description }}

+
+ {% if task_tool.launch_url %} + + {% endif %} +
+

A Tool Run is created before launch. The linked tool receives the task, client, run and return/result URLs so its generated output can be returned to this task without changing the normal checklist workflow.

{% endif %} {% if (task.response_type or 'NONE') != 'NONE' or task.response_required or task.evidence_required %} @@ -156,6 +164,72 @@ +{% if task_tool and task_tool.code != 'NONE' %} +
+
+
+

Tool Output & Evidence

+

Every run is retained against this task. Output returned by the tool is stored with the run and can be downloaded from this task later.

+
+ {{ tool_runs|length }} run{% if tool_runs|length != 1 %}s{% endif %} +
+ + {% if request.query_params.get('tool_saved') %}
Tool result attached to this task.
{% endif %} + {% if request.query_params.get('tool_error') %}
The tool result could not be saved. Check the selected run and output size.
{% endif %} + +
+ {% if tool_runs %} + {% for run in tool_runs %} +
+
+
+
{{ task_tool.name }} · Run #{{ run.id }}
+
Started {{ run.started_at_utc }}{% if run.completed_at_utc %} · Completed {{ run.completed_at_utc }}{% endif %}
+
+ {{ run.status.replace('_',' ').title() }} +
+ {% if run.summary %}

{{ run.summary }}

{% endif %} +
+ {% if run.clean_pass is not none %}Clean Pass: {{ 'Yes' if run.clean_pass else 'No' }}{% endif %} + {% if run.exception_count is not none %}Exceptions: {{ run.exception_count }}{% endif %} + {% if run.output_filename %}Download {{ run.output_filename }}{% endif %} +
+ + {% if run.status == 'started' %} +
+ +
+ + +
+
+ + +
+
+ + +
+
+ + +
+
+ + +
+
+
+ {% endif %} +
+ {% endfor %} + {% else %} +
No tool run has been started for this task yet.
+ {% endif %} +
+
+{% endif %} + {% if task.is_aqmm_task %}
diff --git a/app/modules/services/work_tracker_ui.py b/app/modules/services/work_tracker_ui.py index 43ceb62..7c8a905 100644 --- a/app/modules/services/work_tracker_ui.py +++ b/app/modules/services/work_tracker_ui.py @@ -1,7 +1,7 @@ from __future__ import annotations -from fastapi import APIRouter, Form, Request -from fastapi.responses import RedirectResponse +from fastapi import APIRouter, File, Form, Request, UploadFile +from fastapi.responses import RedirectResponse, Response from sqlalchemy import select from app.core.db.common import CommonSessionLocal @@ -11,7 +11,13 @@ from app.core.templating import templates from app.modules.core.rbac.deps import get_user_permissions, get_user_roles from app.modules.core.iam.models import User from app.modules.core.rbac.permission_guard import require_permission -from app.modules.services.task_tools import get_task_tool +from app.modules.services.task_tools import build_task_tool_launch_url, get_task_tool +from app.modules.services.task_tool_results import ( + get_task_tool_run, + list_task_tool_runs, + record_task_tool_result, + start_task_tool_run, +) from app.modules.services.execution import ( TASK_COMMENT_TYPES, TASK_COMMENT_VISIBILITIES, @@ -445,6 +451,154 @@ def task_edit_page(request: Request, task_id: int): can_manage_fields=_can_bulk_manage_tasks(db, user), can_reassign=_can_assign_staff(db, user), task_tool=get_task_tool(getattr(task, "task_tool_code", "NONE")), + tool_runs=list_task_tool_runs(db, tenant_id=tenant_id, task_id=task.id), + ) + finally: + db.close() + + +@router.post("/tasks/{task_id}/tool/start") +def task_tool_start(request: Request, task_id: int, csrf_token: str = Form(...)): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user = get_current_user(request, db=db) + if not user: + return RedirectResponse(url="/login", status_code=303) + tenant_id = _active_tenant_id(request, user) + branch_id = _active_branch_id(request, user, db) + task = get_task( + db, + tenant_id=tenant_id, + task_id=task_id, + branch_id=branch_id, + assigned_to_user_id=_assigned_user_filter(db, user), + partner_user_id=_partner_visibility_user_id(db, user), + financial_year=_active_financial_year(request), + ) + if not task: + return RedirectResponse(url="/services/work-tracker", status_code=303) + tool = get_task_tool(getattr(task, "task_tool_code", "NONE")) + if tool.code == "NONE" or not tool.launch_url: + return RedirectResponse(url=f"/services/work-tracker/tasks/{task.id}/edit?tool_error=not_configured", status_code=303) + run = start_task_tool_run(db, task=task, user_id=int(user.id)) + launch_url = build_task_tool_launch_url( + tool, + client_id=int(task.client_id), + task_id=int(task.id), + tool_run_id=int(run.id), + ) + return RedirectResponse(url=launch_url, status_code=303) + finally: + db.close() + + +@router.post("/tool-runs/{run_id}/result") +def task_tool_result_submit( + request: Request, + run_id: int, + status: str = Form("completed"), + summary: str = Form(""), + result_json: str = Form(""), + clean_pass: str = Form(""), + exception_count: str = Form(""), + output_file: UploadFile | None = File(None), + csrf_token: str = Form(...), +): + validate_csrf(request, csrf_token) + db = CommonSessionLocal() + try: + user = get_current_user(request, db=db) + if not user: + return RedirectResponse(url="/login", status_code=303) + tenant_id = _active_tenant_id(request, user) + run = get_task_tool_run(db, tenant_id=tenant_id, run_id=run_id) + if not run: + return RedirectResponse(url="/services/work-tracker", status_code=303) + branch_id = _active_branch_id(request, user, db) + task = get_task( + db, + tenant_id=tenant_id, + task_id=int(run.task_instance_id), + branch_id=branch_id, + assigned_to_user_id=_assigned_user_filter(db, user), + partner_user_id=_partner_visibility_user_id(db, user), + financial_year=_active_financial_year(request), + ) + if not task: + return _redirect_denied() + + clean_value = None + if clean_pass.strip().lower() in {"1", "true", "yes", "y", "on"}: + clean_value = True + elif clean_pass.strip().lower() in {"0", "false", "no", "n", "off"}: + clean_value = False + + exception_value = None + if exception_count.strip(): + try: + exception_value = max(0, int(exception_count)) + except ValueError: + return RedirectResponse(url=f"/services/work-tracker/tasks/{task.id}/edit?tool_error=invalid_exception_count", status_code=303) + + file_bytes = None + filename = None + content_type = None + if output_file is not None and getattr(output_file, "filename", None): + file_bytes = output_file.file.read() + filename = str(output_file.filename or "tool-output.bin") + content_type = str(output_file.content_type or "application/octet-stream") + + try: + record_task_tool_result( + db, + run=run, + status=status, + summary=summary, + result_payload=result_json, + clean_pass=clean_value, + exception_count=exception_value, + output_filename=filename, + output_content_type=content_type, + output_bytes=file_bytes, + completed_by_user_id=int(user.id), + ) + except ValueError: + db.rollback() + return RedirectResponse(url=f"/services/work-tracker/tasks/{task.id}/edit?tool_error=result_rejected", status_code=303) + return RedirectResponse(url=f"/services/work-tracker/tasks/{task.id}/edit?tool_saved=1", status_code=303) + finally: + db.close() + + +@router.get("/tool-runs/{run_id}/output") +def task_tool_result_download(request: Request, run_id: int): + db = CommonSessionLocal() + try: + user = get_current_user(request, db=db) + if not user: + return RedirectResponse(url="/login", status_code=303) + tenant_id = _active_tenant_id(request, user) + run = get_task_tool_run(db, tenant_id=tenant_id, run_id=run_id) + if not run or not run.output_bytes: + return RedirectResponse(url="/services/work-tracker", status_code=303) + branch_id = _active_branch_id(request, user, db) + task = get_task( + db, + tenant_id=tenant_id, + task_id=int(run.task_instance_id), + branch_id=branch_id, + assigned_to_user_id=_assigned_user_filter(db, user), + partner_user_id=_partner_visibility_user_id(db, user), + financial_year=_active_financial_year(request), + ) + if not task: + return _redirect_denied() + filename = (run.output_filename or f"task-tool-run-{run.id}.bin").replace('"', "") + return Response( + content=bytes(run.output_bytes), + media_type=run.output_content_type or "application/octet-stream", + headers={"Content-Disposition": f'attachment; filename="{filename}"'}, ) finally: db.close()