From e9e27f914dd799459b7e7dadd88ee8c2b3a69980 Mon Sep 17 00:00:00 2001 From: Jorge Manrubia Date: Fri, 18 Sep 2026 08:51:48 +0200 Subject: [PATCH 1/2] An acknowledgement id is written by the statement that acknowledges MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The trigger read only the row as it was, so a statement could write an acknowledgement id and leave the row exposed. That puts an id in the receipt that no worker ever reported — the half of "with the acknowledgement or never" that the old row cannot see. Three things now have to hold together: nothing recorded before, the row still waiting to be acknowledged, and this same statement acknowledging it. Ack already did all three, so nothing it does changes. --- internal/connector/dispatch_lifecycle_test.go | 43 +++++++++++++++++++ internal/connector/ledger.go | 9 +++- 2 files changed, 50 insertions(+), 2 deletions(-) diff --git a/internal/connector/dispatch_lifecycle_test.go b/internal/connector/dispatch_lifecycle_test.go index 99b6810f3..d74918912 100644 --- a/internal/connector/dispatch_lifecycle_test.go +++ b/internal/connector/dispatch_lifecycle_test.go @@ -549,3 +549,46 @@ func TestOneWriteCannotWithdrawAndComplete(t *testing.T) { assert.Nil(t, withdrawn) assert.Equal(t, "exposed", f.rowContext(ctx, t, 1).Delivery) } + +// An acknowledgement id belongs to the statement that acknowledges. Writing it +// onto a row that stays exposed would leave an id in the receipt that no worker +// ever reported, which is the half of "with the acknowledgement or never" that +// checking only the old row cannot see. +func TestAnAcknowledgementIDIsNotWrittenWithoutTheAcknowledgement(t *testing.T) { + ctx := context.Background() + f := newDispatchFixture(t) + _, err := f.ledger.db.ExecContext(ctx, `UPDATE task_events SET delivery = 'exposed', exposed_at = 'launch' WHERE event_id = 1`) + require.NoError(t, err) + + _, err = f.ledger.db.ExecContext(ctx, `UPDATE task_events SET ack_id = 99 WHERE task_id = ? AND event_id = 1`, f.grant.ID) + + require.Error(t, err) + assert.Contains(t, err.Error(), "an acknowledgement id is written with the acknowledgement, once") + var ackID *int64 + require.NoError(t, f.ledger.db.QueryRowContext(ctx, `SELECT ack_id FROM task_events WHERE task_id = ? AND event_id = 1`, f.grant.ID).Scan(&ackID)) + assert.Nil(t, ackID) + assert.Equal(t, "exposed", f.rowContext(ctx, t, 1).Delivery) +} + +// And the acknowledgement itself still writes one: the rule narrows what may +// write an id, not whether Ack can. +func TestAcknowledgingWritesTheIDItWasGiven(t *testing.T) { + ctx := context.Background() + f := newDispatchFixture(t) + _, err := f.ledger.db.ExecContext(ctx, `UPDATE task_events SET delivery = 'exposed', exposed_at = 'launch' WHERE event_id = 1`) + require.NoError(t, err) + + d, err := f.ledger.Dispatch(ctx, f.grant.Token, adapterAgentID) + require.NoError(t, err) + _, _, err = d.Get(ctx, 1) + require.NoError(t, err) + ackID := int64(99) + _, err = d.Ack(ctx, 1, &ackID) + require.NoError(t, err) + + var got *int64 + require.NoError(t, f.ledger.db.QueryRowContext(ctx, `SELECT ack_id FROM task_events WHERE task_id = ? AND event_id = 1`, f.grant.ID).Scan(&got)) + require.NotNil(t, got) + assert.Equal(t, int64(99), *got) + assert.Equal(t, "delivered", f.rowContext(ctx, t, 1).Delivery) +} diff --git a/internal/connector/ledger.go b/internal/connector/ledger.go index 2d2821f61..975ebea8f 100644 --- a/internal/connector/ledger.go +++ b/internal/connector/ledger.go @@ -722,10 +722,15 @@ BEGIN END; -- The acknowledgement settles with the delivery: the id a worker points at is --- written when it acknowledges, or never. +-- written by the statement that acknowledges, or never. Three things have to +-- hold together — nothing was recorded before, the row was still waiting to be +-- acknowledged, and this same statement acknowledges it. Writing the id while +-- leaving the row exposed would put an id in the receipt that no worker ever +-- reported. CREATE TRIGGER task_events_acknowledgement_settles_once BEFORE UPDATE OF ack_id ON task_events -WHEN NEW.ack_id IS NOT OLD.ack_id AND (OLD.ack_id IS NOT NULL OR OLD.delivery <> 'exposed') +WHEN NEW.ack_id IS NOT OLD.ack_id + AND (OLD.ack_id IS NOT NULL OR OLD.delivery <> 'exposed' OR NEW.delivery <> 'delivered') BEGIN SELECT RAISE(ABORT, 'an acknowledgement id is written with the acknowledgement, once'); END; From bf91e8e0e93e9c86bdd5905de8f5ef991e67ed8f Mon Sep 17 00:00:00 2001 From: Jorge Manrubia Date: Fri, 18 Sep 2026 09:10:55 +0200 Subject: [PATCH 2/2] Replace the acknowledgement trigger in a migration of its own The first version of this edited migration 5, which had already shipped in #736. A ledger already at version 5 would have kept the loose trigger for good: migrate skips what it has applied, so an edit to a shipped migration reaches no existing ledger. The contract at the top of the list says so. Migration 5 is back exactly as it shipped, and migration 6 drops and recreates the trigger. The new test walks the upgrade that actually happens: a ledger built from migrations 1 through 5 alone, opened, and then held to the tighter rule. --- internal/connector/dispatch_lifecycle_test.go | 58 +++++++++++++++++++ internal/connector/ledger.go | 27 ++++++--- 2 files changed, 78 insertions(+), 7 deletions(-) diff --git a/internal/connector/dispatch_lifecycle_test.go b/internal/connector/dispatch_lifecycle_test.go index d74918912..5de670d71 100644 --- a/internal/connector/dispatch_lifecycle_test.go +++ b/internal/connector/dispatch_lifecycle_test.go @@ -2,7 +2,10 @@ package connector import ( "context" + "database/sql" "fmt" + "os" + "path/filepath" "slices" "strings" "testing" @@ -592,3 +595,58 @@ func TestAcknowledgingWritesTheIDItWasGiven(t *testing.T) { assert.Equal(t, int64(99), *got) assert.Equal(t, "delivered", f.rowContext(ctx, t, 1).Delivery) } + +// A ledger born under migration 5 carries that migration's trigger, which +// read only the row as it was. Editing migration 5 would have left every such +// ledger with it, because migrate skips what it has already applied — so the +// replacement is migration 6, and this is the upgrade actually happening. +func TestAnExistingLedgerGetsTheTighterAcknowledgementRule(t *testing.T) { + ctx := context.Background() + path := filepath.Join(t.TempDir(), "state", "connector.db") + require.NoError(t, os.MkdirAll(filepath.Dir(path), 0o700)) + + // A ledger as the previous version wrote it: migrations 1 through 5 and + // nothing after them. + old, err := sql.Open("sqlite", ledgerDSN(path, true)) + require.NoError(t, err) + _, err = old.ExecContext(ctx, `CREATE TABLE schema_migrations (version INTEGER PRIMARY KEY, applied_at TEXT NOT NULL)`) + require.NoError(t, err) + for i := 0; i < 5; i++ { + _, err = old.ExecContext(ctx, migrations[i]) + require.NoError(t, err, "migration %d", i+1) + _, err = old.ExecContext(ctx, `INSERT INTO schema_migrations (version, applied_at) VALUES (?, 'then')`, i+1) + require.NoError(t, err) + } + require.NoError(t, old.Close()) + // The connector's own open makes the file private; a raw sql.Open does + // not, and the privacy check refuses what it finds. + for _, name := range []string{path, path + "-wal", path + "-shm"} { + if _, err := os.Stat(name); err == nil { + require.NoError(t, os.Chmod(name, 0o600)) + } + } + + ledger, err := OpenLedger(path) + require.NoError(t, err) + t.Cleanup(func() { _ = ledger.Close() }) + version, err := ledger.SchemaVersion(ctx) + require.NoError(t, err) + assert.Equal(t, len(migrations), version, "the upgrade ran") + + // The same refusal the fresh-ledger test pins, on a ledger that was not + // born with it. + for _, id := range []int64{1, 2} { + seenRecord(t, ledger, id) + _, err := ledger.Admission().Commit(ctx, admittedVerdict(id, 0, "recording:10304028989")) + require.NoError(t, err) + } + grant, err := ledger.CreateTask(ctx, []int64{1, 2}) + require.NoError(t, err) + _, err = ledger.db.ExecContext(ctx, `UPDATE task_events SET delivery = 'exposed', exposed_at = 'launch' WHERE event_id = 1`) + require.NoError(t, err) + + _, err = ledger.db.ExecContext(ctx, `UPDATE task_events SET ack_id = 99 WHERE task_id = ? AND event_id = 1`, grant.ID) + + require.Error(t, err) + assert.Contains(t, err.Error(), "an acknowledgement id is written with the acknowledgement, once") +} diff --git a/internal/connector/ledger.go b/internal/connector/ledger.go index 975ebea8f..97ae4c4f4 100644 --- a/internal/connector/ledger.go +++ b/internal/connector/ledger.go @@ -722,15 +722,10 @@ BEGIN END; -- The acknowledgement settles with the delivery: the id a worker points at is --- written by the statement that acknowledges, or never. Three things have to --- hold together — nothing was recorded before, the row was still waiting to be --- acknowledged, and this same statement acknowledges it. Writing the id while --- leaving the row exposed would put an id in the receipt that no worker ever --- reported. +-- written when it acknowledges, or never. CREATE TRIGGER task_events_acknowledgement_settles_once BEFORE UPDATE OF ack_id ON task_events -WHEN NEW.ack_id IS NOT OLD.ack_id - AND (OLD.ack_id IS NOT NULL OR OLD.delivery <> 'exposed' OR NEW.delivery <> 'delivered') +WHEN NEW.ack_id IS NOT OLD.ack_id AND (OLD.ack_id IS NOT NULL OR OLD.delivery <> 'exposed') BEGIN SELECT RAISE(ABORT, 'an acknowledgement id is written with the acknowledgement, once'); END; @@ -800,6 +795,24 @@ WHEN (OLD.delivery = 'admitted' AND NEW.delivery IN ('delivered', 'completed')) BEGIN SELECT RAISE(ABORT, 'a worker acknowledges and completes what it pulled; anything else is the dispatcher settling a completed record'); END; +`, + // 6. The acknowledgement id settles with the acknowledgement. + // + // Migration 5 shipped a trigger that read only the row as it was, so a + // statement could write the id and leave the row exposed — an id in the + // receipt that no worker ever reported. A ledger already at version 5 + // keeps that trigger, so replacing it is its own migration rather than an + // edit to one that has shipped. + ` +DROP TRIGGER task_events_acknowledgement_settles_once; + +CREATE TRIGGER task_events_acknowledgement_settles_once +BEFORE UPDATE OF ack_id ON task_events +WHEN NEW.ack_id IS NOT OLD.ack_id + AND (OLD.ack_id IS NOT NULL OR OLD.delivery <> 'exposed' OR NEW.delivery <> 'delivered') +BEGIN + SELECT RAISE(ABORT, 'an acknowledgement id is written with the acknowledgement, once'); +END; `, }