Fixed dropbox file removing issue, UI navigation sticky mode
This commit is contained in:
+2
-4
@@ -27,7 +27,7 @@ from pydantic import BaseModel, EmailStr, Field
|
||||
from starlette.concurrency import run_in_threadpool
|
||||
from httpx import TransportError
|
||||
from supabase import Client, ClientOptions, PostgrestAPIError, create_client
|
||||
from backend import dropbox_storage, media_cleanup, portfolio_content
|
||||
from backend import dropbox_storage, media_cleanup, media_registry, portfolio_content
|
||||
|
||||
|
||||
BASE_DIR = Path(__file__).resolve().parent
|
||||
@@ -333,9 +333,7 @@ def save_media_metadata(upload_id: str, metadata: dict[str, Any]) -> None:
|
||||
client = get_supabase()
|
||||
if client is not None:
|
||||
try:
|
||||
result = client.table("article_media").upsert({"upload_id": upload_id, "metadata": metadata}).execute()
|
||||
if not result.data:
|
||||
raise RuntimeError("Media record was not returned.")
|
||||
media_registry.upsert(client, "article_media", [{"upload_id": upload_id, "metadata": metadata}])
|
||||
except Exception as exc:
|
||||
raise HTTPException(502, "The file uploaded, but its record could not be saved. Please retry.") from exc
|
||||
else:
|
||||
|
||||
@@ -9,6 +9,7 @@ from threading import RLock
|
||||
|
||||
from fastapi import HTTPException
|
||||
from starlette.concurrency import run_in_threadpool
|
||||
from backend import media_registry
|
||||
|
||||
LOCK = RLock()
|
||||
TABLE = "article_media_deletions"
|
||||
@@ -62,9 +63,7 @@ def enqueue(media):
|
||||
client = main.get_supabase()
|
||||
try:
|
||||
if client is not None:
|
||||
result = client.table(TABLE).upsert(list(records.values())).execute()
|
||||
if not result.data:
|
||||
raise RuntimeError("Removal queue was not saved")
|
||||
media_registry.upsert(client, TABLE, list(records.values()))
|
||||
else:
|
||||
save_local_queue(list({**{r["upload_id"]: r for r in pending()}, **records}.values()))
|
||||
except Exception as exc:
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
"""Confirmed, safely retryable writes for immutable upload-ID metadata."""
|
||||
import logging
|
||||
import re
|
||||
import time
|
||||
|
||||
from httpx import TransportError
|
||||
from supabase import PostgrestAPIError
|
||||
|
||||
logger = logging.getLogger("uvicorn.error")
|
||||
TRANSIENT_CODES = {"500", "502", "503", "504", "520", "522", "524"}
|
||||
|
||||
|
||||
class UnconfirmedMediaWrite(RuntimeError):
|
||||
pass
|
||||
|
||||
|
||||
def matches(rows, records):
|
||||
saved = {row["upload_id"]: row["metadata"] for row in (rows or [])}
|
||||
return all(saved.get(record["upload_id"]) == record["metadata"] for record in records)
|
||||
|
||||
|
||||
def upsert(client, table, records):
|
||||
"""Retry only identical upserts by primary key, never article inserts/deletes.
|
||||
|
||||
A timeout can arrive after commit. Read back the exact batch to confirm it,
|
||||
including servers/proxies that return no representation of a successful write.
|
||||
"""
|
||||
for attempt in range(3):
|
||||
try:
|
||||
result = client.table(table).upsert(records, on_conflict="upload_id").retry(False).execute()
|
||||
if matches(result.data, records):
|
||||
return
|
||||
raise UnconfirmedMediaWrite("Media write returned incomplete confirmation")
|
||||
except Exception as exc:
|
||||
code = str(getattr(exc, "code", ""))
|
||||
transient = isinstance(exc, (TransportError, UnconfirmedMediaWrite)) or (
|
||||
isinstance(exc, PostgrestAPIError) and code in TRANSIENT_CODES
|
||||
)
|
||||
safe_code = code if re.fullmatch(r"[A-Za-z0-9_]{1,20}", code) else "unknown"
|
||||
logger.warning("Media registry write: table=%s type=%s code=%s attempt=%s", table, type(exc).__name__, safe_code, attempt + 1)
|
||||
if not transient:
|
||||
raise
|
||||
try:
|
||||
confirmed = client.table(table).select("upload_id,metadata").in_(
|
||||
"upload_id", [record["upload_id"] for record in records],
|
||||
).retry(False).execute()
|
||||
if matches(confirmed.data, records):
|
||||
return
|
||||
except Exception:
|
||||
# Preserve the original write failure if neither operation works.
|
||||
pass
|
||||
if attempt == 2:
|
||||
raise
|
||||
time.sleep(0.25 * (attempt + 1))
|
||||
@@ -216,13 +216,14 @@ class ArticleMediaTests(unittest.TestCase):
|
||||
client.table.side_effect = lambda name: media_table if name == "article_media" else article_table
|
||||
metadata = {"url": "/api/uploads/" + "a" * 32, "name": "image.png", "size": 10,
|
||||
"media_type": "image/png", "storage": "dropbox", "dropbox_file_id": "id:remote_photo"}
|
||||
media_table.upsert.return_value.retry.return_value.execute.return_value.data = [{"upload_id": "a" * 32, "metadata": metadata}]
|
||||
media_query = media_table.select.return_value.eq.return_value.limit.return_value
|
||||
media_query.retry.return_value.execute.return_value.data = [{"metadata": metadata}]
|
||||
article_table.select.return_value.eq.return_value.limit.return_value.execute.return_value.data = []
|
||||
article_table.insert.side_effect = lambda record: MagicMock(execute=lambda: MagicMock(data=[record]))
|
||||
with patch.object(main, "get_supabase", return_value=client):
|
||||
main.save_media_metadata("a" * 32, metadata)
|
||||
media_table.upsert.assert_called_once_with({"upload_id": "a" * 32, "metadata": metadata})
|
||||
media_table.upsert.assert_called_once_with([{"upload_id": "a" * 32, "metadata": metadata}], on_conflict="upload_id")
|
||||
self.assertEqual(main.uploaded_media("a" * 32), metadata)
|
||||
self.assertFalse(main.UPLOAD_DIR.exists())
|
||||
result = self.client.post("/api/posts", headers=self.headers, json={**self.article_payload(), "banner": metadata})
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
import unittest
|
||||
from unittest.mock import MagicMock, patch
|
||||
from httpx import ReadTimeout
|
||||
from supabase import PostgrestAPIError
|
||||
from fastapi import HTTPException
|
||||
|
||||
from backend import main, media_cleanup, media_registry
|
||||
|
||||
|
||||
class MediaRegistryTests(unittest.TestCase):
|
||||
def setUp(self):
|
||||
self.client = MagicMock()
|
||||
self.table = self.client.table.return_value
|
||||
self.write = self.table.upsert.return_value.retry.return_value.execute
|
||||
self.read = self.table.select.return_value.in_.return_value.retry.return_value.execute
|
||||
self.records = [{"upload_id": char * 32, "metadata": {
|
||||
"url": "/api/uploads/" + char * 32, "storage": "dropbox", "dropbox_file_id": "id:file_" + char,
|
||||
"name": "photo.png", "media_type": "image/png", "size": 100,
|
||||
}} for char in ("a", "b")]
|
||||
self.read.return_value.data = []
|
||||
sleep = patch.object(media_registry.time, "sleep")
|
||||
sleep.start(); self.addCleanup(sleep.stop)
|
||||
|
||||
def test_queue_retries_gateway_timeout_with_identical_records(self):
|
||||
self.write.side_effect = [PostgrestAPIError({"code": "504", "message": "Gateway Timeout"}), MagicMock(data=self.records)]
|
||||
by_id = {row["upload_id"]: row["metadata"] for row in self.records}
|
||||
with patch.object(main, "get_supabase", return_value=self.client), patch.object(main, "uploaded_media", side_effect=lambda uid: by_id[uid]):
|
||||
media_cleanup.enqueue(list(by_id.values()))
|
||||
self.assertEqual(self.write.call_count, 2)
|
||||
for call in self.table.upsert.call_args_list:
|
||||
self.assertEqual(call.args[0], self.records)
|
||||
self.assertEqual(call.kwargs, {"on_conflict": "upload_id"})
|
||||
|
||||
def test_committed_timeout_is_confirmed_without_another_write(self):
|
||||
self.write.side_effect = ReadTimeout("response lost after commit")
|
||||
self.read.return_value.data = list(reversed(self.records))
|
||||
media_registry.upsert(self.client, "article_media_deletions", self.records)
|
||||
self.assertEqual(self.write.call_count, 1)
|
||||
|
||||
def test_empty_response_checks_every_record_before_success(self):
|
||||
self.write.side_effect = [MagicMock(data=[]), MagicMock(data=self.records)]
|
||||
self.read.return_value.data = self.records[:1]
|
||||
media_registry.upsert(self.client, "article_media", self.records)
|
||||
self.assertEqual(self.write.call_count, 2)
|
||||
|
||||
def test_empty_response_with_all_records_saved_is_success(self):
|
||||
self.write.return_value.data = []
|
||||
self.read.return_value.data = self.records
|
||||
media_registry.upsert(self.client, "article_media", self.records)
|
||||
self.assertEqual(self.write.call_count, 1)
|
||||
|
||||
def test_exhausted_retries_still_block_article_mutation(self):
|
||||
self.write.side_effect = ReadTimeout("offline")
|
||||
by_id = {row["upload_id"]: row["metadata"] for row in self.records}
|
||||
with patch.object(main, "get_supabase", return_value=self.client), patch.object(main, "uploaded_media", side_effect=lambda uid: by_id[uid]), patch.object(main, "post_by_slug", return_value={"banner": self.records[0]["metadata"], "attachments": [self.records[1]["metadata"]]}):
|
||||
with self.assertRaises(HTTPException) as error:
|
||||
main.delete_post("existing-article")
|
||||
self.assertEqual(error.exception.status_code, 502)
|
||||
self.assertEqual(self.write.call_count, 3)
|
||||
self.table.delete.assert_not_called()
|
||||
|
||||
def test_configuration_errors_are_not_retried(self):
|
||||
self.write.side_effect = PostgrestAPIError({"code": "42501", "message": "permission denied"})
|
||||
with self.assertRaises(PostgrestAPIError):
|
||||
media_registry.upsert(self.client, "article_media", self.records)
|
||||
self.assertEqual(self.write.call_count, 1)
|
||||
self.read.assert_not_called()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user