fixup(pipeline-hardening M4/M5): purge count, shared toast helper, clean imports
- Fix RabbitMQ purge_queue to report the pre-purge message count instead of the post-purge depth (which was always 0). - Add shared attach_toast helper in admin templating; remove duplicated inline _attach_toast implementations in review.py and mq.py. - Keep purge confirmation case-sensitive so the typed safeguard works as intended. - Fix test imports to use known-first-party 'contract_check' per isort config.
This commit is contained in:
parent
c00f3cf9f9
commit
0b2e7665f4
5 changed files with 29 additions and 29 deletions
|
|
@ -12,7 +12,7 @@ from fastapi import APIRouter, Depends, Form, HTTPException, Request, status
|
||||||
from fastapi.responses import HTMLResponse, RedirectResponse
|
from fastapi.responses import HTMLResponse, RedirectResponse
|
||||||
|
|
||||||
from src.contract_check.api.admin.auth import AdminUser, HtmxGuard, require_admin
|
from src.contract_check.api.admin.auth import AdminUser, HtmxGuard, require_admin
|
||||||
from src.contract_check.api.admin.templating import templates
|
from src.contract_check.api.admin.templating import attach_toast, templates
|
||||||
from src.contract_check.core.logging import get_logger
|
from src.contract_check.core.logging import get_logger
|
||||||
from src.contract_check.core.mq.management import (
|
from src.contract_check.core.mq.management import (
|
||||||
MqManagementDisabledError,
|
MqManagementDisabledError,
|
||||||
|
|
@ -36,12 +36,6 @@ def _client() -> RabbitMQManagementClient:
|
||||||
return RabbitMQManagementClient()
|
return RabbitMQManagementClient()
|
||||||
|
|
||||||
|
|
||||||
def _attach_toast(response: RedirectResponse, toast: str) -> None:
|
|
||||||
import json
|
|
||||||
|
|
||||||
response.headers["HX-Trigger"] = json.dumps({"showToast": toast})
|
|
||||||
|
|
||||||
|
|
||||||
def _safe_toast(toast: str) -> str:
|
def _safe_toast(toast: str) -> str:
|
||||||
# Keep toast text short and URL-safe for redirects.
|
# Keep toast text short and URL-safe for redirects.
|
||||||
return toast[:200]
|
return toast[:200]
|
||||||
|
|
@ -184,7 +178,7 @@ async def requeue_dlq(
|
||||||
response = RedirectResponse(
|
response = RedirectResponse(
|
||||||
url=f"/admin/mq/dlq/{queue_name}", status_code=status.HTTP_303_SEE_OTHER
|
url=f"/admin/mq/dlq/{queue_name}", status_code=status.HTTP_303_SEE_OTHER
|
||||||
)
|
)
|
||||||
_attach_toast(response, toast)
|
attach_toast(response, toast)
|
||||||
return response
|
return response
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -221,11 +215,14 @@ async def purge_dlq(
|
||||||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="unknown DLQ")
|
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="unknown DLQ")
|
||||||
|
|
||||||
expected = f"purge {queue_name}"
|
expected = f"purge {queue_name}"
|
||||||
if confirmation.strip().lower() != expected.lower():
|
if confirmation.strip() != expected:
|
||||||
return RedirectResponse(
|
toast = f"Неверное подтверждение. Введите '{expected}'"
|
||||||
url=f"/admin/mq/dlq/{queue_name}/purge/confirm?toast=Неверное+подтверждение.+Введите+%27{expected}%27",
|
response = RedirectResponse(
|
||||||
|
url=f"/admin/mq/dlq/{queue_name}/purge/confirm",
|
||||||
status_code=status.HTTP_303_SEE_OTHER,
|
status_code=status.HTTP_303_SEE_OTHER,
|
||||||
)
|
)
|
||||||
|
attach_toast(response, toast)
|
||||||
|
return response
|
||||||
|
|
||||||
client = _client()
|
client = _client()
|
||||||
try:
|
try:
|
||||||
|
|
@ -252,5 +249,5 @@ async def purge_dlq(
|
||||||
response = RedirectResponse(
|
response = RedirectResponse(
|
||||||
url=f"/admin/mq/dlq/{queue_name}", status_code=status.HTTP_303_SEE_OTHER
|
url=f"/admin/mq/dlq/{queue_name}", status_code=status.HTTP_303_SEE_OTHER
|
||||||
)
|
)
|
||||||
_attach_toast(response, toast)
|
attach_toast(response, toast)
|
||||||
return response
|
return response
|
||||||
|
|
|
||||||
|
|
@ -7,7 +7,6 @@ All mutations are HTMX-driven and protected by the ``HX-Request`` CSRF guard.
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import datetime as dt
|
import datetime as dt
|
||||||
import json
|
|
||||||
import uuid
|
import uuid
|
||||||
from decimal import Decimal
|
from decimal import Decimal
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
|
|
@ -16,7 +15,7 @@ from fastapi import APIRouter, Depends, Form, HTTPException, Request, status
|
||||||
from fastapi.responses import HTMLResponse, RedirectResponse
|
from fastapi.responses import HTMLResponse, RedirectResponse
|
||||||
|
|
||||||
from src.contract_check.api.admin.auth import AdminUser, HtmxGuard, require_admin
|
from src.contract_check.api.admin.auth import AdminUser, HtmxGuard, require_admin
|
||||||
from src.contract_check.api.admin.templating import templates
|
from src.contract_check.api.admin.templating import attach_toast, templates
|
||||||
from src.contract_check.api.deps import AsyncSessionDep, PublisherDep
|
from src.contract_check.api.deps import AsyncSessionDep, PublisherDep
|
||||||
from src.contract_check.core.config import get_settings
|
from src.contract_check.core.config import get_settings
|
||||||
from src.contract_check.core.db.repositories import ReviewQueueRepository
|
from src.contract_check.core.db.repositories import ReviewQueueRepository
|
||||||
|
|
@ -124,7 +123,7 @@ async def list_review_queue(
|
||||||
"conf_max": conf_max,
|
"conf_max": conf_max,
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
_attach_toast(response, request.query_params.get("toast"))
|
attach_toast(response, request.query_params.get("toast"))
|
||||||
return response
|
return response
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -150,15 +149,10 @@ async def review_detail(
|
||||||
"toast": request.query_params.get("toast"),
|
"toast": request.query_params.get("toast"),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
_attach_toast(response, request.query_params.get("toast"))
|
attach_toast(response, request.query_params.get("toast"))
|
||||||
return response
|
return response
|
||||||
|
|
||||||
|
|
||||||
def _attach_toast(response: HTMLResponse, toast: str | None) -> None:
|
|
||||||
if toast:
|
|
||||||
response.headers["HX-Trigger"] = json.dumps({"showToast": toast})
|
|
||||||
|
|
||||||
|
|
||||||
def _review_action_redirect(
|
def _review_action_redirect(
|
||||||
*,
|
*,
|
||||||
document_id: uuid.UUID,
|
document_id: uuid.UUID,
|
||||||
|
|
|
||||||
|
|
@ -8,8 +8,10 @@ datetime rendering consistent.
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import datetime as dt
|
import datetime as dt
|
||||||
|
import json
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
|
from fastapi import Response
|
||||||
from fastapi.templating import Jinja2Templates
|
from fastapi.templating import Jinja2Templates
|
||||||
|
|
||||||
_TEMPLATES_DIR = Path(__file__).resolve().parent / "templates"
|
_TEMPLATES_DIR = Path(__file__).resolve().parent / "templates"
|
||||||
|
|
@ -33,5 +35,11 @@ def _yes_no_none(value: object) -> str:
|
||||||
return "—"
|
return "—"
|
||||||
|
|
||||||
|
|
||||||
|
def attach_toast(response: Response, toast: str | None) -> None:
|
||||||
|
"""Attach an HTMX toast trigger to any response that supports headers."""
|
||||||
|
if toast:
|
||||||
|
response.headers["HX-Trigger"] = json.dumps({"showToast": toast})
|
||||||
|
|
||||||
|
|
||||||
templates.env.filters["dt"] = _format_dt
|
templates.env.filters["dt"] = _format_dt
|
||||||
templates.env.filters["yn"] = _yes_no_none
|
templates.env.filters["yn"] = _yes_no_none
|
||||||
|
|
|
||||||
|
|
@ -240,6 +240,9 @@ class RabbitMQManagementClient:
|
||||||
"""Empty a queue and return the number of messages removed."""
|
"""Empty a queue and return the number of messages removed."""
|
||||||
self._ensure_enabled()
|
self._ensure_enabled()
|
||||||
url = f"/api/queues/{self._vhost}/{quote(queue_name, safe='')}/contents"
|
url = f"/api/queues/{self._vhost}/{quote(queue_name, safe='')}/contents"
|
||||||
|
# Capture the pre-purge depth so we can report how many were removed.
|
||||||
|
before_info = await self.get_queue_info(queue_name)
|
||||||
|
messages_before = max(0, before_info.messages)
|
||||||
try:
|
try:
|
||||||
resp = await self.client.delete(url)
|
resp = await self.client.delete(url)
|
||||||
except httpx.NetworkError as exc:
|
except httpx.NetworkError as exc:
|
||||||
|
|
@ -254,9 +257,7 @@ class RabbitMQManagementClient:
|
||||||
status_code=resp.status_code,
|
status_code=resp.status_code,
|
||||||
body=resp.text[:500],
|
body=resp.text[:500],
|
||||||
)
|
)
|
||||||
# Successful purge returns no body; the queue is now empty.
|
return PurgeResult(messages_removed=messages_before)
|
||||||
info = await self.get_queue_info(queue_name)
|
|
||||||
return PurgeResult(messages_removed=max(0, info.messages))
|
|
||||||
|
|
||||||
async def requeue_batch(
|
async def requeue_batch(
|
||||||
self,
|
self,
|
||||||
|
|
|
||||||
|
|
@ -12,7 +12,7 @@ import httpx
|
||||||
import pytest
|
import pytest
|
||||||
import respx
|
import respx
|
||||||
|
|
||||||
from src.contract_check.core.mq.management import (
|
from contract_check.core.mq.management import (
|
||||||
MqManagementDisabledError,
|
MqManagementDisabledError,
|
||||||
MqManagementNotFoundError,
|
MqManagementNotFoundError,
|
||||||
MqManagementResponseError,
|
MqManagementResponseError,
|
||||||
|
|
@ -23,7 +23,7 @@ from src.contract_check.core.mq.management import (
|
||||||
RabbitMQManagementClient,
|
RabbitMQManagementClient,
|
||||||
RequeueResult,
|
RequeueResult,
|
||||||
)
|
)
|
||||||
from src.contract_check.core.mq.topology import DLQ_FOR
|
from contract_check.core.mq.topology import DLQ_FOR
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
|
|
@ -132,18 +132,18 @@ async def test_peek_messages_empty_non_list(mgmt_client: RabbitMQManagementClien
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.anyio
|
@pytest.mark.anyio
|
||||||
async def test_purge_queue(mgmt_client: RabbitMQManagementClient) -> None:
|
async def test_purge_queue_reports_pre_purge_count(mgmt_client: RabbitMQManagementClient) -> None:
|
||||||
with respx.mock:
|
with respx.mock:
|
||||||
respx.delete("http://rabbitmq:15672/api/queues/%2F/extract.dlq/contents").mock(
|
respx.delete("http://rabbitmq:15672/api/queues/%2F/extract.dlq/contents").mock(
|
||||||
return_value=httpx.Response(204)
|
return_value=httpx.Response(204)
|
||||||
)
|
)
|
||||||
respx.get("http://rabbitmq:15672/api/queues/%2F/extract.dlq").mock(
|
respx.get("http://rabbitmq:15672/api/queues/%2F/extract.dlq").mock(
|
||||||
return_value=httpx.Response(
|
return_value=httpx.Response(
|
||||||
200, json={"name": "extract.dlq", "messages": 0, "state": "running"}
|
200, json={"name": "extract.dlq", "messages": 3, "state": "running"}
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
await mgmt_client.connect()
|
await mgmt_client.connect()
|
||||||
assert await mgmt_client.purge_queue("extract.dlq") == PurgeResult(messages_removed=0)
|
assert await mgmt_client.purge_queue("extract.dlq") == PurgeResult(messages_removed=3)
|
||||||
await mgmt_client.aclose()
|
await mgmt_client.aclose()
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue