feature/mercancias-permisos
This commit is contained in:
@@ -9,7 +9,7 @@ from uuid import uuid4
|
||||
|
||||
from fastapi import APIRouter, File, HTTPException, Query, UploadFile, Depends
|
||||
from sqlalchemy.orm import Session
|
||||
from typing import Dict, Any
|
||||
from typing import Dict, Any, Optional
|
||||
|
||||
from core.celery_app import celery_app
|
||||
from core.database import get_core_db
|
||||
@@ -39,6 +39,50 @@ def _get_redis():
|
||||
return redis.Redis.from_url(url, decode_responses=False)
|
||||
|
||||
|
||||
def _assert_classes_csv_job_access(
|
||||
db: Session,
|
||||
job_id: str,
|
||||
current_user: Dict[str, Any],
|
||||
) -> None:
|
||||
from api.v1.modules.core.tasks_tracking.models import TaskRun
|
||||
|
||||
company_id: Optional[int] = None
|
||||
row = db.query(TaskRun).filter(TaskRun.task_id == job_id).first()
|
||||
if row is not None and row.company_id is not None:
|
||||
company_id = int(row.company_id)
|
||||
|
||||
r = _get_redis()
|
||||
if company_id is None:
|
||||
raw = r.get(f"{CLS_IMPORT_META_PREFIX}{job_id}")
|
||||
if raw:
|
||||
meta = json.loads(raw.decode("utf-8"))
|
||||
cid = meta.get("company_id")
|
||||
if cid is not None:
|
||||
company_id = int(cid)
|
||||
|
||||
if company_id is None and row is not None and row.meta_payload:
|
||||
layout_jid = row.meta_payload.get("layout_import_job_id")
|
||||
if layout_jid:
|
||||
raw2 = r.get(f"{CLS_IMPORT_META_PREFIX}{layout_jid}")
|
||||
if raw2:
|
||||
meta2 = json.loads(raw2.decode("utf-8"))
|
||||
cid2 = meta2.get("company_id")
|
||||
if cid2 is not None:
|
||||
company_id = int(cid2)
|
||||
|
||||
if company_id is None:
|
||||
raise HTTPException(status_code=403, detail="Sin acceso a este job")
|
||||
try:
|
||||
validate_access_to_resource(
|
||||
db, company_id, current_user, ["csv_upload.process"]
|
||||
)
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.error("Classes import: job access validation failed: %s", e)
|
||||
raise HTTPException(status_code=403, detail="Sin acceso a este job") from None
|
||||
|
||||
|
||||
@router.post("/upload", response_model=ImportJobResponse)
|
||||
async def upload_import_file(
|
||||
file: UploadFile = File(...),
|
||||
@@ -110,7 +154,12 @@ async def upload_import_file(
|
||||
|
||||
|
||||
@router.get("/{job_id}/status")
|
||||
async def get_import_status(job_id: str):
|
||||
async def get_import_status(
|
||||
job_id: str,
|
||||
db: Session = Depends(get_core_db),
|
||||
current_user: Dict[str, Any] = Depends(get_current_user),
|
||||
):
|
||||
_assert_classes_csv_job_access(db, job_id, current_user)
|
||||
task_result = celery_app.AsyncResult(job_id)
|
||||
|
||||
if task_result.state == "PENDING":
|
||||
@@ -192,6 +241,7 @@ async def commit_import_job(
|
||||
task_name="classes_insert_valid_rows",
|
||||
task_origin="a76/layouts_csv/classes/commit",
|
||||
args=[job_id],
|
||||
required_permissions=["csv_upload.process"],
|
||||
)
|
||||
return {
|
||||
"status": "committing",
|
||||
@@ -201,5 +251,10 @@ async def commit_import_job(
|
||||
|
||||
|
||||
@router.get("/{job_id}/errors/scan-csv")
|
||||
async def download_scan_errors_csv(job_id: str):
|
||||
async def download_scan_errors_csv(
|
||||
job_id: str,
|
||||
db: Session = Depends(get_core_db),
|
||||
current_user: Dict[str, Any] = Depends(get_current_user),
|
||||
):
|
||||
_assert_classes_csv_job_access(db, job_id, current_user)
|
||||
return download_scan_errors_csv_stream("cls", job_id)
|
||||
|
||||
@@ -29,6 +29,9 @@ def dispatch_tracked_layouts_csv_commit(
|
||||
"""
|
||||
Lee tenant_id / company_id del meta en Redis, valida acceso y despacha la tarea Celery
|
||||
con un task_id nuevo (commit) para que el polling use commit_job_id.
|
||||
|
||||
Si ``required_permissions`` es None, se normaliza a [] para validar acceso a compañía
|
||||
(``validate_access_to_resource(..., None)`` omitiría esa validación por el atajo tipo /me).
|
||||
"""
|
||||
raw = redis_client.get(meta_redis_key)
|
||||
if not raw:
|
||||
@@ -36,22 +39,29 @@ def dispatch_tracked_layouts_csv_commit(
|
||||
meta = json.loads(raw.decode("utf-8"))
|
||||
tenant_id = int(meta["tenant_id"])
|
||||
company_id = meta.get("company_id")
|
||||
if company_id is not None:
|
||||
try:
|
||||
validate_access_to_resource(
|
||||
db, int(company_id), current_user, required_permissions
|
||||
)
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception:
|
||||
raise HTTPException(status_code=403, detail="Sin acceso a este job") from None
|
||||
if company_id is None:
|
||||
raise HTTPException(
|
||||
status_code=403,
|
||||
detail="Job sin compañía asociada o meta incompleto",
|
||||
)
|
||||
effective_permissions: list[str] = (
|
||||
[] if required_permissions is None else required_permissions
|
||||
)
|
||||
try:
|
||||
validate_access_to_resource(
|
||||
db, int(company_id), current_user, effective_permissions
|
||||
)
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception:
|
||||
raise HTTPException(status_code=403, detail="Sin acceso a este job") from None
|
||||
|
||||
commit_id = str(uuid4())
|
||||
track_and_dispatch(
|
||||
db=db,
|
||||
task=celery_task,
|
||||
tenant_id=tenant_id,
|
||||
company_id=int(company_id) if company_id is not None else None,
|
||||
company_id=int(company_id),
|
||||
requested_by_user=current_user.get("preferred_username")
|
||||
or current_user.get("email")
|
||||
or current_user.get("sub"),
|
||||
|
||||
@@ -9,9 +9,9 @@ from uuid import uuid4
|
||||
|
||||
from fastapi import APIRouter, File, HTTPException, Query, UploadFile, Depends
|
||||
from sqlalchemy.orm import Session
|
||||
from typing import Dict, Any
|
||||
|
||||
from typing import Dict, Any, Optional
|
||||
|
||||
from core.celery_app import celery_app
|
||||
from core.database import get_core_db
|
||||
from core.security import get_current_user, validate_access_to_resource
|
||||
from api.v1.modules.core.tasks_tracking import track_and_dispatch
|
||||
@@ -39,6 +39,50 @@ def _get_redis():
|
||||
return redis.Redis.from_url(url, decode_responses=False)
|
||||
|
||||
|
||||
def _assert_parts_csv_job_access(
|
||||
db: Session,
|
||||
job_id: str,
|
||||
current_user: Dict[str, Any],
|
||||
) -> None:
|
||||
from api.v1.modules.core.tasks_tracking.models import TaskRun
|
||||
|
||||
company_id: Optional[int] = None
|
||||
row = db.query(TaskRun).filter(TaskRun.task_id == job_id).first()
|
||||
if row is not None and row.company_id is not None:
|
||||
company_id = int(row.company_id)
|
||||
|
||||
r = _get_redis()
|
||||
if company_id is None:
|
||||
raw = r.get(f"{PART_IMPORT_META_PREFIX}{job_id}")
|
||||
if raw:
|
||||
meta = json.loads(raw.decode("utf-8"))
|
||||
cid = meta.get("company_id")
|
||||
if cid is not None:
|
||||
company_id = int(cid)
|
||||
|
||||
if company_id is None and row is not None and row.meta_payload:
|
||||
layout_jid = row.meta_payload.get("layout_import_job_id")
|
||||
if layout_jid:
|
||||
raw2 = r.get(f"{PART_IMPORT_META_PREFIX}{layout_jid}")
|
||||
if raw2:
|
||||
meta2 = json.loads(raw2.decode("utf-8"))
|
||||
cid2 = meta2.get("company_id")
|
||||
if cid2 is not None:
|
||||
company_id = int(cid2)
|
||||
|
||||
if company_id is None:
|
||||
raise HTTPException(status_code=403, detail="Sin acceso a este job")
|
||||
try:
|
||||
validate_access_to_resource(
|
||||
db, company_id, current_user, ["csv_upload.process"]
|
||||
)
|
||||
except HTTPException:
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.error("Parts import: job access validation failed: %s", e)
|
||||
raise HTTPException(status_code=403, detail="Sin acceso a este job") from None
|
||||
|
||||
|
||||
@router.post("/upload", response_model=ImportJobResponse)
|
||||
async def upload_import_file(
|
||||
file: UploadFile = File(...),
|
||||
@@ -108,8 +152,12 @@ async def upload_import_file(
|
||||
|
||||
|
||||
@router.get("/{job_id}/status")
|
||||
async def get_import_status(job_id: str):
|
||||
from core.celery_app import celery_app
|
||||
async def get_import_status(
|
||||
job_id: str,
|
||||
db: Session = Depends(get_core_db),
|
||||
current_user: Dict[str, Any] = Depends(get_current_user),
|
||||
):
|
||||
_assert_parts_csv_job_access(db, job_id, current_user)
|
||||
task_result = celery_app.AsyncResult(job_id)
|
||||
|
||||
if task_result.state == "PENDING":
|
||||
@@ -191,6 +239,7 @@ async def commit_import_job(
|
||||
task_name="parts_insert_valid_rows",
|
||||
task_origin="a76/layouts_csv/parts/commit",
|
||||
args=[job_id],
|
||||
required_permissions=["csv_upload.process"],
|
||||
)
|
||||
return {
|
||||
"status": "committing",
|
||||
@@ -200,5 +249,10 @@ async def commit_import_job(
|
||||
|
||||
|
||||
@router.get("/{job_id}/errors/scan-csv")
|
||||
async def download_scan_errors_csv(job_id: str):
|
||||
async def download_scan_errors_csv(
|
||||
job_id: str,
|
||||
db: Session = Depends(get_core_db),
|
||||
current_user: Dict[str, Any] = Depends(get_current_user),
|
||||
):
|
||||
_assert_parts_csv_job_access(db, job_id, current_user)
|
||||
return download_scan_errors_csv_stream("part", job_id)
|
||||
|
||||
@@ -30,6 +30,7 @@ router.include_router(
|
||||
enable_filters=True,
|
||||
max_page_size=1000,
|
||||
list_permissions=["goods_parts.view"],
|
||||
get_permissions=["goods_parts.view"],
|
||||
create_permissions=["goods_parts.create"],
|
||||
update_permissions=["goods_parts.edit"],
|
||||
delete_permissions=["goods_parts.delete"],
|
||||
|
||||
Reference in New Issue
Block a user