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' %}
{{ task_tool.description }}
Opening the tool does not automatically mark this task complete. Save the task response/evidence after completing the verification.
+{{ task_tool.description }}
+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.
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.
+{{ run.summary }}
{% endif %} +