From 40ac612767d23ec82009dda5f9657c75c41b9a74 Mon Sep 17 00:00:00 2001 From: "turbolytics.io" Date: Sun, 4 Oct 2026 19:08:01 -0400 Subject: [PATCH 1/2] turbine: end the open transaction before a partition drop The drop ran inside whatever transaction the pipeline's connection had open. Between batches that is an idle tick's progress write, begun before the drop took the pass lock. Its snapshot could predate a window pass that had since deleted closed buckets' rows on the manager's own connection and committed; DuckDB then refused the drop's delete of those rows with "Conflict on tuple deletion", and the pipeline stopped over it (#437). The pass lock could not help: it excludes passes from now on, not the one that committed before the snapshot was taken. dropPartitions now commits the open transaction first, so the drop begins one whose snapshot is after every pass the lock excludes. The progress write is rewritten at every commit, so nothing is lost by committing it early. Test: TestTurbine_PartitionDropSucceedsAfterAPassDeletedTheRows drives a real Turbine with a DuckDB dropper over two connections in exactly that order. On main the run stops with the crash log's error; with the fix it continues. internal/core and internal/cli/run pass under -race. Closes #437. --- CHANGELOG.md | 10 +++++ internal/core/partitiondrop_test.go | 66 +++++++++++++++++++++++++++++ internal/core/turbine.go | 14 ++++++ 3 files changed, 90 insertions(+) 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..f3ba6e0b 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,65 @@ 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) + defer cancel() + done := make(chan error, 1) + go func() { _, err := tb.ConsumeLoop(loopCtx, 0); done <- err }() + 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 { From e6bdda121d45441d6260bba8b8b6d31fb5da490c Mon Sep 17 00:00:00 2001 From: "turbolytics.io" Date: Sun, 4 Oct 2026 19:46:18 -0400 Subject: [PATCH 2/2] test: stop the consume loop before the drop test's cleanups run The test left the loop running and let t.Cleanup close its connection and database under it: a use after free in the ADBC driver, which segfaulted on CI's Linux runner and happened to pass on a Mac. The loop is now cancelled and waited for before anything closes. --- internal/core/partitiondrop_test.go | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/internal/core/partitiondrop_test.go b/internal/core/partitiondrop_test.go index f3ba6e0b..4a0a2c6d 100644 --- a/internal/core/partitiondrop_test.go +++ b/internal/core/partitiondrop_test.go @@ -134,9 +134,15 @@ func TestTurbine_PartitionDropSucceedsAfterAPassDeletedTheRows(t *testing.T) { WithPartitionDropper(NewDuckDBPartitionDropper(pipeline, []string{"w"})), WithStateStore(&nopSaver{}, tx)) loopCtx, cancel := context.WithCancel(ctx) - defer cancel() 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 {