Compare commits

...
4 Commits
Author SHA1 Message Date
mathiasandClaude Opus 4.8 a9be5f285b feat(summarizer): move local fallback off koala to iguana/gemma4-26b
CI / Build & Import (push) Successful in 10s
CI / Lint / Test / Vet (push) Successful in 10s
koala now carries other GPU loads, so the first fallback should not run there.
Change the default chain to koala/phi4-mini → iguana/gemma4-26b → berget/mistral-small:
the local fallback now runs on iguana (M2 Ultra headroom, different host = different
egress IP for the rare fallback fetch). gemma4-26b is the brain-validated homelab
general-purpose model (agentsquad H2/H3 executor) and returned valid summary JSON on
the real prompt in a smoke test (~37s incl. cold-load — fine for a fallback path).

Pure config default (TAPIR_FALLBACK_MODEL); chain mechanism (ADR-022) unchanged.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-11 07:54:56 +02:00
mathiasandClaude Opus 4.8 fe56e2fe01 test: make embedded-postgres per-process so concurrent CI runs don't collide
CI / Lint / Test / Vet (push) Successful in 9s
CI / Build & Import (push) Successful in 10s
A push to main and its version tag fire two CI runs for the same commit. Both ran
`go test ./...`, which starts embedded-postgres on a FIXED port (54329/54330) and
a shared data dir. -p 1 serialises packages WITHIN a run, not across two
concurrent runs — so when the two runs overlapped they collided on the port/data
dir and BOTH failed the Lint/Test job (no image built). Prior commits passed only
because their two runs happened not to overlap.

Derive the port and runtime/data dirs from the PID; share only CachePath so the
PG archive downloads once. Proven: two concurrent `go test` of the store package
now both pass. Unblocks the v0.21.0 (Pillar A) build.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-11 07:26:33 +02:00
mathiasandClaude Opus 4.8 beeb5bc31b feat(gate): foreground caption fetches take priority over the background sweep (ADR-026)
CI / Lint / Test / Vet (push) Failing after 9s
CI / Build & Import (push) Has been skipped
A user waiting on a Summarize click shared the per-IP caption gate equally with
the background firehose, so on a busy IP the click was slow or 429'd. Add a
context-marked priority lane: the web path (engineProcessor.ProcessVideo) marks
its context foreground; the gate serves foreground immediately while background
fetches yield until no foreground is pending. Threaded via a context value (no
new signatures) + a process-wide foregroundPending counter. Clicks are rare, so
the background barely loses throughput; the waiting human gets the cleaner slot.

Drops the credentials probe: ADR-010 and captions.go already settle it — the
timedtext/InnerTube path rejects authenticated requests and the OAuth token does
not authenticate it anyway, so auth cannot help and can hurt. Documented in
ADR-026 rather than built.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 22:10:15 +02:00
mathiasandClaude Opus 4.8 09eb31d1fe fix(scheduler): derive rotation offset from wall-clock, not a reset-on-restart counter
CI / Lint / Test / Vet (push) Successful in 11s
CI / Build & Import (push) Successful in 14s
The lead-user rotation used an in-memory pass counter reset to 0 on every pod
restart, so the first-listed user always re-took the lead after a restart — a
deploy-heavy session re-starved the last user (the pilot stalled at 2 summaries
because each deploy reset his every-other-pass lead before the 2h tick fired).
Derive the offset from wall-clock (floor(now/interval)) so it advances with real
time and is identical across restarts: rotation stays fair however often the pod
bounces.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 21:59:12 +02:00
9 changed files with 180 additions and 15 deletions
+42 -4
View File
@@ -827,10 +827,18 @@ The fix is resilience around it, not replacing it.
**Decision.**
1. **Ordered endpoint chain (`summarizer.NewChain`).** Endpoints are tried in order; the first to
return a *parseable* summary wins. Default chain:
`koala/phi4-mini` (primary, local) → `koala/phi4-14b` (fallback, local) →
`berget/mistral-small` (worst-case, external). All three are reached through the **one** LiteLLM
gateway by alias — the gateway already fronts both llama-swap and berget — so a fallback is a
different alias, not a second client config.
`koala/phi4-mini` (primary, local) → `iguana/gemma4-26b` (fallback, local on a
*different host*) → `berget/mistral-small` (worst-case, external). All three are reached through
the **one** LiteLLM gateway by alias — the gateway already fronts both llama-swap and berget — so
a fallback is a different alias, not a second client config.
**Update 2026-06-11:** the local fallback moved from `koala/phi4-14b` to `iguana/gemma4-26b`.
koala now carries other GPU loads, so keeping the fallback on koala competed with them; iguana
(M2 Ultra) has the headroom, and a different host is also a different egress IP for the rare
fallback fetch. `gemma4-26b` is the brain-validated homelab general-purpose model (agentsquad
H2/H3 executor) and returned valid summary JSON on the real prompt in a smoke test
(~37s incl. cold-load — fine for a path hit only when the fast primary fails). Pure config:
`TAPIR_FALLBACK_MODEL`.
2. **A parse failure advances the chain, same as a transport error.** "Reliably summarized" means
*parseable summary returned*, not *HTTP 200*. This is the behaviour the old Primary→Fallback
shape missed.
@@ -958,6 +966,36 @@ schema change (reuses `transcript_status` from migration 007).
---
## ADR-026 — Foreground caption fetches take priority; the credentials probe is dead
**Status:** Accepted (2026-06-10). **Pillar A of the manual-mode UX work** (Pillar B was
ADR-025). Builds on ADR-014 (the shared per-IP gate).
**Context.** Every caption fetch — the background sweep and the web click-path — shared one
process-wide rate gate equally. So a user waiting on a "Summarize" click competed with the
firehose for both pacing and the scarce pre-429 window; on a busy IP the click was slow or
429'd while the background churned.
**Decision.** A context-marked priority lane. The web path
(`engineProcessor.ProcessVideo`) wraps its context with `ForegroundContext`; the gate gives
foreground fetches a token immediately, while **background fetches yield** — they wait until no
foreground fetch is pending before taking a token. Threaded via a context value (not new
signatures) and a process-wide `foregroundPending` counter. Clicks are rare and bursty, so the
background barely loses throughput; the waiting human gets the next (and cleanest) slot.
**Credentials probe — rejected, not built.** The idea was to fetch captions with the user's
auth in manual mode to dodge 429s. It is a dead end, already settled by ADR-010 and the code:
the caption path is *deliberately anonymous* because the InnerTube/timedtext endpoints **reject
or break on authenticated requests** (`captions.go`: "no OAuth token — it can break the
timedtext endpoint"). The user's OAuth (a Data API credential) does not authenticate InnerTube
at all, and the official `captions.download` is owner-only (403 on third-party). So auth cannot
help here and can actively hurt. No probe needed — building one would only re-confirm the ADR.
**Reversibility.** Context-marker + a yield loop in the gate; removing the marker collapses to
the prior equal-share behaviour. No schema or API change.
---
## Rejected alternatives
Approaches considered during the 2026-06-02 planning + grill session and **deliberately not
+5
View File
@@ -113,6 +113,11 @@ type engineProcessor struct {
}
func (p *engineProcessor) ProcessVideo(ctx context.Context, userID, videoID string) error {
// This is the user-initiated (foreground) path — a click on "Summarize",
// "Try now", or a pasted URL. Mark the context so the caption gate gives it
// priority over the background sweep (ADR-026, Pillar A).
ctx = youtube.ForegroundContext(ctx)
row, err := p.store.GetVideoRow(ctx, userID, videoID)
if err != nil {
return fmt.Errorf("load video %q: %w", videoID, err)
+12 -3
View File
@@ -143,8 +143,18 @@ func runScheduler(
return // disabled
}
pass := 0
// Derive the rotation offset from wall-clock, NOT an in-memory counter. A
// counter reset to 0 on every pod restart always hands the lead to the
// first-listed user — so frequent deploys re-starve whoever is last (exactly
// what happened to the first pilot user during a deploy-heavy session). A
// time-based offset advances with real time and is identical across restarts,
// so the lead rotates fairly regardless of how often the pod bounces.
runPass := func() {
pass := int(time.Now().Unix() / int64(interval/time.Second))
runDiscoveryPass(ctx, pass, lister, runUser, log)
}
runPass()
ticker := time.NewTicker(interval)
defer ticker.Stop()
@@ -153,8 +163,7 @@ func runScheduler(
case <-ctx.Done():
return
case <-ticker.C:
pass++
runDiscoveryPass(ctx, pass, lister, runUser, log)
runPass()
}
}
}
+3 -1
View File
@@ -30,7 +30,9 @@ it** — endpoints and aliases drift, and this file is a snapshot (2026-06-06),
- **Summarizer fallback chain (ADR-022).** The primary alias is the *first* of an ordered chain;
on failure or unparseable output the summarizer advances to the next model. All reached through
the same gateway by alias.
- `TAPIR_FALLBACK_MODEL` — local fallback. **Default `koala/phi4-14b`.** Empty disables it.
- `TAPIR_FALLBACK_MODEL` — local fallback. **Default `iguana/gemma4-26b`** — on iguana, NOT
koala, so the fallback does not compete with koala's other GPU loads (and runs from a different
egress IP). Empty disables it.
- `TAPIR_CLOUD_FALLBACK_MODEL` — worst-case EXTERNAL fallback. **Default `berget/mistral-small`.**
**Set this empty (`""`) for any client/NDA deployment** so content never leaves the local
stack — the chain then contains only local endpoints.
+15 -2
View File
@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"os"
"path/filepath"
"testing"
embeddedpostgres "github.com/fergusstrange/embedded-postgres"
@@ -24,11 +25,22 @@ var _ ports.Sink = (*store.Store)(nil)
var dsn string
func TestMain(m *testing.M) {
const port = 54329
// Port + runtime/data dirs are per-process (PID-derived) so two concurrent
// `go test` invocations — e.g. a push-run and a tag-run firing together in CI —
// don't collide on a fixed port or a shared data dir (which silently failed
// both runs). CachePath is shared so the PG archive is downloaded once, not
// per process. Base 54000 keeps this package's range distinct from web's.
port := uint32(54000 + os.Getpid()%1000)
dsn = fmt.Sprintf("postgres://postgres:postgres@localhost:%d/postgres?sslmode=disable", port)
rt := filepath.Join(os.TempDir(), fmt.Sprintf("tapir-epg-store-%d", os.Getpid()))
pg := embeddedpostgres.NewDatabase(
embeddedpostgres.DefaultConfig().Port(port),
embeddedpostgres.DefaultConfig().
Port(port).
RuntimePath(rt).
DataPath(filepath.Join(rt, "data")).
BinariesPath(filepath.Join(rt, "bin")).
CachePath(filepath.Join(os.TempDir(), "tapir-epg-cache")),
)
if err := pg.Start(); err != nil {
fmt.Fprintf(os.Stderr, "embedded-postgres start: %v\n", err)
@@ -40,6 +52,7 @@ func TestMain(m *testing.M) {
if err := pg.Stop(); err != nil {
fmt.Fprintf(os.Stderr, "embedded-postgres stop: %v\n", err)
}
_ = os.RemoveAll(rt)
os.Exit(code)
}
+47
View File
@@ -2,6 +2,7 @@ package youtube
import (
"context"
"sync/atomic"
"time"
"golang.org/x/time/rate"
@@ -29,9 +30,55 @@ func SetFetchRate(interval time.Duration) {
globalFetchGate = rate.NewLimiter(rate.Every(interval), 1)
}
// foregroundPending counts in-flight foreground (user-initiated) caption fetches.
// The background sweep yields the gate while this is non-zero so a human waiting
// on a click gets the next slot — and, on a near-throttled IP, the pre-429 window
// — instead of competing equally with the firehose (ADR-026, Pillar A). Clicks are
// rare and bursty, so background barely notices; the win to the click is large.
var foregroundPending atomic.Int64
// fgCtxKey marks a context as foreground (user-initiated). Unexported; set via
// ForegroundContext and read via isForeground so only this package owns the key.
type fgCtxKey struct{}
// ForegroundContext marks ctx as a user-initiated (foreground) fetch so the gate
// gives it priority. The web "Summarize"/paste/retry path wraps its context with
// this; the background scheduler leaves it unset.
func ForegroundContext(ctx context.Context) context.Context {
return context.WithValue(ctx, fgCtxKey{}, true)
}
func isForeground(ctx context.Context) bool {
v, _ := ctx.Value(fgCtxKey{}).(bool)
return v
}
// fgYieldPoll is how often a background waiter re-checks whether a foreground
// fetch is still pending. Short enough to feel immediate, long enough not to spin.
const fgYieldPoll = 200 * time.Millisecond
// WaitFetchGate blocks until the process-wide gate allows one timedtext fetch,
// respecting ctx cancellation. Called from httpDo before every live outbound
// caption fetch so the scheduler and the click-path share the same egress budget.
//
// Foreground (user-initiated) fetches take priority: they register as pending and
// acquire a token immediately. Background fetches first yield — they wait until no
// foreground fetch is pending — so a live click is never stuck behind the
// background sweep and gets the cleaner slot against the per-IP limit (ADR-026).
func WaitFetchGate(ctx context.Context) error {
if isForeground(ctx) {
foregroundPending.Add(1)
defer foregroundPending.Add(-1)
return globalFetchGate.Wait(ctx)
}
// Background: defer to any pending foreground fetch before taking a token.
for foregroundPending.Load() > 0 {
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(fgYieldPoll):
}
}
return globalFetchGate.Wait(ctx)
}
+34
View File
@@ -82,3 +82,37 @@ func TestSetFetchRateZeroIsUnlimited(t *testing.T) {
require.NoError(t, WaitFetchGate(context.Background()))
}
}
func TestForegroundContextMarker(t *testing.T) {
require.False(t, isForeground(context.Background()), "plain context is background")
require.True(t, isForeground(ForegroundContext(context.Background())), "marked context is foreground")
}
// TestWaitFetchGateForegroundProceedsImmediately: a foreground fetch acquires a
// token without yielding, even when background callers exist.
func TestWaitFetchGateForegroundProceedsImmediately(t *testing.T) {
SetFetchRate(0) // unlimited limiter — isolate the yield logic from pacing
foregroundPending.Store(0)
t.Cleanup(func() { foregroundPending.Store(0) })
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
require.NoError(t, WaitFetchGate(ForegroundContext(ctx)), "foreground proceeds immediately")
}
// TestWaitFetchGateBackgroundYieldsToForeground: while a foreground fetch is
// pending, a background fetch yields (does not take a token) until the foreground
// clears — proven by a background wait timing out against its own deadline, then
// succeeding once the foreground is done.
func TestWaitFetchGateBackgroundYieldsToForeground(t *testing.T) {
SetFetchRate(0)
foregroundPending.Store(1) // simulate a foreground fetch in flight
t.Cleanup(func() { foregroundPending.Store(0) })
ctx, cancel := context.WithTimeout(context.Background(), 250*time.Millisecond)
defer cancel()
require.Error(t, WaitFetchGate(ctx), "background yields (blocks) while foreground is pending")
foregroundPending.Store(0) // foreground done
require.NoError(t, WaitFetchGate(context.Background()), "background proceeds once foreground clears")
}
+5 -2
View File
@@ -33,7 +33,10 @@ type Config struct {
SummarizerModel string
// FallbackModel is the LOCAL fallback alias tried when the primary fails or
// returns unparseable output (ADR-022). Kept local so content stays on the
// homelab stack. Empty disables it. Default a bigger-context local model.
// homelab stack. Default is an IGUANA model (not koala) so the fallback runs
// on a different host than the koala primary — koala carries other loads, and
// a different host also means a different egress IP for the (rare) fallback.
// Empty disables it.
FallbackModel string
// CloudFallbackModel is the worst-case EXTERNAL fallback alias, tried only
// after every local endpoint has failed (ADR-022). For client deployments set
@@ -148,7 +151,7 @@ func (c Config) DexConfigured() bool { return strings.TrimSpace(c.OIDCIssuer) !=
const (
defaultGatewayURL = "http://koala:30401/v1"
defaultSummarizerModel = "koala/phi4-mini"
defaultFallbackModel = "koala/phi4-14b"
defaultFallbackModel = "iguana/gemma4-26b"
defaultCloudFallbackModel = "berget/mistral-small"
defaultSummaryMaxTokens = 1500
defaultMaxTranscriptChars = 18000
+16 -2
View File
@@ -7,6 +7,7 @@ import (
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"testing"
"time"
@@ -26,10 +27,22 @@ import (
var dsn string
func TestMain(m *testing.M) {
const port = 54330 // distinct from the store package's embedded PG (54329)
// Per-process port + dirs so concurrent `go test` runs (e.g. a push-run and a
// tag-run in CI) never collide on a fixed port or shared data dir. Base 55000
// keeps web's range distinct from the store package (54000). Shared CachePath
// downloads the PG archive once.
port := uint32(55000 + os.Getpid()%1000)
dsn = fmt.Sprintf("postgres://postgres:postgres@localhost:%d/postgres?sslmode=disable", port)
pg := embeddedpostgres.NewDatabase(embeddedpostgres.DefaultConfig().Port(port))
rt := filepath.Join(os.TempDir(), fmt.Sprintf("tapir-epg-web-%d", os.Getpid()))
pg := embeddedpostgres.NewDatabase(
embeddedpostgres.DefaultConfig().
Port(port).
RuntimePath(rt).
DataPath(filepath.Join(rt, "data")).
BinariesPath(filepath.Join(rt, "bin")).
CachePath(filepath.Join(os.TempDir(), "tapir-epg-cache")),
)
if err := pg.Start(); err != nil {
fmt.Fprintf(os.Stderr, "embedded-postgres start: %v\n", err)
os.Exit(1)
@@ -38,6 +51,7 @@ func TestMain(m *testing.M) {
if err := pg.Stop(); err != nil {
fmt.Fprintf(os.Stderr, "embedded-postgres stop: %v\n", err)
}
_ = os.RemoveAll(rt)
os.Exit(code)
}