From 32874ca6787fa3bfa46ed2cc4c6a5b8d995bdc3c Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Fri, 25 Sep 2026 10:09:47 -0400 Subject: [PATCH 1/7] fix: the release page drops internal rule numbers and the miscounted "three" images Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR --- scripts/release_notes.py | 16 ++++++++-------- tests/test_release_notes.py | 10 ++++++++++ 2 files changed, 18 insertions(+), 8 deletions(-) diff --git a/scripts/release_notes.py b/scripts/release_notes.py index 9964a04..24f97fb 100644 --- a/scripts/release_notes.py +++ b/scripts/release_notes.py @@ -210,8 +210,8 @@ def render( ) parts.append( - f"Built from `{short}`. The rollback unit is the immutable `:c-` tag " - f"(rule 145) — these three move together:\n\n```\n" + f"Built from `{short}`. To roll back to this release, pull these " + f"immutable `:c-` tags — the images move together:\n\n```\n" + "\n".join(f"{image}:c-{short}" for image in IMAGES) + "\n```" ) @@ -221,10 +221,10 @@ def render( # truncated to MAX_COMMITS, which is 200 lines of internal build-out # presented to someone who has never seen this project. parts.append( - "---\n\n_First release under rule 148's `vYYYY.MM.DD.HHMM` shape, so " - "there is no predecessor to diff against and no changelog to derive. " - "The description above is README.md's, quoted at publish time. Later " - "releases carry the commits since the previous one._" + "---\n\n_The first release, so there is no earlier one to diff " + "against and no changelog to derive. The description above is " + "README.md's, quoted at publish time. Later releases carry the " + "commits since the previous one._" ) return "\n\n".join(parts) @@ -257,8 +257,8 @@ def cross_checks(tag: str, sha: str) -> list[str]: if not RULE_148.match(tag): notes.append( - f"`{tag}` is not rule 148's `vYYYY.MM.DD.HHMM` shape. Published " - f"anyway — the old `v26.*` tags predate the rule." + f"`{tag}` is not the `vYYYY.MM.DD.HHMM` release-tag shape. " + f"Published anyway — the old `v26.*` tags predate it." ) else: derived = artifact_version("web") diff --git a/tests/test_release_notes.py b/tests/test_release_notes.py index 8bc1862..29c1c73 100644 --- a/tests/test_release_notes.py +++ b/tests/test_release_notes.py @@ -14,6 +14,7 @@ re-implementation of it. """ from __future__ import annotations +import re import subprocess from pathlib import Path @@ -162,6 +163,15 @@ def test_the_first_release_describes_the_product_instead_of_diffing(shaped_histo assert not [ln for ln in body.split("\n") if ln.startswith("- work landing")] +@pytest.mark.parametrize("tag", ["v2026.08.28.2208", "v2026.08.29.1000"]) +def test_the_release_page_cites_no_internal_rule_numbers(shaped_history, tag): + """The release page is read by strangers. "rule 145" names a record in the + operator's own notes, which a reader cannot open — say what the rule means + instead. Covers the first-release overview and the changelog body.""" + body = body_of(notes(tag, cwd=shaped_history)) + assert not re.search(r"\brule\s+\d+", body, re.IGNORECASE), body + + def test_the_overview_is_readmes_words_not_a_second_copy(shaped_history): """Two hand-maintained descriptions of one product drift and nothing catches it. The release page quotes README.md so there is one source.""" -- 2.54.0 From 0fad744bfb98c7bf40a2f84820f3e3a82fd59a40 Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Fri, 25 Sep 2026 10:29:05 -0400 Subject: [PATCH 2/7] =?UTF-8?q?fix:=20a=20restart=20no=20longer=20strands?= =?UTF-8?q?=20downloads=20=E2=80=94=20the=20download=20lane=20closes=20its?= =?UTF-8?q?=20predecessor's=20runs=20as=20interrupted=20and=20frees=20the?= =?UTF-8?q?=20platform=20locks=20at=20boot=20(#4433)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A walk that outlives the 90s stop grace is SIGKILLed: its event never finalizes and its platform lock is held for the 27-min TTL. The 30-min sweep then errored every stranded event and bumped consecutive_failures, backing sources off (and blocking backfills) as if the platform failed. On worker_ready, the process consuming 'download' ends pre-boot pending/running events as skipped with error_type 'interrupted', leaves the source untouched so the next tick resumes it, and releases the serialized platforms' locks. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR --- backend/app/celery_signals.py | 48 +++++++++++++++++- backend/app/services/platform_lock.py | 19 +++++++ backend/app/tasks/maintenance.py | 45 ++++++++++++++++- tests/test_download_lane_boot.py | 71 +++++++++++++++++++++++++++ tests/test_maintenance.py | 51 +++++++++++++++++++ 5 files changed, 232 insertions(+), 2 deletions(-) create mode 100644 tests/test_download_lane_boot.py diff --git a/backend/app/celery_signals.py b/backend/app/celery_signals.py index dec94ea..0afaf42 100644 --- a/backend/app/celery_signals.py +++ b/backend/app/celery_signals.py @@ -23,7 +23,13 @@ import logging from datetime import UTC, datetime from celery.exceptions import SoftTimeLimitExceeded -from celery.signals import task_failure, task_postrun, task_prerun, task_retry +from celery.signals import ( + task_failure, + task_postrun, + task_prerun, + task_retry, + worker_ready, +) from .models import TaskRun from .tasks._sync_engine import sync_session_factory @@ -213,3 +219,43 @@ def _on_retry(sender=None, request=None, reason=None, einfo=None, **_): error_message=str(reason) if reason else None, retry_count=getattr(request, "retries", 0), ) + + +def _consumed_queues(consumer) -> set[str]: + """The queue names this worker process consumes (its `-Q`), or empty when + they can't be read — which makes the boot hook below a no-op, not a guess.""" + try: + return {q.name for q in consumer.task_consumer.queues} + except Exception: # noqa: BLE001 — shape varies across celery versions + return set() + + +@worker_ready.connect +def _on_worker_ready(sender=None, **_): + """The download lane clears what its previous process left behind (#4433). + + A restart SIGKILLs any walk that outlives the stop grace, so its event is + never finalized and its platform lock is never released. Both used to wait + out timers — a 30-min stall sweep that then blamed the source, and a 27-min + lock TTL that stalled every other source on the platform. Only the process + consuming `download` does this; the other lanes booting beside it must not. + Best-effort: a failure here is logged and the worker still starts. + """ + if "download" not in _consumed_queues(sender): + return + booted_at = datetime.now(UTC) + try: + from .services.platform_lock import release_all_platform_locks + from .tasks.maintenance import interrupt_orphaned_download_events + + with sync_session_factory()() as session: + closed = interrupt_orphaned_download_events(session, booted_at=booted_at) + session.commit() + released = release_all_platform_locks() + if closed or released: + log.info( + "download lane boot: closed %d orphaned download event(s) as " + "interrupted, released %d platform lock(s)", closed, released, + ) + except Exception: # noqa: BLE001 — never block the worker from starting + log.exception("download lane boot recovery failed") diff --git a/backend/app/services/platform_lock.py b/backend/app/services/platform_lock.py index d910016..89f36ca 100644 --- a/backend/app/services/platform_lock.py +++ b/backend/app/services/platform_lock.py @@ -60,3 +60,22 @@ def platform_lock(platform: str, *, ttl_seconds: int): except redis.RedisError as exc: # pragma: no cover - broker outage log.warning("platform_lock unavailable for %s: %s", platform, exc) return None + + +def release_all_platform_locks() -> int: + """Drop every serialized platform's lock. Returns how many were held. + + Only for the download lane's boot (#4433). A worker that restarts mid-walk + is SIGKILLed past its stop grace, so its `finally` never releases the lock, + and the TTL keeps every other source on that platform bouncing for up to + 27 minutes after the new worker is ready. At boot no walk of ours can be + running, so a held lock names a dead one. Assumes one download consumer — + the only shape FC deploys; a second replica booting would free a live + walk's lock (not corrupt it: the walk runs on, a second walk may overlap it). + """ + try: + client = _redis() + return int(client.delete(*(f"{_LOCK_PREFIX}{p}" for p in SERIALIZED_PLATFORMS))) + except redis.RedisError as exc: # pragma: no cover - broker outage + log.warning("could not release platform locks at boot: %s", exc) + return 0 diff --git a/backend/app/tasks/maintenance.py b/backend/app/tasks/maintenance.py index 31c0fda..87c7cb8 100644 --- a/backend/app/tasks/maintenance.py +++ b/backend/app/tasks/maintenance.py @@ -7,7 +7,8 @@ from datetime import UTC, datetime, timedelta from pathlib import Path from PIL import Image -from sqlalchemy import Integer, and_, cast, delete, func, or_, select, update +from sqlalchemy import Integer, and_, cast, delete, func, literal, or_, select, update +from sqlalchemy.dialects.postgresql import JSONB from ..celery_app import celery from ..models import ( @@ -731,6 +732,48 @@ def recover_stalled_download_events() -> int: return events_recovered +DOWNLOAD_INTERRUPTED_MESSAGE = ( + "interrupted by a worker restart — the next check picks it up where it left off" +) + + +def interrupt_orphaned_download_events(session, *, booted_at: datetime) -> int: + """Close the download events a restart orphaned, without blaming the source. + + Called when the download lane comes up (#4433). Anything still + pending/running from before this boot belongs to the previous process: + a walk that outlived the 90s stop grace was SIGKILLed, and a queued or + serialize-deferred task is held unacked until Redis redelivers it about an + hour later. Left alone, the 30-min stall sweep would error each one and + bump `consecutive_failures`, backing the source off as if the platform had + failed. + + Instead they end as `skipped` (terminal, not a failure) and the source is + not touched: `last_checked_at` keeps its old value, so the next tick finds + it due and the walk resumes from its checkpoint. A redelivered message + that arrives later finds no pending event and opens a fresh one. + + An event promoted to running after the boot has `started_at` reset to its + real start (download_service), so it is never caught here. Does NOT commit. + """ + now = datetime.now(UTC) + result = session.execute( + update(DownloadEvent) + .where(DownloadEvent.status.in_(["pending", "running"])) + .where(DownloadEvent.started_at < booted_at) + .values( + status="skipped", + finished_at=now, + error=DOWNLOAD_INTERRUPTED_MESSAGE, + metadata_=DownloadEvent.metadata_.op("||")( + literal({"error_type": "interrupted"}, JSONB) + ), + ) + .returning(DownloadEvent.id) + ) + return len(result.all()) + + @celery.task(name="backend.app.tasks.maintenance.recover_stalled_backup_runs") def recover_stalled_backup_runs() -> int: """Flip BackupRun rows stuck in running/restoring past the hard limit diff --git a/tests/test_download_lane_boot.py b/tests/test_download_lane_boot.py new file mode 100644 index 0000000..e98dd4a --- /dev/null +++ b/tests/test_download_lane_boot.py @@ -0,0 +1,71 @@ +"""#4433: only the process consuming `download` clears a restart's leftovers. + +The consolidated container boots four lanes side by side. If the scheduler or +ml lane ran this too, a lane restarting on its own (supervisord restarts a +crashed program) would close the events of walks that are still running. +""" +from __future__ import annotations + +from types import SimpleNamespace + +import pytest + +from backend.app import celery_signals + + +def _consumer(*queues: str): + return SimpleNamespace( + task_consumer=SimpleNamespace(queues=[SimpleNamespace(name=q) for q in queues]) + ) + + +@pytest.fixture +def calls(monkeypatch): + seen: list[str] = [] + + class _Session: + def __enter__(self): + return self + + def __exit__(self, *exc): + return False + + def commit(self): + seen.append("commit") + + monkeypatch.setattr(celery_signals, "sync_session_factory", lambda: _Session) + monkeypatch.setattr( + "backend.app.tasks.maintenance.interrupt_orphaned_download_events", + lambda session, *, booted_at: seen.append("interrupt") or 0, + ) + monkeypatch.setattr( + "backend.app.services.platform_lock.release_all_platform_locks", + lambda: seen.append("release") or 0, + ) + return seen + + +def test_the_download_lane_clears_orphans_and_locks(calls): + celery_signals._on_worker_ready(sender=_consumer("default", "download", "import")) + assert calls == ["interrupt", "commit", "release"] + + +@pytest.mark.parametrize("queues", [("maintenance", "scan"), ("ml",), ("maintenance_long",)]) +def test_other_lanes_leave_downloads_alone(calls, queues): + celery_signals._on_worker_ready(sender=_consumer(*queues)) + assert calls == [] + + +def test_unreadable_queues_do_nothing_rather_than_guess(calls): + celery_signals._on_worker_ready(sender=SimpleNamespace()) + assert calls == [] + + +def test_a_failure_does_not_stop_the_worker_starting(monkeypatch, calls): + def boom(session, *, booted_at): + raise RuntimeError("db down") + + monkeypatch.setattr( + "backend.app.tasks.maintenance.interrupt_orphaned_download_events", boom + ) + celery_signals._on_worker_ready(sender=_consumer("download")) # must not raise diff --git a/tests/test_maintenance.py b/tests/test_maintenance.py index 7bf6ab5..4a1712b 100644 --- a/tests/test_maintenance.py +++ b/tests/test_maintenance.py @@ -718,6 +718,57 @@ def test_recover_stalled_download_skips_fresh(db_sync): assert failures == 0 +def test_download_lane_boot_closes_pre_boot_events_without_blaming_the_source(db_sync): + """#4433: a restart SIGKILLs a walk past its stop grace and strands queued + ones. At the download lane's boot, everything pending/running from before + it ends as `skipped`/interrupted — and the source is left as it was, so it + is due on the next tick instead of backed off as a failure.""" + from sqlalchemy import select + + from backend.app.models import DownloadEvent, Source + from backend.app.tasks.maintenance import ( + DOWNLOAD_INTERRUPTED_MESSAGE, + interrupt_orphaned_download_events, + ) + + sid = _make_source(db_sync, slug="rebooted") + booted_at = datetime.now(UTC) + before = booted_at - timedelta(minutes=5) + walking = DownloadEvent( + source_id=sid, status="running", started_at=before, + metadata_={"live": {"downloaded": 3}}, + ) + queued = DownloadEvent(source_id=sid, status="pending", started_at=before) + finished = DownloadEvent(source_id=sid, status="ok", started_at=before) + # Promoted to running after the boot: download_service resets started_at. + fresh = DownloadEvent( + source_id=sid, status="running", started_at=booted_at + timedelta(seconds=5), + ) + db_sync.add_all([walking, queued, finished, fresh]) + db_sync.commit() + + closed = interrupt_orphaned_download_events(db_sync, booted_at=booted_at) + db_sync.commit() + + assert closed == 2 + db_sync.expire_all() + for ev in (walking, queued): + assert ev.status == "skipped" + assert ev.error == DOWNLOAD_INTERRUPTED_MESSAGE + assert ev.finished_at is not None + assert ev.metadata_["error_type"] == "interrupted" + assert walking.metadata_["live"] == {"downloaded": 3} + assert finished.status == "ok" + assert fresh.status == "running" + src = db_sync.execute( + select(Source.consecutive_failures, Source.last_error, Source.last_checked_at) + .where(Source.id == sid) + ).one() + assert src.consecutive_failures == 0 + assert src.last_error is None + assert src.last_checked_at is None + + def test_recover_stalled_download_flips_stale_pending(db_sync): """A 2-hour-old pending event flips to error AND the source is bumped (consecutive_failures, last_error, last_checked_at) so the next scan -- 2.54.0 From e704c70f32d0268f437e9ddad05b6e01a2f2f60e Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Fri, 25 Sep 2026 10:32:41 -0400 Subject: [PATCH 3/7] fix: the download boot hook imports a model-only module, keeping the membership roster off the fetch path (#4433) Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR --- backend/app/celery_signals.py | 2 +- backend/app/services/download_recovery.py | 57 +++++++++++++++++++++++ backend/app/tasks/maintenance.py | 45 +----------------- tests/test_download_lane_boot.py | 4 +- tests/test_maintenance.py | 2 +- 5 files changed, 62 insertions(+), 48 deletions(-) create mode 100644 backend/app/services/download_recovery.py diff --git a/backend/app/celery_signals.py b/backend/app/celery_signals.py index 0afaf42..1b9e7b2 100644 --- a/backend/app/celery_signals.py +++ b/backend/app/celery_signals.py @@ -246,7 +246,7 @@ def _on_worker_ready(sender=None, **_): booted_at = datetime.now(UTC) try: from .services.platform_lock import release_all_platform_locks - from .tasks.maintenance import interrupt_orphaned_download_events + from .services.download_recovery import interrupt_orphaned_download_events with sync_session_factory()() as session: closed = interrupt_orphaned_download_events(session, booted_at=booted_at) diff --git a/backend/app/services/download_recovery.py b/backend/app/services/download_recovery.py new file mode 100644 index 0000000..e62fc75 --- /dev/null +++ b/backend/app/services/download_recovery.py @@ -0,0 +1,57 @@ +"""Closing the download runs a worker restart orphaned (#4433). + +Its own module, importing nothing but the model, because its caller is the +worker boot hook in `celery_signals` — which every download task imports via +`celery_app`. Living in `tasks.maintenance` put the whole maintenance import +graph, the membership roster included, on the fetch path, which +`test_gated_reason` forbids. +""" +from __future__ import annotations + +from datetime import UTC, datetime + +from sqlalchemy import literal, update +from sqlalchemy.dialects.postgresql import JSONB + +from ..models import DownloadEvent + +DOWNLOAD_INTERRUPTED_MESSAGE = ( + "interrupted by a worker restart — the next check picks it up where it left off" +) + + +def interrupt_orphaned_download_events(session, *, booted_at: datetime) -> int: + """Close the download events a restart orphaned, without blaming the source. + + Called when the download lane comes up (#4433). Anything still + pending/running from before this boot belongs to the previous process: + a walk that outlived the 90s stop grace was SIGKILLed, and a queued or + serialize-deferred task is held unacked until Redis redelivers it about an + hour later. Left alone, the 30-min stall sweep would error each one and + bump `consecutive_failures`, backing the source off as if the platform had + failed. + + Instead they end as `skipped` (terminal, not a failure) and the source is + not touched: `last_checked_at` keeps its old value, so the next tick finds + it due and the walk resumes from its checkpoint. A redelivered message + that arrives later finds no pending event and opens a fresh one. + + An event promoted to running after the boot has `started_at` reset to its + real start (download_service), so it is never caught here. Does NOT commit. + """ + now = datetime.now(UTC) + result = session.execute( + update(DownloadEvent) + .where(DownloadEvent.status.in_(["pending", "running"])) + .where(DownloadEvent.started_at < booted_at) + .values( + status="skipped", + finished_at=now, + error=DOWNLOAD_INTERRUPTED_MESSAGE, + metadata_=DownloadEvent.metadata_.op("||")( + literal({"error_type": "interrupted"}, JSONB) + ), + ) + .returning(DownloadEvent.id) + ) + return len(result.all()) diff --git a/backend/app/tasks/maintenance.py b/backend/app/tasks/maintenance.py index 87c7cb8..31c0fda 100644 --- a/backend/app/tasks/maintenance.py +++ b/backend/app/tasks/maintenance.py @@ -7,8 +7,7 @@ from datetime import UTC, datetime, timedelta from pathlib import Path from PIL import Image -from sqlalchemy import Integer, and_, cast, delete, func, literal, or_, select, update -from sqlalchemy.dialects.postgresql import JSONB +from sqlalchemy import Integer, and_, cast, delete, func, or_, select, update from ..celery_app import celery from ..models import ( @@ -732,48 +731,6 @@ def recover_stalled_download_events() -> int: return events_recovered -DOWNLOAD_INTERRUPTED_MESSAGE = ( - "interrupted by a worker restart — the next check picks it up where it left off" -) - - -def interrupt_orphaned_download_events(session, *, booted_at: datetime) -> int: - """Close the download events a restart orphaned, without blaming the source. - - Called when the download lane comes up (#4433). Anything still - pending/running from before this boot belongs to the previous process: - a walk that outlived the 90s stop grace was SIGKILLed, and a queued or - serialize-deferred task is held unacked until Redis redelivers it about an - hour later. Left alone, the 30-min stall sweep would error each one and - bump `consecutive_failures`, backing the source off as if the platform had - failed. - - Instead they end as `skipped` (terminal, not a failure) and the source is - not touched: `last_checked_at` keeps its old value, so the next tick finds - it due and the walk resumes from its checkpoint. A redelivered message - that arrives later finds no pending event and opens a fresh one. - - An event promoted to running after the boot has `started_at` reset to its - real start (download_service), so it is never caught here. Does NOT commit. - """ - now = datetime.now(UTC) - result = session.execute( - update(DownloadEvent) - .where(DownloadEvent.status.in_(["pending", "running"])) - .where(DownloadEvent.started_at < booted_at) - .values( - status="skipped", - finished_at=now, - error=DOWNLOAD_INTERRUPTED_MESSAGE, - metadata_=DownloadEvent.metadata_.op("||")( - literal({"error_type": "interrupted"}, JSONB) - ), - ) - .returning(DownloadEvent.id) - ) - return len(result.all()) - - @celery.task(name="backend.app.tasks.maintenance.recover_stalled_backup_runs") def recover_stalled_backup_runs() -> int: """Flip BackupRun rows stuck in running/restoring past the hard limit diff --git a/tests/test_download_lane_boot.py b/tests/test_download_lane_boot.py index e98dd4a..ff52b65 100644 --- a/tests/test_download_lane_boot.py +++ b/tests/test_download_lane_boot.py @@ -35,7 +35,7 @@ def calls(monkeypatch): monkeypatch.setattr(celery_signals, "sync_session_factory", lambda: _Session) monkeypatch.setattr( - "backend.app.tasks.maintenance.interrupt_orphaned_download_events", + "backend.app.services.download_recovery.interrupt_orphaned_download_events", lambda session, *, booted_at: seen.append("interrupt") or 0, ) monkeypatch.setattr( @@ -66,6 +66,6 @@ def test_a_failure_does_not_stop_the_worker_starting(monkeypatch, calls): raise RuntimeError("db down") monkeypatch.setattr( - "backend.app.tasks.maintenance.interrupt_orphaned_download_events", boom + "backend.app.services.download_recovery.interrupt_orphaned_download_events", boom ) celery_signals._on_worker_ready(sender=_consumer("download")) # must not raise diff --git a/tests/test_maintenance.py b/tests/test_maintenance.py index 4a1712b..c9df368 100644 --- a/tests/test_maintenance.py +++ b/tests/test_maintenance.py @@ -726,7 +726,7 @@ def test_download_lane_boot_closes_pre_boot_events_without_blaming_the_source(db from sqlalchemy import select from backend.app.models import DownloadEvent, Source - from backend.app.tasks.maintenance import ( + from backend.app.services.download_recovery import ( DOWNLOAD_INTERRUPTED_MESSAGE, interrupt_orphaned_download_events, ) -- 2.54.0 From e2eeb631152994c7bcdd6ace4fc396e83442d4f0 Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Fri, 25 Sep 2026 10:36:32 -0400 Subject: [PATCH 4/7] style: sort the boot hook's imports (#4433) Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR --- backend/app/celery_signals.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/backend/app/celery_signals.py b/backend/app/celery_signals.py index 1b9e7b2..f30c492 100644 --- a/backend/app/celery_signals.py +++ b/backend/app/celery_signals.py @@ -245,8 +245,8 @@ def _on_worker_ready(sender=None, **_): return booted_at = datetime.now(UTC) try: - from .services.platform_lock import release_all_platform_locks from .services.download_recovery import interrupt_orphaned_download_events + from .services.platform_lock import release_all_platform_locks with sync_session_factory()() as session: closed = interrupt_orphaned_download_events(session, booted_at=booted_at) -- 2.54.0 From ca802bd885a96a59f9f93ec0efeb759c0db46c07 Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Fri, 25 Sep 2026 10:56:08 -0400 Subject: [PATCH 5/7] fix: a non-image file lands on the source it was downloaded for, and the undated shells it left are folded into their real posts (#4435) The attachment path looked the post's source up by (artist, platform), taking the artist's lowest-id source. A Discord artist has one source per channel, so every archive or pdf from a later channel became an undated second post on the first channel's source. _post_for_sidecar now takes the downloading source, as upsert_post_record and _apply_sidecar already did. Migration 0114 re-points each shell's attachments to its dated twin and deletes the shell. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR --- .../0114_fold_misfiled_attachment_posts.py | 79 +++++++++++++ backend/app/services/importer.py | 32 ++++-- ...old_misfiled_attachment_posts_migration.py | 106 ++++++++++++++++++ tests/test_importer_attach_in_place.py | 45 ++++++++ 4 files changed, 253 insertions(+), 9 deletions(-) create mode 100644 alembic/versions/0114_fold_misfiled_attachment_posts.py create mode 100644 tests/test_fold_misfiled_attachment_posts_migration.py diff --git a/alembic/versions/0114_fold_misfiled_attachment_posts.py b/alembic/versions/0114_fold_misfiled_attachment_posts.py new file mode 100644 index 0000000..4356898 --- /dev/null +++ b/alembic/versions/0114_fold_misfiled_attachment_posts.py @@ -0,0 +1,79 @@ +"""Fold the undated shell posts a misfiled attachment created into the real post. + +#4435. A non-image file (an archive, a pdf) downloaded for one of an artist's +Discord channels was filed under the artist's FIRST Discord source: the +attachment path looked the source up by (artist, platform), which takes the +lowest id. That created an undated, url-less post there holding only the +attachment, and the message's real post record then created the dated post +under the right source. The importer now uses the source it was downloading +for; this repairs the pairs it left. + +A shell is folded only when all of this holds: + +* it has no date and no url, and nothing synthesized it; +* another post of the same artist, platform and external id HAS a date; +* no image is linked to the shell (it held an attachment and nothing else). + +Its attachments move to the dated post, dropping any the dated post already +has (same sha256, which the per-post unique forbids twice), and the shell is +deleted. A shell with no dated twin is left alone: which channel it belongs to +is not recorded anywhere but the file name. The downgrade does nothing. + +Revision ID: 0114 +Revises: 0113 +Create Date: 2026-09-25 + +""" +import sqlalchemy as sa +from alembic import op + +revision = "0114" +down_revision = "0113" +branch_labels = None +depends_on = None + +_PAIRS = """ + SELECT DISTINCT ON (shell.id) shell.id AS shell_id, real.id AS real_id + FROM post shell + JOIN source ss ON ss.id = shell.source_id + JOIN post real + ON real.artist_id = shell.artist_id + AND real.external_post_id = shell.external_post_id + AND real.id <> shell.id + AND real.post_date IS NOT NULL + JOIN source rs ON rs.id = real.source_id AND rs.platform = ss.platform + WHERE shell.post_date IS NULL + AND shell.post_url IS NULL + AND shell.synthesized_by IS NULL + AND NOT EXISTS (SELECT 1 FROM image_provenance ip WHERE ip.post_id = shell.id) + ORDER BY shell.id, real.id +""" + + +def fold_misfiled_attachment_posts(conn) -> int: + """The data step, on a plain connection, so a test can run it directly. + Returns how many shells were folded.""" + pairs = conn.execute(sa.text(_PAIRS)).all() + for shell_id, real_id in pairs: + conn.execute(sa.text(""" + DELETE FROM post_attachment pa + WHERE pa.post_id = :shell + AND EXISTS ( + SELECT 1 FROM post_attachment keep + WHERE keep.post_id = :real AND keep.sha256 = pa.sha256 + ) + """), {"shell": shell_id, "real": real_id}) + conn.execute( + sa.text("UPDATE post_attachment SET post_id = :real WHERE post_id = :shell"), + {"shell": shell_id, "real": real_id}, + ) + conn.execute(sa.text("DELETE FROM post WHERE id = :shell"), {"shell": shell_id}) + return len(pairs) + + +def upgrade() -> None: + fold_misfiled_attachment_posts(op.get_bind()) + + +def downgrade() -> None: + pass diff --git a/backend/app/services/importer.py b/backend/app/services/importer.py index 8a3b4f4..9e28d07 100644 --- a/backend/app/services/importer.py +++ b/backend/app/services/importer.py @@ -426,11 +426,19 @@ class Importer: return self._upsert_artist(name) if name else None def _post_for_sidecar( - self, source: Path, artist: Artist | None + self, source: Path, artist: Artist | None, + *, source_row: Source | None = None, ) -> Post | None: """If a sidecar sits next to `source`, ensure its Source+Post exist (idempotent) and return the Post — so attachments can link - to the same Post the per-member _apply_sidecar will reuse.""" + to the same Post the per-member _apply_sidecar will reuse. + + `source_row` is the subscription being downloaded, when there is one, + and wins over the (artist, platform) lookup — as it does in + `upsert_post_record`. The lookup takes the artist's FIRST source on the + platform, which is right only while an artist has one: a Discord artist + has one per channel, and every non-image file from a later channel was + filed under the first as an undated second post (#4435).""" sc = find_sidecar(source) if sc is None or artist is None: return None @@ -442,10 +450,13 @@ class Importer: log.warning("sidecar parse failed for %s: %s", sc, exc) return None sd = parse_sidecar(data) - platform = sd.platform or "unknown" - src = self._lookup_source_for_sidecar( - artist_id=artist.id, platform=platform, - ) + if source_row is not None: + src = source_row + else: + platform = sd.platform or "unknown" + src = self._lookup_source_for_sidecar( + artist_id=artist.id, platform=platform, + ) epid = sd.external_post_id or sc.stem return self._find_or_create_post( source_id=src.id if src else None, @@ -543,7 +554,7 @@ class Importer: # nothing silently vanishes, matching extract_archive's # fail-soft contract. artist_use = artist if artist is not None else self._resolve_artist(source) - post = self._post_for_sidecar(source, artist_use) + post = self._post_for_sidecar(source, artist_use, source_row=source_row) self._capture_attachment( source, post=post, artist=artist_use, resolved=True, ) @@ -552,7 +563,7 @@ class Importer: return ImportResult(status="attached", error=reason) artist_use = artist if artist is not None else self._resolve_artist(source) - post = self._post_for_sidecar(source, artist_use) + post = self._post_for_sidecar(source, artist_use, source_row=source_row) member_ids: list[int] = [] # Every member image touched (new + superseded + deduped), so the # from_attachment_id stamp below covers files that already existed in the @@ -1218,7 +1229,10 @@ class Importer: path, artist=artist, source_row=source, ) if not is_supported(path): - post = self._post_for_sidecar(path, artist) if artist else None + post = ( + self._post_for_sidecar(path, artist, source_row=source) + if artist else None + ) return self._capture_attachment( path, post=post, artist=artist, resolved=True, ) diff --git a/tests/test_fold_misfiled_attachment_posts_migration.py b/tests/test_fold_misfiled_attachment_posts_migration.py new file mode 100644 index 0000000..ebc8546 --- /dev/null +++ b/tests/test_fold_misfiled_attachment_posts_migration.py @@ -0,0 +1,106 @@ +"""Migration 0114 (#4435): the undated shell post a misfiled attachment made on +the artist's first Discord source is folded into the real, dated post.""" +import importlib.util +from datetime import UTC, datetime +from pathlib import Path + +import pytest +from sqlalchemy import select + +from backend.app.models import ( + Artist, + ImageProvenance, + Post, + PostAttachment, + Source, +) +from tests.factories import make_image as _img + +pytestmark = pytest.mark.integration + +_MIGRATION = ( + Path(__file__).resolve().parents[1] + / "alembic" / "versions" / "0114_fold_misfiled_attachment_posts.py" +) + + +def _fold(): + spec = importlib.util.spec_from_file_location("m0114", _MIGRATION) + mod = importlib.util.module_from_spec(spec) + spec.loader.exec_module(mod) + return mod.fold_misfiled_attachment_posts + + +def _source(db, artist, channel): + s = Source( + artist_id=artist.id, platform="discord", + url=f"https://discord.com/channels/1/{channel}", + ) + db.add(s) + db.flush() + return s + + +def _post(db, artist, source, epid, when=None): + p = Post( + artist_id=artist.id, source_id=source.id, external_post_id=epid, + post_date=when, + post_url=f"https://discord.com/channels/1/x/{epid}" if when else None, + ) + db.add(p) + db.flush() + return p + + +def _attach(db, post, sha, name): + db.add(PostAttachment( + post_id=post.id, artist_id=post.artist_id, sha256=sha, + path=f"/att/{sha}", original_filename=name, ext=".rar", size_bytes=1, + )) + db.flush() + + +def test_shells_fold_into_their_dated_twin_and_nothing_else_moves(db_sync): + sent = datetime(2024, 12, 26, 1, 42, tzinfo=UTC) + artist = Artist(name="Yellow", slug="yellow") + db_sync.add(artist) + db_sync.flush() + first = _source(db_sync, artist, 100) + second = _source(db_sync, artist, 200) + + # The #4435 shape: shell on the first source, real post on the second. + shell = _post(db_sync, artist, first, "555") + real = _post(db_sync, artist, second, "555", sent) + _attach(db_sync, shell, "a" * 64, "pack.rar") + # A shell whose attachment the real post already has: dropped, not doubled. + shell2 = _post(db_sync, artist, first, "556") + real2 = _post(db_sync, artist, second, "556", sent) + _attach(db_sync, shell2, "b" * 64, "same.rar") + _attach(db_sync, real2, "b" * 64, "same.rar") + # Left alone: no dated twin, and an undated post that holds an image. + lonely = _post(db_sync, artist, first, "557") + _attach(db_sync, lonely, "c" * 64, "lonely.rar") + with_image = _post(db_sync, artist, first, "558") + _post(db_sync, artist, second, "558", sent) + img = _img(db_sync, "d" * 64) + db_sync.add(ImageProvenance(image_record_id=img.id, post_id=with_image.id)) + db_sync.flush() + shell_id, shell2_id = shell.id, shell2.id + + folded = _fold()(db_sync.connection()) + db_sync.expire_all() + + assert folded == 2 + assert db_sync.get(Post, shell_id) is None + assert db_sync.get(Post, shell2_id) is None + owners = dict(db_sync.execute( + select(PostAttachment.original_filename, PostAttachment.post_id) + ).all()) + assert owners["pack.rar"] == real.id + assert owners["lonely.rar"] == lonely.id + kept = db_sync.execute( + select(PostAttachment.id).where(PostAttachment.post_id == real2.id) + ).scalars().all() + assert len(kept) == 1 + assert db_sync.get(Post, lonely.id) is not None + assert db_sync.get(Post, with_image.id) is not None diff --git a/tests/test_importer_attach_in_place.py b/tests/test_importer_attach_in_place.py index ab08a06..47ae24d 100644 --- a/tests/test_importer_attach_in_place.py +++ b/tests/test_importer_attach_in_place.py @@ -435,3 +435,48 @@ def test_attach_in_place_non_media_routes_to_attachment(importer, db_sync): ).scalar_one() assert row.ext == ".txt" assert row.artist_id == artist.id + + +def test_a_non_media_file_lands_on_the_source_being_downloaded(importer, db_sync): + """#4435: a Discord artist has one source per channel. A non-media file + downloaded for the SECOND channel was filed under the artist's first + Discord source (the (artist, platform) lookup takes the lowest id), as an + undated shell post; the real post record then made a second post under + the right source. The attachment must follow the source it was fetched for.""" + from backend.app.models import Post, PostAttachment, Source + + images_root = importer.images_root + artist = Artist(name="Yara", slug="yara") + db_sync.add(artist) + db_sync.flush() + first = Source( + artist_id=artist.id, platform="discord", + url="https://discord.com/channels/1/100", + ) + second = Source( + artist_id=artist.id, platform="discord", + url="https://discord.com/channels/1/200", + ) + db_sync.add_all([first, second]) + db_sync.flush() + + rar = images_root / "yara" / "discord" / "rewards" / "20241226_555_01_pack.rar" + rar.parent.mkdir(parents=True, exist_ok=True) + rar.write_bytes(b"not a real archive, so it is kept as an attachment") + rar.with_suffix(rar.suffix + ".json").write_text( + '{"category": "discord", "id": "555", "message_id": "555"}' + ) + + result = importer.attach_in_place(rar, artist=artist, source=second) + assert result.status == "attached" + + owner = db_sync.execute( + select(Post.source_id, Post.external_post_id) + .join(PostAttachment, PostAttachment.post_id == Post.id) + .where(PostAttachment.original_filename == rar.name) + ).one() + assert owner.source_id == second.id + assert owner.external_post_id == "555" + assert db_sync.execute( + select(func.count()).select_from(Post).where(Post.source_id == first.id) + ).scalar_one() == 0 -- 2.54.0 From 621c2b8315f676724454a9a297a2992b847c4e8f Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Fri, 25 Sep 2026 11:36:17 -0400 Subject: [PATCH 6/7] fix: a post record is marked seen only after it is upserted, and an hourly sweep dates the posts killed walks left undated (#4436) ingest_core marked a post's record key seen when it wrote _post.json, but the record reaches the database (and dates the post and its images) only in phase 3. A walk killed before phase 3 left posts the ledger called recorded and the database never dated; ticks early-out long before reaching them again. Keys now join the media in mark_seen_after_import. date_posts_from_records finds each undated native post's record under its artist's folder and upserts it with the post's own source. Hourly on maintenance_long; an empty query once everything is dated. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR --- backend/app/celery_app.py | 7 ++ backend/app/services/ingest_core.py | 11 ++- backend/app/services/post_record_repair.py | 86 ++++++++++++++++++++++ backend/app/tasks/maintenance.py | 26 +++++++ tests/test_celery_routing.py | 1 + tests/test_patreon_ingester.py | 7 +- tests/test_sidecar_import.py | 48 ++++++++++++ 7 files changed, 182 insertions(+), 4 deletions(-) create mode 100644 backend/app/services/post_record_repair.py diff --git a/backend/app/celery_app.py b/backend/app/celery_app.py index 6843589..63313ed 100644 --- a/backend/app/celery_app.py +++ b/backend/app/celery_app.py @@ -69,6 +69,8 @@ def make_celery() -> Celery: # up behind it (2026-09-24: 7 waiting, "all workers busy for 18 # minutes"). An exact name wins over the glob above. "backend.app.tasks.maintenance.backfill_phash": {"queue": "maintenance_long"}, + # Walks a folder tree per artist with undated posts (#4436). + "backend.app.tasks.maintenance.date_posts_from_records": {"queue": "maintenance_long"}, "backend.app.tasks.backup.*": {"queue": "maintenance_long"}, "backend.app.tasks.admin.*": {"queue": "maintenance_long"}, "backend.app.tasks.library_audit.*": {"queue": "maintenance_long"}, @@ -156,6 +158,11 @@ def make_celery() -> Celery: "schedule": 86400.0, # daily — sweep .part/.partial left by a # download/import killed mid-write (graceful-shutdown fallout) }, + "date-posts-from-records-hourly": { + "task": "backend.app.tasks.maintenance.date_posts_from_records", + "schedule": 3600.0, # an empty query once every native post is + # dated; otherwise upserts the _post.json a killed walk left (#4436) + }, "backfill-phash-daily": { "task": "backend.app.tasks.maintenance.backfill_phash", "schedule": 86400.0, # daily — NULL-only, so a no-op once the diff --git a/backend/app/services/ingest_core.py b/backend/app/services/ingest_core.py index ffc461a..2c8a859 100644 --- a/backend/app/services/ingest_core.py +++ b/backend/app/services/ingest_core.py @@ -289,6 +289,11 @@ class Ingester: # Media handed to phase 3 for import. Marked seen by phase 3 once the # import has run (`mark_seen_after_import`), not here — see there. fetched: list[tuple[str, str]] = [] + # Post-record keys written this walk. Marked with the media, after + # phase 3 has upserted the records — marking them at write time left a + # walk killed before phase 3 with posts the ledger calls recorded that + # the database never dated, and no later tick walks back to them (#4436). + recorded: list[tuple[str, str]] = [] downloaded = 0 errors = 0 quarantined = 0 @@ -342,7 +347,9 @@ class Ingester: written_paths=written, post_record_paths=list(post_records), relink_source_paths=list(relink), - mark_seen_after_import=lambda: self._mark_seen(source_id, fetched), + mark_seen_after_import=lambda: self._mark_seen( + source_id, fetched + recorded, + ), stdout="\n".join(log_lines), stderr="", return_code=return_code, @@ -511,7 +518,7 @@ class Ingester: posts_with_body += 1 if rec.path is not None: post_records.append(str(rec.path)) - self._mark_seen(source_id, [(pkey, ppid)]) + recorded.append((pkey, ppid)) # Per-post handling line in the run stdout (the existing # "Raw stdout" panel) — the downloader already read the # post; we only format its outcome here. post_type beside diff --git a/backend/app/services/post_record_repair.py b/backend/app/services/post_record_repair.py new file mode 100644 index 0000000..7043fd4 --- /dev/null +++ b/backend/app/services/post_record_repair.py @@ -0,0 +1,86 @@ +"""Date the posts whose record reached the disk but never the database (#4436). + +A native walk writes each post's record (`_post.json`, or Discord's +`__post.json`) as it goes, and phase 3 upserts those records +after the walk — which is when the post, and through it its images, get their +date. Until #4436 the walk marked a record's post key seen at write time, so a +walk killed before phase 3 (a restart, a stall) left the post undated with the +ledger saying it was done, and no later tick walked back that far to fix it. + +The record files are still on disk. This finds each undated native post's +record under its artist's folder and upserts it with the post's OWN source — +never the (artist, platform) lookup, which picks the artist's first source and +would misfile a Discord channel's post (#4435). Idempotent: a post that is +already dated is never looked at, and upserting a record twice changes nothing. +""" +from __future__ import annotations + +import json +import logging +from collections import defaultdict +from pathlib import Path + +from sqlalchemy import select +from sqlalchemy.orm import Session + +from ..models import Artist, Post, Source +from ..utils.sidecar import parse_sidecar +from .download_backends import NATIVE_INGESTER_PLATFORMS + +log = logging.getLogger(__name__) + + +def _is_record(path: Path) -> bool: + return path.name == "_post.json" or path.name.endswith("_post.json") + + +def _record_id(path: Path) -> str | None: + try: + data = json.loads(path.read_text("utf-8")) + except (OSError, ValueError): + return None + if not isinstance(data, dict): + return None + return parse_sidecar(data).external_post_id + + +def date_posts_from_records(session: Session, importer, images_root: Path) -> dict: + """Upsert the on-disk record of every undated native post. Returns counts.""" + rows = session.execute( + select(Post.external_post_id, Post.source_id, Post.artist_id, Source.platform) + .join(Source, Source.id == Post.source_id) + .where( + Post.post_date.is_(None), + Post.synthesized_by.is_(None), + Source.platform.in_(NATIVE_INGESTER_PLATFORMS), + ) + ).all() + wanted: dict[tuple[int, str], dict[str, int]] = defaultdict(dict) + for epid, source_id, artist_id, platform in rows: + wanted[(artist_id, platform)][epid] = source_id + + summary = {"undated": len(rows), "dated": 0, "no_record": 0} + for (artist_id, platform), by_epid in wanted.items(): + artist = session.get(Artist, artist_id) + root = Path(images_root) / artist.slug / platform if artist else None + if root is None or not root.is_dir(): + summary["no_record"] += len(by_epid) + continue + found = 0 + for record in root.rglob("*post.json"): + if not _is_record(record): + continue + epid = _record_id(record) + source_id = by_epid.pop(epid, None) if epid else None + if source_id is None: + continue + source = session.get(Source, source_id) + if importer.upsert_post_record(record, artist=artist, source=source): + found += 1 + if not by_epid: + break + summary["dated"] += found + summary["no_record"] += len(by_epid) + if summary["undated"]: + log.info("date_posts_from_records: %s", summary) + return summary diff --git a/backend/app/tasks/maintenance.py b/backend/app/tasks/maintenance.py index 31c0fda..51558c0 100644 --- a/backend/app/tasks/maintenance.py +++ b/backend/app/tasks/maintenance.py @@ -731,6 +731,32 @@ def recover_stalled_download_events() -> int: return events_recovered +@celery.task( + name="backend.app.tasks.maintenance.date_posts_from_records", + soft_time_limit=1500, + time_limit=1800, +) +def date_posts_from_records() -> dict: + """Date undated native posts from the records their walk left on disk + (#4436). Hourly and self-limiting: once every post is dated it is one + empty query.""" + from ..services.importer import Importer + from ..services.post_record_repair import date_posts_from_records as _repair + from ..services.thumbnailer import Thumbnailer + + images_root = IMAGES_ROOT + SessionLocal = _sync_session_factory() + with SessionLocal() as session: + importer = Importer( + session=session, + images_root=images_root, + import_root=images_root, + thumbnailer=Thumbnailer(images_root=images_root), + settings=ImportSettings.load_sync(session), + ) + return _repair(session, importer, images_root) + + @celery.task(name="backend.app.tasks.maintenance.recover_stalled_backup_runs") def recover_stalled_backup_runs() -> int: """Flip BackupRun rows stuck in running/restoring past the hard limit diff --git a/tests/test_celery_routing.py b/tests/test_celery_routing.py index de28fd9..f951594 100644 --- a/tests/test_celery_routing.py +++ b/tests/test_celery_routing.py @@ -28,6 +28,7 @@ def test_quick_maintenance_stays_on_maintenance(): ("backend.app.tasks.translation.translate_posts", "maintenance_long"), ("backend.app.tasks.gpu_queue.enqueue_gpu_backfill", "maintenance"), ("backend.app.tasks.maintenance.backfill_phash", "maintenance_long"), + ("backend.app.tasks.maintenance.date_posts_from_records", "maintenance_long"), ("backend.app.tasks.admin.normalize_tags_task", "maintenance_long"), ("backend.app.tasks.backup.backup_db_task", "maintenance_long"), ("backend.app.tasks.maintenance.vacuum_analyze", "maintenance"), diff --git a/tests/test_patreon_ingester.py b/tests/test_patreon_ingester.py index 8d158bc..b986e53 100644 --- a/tests/test_patreon_ingester.py +++ b/tests/test_patreon_ingester.py @@ -943,8 +943,10 @@ async def test_a_run_that_dies_before_import_leaves_its_media_unmarked( _FakeDownloader(tmp_path)) ing.run(source_id=source_id, campaign_id="c1", artist_slug="ingest", url="https://patreon.com/ingest", mode="tick") - # Phase 3 never ran, so `mark_seen_after_import` never did: only the post key. - assert _count_ledger(sync_engine, source_id) == 1 + # Phase 3 never ran, so `mark_seen_after_import` never did: neither the + # media nor the post record is marked (#4436 — the record, marked at write + # time, left the post undated for good once the run died before phase 3). + assert _count_ledger(sync_engine, source_id) == 0 # The next walk finds the file on disk with no record, and imports it. ing2 = _ingester(sync_engine, tmp_path, _FakeClient([(None, [("p1", [m1])])]), @@ -954,6 +956,7 @@ async def test_a_run_that_dies_before_import_leaves_its_media_unmarked( assert result.written_paths == [str(tmp_path / "p1_1.jpg")] assert result.files_downloaded == 0 # not fetched again assert "on disk but never imported: p1_1.jpg" in result.stdout + assert len(result.post_record_paths) == 1 # the record is written again result.mark_seen_after_import() assert _count_ledger(sync_engine, source_id) == 2 diff --git a/tests/test_sidecar_import.py b/tests/test_sidecar_import.py index 6149624..a3aec25 100644 --- a/tests/test_sidecar_import.py +++ b/tests/test_sidecar_import.py @@ -415,3 +415,51 @@ def test_post_record_redates_images_linked_before_it(importer, import_layout): assert post.post_date != download_time assert rec.effective_date == post.post_date assert rec.earliest_post_date == post.post_date + + +def test_an_undated_post_is_dated_from_the_record_its_walk_left_on_disk( + importer, import_layout, +): + """#4436: a walk killed before phase 3 wrote the message's record to disk + but never upserted it, and marked it seen, so no later tick fixed it. The + repair finds the record under the artist's folder and dates the post — and + its images — under the post's own source.""" + from backend.app.services.post_record_repair import date_posts_from_records + + import_root, images_root = import_layout + artist = Artist(name="Alice", slug="alice") + importer.session.add(artist) + importer.session.flush() + importer.session.add(Source( + artist_id=artist.id, platform="discord", + url="https://discord.com/channels/1/100", + )) + importer.session.flush() + m = import_root / "Alice" / "20240301_123_01_art.jpg" + _split(m, "v") + _sidecar(m, {"category": "discord", "message_id": "123"}) + r = importer.import_one(m) + assert r.status == "imported" + post = importer.session.execute(select(Post)).scalar_one() + assert post.post_date is None and post.source_id is not None + + channel = images_root / "alice" / "discord" / "rewards" + channel.mkdir(parents=True) + (channel / "20240301_123_post.json").write_text(json.dumps({ + "category": "discord", "message_id": "123", "message": "", + "date": "2024-03-01T18:30:00.000000+00:00", + })) + (channel / "20240301_999_post.json").write_text("not json") # skipped, not fatal + + summary = date_posts_from_records(importer.session, importer, images_root) + + assert summary == {"undated": 1, "dated": 1, "no_record": 0} + importer.session.expire_all() + post = importer.session.execute(select(Post)).scalar_one() + rec = importer.session.get(ImageRecord, r.image_id) + assert post.post_date is not None + assert rec.earliest_post_date == post.post_date + # Nothing left undated: the next sweep is an empty query. + assert date_posts_from_records(importer.session, importer, images_root) == { + "undated": 0, "dated": 0, "no_record": 0, + } -- 2.54.0 From 42a60e3e745aa1a06d3118e2ef73cfb6dffcc4e0 Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Fri, 25 Sep 2026 11:41:16 -0400 Subject: [PATCH 7/7] test: the patreon ingester tests run phase 3's seen-marking before counting post keys (#4436) Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR --- tests/test_patreon_ingester.py | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/tests/test_patreon_ingester.py b/tests/test_patreon_ingester.py index b986e53..7ff79c7 100644 --- a/tests/test_patreon_ingester.py +++ b/tests/test_patreon_ingester.py @@ -239,8 +239,8 @@ async def test_tick_downloads_unseen_and_marks_seen(source_id, sync_engine, tmp_ # plan #704: structured run_stats carry the real counts. assert result.run_stats["downloaded_count"] == 2 assert result.posts_processed == 1 - # The media wait for phase 3 to import them; only the post key is in yet. - assert _count_ledger(sync_engine, source_id) == 1 + # The media and the post record both wait for phase 3 to import them (#4436). + assert _count_ledger(sync_engine, source_id) == 0 result.mark_seen_after_import() # 2 media keys + 1 synthetic post key (body/links recaptured per post). assert _count_ledger(sync_engine, source_id) == 3 @@ -648,8 +648,10 @@ async def test_recovery_tier2_disk_still_skips(source_id, sync_engine, tmp_path) assert result.files_downloaded == 0 assert downloader.download_calls == 0 assert result.written_paths == [] - # Disk-skip reconciles the media key + the synthetic post key (recovery - # recaptures the body/links per post) = 2. + # Disk-skip reconciles the media key at once; the synthetic post key + # (recovery recaptures the body/links per post) waits for phase 3 (#4436). + assert _count_ledger(sync_engine, source_id) == 1 + result.mark_seen_after_import() assert _count_ledger(sync_engine, source_id) == 2 @@ -1203,7 +1205,10 @@ async def test_tick_captures_media_less_post_once(source_id, sync_engine, tmp_pa assert result.success is True assert len(result.post_record_paths) == 1 assert downloader.post_records == 1 - # The synthetic `post:ptext` key was marked seen (gates re-capture). + # The synthetic `post:ptext` key is marked once phase 3 has upserted the + # record (#4436), and then gates re-capture. + assert _count_ledger(sync_engine, source_id) == 0 + result.mark_seen_after_import() assert _count_ledger(sync_engine, source_id) == 1 # Second walk: already recorded → gated, no re-write, no new ledger row. @@ -1543,6 +1548,7 @@ async def test_revisits_do_not_feed_the_body_drift_canary( url="https://patreon.com/ingest", mode="tick", revisit_days=30, ) assert first.success is True + first.mark_seen_after_import() # phase 3 ran # Second walk: every post is a revisit, and every body comes back empty. client2 = _FakeClient([(None, posts)], published=published, empty_body=True) -- 2.54.0