Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions .agents/skills/hyperindex/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
19 changes: 19 additions & 0 deletions .agents/skills/hyperindex/references/schema-reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down Expand Up @@ -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.
Expand Down
5 changes: 5 additions & 0 deletions .changes/unreleased/add-record-version-history.yaml
Original file line number Diff line number Diff line change
@@ -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
1 change: 1 addition & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
3 changes: 3 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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: "<id>")`, 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.

Expand Down
18 changes: 18 additions & 0 deletions cmd/hyperindex/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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 != "" {
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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{
Expand Down
8 changes: 8 additions & 0 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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", ""),
Expand Down
38 changes: 36 additions & 2 deletions internal/database/migrations/migrations_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
Expand Down Expand Up @@ -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")
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
DROP TABLE IF EXISTS record_version;
Original file line number Diff line number Diff line change
@@ -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,
Comment on lines +5 to +8

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Remove version rows when purging an actor

When an admin invokes PurgeActor or Tap receives a deleted, deactivated, suspended, or taken-down identity, RecordsRepository.PurgeActorData deletes only record and actor; this independent table has no cascade and neither purge path removes its rows. Because the new public recordHistory resolver reads this table directly, tracked record bodies remain publicly retrievable after an operation documented as removing all indexed data for the DID. Include these rows in the same purge transaction.

Useful? React with 👍 / 👎.

-- Dedupe identity: the CID, `sha256:<hex>` of the body when Tap omits the
-- CID, or `delete:<cid>` 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);
Comment on lines +20 to +22

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve recreated versions after a deletion

When a URI is deleted and later recreated with the same content, ATProto produces the same content-addressed CID, but this URI/CID uniqueness constraint discards both the later create and its subsequent delete:<cid> tombstone. CID-less records have the same problem because identical bodies and all empty-CID tombstones reuse their keys. The resulting history ends at the first delete even while the record exists again, so deduplication needs an event or lifecycle identity that distinguishes redelivery from a legitimate recreate.

AGENTS.md reference: AGENTS.md:L94-L94

Useful? React with 👍 / 👎.

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);
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
DROP TABLE IF EXISTS record_version;
23 changes: 23 additions & 0 deletions internal/database/migrations/sqlite/015_add_record_version.up.sql
Original file line number Diff line number Diff line change
@@ -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:<hex>` of the body when Tap omits the
-- CID, or `delete:<cid>` 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);
Loading
Loading