Compare commits

...
6 Commits
Author SHA1 Message Date
ik d6bc5ce85c fix: alert on degraded community scheduler health
CI / backend-and-migrations (push) Canceled after 0s
CI / astro-build (push) Canceled after 0s
CI / dependency-audit (push) Canceled after 0s
CI / compose-e2e (push) Canceled after 0s
2026-09-11 07:44:30 +07:00
ik 5b2db1cf48 fix: reject idempotency key payload conflicts 2026-09-11 07:44:07 +07:00
ik 42bdec6c2f fix: isolate scheduler validation in bootstrap 2026-09-11 07:43:14 +07:00
ik 531f1aa82f fix: replay idempotent upload token safely 2026-09-11 07:42:41 +07:00
ik 486e4c4673 fix: wire report idempotency and bootstrap scheduler check 2026-09-11 07:41:55 +07:00
ik bea3b9ccbb fix: align recovery schema idempotency and record counts 2026-09-11 07:40:58 +07:00
7 changed files with 79 additions and 22 deletions
@@ -0,0 +1,30 @@
"""add recovery idempotency and import history columns
Revision ID: 20260910_recovery
Revises: 48094a7d1b92
"""
from alembic import op
import sqlalchemy as sa
revision = "20260910_recovery"
down_revision = "48094a7d1b92"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.add_column("submission_attempt", sa.Column("idempotency_key", sa.String(128), nullable=True))
op.add_column("submission_attempt", sa.Column("catch_report_id", sa.Uuid(), nullable=True))
op.add_column("submission_attempt", sa.Column("payload_hash", sa.String(64), nullable=True))
op.create_foreign_key("fk_submission_attempt_report", "submission_attempt", "catch_report", ["catch_report_id"], ["id"])
op.create_index("ix_submission_attempt_idempotency_key", "submission_attempt", ["idempotency_key"], unique=True)
op.add_column("import_record_event", sa.Column("changes", sa.JSON(), nullable=True))
op.add_column("import_record_event", sa.Column("provenance", sa.JSON(), nullable=True))
def downgrade() -> None:
op.drop_column("import_record_event", "provenance")
op.drop_column("import_record_event", "changes")
op.drop_index("ix_submission_attempt_idempotency_key", table_name="submission_attempt")
op.drop_constraint("fk_submission_attempt_report", "submission_attempt", type_="foreignkey")
op.drop_column("submission_attempt", "catch_report_id")
op.drop_column("submission_attempt", "payload_hash")
op.drop_column("submission_attempt", "idempotency_key")
+19 -14
View File
@@ -5,6 +5,7 @@ from datetime import datetime, timedelta, timezone
from ipaddress import IPv4Address, IPv6Address, IPv4Network, IPv6Network from ipaddress import IPv4Address, IPv6Address, IPv4Network, IPv6Network
import hashlib import hashlib
import hmac import hmac
import json
import logging import logging
import secrets import secrets
import time as time_module import time as time_module
@@ -287,14 +288,9 @@ def records(
query = query.join(CatchReport.waterbody).where(Waterbody.slug == waterbody) query = query.join(CatchReport.waterbody).where(Waterbody.slug == waterbody)
if category: if category:
query = query.where(CatchReport.raw_payload["category"].as_string() == category) query = query.where(CatchReport.raw_payload["category"].as_string() == category)
# Count total before pagination # Count from the exact same filtered query (before pagination).
total = db.scalar(select(func.count()).select_from(CatchReport).where(CatchReport.source_type == SourceType.official_record)) or 0 count_query = query.with_only_columns(func.count(CatchReport.id), maintain_column_froms=True).order_by(None)
if fish: total = db.scalar(count_query) or 0
total = db.scalar(select(func.count()).select_from(CatchReport).join(CatchReport.fish).where(Fish.slug == fish, CatchReport.source_type == SourceType.official_record)) or 0
if waterbody:
total = db.scalar(select(func.count()).select_from(CatchReport).join(CatchReport.waterbody).where(Waterbody.slug == waterbody, CatchReport.source_type == SourceType.official_record)) or 0
if category:
total = db.scalar(select(func.count()).select_from(CatchReport).where(CatchReport.source_type == SourceType.official_record, CatchReport.raw_payload["category"].as_string() == category)) or 0
items = list(db.scalars(query.order_by(CatchReport.caught_at.desc(), CatchReport.weight_g.desc(), CatchReport.id.desc()).offset(offset).limit(limit))) items = list(db.scalars(query.order_by(CatchReport.caught_at.desc(), CatchReport.weight_g.desc(), CatchReport.id.desc()).offset(offset).limit(limit)))
return PaginatedOfficialRecordOut( return PaginatedOfficialRecordOut(
items=[OfficialRecordOut(id=r.id, fish=r.fish.name_ru, weight_g=r.weight_g, waterbody=r.waterbody.name_ru, bait=r.bait.name if r.bait else None, player_name=r.player_name, record_date=r.caught_at, category=(r.raw_payload or {}).get("category"), region=(r.raw_payload or {}).get("region"), source_url=r.source_url) for r in items], items=[OfficialRecordOut(id=r.id, fish=r.fish.name_ru, weight_g=r.weight_g, waterbody=r.waterbody.name_ru, bait=r.bait.name if r.bait else None, player_name=r.player_name, record_date=r.caught_at, category=(r.raw_payload or {}).get("category"), region=(r.raw_payload or {}).get("region"), source_url=r.source_url) for r in items],
@@ -465,6 +461,7 @@ def create_catch_report(
) -> CatchReportAccepted: ) -> CatchReportAccepted:
if payload.website: if payload.website:
raise HTTPException(status_code=400, detail="invalid submission") raise HTTPException(status_code=400, detail="invalid submission")
payload_hash = hashlib.sha256(json.dumps(payload.model_dump(mode="json"), sort_keys=True, separators=(",", ":")).encode()).hexdigest()
# A05: Server-side idempotency — check BEFORE rate limit to avoid polluting table # A05: Server-side idempotency — check BEFORE rate limit to avoid polluting table
if idempotency_key: if idempotency_key:
key_hash = hmac.new(settings.rate_limit_secret.encode(), idempotency_key.encode(), hashlib.sha256).hexdigest() key_hash = hmac.new(settings.rate_limit_secret.encode(), idempotency_key.encode(), hashlib.sha256).hexdigest()
@@ -480,10 +477,17 @@ def create_catch_report(
if existing is not None: if existing is not None:
# Return 200 with idempotent flag — client can retry safely # Return 200 with idempotent flag — client can retry safely
logger.info("idempotent hit", extra={"idempotency_key": idempotency_key[:8]}) logger.info("idempotent hit", extra={"idempotency_key": idempotency_key[:8]})
return JSONResponse( report = existing.catch_report
status_code=200, if existing.payload_hash and not hmac.compare_digest(existing.payload_hash, payload_hash):
content={"id": "00000000-0000-0000-0000-000000000000", "moderation_status": "pending", "screenshot_upload_token": "", "idempotent": True}, raise HTTPException(status_code=409, detail="Idempotency-Key was already used with different payload")
) if report is None:
raise HTTPException(status_code=409, detail="idempotency record is incomplete; retry with a new key")
# Re-derive the one-time upload token from the idempotency key;
# only its hash is persisted, so the secret is never stored.
replay_token = hmac.new(settings.rate_limit_secret.encode(), (key_hash + ":upload").encode(), hashlib.sha256).hexdigest()
if not hmac.compare_digest(hashlib.sha256(replay_token.encode()).hexdigest(), report.screenshot_upload_token_hash or ""):
raise HTTPException(status_code=409, detail="idempotency record token mismatch; retry with a new key")
return JSONResponse(status_code=200, content={"id": str(report.id), "moderation_status": report.moderation_status.value, "screenshot_upload_token": replay_token, "idempotent": True})
logger.info("idempotency check miss", extra={"idempotency_key": idempotency_key[:8]}) logger.info("idempotency check miss", extra={"idempotency_key": idempotency_key[:8]})
_check_rate_limit(request, db) _check_rate_limit(request, db)
fish = db.scalar(select(Fish).where(Fish.slug == payload.fish_slug)) fish = db.scalar(select(Fish).where(Fish.slug == payload.fish_slug))
@@ -501,14 +505,15 @@ def create_catch_report(
if bait is None: if bait is None:
bait = Bait(name=payload.bait_name.strip(), normalized_name=key, kind=BaitKind.unknown) bait = Bait(name=payload.bait_name.strip(), normalized_name=key, kind=BaitKind.unknown)
db.add(bait) db.add(bait)
upload_token = secrets.token_urlsafe(32) upload_token = (hmac.new(settings.rate_limit_secret.encode(), (key_hash + ":upload").encode(), hashlib.sha256).hexdigest()
if idempotency_key else secrets.token_urlsafe(32))
report = CatchReport(fish=fish, spot=spot, waterbody=waterbody, bait=bait, weight_g=payload.weight_g, fishing_method=payload.fishing_method, rig_type=payload.rig_type, retrieve_method=payload.retrieve_method, retrieve_speed=payload.retrieve_speed, caught_at=payload.caught_at, reported_at=datetime.now(timezone.utc), player_name=payload.player_name, source_type=SourceType.user, source_url=payload.source_url, source_confidence=60, moderation_status=ModerationStatus.pending, raw_payload={"comment": payload.comment} if payload.comment else None, screenshot_upload_token_hash=hashlib.sha256(upload_token.encode()).hexdigest()) report = CatchReport(fish=fish, spot=spot, waterbody=waterbody, bait=bait, weight_g=payload.weight_g, fishing_method=payload.fishing_method, rig_type=payload.rig_type, retrieve_method=payload.retrieve_method, retrieve_speed=payload.retrieve_speed, caught_at=payload.caught_at, reported_at=datetime.now(timezone.utc), player_name=payload.player_name, source_type=SourceType.user, source_url=payload.source_url, source_confidence=60, moderation_status=ModerationStatus.pending, raw_payload={"comment": payload.comment} if payload.comment else None, screenshot_upload_token_hash=hashlib.sha256(upload_token.encode()).hexdigest())
db.add(report) db.add(report)
db.commit() db.commit()
# Store idempotency key if provided # Store idempotency key if provided
if idempotency_key: if idempotency_key:
key_hash = hmac.new(settings.rate_limit_secret.encode(), idempotency_key.encode(), hashlib.sha256).hexdigest() key_hash = hmac.new(settings.rate_limit_secret.encode(), idempotency_key.encode(), hashlib.sha256).hexdigest()
db.add(SubmissionAttempt(client_hash="", idempotency_key=key_hash, created_at=datetime.now(timezone.utc))) db.add(SubmissionAttempt(client_hash="", idempotency_key=key_hash, catch_report_id=report.id, payload_hash=payload_hash, created_at=datetime.now(timezone.utc)))
db.commit() db.commit()
logger.info("idempotency key stored", extra={"idempotency_key": idempotency_key[:8]}) logger.info("idempotency key stored", extra={"idempotency_key": idempotency_key[:8]})
return CatchReportAccepted(id=report.id, moderation_status=report.moderation_status.value, screenshot_upload_token=upload_token, idempotent=False) return CatchReportAccepted(id=report.id, moderation_status=report.moderation_status.value, screenshot_upload_token=upload_token, idempotent=False)
+3
View File
@@ -137,7 +137,10 @@ class SubmissionAttempt(Base):
id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4) id: Mapped[uuid.UUID] = mapped_column(primary_key=True, default=uuid.uuid4)
client_hash: Mapped[str] = mapped_column(String(64), index=True) client_hash: Mapped[str] = mapped_column(String(64), index=True)
idempotency_key: Mapped[str | None] = mapped_column(String(128), index=True) idempotency_key: Mapped[str | None] = mapped_column(String(128), index=True)
catch_report_id: Mapped[uuid.UUID | None] = mapped_column(ForeignKey("catch_report.id"), nullable=True)
payload_hash: Mapped[str | None] = mapped_column(String(64), nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), index=True) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), index=True)
catch_report: Mapped[CatchReport | None] = relationship()
class DataSource(Base): class DataSource(Base):
+5
View File
@@ -439,3 +439,8 @@ def test_catch_report_idempotency_key_prevents_duplicates(monkeypatch) -> None:
second = client.post("/api/v1/catch-reports", json=payload, headers=headers) second = client.post("/api/v1/catch-reports", json=payload, headers=headers)
assert second.status_code == 200, f"Expected 200, got {second.status_code}. Response: {second.json()}" assert second.status_code == 200, f"Expected 200, got {second.status_code}. Response: {second.json()}"
assert second.json()["idempotent"] is True assert second.json()["idempotent"] is True
assert second.json()["id"] == report_id
assert second.json()["screenshot_upload_token"] == first.json()["screenshot_upload_token"]
changed = dict(payload, weight_g=7800)
conflict = client.post("/api/v1/catch-reports", json=changed, headers=headers)
assert conflict.status_code == 409
+8 -1
View File
@@ -12,8 +12,14 @@ export const POST: APIRoute = async ({ request, redirect, cookies }) => {
const payload = { fish_slug: String(form.get("fish_slug") || ""), waterbody_slug: String(form.get("waterbody_slug") || ""), x: Number(form.get("x")), y: Number(form.get("y")), weight_g: Number(form.get("weight_g")), bait_name: text("bait_name"), fishing_method: text("fishing_method"), retrieve_method: text("retrieve_method"), retrieve_speed: number("retrieve_speed"), player_name: text("player_name"), comment: text("comment"), website: String(form.get("website") || "") }; const payload = { fish_slug: String(form.get("fish_slug") || ""), waterbody_slug: String(form.get("waterbody_slug") || ""), x: Number(form.get("x")), y: Number(form.get("y")), weight_g: Number(form.get("weight_g")), bait_name: text("bait_name"), fishing_method: text("fishing_method"), retrieve_method: text("retrieve_method"), retrieve_speed: number("retrieve_speed"), player_name: text("player_name"), comment: text("comment"), website: String(form.get("website") || "") };
let createdId: string | null = null; let createdId: string | null = null;
let uploadToken: string | null = null; let uploadToken: string | null = null;
// Keep one key across retries so a timeout cannot create a duplicate report.
let idempotencyKey = cookies.get("rf4-idempotency-key")?.value;
if (!idempotencyKey) {
idempotencyKey = crypto.randomUUID();
cookies.set("rf4-idempotency-key", idempotencyKey, { httpOnly: true, sameSite: "strict", secure: import.meta.env.PROD, path: "/", maxAge: 900 });
}
try { try {
const response = await fetch(`${base}/api/v1/catch-reports`, { method: "POST", headers: { "content-type": "application/json", "X-Forwarded-For": request.headers.get("x-forwarded-for") || request.headers.get("x-real-ip") || "unknown" }, body: JSON.stringify(payload), signal: AbortSignal.timeout(30_000) }); const response = await fetch(`${base}/api/v1/catch-reports`, { method: "POST", headers: { "content-type": "application/json", "Idempotency-Key": idempotencyKey, "X-Forwarded-For": request.headers.get("x-forwarded-for") || request.headers.get("x-real-ip") || "unknown" }, body: JSON.stringify(payload), signal: AbortSignal.timeout(30_000) });
if (!response.ok) { if (!response.ok) {
if (response.status === 429) return redirect("/report?state=rate_limited", 303); if (response.status === 429) return redirect("/report?state=rate_limited", 303);
if (response.status >= 500) return redirect("/report?state=server_error", 303); if (response.status >= 500) return redirect("/report?state=server_error", 303);
@@ -32,6 +38,7 @@ export const POST: APIRoute = async ({ request, redirect, cookies }) => {
return redirect(`/report?state=screenshot_error&report_id=${encodeURIComponent(created.id)}`, 303); return redirect(`/report?state=screenshot_error&report_id=${encodeURIComponent(created.id)}`, 303);
} }
} }
cookies.delete("rf4-idempotency-key", { path: "/" });
return redirect("/report?state=sent", 303); return redirect("/report?state=sent", 303);
} }
catch (err) { catch (err) {
+12
View File
@@ -45,6 +45,18 @@ done
ready_payload=$(curl -fsS --max-time 10 "$base_url/ready" 2>/dev/null) || ready_payload="" ready_payload=$(curl -fsS --max-time 10 "$base_url/ready" 2>/dev/null) || ready_payload=""
printf '%s' "$ready_payload" | grep -Eq '"status"[[:space:]]*:[[:space:]]*"ready"' || fail "readiness" printf '%s' "$ready_payload" | grep -Eq '"status"[[:space:]]*:[[:space:]]*"ready"' || fail "readiness"
# Import health is intentionally non-blocking for /ready, but must still alert.
# Parse the nested component instead of accepting any unrelated "ready" field.
community_status=$(printf '%s' "$ready_payload" | python3 -c 'import json,sys
try:
payload=json.load(sys.stdin)
print(payload.get("components",{}).get("community_scheduler",{}).get("status", "unknown"))
except Exception:
print("unknown")')
case "$community_status" in
ready) : ;;
*) fail "community-scheduler:${community_status}" ;;
esac
disk_used=$(df -Pk "$repo" | awk 'NR==2 {gsub(/%/, "", $5); print $5}') disk_used=$(df -Pk "$repo" | awk 'NR==2 {gsub(/%/, "", $5); print $5}')
case "$disk_used" in case "$disk_used" in
+2 -7
View File
@@ -35,18 +35,13 @@ docker compose --env-file .env.production.example -f compose.production.yaml -f
} }
echo "Caddy configuration valid ✓" echo "Caddy configuration valid ✓"
# A10: Check scheduler can start without external network # A10: Check scheduler module can load without running real sources
echo "Validating community scheduler..." echo "Validating community scheduler..."
$compose up --build -d --wait community-scheduler $compose run --rm --no-deps community-scheduler python -c "import app.community_scheduler; print('scheduler module loads OK')" >/dev/null 2>&1 || {
# Scheduler runs in a loop, check it's healthy by verifying the process is running
$compose exec -T community-scheduler python -c "from app.community_scheduler import schedule_interval; print('scheduler module loads OK')" >/dev/null 2>&1 || {
echo "ERROR: Community scheduler failed to start" >&2 echo "ERROR: Community scheduler failed to start" >&2
$compose logs --no-color community-scheduler >&2
exit 1 exit 1
} }
echo "Community scheduler valid ✓" echo "Community scheduler valid ✓"
# Stop scheduler to free resources
$compose stop community-scheduler
# A10: Extract Alembic revision ID programmatically, handle multiple heads # A10: Extract Alembic revision ID programmatically, handle multiple heads
ALEMBIC_HEADS_OUTPUT=$($compose exec -T api alembic heads 2>/dev/null || true) ALEMBIC_HEADS_OUTPUT=$($compose exec -T api alembic heads 2>/dev/null || true)