diff --git a/.agents/skills/hyperindex/SKILL.md b/.agents/skills/hyperindex/SKILL.md index feca7c43..c2721a40 100644 --- a/.agents/skills/hyperindex/SKILL.md +++ b/.agents/skills/hyperindex/SKILL.md @@ -233,6 +233,22 @@ Variables: Use typed collection queries when the caller needs collection-specific filters, sorting, exact totals, or typed fields. Use `recordTimeline` when the primary requirement is one stable newest-first page across selected collections. +## Record version history + +For collections the deployment tracks, `recordHistory(uri: "at://…")` returns the observed versions of one record, oldest first, 100 per page by default (`first` up to 500; pass the last `id` as `after` for more). `value` holds the record body, or null for a delete. Use it to show how a record changed, for example who renamed an observation: + +```graphql +{ + recordHistory(uri: "at://did:plc:example/app.gainforest.dwc.occurrence/3kabc") { + action + observedAt + value + } +} +``` + +The first entry may be a `baseline`: the version indexed when history was switched on. Records that were never tracked return an empty list. + ## External labeler filtering Use external label filtering only after confirming that the target endpoint exposes external label support. diff --git a/.agents/skills/hyperindex/references/schema-reference.md b/.agents/skills/hyperindex/references/schema-reference.md index 84a1e067..13e762e0 100644 --- a/.agents/skills/hyperindex/references/schema-reference.md +++ b/.agents/skills/hyperindex/references/schema-reference.md @@ -61,6 +61,7 @@ Use `https://api.indexer.hypercerts.dev/stats` for production and `https://dev.a | `orgHypercertsWorkscopeTag` | after: `String`, before: `String`, first: `Int`, last: `Int`, sortBy: `OrgHypercertsWorkscopeTagSortField`, sortDirection: `SortDirection`, where: `OrgHypercertsWorkscopeTagWhereInput` | Query org.hypercerts.workscope.tag records | | `orgHypercertsWorkscopeTagByUri` | uri: `String!` | Get a single org.hypercerts.workscope.tag by AT-URI | | `records` | after: `String`, before: `String`, collection: `String!`, first: `Int`, last: `Int` | Query records from any collection (useful for collections without lexicon schemas) | +| `recordHistory` | after: `String`, first: `Int`, uri: `String!` | Versions the indexer has observed for one record, oldest first, one page at a time. | | `recordTimeline` | after: `String`, first: `Int`, where: `RecordTimelineWhereInput!` | Query a newest-first page of current records across selected collections, optionally filtered by author DIDs. | | `externalLabels` | activeOnly: `Boolean`, sources: `[String!]`, subjects: `[String!]!`, values: `[String!]` | Query locally ingested external ATProto labels by DID or AT-URI subject. | | `search` | after: `String`, collection: `String`, first: `Int`, query: `String!` | Search records by text content | @@ -329,6 +330,24 @@ External label predicates bound by the containing filter field. | `has` | `ExternalLabelPredicateInput` | Keep records whose bound label subject has a matching external label. | | `none` | `ExternalLabelPredicateInput` | Keep records whose bound label subject does not have a matching external label. | +## Record version history + +`recordHistory(uri: String!, first: Int = 100, after: String)` returns `[RecordVersion!]!`, oldest first. `first` must be 1 to 500; for the next page pass the last version's `id` as `after`. Versions are recorded for collections the deployment tracks (`RECORD_HISTORY_COLLECTIONS`; production tracks `app.gainforest.dwc.occurrence`); records that were never tracked return an empty list, and history recorded before a collection was removed stays queryable. A `baseline` version is what was indexed when history was switched on; earlier edits are not available. + +### `RecordVersion` + +| Field | Type | Description | +| --- | --- | --- | +| `action` | `String!` | `baseline`, `create`, `update`, or `delete`. | +| `cid` | `String!` | CID of this version (for deletes, the version that was deleted). | +| `collection` | `String!` | Collection NSID. | +| `did` | `String!` | Repository DID. | +| `id` | `String!` | Monotonic version id; later versions have larger ids. | +| `live` | `Boolean` | True when seen on the live stream, false for a resync delivery, null for baseline rows. | +| `observedAt` | `String!` | When the indexer observed this version (RFC 3339). | +| `uri` | `String!` | Record AT-URI. | +| `value` | `JSON` | The record body at this version; null for deletes. | + ## External label support Production exposes external ATProto labels through the root `externalLabels` query, each generated record type's virtual `externalLabels` field, `where.externalLabels.has` / `where.externalLabels.none` predicates for record AT-URI labels, and `where.authorLabels.has` / `where.authorLabels.none` predicates for author DID labels on typed list queries. diff --git a/.changes/unreleased/add-record-version-history.yaml b/.changes/unreleased/add-record-version-history.yaml new file mode 100644 index 00000000..fe23a578 --- /dev/null +++ b/.changes/unreleased/add-record-version-history.yaml @@ -0,0 +1,5 @@ +kind: added +body: Add opt-in record version history. Collections listed in RECORD_HISTORY_COLLECTIONS (exact NSIDs or prefix.* patterns, Tap mode) keep every observed version and delete in a new record_version table, seeded on startup with a baseline row per existing record, and the new root recordHistory(uri, first, after) GraphQL query pages through them oldest first. History is off by default; migration 015 adds the table. +time: 2026-09-23T03:10:00.000000+02:00 +custom: + Affects: operator diff --git a/AGENTS.md b/AGENTS.md index 16e7f447..d469f82b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -91,6 +91,7 @@ Run verification based on what changed. - `SECRET_KEY_BASE` must be at least 64 characters. - `TAP_ENABLED=true` switches record ingestion to Tap mode. After Tap becomes healthy, Hyperindex runs one bounded background pass that fills missing `actor.handle` values from Tap's local `/info/:did` metadata without delaying API startup. - `LABELER_SUBSCRIBE_ENABLED=true` with `LABELER_SUBSCRIBE_URLS` starts optional external `com.atproto.label.subscribeLabels` ingestion. +- `RECORD_HISTORY_COLLECTIONS` (Tap mode) keeps every version of matching records in `record_version`, exposed by the root `recordHistory(uri)` query. Startup seeds a `baseline` row for existing records that have no history yet; the version is written before the current-state upsert so Tap redelivery keeps both consistent. - Migrations run automatically on startup. - Be careful with `ALLOWED_ORIGINS`: current code allows all origins when unset, even if older prose suggests stricter defaults. diff --git a/README.md b/README.md index 82249f49..94770506 100644 --- a/README.md +++ b/README.md @@ -116,6 +116,9 @@ TAP_SIGNAL_COLLECTION=app.bsky.feed.post docker compose -f docker-compose.tap.ym | `TAP_DISABLE_ACKS` | Disable ack-based delivery (useful for debugging) | `false` | | `TAP_SIGNAL_COLLECTION` | Collection NSID for auto-discovery of repos | *(empty)* | | `TAP_COLLECTION_FILTERS` | Comma-separated collection NSIDs for Tap sidecar record filtering; set independently from legacy `JETSTREAM_COLLECTIONS` | *(empty)* | +| `RECORD_HISTORY_COLLECTIONS` | Comma-separated collection NSIDs or `prefix.*` patterns whose versions are kept in `record_version` and served, paged, by the `recordHistory` GraphQL query (Tap mode only) | *(empty: no history)* | + +**Record version history.** Records in collections listed in `RECORD_HISTORY_COLLECTIONS` keep an append-only history: one row per distinct version (CID) the indexer observes, plus delete tombstones. On startup, every existing record in a newly listed collection gets a `baseline` row holding its current version, so history starts from what was indexed when it was switched on. Earlier versions are not recoverable because repositories only keep the latest one. Query it with `recordHistory(uri: "at://…", first: 100, after: "")`, one page of up to 500 versions at a time, oldest first; pass the last version's `id` as `after` for the next page. Versions are told apart by CID (or a hash of the body when Tap sends none), so repeats are stored once. Removing a collection from the variable stops recording; history already stored stays queryable. Tap Docker deployments also require `ADMIN_API_KEY` in `.env` because Hyperindex requires admin authentication at startup. `TAP_COLLECTION_FILTERS` is read by the Tap sidecar only; legacy `JETSTREAM_COLLECTIONS` remains part of Jetstream mode and is not used as a Tap filtering fallback. diff --git a/cmd/hyperindex/main.go b/cmd/hyperindex/main.go index 0c23001e..6d92b292 100644 --- a/cmd/hyperindex/main.go +++ b/cmd/hyperindex/main.go @@ -68,6 +68,7 @@ type services struct { labelDefinitions *repositories.LabelDefinitionsRepository labelPreferences *repositories.LabelPreferencesRepository reports *repositories.ReportsRepository + recordVersions *repositories.RecordVersionsRepository } // backgroundServices tracks cancellable background goroutines for clean shutdown. @@ -200,6 +201,7 @@ func initServices(cfg *config.Config) (*services, error) { labelDefinitions: repositories.NewLabelDefinitionsRepository(db), labelPreferences: repositories.NewLabelPreferencesRepository(db), reports: repositories.NewReportsRepository(db), + recordVersions: repositories.NewRecordVersionsRepository(db), } if cfg.PLCDirectoryURL != "" { @@ -802,6 +804,7 @@ func setupGraphQL(r *chi.Mux, cfg *config.Config, svc *services, pubsub *subscri Actors: svc.actors, Lexicons: svc.lexicons, ExternalLabels: svc.externalLabels, + RecordVersions: svc.recordVersions, } graphqlHandler, err := hgraphql.NewHandler(registry, repos) @@ -998,6 +1001,21 @@ func startTap( // Create handler that stores records and publishes to subscriptions. handler := tap.NewIndexHandler(svc.records, svc.actors, svc.activity, pubsub) + if historyCollections := repositories.NewCollectionMatcher(cfg.RecordHistoryCollections); !historyCollections.Empty() { + // Seed the current version of every opted-in record before live events + // start, so each history begins with what was indexed at switch-on. + // History is only switched on once seeding succeeds: recording events + // for records without a baseline would make the seed skip them later. + seedCtx, cancelSeed := context.WithTimeout(context.Background(), 10*time.Minute) + seeded, err := svc.recordVersions.SeedBaseline(seedCtx, historyCollections) + cancelSeed() + if err != nil { + slog.Error("Record version history stays off until the next start: baseline seeding failed", "collections", cfg.RecordHistoryCollections, "error", err) + } else { + handler.WithRecordHistory(svc.recordVersions, historyCollections) + slog.Info("Record version history enabled", "collections", cfg.RecordHistoryCollections, "baseline_seeded", seeded) + } + } // Create consumer. consumer := tap.NewConsumer(tap.ConsumerConfig{ diff --git a/internal/config/config.go b/internal/config/config.go index 857422e4..990361a4 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -63,6 +63,11 @@ type Config struct { TapDisableAcks bool // Fire-and-forget mode (default: false) TapEnabled bool // Use Tap instead of Jetstream+Backfill (default: false) + // Record version history (Tap mode): comma-separated collection NSIDs or + // `prefix.*` patterns whose every version is kept in record_version. + // Empty (the default) keeps no history. + RecordHistoryCollections string + // Labeler subscriptions LabelerSubscribeEnabled bool LabelerSubscribeURLs string // Comma-separated com.atproto.label.subscribeLabels websocket URLs @@ -127,6 +132,9 @@ func Load() (*Config, error) { TapDisableAcks: getEnvBool("TAP_DISABLE_ACKS", false), TapEnabled: getEnvBool("TAP_ENABLED", false), + // Record version history + RecordHistoryCollections: getEnv("RECORD_HISTORY_COLLECTIONS", ""), + // Labeler subscriptions LabelerSubscribeEnabled: getEnvBool("LABELER_SUBSCRIBE_ENABLED", false), LabelerSubscribeURLs: getEnv("LABELER_SUBSCRIBE_URLS", ""), diff --git a/internal/database/migrations/migrations_test.go b/internal/database/migrations/migrations_test.go index 4b6e1c5c..c85a51fe 100644 --- a/internal/database/migrations/migrations_test.go +++ b/internal/database/migrations/migrations_test.go @@ -219,8 +219,11 @@ func TestMigrations_RunAndRollbackPostgres(t *testing.T) { assertPostgresSequenceExists(ctx, t, exec, schemaName, "indexing_activity_id_seq") assertPostgresIndexNotExists(ctx, t, exec, schemaName, "idx_record_json_gin") - if err := migrations.Rollback(ctx, exec); err != nil { - t.Fatalf("Rollback() returned error: %v", err) + // 015 (record_version) is the newest migration; roll it back, then 011. + for _, version := range []string{"015", "011"} { + if err := migrations.Rollback(ctx, exec); err != nil { + t.Fatalf("Rollback(%s) returned error: %v", version, err) + } } assertPostgresIndexExists(ctx, t, exec, schemaName, "idx_record_json_gin") } @@ -499,3 +502,34 @@ func assertPostgresSequenceExists(ctx context.Context, t *testing.T, exec *postg t.Errorf("postgres sequence %q count = %d, want 1", sequenceName, count) } } + +func TestMigrations_RecordVersionUpAndDown(t *testing.T) { + exec := newTestExecutor(t) + ctx := context.Background() + if err := migrations.Run(ctx, exec); err != nil { + t.Fatalf("Run() error = %v", err) + } + indexCount := func() int { + t.Helper() + var count int + if err := exec.DB().QueryRowContext(ctx, `SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' AND name = 'idx_record_version_uri_key'`).Scan(&count); err != nil { + t.Fatalf("query record_version index: %v", err) + } + return count + } + if indexCount() != 1 { + t.Fatal("migration 015 did not create record_version") + } + if err := migrations.Rollback(ctx, exec); err != nil { + t.Fatalf("Rollback(015) error = %v", err) + } + if indexCount() != 0 { + t.Fatal("migration 015 schema remained after rollback") + } + if err := migrations.Run(ctx, exec); err != nil { + t.Fatalf("re-Run() error = %v", err) + } + if indexCount() != 1 { + t.Fatal("migration 015 was not re-applied") + } +} diff --git a/internal/database/migrations/postgres/015_add_record_version.down.sql b/internal/database/migrations/postgres/015_add_record_version.down.sql new file mode 100644 index 00000000..c041ca04 --- /dev/null +++ b/internal/database/migrations/postgres/015_add_record_version.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS record_version; diff --git a/internal/database/migrations/postgres/015_add_record_version.up.sql b/internal/database/migrations/postgres/015_add_record_version.up.sql new file mode 100644 index 00000000..07a577c6 --- /dev/null +++ b/internal/database/migrations/postgres/015_add_record_version.up.sql @@ -0,0 +1,25 @@ +-- Append-only version history for opted-in collections +-- (RECORD_HISTORY_COLLECTIONS). One row per distinct record version (CID) the +-- indexer observes, plus delete tombstones. `baseline` rows seed the version +-- that was current when history was switched on for a collection. +CREATE TABLE IF NOT EXISTS record_version ( + id BIGSERIAL PRIMARY KEY, + uri TEXT NOT NULL, + cid TEXT NOT NULL, + -- Dedupe identity: the CID, `sha256:` of the body when Tap omits the + -- CID, or `delete:` for a tombstone of that version. + version_key TEXT NOT NULL, + did TEXT NOT NULL, + collection TEXT NOT NULL, + action TEXT NOT NULL CHECK (action IN ('baseline', 'create', 'update', 'delete')), + json JSONB, + live BOOLEAN, + observed_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW() +); + +-- A version is recorded once, however often Tap redelivers or resyncs it. +CREATE UNIQUE INDEX IF NOT EXISTS idx_record_version_uri_key +ON record_version(uri, version_key); +CREATE INDEX IF NOT EXISTS idx_record_version_uri_id ON record_version(uri, id); +CREATE INDEX IF NOT EXISTS idx_record_version_collection_observed +ON record_version(collection, observed_at DESC); diff --git a/internal/database/migrations/sqlite/015_add_record_version.down.sql b/internal/database/migrations/sqlite/015_add_record_version.down.sql new file mode 100644 index 00000000..c041ca04 --- /dev/null +++ b/internal/database/migrations/sqlite/015_add_record_version.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS record_version; diff --git a/internal/database/migrations/sqlite/015_add_record_version.up.sql b/internal/database/migrations/sqlite/015_add_record_version.up.sql new file mode 100644 index 00000000..399127fd --- /dev/null +++ b/internal/database/migrations/sqlite/015_add_record_version.up.sql @@ -0,0 +1,23 @@ +-- Append-only version history for opted-in collections +-- (RECORD_HISTORY_COLLECTIONS). See the PostgreSQL migration for details. +CREATE TABLE IF NOT EXISTS record_version ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + uri TEXT NOT NULL, + cid TEXT NOT NULL, + -- Dedupe identity: the CID, `sha256:` of the body when Tap omits the + -- CID, or `delete:` for a tombstone of that version. + version_key TEXT NOT NULL, + did TEXT NOT NULL, + collection TEXT NOT NULL, + action TEXT NOT NULL CHECK (action IN ('baseline', 'create', 'update', 'delete')), + json TEXT, + live INTEGER, + observed_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%SZ', 'now')) +); + +-- A version is recorded once, however often Tap redelivers or resyncs it. +CREATE UNIQUE INDEX IF NOT EXISTS idx_record_version_uri_key +ON record_version(uri, version_key); +CREATE INDEX IF NOT EXISTS idx_record_version_uri_id ON record_version(uri, id); +CREATE INDEX IF NOT EXISTS idx_record_version_collection_observed +ON record_version(collection, observed_at DESC); diff --git a/internal/database/repositories/record_versions.go b/internal/database/repositories/record_versions.go new file mode 100644 index 00000000..e209f8d0 --- /dev/null +++ b/internal/database/repositories/record_versions.go @@ -0,0 +1,270 @@ +package repositories + +import ( + "context" + "crypto/sha256" + "database/sql" + "encoding/hex" + "fmt" + "strings" + "time" + + "github.com/GainForest/hyperindex/internal/database" +) + +// Record version actions stored in record_version.action. +const ( + // RecordVersionBaseline is the version that was current when history was + // switched on for a collection (seeded from the record table). + RecordVersionBaseline = "baseline" + RecordVersionCreate = "create" + RecordVersionUpdate = "update" + RecordVersionDelete = "delete" + + // MaxRecordHistoryPageSize bounds one recordHistory response. + MaxRecordHistoryPageSize = 500 +) + +// RecordVersion is one observed version of a record, or a delete tombstone. +type RecordVersion struct { + ID int64 + URI string + CID string + DID string + Collection string + Action string + // JSON is the record body; nil for delete tombstones. + JSON *string + // Live reports whether Tap delivered the event from the live stream + // (false for resync/backfill deliveries); nil for baseline rows. + Live *bool + ObservedAt time.Time +} + +// RecordVersionWrite is a version to append. +type RecordVersionWrite struct { + URI string + CID string + DID string + Collection string + Action string + JSON *string + Live *bool +} + +// CollectionMatcher selects which collections keep version history. Entries +// are exact NSIDs or `prefix.*` patterns (e.g. `app.gainforest.dwc.*`). +type CollectionMatcher struct { + exact map[string]struct{} + prefixes []string +} + +// NewCollectionMatcher parses a comma-separated collection list. +func NewCollectionMatcher(list string) CollectionMatcher { + m := CollectionMatcher{exact: map[string]struct{}{}} + for _, raw := range strings.Split(list, ",") { + entry := strings.TrimSpace(raw) + if entry == "" { + continue + } + if strings.HasSuffix(entry, ".*") { + m.prefixes = append(m.prefixes, strings.TrimSuffix(entry, "*")) + continue + } + m.exact[entry] = struct{}{} + } + return m +} + +// Empty reports whether history is switched off. +func (m CollectionMatcher) Empty() bool { + return len(m.exact) == 0 && len(m.prefixes) == 0 +} + +// Matches reports whether a collection keeps version history. +func (m CollectionMatcher) Matches(collection string) bool { + if _, ok := m.exact[collection]; ok { + return true + } + for _, prefix := range m.prefixes { + if strings.HasPrefix(collection, prefix) { + return true + } + } + return false +} + +// sqlPredicate renders the matcher as a WHERE predicate on `column`, with +// placeholders starting after `base`. +func (m CollectionMatcher) sqlPredicate(db database.Executor, column string, base int) (string, []database.Value) { + var parts []string + var params []database.Value + for collection := range m.exact { + params = append(params, database.Text(collection)) + parts = append(parts, fmt.Sprintf("%s = %s", column, db.Placeholder(base+len(params)))) + } + for _, prefix := range m.prefixes { + params = append(params, database.Text(escapeLike(prefix)+"%")) + parts = append(parts, fmt.Sprintf("%s LIKE %s ESCAPE '\\'", column, db.Placeholder(base+len(params)))) + } + if len(parts) == 0 { + return "1 = 0", nil + } + return "(" + strings.Join(parts, " OR ") + ")", params +} + +func escapeLike(value string) string { + return strings.NewReplacer(`\`, `\\`, `%`, `\%`, `_`, `\_`).Replace(value) +} + +// RecordVersionsRepository stores and reads record version history. +type RecordVersionsRepository struct { + db database.Executor +} + +// NewRecordVersionsRepository creates a record version repository. +func NewRecordVersionsRepository(db database.Executor) *RecordVersionsRepository { + return &RecordVersionsRepository{db: db} +} + +// versionKey is the dedupe identity of a version: its CID; for Tap events +// without a CID, a hash of the body; for tombstones, the deleted CID. +func versionKey(v RecordVersionWrite) string { + if v.Action == RecordVersionDelete { + return "delete:" + v.CID + } + if v.CID != "" { + return v.CID + } + body := "" + if v.JSON != nil { + body = *v.JSON + } + sum := sha256.Sum256([]byte(body)) + return "sha256:" + hex.EncodeToString(sum[:]) +} + +// Append records one version. A version already stored for the URI (same +// version key) is left untouched, so redelivered and resynced events, +// including repeated deletes, are no-ops. +func (r *RecordVersionsRepository) Append(ctx context.Context, v RecordVersionWrite) error { + if v.URI == "" || v.DID == "" || v.Collection == "" { + return fmt.Errorf("record version requires uri, did and collection") + } + switch v.Action { + case RecordVersionCreate, RecordVersionUpdate, RecordVersionDelete: + default: + return fmt.Errorf("unsupported record version action %q", v.Action) + } + jsonPlaceholder := r.db.Placeholder(7) + if r.db.Dialect() == database.PostgreSQL { + jsonPlaceholder += "::jsonb" + } + sqlStr := fmt.Sprintf(`INSERT INTO record_version (uri, cid, version_key, did, collection, action, json, live) + VALUES (%s, %s, %s, %s, %s, %s, %s, %s) + ON CONFLICT DO NOTHING`, + r.db.Placeholder(1), r.db.Placeholder(2), r.db.Placeholder(3), r.db.Placeholder(4), r.db.Placeholder(5), + r.db.Placeholder(6), jsonPlaceholder, r.db.Placeholder(8)) + live := database.Value(database.Null()) + if v.Live != nil { + live = database.Bool(*v.Live) + } + _, err := r.db.Exec(ctx, sqlStr, []database.Value{ + database.Text(v.URI), + database.Text(v.CID), + database.Text(versionKey(v)), + database.Text(v.DID), + database.Text(v.Collection), + database.Text(v.Action), + database.NullableText(v.JSON), + live, + }) + if err != nil { + return fmt.Errorf("append record version for %s: %w", v.URI, err) + } + return nil +} + +// SeedBaseline adds a `baseline` version for every current record in a +// matched collection that has no history yet, stamped with the record's +// indexed_at. Idempotent: records that already have versions are skipped, so +// running it on every startup only seeds newly opted-in collections. +func (r *RecordVersionsRepository) SeedBaseline(ctx context.Context, matcher CollectionMatcher) (int64, error) { + if matcher.Empty() { + return 0, nil + } + predicate, params := matcher.sqlPredicate(r.db, "rec.collection", 0) + observedAt := "rec.indexed_at" + if r.db.Dialect() == database.SQLite { + observedAt = "strftime('%Y-%m-%dT%H:%M:%SZ', rec.indexed_at)" + } + sqlStr := fmt.Sprintf(`INSERT INTO record_version (uri, cid, version_key, did, collection, action, json, live, observed_at) + SELECT rec.uri, rec.cid, CASE WHEN rec.cid = '' THEN 'baseline' ELSE rec.cid END, rec.did, rec.collection, '%s', rec.json, NULL, %s + FROM record rec + WHERE %s + AND NOT EXISTS (SELECT 1 FROM record_version v WHERE v.uri = rec.uri) + ON CONFLICT DO NOTHING`, RecordVersionBaseline, observedAt, predicate) + result, err := r.db.Exec(ctx, sqlStr, params) + if err != nil { + return 0, fmt.Errorf("seed record version baseline: %w", err) + } + inserted, err := result.RowsAffected() + if err != nil { + return 0, fmt.Errorf("count seeded record versions: %w", err) + } + return inserted, nil +} + +// ListByURI returns up to `limit` of a record's versions oldest first, +// starting after version id `afterID` (0 for the first page). Callers page by +// passing the last returned id. +func (r *RecordVersionsRepository) ListByURI(ctx context.Context, uri string, afterID int64, limit int) ([]RecordVersion, error) { + if limit <= 0 || limit > MaxRecordHistoryPageSize { + return nil, fmt.Errorf("record history page size must be between 1 and %d", MaxRecordHistoryPageSize) + } + jsonExpr, liveExpr, observedExpr := "json", "live", "observed_at" + if r.db.Dialect() == database.PostgreSQL { + jsonExpr, liveExpr, observedExpr = "json::text", "CASE WHEN live IS NULL THEN NULL WHEN live THEN 1 ELSE 0 END", "observed_at::text" + } + sqlStr := fmt.Sprintf(`SELECT id, uri, cid, did, collection, action, %s, %s, %s + FROM record_version + WHERE uri = %s AND id > %s + ORDER BY id ASC + LIMIT %d`, jsonExpr, liveExpr, observedExpr, r.db.Placeholder(1), r.db.Placeholder(2), limit) + rows, err := r.db.DB().QueryContext(ctx, sqlStr, r.db.ConvertParams([]database.Value{database.Text(uri), database.Int(afterID)})...) + if err != nil { + return nil, fmt.Errorf("list record versions for %s: %w", uri, err) + } + defer func() { _ = rows.Close() }() + + versions := make([]RecordVersion, 0) + for rows.Next() { + var ( + v RecordVersion + jsonText sql.NullString + live sql.NullInt64 + observedAt string + ) + if err := rows.Scan(&v.ID, &v.URI, &v.CID, &v.DID, &v.Collection, &v.Action, &jsonText, &live, &observedAt); err != nil { + return nil, fmt.Errorf("scan record version: %w", err) + } + if jsonText.Valid { + text := jsonText.String + v.JSON = &text + } + if live.Valid { + value := live.Int64 != 0 + v.Live = &value + } + parsed, err := parseDBTime(observedAt) + if err != nil { + return nil, fmt.Errorf("parse record_version.observed_at: %w", err) + } + v.ObservedAt = parsed + versions = append(versions, v) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate record versions: %w", err) + } + return versions, nil +} diff --git a/internal/database/repositories/record_versions_test.go b/internal/database/repositories/record_versions_test.go new file mode 100644 index 00000000..dcfa3f95 --- /dev/null +++ b/internal/database/repositories/record_versions_test.go @@ -0,0 +1,192 @@ +package repositories_test + +import ( + "context" + "strconv" + "testing" + "time" + + "github.com/GainForest/hyperindex/internal/database/repositories" + "github.com/GainForest/hyperindex/internal/testutil" +) + +// recordVersionTestDBs returns SQLite always, plus PostgreSQL when DATABASE_URL +// points at a safe test database. +func recordVersionTestDBs(t *testing.T) map[string]*testutil.TestDB { + t.Helper() + dbs := map[string]*testutil.TestDB{"sqlite": testutil.SetupTestDB(t)} + if url, ok := safePostgresTestDatabaseURL(t); ok { + dbs["postgres"] = testutil.SetupTestDBWithURL(t, url) + } + return dbs +} + +func boolPtr(value bool) *bool { return &value } + +func TestCollectionMatcher(t *testing.T) { + m := repositories.NewCollectionMatcher(" app.gainforest.dwc.occurrence , org.hypercerts.* ,,") + for collection, want := range map[string]bool{ + "app.gainforest.dwc.occurrence": true, + "app.gainforest.dwc.dataset": false, + "org.hypercerts.collection": true, + "org.hypercertsx.collection": false, + "app.bsky.feed.post": false, + } { + if got := m.Matches(collection); got != want { + t.Errorf("Matches(%q) = %v, want %v", collection, got, want) + } + } + if !repositories.NewCollectionMatcher(" ").Empty() { + t.Error("blank list should disable history") + } +} + +func TestRecordVersions_AppendDedupesAndListsOldestFirst(t *testing.T) { + for name, db := range recordVersionTestDBs(t) { + t.Run(name, func(t *testing.T) { + ctx := context.Background() + repo := db.RecordVersions + suffix := strconv.FormatInt(time.Now().UnixNano(), 36) + uri := "at://did:plc:alice/app.gainforest.dwc.occurrence/" + suffix + write := func(action, cid string, body *string) { + t.Helper() + if err := repo.Append(ctx, repositories.RecordVersionWrite{ + URI: uri, CID: cid, DID: "did:plc:alice", Collection: "app.gainforest.dwc.occurrence", + Action: action, JSON: body, Live: boolPtr(true), + }); err != nil { + t.Fatalf("Append(%s, %s): %v", action, cid, err) + } + } + write(repositories.RecordVersionCreate, "cid1", strPtr(`{"scientificName":"Hirundo rustica","identifiedBy":"ai:gemini-3.8-flash"}`)) + write(repositories.RecordVersionCreate, "cid1", strPtr(`{"scientificName":"Hirundo rustica","identifiedBy":"ai:gemini-3.8-flash"}`)) // redelivery + write(repositories.RecordVersionUpdate, "cid2", strPtr(`{"scientificName":"Hirundo tahitica","identifiedBy":"Maria"}`)) + write(repositories.RecordVersionDelete, "cid2", nil) + write(repositories.RecordVersionDelete, "cid2", nil) // redelivered delete + + versions, err := repo.ListByURI(ctx, uri, 0, 500) + if err != nil { + t.Fatalf("ListByURI: %v", err) + } + if len(versions) != 3 { + t.Fatalf("got %d versions, want 3 (redelivery deduped): %+v", len(versions), versions) + } + wantActions := []string{"create", "update", "delete"} + for i, v := range versions { + if v.Action != wantActions[i] { + t.Errorf("version %d action = %q, want %q", i, v.Action, wantActions[i]) + } + if v.ObservedAt.IsZero() { + t.Errorf("version %d has no observedAt", i) + } + } + if versions[1].JSON == nil || versions[1].CID != "cid2" { + t.Errorf("update version = %+v", versions[1]) + } + if versions[2].JSON != nil { + t.Errorf("delete tombstone should have no body, got %q", *versions[2].JSON) + } + if versions[0].Live == nil || !*versions[0].Live { + t.Errorf("live flag not kept: %+v", versions[0].Live) + } + if err := repo.Append(ctx, repositories.RecordVersionWrite{URI: uri, DID: "did:plc:alice", Collection: "x", Action: "baseline"}); err == nil { + t.Error("Append should refuse the baseline action; baselines come from SeedBaseline") + } + }) + } +} + +func TestRecordVersions_SeedBaselineIsIdempotentAndScoped(t *testing.T) { + for name, db := range recordVersionTestDBs(t) { + t.Run(name, func(t *testing.T) { + ctx := context.Background() + suffix := strconv.FormatInt(time.Now().UnixNano(), 36) + tracked := "at://did:plc:bob/app.gainforest.dwc.occurrence/" + suffix + untracked := "at://did:plc:bob/app.bsky.feed.post/" + suffix + withHistory := "at://did:plc:bob/app.gainforest.dwc.occurrence/h" + suffix + if err := db.Records.BatchInsert(ctx, []*repositories.Record{ + {URI: tracked, CID: "c-tracked", DID: "did:plc:bob", Collection: "app.gainforest.dwc.occurrence", JSON: `{"scientificName":"Quercus robur"}`}, + {URI: untracked, CID: "c-post", DID: "did:plc:bob", Collection: "app.bsky.feed.post", JSON: `{"text":"hi"}`}, + {URI: withHistory, CID: "c-new", DID: "did:plc:bob", Collection: "app.gainforest.dwc.occurrence", JSON: `{"scientificName":"Quercus petraea"}`}, + }); err != nil { + t.Fatalf("BatchInsert: %v", err) + } + if err := db.RecordVersions.Append(ctx, repositories.RecordVersionWrite{ + URI: withHistory, CID: "c-new", DID: "did:plc:bob", Collection: "app.gainforest.dwc.occurrence", + Action: repositories.RecordVersionCreate, JSON: strPtr(`{"scientificName":"Quercus petraea"}`), Live: boolPtr(true), + }); err != nil { + t.Fatalf("Append: %v", err) + } + + matcher := repositories.NewCollectionMatcher("app.gainforest.*") + if _, err := db.RecordVersions.SeedBaseline(ctx, matcher); err != nil { + t.Fatalf("SeedBaseline: %v", err) + } + if _, err := db.RecordVersions.SeedBaseline(ctx, matcher); err != nil { + t.Fatalf("SeedBaseline (second run): %v", err) + } + + assertVersions := func(uri string, wantActions ...string) { + t.Helper() + versions, err := db.RecordVersions.ListByURI(ctx, uri, 0, 500) + if err != nil { + t.Fatalf("ListByURI(%s): %v", uri, err) + } + if len(versions) != len(wantActions) { + t.Fatalf("%s: got %d versions, want %v", uri, len(versions), wantActions) + } + for i, want := range wantActions { + if versions[i].Action != want { + t.Errorf("%s version %d = %q, want %q", uri, i, versions[i].Action, want) + } + } + } + assertVersions(tracked, "baseline") + assertVersions(untracked) + assertVersions(withHistory, "create") + + versions, _ := db.RecordVersions.ListByURI(ctx, tracked, 0, 500) + if versions[0].CID != "c-tracked" || versions[0].JSON == nil || versions[0].Live != nil { + t.Errorf("baseline = %+v", versions[0]) + } + }) + } +} + +func TestRecordVersions_CIDLessVersionsStayDistinctAndPages(t *testing.T) { + for name, db := range recordVersionTestDBs(t) { + t.Run(name, func(t *testing.T) { + ctx := context.Background() + repo := db.RecordVersions + uri := "at://did:plc:carol/app.gainforest.dwc.occurrence/" + strconv.FormatInt(time.Now().UnixNano(), 36) + for _, body := range []string{`{"scientificName":"A a"}`, `{"scientificName":"B b"}`, `{"scientificName":"A a"}`, `{"scientificName":"C c"}`} { + body := body + if err := repo.Append(ctx, repositories.RecordVersionWrite{ + URI: uri, DID: "did:plc:carol", Collection: "app.gainforest.dwc.occurrence", + Action: repositories.RecordVersionUpdate, JSON: &body, + }); err != nil { + t.Fatalf("Append: %v", err) + } + } + all, err := repo.ListByURI(ctx, uri, 0, 500) + if err != nil { + t.Fatalf("ListByURI: %v", err) + } + // Without a CID, versions are told apart by content: A, B, C (the + // repeated A body is the same version). + if len(all) != 3 { + t.Fatalf("got %d CID-less versions, want 3", len(all)) + } + page1, err := repo.ListByURI(ctx, uri, 0, 2) + if err != nil || len(page1) != 2 { + t.Fatalf("page 1 = %d versions (err %v), want 2", len(page1), err) + } + page2, err := repo.ListByURI(ctx, uri, page1[1].ID, 2) + if err != nil || len(page2) != 1 || page2[0].ID != all[2].ID { + t.Fatalf("page 2 = %+v (err %v), want the third version", page2, err) + } + if _, err := repo.ListByURI(ctx, uri, 0, 0); err == nil { + t.Fatal("ListByURI should reject a zero page size") + } + }) + } +} diff --git a/internal/graphql/recordhistory/recordhistory.go b/internal/graphql/recordhistory/recordhistory.go new file mode 100644 index 00000000..a584713e --- /dev/null +++ b/internal/graphql/recordhistory/recordhistory.go @@ -0,0 +1,80 @@ +// Package recordhistory exposes the append-only record version history kept +// for collections opted in with RECORD_HISTORY_COLLECTIONS. +package recordhistory + +import ( + "encoding/json" + "log/slog" + "strconv" + "time" + + "github.com/graphql-go/graphql" + + "github.com/GainForest/hyperindex/internal/database/repositories" + "github.com/GainForest/hyperindex/internal/graphql/types" +) + +// DefaultPageSize is the number of versions returned when `first` is omitted. +const DefaultPageSize = 100 + +// Type is one observed version of a record, or a delete tombstone. +var Type = graphql.NewObject(graphql.ObjectConfig{ + Name: "RecordVersion", + Description: "One observed version of a record (oldest first), or a delete tombstone. Only kept for collections the indexer is configured to track.", + Fields: graphql.Fields{ + "id": &graphql.Field{ + Type: graphql.NewNonNull(graphql.String), + Description: "Monotonic version id; later versions have larger ids.", + }, + "uri": &graphql.Field{Type: graphql.NewNonNull(graphql.String), Description: "Record AT-URI."}, + "cid": &graphql.Field{Type: graphql.NewNonNull(graphql.String), Description: "CID of this version (for deletes, the version that was deleted)."}, + "did": &graphql.Field{Type: graphql.NewNonNull(graphql.String), Description: "Repository DID."}, + "collection": &graphql.Field{Type: graphql.NewNonNull(graphql.String), Description: "Collection NSID."}, + "action": &graphql.Field{ + Type: graphql.NewNonNull(graphql.String), + Description: "baseline (current when history was switched on), create, update, or delete.", + }, + "value": &graphql.Field{ + Type: types.JSONScalar, + Description: "The record body at this version; null for deletes.", + }, + "live": &graphql.Field{ + Type: graphql.Boolean, + Description: "True when seen on the live stream, false for a resync delivery, null for baseline rows.", + }, + "observedAt": &graphql.Field{ + Type: graphql.NewNonNull(graphql.String), + Description: "When the indexer observed this version (RFC 3339). For baseline rows, when the version was indexed.", + }, + }, +}) + +// ToGraphQL converts versions to GraphQL result maps. +func ToGraphQL(versions []repositories.RecordVersion) []map[string]interface{} { + out := make([]map[string]interface{}, 0, len(versions)) + for _, v := range versions { + var value interface{} + if v.JSON != nil { + if err := json.Unmarshal([]byte(*v.JSON), &value); err != nil { + slog.Warn("Record version has invalid JSON", "uri", v.URI, "id", v.ID, "error", err) + value = nil + } + } + var live interface{} + if v.Live != nil { + live = *v.Live + } + out = append(out, map[string]interface{}{ + "id": strconv.FormatInt(v.ID, 10), + "uri": v.URI, + "cid": v.CID, + "did": v.DID, + "collection": v.Collection, + "action": v.Action, + "value": value, + "live": live, + "observedAt": v.ObservedAt.UTC().Format(time.RFC3339Nano), + }) + } + return out +} diff --git a/internal/graphql/resolver/context.go b/internal/graphql/resolver/context.go index 8e263d44..2f6e5ed1 100644 --- a/internal/graphql/resolver/context.go +++ b/internal/graphql/resolver/context.go @@ -21,6 +21,7 @@ type Repositories struct { Actors *repositories.ActorsRepository Lexicons *repositories.LexiconsRepository ExternalLabels *repositories.ExternalLabelsRepository + RecordVersions *repositories.RecordVersionsRepository } // NewRepositories creates a new Repositories from a database executor. @@ -30,6 +31,7 @@ func NewRepositories(db database.Executor) *Repositories { Actors: repositories.NewActorsRepository(db), Lexicons: repositories.NewLexiconsRepository(db), ExternalLabels: repositories.NewExternalLabelsRepository(db), + RecordVersions: repositories.NewRecordVersionsRepository(db), } } diff --git a/internal/graphql/schema/builder.go b/internal/graphql/schema/builder.go index 6a10426f..77c7e8bb 100644 --- a/internal/graphql/schema/builder.go +++ b/internal/graphql/schema/builder.go @@ -10,6 +10,7 @@ import ( "fmt" "log/slog" "maps" + "strconv" "strings" "time" "unicode/utf8" @@ -22,6 +23,7 @@ import ( "github.com/GainForest/hyperindex/internal/graphql/certifiedprofiles" "github.com/GainForest/hyperindex/internal/graphql/externallabels" "github.com/GainForest/hyperindex/internal/graphql/query" + "github.com/GainForest/hyperindex/internal/graphql/recordhistory" "github.com/GainForest/hyperindex/internal/graphql/resolver" "github.com/GainForest/hyperindex/internal/graphql/types" "github.com/GainForest/hyperindex/internal/lexicon" @@ -662,6 +664,28 @@ func (b *Builder) buildQueryType() *graphql.Object { Resolve: b.createRecordTimelineResolver(), } + // Add version history lookup for records in history-enabled collections. + fields["recordHistory"] = &graphql.Field{ + Type: graphql.NewNonNull(graphql.NewList(graphql.NewNonNull(recordhistory.Type))), + Description: "Versions the indexer has observed for one record, oldest first, one page at a time. Recorded only while the record's collection is configured for history (RECORD_HISTORY_COLLECTIONS); history recorded earlier stays queryable after a collection is removed.", + Args: graphql.FieldConfigArgument{ + "uri": &graphql.ArgumentConfig{ + Type: graphql.NewNonNull(graphql.String), + Description: "Record AT-URI.", + }, + "first": &graphql.ArgumentConfig{ + Type: graphql.Int, + DefaultValue: recordhistory.DefaultPageSize, + Description: "Maximum versions to return, 1 to 500 (default 100).", + }, + "after": &graphql.ArgumentConfig{ + Type: graphql.String, + Description: "Return versions after this version id (the `id` of the last version on the previous page).", + }, + }, + Resolve: b.createRecordHistoryResolver(), + } + // Add external label lookup by subject DID or AT-URI. fields["externalLabels"] = &graphql.Field{ Type: graphql.NewNonNull(graphql.NewList(graphql.NewNonNull(externallabels.Type))), @@ -1653,6 +1677,36 @@ func (b *Builder) resolveRecordConnection( } // createExternalLabelsResolver creates a resolver for generic external label subject lookups. +func (b *Builder) createRecordHistoryResolver() graphql.FieldResolveFn { + return func(p graphql.ResolveParams) (interface{}, error) { + repos := resolver.GetRepositories(p.Context) + if repos == nil || repos.RecordVersions == nil { + return []map[string]interface{}{}, nil + } + uri, _ := p.Args["uri"].(string) + if strings.TrimSpace(uri) == "" { + return []map[string]interface{}{}, nil + } + first, _ := p.Args["first"].(int) + if first < 1 || first > repositories.MaxRecordHistoryPageSize { + return nil, fmt.Errorf("first must be between 1 and %d", repositories.MaxRecordHistoryPageSize) + } + var afterID int64 + if after, ok := p.Args["after"].(string); ok && strings.TrimSpace(after) != "" { + parsed, err := strconv.ParseInt(strings.TrimSpace(after), 10, 64) + if err != nil || parsed < 0 { + return nil, fmt.Errorf("after must be a version id from a previous recordHistory page") + } + afterID = parsed + } + versions, err := repos.RecordVersions.ListByURI(p.Context, uri, afterID, first) + if err != nil { + return nil, fmt.Errorf("failed to query record history: %w", err) + } + return recordhistory.ToGraphQL(versions), nil + } +} + func (b *Builder) createExternalLabelsResolver() graphql.FieldResolveFn { return func(p graphql.ResolveParams) (interface{}, error) { repos := resolver.GetRepositories(p.Context) diff --git a/internal/graphql/schema/record_history_test.go b/internal/graphql/schema/record_history_test.go new file mode 100644 index 00000000..477edf9b --- /dev/null +++ b/internal/graphql/schema/record_history_test.go @@ -0,0 +1,87 @@ +package schema + +import ( + "context" + "testing" + + "github.com/graphql-go/graphql" + + "github.com/GainForest/hyperindex/internal/database/repositories" + "github.com/GainForest/hyperindex/internal/graphql/resolver" + "github.com/GainForest/hyperindex/internal/testutil" +) + +func TestRecordHistoryGraphQL(t *testing.T) { + const uri = "at://did:plc:alice/app.gainforest.dwc.occurrence/occ1" + schema := buildExternalLabelsTestSchema(t) + db := testutil.SetupTestDB(t) + ctx := context.Background() + + body1 := `{"scientificName":"Hirundo rustica","identifiedBy":"ai:gemini-3.8-flash"}` + body2 := `{"scientificName":"Hirundo tahitica","identifiedBy":"Maria","previousIdentifications":"Hirundo rustica (ai:gemini-3.8-flash, 2026-09-22)"}` + live := true + for _, v := range []repositories.RecordVersionWrite{ + {URI: uri, CID: "cid1", DID: "did:plc:alice", Collection: "app.gainforest.dwc.occurrence", Action: repositories.RecordVersionCreate, JSON: &body1, Live: &live}, + {URI: uri, CID: "cid2", DID: "did:plc:alice", Collection: "app.gainforest.dwc.occurrence", Action: repositories.RecordVersionUpdate, JSON: &body2, Live: &live}, + } { + if err := db.RecordVersions.Append(ctx, v); err != nil { + t.Fatalf("Append: %v", err) + } + } + ctx = resolver.WithRepositories(ctx, &resolver.Repositories{Records: db.Records, RecordVersions: db.RecordVersions}) + + result := graphql.Do(graphql.Params{ + Schema: *schema, + RequestString: `{ recordHistory(uri: "` + uri + `") { id cid action live observedAt value } missing: recordHistory(uri: "at://did:plc:none/x.y.z/1") { id } }`, + Context: ctx, + }) + if len(result.Errors) > 0 { + t.Fatalf("GraphQL errors: %v", result.Errors) + } + data := result.Data.(map[string]interface{}) + versions := data["recordHistory"].([]interface{}) + if len(versions) != 2 { + t.Fatalf("recordHistory length = %d, want 2", len(versions)) + } + first := versions[0].(map[string]interface{}) + second := versions[1].(map[string]interface{}) + if first["action"] != "create" || second["action"] != "update" || second["cid"] != "cid2" || first["live"] != true { + t.Fatalf("versions = %+v", versions) + } + value := second["value"].(map[string]interface{}) + if value["scientificName"] != "Hirundo tahitica" || value["previousIdentifications"] == nil { + t.Fatalf("second value = %+v", value) + } + if first["observedAt"] == "" { + t.Fatal("observedAt missing") + } + if missing := data["missing"].([]interface{}); len(missing) != 0 { + t.Fatalf("unknown uri returned %v", missing) + } + + // Paging: the second page starts after the first version's id. + firstID := first["id"].(string) + paged := graphql.Do(graphql.Params{ + Schema: *schema, + RequestString: `{ recordHistory(uri: "` + uri + `", first: 1, after: "` + firstID + `") { cid } }`, + Context: ctx, + }) + if len(paged.Errors) > 0 { + t.Fatalf("paged GraphQL errors: %v", paged.Errors) + } + page := paged.Data.(map[string]interface{})["recordHistory"].([]interface{}) + if len(page) != 1 || page[0].(map[string]interface{})["cid"] != "cid2" { + t.Fatalf("second page = %v, want cid2", page) + } + + for _, bad := range []string{`first: 0`, `first: 501`, `after: "nope"`} { + res := graphql.Do(graphql.Params{ + Schema: *schema, + RequestString: `{ recordHistory(uri: "` + uri + `", ` + bad + `) { id } }`, + Context: ctx, + }) + if len(res.Errors) == 0 { + t.Errorf("recordHistory(%s) should fail", bad) + } + } +} diff --git a/internal/tap/handler.go b/internal/tap/handler.go index 55697a6b..a6de0257 100644 --- a/internal/tap/handler.go +++ b/internal/tap/handler.go @@ -2,6 +2,8 @@ package tap import ( "context" + "database/sql" + "errors" "fmt" "log/slog" "strings" @@ -17,6 +19,22 @@ type IndexHandler struct { actors *repositories.ActorsRepository activity *repositories.IndexingActivityRepository // records indexing activity pubsub *subscription.PubSub + + // Optional version history for opted-in collections (RECORD_HISTORY_COLLECTIONS). + versions *repositories.RecordVersionsRepository + historyCollections repositories.CollectionMatcher +} + +// WithRecordHistory enables append-only version history for the collections +// the matcher selects. Every distinct version and delete is recorded. +func (h *IndexHandler) WithRecordHistory(versions *repositories.RecordVersionsRepository, collections repositories.CollectionMatcher) *IndexHandler { + h.versions = versions + h.historyCollections = collections + return h +} + +func (h *IndexHandler) keepsHistory(collection string) bool { + return h.versions != nil && h.historyCollections.Matches(collection) } // NewIndexHandler creates a new IndexHandler. @@ -55,6 +73,29 @@ func (h *IndexHandler) HandleRecord(ctx context.Context, event *RecordEvent) err slog.Debug("Failed to upsert actor", "did", event.DID, "error", err) } + // Record the version before the current-state write: if either write + // fails, Tap redelivers and the version insert is an idempotent no-op. + if h.keepsHistory(event.Collection) { + action := repositories.RecordVersionCreate + if event.Action == ActionUpdate { + action = repositories.RecordVersionUpdate + } + body := string(event.Record) + live := event.Live + setEventPhase(ctx, "record_versions.append") + if err := h.versions.Append(ctx, repositories.RecordVersionWrite{ + URI: uri, + CID: event.CID, + DID: event.DID, + Collection: event.Collection, + Action: action, + JSON: &body, + Live: &live, + }); err != nil { + return fmt.Errorf("failed to record version history: %w", err) + } + } + // Store record setEventPhase(ctx, "records.insert") result, err := h.records.Insert(ctx, uri, event.CID, event.DID, event.Collection, string(event.Record)) @@ -91,6 +132,32 @@ func (h *IndexHandler) HandleRecord(ctx context.Context, event *RecordEvent) err } case ActionDelete: + // Tombstone the version being deleted before removing it. If either + // write fails the event is retried: the tombstone insert is idempotent + // (keyed by the deleted CID) and the record is still there to delete. + if h.keepsHistory(event.Collection) { + setEventPhase(ctx, "record_versions.tombstone") + current, err := h.records.GetByURI(ctx, uri) + switch { + case err == nil: + live := event.Live + if err := h.versions.Append(ctx, repositories.RecordVersionWrite{ + URI: uri, + CID: current.CID, + DID: event.DID, + Collection: event.Collection, + Action: repositories.RecordVersionDelete, + Live: &live, + }); err != nil { + return fmt.Errorf("failed to record delete in version history: %w", err) + } + case errors.Is(err, sql.ErrNoRows): + // Nothing indexed to delete (or already deleted and tombstoned). + default: + return fmt.Errorf("failed to read record before delete: %w", err) + } + } + setEventPhase(ctx, "records.delete") if err := h.records.Delete(ctx, uri); err != nil { return fmt.Errorf("failed to delete record: %w", err) diff --git a/internal/tap/handler_history_test.go b/internal/tap/handler_history_test.go new file mode 100644 index 00000000..a166fd2a --- /dev/null +++ b/internal/tap/handler_history_test.go @@ -0,0 +1,75 @@ +package tap_test + +import ( + "context" + "encoding/json" + "testing" + + "github.com/GainForest/hyperindex/internal/database/repositories" + "github.com/GainForest/hyperindex/internal/tap" +) + +func TestIndexHandler_RecordHistory(t *testing.T) { + handler, db, _ := setupHandler(t) + handler.WithRecordHistory(db.RecordVersions, repositories.NewCollectionMatcher("app.gainforest.dwc.occurrence")) + ctx := context.Background() + + occurrence := func(action tap.ActionType, cid, body string, live bool) *tap.RecordEvent { + return &tap.RecordEvent{ + Live: live, DID: "did:plc:alice", Collection: "app.gainforest.dwc.occurrence", RKey: "occ1", + Action: action, CID: cid, Record: json.RawMessage(body), + } + } + events := []*tap.RecordEvent{ + occurrence(tap.ActionCreate, "cid1", `{"scientificName":"Hirundo rustica","identifiedBy":"ai:gemini-3.8-flash"}`, true), + occurrence(tap.ActionCreate, "cid1", `{"scientificName":"Hirundo rustica","identifiedBy":"ai:gemini-3.8-flash"}`, false), // resync of the same version + occurrence(tap.ActionUpdate, "cid2", `{"scientificName":"Hirundo tahitica","identifiedBy":"Maria"}`, true), + {Live: true, DID: "did:plc:alice", Collection: "app.gainforest.dwc.occurrence", RKey: "occ1", Action: tap.ActionDelete}, + {Live: true, DID: "did:plc:alice", Collection: "app.gainforest.dwc.occurrence", RKey: "occ1", Action: tap.ActionDelete}, // redelivered delete + // Collections without history are indexed as before and leave no versions. + {Live: true, DID: "did:plc:alice", Collection: "app.bsky.feed.post", RKey: "p1", Action: tap.ActionCreate, CID: "cidp", Record: json.RawMessage(`{"text":"hi"}`)}, + } + for i, event := range events { + if err := handler.HandleRecord(ctx, event); err != nil { + t.Fatalf("event %d: HandleRecord: %v", i, err) + } + } + + versions, err := db.RecordVersions.ListByURI(ctx, "at://did:plc:alice/app.gainforest.dwc.occurrence/occ1", 0, 500) + if err != nil { + t.Fatalf("ListByURI: %v", err) + } + var actions []string + for _, v := range versions { + actions = append(actions, v.Action+":"+v.CID) + } + want := []string{"create:cid1", "update:cid2", "delete:cid2"} + if len(actions) != len(want) { + t.Fatalf("versions = %v, want %v", actions, want) + } + for i := range want { + if actions[i] != want[i] { + t.Fatalf("versions = %v, want %v", actions, want) + } + } + + posts, err := db.RecordVersions.ListByURI(ctx, "at://did:plc:alice/app.bsky.feed.post/p1", 0, 500) + if err != nil || len(posts) != 0 { + t.Fatalf("untracked collection versions = %v (err %v), want none", posts, err) + } +} + +func TestIndexHandler_NoHistoryByDefault(t *testing.T) { + handler, db, _ := setupHandler(t) + ctx := context.Background() + if err := handler.HandleRecord(ctx, &tap.RecordEvent{ + Live: true, DID: "did:plc:alice", Collection: "app.gainforest.dwc.occurrence", RKey: "occ2", + Action: tap.ActionCreate, CID: "cid1", Record: json.RawMessage(`{"scientificName":"Quercus robur"}`), + }); err != nil { + t.Fatalf("HandleRecord: %v", err) + } + versions, err := db.RecordVersions.ListByURI(ctx, "at://did:plc:alice/app.gainforest.dwc.occurrence/occ2", 0, 500) + if err != nil || len(versions) != 0 { + t.Fatalf("versions without WithRecordHistory = %v (err %v), want none", versions, err) + } +} diff --git a/internal/testutil/db.go b/internal/testutil/db.go index a889aefe..d02304a3 100644 --- a/internal/testutil/db.go +++ b/internal/testutil/db.go @@ -27,6 +27,7 @@ type TestDB struct { LabelDefinitions *repositories.LabelDefinitionsRepository LabelPreferences *repositories.LabelPreferencesRepository Reports *repositories.ReportsRepository + RecordVersions *repositories.RecordVersionsRepository } // SetupTestDB creates an in-memory SQLite database with all migrations applied. @@ -73,6 +74,7 @@ func SetupTestDBWithURL(t *testing.T, databaseURL string) *TestDB { LabelDefinitions: repositories.NewLabelDefinitionsRepository(exec), LabelPreferences: repositories.NewLabelPreferencesRepository(exec), Reports: repositories.NewReportsRepository(exec), + RecordVersions: repositories.NewRecordVersionsRepository(exec), } t.Cleanup(func() {