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>
174 lines
5.1 KiB
Go
174 lines
5.1 KiB
Go
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)
|
|
}
|