diff --git a/CHANGELOG.md b/CHANGELOG.md index 0aa4623d..668f7577 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -35,6 +35,16 @@ ### Fixed +- **A `partition_owned` worker no longer crashes dropping the partitions + a failed session took from it.** The drop ran inside whatever transaction + the pipeline's connection had open, begun by an idle tick's progress + write before the drop took its locks. Its snapshot could predate a window + pass that had since deleted closed buckets' rows on the manager's + connection and committed, and DuckDB refused the drop's delete with + "Conflict on tuple deletion"; the pipeline stopped over it, and in a + two-worker stack a network blip crashed one worker while the other + stalled (#436). The drop now ends that transaction first, so it begins one + whose snapshot is after every pass its lock excludes. (#437) - **A windowed worker that loses its group session no longer stops closing windows after it rejoins.** A lost partition stayed in the worker's watermark minimum until it was assigned back; when the rejoin gave it to diff --git a/internal/core/partitiondrop_test.go b/internal/core/partitiondrop_test.go index f7936501..4a0a2c6d 100644 --- a/internal/core/partitiondrop_test.go +++ b/internal/core/partitiondrop_test.go @@ -4,8 +4,12 @@ import ( "context" "testing" + "time" + + "github.com/apache/arrow-adbc/go/adbc" "github.com/apache/arrow-go/v18/arrow/array" "github.com/turbolytics/sql-flow/internal/coverage" + "github.com/turbolytics/sql-flow/internal/duckdb" "github.com/zeebo/assert" ) @@ -83,3 +87,71 @@ func TestDuckDBPartitionDropper_DropAllEmptiesEveryOwnedTable(t *testing.T) { assert.Equal(t, int64(0), rdr.Record().Column(0).(*array.Int64).Value(0)) assert.Equal(t, int64(1), rdr.Record().Column(1).(*array.Int64).Value(0)) } + +// The drop runs on the pipeline's connection, whose transaction may have +// been open since before the manager's pass deleted the same rows on its +// own connection and committed. DuckDB then refuses the drop's delete with +// "Conflict on tuple deletion", and the pipeline stopped over it (#437). +// The drop must succeed whatever the pipeline's transaction has seen. +// +// pipeline conn (autocommit off) manager conn +// ---------------------------------------------------------- +// INSERT ... ; COMMIT +// UPDATE progress (opens the next tx) +// DELETE closed rows; COMMIT +// lose partitions -> drop <-- conflict, unless the tx ends first +func TestTurbine_PartitionDropSucceedsAfterAPassDeletedTheRows(t *testing.T) { + coverage.Covers(t, "source.kafka") + ctx := context.Background() + db, err := duckdb.OpenPath(ctx, "") + assert.NoError(t, err) + t.Cleanup(func() { db.Close() }) + pipeline, err := db.Connect(ctx) + assert.NoError(t, err) + t.Cleanup(func() { pipeline.Close() }) + manager, err := db.Connect(ctx) + assert.NoError(t, err) + t.Cleanup(func() { manager.Close() }) + + exec(t, pipeline, `CREATE TABLE w (minute TIMESTAMPTZ, kafka_partition INTEGER, n INTEGER)`) + exec(t, pipeline, `CREATE TABLE progress (seen TIMESTAMPTZ)`) + exec(t, pipeline, `INSERT INTO progress VALUES (now())`) + assert.NoError(t, pipeline.(adbc.PostInitOptions).SetOption(adbc.OptionKeyAutoCommit, adbc.OptionValueDisabled)) + tx := pipeline.(stateTx) + + // A batch: rows for partitions 3 and 4, committed. + exec(t, pipeline, `INSERT INTO w VALUES (TIMESTAMPTZ '2026-10-04 19:33:00+00', 3, 1), (TIMESTAMPTZ '2026-10-04 19:33:00+00', 4, 1)`) + assert.NoError(t, tx.Commit(ctx)) + // An idle tick's progress write: opens the pipeline's next transaction. + exec(t, pipeline, `UPDATE progress SET seen = now()`) + // The manager's pass closes the minute: deletes its rows, commits. + exec(t, manager, `DELETE FROM w WHERE minute = TIMESTAMPTZ '2026-10-04 19:33:00+00'`) + + // The session is lost: the turbine drops the partitions. + src := newOwnerSource() + w := NewWatermarks([]WindowSpec{ownedSpec}, nil) + tb := newWindowedTurbine(src, &fakeHandler{}, &fakeSink{}, 100, w, + WithPartitionDropper(NewDuckDBPartitionDropper(pipeline, []string{"w"})), + WithStateStore(&nopSaver{}, tx)) + loopCtx, cancel := context.WithCancel(ctx) + done := make(chan error, 1) + go func() { _, err := tb.ConsumeLoop(loopCtx, 0); done <- err }() + // The loop must be gone before the cleanups close its connection and + // database: a close under a running loop is a use after free in the + // ADBC driver, which segfaulted on CI. + defer func() { + cancel() + <-done + }() + src.assigned(map[string][]int32{"t": {3, 4, 5}}) + src.lost(map[string][]int32{"t": {3, 4, 5}}) // blocks until dropped + select { + case err := <-done: + t.Fatalf("the run stopped: %v", err) + case <-time.After(200 * time.Millisecond): + } +} + +type nopSaver struct{} + +func (nopSaver) Save(context.Context, *Marks) error { return nil } diff --git a/internal/core/turbine.go b/internal/core/turbine.go index 3dc161a8..70ef9ac8 100644 --- a/internal/core/turbine.go +++ b/internal/core/turbine.go @@ -916,6 +916,20 @@ func (t *Turbine) dropPartitions(ctx context.Context, parts map[string][]int32) t.lock.Lock() defer t.lock.Unlock() + // The drop runs between batches, so the connection's open transaction, + // if any, holds only an idle tick's progress write, begun before this + // lock was taken. Its snapshot can predate a window pass that has since + // deleted closed buckets' rows on the manager's connection and committed; + // DuckDB then refuses the drop's delete of those rows as a conflict on + // tuple deletion, and the pipeline stopped over it (#437). Committed + // here, so the drop begins a transaction whose snapshot is after every + // pass this lock excludes. The progress write is rewritten every commit, + // so committing it early loses nothing. + if t.dropper != nil && t.stateTx != nil { + if err := t.stateTx.Commit(ctx); err != nil { + return errs.Wrap(errs.CodeStateCommitFailed, err, "ending the transaction before a partition drop") + } + } fail := func(err error) error { if t.stateTx != nil { if rbErr := t.stateTx.Rollback(context.WithoutCancel(ctx)); rbErr != nil {