Skip to content

Commit 91a91cb

Browse files
committed
fix: repair missing DuckDB mirror sessions
Incremental DuckDB pushes could prune a cached fingerprint after discovering a missing mirror row, then advance the watermark without adding the unchanged local session back to the push candidate set. That left the mirror missing the row until a full push or later local modification. Treat pruned fingerprints that still have in-scope local sessions as repair candidates before fingerprint filtering, so the normal push path reinserts the session and dependents. The regression test deletes a mirrored row after a successful full push and verifies the next incremental push restores it.
1 parent 8c2d6e6 commit 91a91cb

2 files changed

Lines changed: 73 additions & 8 deletions

File tree

internal/duckdb/sync.go

Lines changed: 43 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -270,13 +270,11 @@ func (s *Sync) Push(
270270
if err != nil {
271271
return result, fmt.Errorf("listing local sessions: %w", err)
272272
}
273-
result.Diagnostics.LocalSessions = countPushSessions(allLocalSessions)
274-
sessionFingerprints, err := s.sessionFingerprints(ctx, sessions)
275-
if err != nil {
276-
return result, err
273+
allLocalSessionByID := make(map[string]db.Session, len(allLocalSessions))
274+
for _, sess := range allLocalSessions {
275+
allLocalSessionByID[sess.ID] = sess
277276
}
278-
candidateSessions := append([]db.Session(nil), sessions...)
279-
result.Diagnostics.CandidateSessions = countPushSessions(candidateSessions)
277+
result.Diagnostics.LocalSessions = countPushSessions(allLocalSessions)
280278
priorFingerprints := map[string]string{}
281279
if !full {
282280
priorFingerprints, err = readSyncFingerprintsWithKey(
@@ -305,7 +303,20 @@ func (s *Sync) Push(
305303
}
306304
result.Diagnostics.DeletedStaleSessions = len(staleIDs)
307305
if !full {
308-
pruneMissingMirrorFingerprints(priorFingerprints, mirrorSessionIDs)
306+
missingMirrorIDs := pruneMissingMirrorFingerprints(
307+
priorFingerprints, mirrorSessionIDs,
308+
)
309+
sessions = appendMissingMirrorRepairCandidates(
310+
sessions, sessionByID, allLocalSessionByID, missingMirrorIDs,
311+
)
312+
}
313+
sessionFingerprints, err := s.sessionFingerprints(ctx, sessions)
314+
if err != nil {
315+
return result, err
316+
}
317+
candidateSessions := append([]db.Session(nil), sessions...)
318+
result.Diagnostics.CandidateSessions = countPushSessions(candidateSessions)
319+
if !full {
309320
sessions = filterUnchangedSessions(sessions, priorFingerprints, sessionFingerprints)
310321
result.Diagnostics.SkippedUnchangedSessions = skippedPushSessions(
311322
candidateSessions, sessions,
@@ -884,12 +895,36 @@ func filterUnchangedSessions(
884895
func pruneMissingMirrorFingerprints(
885896
priorFingerprints map[string]string,
886897
mirrorSessionIDs map[string]bool,
887-
) {
898+
) []string {
899+
var missing []string
888900
for id := range priorFingerprints {
889901
if !mirrorSessionIDs[id] {
890902
delete(priorFingerprints, id)
903+
missing = append(missing, id)
904+
}
905+
}
906+
sort.Strings(missing)
907+
return missing
908+
}
909+
910+
func appendMissingMirrorRepairCandidates(
911+
sessions []db.Session,
912+
sessionByID map[string]db.Session,
913+
allLocalSessionByID map[string]db.Session,
914+
missingMirrorIDs []string,
915+
) []db.Session {
916+
for _, id := range missingMirrorIDs {
917+
sess, ok := allLocalSessionByID[id]
918+
if !ok {
919+
continue
920+
}
921+
if _, ok := sessionByID[id]; ok {
922+
continue
891923
}
924+
sessionByID[id] = sess
925+
sessions = append(sessions, sess)
892926
}
927+
return sessions
893928
}
894929

895930
func previousLocalSyncTimestamp(value string) (string, error) {

internal/duckdb/sync_fastpath_test.go

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -392,6 +392,36 @@ func TestSyncCheckpointFailureAfterHardDeleteDoesNotKeepStaleFingerprint(t *test
392392
assert.NotContains(t, fingerprints, fixture.betaID)
393393
}
394394

395+
func TestSyncMissingMirrorRowWithUnchangedLocalSessionIsRepaired(t *testing.T) {
396+
ctx := context.Background()
397+
local := newLocalDB(t)
398+
fixture := seedDuckDBSyncFixture(t, local)
399+
syncer := newInMemoryTestSync(t, local, SyncOptions{})
400+
policy := &checkpointSpy{}
401+
syncer.maintenance = policy
402+
403+
first, err := syncer.Push(ctx, true, nil)
404+
require.NoError(t, err)
405+
require.Equal(t, 2, first.SessionsPushed)
406+
assert.Equal(t, 1, policy.calls)
407+
require.NoError(t, syncer.withDuckTx(ctx, "test delete mirror session", func(tx *sql.Tx) error {
408+
return syncer.deleteMirrorSession(ctx, tx, fixture.betaID)
409+
}))
410+
assertDuckDBCountWhere(t, syncer.DB(), "sessions", "id = ?", fixture.betaID, 0)
411+
412+
second, err := syncer.Push(ctx, false, nil)
413+
require.NoError(t, err)
414+
415+
assert.Equal(t, 1, second.SessionsPushed)
416+
assert.Equal(t, 1, second.MessagesPushed)
417+
assert.Equal(t, 2, policy.calls)
418+
assertDuckDBCountWhere(t, syncer.DB(), "sessions", "id = ?", fixture.betaID, 1)
419+
assertDuckDBCountWhere(t, syncer.DB(), "messages", "session_id = ?", fixture.betaID, 1)
420+
fingerprints, err := readSyncFingerprintsWithKey(local, lastPushBoundaryStateKey)
421+
require.NoError(t, err)
422+
assert.Contains(t, fingerprints, fixture.betaID)
423+
}
424+
395425
func TestSyncCheckpointPolicySkipsRemoteQuackTargets(t *testing.T) {
396426
ctx := context.Background()
397427
local := newLocalDB(t)

0 commit comments

Comments
 (0)