fix(11-08): refuse unmanaged after-commit work
This commit is contained in:
@@ -337,8 +337,8 @@ func TestSyncGates(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// TestSyncAfterCommit covers D-20: a write syncs only after it commits
|
// TestSyncAfterCommit covers D-20: a write syncs only after it commits
|
||||||
// (lagoon.Transaction, a single statement's implicit transaction, a plain
|
// (lagoon.Transaction and a single statement's implicit transaction), a
|
||||||
// gorm transaction through its handle), a rollback syncs nothing, and the
|
// foreign GORM transaction is refused, a rollback syncs nothing, and the
|
||||||
// document is built from a reload on a clean statement.
|
// document is built from a reload on a clean statement.
|
||||||
func TestSyncAfterCommit(t *testing.T) {
|
func TestSyncAfterCommit(t *testing.T) {
|
||||||
env := newSyncEnv(t, nil, true)
|
env := newSyncEnv(t, nil, true)
|
||||||
@@ -403,20 +403,22 @@ func TestSyncAfterCommit(t *testing.T) {
|
|||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
t.Run("plain_gorm_transaction_syncs_through_the_tx", func(t *testing.T) {
|
t.Run("plain_gorm_transaction_rollback_sends_nothing", func(t *testing.T) {
|
||||||
|
rollback := errors.New("roll back plain transaction")
|
||||||
err := env.gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
err := env.gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||||
if err := tx.Create(&Doc{Title: "plain-tx", OwnerID: 1}).Error; err != nil {
|
if err := tx.Create(&Doc{Title: "plain-tx", OwnerID: 1}).Error; err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
calls := env.eng.take()
|
if calls := env.eng.take(); len(calls) != 0 {
|
||||||
if len(calls) != 1 || calls[0].docs[0]["title"] != "plain-tx" {
|
t.Errorf("inside a plain transaction: calls = %+v, want none", calls)
|
||||||
t.Errorf("inside a plain transaction: calls = %+v, want an immediate sync from the uncommitted row", calls)
|
|
||||||
}
|
}
|
||||||
// The transaction is still usable.
|
return rollback
|
||||||
return tx.Create(&Plain{Name: "after"}).Error
|
|
||||||
})
|
})
|
||||||
if err != nil {
|
if !errors.Is(err, rollback) {
|
||||||
t.Fatal(err)
|
t.Fatalf("err = %v", err)
|
||||||
|
}
|
||||||
|
if calls := env.eng.take(); len(calls) != 0 {
|
||||||
|
t.Fatalf("after rollback: calls = %+v, want none", calls)
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ Postgres data layer: the shared GORM connection, per-plugin migrations, model he
|
|||||||
|
|
||||||
- One shared pool: `lagoon.Open`, `lagoon.Use` and `lagoon.OpenFromApp` return a `*sql.DB` and a `*gorm.DB` built on that same pool; `lagoon.Publish` makes both available on the `backpack.App`.
|
- One shared pool: `lagoon.Open`, `lagoon.Use` and `lagoon.OpenFromApp` return a `*sql.DB` and a `*gorm.DB` built on that same pool; `lagoon.Publish` makes both available on the `backpack.App`.
|
||||||
- Database-ready hooks: `lagoon.OnDatabase` runs a callback with the pool and GORM handle as soon as the database is published, immediately when it already is, otherwise when `lagoon.Publish` runs. Plugins register GORM callbacks through it from Boot, which runs before the `serve` command publishes the database.
|
- Database-ready hooks: `lagoon.OnDatabase` runs a callback with the pool and GORM handle as soon as the database is published, immediately when it already is, otherwise when `lagoon.Publish` runs. Plugins register GORM callbacks through it from Boot, which runs before the `serve` command publishes the database.
|
||||||
- After-commit work: `lagoon.Transaction` runs a function in a transaction and then the callbacks registered with `lagoon.AfterCommit`, in order, only after the commit succeeds; a nested `lagoon.Transaction` is a savepoint whose callbacks are dropped with it when it fails. A single-statement write for which GORM opens its own transaction runs its `lagoon.AfterCommit` callbacks from the `lagoon:after_commit` GORM callback (`lagoon.AfterCommitCallback`) once GORM commits, and never when the write fails. Outside both, including inside a plain GORM `Transaction`, `lagoon.AfterCommit` runs the callback immediately on that transaction's connection. The handle a callback receives always has an empty statement on the connection its work belongs to, so a query through it, even one that starts with `WithContext`, never continues from the written model's statement. A panicking callback is logged and never turns a committed write into an error.
|
- After-commit work: `lagoon.Transaction` runs a function in a transaction and then the callbacks registered with `lagoon.AfterCommit`, in order, only after the commit succeeds; a nested `lagoon.Transaction` is a savepoint whose callbacks are dropped with it when it fails. A single-statement write for which GORM opens its own implicit transaction runs its callbacks from `lagoon:after_commit` once GORM commits, and never when the write fails. A callback registered inside a foreign plain GORM transaction is unsafe because Lagoon cannot observe its commit, so `lagoon.AfterCommit` warns and skips it. Outside a transaction, callbacks run immediately. The handle a supported callback receives always has an empty statement on the connection its work belongs to. A panicking callback is logged and never turns a committed write into an error.
|
||||||
- Database check at connect time: `lagoon.CheckLocale` refuses a database whose default collation is not the ICU `pl-PL` locale, so ordering matches the database default without per-query `COLLATE`.
|
- Database check at connect time: `lagoon.CheckLocale` refuses a database whose default collation is not the ICU `pl-PL` locale, so ordering matches the database default without per-query `COLLATE`.
|
||||||
- Per-plugin migrations: `lagoon.Migrate` runs the framework's `system_files` set (`attach.Migrations`), backend admin identity set (`lagoon.BackendAdminMigrations`) and job-queue set (`lagoon.QueueMigrations`: River's schema pinned at `lagoon.RiverSchemaVersion`, then the `lagoon.JobsTable` record table, under the `lagoon.QueueHistoryID` history), then every `pact.HasMigrations` set in plugin activation order, each in its own `summer_migrations_<plugin_id>` history table (`lagoon.HistoryTableName`). `lagoon.RollbackLast` and `lagoon.Status` cover rollback and history.
|
- Per-plugin migrations: `lagoon.Migrate` runs the framework's `system_files` set (`attach.Migrations`), backend admin identity set (`lagoon.BackendAdminMigrations`) and job-queue set (`lagoon.QueueMigrations`: River's schema pinned at `lagoon.RiverSchemaVersion`, then the `lagoon.JobsTable` record table, under the `lagoon.QueueHistoryID` history), then every `pact.HasMigrations` set in plugin activation order, each in its own `summer_migrations_<plugin_id>` history table (`lagoon.HistoryTableName`). `lagoon.RollbackLast` and `lagoon.Status` cover rollback and history.
|
||||||
- Mass assignment: `lagoon.Fill` copies only allow-listed keys onto a model by GORM column name and silently drops the rest, logging each dropped key once outside production. A `json.Number` (from a decoder using `UseNumber`) fills integer, unsigned and float fields. A value that does not fit its column (a fraction, an exponent or an overflow for an integer field, or a value of the wrong type) is a `lagoon.FillTypeError` naming the key, so a caller can answer it as a validation failure on that field. `lagoon.HasFillable` and `lagoon.HasHidden` are the Go forms of `$fillable` and `$hidden`.
|
- Mass assignment: `lagoon.Fill` copies only allow-listed keys onto a model by GORM column name and silently drops the rest, logging each dropped key once outside production. A `json.Number` (from a decoder using `UseNumber`) fills integer, unsigned and float fields. A value that does not fit its column (a fraction, an exponent or an overflow for an integer field, or a value of the wrong type) is a `lagoon.FillTypeError` naming the key, so a caller can answer it as a validation failure on that field. `lagoon.HasFillable` and `lagoon.HasHidden` are the Go forms of `$fillable` and `$hidden`.
|
||||||
|
|||||||
@@ -52,6 +52,9 @@ func Transaction(ctx context.Context, gdb *gorm.DB, fn func(ctx context.Context,
|
|||||||
ctx = context.Background()
|
ctx = context.Background()
|
||||||
}
|
}
|
||||||
if parent, ok := ctx.Value(afterCommitKey{}).(*afterCommitBuffer); ok && parent != nil {
|
if parent, ok := ctx.Value(afterCommitKey{}).(*afterCommitBuffer); ok && parent != nil {
|
||||||
|
if !transactionalHandle(gdb) {
|
||||||
|
return fmt.Errorf("lagoon: nested transaction requires the parent transaction handle")
|
||||||
|
}
|
||||||
child := &afterCommitBuffer{}
|
child := &afterCommitBuffer{}
|
||||||
childCtx := context.WithValue(ctx, afterCommitKey{}, child)
|
childCtx := context.WithValue(ctx, afterCommitKey{}, child)
|
||||||
err := gdb.WithContext(childCtx).Transaction(func(tx *gorm.DB) error {
|
err := gdb.WithContext(childCtx).Transaction(func(tx *gorm.DB) error {
|
||||||
@@ -79,8 +82,9 @@ func Transaction(ctx context.Context, gdb *gorm.DB, fn func(ctx context.Context,
|
|||||||
// single-statement write for which GORM opened its own transaction (for
|
// single-statement write for which GORM opened its own transaction (for
|
||||||
// example from a GORM create callback) it runs after that commit through the
|
// example from a GORM create callback) it runs after that commit through the
|
||||||
// AfterCommitCallback callback, and not at all when the write fails.
|
// AfterCommitCallback callback, and not at all when the write fails.
|
||||||
// Anywhere else, including inside a plain gorm Transaction, fn runs
|
// Inside a foreign plain GORM transaction it logs a warning and refuses to
|
||||||
// immediately on db's connection.
|
// run, because Lagoon cannot know whether that transaction will commit.
|
||||||
|
// Anywhere else, fn runs immediately on db's connection.
|
||||||
//
|
//
|
||||||
// The handle fn receives always has an empty statement on the connection
|
// The handle fn receives always has an empty statement on the connection
|
||||||
// the work belongs to (the pool after a commit, the open transaction
|
// the work belongs to (the pool after a commit, the open transaction
|
||||||
@@ -113,10 +117,22 @@ func AfterCommit(ctx context.Context, db *gorm.DB, fn func(ctx context.Context,
|
|||||||
db.InstanceSet(statementBufferKey, buf)
|
db.InstanceSet(statementBufferKey, buf)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
if transactionalHandle(db) {
|
||||||
|
slog.Default().Warn("lagoon: after-commit callback skipped inside unmanaged transaction")
|
||||||
|
return
|
||||||
|
}
|
||||||
}
|
}
|
||||||
runAfterCommit(ctx, cleanHandle(db, ctx), []func(context.Context, *gorm.DB){fn})
|
runAfterCommit(ctx, cleanHandle(db, ctx), []func(context.Context, *gorm.DB){fn})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func transactionalHandle(db *gorm.DB) bool {
|
||||||
|
if db == nil || db.Statement == nil || db.Statement.ConnPool == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
_, ok := db.Statement.ConnPool.(gorm.TxCommitter)
|
||||||
|
return ok
|
||||||
|
}
|
||||||
|
|
||||||
func bufferFrom(ctx context.Context, db *gorm.DB) *afterCommitBuffer {
|
func bufferFrom(ctx context.Context, db *gorm.DB) *afterCommitBuffer {
|
||||||
if buf, ok := ctx.Value(afterCommitKey{}).(*afterCommitBuffer); ok && buf != nil {
|
if buf, ok := ctx.Value(afterCommitKey{}).(*afterCommitBuffer); ok && buf != nil {
|
||||||
return buf
|
return buf
|
||||||
|
|||||||
@@ -1,10 +1,13 @@
|
|||||||
package lagoon
|
package lagoon
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"errors"
|
"errors"
|
||||||
|
"log/slog"
|
||||||
"reflect"
|
"reflect"
|
||||||
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
@@ -191,21 +194,25 @@ func TestTransactionAfterCommit(t *testing.T) {
|
|||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
t.Run("plain_gorm_transaction_runs_now", func(t *testing.T) {
|
t.Run("plain_gorm_transaction_is_refused", func(t *testing.T) {
|
||||||
var got *gorm.DB
|
var logs bytes.Buffer
|
||||||
|
previous := slog.Default()
|
||||||
|
slog.SetDefault(slog.New(slog.NewTextHandler(&logs, &slog.HandlerOptions{Level: slog.LevelWarn})))
|
||||||
|
defer slog.SetDefault(previous)
|
||||||
|
ran := false
|
||||||
err := gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
err := gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||||
AfterCommit(ctx, tx, func(_ context.Context, d *gorm.DB) { got = d })
|
AfterCommit(ctx, tx, func(context.Context, *gorm.DB) { ran = true })
|
||||||
if got == nil {
|
|
||||||
t.Fatal("callback did not run immediately")
|
|
||||||
}
|
|
||||||
if got.Statement.ConnPool != tx.Statement.ConnPool {
|
|
||||||
t.Error("callback did not run on the transaction's connection")
|
|
||||||
}
|
|
||||||
return nil
|
return nil
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
|
if ran {
|
||||||
|
t.Fatal("callback ran inside an unmanaged transaction")
|
||||||
|
}
|
||||||
|
if !strings.Contains(logs.String(), "after-commit callback skipped inside unmanaged transaction") {
|
||||||
|
t.Fatalf("warning = %q", logs.String())
|
||||||
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
// The handle a callback receives has an empty statement on the write's
|
// The handle a callback receives has an empty statement on the write's
|
||||||
@@ -260,11 +267,14 @@ func TestTransactionAfterCommit(t *testing.T) {
|
|||||||
}
|
}
|
||||||
mu.Lock()
|
mu.Lock()
|
||||||
defer mu.Unlock()
|
defer mu.Unlock()
|
||||||
for _, name := range []string{"clean-implicit", "clean-plain-tx", "clean-lagoon-tx"} {
|
for _, name := range []string{"clean-implicit", "clean-lagoon-tx"} {
|
||||||
if got := labels[name]; got != "other" {
|
if got := labels[name]; got != "other" {
|
||||||
t.Errorf("%s: query through the callback handle read %q, want the lagoon_ac_others row \"other\"", name, got)
|
t.Errorf("%s: query through the callback handle read %q, want the lagoon_ac_others row \"other\"", name, got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if _, ok := labels["clean-plain-tx"]; ok {
|
||||||
|
t.Fatal("unmanaged transaction callback ran")
|
||||||
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
t.Run("outside_transaction_runs_now", func(t *testing.T) {
|
t.Run("outside_transaction_runs_now", func(t *testing.T) {
|
||||||
@@ -311,6 +321,21 @@ func TestTransactionEdges(t *testing.T) {
|
|||||||
if got == nil {
|
if got == nil {
|
||||||
t.Fatal("callback of a Transaction with a nil ctx did not run with a context")
|
t.Fatal("callback of a Transaction with a nil ctx did not run with a context")
|
||||||
}
|
}
|
||||||
|
var nestedRan bool
|
||||||
|
var nestedEntered bool
|
||||||
|
err = Transaction(t.Context(), gdb, func(ctx context.Context, _ *gorm.DB) error {
|
||||||
|
return Transaction(ctx, gdb, func(ctx context.Context, tx *gorm.DB) error {
|
||||||
|
nestedEntered = true
|
||||||
|
AfterCommit(ctx, tx, func(context.Context, *gorm.DB) { nestedRan = true })
|
||||||
|
return nil
|
||||||
|
})
|
||||||
|
})
|
||||||
|
if err == nil || !strings.Contains(err.Error(), "parent transaction handle") {
|
||||||
|
t.Fatalf("nested Transaction with root handle err = %v", err)
|
||||||
|
}
|
||||||
|
if nestedEntered || nestedRan {
|
||||||
|
t.Fatal("nested root handle entered work or ran callbacks")
|
||||||
|
}
|
||||||
if err := registerAfterCommit(gdb); err != nil {
|
if err := registerAfterCommit(gdb); err != nil {
|
||||||
t.Fatalf("registering the after-commit callback twice: %v", err)
|
t.Fatalf("registering the after-commit callback twice: %v", err)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user