Skip to content

Commit 8c2d6e6

Browse files
committed
fix: prune DuckDB fingerprints for missing mirror rows
A checkpoint failure can happen after hard-deleted DuckDB rows have already been removed but before the push boundary fingerprints are finalized. On the next retry those rows are no longer reported as stale, so old cached fingerprints could be copied forward and later suppress a restored or missing-row session push. Delete reconciliation now returns the mirror session IDs that remain after stale cleanup, and incremental skip filtering runs only after cached fingerprints have been pruned against that mirror state. A regression covers the checkpoint-failure retry path without advancing the watermark.
1 parent 25d39eb commit 8c2d6e6

3 files changed

Lines changed: 72 additions & 14 deletions

File tree

internal/duckdb/push.go

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -316,42 +316,45 @@ func (s *Sync) replaceSessionDependents(
316316
func (s *Sync) deleteHardDeletedMirrorSessions(
317317
ctx context.Context, tx *sql.Tx, localSessions []db.Session,
318318
machine string, projects, excludeProjects []string,
319-
) ([]string, error) {
319+
) ([]string, map[string]bool, error) {
320320
localIDs := make(map[string]bool, len(localSessions))
321321
for _, sess := range localSessions {
322322
localIDs[sess.ID] = true
323323
}
324+
mirrorIDs := make(map[string]bool)
324325
rows, err := tx.QueryContext(ctx,
325326
`SELECT id, project FROM sessions WHERE machine = ?`,
326327
machine,
327328
)
328329
if err != nil {
329-
return nil, fmt.Errorf("listing duckdb sessions for deletion reconciliation: %w", err)
330+
return nil, nil, fmt.Errorf("listing duckdb sessions for deletion reconciliation: %w", err)
330331
}
331332
defer rows.Close()
332333
var stale []string
333334
for rows.Next() {
334335
var id, project string
335336
if err := rows.Scan(&id, &project); err != nil {
336-
return nil, fmt.Errorf("scanning duckdb session for deletion reconciliation: %w", err)
337+
return nil, nil, fmt.Errorf("scanning duckdb session for deletion reconciliation: %w", err)
337338
}
339+
mirrorIDs[id] = true
338340
if !projectInSyncScope(project, projects, excludeProjects) {
339341
continue
340342
}
341343
if !localIDs[id] {
342344
stale = append(stale, id)
345+
delete(mirrorIDs, id)
343346
}
344347
}
345348
if err := rows.Err(); err != nil {
346-
return nil, err
349+
return nil, nil, err
347350
}
348351
sort.Strings(stale)
349352
for _, id := range stale {
350353
if err := s.deleteMirrorSession(ctx, tx, id); err != nil {
351-
return nil, err
354+
return nil, nil, err
352355
}
353356
}
354-
return stale, nil
357+
return stale, mirrorIDs, nil
355358
}
356359

357360
func (s *Sync) deleteMirrorSession(

internal/duckdb/sync.go

Lines changed: 23 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -286,19 +286,13 @@ func (s *Sync) Push(
286286
if err != nil {
287287
return result, err
288288
}
289-
sessions = filterUnchangedSessions(sessions, priorFingerprints, sessionFingerprints)
290-
result.Diagnostics.SkippedUnchangedSessions = skippedPushSessions(
291-
candidateSessions, sessions,
292-
)
293289
}
294-
sort.Slice(sessions, func(i, j int) bool {
295-
return sessions[i].ID < sessions[j].ID
296-
})
297290

298291
var staleIDs []string
292+
var mirrorSessionIDs map[string]bool
299293
err = s.withDuckTx(ctx, "delete hard-deleted sessions", func(tx *sql.Tx) error {
300294
var txErr error
301-
staleIDs, txErr = s.deleteHardDeletedMirrorSessions(
295+
staleIDs, mirrorSessionIDs, txErr = s.deleteHardDeletedMirrorSessions(
302296
ctx, tx, allLocalSessions, s.machine, s.projects, s.excludeProjects,
303297
)
304298
return txErr
@@ -310,6 +304,16 @@ func (s *Sync) Push(
310304
delete(priorFingerprints, id)
311305
}
312306
result.Diagnostics.DeletedStaleSessions = len(staleIDs)
307+
if !full {
308+
pruneMissingMirrorFingerprints(priorFingerprints, mirrorSessionIDs)
309+
sessions = filterUnchangedSessions(sessions, priorFingerprints, sessionFingerprints)
310+
result.Diagnostics.SkippedUnchangedSessions = skippedPushSessions(
311+
candidateSessions, sessions,
312+
)
313+
}
314+
sort.Slice(sessions, func(i, j int) bool {
315+
return sessions[i].ID < sessions[j].ID
316+
})
313317

314318
pushed := make([]db.Session, 0, len(sessions))
315319
for start := 0; start < len(sessions); start += duckSessionPushBatchSize {
@@ -877,6 +881,17 @@ func filterUnchangedSessions(
877881
return out
878882
}
879883

884+
func pruneMissingMirrorFingerprints(
885+
priorFingerprints map[string]string,
886+
mirrorSessionIDs map[string]bool,
887+
) {
888+
for id := range priorFingerprints {
889+
if !mirrorSessionIDs[id] {
890+
delete(priorFingerprints, id)
891+
}
892+
}
893+
}
894+
880895
func previousLocalSyncTimestamp(value string) (string, error) {
881896
if value == "" {
882897
return "", nil

internal/duckdb/sync_fastpath_test.go

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -352,6 +352,46 @@ func TestSyncCheckpointFailureDoesNotAdvanceWatermark(t *testing.T) {
352352
assert.Empty(t, watermark)
353353
}
354354

355+
func TestSyncCheckpointFailureAfterHardDeleteDoesNotKeepStaleFingerprint(t *testing.T) {
356+
ctx := context.Background()
357+
local := newLocalDB(t)
358+
fixture := seedDuckDBSyncFixture(t, local)
359+
syncer := newInMemoryTestSync(t, local, SyncOptions{})
360+
361+
first, err := syncer.Push(ctx, true, nil)
362+
require.NoError(t, err)
363+
require.Equal(t, 2, first.SessionsPushed)
364+
firstWatermark, err := local.GetSyncState(lastPushStateKey)
365+
require.NoError(t, err)
366+
fingerprints, err := readSyncFingerprintsWithKey(local, lastPushBoundaryStateKey)
367+
require.NoError(t, err)
368+
require.Contains(t, fingerprints, fixture.betaID)
369+
370+
require.NoError(t, local.SoftDeleteSession(fixture.betaID))
371+
deleted, err := local.DeleteSessionIfTrashed(fixture.betaID)
372+
require.NoError(t, err)
373+
require.EqualValues(t, 1, deleted)
374+
policy := &checkpointSpy{err: errors.New("checkpoint failed")}
375+
syncer.maintenance = policy
376+
377+
_, err = syncer.Push(ctx, false, nil)
378+
require.ErrorContains(t, err, "checkpoint failed")
379+
assert.Equal(t, 1, policy.calls)
380+
assertDuckDBCountWhere(t, syncer.DB(), "sessions", "id = ?", fixture.betaID, 0)
381+
watermarkAfterFailure, err := local.GetSyncState(lastPushStateKey)
382+
require.NoError(t, err)
383+
assert.Equal(t, firstWatermark, watermarkAfterFailure)
384+
385+
policy.err = nil
386+
retry, err := syncer.Push(ctx, false, nil)
387+
require.NoError(t, err)
388+
assert.Zero(t, retry.SessionsPushed)
389+
assert.Equal(t, 1, policy.calls, "retry should only repair local sync state")
390+
fingerprints, err = readSyncFingerprintsWithKey(local, lastPushBoundaryStateKey)
391+
require.NoError(t, err)
392+
assert.NotContains(t, fingerprints, fixture.betaID)
393+
}
394+
355395
func TestSyncCheckpointPolicySkipsRemoteQuackTargets(t *testing.T) {
356396
ctx := context.Background()
357397
local := newLocalDB(t)

0 commit comments

Comments
 (0)