Files
mathias 40808f2d4b
CI / Lint / Test / Vet (push) Successful in 10s
CI / Build & Import (push) Successful in 10s
docs(onboard): BDD scenarios + architecture for the ADR-028 burst
- connect_account.feature: refine the burst scenario to "best recent" (likely-good
  selection, skips too-short/too-long) and add a scenario for the stronger model.
- scenario_coverage: remap to TestOnboardBurstVideoIDs + add the model scenario
  (the old NewestUnsummarizedVideoIDs test was removed).
- architecture.md: new "Connect-time onboarding burst" subsection — junk-avoiding
  selection over persisted duration, the burst-only stronger-model chain, and why
  has-captions/cached-first are not selection signals.
2026-06-11 19:54:29 +02:00

17 KiB

Tapir — Architecture (C4 + sequences)

Diagrams describe current assumptions, not a frozen design. They are diffable Mermaid so they live in version control and render in Gitea. When a decision here changes, record it in DECISIONS.md and update the relevant diagram in the same commit.

Altitude: C4 Level 1 (Context) and Level 2 (Container), plus sequence diagrams for the two behaviours that carry the most design risk. Clean Architecture layering is described after the diagrams.


C4 L1 — System Context

Who and what Tapir talks to. (Brain and external AI are dashed — both optional.)

graph TB
    user["User<br/>(maintainer; later, trusted users)"]
    tapir["Tapir<br/>watches subscriptions,<br/>summarizes new videos"]
    yt["YouTube<br/>(Data API + captions)"]
    vimeo["Vimeo<br/>(API + text tracks)"]
    local["Local AI stack<br/>(LiteLLM / piguard alias)"]
    byo["User's BYO AI<br/>(Claude / OpenAI / Gemini)"]
    brain["brain<br/>(brain-mcp)"]

    user -->|connects accounts,<br/>reads summaries| tapir
    tapir -->|subscriptions,<br/>new-video events,<br/>captions| yt
    tapir -->|subscriptions,<br/>text tracks| vimeo
    tapir -->|summarize<br/>PRIMARY| local
    tapir -.->|summarize<br/>FALLBACK, opt-in| byo
    tapir -.->|optional sink:<br/>brain_ingest| brain

    classDef opt stroke-dasharray: 5 5;
    class byo,brain opt;

C4 L2 — Containers

Inside Tapir. The engine is provider- and sink-agnostic; everything external is an adapter behind an interface (Clean Architecture ports & adapters).

graph TB
    subgraph tapir["Tapir (Go)"]
        http["tapir serve<br/>(HTMX+Templ web surface:<br/>read summaries, connect,<br/>account, summarize)"]
        watcher["Watcher<br/>detects new videos<br/>(WebSub + poll)"]
        engine["Summarization engine<br/>(use-case core)"]
        resolver["Transcript resolver<br/>(captions-first)"]
        router["AI router<br/>(llm.Router:<br/>Primary -> Fallback)"]
        store[("User store<br/>(Postgres,<br/>per-user isolated)")]

        subgraph sinks["Sink adapters (Sink interface)"]
            sink_store["Store sink<br/>(primary)"]
            sink_brain["Brain sink<br/>(HTTP brain-mcp)"]
        end

        subgraph providers["Provider adapters (VideoSource interface)"]
            p_yt["YouTube adapter"]
            p_vimeo["Vimeo adapter"]
        end
    end

    http --> store
    watcher --> p_yt
    watcher --> p_vimeo
    watcher -->|NewVideo event| engine
    engine --> resolver
    resolver --> p_yt
    resolver --> p_vimeo
    engine --> router
    engine --> sink_store
    engine -.-> sink_brain
    sink_store --> store

    classDef opt stroke-dasharray: 5 5;
    class sink_brain opt;

Interfaces that keep the engine pure (the ports):

  • VideoSource — list subscriptions, detect new videos, fetch transcript. Implemented by YouTube and Vimeo adapters.
  • Summarizer — turn a transcript + context into highlights/takeaways. Backed by the llm.Router (Primary local, Fallback BYO).
  • Sink — deliver a summary. Implemented by the store sink (primary) and brain sink (optional). New sinks add an implementation, nothing else.
  • SecretStore — fetch/store per-user OAuth tokens and BYO keys. Backed by ESO/1Password.

The engine depends only on these interfaces, never on YouTube, brain, or a concrete model. This is what makes "standalone vs homelab" a configuration of which adapters are wired, not two codebases (ADR-003).


Web surface — tapir serve (Stage 1, ADR-011 → ADR-012)

A later transport added over the unchanged engine/ports/sinks core (ADR-003): tapir serve is an HTMX+Templ reader/writer (internal/web) over the existing store. It added no business logic to the engine — it reads the store and, for one action, kicks the existing engine. ADR-011 shipped it single-user; ADR-012 opened multi-user with DB-enforced (RLS) isolation.

graph TB
    browser["Browser<br/>(Dex-authenticated user)"]
    subgraph web["internal/web (tapir serve)"]
        oidc["oidc<br/>Dex OIDC session<br/>(authenticate-only)"]
        gate["registration gate<br/>new subject -> /register"]
        pages["summary list + detail<br/>(read) + actions"]
        connect["/oauth/youtube/callback<br/>per-user token connect"]
        account["account<br/>(disconnect, delete)"]
        summarize["Summarize button<br/>-> background goroutine"]
    end
    store[("store<br/>(Postgres, RLS per user)")]
    engine["Summarization engine<br/>(unchanged core)"]
    secrets["SecretStore<br/>(per-user token refs)"]

    browser --> oidc
    oidc --> gate
    gate --> pages
    pages --> store
    connect --> secrets
    connect --> store
    account --> store
    account --> secrets
    summarize -->|background| engine
    summarize -->|HTMX status poll| store
    engine --> store
  • Dex OIDC session layer (internal/web/oidc) — authenticate-only (ADR-012). It proves who; authorization/isolation is the DB's job (RLS), not the session's.
  • Registration gate — a Dex subject with no users row is routed to /register, which creates the users row + the user_identities mapping (migration 004). Returning subjects pass straight through.
  • Web-initiated YouTube connect/oauth/youtube/connect/oauth/youtube/callback persists a per-user refresh-token ref (youtube/<userID>/refresh_token) via SecretStore and a video_connections row (ADR-006, migration 005). Distinct from the CLI tapir auth.
  • Account management/account offers disconnect and delete account. Delete removes only Tapir-side state (cascade across the user's tables + secret refs); the shared Dex identity is left intact (ADR-013).
  • Immediate summarization — the web "Summarize" button (POST /v/{id}/summarize) fires the engine in a background goroutine inside serve; the page HTMX-polls /v/{id}/status, showing a Charmbracelet spinner while in-flight (and an honest "queued/waiting" state under rate-limiting — ADR-014).
  • Summarization modeusers.auto_summarize (migration 006). Default is true for new users (migration 011, ADR-018); all existing rows were back-filled via migration 012. Auto: new videos published within the recency window (TAPIR_AUTO_SUMMARIZE_WINDOW, default ~7d, ADR-020) are summarized automatically; older videos are discovered and listed but wait for an explicit "Summarize". Manual: new videos appear unsummarized; the button sets videos.summarize_requested, which the next tapir run processes and clears. A manual request bypasses the recency bound. Both the click path and the batch tapir run drive the same unchanged engine.
  • List surface (ADR-020) — the list reads ListVideos ordered summarized-first, then published_at DESC NULLS LAST. The web layer collapses the noise so summaries are not buried: un-summarized videos older than the recency window fold into one "Show N older videos" disclosure, and caption-less videos collapse to a single count line. Copy surfaces scarcity honestly (queue counts, gradual-fill note) — it never implies the feed is fuller than it is.

The engine, ports, and sink adapters are untouched by all of the above — the web surface only reads the store and triggers the existing engine. Adding it changed wiring, not the core (ADR-003).


In-process scheduler (ADR-018)

cmdServe launches a background goroutine when TAPIR_DISCOVERY_INTERVAL > 0. On each tick it calls store.ListAllUsers (un-RLS'd admin query), builds a per-user runner.Runner, and calls RunOnce for each registered user in sequence.

sequenceDiagram
    participant S as Scheduler goroutine
    participant DB as Postgres (RLS)
    participant YT as YouTube timedtext
    participant LLM as LiteLLM gateway

    loop every TAPIR_DISCOVERY_INTERVAL
        S->>DB: ListAllUsers() [un-RLS'd]
        loop per user
            S->>DB: GetAutoSummarize(userID)
            S->>YT: ListSubscriptions + NewVideos
            Note over S,DB: auto: skip videos published before<br/>TAPIR_AUTO_SUMMARIZE_WINDOW (ADR-020);<br/>older ones listed, await manual request
            Note over S,YT: WaitFetchGate(ctx) throttles<br/>all fetches to TAPIR_FETCH_RATE
            alt transcript available
                S->>LLM: Summarize
                S->>DB: Deliver(summary)
            else 429
                S->>DB: SetTranscriptStatus(rate_limited)
            end
        end
    end

Single-replica constraint (load-bearing). The scheduler lives in the web process; replicas: 1 in the k3s deployment manifest is not cosmetic — running tapir serve at

1 replica makes every replica run the full discovery loop, causing every registered user to be fetched in parallel from the same egress IP (429s + duplicate work). Do not scale serve past 1 replica without first moving discovery to a k8s CronJob or adding leader election. The process logs a Warn at startup when scheduled discovery is enabled as a reminder.


Process-wide timedtext rate gate

internal/adapters/youtube/gate.go (ADR-014 item 2): a single rate.Limiter (golang.org/x/time/rate) shared across all Adapter instances. Every httpDo call for a caption fetch passes through WaitFetchGate(ctx) before hitting YouTube. This serialises the scheduler loop AND the web click-path through the same per-egress-IP budget. Configured via TAPIR_FETCH_RATE (Go duration, default 2s). Setting it to 0 disables the gate (dev/tests only).

This is the precondition that makes scheduled auto-summarize safe: without the gate, a multi-user scheduler pass could fire many concurrent timedtext requests from the same IP within seconds, triggering 429s for all users.

Two-path summarisation model

Both paths share globalFetchGate — rate limiting is respected in both, not routed around.

Path Trigger Order Rationale
Foreground User clicks "Summarize" on any non-summarized card (POST /v/{id}/retry-now for rate-limited; POST /v/{id}/summarize for pending) Single chosen video On-demand value: user picks a specific video to read now — bypasses the recency bound
Background batch Scheduled discovery pass every TAPIR_DISCOVERY_INTERVAL Newest-first across all channels (see below), bounded to the recency window (ADR-020) Onboarding prioritisation within bounded load: recent videos auto-fill; the older back-catalogue stays on-demand

The rationale for both paths is onboarding prioritisation under an honest, bounded load — a new user gets summaries of their most recent videos automatically, while the older back-catalogue is listed but summarised only on demand, so it never re-drives the shared rate gate every cycle.

Newest-first batch ordering (ADR-018)

Within each scheduled pass, RunOnce uses a three-phase structure:

  1. Discover + persist: walk all channels, UpsertVideo every candidate (so it appears in the list), apply pre-filters (seen/manual/backoff/recency), collect surviving candidates. The recency pre-filter (ADR-020) drops auto-mode videos published before now - TAPIR_AUTO_SUMMARIZE_WINDOW unless they are explicitly requested; an undated video is never aged out. They remain persisted/listed — only auto-summarisation is skipped.
  2. Sort: order candidates published_at DESC, NULLS LAST, discovery_pos ASC. Videos with no publish date (schema 001: nullable) sort after all dated content. The sort is in-memory (slices.SortStableFunc) — at current scale this is fine.
  3. Process: feed candidates to the engine in sorted order through globalFetchGate.

Before (per-channel inline): [chanA-old, chanA-mid, chanB-new, chanB-null] After (newest-first): [chanB-new, chanA-mid, chanA-old, chanB-null]

The set of processed videos now also excludes auto-mode back-catalogue beyond the recency window (those stay listed, summarised on demand); within the processed set, only order changes.

Connect-time onboarding burst (ADR-018 → ADR-028)

On a successful YouTube connect, ConnectHandler enqueues a connect-triggered discovery pass; the discoveryTrigger runs that pass and then fires the onboarding burst — a third entry path that summarises up to TAPIR_ONBOARD_SUMMARIZE_COUNT (default 3, hard-capped) of the new user's videos so the first session is not empty. The burst still flows through globalFetchGate (it is not a throughput change); ADR-028 sharpened which videos and which model:

  • Selection is OnboardBurstVideoIDs, not pure newest-first. It keeps newest-first order but excludes a video whose known duration is outside [TAPIR_MIN_VIDEO_SECONDS, TAPIR_ONBOARD_MAX_VIDEO_SECONDS] (drops Shorts and multi-hour livestream VODs). An unknown (NULL) duration is degrade-open — kept, but ranked after known-good rows. The connect-triggered discovery pass runs before the burst, and ADR-023's videos.list enrichment now persists duration_s (instead of discarding it after the Shorts filter), so a fresh user's candidates carry a duration in time for selection.
  • Model: the burst runs through a dedicated summarizer chain led by TAPIR_ONBOARD_SUMMARIZER_MODEL (default iguana/gemma4-26b, the stronger local model), with the standard ADR-022 chain following as fallback. This is a wiring choice — a second engineProcessor over the same store / transcript cache / sink; the engine and ports are unchanged. Empty / equal-to-primary collapses it back onto the shared processor.

has-captions is deliberately not a selection signal — it is only knowable after a gate fetch (or a ~0-probability cache hit at pilot scale), so the burst can avoid known-junk but cannot promise captions. Cached-transcript-first selection was investigated and rejected (ADR-028: ~3% cross-user overlap).


Sequence — core use case: new video summarized

sequenceDiagram
    participant Src as VideoSource<br/>(YouTube/Vimeo)
    participant W as Watcher
    participant E as Engine
    participant R as Transcript resolver
    participant AI as AI router
    participant S as Sink(s)

    Src->>W: new upload (WebSub push / poll)
    W->>E: NewVideo{user, channel, videoID}
    E->>R: resolve transcript(videoID)
    R->>Src: fetch captions
    alt captions available
        Src-->>R: captions text
        R-->>E: Transcript{source: captions}
        E->>AI: summarize(transcript, userContext)
        AI-->>E: Summary{highlights, takeaways}
        E->>S: deliver(summary)
        S-->>E: ok
    else no transcript
        Src-->>R: none
        R-->>E: NoTranscript
        E->>S: deliver(skipped: no transcript)
    end

Sequence — AI routing (local-first, BYO fallback)

sequenceDiagram
    participant E as Engine
    participant R as llm.Router
    participant P as Primary<br/>(local: LiteLLM/piguard)
    participant F as Fallback<br/>(user BYO key)

    E->>R: summarize(transcript)
    R->>P: complete(prompt)
    alt local succeeds
        P-->>R: summary
        R-->>E: summary (fallback_used = false)
    else local fails (error/timeout/unavailable)
        P-->>R: error
        alt user has BYO configured
            R->>F: complete(prompt)
            F-->>R: summary
            R-->>E: summary (fallback_used = true)
        else no BYO
            R-->>E: error (queued for retry)
        end
    end

"Reliably" (VISION / the BYO trigger) is operationalized as: Primary returned without error within timeout. Richer quality scoring can layer on later without changing the interface. fallback_used is recorded per summary so "is the local stack good enough?" becomes a query, not a guess.


Clean Architecture layering

graph LR
    subgraph domain["Domain (entities)"]
        d["User, Subscription, Video,<br/>Transcript, Summary"]
    end
    subgraph usecase["Use cases (engine)"]
        u["SummarizeNewVideo,<br/>ConnectAccount,<br/>RouteAI"]
    end
    subgraph ports["Ports (interfaces)"]
        po["VideoSource, Summarizer,<br/>Sink, SecretStore"]
    end
    subgraph adapters["Adapters (infra)"]
        a["YouTube, Vimeo, llm.Router,<br/>Store sink, Brain sink,<br/>ESO secret store, HTTP, Postgres"]
    end

    a --> po
    po --> u
    u --> d

Dependencies point inward only. Domain knows nothing of YouTube, brain, Postgres, or any model. Adapters are swappable; tests target the use-case core through fake adapters (see the Gherkin features in docs/use-cases/).


Out of scope in these diagrams (deferred per ADRs)

  • Audio-download + speech-to-text resolver (ADR-007) — would be an additional VideoSource fallback path, drawn when built.
  • Per-user isolation is live, not deferred: Postgres RLS FORCEd on every user-owned table (ADR-012, migration 003), realising ADR-002's per-tenant intent at the DB layer. The coarser multi-tenant primitives (per-namespace NetworkPolicy, Kyverno, tenant label) remain a Stage-2 hardening item, not exercised yet.
  • Public SaaS surface (sign-up, billing) — Future C, not built (ADR-008).