Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
72 changes: 72 additions & 0 deletions internal/core/partitiondrop_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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 }
14 changes: 14 additions & 0 deletions internal/core/turbine.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading