feat: add Postgres store with append-only schema and migrations
internal/store connects via pgx and runs golang-migrate migrations embedded in the binary (go:embed), so Deklarix stays a single binary despite the move to Postgres. Schema covers the five MVP tables (submission, asset, extraction, finding, evidence_package, participant). extraction, finding and evidence_package are append-only by design: a Postgres trigger rejects UPDATE/DELETE outright, since a corrigible evidence archive isn't an evidence archive. Corrections to a finding are new rows whose supersedes column points at the row they replace (set at INSERT time on the new row, since the trigger blocks UPDATE on the old one) — "currently valid" findings are the ones no other row supersedes. scripts/test.sh now spins up a disposable Postgres container so the store's integration tests (including the append-only guarantee) actually run on every test.sh/release.sh invocation instead of silently skipping for lack of DATABASE_URL. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
48
internal/store/migrate.go
Normal file
48
internal/store/migrate.go
Normal file
@@ -0,0 +1,48 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"embed"
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
"github.com/golang-migrate/migrate/v4"
|
||||
pgxmigrate "github.com/golang-migrate/migrate/v4/database/pgx/v5"
|
||||
"github.com/golang-migrate/migrate/v4/source/iofs"
|
||||
_ "github.com/jackc/pgx/v5/stdlib"
|
||||
)
|
||||
|
||||
//go:embed migrations/*.sql
|
||||
var migrationsFS embed.FS
|
||||
|
||||
// Migrate wendet alle ausstehenden Migrationen aus migrations/ an.
|
||||
// Die Migrationen sind im Binary eingebettet (go:embed), damit Deklarix
|
||||
// weiterhin als einzelnes Binary lauffähig bleibt.
|
||||
func Migrate(databaseURL string) error {
|
||||
db, err := sql.Open("pgx", databaseURL)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: open db for migration: %w", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
driver, err := pgxmigrate.WithInstance(db, &pgxmigrate.Config{})
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: migration driver: %w", err)
|
||||
}
|
||||
|
||||
source, err := iofs.New(migrationsFS, "migrations")
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: migration source: %w", err)
|
||||
}
|
||||
|
||||
m, err := migrate.NewWithInstance("iofs", source, "pgx5", driver)
|
||||
if err != nil {
|
||||
return fmt.Errorf("store: migrate init: %w", err)
|
||||
}
|
||||
defer m.Close()
|
||||
|
||||
if err := m.Up(); err != nil && !errors.Is(err, migrate.ErrNoChange) {
|
||||
return fmt.Errorf("store: migrate up: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
11
internal/store/migrations/0001_init.down.sql
Normal file
11
internal/store/migrations/0001_init.down.sql
Normal file
@@ -0,0 +1,11 @@
|
||||
DROP TRIGGER IF EXISTS evidence_package_append_only ON evidence_package;
|
||||
DROP TRIGGER IF EXISTS finding_append_only ON finding;
|
||||
DROP TRIGGER IF EXISTS extraction_append_only ON extraction;
|
||||
DROP FUNCTION IF EXISTS forbid_update_delete();
|
||||
|
||||
DROP TABLE IF EXISTS participant;
|
||||
DROP TABLE IF EXISTS evidence_package;
|
||||
DROP TABLE IF EXISTS finding;
|
||||
DROP TABLE IF EXISTS extraction;
|
||||
DROP TABLE IF EXISTS asset;
|
||||
DROP TABLE IF EXISTS submission;
|
||||
89
internal/store/migrations/0001_init.up.sql
Normal file
89
internal/store/migrations/0001_init.up.sql
Normal file
@@ -0,0 +1,89 @@
|
||||
CREATE EXTENSION IF NOT EXISTS pgcrypto;
|
||||
|
||||
CREATE TABLE submission (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
platform TEXT NOT NULL CHECK (platform IN ('instagram', 'tiktok', 'youtube', 'linkedin')),
|
||||
post_type TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'draft' CHECK (status IN ('draft', 'checked', 'published', 'archived')),
|
||||
caption TEXT,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE TABLE asset (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
submission_id UUID NOT NULL REFERENCES submission (id),
|
||||
kind TEXT NOT NULL CHECK (kind IN ('image', 'video', 'file')),
|
||||
path TEXT NOT NULL,
|
||||
sha256 TEXT NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
-- Stufe 1: striktes JSON aus der Claude-Extraktion. Append-only, siehe Trigger unten.
|
||||
CREATE TABLE extraction (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
submission_id UUID NOT NULL REFERENCES submission (id),
|
||||
payload JSONB NOT NULL,
|
||||
model_version TEXT NOT NULL,
|
||||
prompt_version TEXT NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
-- Stufe 2: Ergebnis pro Regel. Append-only, Korrektur = neue Zeile.
|
||||
-- supersedes zeigt auf die alte Zeile, die diese Zeile ersetzt (gesetzt
|
||||
-- beim INSERT der Korrektur, nie per UPDATE — der Trigger würde das
|
||||
-- verbieten). "Aktuell gültig" = Zeilen, auf die kein supersedes zeigt,
|
||||
-- siehe Index unten für den Anti-Join.
|
||||
CREATE TABLE finding (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
submission_id UUID NOT NULL REFERENCES submission (id),
|
||||
extraction_id UUID REFERENCES extraction (id),
|
||||
rule_id TEXT NOT NULL,
|
||||
rule_version INTEGER NOT NULL,
|
||||
severity TEXT NOT NULL CHECK (severity IN ('niedrig', 'mittel', 'hoch')),
|
||||
message TEXT NOT NULL,
|
||||
supersedes UUID REFERENCES finding (id),
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE INDEX finding_supersedes_idx ON finding (supersedes) WHERE supersedes IS NOT NULL;
|
||||
|
||||
-- Dossier + Beweiskette (Hash, RFC-3161-Token). Append-only.
|
||||
CREATE TABLE evidence_package (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
submission_id UUID NOT NULL REFERENCES submission (id),
|
||||
dossier_path TEXT NOT NULL,
|
||||
sha256 TEXT NOT NULL,
|
||||
timestamp_token BYTEA NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE TABLE participant (
|
||||
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||
submission_id UUID NOT NULL REFERENCES submission (id),
|
||||
role TEXT NOT NULL CHECK (role IN ('creator', 'agentur', 'marke', 'kanzlei')),
|
||||
name TEXT NOT NULL,
|
||||
vorgegeben BOOLEAN NOT NULL DEFAULT false,
|
||||
freigegeben BOOLEAN NOT NULL DEFAULT false,
|
||||
approved_at TIMESTAMPTZ,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
-- Ein Beweisarchiv, in dem man Zeilen ändern kann, ist kein Beweisarchiv.
|
||||
CREATE FUNCTION forbid_update_delete() RETURNS TRIGGER AS $$
|
||||
BEGIN
|
||||
RAISE EXCEPTION 'append-only table: % on % is not allowed', TG_OP, TG_TABLE_NAME;
|
||||
END;
|
||||
$$ LANGUAGE plpgsql;
|
||||
|
||||
CREATE TRIGGER extraction_append_only
|
||||
BEFORE UPDATE OR DELETE ON extraction
|
||||
FOR EACH ROW EXECUTE FUNCTION forbid_update_delete();
|
||||
|
||||
CREATE TRIGGER finding_append_only
|
||||
BEFORE UPDATE OR DELETE ON finding
|
||||
FOR EACH ROW EXECUTE FUNCTION forbid_update_delete();
|
||||
|
||||
CREATE TRIGGER evidence_package_append_only
|
||||
BEFORE UPDATE OR DELETE ON evidence_package
|
||||
FOR EACH ROW EXECUTE FUNCTION forbid_update_delete();
|
||||
31
internal/store/store.go
Normal file
31
internal/store/store.go
Normal file
@@ -0,0 +1,31 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
// Store hält den Verbindungspool zur Postgres-Datenbank.
|
||||
type Store struct {
|
||||
Pool *pgxpool.Pool
|
||||
}
|
||||
|
||||
// Open baut den Verbindungspool auf und prüft ihn mit einem Ping.
|
||||
func Open(ctx context.Context, databaseURL string) (*Store, error) {
|
||||
pool, err := pgxpool.New(ctx, databaseURL)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("store: open 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
|
||||
}
|
||||
|
||||
// Close gibt den Verbindungspool frei.
|
||||
func (s *Store) Close() {
|
||||
s.Pool.Close()
|
||||
}
|
||||
86
internal/store/store_test.go
Normal file
86
internal/store/store_test.go
Normal file
@@ -0,0 +1,86 @@
|
||||
package store_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/netcell-it/deklarix/internal/store"
|
||||
)
|
||||
|
||||
// Diese Tests brauchen eine laufende Postgres-Instanz und werden ohne
|
||||
// DATABASE_URL übersprungen, statt eine Verbindung vorzutäuschen.
|
||||
func testDatabaseURL(t *testing.T) string {
|
||||
t.Helper()
|
||||
url := os.Getenv("DATABASE_URL")
|
||||
if url == "" {
|
||||
t.Skip("DATABASE_URL nicht gesetzt, überspringe Store-Integrationstest")
|
||||
}
|
||||
return url
|
||||
}
|
||||
|
||||
func TestMigrateAndOpen(t *testing.T) {
|
||||
url := testDatabaseURL(t)
|
||||
|
||||
if err := store.Migrate(url); err != nil {
|
||||
t.Fatalf("Migrate: %v", err)
|
||||
}
|
||||
|
||||
s, err := store.Open(context.Background(), url)
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
var tableCount int
|
||||
err = s.Pool.QueryRow(context.Background(), `
|
||||
SELECT count(*) FROM information_schema.tables
|
||||
WHERE table_schema = 'public' AND table_name = ANY($1)
|
||||
`, []string{"submission", "asset", "extraction", "finding", "evidence_package", "participant"}).Scan(&tableCount)
|
||||
if err != nil {
|
||||
t.Fatalf("query tables: %v", err)
|
||||
}
|
||||
if tableCount != 6 {
|
||||
t.Fatalf("expected 6 tables, got %d", tableCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFindingIsAppendOnly(t *testing.T) {
|
||||
url := testDatabaseURL(t)
|
||||
|
||||
if err := store.Migrate(url); err != nil {
|
||||
t.Fatalf("Migrate: %v", err)
|
||||
}
|
||||
|
||||
s, err := store.Open(context.Background(), url)
|
||||
if err != nil {
|
||||
t.Fatalf("Open: %v", err)
|
||||
}
|
||||
defer s.Close()
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
var submissionID string
|
||||
err = s.Pool.QueryRow(ctx, `
|
||||
INSERT INTO submission (platform, post_type) VALUES ('instagram', 'reel')
|
||||
RETURNING id
|
||||
`).Scan(&submissionID)
|
||||
if err != nil {
|
||||
t.Fatalf("insert submission: %v", err)
|
||||
}
|
||||
|
||||
var findingID string
|
||||
err = s.Pool.QueryRow(ctx, `
|
||||
INSERT INTO finding (submission_id, rule_id, rule_version, severity, message)
|
||||
VALUES ($1, 'WK-004', 3, 'hoch', 'Testfeststellung')
|
||||
RETURNING id
|
||||
`, submissionID).Scan(&findingID)
|
||||
if err != nil {
|
||||
t.Fatalf("insert finding: %v", err)
|
||||
}
|
||||
|
||||
_, err = s.Pool.Exec(ctx, `UPDATE finding SET message = 'geändert' WHERE id = $1`, findingID)
|
||||
if err == nil {
|
||||
t.Fatal("expected UPDATE on finding to be rejected, but it succeeded")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user