diff --git a/pkg/atproto/atproto.go b/pkg/atproto/atproto.go index 02960fe5..30b5137b 100644 --- a/pkg/atproto/atproto.go +++ b/pkg/atproto/atproto.go @@ -349,6 +349,19 @@ func (atsync *ATProtoSynchronizer) directory(cached bool) identity.Directory { return atsync.PLCDirectory } +// resetIdentCache throws away every cached identity resolution. For when the +// node has knowingly skipped firehose it can never replay: any number of +// #identity events may sit in that gap, and a cache entry from before it +// agrees with the equally-stale repo row — hiding exactly the drift the next +// sweep's head check is being asked to find. Starting the cache over makes +// that sweep's resolutions authoritative. (A restart gets this for free; the +// cache lives in process memory.) +func (atsync *ATProtoSynchronizer) resetIdentCache() { + atsync.dirMu.Lock() + defer atsync.dirMu.Unlock() + atsync.CachedPLCDirectory = nil +} + // purgeIdentCache drops a DID's cached identity, so the next cached resolve // re-reads the directory instead of serving an entry we just proved stale. func (atsync *ATProtoSynchronizer) purgeIdentCache(ctx context.Context, did string) { diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index 3742d0e5..9cf62bf2 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -263,6 +263,13 @@ func (atsync *ATProtoSynchronizer) dropStaleCursor(ctx context.Context, cursor * if !cursor.dropIfStale(ctx, atsync.firehoseReplayWindow()) { return } + // The skipped span may hold #identity events as well as commits. Commits + // the head check finds by asking hosts, but identity drift it can only see + // by resolving — and a cached resolution from before the gap agrees with + // the equally-stale repo row, hiding the change until the cache entry ages + // out. Starting the cache over makes the healing sweep's resolutions + // authoritative. + atsync.resetIdentCache() if atsync.StatefulDB == nil { // No index to sweep (tests that exercise the firehose on its own). return diff --git a/pkg/atproto/headcheck_test.go b/pkg/atproto/headcheck_test.go index 473fe572..0eb17608 100644 --- a/pkg/atproto/headcheck_test.go +++ b/pkg/atproto/headcheck_test.go @@ -5,6 +5,8 @@ import ( "fmt" "testing" + "github.com/bluesky-social/indigo/atproto/identity" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/devenv" "stream.place/streamplace/pkg/placestream" @@ -130,3 +132,40 @@ func TestSweepRefreshesDriftedIdentity(t *testing.T) { require.Equal(t, indexed.Version, healed.Version, "identity refresh left the sync state alone") require.Equal(t, indexed.BackfillDone, healed.BackfillDone) } + +// TestResetIdentCacheDropsCachedResolutions proves the reset is real: a warm +// cache serves hits until resetIdentCache (what dropping a stale cursor +// calls), after which the next resolve goes back to the directory. Combined +// with TestSweepRefreshesDriftedIdentity above — a resolution that disagrees +// with the row heals it — this is the full path by which an identity change +// hidden inside a skipped replay reaches the row: the reset makes the healing +// sweep's resolutions fresh, and a fresh resolution that disagrees with the +// row rewrites it. +// +// (A true end-to-end rename is not observable in the dev environment: handle +// verification cannot complete there, so every account resolves as +// handle.invalid and a rename changes nothing the directory reports.) +func TestResetIdentCacheDropsCachedResolutions(t *testing.T) { + dev := devenv.WithDevEnv(t) + ctx := context.Background() + atsync, _ := backfillTestSynchronizer(t, dev) + user := dev.CreateAccount(t) + + _, err := atsync.resolveIdent(ctx, user.DID, true) + require.NoError(t, err) + did, err := syntax.ParseDID(user.DID) + require.NoError(t, err) + cd, ok := atsync.directory(true).(*identity.CacheDirectory) + require.True(t, ok) + _, hit, err := cd.LookupDIDWithCacheState(ctx, did) + require.NoError(t, err) + require.True(t, hit, "the resolve above warmed the cache") + + atsync.resetIdentCache() + + cd, ok = atsync.directory(true).(*identity.CacheDirectory) + require.True(t, ok) + _, hit, err = cd.LookupDIDWithCacheState(ctx, did) + require.NoError(t, err) + require.False(t, hit, "the reset cache resolves from the directory again") +}