feat: an open grouping — a later drop joins its post (milestone 388 step E3)
CI / extension-version (push) Successful in 3s
CI / lint (push) Successful in 3s
Build images / sign-extension (push) Successful in 4s
Build images / build-agent (push) Successful in 7s
CI / frontend-build (push) Successful in 22s
CI / backend-lint-and-test (push) Successful in 31s
Build images / build-web (push) Successful in 1m6s
Build images / smoke-web (push) Skipped
Build images / build-ml (push) Successful in 1m59s
Build images / promote (push) Skipped
CI / integration (push) Failing after 2m7s

A synthetic post is no longer sealed at creation. A creator who adds two more
variants the next day extends the existing post, its body grows with the new
messages, and no rival post appears. That is what makes chat capture read as
content trickling in rather than as a stream of separate arrivals.

The sweep now runs two passes per source and the ORDER is load-bearing: offer
new messages to still-open groups BEFORE founding new ones, because whichever
runs first claims a message.

E3's three named problems, each answered rather than discovered later:

**Bridging.** A candidate near two groups joins NEITHER. Nearest-wins would
silently make an arbitrary choice between two posts the operator may already
have seen; merging them is worse still, because a merge rewrites history and
anything pointing at the absorbed post dangles. Leaving it to found its own
group is the recoverable failure. AMBIGUITY_MARGIN is a module constant and
deliberately not a setting — it is not a quality dial anyone would tune toward
a better feed, and exposing it would invite turning it to zero, which is
exactly the silent arbitrary choice it prevents.

**Re-surfacing without thrashing.** A grouping has two dates, and which one
orders the feed is a real decision, so the feed orders by neither directly.
Ordering by when the drop STARTED buries a group that grows a week later under
a week of other posts — defeating the point of keeping it open. Ordering by
every growth lets a group gaining one image a day live permanently at the top,
so chat out-competes authored posts for the front page — the opposite of "post
pacing stays front and centre". Instead `resurfaced_at` moves only when growth
clears BOTH a minimum-images bar and a cooldown, so a drip-feed updates in
place and a genuine second wave resurfaces exactly once. It is NULL on every
ordinary post, so the sort key COALESCEs through it without moving anything
that is not a grouping.

**Reopening forever.** Groups close after a quiet period — artists reuse
characters for years, and a group left open indefinitely will eventually
absorb something it shouldn't. Openness is DERIVED, not stored: a group is
open if it grew (or started) within the window. Lowering the setting closes
old groups and raising it reopens them, with nothing to repair either way; a
stored closed_at would have needed a sweep to set it and a repair path to ever
change the policy.

Rule 89 is satisfied structurally rather than by a parallel mechanism:
celery_signals writes a TaskRun for every task, which already supplies
duration, the 5-minute stalled-run recovery, and retention pruning. What this
step owed on top of that was a wall-clock limit (present) and idempotence —
re-running the joiner adds nothing, asserted directly rather than left to the
unique (image, post) constraint to catch.

Two bugs fixed in the writing, one of which my own test would have hit:

* `assign_to_group` sorted bare (distance, Post) tuples, which falls through
  to comparing Posts when two distances tie — and a perfectly symmetric
  bridge, the exact case the function exists for, would have raised TypeError
  instead of declining to choose. Now keyed on the distance alone.
* The cursor was still built from `post_date or downloaded_at` while the
  ORDER BY had gained `resurfaced_at`. Two expressions that disagree at a page
  boundary don't error, they silently skip or repeat rows; both sites now go
  through one `_post_sort_value`, and a test pages through one row at a time
  to prove the walk matches the whole list.

Image linking is now one shared helper rather than written twice, because
creation and joining would otherwise be free to drift on exactly the detail
(which post owns the image) that makes a grouping reversible.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LNXXULQDjVZmbuNa2G9mD9
This commit is contained in:
2026-09-10 11:30:17 -04:00
co-authored by Claude Opus 5
parent 7071c87cd6
commit 1e45e2c56c
11 changed files with 974 additions and 38 deletions
+327 -29
View File
@@ -245,6 +245,10 @@ async def _synthesize(
"max_distance": max_distance,
"window_minutes": window_minutes,
"grouped_at": datetime.now(UTC).isoformat(),
# Growth accumulated since the post last moved in the feed (E3).
# Seeded here so the joiner never meets its absence — creation IS
# a surfacing, so the count starts at zero.
"images_since_surface": 0,
},
)
session.add(post)
@@ -256,31 +260,52 @@ async def _synthesize(
.values(absorbed_by_post_id=post.id)
)
# Link every member image to the synthetic post via provenance. The feed
# and detail views already union provenance with primary_post_id
# (post_feed_service._thumbnails_for), so this alone makes the drop's
# images show up under the post FC wrote — no second render path.
#
# primary_post_id is deliberately NOT rewritten: the message post remains
# the image's true origin, and the synthetic post is an ADDITIONAL claim on
# it, which is what keeps the grouping reversible.
image_rows = (await session.execute(
select(ImageRecord.id).where(ImageRecord.primary_post_id.in_([m.id for m in members]))
)).scalars().all()
if image_rows:
await session.execute(
pg_insert(ImageProvenance)
.values([
{"image_record_id": iid, "post_id": post.id, "source_id": source.id}
for iid in image_rows
])
# (image, post) is unique; a re-run that raced itself is a no-op
# rather than an IntegrityError that loses the whole sweep.
.on_conflict_do_nothing(constraint="uq_image_provenance_image_post")
)
await _link_member_images(
session, post_id=post.id, member_ids=[m.id for m in members],
source_id=source.id,
)
return post
async def _link_member_images(
session: AsyncSession, *, post_id: int, member_ids: list[int], source_id: int | None,
) -> int:
"""Attach every member's images to the synthetic post. Returns how many.
The feed and detail views already union provenance with `primary_post_id`
(post_feed_service._thumbnails_for), so this alone makes the drop's images
show up under the post FC wrote — no second render path.
`primary_post_id` is deliberately NOT rewritten: the message post remains
the image's true origin, and the synthetic post is an ADDITIONAL claim on
it, which is what keeps the grouping reversible.
Shared by creation and by the E3 joiner rather than written twice, because
the two would otherwise be free to drift on exactly the detail (which post
owns the image) that makes a grouping reversible.
"""
if not member_ids:
return 0
image_rows = (await session.execute(
select(ImageRecord.id).where(ImageRecord.primary_post_id.in_(member_ids))
)).scalars().all()
if not image_rows:
return 0
await session.execute(
pg_insert(ImageProvenance)
.values([
{"image_record_id": iid, "post_id": post_id, "source_id": source_id}
for iid in image_rows
])
# (image, post) is unique. This is a BACKSTOP against a re-run that
# raced itself, not the correctness argument: callers only ever pass
# members that were unabsorbed a moment ago, so a conflict here means
# concurrency, not a logic error.
.on_conflict_do_nothing(constraint="uq_image_provenance_image_post")
)
return len(image_rows)
async def group_source(
session: AsyncSession,
source: Source,
@@ -329,11 +354,265 @@ async def group_source(
return created
# ---------------------------------------------------------------------------
# E3: an open grouping — a later drop joins its post and updates it.
# ---------------------------------------------------------------------------
#
# A synthetic post is not sealed at creation. A creator who adds two more
# variants the next day extends the existing post rather than starting a new
# one, and its body grows with the new messages. That is what makes chat
# capture read as content TRICKLING IN rather than as a stream of separate
# arrivals.
#
# Three hard problems, each answered deliberately below: bridging (a candidate
# near two groups), re-surfacing without thrashing the feed, and groups that
# stay open forever.
# How much closer the nearest group must be than the runner-up before a
# candidate is assigned to it at all.
#
# NOT a setting, deliberately. It is not a quality dial the operator would tune
# toward a better feed — it expresses "these two are too close to call", and
# exposing it would invite turning it to zero, which is precisely the silent
# arbitrary choice it exists to prevent. When a candidate is genuinely between
# two groups the recoverable answer is to leave it out and let it start its
# own; the unrecoverable one is to merge, because a merge rewrites history —
# two posts the operator may already have seen become one, and anything
# pointing at the absorbed post dangles.
AMBIGUITY_MARGIN = 0.02
def should_resurface(
*,
images_since_surface: int,
last_surface_at: datetime,
now: datetime,
min_images: int,
cooldown: timedelta,
) -> bool:
"""Has this group grown enough, and waited long enough, to move in the feed?
An updated post SHOULD be visible — that is the point of keeping it open —
but a group gaining one image a day must not sit permanently at the top.
Both conditions have to hold: enough new images that the update is worth an
interruption, and enough time since the last one that a steady drip cannot
chain bumps together.
"""
if images_since_surface < min_images:
return False
return now - last_surface_at >= cooldown
async def _group_seed(session: AsyncSession, post_id: int) -> list[float] | None:
"""The embedding a group is measured against — its FIRST member's image.
Derived rather than stored, and derived by the same definition `build_groups`
used (earliest member, lowest-id embedded image). Storing it at creation
would have meant a backfill for groups already written and two definitions
free to disagree; this way there is one.
"""
sort_key = func.coalesce(Post.post_date, Post.downloaded_at)
return (await session.execute(
select(ImageRecord.siglip_embedding)
.join(Post, ImageRecord.primary_post_id == Post.id)
.where(
Post.absorbed_by_post_id == post_id,
ImageRecord.siglip_embedding.is_not(None),
)
.order_by(sort_key, Post.id, ImageRecord.id)
.limit(1)
)).scalar_one_or_none()
async def open_groups(
session: AsyncSession, source_id: int, *, now: datetime, close_after: timedelta,
) -> list[tuple[Post, list[float]]]:
"""This source's synthetic posts that are still accepting members.
Openness is DERIVED, not stored: a group is open if it grew — or started —
within `close_after`. A group left open forever would eventually absorb
something it shouldn't, because artists reuse characters for years; a
stored `closed_at` would need a sweep to set it and a repair path to ever
change the policy. This way the policy IS the query.
"""
cutoff = now - close_after
posts = (await session.execute(
select(Post).where(
Post.source_id == source_id,
Post.synthesized_by == DROP_GROUPER,
# A synthetic post that was itself absorbed is not a thing today
# (nothing absorbs one), but joining into one would nest groups.
Post.absorbed_by_post_id.is_(None),
func.coalesce(Post.last_grew_at, Post.post_date, Post.downloaded_at) >= cutoff,
)
)).scalars().all()
out: list[tuple[Post, list[float]]] = []
for post in posts:
seed = await _group_seed(session, post.id)
if seed is not None:
out.append((post, seed))
return out
def assign_to_group(
embedding: list[float],
groups: list[tuple[Post, list[float]]],
*,
max_distance: float,
) -> Post | None:
"""Pick the one group this image belongs to, or None to leave it alone.
Returns None in two cases that mean different things and are deliberately
treated the same: nothing is close enough (so E2 will start a new group
from it), or two groups are BOTH close and too near each other to choose
between (so E2 will start a new group from it). The second is the bridging
case, and letting it start its own post is the recoverable failure —
merging two existing posts is not.
"""
# key= on the distance ALONE. Sorting bare tuples falls through to the
# second element when two distances tie, and `Post` has no ordering — so a
# perfectly symmetric bridge (the exact case this function exists for)
# would raise TypeError instead of declining to choose.
scored = sorted(
((cosine_distance(seed, embedding), post) for post, seed in groups),
key=lambda pair: pair[0],
)
within = [(d, p) for d, p in scored if d <= max_distance]
if not within:
return None
if len(within) >= 2 and (within[1][0] - within[0][0]) < AMBIGUITY_MARGIN:
return None
return within[0][1]
async def _absorb_into(
session: AsyncSession,
*,
group: Post,
member_ids: list[int],
source_id: int | None,
now: datetime,
min_images: int,
cooldown: timedelta,
) -> int:
"""Extend an existing synthetic post with new members. Returns images added."""
members = (await session.execute(
select(Post)
.where(Post.id.in_(member_ids))
.order_by(func.coalesce(Post.post_date, Post.downloaded_at), Post.id)
)).scalars().all()
if not members:
return 0
added_images = await _link_member_images(
session, post_id=group.id, member_ids=[m.id for m in members],
source_id=source_id,
)
await session.execute(
update(Post)
.where(Post.id.in_([m.id for m in members]))
.values(absorbed_by_post_id=group.id)
)
# The new messages' text joins the body, in arrival order, exactly as at
# creation — the post's body is the drop's text and the drop just grew.
new_text = "\n\n".join(
m.description.strip() for m in members if m.description and m.description.strip()
)
if new_text:
group.description = f"{group.description}\n\n{new_text}" if group.description else new_text
details = dict(group.synthesis_details or {})
existing_ids = list(details.get("member_post_ids") or [])
details["member_post_ids"] = existing_ids + [
m.id for m in members if m.id not in existing_ids
]
details["message_count"] = len(details["member_post_ids"])
since = int(details.get("images_since_surface") or 0) + added_images
details["last_grew_at"] = now.isoformat()
# The feed position moves only when the anti-thrash rule fires. Measured
# from the last time the post actually MOVED (resurfaced_at), falling back
# to when the drop started — creation is itself a surfacing.
last_surface = group.resurfaced_at or group.post_date or group.downloaded_at
if should_resurface(
images_since_surface=since, last_surface_at=last_surface, now=now,
min_images=min_images, cooldown=cooldown,
):
group.resurfaced_at = now
since = 0
details["images_since_surface"] = since
group.synthesis_details = details
group.last_grew_at = now
return added_images
async def join_open_groups(
session: AsyncSession,
source: Source,
*,
max_distance: float,
window_minutes: float,
close_after_hours: float,
resurface_min_images: int,
resurface_cooldown_hours: float,
now: datetime | None = None,
) -> int:
"""Offer this source's ungrouped messages to its open groups.
Runs BEFORE `group_source` in the sweep: a message that belongs to an
existing drop must join it rather than found a rival post, and whichever
runs first wins that message.
"""
now = now or datetime.now(UTC)
window = timedelta(minutes=window_minutes)
groups = await open_groups(
session, source.id, now=now, close_after=timedelta(hours=close_after_hours),
)
if not groups:
return 0
# Same quarantine as E2: a drop still arriving is left for the next run.
rows = (await session.execute(
_candidate_stmt(source.id, not_after=now - window)
.limit(MAX_CANDIDATES_PER_SOURCE)
)).all()
if not rows:
return 0
claimed: dict[int, list[int]] = {}
for post_id, _occurred_at, embedding in rows:
target = assign_to_group(embedding, groups, max_distance=max_distance)
if target is not None:
claimed.setdefault(target.id, []).append(post_id)
by_id = {post.id: post for post, _seed in groups}
joined = 0
for group_id, member_ids in claimed.items():
joined += await _absorb_into(
session, group=by_id[group_id], member_ids=member_ids,
source_id=source.id, now=now,
min_images=resurface_min_images,
cooldown=timedelta(hours=resurface_cooldown_hours),
)
return joined
async def sweep(session: AsyncSession, *, now: datetime | None = None) -> dict:
"""Group every enabled Discord source. No-op when the switch is off."""
"""Group every enabled Discord source. No-op when the switch is off.
Two passes per source, and the ORDER is load-bearing: offer new messages to
the groups that are still open (E3) BEFORE founding new ones (E2). Whichever
runs first claims a message, and a variant that belongs to yesterday's drop
must extend that post rather than found a rival to it.
"""
settings = await MLSettings.load(session)
if not settings.discord_grouping_enabled:
return {"enabled": False, "sources": 0, "posts_created": 0}
return {
"enabled": False, "sources": 0, "posts_created": 0, "images_joined": 0,
}
sources = (await session.execute(
select(Source).where(
@@ -342,16 +621,35 @@ async def sweep(session: AsyncSession, *, now: datetime | None = None) -> dict:
)
)).scalars().all()
max_distance = float(settings.discord_group_max_distance)
window_minutes = float(settings.discord_group_window_minutes)
created = 0
joined = 0
for source in sources:
joined += await join_open_groups(
session, source,
max_distance=max_distance,
window_minutes=window_minutes,
close_after_hours=float(settings.discord_group_close_after_hours),
resurface_min_images=int(settings.discord_group_resurface_min_images),
resurface_cooldown_hours=float(
settings.discord_group_resurface_cooldown_hours
),
now=now,
)
created += await group_source(
session, source,
max_distance=float(settings.discord_group_max_distance),
window_minutes=float(settings.discord_group_window_minutes),
max_distance=max_distance,
window_minutes=window_minutes,
now=now,
)
log.info(
"discord drop grouping: %d source(s), %d synthetic post(s) created",
len(sources), created,
"discord drop grouping: %d source(s), %d synthetic post(s) created, "
"%d image(s) joined to open groups",
len(sources), created, joined,
)
return {"enabled": True, "sources": len(sources), "posts_created": created}
return {
"enabled": True, "sources": len(sources),
"posts_created": created, "images_joined": joined,
}
+43 -4
View File
@@ -41,8 +41,40 @@ THUMBNAIL_LIMIT = 6
def _sort_key():
"""Postgres COALESCE expression used in ORDER BY and WHERE clauses."""
return func.coalesce(Post.post_date, Post.downloaded_at)
"""Postgres COALESCE expression used in ORDER BY and WHERE clauses.
`resurfaced_at` leads (milestone 388 E3). A synthetic post stays OPEN — a
creator who adds variants the next day extends the existing post — so such
a post has two dates, and which one orders the feed is a real decision:
* ordering by when the drop STARTED buries a group that grows a week later
under a week of other posts, so the operator never sees the new content —
which defeats keeping the group open at all;
* ordering by every growth lets a group that gains one image a day sit
permanently at the top, so chat out-competes authored posts for the front
page — the opposite of "post pacing stays front and centre".
So the feed orders by neither directly. `resurfaced_at` moves only when the
anti-thrash rule fires (discord_grouping.should_resurface: enough new
images AND enough time since the last move), which means a drip-feed
updates IN PLACE and a genuine second wave resurfaces exactly once.
It is NULL on every ordinary post, so this COALESCE cannot move anything
that is not a grouping. Used identically in ORDER BY and in the cursor's
WHERE, which is what keeps pagination stable across the change.
"""
return func.coalesce(Post.resurfaced_at, Post.post_date, Post.downloaded_at)
def _post_sort_value(post: Post):
"""The Python twin of `_sort_key()`, for building a cursor from a loaded row.
Kept next to it on purpose: these two are one expression in two languages,
and the failure when they disagree is not an error but a quiet one — rows
skipped or repeated at page boundaries, which reads as a backend bug
anywhere but here.
"""
return post.resurfaced_at or post.post_date or post.downloaded_at
class PostFeedService:
@@ -134,7 +166,10 @@ class PostFeedService:
# Far edge in the travel direction: oldest row going older,
# newest row going newer (rows is descending for display).
edge_post = rows[-1][0] if direction == "older" else rows[0][0]
edge_key = edge_post.post_date or edge_post.downloaded_at
# Must match _sort_key() exactly, including resurfaced_at's
# precedence: a cursor built from a different expression than the
# ORDER BY silently skips or repeats rows at every page boundary.
edge_key = _post_sort_value(edge_post)
next_cursor = encode_cursor(edge_key, edge_post.id)
post_ids = [p.id for p, _, _ in rows]
@@ -168,7 +203,7 @@ class PostFeedService:
if anchor is None:
return None
anchor_post, anchor_artist, anchor_source = anchor
anchor_key = anchor_post.post_date or anchor_post.downloaded_at
anchor_key = _post_sort_value(anchor_post)
anchor_cursor = encode_cursor(anchor_key, anchor_post.id)
older = await self.scroll(
@@ -414,6 +449,10 @@ class PostFeedService:
# keys are always present so the frontend never branches on absence.
"synthesized_by": post.synthesized_by,
"synthesis": post.synthesis_details,
# #388 E3. A grouping stays open, so the card can say "updated N
# ago" — which is the whole signal that chat content is trickling
# in. NULL means it has not grown since it was created.
"last_grew_at": post.last_grew_at.isoformat() if post.last_grew_at else None,
# Non-null on a chat message a synthetic post absorbed. The feed
# filters these out, but `around`/`get_post` still reach them, and
# the UI uses this to explain why a post it linked to is not in the