Merge pull request 'feature/table-celery-tasks' (#247) from feature/table-celery-tasks into development
Reviewed-on: ADUANASOFT/anexo76#247
This commit is contained in:
@@ -6,6 +6,7 @@ from .user_tenant.routes import router as user_tenant_router
|
||||
from .users.routes import router as users_router
|
||||
from .dashboard.routes import router as dashboard_router
|
||||
from .help_center.routes import router as help_center_router
|
||||
from .tasks_tracking.routes import router as tasks_tracking_router
|
||||
from fastapi import APIRouter
|
||||
|
||||
router = APIRouter()
|
||||
@@ -18,3 +19,4 @@ router.include_router(licenses_router, prefix="/core", tags=["core / licenses"])
|
||||
router.include_router(permissions_router, prefix="/core", tags=["core / permissions"])
|
||||
router.include_router(dashboard_router, prefix="/core", tags=["core / dashboard"])
|
||||
router.include_router(help_center_router, prefix="/core", tags=["core / help-center"])
|
||||
router.include_router(tasks_tracking_router, prefix="/core", tags=["core / tasks"])
|
||||
|
||||
4
backend/api/v1/modules/core/tasks_tracking/__init__.py
Normal file
4
backend/api/v1/modules/core/tasks_tracking/__init__.py
Normal file
@@ -0,0 +1,4 @@
|
||||
from .dispatch import track_and_dispatch
|
||||
from .service import TaskTrackerService
|
||||
|
||||
__all__ = ["track_and_dispatch", "TaskTrackerService"]
|
||||
36
backend/api/v1/modules/core/tasks_tracking/dispatch.py
Normal file
36
backend/api/v1/modules/core/tasks_tracking/dispatch.py
Normal file
@@ -0,0 +1,36 @@
|
||||
from typing import Any
|
||||
|
||||
from celery import Task
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from .service import TaskTrackerService
|
||||
|
||||
|
||||
def track_and_dispatch(
|
||||
*,
|
||||
db: Session,
|
||||
task: Task,
|
||||
tenant_id: int,
|
||||
task_name: str,
|
||||
task_group: str,
|
||||
company_id: int | None = None,
|
||||
requested_by_user: str | None = None,
|
||||
task_origin: str | None = None,
|
||||
args: list[Any] | None = None,
|
||||
kwargs: dict[str, Any] | None = None,
|
||||
task_id: str | None = None,
|
||||
meta_payload: dict[str, Any] | None = None,
|
||||
):
|
||||
celery_task = task.apply_async(args=args or [], kwargs=kwargs or {}, task_id=task_id)
|
||||
tracker = TaskTrackerService(db)
|
||||
tracker.register_dispatch(
|
||||
task_id=celery_task.id,
|
||||
tenant_id=tenant_id,
|
||||
company_id=company_id,
|
||||
requested_by_user=requested_by_user,
|
||||
task_name=task_name,
|
||||
task_group=task_group,
|
||||
task_origin=task_origin,
|
||||
meta_payload=meta_payload,
|
||||
)
|
||||
return celery_task
|
||||
62
backend/api/v1/modules/core/tasks_tracking/models.py
Normal file
62
backend/api/v1/modules/core/tasks_tracking/models.py
Normal file
@@ -0,0 +1,62 @@
|
||||
import enum
|
||||
from datetime import datetime
|
||||
|
||||
from sqlalchemy import DateTime, ForeignKey, Index, Integer, String, Text, func
|
||||
from sqlalchemy.dialects.postgresql import JSONB
|
||||
from sqlalchemy.orm import Mapped, mapped_column
|
||||
|
||||
from core.database import Base
|
||||
|
||||
|
||||
class TaskStatus(str, enum.Enum):
|
||||
PENDING = "pending"
|
||||
ACTIVE = "active"
|
||||
COMPLETED = "completed"
|
||||
FAILED = "failed"
|
||||
|
||||
|
||||
class TaskRun(Base):
|
||||
__tablename__ = "task_runs"
|
||||
__table_args__ = (
|
||||
Index("ix_core_task_runs_task_id", "task_id", unique=True),
|
||||
Index("ix_core_task_runs_tenant_updated", "tenant_id", "updated_at"),
|
||||
Index("ix_core_task_runs_tenant_status_updated", "tenant_id", "status", "updated_at"),
|
||||
Index("ix_core_task_runs_tenant_group_updated", "tenant_id", "task_group", "updated_at"),
|
||||
Index("ix_core_task_runs_tenant_company_updated", "tenant_id", "company_id", "updated_at"),
|
||||
{"schema": "core"},
|
||||
)
|
||||
|
||||
id: Mapped[int] = mapped_column(Integer, primary_key=True)
|
||||
task_id: Mapped[str] = mapped_column(String(255), nullable=False)
|
||||
tenant_id: Mapped[int] = mapped_column(
|
||||
Integer, ForeignKey("core.tenants.id"), nullable=False, index=True
|
||||
)
|
||||
company_id: Mapped[int | None] = mapped_column(
|
||||
Integer, ForeignKey("a76.company.id"), nullable=True, index=True
|
||||
)
|
||||
requested_by_user: Mapped[str | None] = mapped_column(String(255), nullable=True)
|
||||
task_name: Mapped[str] = mapped_column(String(255), nullable=False)
|
||||
task_group: Mapped[str] = mapped_column(String(100), nullable=False)
|
||||
task_origin: Mapped[str | None] = mapped_column(String(255), nullable=True)
|
||||
status: Mapped[TaskStatus] = mapped_column(
|
||||
String(20), nullable=False, default=TaskStatus.PENDING.value
|
||||
)
|
||||
celery_state_raw: Mapped[str] = mapped_column(String(30), nullable=False, default="PENDING")
|
||||
progress_current: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
progress_total: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
progress_percent: Mapped[float | None] = mapped_column(nullable=True)
|
||||
progress_message: Mapped[str | None] = mapped_column(String(500), nullable=True)
|
||||
retries: Mapped[int | None] = mapped_column(Integer, nullable=True)
|
||||
exception_type: Mapped[str | None] = mapped_column(String(255), nullable=True)
|
||||
exception_message: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
traceback_excerpt: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
result_summary: Mapped[dict | None] = mapped_column(JSONB, nullable=True)
|
||||
meta_payload: Mapped[dict | None] = mapped_column(JSONB, nullable=True)
|
||||
started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
|
||||
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), nullable=False, server_default=func.now()
|
||||
)
|
||||
updated_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), nullable=False, server_default=func.now(), onupdate=func.now()
|
||||
)
|
||||
150
backend/api/v1/modules/core/tasks_tracking/routes.py
Normal file
150
backend/api/v1/modules/core/tasks_tracking/routes.py
Normal file
@@ -0,0 +1,150 @@
|
||||
from typing import Any
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from core.database import get_core_db
|
||||
from core.security import get_current_user, get_tenant_from_token
|
||||
|
||||
from .models import TaskRun, TaskStatus
|
||||
from .schemas import TaskCatalogsResponse, TaskRunDetail, TaskRunListItem, TaskRunsResponse, TaskSyncRequest
|
||||
from .service import TaskTrackerService
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
||||
def _map_row(row: TaskRun) -> TaskRunListItem:
|
||||
pct = row.progress_percent
|
||||
if row.status == TaskStatus.COMPLETED.value and (pct is None or pct == 0):
|
||||
pct = 100.0
|
||||
return TaskRunListItem(
|
||||
task_id=row.task_id,
|
||||
task_name=row.task_name,
|
||||
task_group=row.task_group,
|
||||
task_origin=row.task_origin,
|
||||
status=row.status,
|
||||
celery_state_raw=row.celery_state_raw,
|
||||
progress={
|
||||
"current": row.progress_current,
|
||||
"total": row.progress_total,
|
||||
"percent": pct,
|
||||
"message": row.progress_message,
|
||||
},
|
||||
retries=row.retries,
|
||||
error=(
|
||||
{"type": row.exception_type, "message": row.exception_message}
|
||||
if row.exception_type or row.exception_message
|
||||
else None
|
||||
),
|
||||
tenant_id=row.tenant_id,
|
||||
company_id=row.company_id,
|
||||
requested_by_user=row.requested_by_user,
|
||||
started_at=row.started_at,
|
||||
finished_at=row.finished_at,
|
||||
created_at=row.created_at,
|
||||
updated_at=row.updated_at,
|
||||
)
|
||||
|
||||
|
||||
@router.get("/tasks", response_model=TaskRunsResponse)
|
||||
def list_tasks(
|
||||
page: int = Query(1, ge=1),
|
||||
page_size: int = Query(50, ge=1, le=200),
|
||||
status: list[str] | None = Query(None),
|
||||
task_group: list[str] | None = Query(None),
|
||||
task_name: list[str] | None = Query(None),
|
||||
company_id: int | None = Query(None),
|
||||
search: str | None = Query(None),
|
||||
sync_active: bool = Query(False),
|
||||
current_user: dict[str, Any] = Depends(get_current_user),
|
||||
db: Session = Depends(get_core_db),
|
||||
):
|
||||
tenant_id = get_tenant_from_token(current_user)
|
||||
if not tenant_id:
|
||||
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
|
||||
|
||||
tracker = TaskTrackerService(db)
|
||||
if sync_active:
|
||||
tracker.sync_active_tasks(tenant_id=tenant_id)
|
||||
|
||||
rows, total = tracker.list_tasks(
|
||||
tenant_id=tenant_id,
|
||||
page=page,
|
||||
page_size=page_size,
|
||||
status=status,
|
||||
task_group=task_group,
|
||||
task_name=task_name,
|
||||
company_id=company_id,
|
||||
search=search,
|
||||
)
|
||||
items = [_map_row(row) for row in rows]
|
||||
return TaskRunsResponse(
|
||||
items=items,
|
||||
total=total,
|
||||
page=page,
|
||||
page_size=page_size,
|
||||
has_next=(page * page_size) < total,
|
||||
)
|
||||
|
||||
|
||||
@router.get("/tasks/{task_id}", response_model=TaskRunDetail)
|
||||
def get_task_detail(
|
||||
task_id: str,
|
||||
sync: bool = Query(True),
|
||||
current_user: dict[str, Any] = Depends(get_current_user),
|
||||
db: Session = Depends(get_core_db),
|
||||
):
|
||||
tenant_id = get_tenant_from_token(current_user)
|
||||
if not tenant_id:
|
||||
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
|
||||
|
||||
row = db.query(TaskRun).filter(TaskRun.task_id == task_id, TaskRun.tenant_id == tenant_id).first()
|
||||
if not row:
|
||||
raise HTTPException(status_code=404, detail="Task not found")
|
||||
|
||||
tracker = TaskTrackerService(db)
|
||||
if sync:
|
||||
row = tracker.sync_task(row)
|
||||
|
||||
item = _map_row(row)
|
||||
return TaskRunDetail(
|
||||
**item.model_dump(),
|
||||
traceback_excerpt=row.traceback_excerpt,
|
||||
result_summary=row.result_summary,
|
||||
meta_payload=row.meta_payload,
|
||||
)
|
||||
|
||||
|
||||
@router.post("/tasks/sync")
|
||||
def sync_tasks(
|
||||
body: TaskSyncRequest,
|
||||
current_user: dict[str, Any] = Depends(get_current_user),
|
||||
db: Session = Depends(get_core_db),
|
||||
):
|
||||
tenant_id = get_tenant_from_token(current_user)
|
||||
if not tenant_id:
|
||||
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
|
||||
tracker = TaskTrackerService(db)
|
||||
updated = tracker.sync_active_tasks(tenant_id=tenant_id, task_ids=body.task_ids)
|
||||
return {"updated": updated}
|
||||
|
||||
|
||||
@router.get("/tasks/catalogs", response_model=TaskCatalogsResponse)
|
||||
def get_catalogs(
|
||||
current_user: dict[str, Any] = Depends(get_current_user),
|
||||
db: Session = Depends(get_core_db),
|
||||
):
|
||||
tenant_id = get_tenant_from_token(current_user)
|
||||
if not tenant_id:
|
||||
raise HTTPException(status_code=400, detail="Tenant ID not found in token")
|
||||
|
||||
groups = (
|
||||
db.query(TaskRun.task_group).filter(TaskRun.tenant_id == tenant_id).distinct().order_by(TaskRun.task_group).all()
|
||||
)
|
||||
names = db.query(TaskRun.task_name).filter(TaskRun.tenant_id == tenant_id).distinct().order_by(TaskRun.task_name).all()
|
||||
statuses = db.query(TaskRun.status).filter(TaskRun.tenant_id == tenant_id).distinct().order_by(TaskRun.status).all()
|
||||
return TaskCatalogsResponse(
|
||||
task_groups=[g[0] for g in groups if g[0]],
|
||||
task_names=[n[0] for n in names if n[0]],
|
||||
statuses=[s[0] for s in statuses if s[0]],
|
||||
)
|
||||
59
backend/api/v1/modules/core/tasks_tracking/schemas.py
Normal file
59
backend/api/v1/modules/core/tasks_tracking/schemas.py
Normal file
@@ -0,0 +1,59 @@
|
||||
from datetime import datetime
|
||||
from typing import Any
|
||||
|
||||
from pydantic import BaseModel
|
||||
|
||||
|
||||
class TaskProgress(BaseModel):
|
||||
current: int | None = None
|
||||
total: int | None = None
|
||||
percent: float | None = None
|
||||
message: str | None = None
|
||||
|
||||
|
||||
class TaskError(BaseModel):
|
||||
type: str | None = None
|
||||
message: str | None = None
|
||||
|
||||
|
||||
class TaskRunListItem(BaseModel):
|
||||
task_id: str
|
||||
task_name: str
|
||||
task_group: str
|
||||
task_origin: str | None = None
|
||||
status: str
|
||||
celery_state_raw: str
|
||||
progress: TaskProgress
|
||||
retries: int | None = None
|
||||
error: TaskError | None = None
|
||||
tenant_id: int
|
||||
company_id: int | None = None
|
||||
requested_by_user: str | None = None
|
||||
started_at: datetime | None = None
|
||||
finished_at: datetime | None = None
|
||||
created_at: datetime
|
||||
updated_at: datetime
|
||||
|
||||
|
||||
class TaskRunDetail(TaskRunListItem):
|
||||
traceback_excerpt: str | None = None
|
||||
result_summary: dict[str, Any] | None = None
|
||||
meta_payload: dict[str, Any] | None = None
|
||||
|
||||
|
||||
class TaskRunsResponse(BaseModel):
|
||||
items: list[TaskRunListItem]
|
||||
total: int
|
||||
page: int
|
||||
page_size: int
|
||||
has_next: bool
|
||||
|
||||
|
||||
class TaskSyncRequest(BaseModel):
|
||||
task_ids: list[str] | None = None
|
||||
|
||||
|
||||
class TaskCatalogsResponse(BaseModel):
|
||||
task_groups: list[str]
|
||||
task_names: list[str]
|
||||
statuses: list[str]
|
||||
242
backend/api/v1/modules/core/tasks_tracking/service.py
Normal file
242
backend/api/v1/modules/core/tasks_tracking/service.py
Normal file
@@ -0,0 +1,242 @@
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
from celery.result import AsyncResult
|
||||
from sqlalchemy import asc, desc, func, or_
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from core.celery_app import celery_app
|
||||
|
||||
from .models import TaskRun, TaskStatus
|
||||
|
||||
# No usar valid_rows/total_rows del resultado final: en SUCCESS distorsiona el % (errores de scan).
|
||||
_CELERY_TERMINAL_RAW = frozenset({"SUCCESS", "FAILURE", "REVOKED", "REJECTED"})
|
||||
|
||||
|
||||
def _progress_from_row_count_dict(d: dict[str, Any]) -> tuple[int, int] | None:
|
||||
total = d.get("total_rows")
|
||||
if not isinstance(total, (int, float)) or total <= 0:
|
||||
return None
|
||||
valid = d.get("valid_rows")
|
||||
if isinstance(valid, (int, float)):
|
||||
return int(valid), int(total)
|
||||
processed = d.get("processed_rows")
|
||||
if isinstance(processed, (int, float)):
|
||||
return int(processed), int(total)
|
||||
return None
|
||||
|
||||
|
||||
def normalize_celery_state(state: str | None) -> TaskStatus:
|
||||
raw = (state or "PENDING").upper()
|
||||
if raw in {"STARTED", "PROGRESS", "PROCESSING", "RETRY"}:
|
||||
return TaskStatus.ACTIVE
|
||||
if raw == "SUCCESS":
|
||||
return TaskStatus.COMPLETED
|
||||
if raw in {"FAILURE", "REVOKED", "REJECTED"}:
|
||||
return TaskStatus.FAILED
|
||||
return TaskStatus.PENDING
|
||||
|
||||
|
||||
def _extract_progress(result: AsyncResult, raw_state_upper: str) -> tuple[int | None, int | None, float | None, str | None]:
|
||||
payload = result.info if isinstance(result.info, dict) else {}
|
||||
current = payload.get("current")
|
||||
total = payload.get("total")
|
||||
message = payload.get("status")
|
||||
|
||||
if current is None and isinstance(result.result, dict):
|
||||
current = result.result.get("current")
|
||||
if total is None and isinstance(result.result, dict):
|
||||
total = result.result.get("total")
|
||||
if message is None and isinstance(result.result, dict):
|
||||
message = result.result.get("status")
|
||||
|
||||
percent = None
|
||||
if isinstance(current, (int, float)) and isinstance(total, (int, float)) and total > 0:
|
||||
percent = min(100.0, max(0.0, (float(current) / float(total)) * 100.0))
|
||||
|
||||
raw = (raw_state_upper or "PENDING").upper()
|
||||
if percent is None and raw not in _CELERY_TERMINAL_RAW:
|
||||
for src in (payload, result.result if isinstance(result.result, dict) else {}):
|
||||
if not isinstance(src, dict):
|
||||
continue
|
||||
pair = _progress_from_row_count_dict(src)
|
||||
if pair is None:
|
||||
continue
|
||||
cur_i, tot_i = pair
|
||||
current, total = cur_i, tot_i
|
||||
percent = min(100.0, max(0.0, (float(cur_i) / float(tot_i)) * 100.0))
|
||||
break
|
||||
|
||||
return (
|
||||
int(current) if isinstance(current, (int, float)) else None,
|
||||
int(total) if isinstance(total, (int, float)) else None,
|
||||
percent,
|
||||
str(message) if message is not None else None,
|
||||
)
|
||||
|
||||
|
||||
def _extract_failure(result: AsyncResult) -> tuple[str | None, str | None, str | None]:
|
||||
exception_type = None
|
||||
exception_message = None
|
||||
traceback_excerpt = None
|
||||
|
||||
err = result.result
|
||||
if isinstance(err, Exception):
|
||||
exception_type = type(err).__name__
|
||||
exception_message = str(err)
|
||||
elif isinstance(err, dict):
|
||||
exception_type = err.get("exc_type")
|
||||
exception_message = err.get("exc_message") or err.get("error") or err.get("message")
|
||||
elif err is not None:
|
||||
exception_message = str(err)
|
||||
|
||||
tb = getattr(result, "traceback", None)
|
||||
if isinstance(tb, str):
|
||||
traceback_excerpt = tb[-4000:]
|
||||
|
||||
return exception_type, exception_message, traceback_excerpt
|
||||
|
||||
|
||||
class TaskTrackerService:
|
||||
def __init__(self, db: Session):
|
||||
self.db = db
|
||||
|
||||
def register_dispatch(
|
||||
self,
|
||||
*,
|
||||
task_id: str,
|
||||
tenant_id: int,
|
||||
task_name: str,
|
||||
task_group: str,
|
||||
company_id: int | None = None,
|
||||
requested_by_user: str | None = None,
|
||||
task_origin: str | None = None,
|
||||
meta_payload: dict[str, Any] | None = None,
|
||||
) -> TaskRun:
|
||||
current = self.db.query(TaskRun).filter(TaskRun.task_id == task_id).first()
|
||||
if current:
|
||||
return current
|
||||
|
||||
row = TaskRun(
|
||||
task_id=task_id,
|
||||
tenant_id=tenant_id,
|
||||
company_id=company_id,
|
||||
requested_by_user=requested_by_user,
|
||||
task_name=task_name,
|
||||
task_group=task_group,
|
||||
task_origin=task_origin,
|
||||
status=TaskStatus.PENDING.value,
|
||||
celery_state_raw="PENDING",
|
||||
progress_current=0,
|
||||
progress_total=100,
|
||||
progress_percent=0.0,
|
||||
progress_message="Queued",
|
||||
retries=0,
|
||||
meta_payload=meta_payload,
|
||||
started_at=None,
|
||||
finished_at=None,
|
||||
)
|
||||
self.db.add(row)
|
||||
self.db.commit()
|
||||
self.db.refresh(row)
|
||||
return row
|
||||
|
||||
def sync_task(self, task_run: TaskRun) -> TaskRun:
|
||||
async_result = celery_app.AsyncResult(task_run.task_id)
|
||||
raw_state = (async_result.state or "PENDING").upper()
|
||||
normalized = normalize_celery_state(raw_state)
|
||||
now = datetime.now(timezone.utc)
|
||||
|
||||
task_run.celery_state_raw = raw_state
|
||||
task_run.status = normalized.value
|
||||
task_run.retries = int(getattr(async_result, "retries", 0) or 0)
|
||||
|
||||
current, total, percent, message = _extract_progress(async_result, raw_state)
|
||||
if current is not None:
|
||||
task_run.progress_current = current
|
||||
if total is not None:
|
||||
task_run.progress_total = total
|
||||
if percent is not None:
|
||||
task_run.progress_percent = percent
|
||||
if message:
|
||||
task_run.progress_message = message
|
||||
|
||||
if normalized == TaskStatus.ACTIVE and task_run.started_at is None:
|
||||
task_run.started_at = now
|
||||
|
||||
if normalized == TaskStatus.COMPLETED:
|
||||
if task_run.started_at is None:
|
||||
task_run.started_at = now
|
||||
task_run.finished_at = now
|
||||
task_run.exception_type = None
|
||||
task_run.exception_message = None
|
||||
task_run.traceback_excerpt = None
|
||||
if isinstance(async_result.result, dict):
|
||||
task_run.result_summary = async_result.result
|
||||
else:
|
||||
task_run.result_summary = {"result": str(async_result.result)}
|
||||
# register_dispatch seeds progress_percent=0; Celery SUCCESS often has no current/total meta
|
||||
if percent is None:
|
||||
task_run.progress_percent = 100.0
|
||||
|
||||
if normalized == TaskStatus.FAILED:
|
||||
if task_run.started_at is None:
|
||||
task_run.started_at = now
|
||||
task_run.finished_at = now
|
||||
etype, emsg, tb = _extract_failure(async_result)
|
||||
task_run.exception_type = etype
|
||||
task_run.exception_message = emsg
|
||||
task_run.traceback_excerpt = tb
|
||||
|
||||
self.db.add(task_run)
|
||||
self.db.commit()
|
||||
self.db.refresh(task_run)
|
||||
return task_run
|
||||
|
||||
def sync_active_tasks(self, tenant_id: int, task_ids: list[str] | None = None) -> int:
|
||||
query = self.db.query(TaskRun).filter(
|
||||
TaskRun.tenant_id == tenant_id, TaskRun.status.in_([TaskStatus.PENDING.value, TaskStatus.ACTIVE.value])
|
||||
)
|
||||
if task_ids:
|
||||
query = query.filter(TaskRun.task_id.in_(task_ids))
|
||||
rows = query.limit(200).all()
|
||||
for row in rows:
|
||||
self.sync_task(row)
|
||||
return len(rows)
|
||||
|
||||
def list_tasks(
|
||||
self,
|
||||
*,
|
||||
tenant_id: int,
|
||||
page: int,
|
||||
page_size: int,
|
||||
status: list[str] | None = None,
|
||||
task_group: list[str] | None = None,
|
||||
task_name: list[str] | None = None,
|
||||
company_id: int | None = None,
|
||||
search: str | None = None,
|
||||
order: str = "desc",
|
||||
) -> tuple[list[TaskRun], int]:
|
||||
query = self.db.query(TaskRun).filter(TaskRun.tenant_id == tenant_id)
|
||||
if status:
|
||||
query = query.filter(TaskRun.status.in_(status))
|
||||
if task_group:
|
||||
query = query.filter(TaskRun.task_group.in_(task_group))
|
||||
if task_name:
|
||||
query = query.filter(TaskRun.task_name.in_(task_name))
|
||||
if company_id is not None:
|
||||
# Tareas sin company_id (p. ej. reportes solo por tenant) deben seguir visibles
|
||||
query = query.filter(or_(TaskRun.company_id == company_id, TaskRun.company_id.is_(None)))
|
||||
if search:
|
||||
term = f"%{search}%"
|
||||
query = query.filter(
|
||||
(TaskRun.task_id.ilike(term))
|
||||
| (TaskRun.task_name.ilike(term))
|
||||
| (TaskRun.task_origin.ilike(term))
|
||||
| (TaskRun.exception_message.ilike(term))
|
||||
)
|
||||
|
||||
total = query.with_entities(func.count(TaskRun.id)).scalar() or 0
|
||||
order_expr = asc(TaskRun.updated_at) if order == "asc" else desc(TaskRun.updated_at)
|
||||
items = query.order_by(order_expr).offset((page - 1) * page_size).limit(page_size).all()
|
||||
return items, total
|
||||
Reference in New Issue
Block a user