- deferred_bindings migration set under summercms.deferred with backend_user_id - lagoon.DeferredBind/Unbind/Bindings/Forget/Slaves scoped by DeferredKey - lagoon.PurgeDeferred with SKIP LOCKED batches and after-commit blob deletes - attach.Store with the ported image guard, extension and MIME limits - attach.Relation, attach.HasRelations, attach.BlobKeys, File.ThumbKey - lagoon README and attachments docs
254 lines
8.1 KiB
Go
254 lines
8.1 KiB
Go
package lagoon
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"reflect"
|
|
"strconv"
|
|
"time"
|
|
|
|
"git.golem15.com/golem15/summercms/modules/lagoon/attach"
|
|
"gocloud.dev/blob"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
"gorm.io/gorm/schema"
|
|
)
|
|
|
|
// purgeBatch is how many bindings one PurgeDeferred transaction processes.
|
|
const purgeBatch = 500
|
|
|
|
// PurgeOptions configures PurgeDeferred. Before is the cut-off: bindings
|
|
// created before it are expired. Models resolves the slave_type of a
|
|
// binding for a child created under deferral to a model value of its Go
|
|
// type (keyed by MorphType); a nil Models, or a type it does not know,
|
|
// leaves such bindings in place.
|
|
type PurgeOptions struct {
|
|
Before time.Time
|
|
Models func(slaveType string) (any, bool)
|
|
}
|
|
|
|
// PurgeResult counts what PurgeDeferred did: Bindings deleted, unattached
|
|
// system_files rows (Files) and created child rows (Children) deleted, and
|
|
// bindings Skipped because their slave type could not be resolved.
|
|
type PurgeResult struct {
|
|
Bindings int
|
|
Files int
|
|
Children int
|
|
Skipped int
|
|
}
|
|
|
|
// PurgeDeferred removes expired deferred bindings, WinterCMS's
|
|
// DeferredBinding::cleanUp. It processes bindings created before
|
|
// opts.Before in id order, in batches of 500, each in its own
|
|
// lagoon.Transaction that locks the batch FOR UPDATE SKIP LOCKED, so a save
|
|
// that is committing the same session's bindings is never disturbed.
|
|
//
|
|
// For each binding: a bind of a system_files row deletes that row when it is
|
|
// still unattached (empty or NULL attachment_id) and queues its blob and
|
|
// thumbnails; a bind whose DeferredEnvelope has Created set deletes the
|
|
// child row through GORM (model hooks and soft delete apply) after resolving
|
|
// its type through opts.Models, and is skipped (left in place, counted in
|
|
// Skipped and logged once per type) when the type cannot be resolved; every
|
|
// other binding (an unbind, or a bind that only linked an existing record)
|
|
// is just deleted, and its slave is kept. Processed bindings are deleted in
|
|
// the same transaction, and the queued blob keys are deleted through
|
|
// attach.DeleteKeys only after the transaction commits (lagoon.AfterCommit).
|
|
//
|
|
// A nil bucket is an error, before any delete, when an expired binding
|
|
// points at a file.
|
|
func PurgeDeferred(ctx context.Context, db *gorm.DB, bucket *blob.Bucket, opts PurgeOptions) (PurgeResult, error) {
|
|
var res PurgeResult
|
|
if ctx == nil {
|
|
ctx = context.Background()
|
|
}
|
|
if db == nil {
|
|
return res, fmt.Errorf("lagoon: gorm db is nil")
|
|
}
|
|
if opts.Before.IsZero() {
|
|
return res, fmt.Errorf("lagoon: purge cut-off is zero")
|
|
}
|
|
if bucket == nil {
|
|
var files int64
|
|
err := db.Session(&gorm.Session{NewDB: true, Context: ctx}).
|
|
Model(&DeferredBinding{}).
|
|
Where("created_at < ? AND is_bind AND slave_type = ?", opts.Before, DeferredFileType).
|
|
Count(&files).Error
|
|
if err != nil {
|
|
return res, fmt.Errorf("lagoon: purge deferred bindings: %w", err)
|
|
}
|
|
if files > 0 {
|
|
return res, fmt.Errorf("lagoon: purge deferred bindings: %d expired file bindings need the uploads bucket, which is nil", files)
|
|
}
|
|
}
|
|
warned := map[string]bool{}
|
|
var cursor uint
|
|
for {
|
|
var batch PurgeResult
|
|
var last uint
|
|
var n int
|
|
err := Transaction(ctx, db, func(ctx context.Context, tx *gorm.DB) error {
|
|
batch = PurgeResult{}
|
|
var rows []DeferredBinding
|
|
err := tx.Session(&gorm.Session{NewDB: true}).
|
|
Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}).
|
|
Where("created_at < ? AND id > ?", opts.Before, cursor).
|
|
Order("id").
|
|
Limit(purgeBatch).
|
|
Find(&rows).Error
|
|
if err != nil {
|
|
return fmt.Errorf("lagoon: purge deferred bindings: %w", err)
|
|
}
|
|
n = len(rows)
|
|
if n == 0 {
|
|
return nil
|
|
}
|
|
last = rows[n-1].ID
|
|
var done []uint
|
|
var keys []string
|
|
for _, row := range rows {
|
|
processed, err := purgeBinding(ctx, tx, bucket, opts, row, &batch, &keys, warned)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if processed {
|
|
done = append(done, row.ID)
|
|
}
|
|
}
|
|
if len(done) > 0 {
|
|
if err := tx.Session(&gorm.Session{NewDB: true}).Where("id IN ?", done).Delete(&DeferredBinding{}).Error; err != nil {
|
|
return fmt.Errorf("lagoon: purge deferred bindings: %w", err)
|
|
}
|
|
}
|
|
batch.Bindings = len(done)
|
|
if len(keys) > 0 {
|
|
AfterCommit(ctx, tx, func(ctx context.Context, _ *gorm.DB) {
|
|
if err := attach.DeleteKeys(context.WithoutCancel(ctx), bucket, keys); err != nil {
|
|
slog.Default().WarnContext(ctx, "lagoon: purge could not delete file blobs", slog.String("error", err.Error()))
|
|
}
|
|
})
|
|
}
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return res, err
|
|
}
|
|
if n == 0 {
|
|
return res, nil
|
|
}
|
|
res.Bindings += batch.Bindings
|
|
res.Files += batch.Files
|
|
res.Children += batch.Children
|
|
res.Skipped += batch.Skipped
|
|
cursor = last
|
|
}
|
|
}
|
|
|
|
// purgeBinding handles one expired binding inside the batch transaction and
|
|
// reports whether it may be deleted.
|
|
func purgeBinding(ctx context.Context, tx *gorm.DB, bucket *blob.Bucket, opts PurgeOptions, row DeferredBinding, res *PurgeResult, keys *[]string, warned map[string]bool) (bool, error) {
|
|
if !row.IsBind {
|
|
return true, nil
|
|
}
|
|
if row.SlaveType == DeferredFileType {
|
|
id, err := strconv.ParseUint(row.SlaveID, 10, 64)
|
|
if err != nil || id == 0 {
|
|
return true, nil
|
|
}
|
|
if bucket == nil {
|
|
return false, fmt.Errorf("lagoon: purge deferred binding %d: the uploads bucket is nil", row.ID)
|
|
}
|
|
var f attach.File
|
|
err = tx.Session(&gorm.Session{NewDB: true}).
|
|
Where("id = ? AND (attachment_id IS NULL OR attachment_id = '')", id).
|
|
Take(&f).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return true, nil
|
|
}
|
|
if err != nil {
|
|
return false, fmt.Errorf("lagoon: purge file %d: %w", id, err)
|
|
}
|
|
if err := tx.Session(&gorm.Session{NewDB: true}).Where("id = ?", f.ID).Delete(&attach.File{}).Error; err != nil {
|
|
return false, fmt.Errorf("lagoon: purge file %d: %w", id, err)
|
|
}
|
|
*keys = append(*keys, attach.BlobKeys(f)...)
|
|
res.Files++
|
|
return true, nil
|
|
}
|
|
env, err := row.Envelope()
|
|
if err != nil {
|
|
slog.Default().WarnContext(ctx, "lagoon: purge treats an unreadable deferred binding as a plain link",
|
|
slog.Uint64("binding_id", uint64(row.ID)), slog.String("error", err.Error()))
|
|
return true, nil
|
|
}
|
|
if !env.Created {
|
|
return true, nil
|
|
}
|
|
var model any
|
|
ok := false
|
|
if opts.Models != nil {
|
|
model, ok = opts.Models(row.SlaveType)
|
|
}
|
|
if !ok || model == nil {
|
|
if !warned[row.SlaveType] {
|
|
warned[row.SlaveType] = true
|
|
slog.Default().WarnContext(ctx, "lagoon: purge cannot resolve the model of a deferred child; its bindings are kept",
|
|
slog.String("slave_type", row.SlaveType))
|
|
}
|
|
res.Skipped++
|
|
return false, nil
|
|
}
|
|
deleted, err := deleteCreatedChild(tx, model, row.SlaveID)
|
|
if err != nil {
|
|
return false, fmt.Errorf("lagoon: purge %s %s: %w", row.SlaveType, row.SlaveID, err)
|
|
}
|
|
if deleted {
|
|
res.Children++
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
// deleteCreatedChild loads the row with primary key id into a fresh value of
|
|
// model's type and deletes it through GORM. A missing row is not an error.
|
|
func deleteCreatedChild(tx *gorm.DB, model any, id string) (bool, error) {
|
|
typ := reflect.TypeOf(model)
|
|
for typ.Kind() == reflect.Pointer {
|
|
typ = typ.Elem()
|
|
}
|
|
if typ.Kind() != reflect.Struct {
|
|
return false, fmt.Errorf("model %T is not a struct", model)
|
|
}
|
|
fresh := reflect.New(typ).Interface()
|
|
stmt := &gorm.Statement{DB: tx}
|
|
if err := stmt.Parse(fresh); err != nil {
|
|
return false, err
|
|
}
|
|
pk := stmt.Schema.PrioritizedPrimaryField
|
|
if pk == nil {
|
|
return false, fmt.Errorf("model %T has no primary key", model)
|
|
}
|
|
var key any = id
|
|
switch pk.DataType {
|
|
case schema.Int, schema.Uint:
|
|
n, err := strconv.ParseInt(id, 10, 64)
|
|
if err != nil {
|
|
// slave_id is not a key of this integer primary key: nothing to delete.
|
|
return false, nil
|
|
}
|
|
key = n
|
|
}
|
|
q := tx.Session(&gorm.Session{NewDB: true})
|
|
err := q.Where(clause.Eq{Column: clause.Column{Table: clause.CurrentTable, Name: pk.DBName}, Value: key}).Take(fresh).Error
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return false, nil
|
|
}
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
if err := tx.Session(&gorm.Session{NewDB: true}).Delete(fresh).Error; err != nil {
|
|
return false, err
|
|
}
|
|
return true, nil
|
|
}
|