Files
tapir/internal/adapters/store/store.go
T
mathiasandClaude Opus 4.8 2695b5d91e
CI / Lint / Test / Vet (push) Successful in 10s
CI / Build & Import (push) Failing after 1s
CI / Mirror to GitHub (push) Has been skipped
feat(adapters): add Postgres store sink with durable dedup
Implements ports.Sink over Postgres (pgx/v5 + pgxpool, DSN from env per
estate convention). This is the primary sink (ADR-003) and the source of
the engine's durable, cross-restart dedup — the in-engine processed map is
process-lifetime only.

- Migrations (golang-migrate, NNN_name.up/down.sql per estate convention,
  applied from an embedded FS on New): users, videos, transcripts,
  summaries, sink_deliveries. Every user-owned table carries user_id
  (Stage-0 per-user isolation promise, data-model.md). summaries has
  UNIQUE(user_id, video_id) — at most one summary per video; highlights /
  takeaways are jsonb.
- Deliver upserts the summary idempotently on (user_id, video_id)
  (ON CONFLICT DO UPDATE) inside one tx with its sink_delivery row. Re-
  delivering the same summary updates in place, never duplicates or errors.
- Dedup reads (store methods, not a new port): HasSummary(ctx,userID,
  videoID) and SeenVideoIDs(ctx,userID) — both user_id-scoped, so one
  user never sees another's videos.

summaries.video_id is intentionally not FK-constrained to videos at Stage 0:
the sink receives only a Summary, so the dedup key stands alone; video-row
persistence is the engine/source's concern, deferred.

Tested against a real in-process Postgres via embedded-postgres (real SQL:
constraints, ON CONFLICT, jsonb, user_id scoping) — no docker, no live
cluster, no creds, fully offline.

Deps: golang-migrate/migrate/v4 and jackc/pgx/v5 (runtime),
fergusstrange/embedded-postgres + stretchr/testify (test-only). go mod tidy
raised the go directive to 1.25.0 (minimum required by the dep graph;
estate elsewhere already runs 1.26.1).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-02 20:04:04 +02:00

202 lines
6.5 KiB
Go

// Package store is the Postgres-backed implementation of ports.Sink: the user's
// own durable store of summaries. It is the primary sink (ADR-003) and the source
// of the engine's durable, cross-restart dedup (the in-engine map is process-
// lifetime only). Connection is pgx/v5 + pgxpool with the DSN from the caller;
// schema is applied via golang-migrate from embedded migrations.
//
// Per-user isolation (docs/data-model.md) is a Stage-0 promise: every row carries
// user_id and every read is scoped by it, even though the demo has one user.
package store
import (
"context"
"database/sql"
"embed"
"encoding/json"
"errors"
"fmt"
"github.com/golang-migrate/migrate/v4"
migratepgx "github.com/golang-migrate/migrate/v4/database/pgx/v5"
"github.com/golang-migrate/migrate/v4/source/iofs"
"github.com/jackc/pgx/v5/pgxpool"
_ "github.com/jackc/pgx/v5/stdlib" // register the "pgx" database/sql driver for migrate
"gitea.d-ma.be/mathias/tapir/internal/domain"
)
//go:embed migrations/*.sql
var migrationsFS embed.FS
// Store persists summaries to Postgres and answers durable dedup queries.
type Store struct {
pool *pgxpool.Pool
}
// New connects a pool to dsn, applies all pending migrations, and verifies the
// connection. The caller owns the lifetime: call Close when done.
func New(ctx context.Context, dsn string) (*Store, error) {
if dsn == "" {
return nil, errors.New("store: empty DSN")
}
if err := Migrate(dsn); err != nil {
return nil, err
}
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
return nil, fmt.Errorf("store: connect pool: %w", err)
}
if err := pool.Ping(ctx); err != nil {
pool.Close()
return nil, fmt.Errorf("store: ping: %w", err)
}
return &Store{pool: pool}, nil
}
// Migrate applies all pending up-migrations against dsn. It opens its own
// short-lived connection (golang-migrate uses database/sql) and closes it before
// returning, so it can run before the pool is created or be invoked standalone.
func Migrate(dsn string) error {
db, err := sql.Open("pgx", dsn)
if err != nil {
return fmt.Errorf("store: open migrate db: %w", err)
}
defer func() { _ = db.Close() }()
drv, err := migratepgx.WithInstance(db, &migratepgx.Config{})
if err != nil {
return fmt.Errorf("store: migrate driver: %w", err)
}
src, err := iofs.New(migrationsFS, "migrations")
if err != nil {
return fmt.Errorf("store: migrate source: %w", err)
}
m, err := migrate.NewWithInstance("iofs", src, "pgx", drv)
if err != nil {
return fmt.Errorf("store: migrator: %w", err)
}
if err := m.Up(); err != nil && !errors.Is(err, migrate.ErrNoChange) {
return fmt.Errorf("store: migrate up: %w", err)
}
return nil
}
// Close releases the connection pool.
func (s *Store) Close() {
s.pool.Close()
}
// Name identifies this sink in delivery records.
func (s *Store) Name() string { return "store" }
// Deliver upserts the summary idempotently on (user_id, video_id) and records the
// store delivery. Re-delivering the same summary updates in place — it never
// errors or duplicates. The whole write is one transaction so a summary and its
// delivery row stay consistent.
func (s *Store) Deliver(ctx context.Context, sum domain.Summary) error {
highlights, err := marshalList(sum.Highlights)
if err != nil {
return fmt.Errorf("store: marshal highlights: %w", err)
}
takeaways, err := marshalList(sum.Takeaways)
if err != nil {
return fmt.Errorf("store: marshal takeaways: %w", err)
}
tx, err := s.pool.Begin(ctx)
if err != nil {
return fmt.Errorf("store: begin: %w", err)
}
defer tx.Rollback(ctx) //nolint:errcheck // no-op after Commit
// Ensure the owning user exists (FK target). The store sink receives only a
// Summary, so a minimal user row is enough at Stage 0.
if _, err := tx.Exec(ctx,
`INSERT INTO users (id) VALUES ($1) ON CONFLICT (id) DO NOTHING`,
sum.UserID); err != nil {
return fmt.Errorf("store: upsert user: %w", err)
}
var summaryID string
if err := tx.QueryRow(ctx,
`INSERT INTO summaries
(user_id, video_id, summary, highlights, takeaways, ai_provider, ai_model, fallback_used)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (user_id, video_id) DO UPDATE SET
summary = EXCLUDED.summary,
highlights = EXCLUDED.highlights,
takeaways = EXCLUDED.takeaways,
ai_provider = EXCLUDED.ai_provider,
ai_model = EXCLUDED.ai_model,
fallback_used = EXCLUDED.fallback_used
RETURNING id`,
sum.UserID, sum.VideoID, sum.Summary, highlights, takeaways,
sum.AIProvider, sum.AIModel, sum.FallbackUsed,
).Scan(&summaryID); err != nil {
return fmt.Errorf("store: upsert summary: %w", err)
}
if _, err := tx.Exec(ctx,
`INSERT INTO sink_deliveries (summary_id, sink, status)
VALUES ($1, 'store', 'delivered')
ON CONFLICT (summary_id, sink) DO UPDATE SET
status = 'delivered',
detail = NULL,
updated_at = NOW()`,
summaryID); err != nil {
return fmt.Errorf("store: record delivery: %w", err)
}
if err := tx.Commit(ctx); err != nil {
return fmt.Errorf("store: commit: %w", err)
}
return nil
}
// HasSummary reports whether a summary already exists for (userID, videoID).
// This is the per-video durable dedup check.
func (s *Store) HasSummary(ctx context.Context, userID, videoID string) (bool, error) {
var exists bool
if err := s.pool.QueryRow(ctx,
`SELECT EXISTS(SELECT 1 FROM summaries WHERE user_id = $1 AND video_id = $2)`,
userID, videoID).Scan(&exists); err != nil {
return false, fmt.Errorf("store: has summary: %w", err)
}
return exists, nil
}
// SeenVideoIDs returns the set of video IDs that already have a summary for the
// user. The watcher uses it to skip re-summarizing across restarts. Scoped by
// user_id, so one user never sees another's videos.
func (s *Store) SeenVideoIDs(ctx context.Context, userID string) (map[string]bool, error) {
rows, err := s.pool.Query(ctx,
`SELECT video_id FROM summaries WHERE user_id = $1`, userID)
if err != nil {
return nil, fmt.Errorf("store: seen video ids: %w", err)
}
defer rows.Close()
seen := make(map[string]bool)
for rows.Next() {
var id string
if err := rows.Scan(&id); err != nil {
return nil, fmt.Errorf("store: scan video id: %w", err)
}
seen[id] = true
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("store: iterate video ids: %w", err)
}
return seen, nil
}
// marshalList renders a string slice as a JSON array, normalising nil to "[]" so
// the jsonb columns never hold SQL/JSON null.
func marshalList(xs []string) ([]byte, error) {
if xs == nil {
xs = []string{}
}
return json.Marshal(xs)
}