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>
202 lines
6.5 KiB
Go
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)
|
|
}
|