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>
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
DROP TABLE IF EXISTS sink_deliveries;
|
||||
DROP TABLE IF EXISTS summaries;
|
||||
DROP TABLE IF EXISTS transcripts;
|
||||
DROP TABLE IF EXISTS videos;
|
||||
DROP TABLE IF EXISTS users;
|
||||
@@ -0,0 +1,73 @@
|
||||
-- Tapir initial schema (Stage 0). Per-user isolation: every user-owned table
|
||||
-- carries user_id even though Stage 0 has a single user (docs/data-model.md).
|
||||
-- Scoped to Stage 0 dedup + delivery; Future B/C concerns are out of scope.
|
||||
|
||||
CREATE TABLE users (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
display_name TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
);
|
||||
|
||||
-- One row per (user, video) — per-user isolation, not a global deduped table
|
||||
-- (DECISIONS.md "Rejected alternatives"). subscription_id has no FK at Stage 0:
|
||||
-- the subscriptions table is not part of the store-sink slice.
|
||||
CREATE TABLE videos (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
subscription_id UUID,
|
||||
provider TEXT NOT NULL,
|
||||
provider_video_id TEXT NOT NULL,
|
||||
title TEXT,
|
||||
duration_s INTEGER,
|
||||
published_at TIMESTAMPTZ,
|
||||
url TEXT,
|
||||
seen_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
CONSTRAINT videos_user_provider_video_unique UNIQUE (user_id, provider, provider_video_id)
|
||||
);
|
||||
|
||||
CREATE INDEX idx_videos_user_id ON videos(user_id);
|
||||
|
||||
-- At most one transcript per video. source = 'none' records "checked, no usable
|
||||
-- transcript" so the watcher does not reprocess (ADR-007); content is NULL then.
|
||||
CREATE TABLE transcripts (
|
||||
video_id UUID PRIMARY KEY REFERENCES videos(id) ON DELETE CASCADE,
|
||||
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
source TEXT NOT NULL,
|
||||
language TEXT,
|
||||
content TEXT,
|
||||
resolved_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
);
|
||||
|
||||
CREATE INDEX idx_transcripts_user_id ON transcripts(user_id);
|
||||
|
||||
-- At most one summary per video (UNIQUE on (user_id, video_id)). video_id is not
|
||||
-- FK-constrained to videos at Stage 0: the store sink receives only a Summary
|
||||
-- (ports.Sink.Deliver), so the durable dedup key (user_id, video_id) stands on
|
||||
-- its own; video-row persistence is the engine/source's concern, deferred.
|
||||
CREATE TABLE summaries (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE,
|
||||
video_id UUID NOT NULL,
|
||||
summary TEXT NOT NULL,
|
||||
highlights JSONB NOT NULL DEFAULT '[]'::jsonb,
|
||||
takeaways JSONB NOT NULL DEFAULT '[]'::jsonb,
|
||||
ai_provider TEXT,
|
||||
ai_model TEXT,
|
||||
fallback_used BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
CONSTRAINT summaries_user_video_unique UNIQUE (user_id, video_id)
|
||||
);
|
||||
|
||||
CREATE INDEX idx_summaries_user_id ON summaries(user_id);
|
||||
|
||||
-- One row per (summary, sink) attempt. "Also sent to brain" lives here as a
|
||||
-- delivery row with sink = 'brain'; no brain-specific tables (data-model.md).
|
||||
CREATE TABLE sink_deliveries (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
summary_id UUID NOT NULL REFERENCES summaries(id) ON DELETE CASCADE,
|
||||
sink TEXT NOT NULL,
|
||||
status TEXT NOT NULL,
|
||||
detail TEXT,
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
CONSTRAINT sink_deliveries_summary_sink_unique UNIQUE (summary_id, sink)
|
||||
);
|
||||
@@ -0,0 +1,201 @@
|
||||
// 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)
|
||||
}
|
||||
@@ -0,0 +1,173 @@
|
||||
package store_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
embeddedpostgres "github.com/fergusstrange/embedded-postgres"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"gitea.d-ma.be/mathias/tapir/internal/adapters/store"
|
||||
"gitea.d-ma.be/mathias/tapir/internal/domain"
|
||||
"gitea.d-ma.be/mathias/tapir/internal/ports"
|
||||
)
|
||||
|
||||
// Static check: Store satisfies the Sink port.
|
||||
var _ ports.Sink = (*store.Store)(nil)
|
||||
|
||||
// dsn points at the in-process Postgres started in TestMain. Tests run against
|
||||
// real SQL (constraints, ON CONFLICT, jsonb) — not a mock — without docker or
|
||||
// live-cluster credentials (embedded-postgres downloads its own PG binary).
|
||||
var dsn string
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
const port = 54329
|
||||
dsn = fmt.Sprintf("postgres://postgres:postgres@localhost:%d/postgres?sslmode=disable", port)
|
||||
|
||||
pg := embeddedpostgres.NewDatabase(
|
||||
embeddedpostgres.DefaultConfig().Port(port),
|
||||
)
|
||||
if err := pg.Start(); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "embedded-postgres start: %v\n", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
code := m.Run()
|
||||
|
||||
if err := pg.Stop(); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "embedded-postgres stop: %v\n", err)
|
||||
}
|
||||
os.Exit(code)
|
||||
}
|
||||
|
||||
// uuids — fixed so tests are deterministic. user_id/video_id are UUID columns.
|
||||
const (
|
||||
userA = "11111111-1111-1111-1111-111111111111"
|
||||
userB = "22222222-2222-2222-2222-222222222222"
|
||||
videoX = "aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa"
|
||||
videoY = "bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb"
|
||||
)
|
||||
|
||||
func newStore(t *testing.T) *store.Store {
|
||||
t.Helper()
|
||||
s, err := store.New(context.Background(), dsn)
|
||||
require.NoError(t, err, "migrations must apply clean and the pool must connect")
|
||||
t.Cleanup(s.Close)
|
||||
return s
|
||||
}
|
||||
|
||||
// rawPool is a direct connection for test introspection (counts, truncation),
|
||||
// kept out of the production Store API. Closed via t.Cleanup.
|
||||
func rawPool(t *testing.T) *pgxpool.Pool {
|
||||
t.Helper()
|
||||
p, err := pgxpool.New(context.Background(), dsn)
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(p.Close)
|
||||
return p
|
||||
}
|
||||
|
||||
// resetDB truncates between tests so each starts from a known state. Schema is
|
||||
// shared across the run (migrations are idempotent via New).
|
||||
func resetDB(t *testing.T, p *pgxpool.Pool) {
|
||||
t.Helper()
|
||||
_, err := p.Exec(context.Background(),
|
||||
`TRUNCATE sink_deliveries, summaries, transcripts, videos, users CASCADE`)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
func summary(userID, videoID, text string) domain.Summary {
|
||||
return domain.Summary{
|
||||
UserID: userID,
|
||||
VideoID: videoID,
|
||||
Summary: text,
|
||||
Highlights: []string{"h1", "h2"},
|
||||
Takeaways: []string{"t1"},
|
||||
AIProvider: "local",
|
||||
AIModel: "qwen",
|
||||
}
|
||||
}
|
||||
|
||||
func TestName(t *testing.T) {
|
||||
require.Equal(t, "store", newStore(t).Name())
|
||||
}
|
||||
|
||||
func TestDeliverInsertsSummary(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
s := newStore(t)
|
||||
resetDB(t, rawPool(t))
|
||||
|
||||
require.NoError(t, s.Deliver(ctx, summary(userA, videoX, "first")))
|
||||
|
||||
ok, err := s.HasSummary(ctx, userA, videoX)
|
||||
require.NoError(t, err)
|
||||
require.True(t, ok)
|
||||
}
|
||||
|
||||
func TestDeliverIsIdempotentOnUserVideo(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
s := newStore(t)
|
||||
p := rawPool(t)
|
||||
resetDB(t, p)
|
||||
|
||||
require.NoError(t, s.Deliver(ctx, summary(userA, videoX, "first")))
|
||||
// Re-deliver the same (user, video): must update in place, not duplicate or error.
|
||||
require.NoError(t, s.Deliver(ctx, summary(userA, videoX, "second")))
|
||||
|
||||
var count int
|
||||
require.NoError(t, p.QueryRow(ctx,
|
||||
`SELECT count(*) FROM summaries WHERE user_id = $1 AND video_id = $2`,
|
||||
userA, videoX).Scan(&count))
|
||||
require.Equal(t, 1, count, "second delivery must update, not duplicate")
|
||||
|
||||
var text string
|
||||
require.NoError(t, p.QueryRow(ctx,
|
||||
`SELECT summary FROM summaries WHERE user_id = $1 AND video_id = $2`,
|
||||
userA, videoX).Scan(&text))
|
||||
require.Equal(t, "second", text, "second delivery must overwrite the summary text")
|
||||
|
||||
// Exactly one store-delivery row for the summary (ON CONFLICT update).
|
||||
var deliveries int
|
||||
require.NoError(t, p.QueryRow(ctx,
|
||||
`SELECT count(*) FROM sink_deliveries d
|
||||
JOIN summaries m ON m.id = d.summary_id
|
||||
WHERE m.user_id = $1 AND m.video_id = $2 AND d.sink = 'store'`,
|
||||
userA, videoX).Scan(&deliveries))
|
||||
require.Equal(t, 1, deliveries)
|
||||
}
|
||||
|
||||
func TestSeenVideoIDsReturnsUsersSet(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
s := newStore(t)
|
||||
resetDB(t, rawPool(t))
|
||||
|
||||
require.NoError(t, s.Deliver(ctx, summary(userA, videoX, "x")))
|
||||
require.NoError(t, s.Deliver(ctx, summary(userA, videoY, "y")))
|
||||
|
||||
seen, err := s.SeenVideoIDs(ctx, userA)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, map[string]bool{videoX: true, videoY: true}, seen)
|
||||
}
|
||||
|
||||
func TestPerUserIsolation(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
s := newStore(t)
|
||||
resetDB(t, rawPool(t))
|
||||
|
||||
require.NoError(t, s.Deliver(ctx, summary(userA, videoX, "a-owns-this")))
|
||||
|
||||
// User B must not see user A's video, by either dedup read.
|
||||
seenB, err := s.SeenVideoIDs(ctx, userB)
|
||||
require.NoError(t, err)
|
||||
require.Empty(t, seenB, "user B must not see user A's videos")
|
||||
|
||||
hasB, err := s.HasSummary(ctx, userB, videoX)
|
||||
require.NoError(t, err)
|
||||
require.False(t, hasB, "the same video_id under another user must be invisible")
|
||||
|
||||
hasA, err := s.HasSummary(ctx, userA, videoX)
|
||||
require.NoError(t, err)
|
||||
require.True(t, hasA)
|
||||
}
|
||||
Reference in New Issue
Block a user