Skip to content

Pipeline v2 core: platform-authoritative drop state (opt-in) - #52

Open
aalejandrofer wants to merge 60 commits into
masterfrom
feat/pipeline-v2
Open

aalejandrofer wants to merge 60 commits into
masterfrom
feat/pipeline-v2

Conversation

@aalejandrofer

@aalejandrofer aalejandrofer commented Sep 28, 2026 •

Copy link
Copy Markdown
Owner

Pipeline v2 core (plan 1 of 2) for #50, with streamer priority from #49 folded in. It's headless and opt-in: the default stays v1, and internal/watcher is untouched.

What

  • Per-drop state from the platform. A new drop_state table (migration 0016) holds one row per account and drop: eligible, accruing, claimable, claimed, or blocked with a reason. The reconciler fills it from each campaign's platform details (Twitch DropCampaignDetails self, Kick /drops/progress). A claim the platform has confirmed is never demoted.
  • New packages under internal/pipeline/:
    • dropstate: pure state transitions.
    • reconcile: syncs drop state from the platform.
    • planner: pure; decides which channel to watch and which drops it serves, with streamer priority, a 10-minute swap hold and force-watch.
    • session: runs one watch and reports progress, stream-down and per-drop stalls.
    • claimer: turns typed claim results into state, with backoff 1m, 5m, 30m, then blocks after 5 failures.
    • loop: one goroutine per account owns every state change.
  • Platform capabilities. DropProgressSource, DropClaimer (typed ClaimResult) and ChannelProber, with Twitch and Kick adapters. Kick's required minutes now come from the platform.
  • Backfill. A one-time step copies claims, manual marks and ghost-skips into drop_state. The old kv keys stay, so rolling back to v1 is safe.
  • Streamer priority table. account_streamer_priority is seeded from account_channels.
  • Scheduler reads snapshots and discoveries through interfaces, so both v1 and v2 runners report state.
  • Opt-in. Enable with GRUB_PIPELINE=v2 or kv pipeline_override:<accountID>=v2. If a backend lacks the v2 capabilities (for example Twitch BrowserBackend), that account falls back to v1.

Why

Fixes the "claimed on the platform, shown unclaimed, mining stuck" bug and the other cases listed in spec §2 (P1 to P12): lost claim responses, endless claim retries, Kick link-required claim loops, Kick minutes that never reach required, the ghost-skip deadlock, and full reloads on manual mark. There's now one state model, and every drop has a reason for being mined or not.

  • Spec: docs/superpowers/specs/2026-09-27-pipeline-v2-design.md (see §10 for refinements made during planning)
  • Plan: docs/superpowers/plans/2026-09-27-pipeline-v2-core.md

Not in this PR

Update: merged with v1.4.0

v1.4.0 (TV-client Twitch login, channel-first discovery) is merged in. The two conflicts were twitch/claim.go, where the claim split now also keeps master's null-claim check, and CHANGELOG. On top of the merge:

  • TV accounts find campaigns in v2. The whitelisted game names now reach the session, which v1.4.0's channel-first discovery needs.
  • TV accounts read progress from Inventory. Twitch hides campaign details from the TV client, so v2 never calls it for TV sessions. Android sessions keep the details path.
  • Subscription-only drops are no longer mined in v2 (Bot trying to farm Sub rewards #47 rule, applied on both paths).
  • Every v2 Twitch call sends the account's own Client-Id. v2's three new methods now bind the session's client profile themselves.
  • Known limit on TV: a campaign that is finished and was claimed outside the app drops out of Inventory, and TV has no details query to fall back on. v1.4.0 has the same blind spot.

Live verification now also needs a TV-client Twitch account. The Twitch details probe only matters for Android sessions.

Before merge or tag

This touches accrual and claims, so no tag until a live drop is verified on staging.

  • Live probes, behind the live build tag: Twitch DropCampaignDetails returns self with an Android token and a TV-client token (Bot trying to farm Sub rewards #47). Kick keeps claimed rewards in /drops/progress.
  • Staging: one Twitch and one Kick account on v2 with a live drop. Minutes accrue, the claim lands, and a drop claimed on the website flips to claimed within one reconcile.
  • Watch on staging: a fresh Kick reward that isn't listed yet may count as a stall.

Known follow-up

internal/canary TestKickProbe_OK races under -race (kick_test.go:79/133). It already fails on master and this PR doesn't touch it.

Checklist

  • go build ./..., go vet ./..., go test ./... pass
  • gofmt -l . is clean
  • UI changes follow docs/DESIGN.md (no UI changes)
  • No secrets / personal infra added

🤖 Generated with Claude Code

aalejandrofer and others added 30 commits September 27, 2026 03:45
Design for #50 (folds in #49 streamer priority): platform-authoritative
per-drop state, pure planner, session/claimer split, v1/v2 flag rollout.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ed retention

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
…ub_only

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
…resis

Refs #49

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
…e a stall

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
… adapter

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ed_units

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ry, drop stale events

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…faces

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…, unlinked, allow-list, integrity wall)

- C1: unknown observations adopt the catalog requirement when the row has
  none; backfill seeds required minutes on skip, claim and mark rows.
- I2: backfill applies manual marks after claim history so they stay user claims.
- I3: campaign link block uses the new unlinked reason, lifted on the next
  definite answer; claim-level needs_link keeps its 24h hold.
- I4: skip priority probing for restricted campaigns with no known allow-list.
- I5: Twitch self:null is Known=false unless the in-progress inventory has it.
- I6: integrity wall stops the session, reports auth_required, Run exits nil.
- stopSession waits up to 15s for the session goroutine; same-channel
  restarts keep the PubSub subscription.
- Claimable to Accruing/Eligible resets FailCount.
- session_test.go back to mode 644.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Resolves claim.go (keep claimStatus split, add master's null claimDropRewards
check) and CHANGELOG (v2 entries stay under Unreleased above 1.4.0).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
TV-client Twitch sessions discover channel-first and need Session.Games;
v2 never set it, so TV accounts discovered nothing.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
DropCampaignDetails returns dropCampaign:null for TV-client tokens, so
v2 DropProgress skips it for TV sessions and builds observations from
Inventory (minutes, requirement, claimed, instance id). v1 inventory()
is untouched.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
DropProgress, ClaimDrop and ProbeChannels can run before any v1 entry
point bound the token, so a TV session went out under the Android
Client-Id.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Dashboard HEARTBEATS/HR counts kind=heartbeat log lines; v2 logged none,
so v2 accounts showed 0. session.Config gains AccountID for the field.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
aalejandrofer and others added 30 commits September 28, 2026 02:59
DropCampaignDetails returns self: null for claimed campaigns, so a drop
claimed outside GrubDrops whose campaign left Inventory stayed Eligible and
was watched, stalled and cooled forever. On stall, the loop now sends one
claim per served Eligible drop absent from the last inventory read:
OK/ALREADY_CLAIMED marks it claimed, NeedsLink blocks needs_link, failure
blocks not_enrolled (1h re-check, FailCount untouched). Twitch only
(StallClaimProbe); Kick keeps claimed rewards listed. Loop passes AccountID
to the session heartbeat log. Spec §9 results + §10 item 9, changelog.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
… is failure)

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
v1->v2 account switches leave drop_state with no row yet for a drop v1
already claimed. load()'s bridge only touches rows that already exist,
so the first reconcile creates the row from the platform's silent
answer and nothing ever corrects it: Twitch gives no definite signal
for a drop claimed outside the app whose campaign left the Inventory,
so it stayed eligible and got watched forever (22 Rust drops on prod).

Extract bridgeClaimHistory() and call it from both load() and the end
of reconcile(), reusing dropstate.MarkCollected + commit() so a later
definite platform answer still overwrites a user mark, and a row
already Claimed is a no-op.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
AvailableDrops returns streamers' own game-less channel campaigns and
other games' campaigns; listByChannels persisted them all (753 'no game'
rows flooding /drops). Keep only the game being walked; Inventory
campaigns are unaffected.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Twitch hides ViewerDropsDashboard/DropCampaignDetails from TV tokens and
AvailableDrops returns one campaign per channel, so a global badge
campaign masked every other campaign of a game (R6 S2 watch drops hidden
from TV accounts). Non-TV backends now publish their campaign list +
allow-lists to one process-wide twitch.Catalog (link flags made
account-neutral); TV backends merge fresh (45m TTL), active, whitelisted
campaigns they didn't find themselves. Inventory stays authoritative.
Watch/claim requests keep each session's own client.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
/drops tabs already group by ends_at, not stored status; this test pins
that an ended campaign with status=active lands in Past, never Current.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…any Known observation

Prod: a TV-login account mined a drop whose campaign reported
self.isAccountConnected=false, because dropstate.Apply re-derived
Blocked/Unlinked on the next Known observation regardless of source. The
session's own Progress events (loop.applyProgress -> Apply with Known=true)
lifted the block between reconciles just like a real reconcile did, so the
planner mined a drop whose claim could never succeed until the next sync
re-blocked it.

- dropstate.Apply: Blocked/Unlinked now behaves like Blocked/UserSkip - only
  a claimed observation moves it; every other Known observation leaves it
  alone. Lifting it is the reconciler's job.
- reconcile.Run: when the campaign is linked and the prev row is
  Blocked/Unlinked, lift it first with dropstate.Retry (re-derives from
  minutes/required) before folding in the observation. Still-unlinked rows
  keep BlockLink after Apply as before.
- Updated the superseded I3 test (Apply re-derives Unlinked immediately on
  any Known observation) to the new rule; the reconcile-level lift test it
  was standing in for already exists and continues to pass unchanged
  (TestRun_LinkedAfterUnlinkedUnblocksImmediately).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Twitch's TV OAuth client (every login since #48/v1.4.0) can't see the
drops dashboard, so those accounts only find campaigns for whitelisted
games. Add a banner on the dashboard and /drops when at least one
enabled Twitch account has a TV-client session, pointing at the game
whitelist (/priority); it renders as a warning when that account's
effective whitelist (account games, falling back to global) is empty.

Reuses the existing dashAlert banner component (now factored into a
shared _dash_alerts.html partial so /drops gets it too) instead of a
new one, per the design system's flat/no-new-component rule.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
- Channel add/remove on /drops now writes both account_channels (v1) and
  account_streamer_priority (v2), since v2 reads the priority table
  exclusively and account_channels was only ever copied there once by the
  0016 migration seed. New adds rank after the account's current max.
- Add sqlc AddStreamerPriority/RemoveStreamerPriority queries.
- The v2 branch in build() now only passes ListForceChannels rows when the
  account's force-watch toggle is enabled, via a shared forceWatchEnabled
  helper reused by v1's forceWatchStore.Next (no v1 behaviour change).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Move tvDiscoveryAlert computation out of collectPage to the full-page
handler (page()). The TV discovery banner decrypts every enabled Twitch
account's session to check for TV-client logins, which was running
unnecessarily on every polled HTMX partial (cards/telemetry, every 10s).
collectPage now skips the expensive session scan entirely; only the full
dashboard render adds the alert. The /drops page handler still calls
tvDiscoveryAlert directly and is unaffected.

Test: collectPage no longer includes tv_discovery alerts; full-page
banner computation is verified separately via tvDiscoveryAlert function.

Co-Authored-By: Claude Haiku 4.5 <noreply@anthropic.com>
Updated warning message for clarity.
The Required==0 branch set Claimable directly, ignoring a user skip or a
campaign link block. Both branches now return early for those rows.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Directory pages are top-N and churn; a refresh that omitted the channel
being mined made v2 swap away and back every LiveEvery. refreshLive now
re-probes just that channel for the campaigns it serves and keeps it
when still live. No prober: unchanged.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
An account added after finishing its drops has every row Eligible with
no platform record, and v2 only learned they were claimed via two
stalls per drop. After each reconcile, queue up to 5 such drops (active
in-scope campaign, Required>0, Known=false, unprobed this run) and send
one claim each, one every 3s off a loop timer so the loop stays
responsive. OK/Already -> claimed, NeedsLink -> needs_link, Failed ->
row unchanged. Rename StallClaimProbe to ClaimProbe (gates both;
still Twitch-only in main.go).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…rces

- listByChannels now matches vc.Game.ID against the directory stream's
  game ID first, falling back to the display-name slug compare only
  when either ID is empty. A campaign whose game display name diverges
  from the directory's was wrongly dropped even though both carry the
  same game ID.
- Catalog.Publish/snapshot now evict sources older than catalogTTL
  instead of merely skipping them, so a Reload-recycled backend (keyed
  by pointer when AccountID is unset) doesn't grow the sources map
  without bound across the process lifetime.
- "tv discovery: merged shared catalog" now logs at INFO only when it
  actually merged something; a no-op merge logs at DEBUG.
- Added a concurrent Publish/snapshot test (run under -race).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Reload built every EntryBuilder serially, so one slow account (e.g. a
browser-sidecar cold start) held up every other account's rebuild.
buildEntriesConcurrently now runs builders with a cap of 4 workers
while preserving the original entry order (indexed by builder
position, not completion order).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ence

A Kick account with no live channel for any campaign exhausts every
campaign, idles, and immediately re-arms on the next tick, cycling
pick_stream→pick_campaign→sleeping roughly every TickInterval and
spamming state-change logs. nextKickIdleWait ramps the idle
re-discovery wait 30s -> 60s -> 120s, then holds at the same
recheckInterval (5m) cap Twitch already uses, and resets to the 30s
floor the moment a pick succeeds. Twitch's fixed recheckInterval cadence
is unchanged.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
pipelineModeFor now returns v2 unless GRUB_PIPELINE=v1 or a
per-account kv override says v1. build()'s existing automatic
fallback to v1 when loop.New fails (e.g. no BrowserBackend for Kick)
is unchanged.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
README config table gains GRUB_PIPELINE (default v2, v1 = legacy
fallback); spec section 6.5 records the default flip. CHANGELOG
[Unreleased] logs the batch B changes: v2 default, parallel Reload
builds, Kick idle-sleep backoff, and the catalog game-ID-match +
eviction fixes.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
TV-minted Twitch tokens can't read the campaign list, but Inventory lists
every campaign the viewer is enrolled in, and watching a counting channel
enrolls them. When Plan() says Idle, the v2 loop now watches the top live
DROPS_ENABLED channel of the next whitelisted game (round-robin over Games,
via ListEligibleChannels on a synthetic open campaign -> Twitch game
directory) for EnrollWatch (10m), then marks the game probed (6h cooldown),
reconciles and refreshes channels. Any non-Idle decision preempts the
enroll watch at once; a down channel is cooled and the game retried on its
next channel. Planner stays pure: enrolling is loop-side only.

Snapshot state "discovering" with channel + game. Wired in main.go for
twitch sessions with ClientID == twitch.ClientTV; Kick and Android off.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Critical review finding on batch B item 3: idleWait was reset whenever
a step returned nil and the state wasn't Sleeping/AwaitingConnect.
pickCampaign->StatePickStream and pickStream's no-live-channel
path->StatePickCampaign both return nil too, so for a Kick account
with matched-but-offline campaigns idleWait was zeroed every round and
the ramp never climbed past its first step -- re-running discovery
every 30s instead of the old flat 5m, worse than before.

idleWait now resets only when the watcher actually reaches
StateWatching (a successful stream pick), not on intermediate
discovery states. Twitch's fixed recheckInterval cadence is unchanged.

Added a Run()-level regression test
(TestKickIdleWait_RampsWhileOfflineAndResetsOnlyAfterWatching) that
drives a real Watcher through repeated
pick_campaign->pick_stream(no live)->sleeping rounds and asserts the
idle wait climbs step to step and caps, then resets to the floor only
after a round reaches StateWatching. Verified revert-proof: restoring
the old reset condition fails the ramp assertions (every gap collapses
to ~the floor step).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…, reset state on exit

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

Status: In Progress

Development

Successfully merging this pull request may close these issues.

1 participant