feat(12.2-01): add deferred bindings, guarded upload store and purge
- 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
This commit is contained in:
@@ -15,7 +15,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`.
|
||||
- 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 nested `lagoon.Transaction` must be given the outer transaction's handle: given a root handle it returns an error without running its function, rather than open an independent transaction whose callbacks would wait on the outer one. 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.
|
||||
- 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`), `deferred_bindings` set (`lagoon.DeferredBindingMigrations`, under the `lagoon.DeferredHistoryID` history) 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`.
|
||||
- Validation: `lagoon.Validate` accepts Laravel-style rule strings (`required`, `nullable`, `integer`, `numeric`, `between`, `min`, `max`, `in`, `unique`, `boolean`, `email`, `confirmed`, `different`, `mimes`) and returns a field-to-messages map, translated through phrasebook when a translator is given. Unknown rule tokens are an error. A failed numeric range reports the bound that failed: the `min` message below the lower bound, the `max` message above the upper one, and the numeric `between` message when the bound came from `between`.
|
||||
- Request validation: `lagoon.ValidateRequest` reproduces Laravel 9 request validation for ported API endpoints, so a 422 body matches the PHP one message for message. It takes the decoded input and an ordered `lagoon.RequestRule` table (attribute names may hold `*` wildcards, expanded against the input to `posts.0.title`), runs the rules of each attribute in order and stops an attribute after a failed implicit rule (`required`, `present`, `filled`, `accepted`) or, under `bail`, after any failure. A non-implicit rule is skipped for an absent attribute, a blank string, a null value under `nullable` and an absent key under `sometimes`. Supported rules: `required`, `present`, `filled`, `accepted`, `nullable`, `sometimes`, `bail`, `array`, `string`, `integer`, `numeric`, `boolean`, `email` (PHP `FILTER_VALIDATE_EMAIL`, WinterCMS's default), `url`, `date`, `after`, `after_or_equal`, `before`, `before_or_equal` (a date, a relative word such as `tomorrow`, or another field), `exists:table,column`, `regex`, `not_regex`, `in`, `not_in`, `file`, `image`, `mimes`, `min`, `max`, `size` and `between`, plus closure rules built with `lagoon.CustomRule`. The size rules compare the number under `numeric` or `integer` (exactly, as decimals), the element count of an array, kilobytes of a `lagoon.UploadedFile`, and otherwise the length in characters, and pick the matching message. Messages come from the `lagoon::validation` catalog in the request locale; `lagoon.ErrorKeys` gives the attribute order of PHP's message bag.
|
||||
@@ -23,8 +23,9 @@ Postgres data layer: the shared GORM connection, per-plugin migrations, model he
|
||||
- Pagination: `lagoon.Paginate` builds a `lagoon.Page` with `data` and `meta` (`current_page`, `last_page`, `per_page`, `total`).
|
||||
- Column types: `lagoon.Encrypted` stores AES-256-GCM ciphertext under a key derived from `app.key`, decrypts with previous keys during rotation, and always redacts itself in JSON and string output; `lagoon.Jsonable` stores JSON as TEXT and keeps SQL NULL distinct from an empty value.
|
||||
- Lifecycle and relations: hook interfaces matching GORM's native method names (`lagoon.HasBeforeCreate`, `lagoon.HasBeforeSave`, `lagoon.HasBeforeDelete`, `lagoon.HasAfterDelete`) plus `lagoon.HasBeforeValidate`; `lagoon.WithSoftDeleteCascade` runs a cascade inside the parent delete; `lagoon.RegisterJoinTable` wires pivot models with business columns.
|
||||
- Deferred binding: WinterCMS's `deferred_bindings` table holds the uploads and related-record changes of a form whose record is not saved yet. Every operation takes a `lagoon.DeferredKey` (the form's session key, the owning backend admin's id and the master record's morph type from `lagoon.MorphType`) and never reads or changes another admin's rows, since each row stores `backend_user_id`. `lagoon.DeferredBind` and `lagoon.DeferredUnbind` port WinterCMS's duplicate and cancel rules: a repeated bind writes nothing, and an unbind of a slave with a pending bind deletes that bind and returns it so the caller can remove what it created. `lagoon.DeferredBindings` reads and locks a session's bindings for the save that commits them, `lagoon.DeferredForget` deletes them once applied, and `lagoon.DeferredSlaves` is the subquery a list uses to include pending rows. A child created under deferral carries the `lagoon.DeferredEnvelope` (`{"created":true,"pivot":{...}}`) in `pivot_data`. `lagoon.PurgeDeferred` removes expired bindings: it deletes an unattached `system_files` row a bind points at, and its blobs only after the commit, deletes a child only when its binding carries the created envelope, keeps records that were only linked, and locks each batch with `FOR UPDATE SKIP LOCKED`.
|
||||
- Imports from Laravel: `lagoon.DecryptLaravelPayload` decrypts Laravel `encrypted` payloads with the old application key, for one-off data imports.
|
||||
- Attachments (`attach`): the `attach.File` model for `system_files` rows, WinterCMS-compatible partitioned storage keys (`attach.BlobKey`, `attach.PartitionDirectory`), public URLs (`attach.PublicURL` for any key, `attach.File.URL` for an original, matching WinterCMS's `File::getPath()` under the WinterCMS layout), on-demand thumbnails through `attach.File.Thumb` for JPEG, PNG, GIF and WebP originals (a WebP original's thumbnail is JPEG bytes under its `.webp` name, since WebP cannot be encoded; a missing, undecodable or oversized original gets WinterCMS's broken-image picture, `attach.BrokenImagePNG`, as its thumbnail, as `File::makeThumb` does), static serving with an optional `is_public` gate (`attach.StaticHandlerPublic`), and a two-phase delete that removes blobs only after the database transaction commits (`attach.DeleteForOwner`, `attach.DeleteKeys`).
|
||||
- Attachments (`attach`): the `attach.File` model for `system_files` rows, WinterCMS-compatible partitioned storage keys (`attach.BlobKey`, `attach.PartitionDirectory`), public URLs (`attach.PublicURL` for any key, `attach.File.URL` for an original, matching WinterCMS's `File::getPath()` under the WinterCMS layout), on-demand thumbnails through `attach.File.Thumb` for JPEG, PNG, GIF and WebP originals (a WebP original's thumbnail is JPEG bytes under its `.webp` name, since WebP cannot be encoded; a missing, undecodable or oversized original gets WinterCMS's broken-image picture, `attach.BrokenImagePNG`, as its thumbnail, as `File::makeThumb` does), storing uploads through `attach.Store` (a server-generated disk name, an extension allow-list with `attach.DefaultImageExtensions` and `attach.DefaultFileExtensions` as defaults, a MIME filter, a size limit enforced while streaming and, in image mode, the `attach.IsAllowedImage` content guard), attachment relation declarations (`attach.Relation`, `attach.HasRelations`), static serving with an optional `is_public` gate (`attach.StaticHandlerPublic`), and a two-phase delete that removes blobs only after the database transaction commits (`attach.DeleteForOwner`, `attach.DeleteKeys`).
|
||||
|
||||
## Usage
|
||||
|
||||
@@ -127,6 +128,21 @@ func (p *Plugin) Migrations() []*gormigrate.Migration {
|
||||
| `lagoon.QueueHistoryID` | History id of the job-queue set, `summercms.conga`. |
|
||||
| `lagoon.JobsTable` | Name of the job record table, `summer_jobs`. |
|
||||
| `lagoon.RiverSchemaVersion` | The pinned River schema version, 7. |
|
||||
| `lagoon.DeferredBindingMigrations` | Creates WinterCMS's `deferred_bindings` table plus the `backend_user_id` owner column. |
|
||||
| `lagoon.DeferredHistoryID` | History id of the deferred-binding set, `summercms.deferred`. |
|
||||
| `lagoon.DeferredBinding` | The `deferred_bindings` row model; `lagoon.DeferredBinding.Envelope` decodes its `pivot_data`. |
|
||||
| `lagoon.DeferredKey` | Session key, admin id and master type that scope every deferred-binding operation. |
|
||||
| `lagoon.DeferredEnvelope` | The framework's `pivot_data` shape: `Created` marks a child created under deferral, `Pivot` holds pivot values. |
|
||||
| `lagoon.DeferredFileType` | The `slave_type` of a binding that points at a `system_files` row. |
|
||||
| `lagoon.MorphType` | The `master_type` or `slave_type` string of a model: its `attach.Owner` morph name, else its table name. |
|
||||
| `lagoon.DeferredBind` | Records a pending bind; a repeat writes nothing and a pending unbind of the same slave is cancelled. |
|
||||
| `lagoon.DeferredUnbind` | Records a pending unbind, or cancels a pending bind of the same slave and returns it. |
|
||||
| `lagoon.DeferredBindings` | Reads and locks a session's bindings for the given relation fields, in id order. |
|
||||
| `lagoon.DeferredForget` | Deletes applied bindings by id. |
|
||||
| `lagoon.DeferredSlaves` | Subquery of a session's bound or unbound slave ids, for list queries. |
|
||||
| `lagoon.PurgeDeferred` | Removes expired bindings, their unattached files (blobs after commit) and the children created under deferral. |
|
||||
| `lagoon.PurgeOptions` | Cut-off time and created-child model resolver for `lagoon.PurgeDeferred`. |
|
||||
| `lagoon.PurgeResult` | Counts of deleted bindings, files and children and of skipped bindings. |
|
||||
| `lagoon.RuntimeCommands` | Returns the migrate, migrate:rollback, migrate:status and key:generate commands. |
|
||||
| `lagoon.KeyGenerateCommand` | Returns the key:generate command on its own. |
|
||||
| `lagoon.LoadAppKey` | Decodes `app.key` and `app.previous_keys`. |
|
||||
@@ -162,6 +178,22 @@ func (p *Plugin) Migrations() []*gormigrate.Migration {
|
||||
| `attach.DeleteForOwner` | Deletes an owner's attachment rows in a transaction and reports their blob keys. |
|
||||
| `attach.DeleteKeys` | Deletes blobs, including thumbnails, after the transaction commits. |
|
||||
| `attach.Migrations` | Creates the `system_files` table. |
|
||||
| `attach.Store` | Stores an upload as an unattached `system_files` row with `sort_order` equal to its id, after the type, size and image checks. |
|
||||
| `attach.Upload` | The client file name, body and public flag of one upload. |
|
||||
| `attach.Limits` | Size limit, allowed extensions, allowed MIME types and image mode for `attach.Store`. |
|
||||
| `attach.ErrTooLarge` | Returned by `attach.Store` for a body over `MaxBytes`. |
|
||||
| `attach.ErrFileType` | Returned by `attach.Store` for a missing, malformed or disallowed extension. |
|
||||
| `attach.ErrMIMEType` | Returned by `attach.Store` when the content type matches no `MIMETypes` entry. |
|
||||
| `attach.ErrNotImage` | Returned by `attach.Store` in image mode for content that is not an allowed image. |
|
||||
| `attach.DefaultImageExtensions` | Image-mode extensions when none are given: jpg, jpeg, png, gif, webp. |
|
||||
| `attach.DefaultFileExtensions` | File-mode extensions when none are given: WinterCMS's default list without the script-capable types. |
|
||||
| `attach.AllowedImageMIMEs` | The sniffed content types the image guard accepts. |
|
||||
| `attach.IsAllowedImage` | The image guard: sniffed type, decoded header and a pixel ceiling, failing closed. |
|
||||
| `attach.MaxImagePixels` | The image guard's pixel ceiling, 4096 by 4096. |
|
||||
| `attach.Relation` | One attachOne or attachMany relation: `Name`, `Many` and `Public`. |
|
||||
| `attach.HasRelations` | Implemented by an owner model that declares its attachment relations. |
|
||||
| `attach.BlobKeys` | The original's key and the thumbnail prefix of a file, for `attach.DeleteKeys`. |
|
||||
| `attach.File.ThumbKey` | Blob key of a lazily generated thumbnail, for serving a protected file's thumbnail without a public URL. |
|
||||
|
||||
## Configuration
|
||||
|
||||
@@ -191,7 +223,7 @@ storage:
|
||||
|
||||
| Command | Flags | Description |
|
||||
|---------|-------|-------------|
|
||||
| `migrate` | none | Runs the framework migrations, then each plugin's migrations in dependency order. |
|
||||
| `migrate` | none | Runs the framework migrations (`system_files`, the backend admin tables, `deferred_bindings` and the job queue), then each plugin's migrations in dependency order. |
|
||||
| `migrate:rollback` | `--plugin <id>` | Rolls back the last migration of the given plugin; without the flag, of the last activated plugin that has migrations. |
|
||||
| `migrate:status` | none | Prints a table of plugin, history table and applied migration IDs. |
|
||||
| `key:generate` | none | Prints a fresh base64 32-byte key for `app.key`; writes nothing. |
|
||||
|
||||
@@ -107,7 +107,10 @@ func init() {
|
||||
Register(&File{})
|
||||
}
|
||||
|
||||
func blobKeysFor(f File) []string {
|
||||
// BlobKeys returns the blob keys of f: the original's partitioned key and
|
||||
// the thumb_<id>_ prefix of its thumbnails. DeleteKeys treats the second as
|
||||
// a prefix, so passing both removes the original and every thumbnail.
|
||||
func BlobKeys(f File) []string {
|
||||
part := PartitionDirectory(f.DiskName)
|
||||
return []string{
|
||||
part + f.DiskName,
|
||||
@@ -141,7 +144,7 @@ func DeleteForOwner(tx *gorm.DB, owner Owner, ownerID string, afterCommit func(b
|
||||
}
|
||||
var keys []string
|
||||
for _, f := range files {
|
||||
keys = append(keys, blobKeysFor(f)...)
|
||||
keys = append(keys, BlobKeys(f)...)
|
||||
}
|
||||
if err := tx.Where("attachment_type = ? AND attachment_id = ?", morph, ownerID).Delete(&File{}).Error; err != nil {
|
||||
return err
|
||||
|
||||
@@ -86,7 +86,7 @@ func TestDeleteKeysThumbPrefixIsIDDelimited(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if err := DeleteKeys(ctx, bucket, blobKeysFor(File{ID: 4, DiskName: disk})); err != nil {
|
||||
if err := DeleteKeys(ctx, bucket, BlobKeys(File{ID: 4, DiskName: disk})); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, key := range doomed {
|
||||
|
||||
47
modules/lagoon/attach/guard.go
Normal file
47
modules/lagoon/attach/guard.go
Normal file
@@ -0,0 +1,47 @@
|
||||
package attach
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"image"
|
||||
"net/http"
|
||||
"slices"
|
||||
)
|
||||
|
||||
// AllowedImageMIMEs are the content types IsAllowedImage accepts, as
|
||||
// sniffed from the bytes: JPEG, PNG, GIF and WebP, the formats the
|
||||
// thumbnailer decodes.
|
||||
var AllowedImageMIMEs = []string{"image/jpeg", "image/png", "image/gif", "image/webp"}
|
||||
|
||||
// MaxImagePixels is the largest image (width times height) IsAllowedImage
|
||||
// accepts, the same ceiling File.Thumb applies before decoding, so every
|
||||
// accepted image can be thumbnailed.
|
||||
const MaxImagePixels = maxThumbSourcePixels
|
||||
|
||||
// IsAllowedImage reports whether data is a JPEG, PNG, GIF or WebP image:
|
||||
// the bytes must sniff as one of AllowedImageMIMEs (http.DetectContentType,
|
||||
// independent of any file name or client header), the header must decode
|
||||
// through image.DecodeConfig as that format with a positive width and
|
||||
// height, and the image must not exceed MaxImagePixels. Decoding the header
|
||||
// rejects a polyglot whose first bytes alone look right. It fails closed:
|
||||
// empty or unreadable content is refused. data may be a prefix of the file
|
||||
// as long as it holds the image header.
|
||||
func IsAllowedImage(data []byte) bool {
|
||||
if len(data) == 0 {
|
||||
return false
|
||||
}
|
||||
if !slices.Contains(AllowedImageMIMEs, http.DetectContentType(data)) {
|
||||
return false
|
||||
}
|
||||
cfg, format, err := image.DecodeConfig(bytes.NewReader(data))
|
||||
if err != nil || cfg.Width <= 0 || cfg.Height <= 0 {
|
||||
return false
|
||||
}
|
||||
if int64(cfg.Width)*int64(cfg.Height) > MaxImagePixels {
|
||||
return false
|
||||
}
|
||||
switch format {
|
||||
case "jpeg", "png", "gif", "webp":
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
21
modules/lagoon/attach/relation.go
Normal file
21
modules/lagoon/attach/relation.go
Normal file
@@ -0,0 +1,21 @@
|
||||
package attach
|
||||
|
||||
// Relation declares one WinterCMS attachOne or attachMany relation of an
|
||||
// Owner model. Name is the relation name stored in system_files.field, Many
|
||||
// is true for attachMany (false for attachOne), and Public decides the
|
||||
// is_public flag of files stored through the relation: a protected relation
|
||||
// (Public false) stores is_public=false rows, for which the framework never
|
||||
// builds a public URL.
|
||||
type Relation struct {
|
||||
Name string
|
||||
Many bool
|
||||
Public bool
|
||||
}
|
||||
|
||||
// HasRelations is implemented by an Owner model that declares its
|
||||
// attachment relations, the Go form of WinterCMS's $attachOne and
|
||||
// $attachMany arrays. A form field that edits attachments must name one of
|
||||
// the declared relations.
|
||||
type HasRelations interface {
|
||||
AttachRelations() []Relation
|
||||
}
|
||||
281
modules/lagoon/attach/store.go
Normal file
281
modules/lagoon/attach/store.go
Normal file
@@ -0,0 +1,281 @@
|
||||
package attach
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"mime"
|
||||
"net/http"
|
||||
"path"
|
||||
"regexp"
|
||||
"slices"
|
||||
"strings"
|
||||
|
||||
"gocloud.dev/blob"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// sniffBytes is how much of an upload Store reads ahead for the content
|
||||
// sniff and the image guard.
|
||||
const sniffBytes = 1 << 20
|
||||
|
||||
var (
|
||||
// ErrTooLarge is returned by Store when the body exceeds Limits.MaxBytes.
|
||||
ErrTooLarge = errors.New("attach: file is too large")
|
||||
// ErrFileType is returned by Store when the file name's extension is
|
||||
// missing, malformed or not in the allowed extension list.
|
||||
ErrFileType = errors.New("attach: file type is not allowed")
|
||||
// ErrMIMEType is returned by Store when the content type matches none
|
||||
// of Limits.MIMETypes.
|
||||
ErrMIMEType = errors.New("attach: file content type is not allowed")
|
||||
// ErrNotImage is returned by Store in image mode when the bytes are not
|
||||
// an image IsAllowedImage accepts.
|
||||
ErrNotImage = errors.New("attach: file is not an allowed image")
|
||||
)
|
||||
|
||||
// DefaultImageExtensions is the extension list of an image upload when
|
||||
// Limits.Extensions is empty: jpg, jpeg, png, gif and webp, the formats
|
||||
// IsAllowedImage and the thumbnailer handle. WinterCMS's image list also
|
||||
// has avif, bmp and svg; they are left out because nothing here decodes
|
||||
// them and svg can carry script.
|
||||
var DefaultImageExtensions = []string{"jpg", "jpeg", "png", "gif", "webp"}
|
||||
|
||||
// DefaultFileExtensions is the extension list of a file upload when
|
||||
// Limits.Extensions is empty: WinterCMS's default list (winter/storm
|
||||
// Filesystem\Definitions::defaultExtensions) minus the script-capable types
|
||||
// svg, js, map, css, less, scss, swf and xml. The final list is avi, avif,
|
||||
// bmp, doc, docx, eot, flv, gif, ico, ics, jpeg, jpg, mkv, mov, mp3, mp4,
|
||||
// mpeg, ods, odt, ogg, pdf, png, ppt, pptx, rar, ttf, txt, wav, webm, webp,
|
||||
// wmv, woff, woff2, xls, xlsx and zip.
|
||||
var DefaultFileExtensions = []string{
|
||||
"avi", "avif", "bmp", "doc", "docx", "eot", "flv", "gif", "ico", "ics",
|
||||
"jpeg", "jpg", "mkv", "mov", "mp3", "mp4", "mpeg", "ods", "odt", "ogg",
|
||||
"pdf", "png", "ppt", "pptx", "rar", "ttf", "txt", "wav", "webm", "webp",
|
||||
"wmv", "woff", "woff2", "xls", "xlsx", "zip",
|
||||
}
|
||||
|
||||
var extPattern = regexp.MustCompile(`^[a-z0-9]{1,10}$`)
|
||||
|
||||
// Upload is one file to store. FileName is the client's file name: only its
|
||||
// extension and base name are used (for the allowed-type check and the
|
||||
// file_name column); no part of it reaches a blob key. Body is read once,
|
||||
// to the end or to the size limit. Public sets the row's is_public flag.
|
||||
type Upload struct {
|
||||
FileName string
|
||||
Body io.Reader
|
||||
Public bool
|
||||
}
|
||||
|
||||
// Limits restricts what Store accepts.
|
||||
//
|
||||
// MaxBytes is the largest body in bytes; 0 means no limit of its own (the
|
||||
// caller's request body cap still applies). Extensions lists the allowed
|
||||
// lower-case extensions without the dot; empty means DefaultImageExtensions
|
||||
// when Image is set, else DefaultFileExtensions. MIMETypes, when not empty,
|
||||
// must match the stored content type: an entry containing a slash is a MIME
|
||||
// pattern such as "image/png" or "image/*", an entry without one is an
|
||||
// extension. Image applies the image guard (IsAllowedImage) to the content.
|
||||
type Limits struct {
|
||||
MaxBytes int64
|
||||
Extensions []string
|
||||
MIMETypes []string
|
||||
Image bool
|
||||
}
|
||||
|
||||
// Store saves an upload as an unattached system_files row.
|
||||
//
|
||||
// It accepts the client extension, lower-cased, only when it matches
|
||||
// [a-z0-9]{1,10} and is allowed by Limits (else ErrFileType). It reads up to
|
||||
// 1 MiB ahead to sniff the content type from the bytes; in image mode those
|
||||
// bytes must pass IsAllowedImage (else ErrNotImage), and Limits.MIMETypes is
|
||||
// checked against the sniffed type, or the extension's registered type when
|
||||
// the sniff only says application/octet-stream (else ErrMIMEType). The body
|
||||
// is then streamed into bucket at BlobKey of a server-generated disk name (22
|
||||
// random lowercase hex characters, a dot and the extension); a body longer
|
||||
// than Limits.MaxBytes aborts the write, deletes the key and returns
|
||||
// ErrTooLarge. Finally it inserts the row with empty attachment columns,
|
||||
// is_public from Upload.Public, the byte size and the content type, and sets
|
||||
// sort_order to the new id as WinterCMS's Sortable trait does. When the row
|
||||
// cannot be written the blob is deleted again.
|
||||
//
|
||||
// db may be a transaction. The blob is written before the row, so a caller
|
||||
// whose transaction rolls back after Store returned must delete the
|
||||
// returned file's BlobKeys itself.
|
||||
func Store(ctx context.Context, db *gorm.DB, bucket *blob.Bucket, in Upload, lim Limits) (*File, error) {
|
||||
if ctx == nil {
|
||||
ctx = context.Background()
|
||||
}
|
||||
if db == nil {
|
||||
return nil, fmt.Errorf("attach: store db is nil")
|
||||
}
|
||||
if bucket == nil {
|
||||
return nil, fmt.Errorf("attach: bucket is nil")
|
||||
}
|
||||
if in.Body == nil {
|
||||
return nil, fmt.Errorf("attach: upload body is nil")
|
||||
}
|
||||
if lim.MaxBytes < 0 {
|
||||
return nil, fmt.Errorf("attach: negative size limit %d", lim.MaxBytes)
|
||||
}
|
||||
name := clientBaseName(in.FileName)
|
||||
ext := strings.ToLower(strings.TrimPrefix(path.Ext(name), "."))
|
||||
if !extPattern.MatchString(ext) || !slices.Contains(allowedExtensions(lim), ext) {
|
||||
return nil, fmt.Errorf("%w: %q", ErrFileType, ext)
|
||||
}
|
||||
|
||||
br := bufio.NewReaderSize(in.Body, sniffBytes)
|
||||
head, err := br.Peek(sniffBytes)
|
||||
if err != nil && !errors.Is(err, io.EOF) && !errors.Is(err, bufio.ErrBufferFull) {
|
||||
return nil, fmt.Errorf("attach: read upload: %w", err)
|
||||
}
|
||||
if lim.MaxBytes > 0 && int64(len(head)) > lim.MaxBytes {
|
||||
return nil, ErrTooLarge
|
||||
}
|
||||
if lim.Image && !IsAllowedImage(head) {
|
||||
return nil, ErrNotImage
|
||||
}
|
||||
contentType := baseMediaType(http.DetectContentType(head))
|
||||
if contentType == "application/octet-stream" {
|
||||
if byExt := baseMediaType(mime.TypeByExtension("." + ext)); byExt != "" {
|
||||
contentType = byExt
|
||||
}
|
||||
}
|
||||
if len(lim.MIMETypes) > 0 && !mimeAllowed(lim.MIMETypes, contentType, ext) {
|
||||
return nil, fmt.Errorf("%w: %s", ErrMIMEType, contentType)
|
||||
}
|
||||
|
||||
diskName, err := newDiskName(ext)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
key := BlobKey(diskName)
|
||||
size, err := writeBlob(ctx, bucket, key, br, contentType, lim.MaxBytes)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
public := in.Public
|
||||
f := &File{
|
||||
DiskName: diskName,
|
||||
FileName: name,
|
||||
FileSize: size,
|
||||
ContentType: contentType,
|
||||
IsPublic: &public,
|
||||
}
|
||||
q := db.Session(&gorm.Session{NewDB: true, Context: ctx})
|
||||
err = q.Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Create(f).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
f.SortOrder = int(f.ID)
|
||||
return tx.Model(&File{}).Where("id = ?", f.ID).Update("sort_order", f.SortOrder).Error
|
||||
})
|
||||
if err != nil {
|
||||
_ = deleteKey(context.WithoutCancel(ctx), bucket, key)
|
||||
return nil, fmt.Errorf("attach: store row: %w", err)
|
||||
}
|
||||
return f, nil
|
||||
}
|
||||
|
||||
// writeBlob streams r into key and returns the byte count. With limit > 0 a
|
||||
// body of more than limit bytes aborts the write and deletes the key.
|
||||
func writeBlob(ctx context.Context, bucket *blob.Bucket, key string, r io.Reader, contentType string, limit int64) (int64, error) {
|
||||
writeCtx, cancel := context.WithCancel(ctx)
|
||||
defer cancel()
|
||||
w, err := bucket.NewWriter(writeCtx, key, &blob.WriterOptions{ContentType: contentType})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("attach: blob writer: %w", err)
|
||||
}
|
||||
src := r
|
||||
if limit > 0 {
|
||||
src = io.LimitReader(r, limit+1)
|
||||
}
|
||||
n, copyErr := io.Copy(w, src)
|
||||
if copyErr == nil && limit > 0 && n > limit {
|
||||
copyErr = ErrTooLarge
|
||||
}
|
||||
if copyErr != nil {
|
||||
// Cancelling the writer's context before Close discards the write.
|
||||
cancel()
|
||||
_ = w.Close()
|
||||
_ = deleteKey(context.WithoutCancel(ctx), bucket, key)
|
||||
if errors.Is(copyErr, ErrTooLarge) {
|
||||
return 0, ErrTooLarge
|
||||
}
|
||||
return 0, fmt.Errorf("attach: write upload: %w", copyErr)
|
||||
}
|
||||
if err := w.Close(); err != nil {
|
||||
_ = deleteKey(context.WithoutCancel(ctx), bucket, key)
|
||||
return 0, fmt.Errorf("attach: write upload: %w", err)
|
||||
}
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// clientBaseName is the last element of a client file name, with either
|
||||
// slash style treated as a separator.
|
||||
func clientBaseName(name string) string {
|
||||
name = strings.ReplaceAll(name, `\`, "/")
|
||||
if i := strings.LastIndex(name, "/"); i >= 0 {
|
||||
name = name[i+1:]
|
||||
}
|
||||
return strings.TrimSpace(name)
|
||||
}
|
||||
|
||||
func allowedExtensions(lim Limits) []string {
|
||||
if len(lim.Extensions) == 0 {
|
||||
if lim.Image {
|
||||
return DefaultImageExtensions
|
||||
}
|
||||
return DefaultFileExtensions
|
||||
}
|
||||
out := make([]string, 0, len(lim.Extensions))
|
||||
for _, e := range lim.Extensions {
|
||||
out = append(out, strings.ToLower(strings.TrimPrefix(strings.TrimSpace(e), ".")))
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func baseMediaType(ct string) string {
|
||||
if ct == "" {
|
||||
return ""
|
||||
}
|
||||
mt, _, err := mime.ParseMediaType(ct)
|
||||
if err != nil {
|
||||
return strings.ToLower(strings.TrimSpace(strings.SplitN(ct, ";", 2)[0]))
|
||||
}
|
||||
return mt
|
||||
}
|
||||
|
||||
// mimeAllowed reports whether contentType or ext matches one of patterns.
|
||||
func mimeAllowed(patterns []string, contentType, ext string) bool {
|
||||
for _, p := range patterns {
|
||||
p = strings.ToLower(strings.TrimSpace(p))
|
||||
if p == "" {
|
||||
continue
|
||||
}
|
||||
if !strings.Contains(p, "/") {
|
||||
if strings.TrimPrefix(p, ".") == ext {
|
||||
return true
|
||||
}
|
||||
continue
|
||||
}
|
||||
pType, pSub, _ := strings.Cut(p, "/")
|
||||
cType, cSub, _ := strings.Cut(contentType, "/")
|
||||
if (pType == "*" || pType == cType) && (pSub == "*" || pSub == cSub) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func newDiskName(ext string) (string, error) {
|
||||
raw := make([]byte, 11)
|
||||
if _, err := rand.Read(raw); err != nil {
|
||||
return "", fmt.Errorf("attach: disk name: %w", err)
|
||||
}
|
||||
return hex.EncodeToString(raw) + "." + ext, nil
|
||||
}
|
||||
126
modules/lagoon/attach/store_test.go
Normal file
126
modules/lagoon/attach/store_test.go
Normal file
@@ -0,0 +1,126 @@
|
||||
package attach_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"image"
|
||||
"image/png"
|
||||
"io"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.golem15.com/golem15/summercms/modules/lagoon"
|
||||
"git.golem15.com/golem15/summercms/modules/lagoon/attach"
|
||||
"gocloud.dev/blob"
|
||||
"gocloud.dev/blob/memblob"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func smokePNG(t *testing.T) []byte {
|
||||
t.Helper()
|
||||
var buf bytes.Buffer
|
||||
if err := png.Encode(&buf, image.NewRGBA(image.Rect(0, 0, 4, 3))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return buf.Bytes()
|
||||
}
|
||||
|
||||
func bucketKeys(t *testing.T, bucket *blob.Bucket) []string {
|
||||
t.Helper()
|
||||
var keys []string
|
||||
iter := bucket.List(nil)
|
||||
for {
|
||||
obj, err := iter.Next(context.Background())
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
keys = append(keys, obj.Key)
|
||||
}
|
||||
return keys
|
||||
}
|
||||
|
||||
// TestStoreSmoke stores a guarded PNG and a body of exactly MaxBytes.
|
||||
func TestStoreSmoke(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("requires testcontainers postgres")
|
||||
}
|
||||
ctx := t.Context()
|
||||
gdb := attachGorm(t)
|
||||
if err := lagoon.Migrate(gdb, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
bucket := memblob.OpenBucket(nil)
|
||||
t.Cleanup(func() { _ = bucket.Close() })
|
||||
|
||||
data := smokePNG(t)
|
||||
f, err := attach.Store(ctx, gdb, bucket, attach.Upload{FileName: `C:\photos\Cover.PNG`, Body: bytes.NewReader(data), Public: false}, attach.Limits{Image: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if f.ID == 0 || f.SortOrder != int(f.ID) {
|
||||
t.Fatalf("sort_order %d, id %d", f.SortOrder, f.ID)
|
||||
}
|
||||
if !strings.HasSuffix(f.DiskName, ".png") || len(f.DiskName) != 26 {
|
||||
t.Fatalf("disk name %q", f.DiskName)
|
||||
}
|
||||
if f.FileName != "Cover.PNG" || f.ContentType != "image/png" || f.FileSize != int64(len(data)) || f.Public() {
|
||||
t.Fatalf("row %+v", f)
|
||||
}
|
||||
var stored attach.File
|
||||
if err := gdb.First(&stored, f.ID).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if stored.SortOrder != int(f.ID) || stored.Public() || stored.AttachmentID != "" {
|
||||
t.Fatalf("stored row %+v", stored)
|
||||
}
|
||||
got, err := bucket.ReadAll(ctx, attach.BlobKey(f.DiskName))
|
||||
if err != nil || !bytes.Equal(got, data) {
|
||||
t.Fatalf("blob %v", err)
|
||||
}
|
||||
|
||||
exact := bytes.Repeat([]byte("a"), 64)
|
||||
g, err := attach.Store(ctx, gdb, bucket, attach.Upload{FileName: "notes.txt", Body: bytes.NewReader(exact)}, attach.Limits{MaxBytes: 64})
|
||||
if err != nil {
|
||||
t.Fatalf("exactly MaxBytes: %v", err)
|
||||
}
|
||||
if g.FileSize != 64 || g.Public() || g.ContentType != "text/plain" {
|
||||
t.Fatalf("row %+v", g)
|
||||
}
|
||||
}
|
||||
|
||||
// TestStoreSmokeRefusals covers the refusals that happen before any row is
|
||||
// written: an SVG in image mode and bodies of MaxBytes+1 bytes.
|
||||
func TestStoreSmokeRefusals(t *testing.T) {
|
||||
ctx := t.Context()
|
||||
bucket := memblob.OpenBucket(nil)
|
||||
t.Cleanup(func() { _ = bucket.Close() })
|
||||
db := &gorm.DB{}
|
||||
|
||||
svg := []byte(`<svg xmlns="http://www.w3.org/2000/svg"><script>alert(1)</script></svg>`)
|
||||
_, err := attach.Store(ctx, db, bucket, attach.Upload{FileName: "x.png", Body: bytes.NewReader(svg)}, attach.Limits{Image: true})
|
||||
if !errors.Is(err, attach.ErrNotImage) {
|
||||
t.Fatalf("svg bytes: %v", err)
|
||||
}
|
||||
_, err = attach.Store(ctx, db, bucket, attach.Upload{FileName: "x.svg", Body: bytes.NewReader(svg)}, attach.Limits{Image: true})
|
||||
if !errors.Is(err, attach.ErrFileType) {
|
||||
t.Fatalf("svg extension: %v", err)
|
||||
}
|
||||
|
||||
_, err = attach.Store(ctx, db, bucket, attach.Upload{FileName: "notes.txt", Body: bytes.NewReader(bytes.Repeat([]byte("a"), 65))}, attach.Limits{MaxBytes: 64})
|
||||
if !errors.Is(err, attach.ErrTooLarge) {
|
||||
t.Fatalf("small body: %v", err)
|
||||
}
|
||||
// Past the 1 MiB read-ahead the limit is enforced while streaming.
|
||||
const limit = 2 << 20
|
||||
_, err = attach.Store(ctx, db, bucket, attach.Upload{FileName: "notes.txt", Body: bytes.NewReader(bytes.Repeat([]byte("a"), limit+1))}, attach.Limits{MaxBytes: limit})
|
||||
if !errors.Is(err, attach.ErrTooLarge) {
|
||||
t.Fatalf("streamed body: %v", err)
|
||||
}
|
||||
if keys := bucketKeys(t, bucket); len(keys) != 0 {
|
||||
t.Fatalf("blobs left behind: %v", keys)
|
||||
}
|
||||
}
|
||||
@@ -126,6 +126,19 @@ var encodeImage = defaultEncodeImage
|
||||
|
||||
// Thumb returns the public URL of a lazily generated thumbnail. The second
|
||||
// call for the same dimensions hits the existing blob and does not resize.
|
||||
// It is PublicURL of ThumbKey; see ThumbKey for the generation rules.
|
||||
func (f *File) Thumb(ctx context.Context, bucket *blob.Bucket, w, h int, mode string) (string, error) {
|
||||
key, err := f.ThumbKey(ctx, bucket, w, h, mode)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return PublicURL(key), nil
|
||||
}
|
||||
|
||||
// ThumbKey returns the blob key of a lazily generated thumbnail, generating
|
||||
// it on first use, so a caller that must not emit a public URL (a protected
|
||||
// file) can stream the thumbnail itself. The second call for the same
|
||||
// dimensions hits the existing blob and does not resize.
|
||||
//
|
||||
// As WinterCMS's File::makeThumb does, an original that is missing from
|
||||
// the bucket, does not decode, or declares more than 4096 by 4096 pixels
|
||||
@@ -133,7 +146,7 @@ var encodeImage = defaultEncodeImage
|
||||
// thumbnail, and the failure is logged at warn level instead of returned:
|
||||
// one unusable upload never fails the listings that show it. An invalid
|
||||
// mode or size, a storage error and an encode failure are still errors.
|
||||
func (f *File) Thumb(ctx context.Context, bucket *blob.Bucket, w, h int, mode string) (string, error) {
|
||||
func (f *File) ThumbKey(ctx context.Context, bucket *blob.Bucket, w, h int, mode string) (string, error) {
|
||||
if f == nil {
|
||||
return "", fmt.Errorf("attach: file is nil")
|
||||
}
|
||||
@@ -162,7 +175,7 @@ func (f *File) Thumb(ctx context.Context, bucket *blob.Bucket, w, h int, mode st
|
||||
return "", fmt.Errorf("attach: thumb exists: %w", err)
|
||||
}
|
||||
if exists {
|
||||
return PublicURL(thumbKey), nil
|
||||
return thumbKey, nil
|
||||
}
|
||||
origKey := part + f.DiskName
|
||||
r, err := bucket.NewReader(ctx, origKey, nil)
|
||||
@@ -219,7 +232,7 @@ func (f *File) Thumb(ctx context.Context, bucket *blob.Bucket, w, h int, mode st
|
||||
}
|
||||
return "", closeErr
|
||||
}
|
||||
return PublicURL(thumbKey), nil
|
||||
return thumbKey, nil
|
||||
}
|
||||
|
||||
var errOriginalTooLarge = errors.New("original image is too large")
|
||||
@@ -239,12 +252,12 @@ var BrokenImagePNG = func() []byte {
|
||||
}()
|
||||
|
||||
// brokenThumb is the catch branch of WinterCMS's File::makeThumb: log the
|
||||
// reason, store BrokenImagePNG under the thumbnail key and return its URL.
|
||||
// reason, store BrokenImagePNG under the thumbnail key and return that key.
|
||||
func brokenThumb(ctx context.Context, bucket *blob.Bucket, f *File, thumbKey string, reason error) (string, error) {
|
||||
slog.Default().WarnContext(ctx, "attach: thumbnail original is unusable, storing the broken-image picture",
|
||||
slog.Uint64("file_id", uint64(f.ID)), slog.String("error", reason.Error()))
|
||||
if err := bucket.WriteAll(ctx, thumbKey, BrokenImagePNG, &blob.WriterOptions{ContentType: "image/png"}); err != nil {
|
||||
return "", fmt.Errorf("attach: broken-image thumb: %w", err)
|
||||
}
|
||||
return PublicURL(thumbKey), nil
|
||||
return thumbKey, nil
|
||||
}
|
||||
|
||||
294
modules/lagoon/deferred.go
Normal file
294
modules/lagoon/deferred.go
Normal file
@@ -0,0 +1,294 @@
|
||||
package lagoon
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"git.golem15.com/golem15/summercms/modules/lagoon/attach"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
// DeferredFileType is the slave_type of a deferred binding that points at a
|
||||
// system_files row: the table name, never a PHP class string.
|
||||
const DeferredFileType = "system_files"
|
||||
|
||||
// DeferredBinding is one row of WinterCMS's deferred_bindings table: a
|
||||
// pending bind (IsBind true) or unbind of the slave SlaveType/SlaveID to the
|
||||
// MasterField relation of a MasterType record that one admin's form session
|
||||
// (SessionKey, BackendUserID) has not saved yet. PivotData holds the
|
||||
// DeferredEnvelope JSON or is nil.
|
||||
type DeferredBinding struct {
|
||||
ID uint `gorm:"column:id;primaryKey"`
|
||||
MasterType string `gorm:"column:master_type"`
|
||||
MasterField string `gorm:"column:master_field"`
|
||||
SlaveType string `gorm:"column:slave_type"`
|
||||
SlaveID string `gorm:"column:slave_id"`
|
||||
PivotData *string `gorm:"column:pivot_data"`
|
||||
SessionKey string `gorm:"column:session_key"`
|
||||
IsBind bool `gorm:"column:is_bind"`
|
||||
BackendUserID uint `gorm:"column:backend_user_id"`
|
||||
CreatedAt time.Time `gorm:"column:created_at"`
|
||||
UpdatedAt time.Time `gorm:"column:updated_at"`
|
||||
}
|
||||
|
||||
// TableName is deferred_bindings.
|
||||
func (DeferredBinding) TableName() string { return "deferred_bindings" }
|
||||
|
||||
// DeferredKey identifies one admin's pending work on one master type: the
|
||||
// form's session key, the backend admin who owns it and the master record's
|
||||
// morph type (MorphType). Every deferred-binding operation is scoped by all
|
||||
// three, so a session key used by another admin, or against another model,
|
||||
// matches nothing.
|
||||
type DeferredKey struct {
|
||||
SessionKey string
|
||||
AdminID uint
|
||||
MasterType string
|
||||
}
|
||||
|
||||
func (k DeferredKey) validate() error {
|
||||
if strings.TrimSpace(k.SessionKey) == "" {
|
||||
return fmt.Errorf("lagoon: deferred binding session key is empty")
|
||||
}
|
||||
if k.AdminID == 0 {
|
||||
return fmt.Errorf("lagoon: deferred binding admin id is zero")
|
||||
}
|
||||
if strings.TrimSpace(k.MasterType) == "" {
|
||||
return fmt.Errorf("lagoon: deferred binding master type is empty")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeferredEnvelope is the framework's pivot_data shape,
|
||||
// {"created":true,"pivot":{...}}. Created marks a slave row that was created
|
||||
// under deferral, which PurgeDeferred deletes with an expired binding; Pivot
|
||||
// holds the pivot values of a deferred belongsToMany link. A binding without
|
||||
// the envelope (nil pivot_data, or JSON without these keys) is a plain link.
|
||||
type DeferredEnvelope struct {
|
||||
Created bool `json:"created,omitempty"`
|
||||
Pivot map[string]any `json:"pivot,omitempty"`
|
||||
}
|
||||
|
||||
// Envelope decodes the binding's pivot_data. Nil or empty pivot_data is the
|
||||
// zero envelope (a plain link); invalid JSON is an error.
|
||||
func (b DeferredBinding) Envelope() (DeferredEnvelope, error) {
|
||||
var env DeferredEnvelope
|
||||
if b.PivotData == nil || strings.TrimSpace(*b.PivotData) == "" {
|
||||
return env, nil
|
||||
}
|
||||
if err := json.Unmarshal([]byte(*b.PivotData), &env); err != nil {
|
||||
return DeferredEnvelope{}, fmt.Errorf("lagoon: deferred binding %d pivot_data: %w", b.ID, err)
|
||||
}
|
||||
return env, nil
|
||||
}
|
||||
|
||||
// MorphType is the master_type or slave_type string of model: its
|
||||
// attach.Owner MorphName when it implements attach.Owner, else its GORM table
|
||||
// name under db's naming strategy. An empty result is an error.
|
||||
func MorphType(db *gorm.DB, model any) (string, error) {
|
||||
if model == nil {
|
||||
return "", fmt.Errorf("lagoon: morph type of a nil model")
|
||||
}
|
||||
if owner, ok := model.(attach.Owner); ok {
|
||||
if name := strings.TrimSpace(owner.MorphName()); name != "" {
|
||||
return name, nil
|
||||
}
|
||||
return "", fmt.Errorf("lagoon: morph type of %T is empty", model)
|
||||
}
|
||||
if db == nil {
|
||||
return "", fmt.Errorf("lagoon: gorm db is nil")
|
||||
}
|
||||
stmt := &gorm.Statement{DB: db}
|
||||
if err := stmt.Parse(model); err != nil {
|
||||
return "", fmt.Errorf("lagoon: morph type of %T: %w", model, err)
|
||||
}
|
||||
if stmt.Schema == nil || strings.TrimSpace(stmt.Schema.Table) == "" {
|
||||
return "", fmt.Errorf("lagoon: morph type of %T is empty", model)
|
||||
}
|
||||
return stmt.Schema.Table, nil
|
||||
}
|
||||
|
||||
func deferredSession(ctx context.Context, tx *gorm.DB) *gorm.DB {
|
||||
if ctx == nil {
|
||||
ctx = context.Background()
|
||||
}
|
||||
return tx.Session(&gorm.Session{NewDB: true, Context: ctx})
|
||||
}
|
||||
|
||||
func checkDeferredArgs(tx *gorm.DB, key DeferredKey, field, slaveType, slaveID string) error {
|
||||
if tx == nil {
|
||||
return fmt.Errorf("lagoon: gorm db is nil")
|
||||
}
|
||||
if err := key.validate(); err != nil {
|
||||
return err
|
||||
}
|
||||
if strings.TrimSpace(field) == "" {
|
||||
return fmt.Errorf("lagoon: deferred binding field is empty")
|
||||
}
|
||||
if strings.TrimSpace(slaveType) == "" {
|
||||
return fmt.Errorf("lagoon: deferred binding slave type is empty")
|
||||
}
|
||||
if strings.TrimSpace(slaveID) == "" {
|
||||
return fmt.Errorf("lagoon: deferred binding slave id is empty")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// findBinding is WinterCMS's DeferredBinding::findBindingRecord, scoped to
|
||||
// the owning admin and locked for the rest of the transaction.
|
||||
func findBinding(ctx context.Context, tx *gorm.DB, key DeferredKey, field, slaveType, slaveID string) (*DeferredBinding, error) {
|
||||
var row DeferredBinding
|
||||
err := deferredSession(ctx, tx).
|
||||
Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Where("master_type = ? AND master_field = ? AND slave_type = ? AND slave_id = ? AND session_key = ? AND backend_user_id = ?",
|
||||
key.MasterType, field, slaveType, slaveID, key.SessionKey, key.AdminID).
|
||||
Order("id").
|
||||
Take(&row).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("lagoon: find deferred binding: %w", err)
|
||||
}
|
||||
return &row, nil
|
||||
}
|
||||
|
||||
func insertBinding(ctx context.Context, tx *gorm.DB, key DeferredKey, field, slaveType, slaveID string, bind bool, env *DeferredEnvelope) error {
|
||||
row := DeferredBinding{
|
||||
MasterType: key.MasterType,
|
||||
MasterField: field,
|
||||
SlaveType: slaveType,
|
||||
SlaveID: slaveID,
|
||||
SessionKey: key.SessionKey,
|
||||
IsBind: bind,
|
||||
BackendUserID: key.AdminID,
|
||||
}
|
||||
if env != nil && (env.Created || len(env.Pivot) > 0) {
|
||||
raw, err := json.Marshal(env)
|
||||
if err != nil {
|
||||
return fmt.Errorf("lagoon: deferred binding envelope: %w", err)
|
||||
}
|
||||
s := string(raw)
|
||||
row.PivotData = &s
|
||||
}
|
||||
if err := deferredSession(ctx, tx).Create(&row).Error; err != nil {
|
||||
return fmt.Errorf("lagoon: insert deferred binding: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func deleteBinding(ctx context.Context, tx *gorm.DB, id uint) error {
|
||||
if err := deferredSession(ctx, tx).Where("id = ?", id).Delete(&DeferredBinding{}).Error; err != nil {
|
||||
return fmt.Errorf("lagoon: delete deferred binding: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeferredBind records that slaveType/slaveID is to be bound to the field
|
||||
// relation of key's unsaved master, with env (nil for a plain link) as its
|
||||
// pivot_data. It ports WinterCMS's DeferredBinding::beforeCreate: a second
|
||||
// bind of the same slave in the same session writes nothing, and a bind
|
||||
// that meets a pending unbind of the same slave cancels it (the unbind row
|
||||
// is deleted and no bind row is written). An empty session key, a zero
|
||||
// AdminID or an empty MasterType is an error.
|
||||
func DeferredBind(ctx context.Context, tx *gorm.DB, key DeferredKey, field, slaveType, slaveID string, env *DeferredEnvelope) error {
|
||||
if err := checkDeferredArgs(tx, key, field, slaveType, slaveID); err != nil {
|
||||
return err
|
||||
}
|
||||
existing, err := findBinding(ctx, tx, key, field, slaveType, slaveID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if existing != nil {
|
||||
if existing.IsBind {
|
||||
return nil
|
||||
}
|
||||
return deleteBinding(ctx, tx, existing.ID)
|
||||
}
|
||||
return insertBinding(ctx, tx, key, field, slaveType, slaveID, true, env)
|
||||
}
|
||||
|
||||
// DeferredUnbind records that slaveType/slaveID is to be unbound from the
|
||||
// field relation of key's master. A second unbind writes nothing. An unbind
|
||||
// that meets a pending bind of the same slave cancels the pair: the bind
|
||||
// row is deleted, no unbind row is written, and the cancelled bind is
|
||||
// returned so the caller can remove the slave it created (a pending upload
|
||||
// or a child created under deferral). Otherwise it returns nil.
|
||||
func DeferredUnbind(ctx context.Context, tx *gorm.DB, key DeferredKey, field, slaveType, slaveID string) (*DeferredBinding, error) {
|
||||
if err := checkDeferredArgs(tx, key, field, slaveType, slaveID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
existing, err := findBinding(ctx, tx, key, field, slaveType, slaveID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if existing != nil {
|
||||
if !existing.IsBind {
|
||||
return nil, nil
|
||||
}
|
||||
if err := deleteBinding(ctx, tx, existing.ID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return existing, nil
|
||||
}
|
||||
return nil, insertBinding(ctx, tx, key, field, slaveType, slaveID, false, nil)
|
||||
}
|
||||
|
||||
// DeferredBindings returns key's bindings whose master_field is one of
|
||||
// fields, in id order, locked FOR UPDATE so two saves with the same session
|
||||
// key are serialized. No fields means no bindings.
|
||||
func DeferredBindings(ctx context.Context, tx *gorm.DB, key DeferredKey, fields []string) ([]DeferredBinding, error) {
|
||||
if tx == nil {
|
||||
return nil, fmt.Errorf("lagoon: gorm db is nil")
|
||||
}
|
||||
if err := key.validate(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(fields) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
var rows []DeferredBinding
|
||||
err := deferredSession(ctx, tx).
|
||||
Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Where("session_key = ? AND backend_user_id = ? AND master_type = ? AND master_field IN ?",
|
||||
key.SessionKey, key.AdminID, key.MasterType, fields).
|
||||
Order("id").
|
||||
Find(&rows).Error
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("lagoon: deferred bindings: %w", err)
|
||||
}
|
||||
return rows, nil
|
||||
}
|
||||
|
||||
// DeferredForget deletes exactly the bindings with the given ids, after
|
||||
// their work has been applied. The caller passes only ids it read for its
|
||||
// own key through DeferredBindings.
|
||||
func DeferredForget(ctx context.Context, tx *gorm.DB, ids []uint) error {
|
||||
if tx == nil {
|
||||
return fmt.Errorf("lagoon: gorm db is nil")
|
||||
}
|
||||
if len(ids) == 0 {
|
||||
return nil
|
||||
}
|
||||
if err := deferredSession(ctx, tx).Where("id IN ?", ids).Delete(&DeferredBinding{}).Error; err != nil {
|
||||
return fmt.Errorf("lagoon: forget deferred bindings: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeferredSlaves returns a subquery selecting the slave_id values of key's
|
||||
// bindings for field and slaveType in one direction (bind true for pending
|
||||
// binds, false for pending unbinds), for use as CAST(pk AS TEXT) IN (?) in
|
||||
// a list query (WinterCMS's withDeferred). slave_id is a text column, so the
|
||||
// primary key must be cast to text for the comparison.
|
||||
func DeferredSlaves(tx *gorm.DB, key DeferredKey, field, slaveType string, bind bool) *gorm.DB {
|
||||
return tx.Session(&gorm.Session{NewDB: true}).
|
||||
Model(&DeferredBinding{}).
|
||||
Select("slave_id").
|
||||
Where("session_key = ? AND backend_user_id = ? AND master_type = ? AND master_field = ? AND slave_type = ? AND is_bind = ?",
|
||||
key.SessionKey, key.AdminID, key.MasterType, field, slaveType, bind)
|
||||
}
|
||||
55
modules/lagoon/deferred_migrations.go
Normal file
55
modules/lagoon/deferred_migrations.go
Normal file
@@ -0,0 +1,55 @@
|
||||
package lagoon
|
||||
|
||||
import (
|
||||
"github.com/go-gormigrate/gormigrate/v2"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// DeferredHistoryID is the history id of the deferred-binding migration set;
|
||||
// its gormigrate table is summer_migrations_summercms_deferred.
|
||||
const DeferredHistoryID = "summercms.deferred"
|
||||
|
||||
// DeferredBindingMigrations creates WinterCMS's deferred_bindings table,
|
||||
// folded from winter/storm's 2013_10_01_000001_Db_Deferred_Bindings.php and
|
||||
// 2021_01_19_000001_Db_Add_Pivot_Data_To_Deferred_Bindings.php. One column is
|
||||
// added: backend_user_id, the backend admin who owns the binding, so a
|
||||
// session key replayed by another admin matches nothing. master_type and
|
||||
// slave_type hold morph type strings (see MorphType), never PHP class names
|
||||
// of framework internals. Rolling the set back drops the table.
|
||||
var DeferredBindingMigrations = []*gormigrate.Migration{
|
||||
{
|
||||
ID: "202610020001_create_deferred_bindings",
|
||||
Migrate: func(tx *gorm.DB) error {
|
||||
stmts := []string{
|
||||
`CREATE TABLE deferred_bindings (
|
||||
id SERIAL PRIMARY KEY,
|
||||
master_type TEXT NOT NULL,
|
||||
master_field TEXT NOT NULL,
|
||||
slave_type TEXT NOT NULL,
|
||||
slave_id TEXT NOT NULL,
|
||||
pivot_data TEXT,
|
||||
session_key TEXT NOT NULL,
|
||||
is_bind BOOLEAN NOT NULL DEFAULT TRUE,
|
||||
backend_user_id INTEGER NOT NULL,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
|
||||
)`,
|
||||
`CREATE INDEX deferred_bindings_master_type_index ON deferred_bindings (master_type)`,
|
||||
`CREATE INDEX deferred_bindings_master_field_index ON deferred_bindings (master_field)`,
|
||||
`CREATE INDEX deferred_bindings_slave_type_index ON deferred_bindings (slave_type)`,
|
||||
`CREATE INDEX deferred_bindings_slave_id_index ON deferred_bindings (slave_id)`,
|
||||
`CREATE INDEX deferred_bindings_session_lookup_index ON deferred_bindings (session_key, backend_user_id, master_type)`,
|
||||
`CREATE INDEX deferred_bindings_created_at_index ON deferred_bindings (created_at)`,
|
||||
}
|
||||
for _, stmt := range stmts {
|
||||
if err := tx.Exec(stmt).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
},
|
||||
Rollback: func(tx *gorm.DB) error {
|
||||
return tx.Exec("DROP TABLE IF EXISTS deferred_bindings").Error
|
||||
},
|
||||
},
|
||||
}
|
||||
116
modules/lagoon/deferred_test.go
Normal file
116
modules/lagoon/deferred_test.go
Normal file
@@ -0,0 +1,116 @@
|
||||
package lagoon
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"image"
|
||||
"image/png"
|
||||
"strconv"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.golem15.com/golem15/summercms/modules/lagoon/attach"
|
||||
"gocloud.dev/blob/memblob"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// TestDeferredUploadPurgeTracer runs the 12.2 storage path end to end: a
|
||||
// guarded upload held by one admin's deferred binding, then purged with its
|
||||
// row and, after commit, its blob.
|
||||
func TestDeferredUploadPurgeTracer(t *testing.T) {
|
||||
db, _ := dedicatedDB(t, "lagoon_deferred_tracer")
|
||||
gdb, err := Use(t.Context(), db)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := Migrate(gdb, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx := t.Context()
|
||||
bucket := memblob.OpenBucket(nil)
|
||||
t.Cleanup(func() { _ = bucket.Close() })
|
||||
|
||||
var img bytes.Buffer
|
||||
if err := png.Encode(&img, image.NewRGBA(image.Rect(0, 0, 2, 2))); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
f, err := attach.Store(ctx, gdb, bucket, attach.Upload{FileName: "cover.png", Body: &img, Public: true}, attach.Limits{Image: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
key := DeferredKey{SessionKey: "tracer-session-key-0123456789abcdef", AdminID: 7, MasterType: "acme_posts"}
|
||||
slaveID := strconv.FormatUint(uint64(f.ID), 10)
|
||||
if err := Transaction(ctx, gdb, func(ctx context.Context, tx *gorm.DB) error {
|
||||
if err := DeferredBind(ctx, tx, key, "cover", DeferredFileType, slaveID, nil); err != nil {
|
||||
return err
|
||||
}
|
||||
// A repeated bind writes nothing.
|
||||
return DeferredBind(ctx, tx, key, "cover", DeferredFileType, slaveID, nil)
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var rows []DeferredBinding
|
||||
if err := gdb.Find(&rows).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(rows) != 1 || !rows[0].IsBind || rows[0].BackendUserID != 7 || rows[0].MasterType != "acme_posts" {
|
||||
t.Fatalf("bindings %+v", rows)
|
||||
}
|
||||
// Another admin's key sees nothing.
|
||||
other := key
|
||||
other.AdminID = 8
|
||||
var seen []DeferredBinding
|
||||
if err := Transaction(ctx, gdb, func(ctx context.Context, tx *gorm.DB) error {
|
||||
var err error
|
||||
seen, err = DeferredBindings(ctx, tx, other, []string{"cover"})
|
||||
return err
|
||||
}); err != nil || len(seen) != 0 {
|
||||
t.Fatalf("foreign admin saw %v (%v)", seen, err)
|
||||
}
|
||||
|
||||
if err := gdb.Exec(`UPDATE deferred_bindings SET created_at = NOW() - INTERVAL '6 days'`).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
res, err := PurgeDeferred(ctx, gdb, bucket, PurgeOptions{Before: time.Now().Add(-5 * 24 * time.Hour)})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if res.Bindings != 1 || res.Files != 1 || res.Skipped != 0 {
|
||||
t.Fatalf("purge result %+v", res)
|
||||
}
|
||||
var count int64
|
||||
if err := gdb.Model(&DeferredBinding{}).Count(&count).Error; err != nil || count != 0 {
|
||||
t.Fatalf("bindings left %d (%v)", count, err)
|
||||
}
|
||||
if err := gdb.Model(&attach.File{}).Where("id = ?", f.ID).Count(&count).Error; err != nil || count != 0 {
|
||||
t.Fatalf("file rows left %d (%v)", count, err)
|
||||
}
|
||||
exists, err := bucket.Exists(ctx, attach.BlobKey(f.DiskName))
|
||||
if err != nil || exists {
|
||||
t.Fatalf("blob still exists=%v (%v)", exists, err)
|
||||
}
|
||||
|
||||
// A bind followed by an unbind of the same slave cancels the pair and
|
||||
// hands the cancelled bind back.
|
||||
var cancelled *DeferredBinding
|
||||
if err := Transaction(ctx, gdb, func(ctx context.Context, tx *gorm.DB) error {
|
||||
env := &DeferredEnvelope{Created: true}
|
||||
if err := DeferredBind(ctx, tx, key, "comments", "acme_comments", "41", env); err != nil {
|
||||
return err
|
||||
}
|
||||
var err error
|
||||
cancelled, err = DeferredUnbind(ctx, tx, key, "comments", "acme_comments", "41")
|
||||
return err
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if cancelled == nil || !cancelled.IsBind {
|
||||
t.Fatalf("cancelled bind %+v", cancelled)
|
||||
}
|
||||
if env, err := cancelled.Envelope(); err != nil || !env.Created {
|
||||
t.Fatalf("envelope %+v (%v)", env, err)
|
||||
}
|
||||
if err := gdb.Model(&DeferredBinding{}).Count(&count).Error; err != nil || count != 0 {
|
||||
t.Fatalf("bindings after cancel %d (%v)", count, err)
|
||||
}
|
||||
}
|
||||
@@ -57,9 +57,10 @@ func migrator(gdb *gorm.DB, pluginID string, migrations []*gormigrate.Migration)
|
||||
}, migrations), nil
|
||||
}
|
||||
|
||||
// Migrate runs the framework-owned system_files, backend-admin and job-queue
|
||||
// (QueueMigrations) sets first, then each plugin's HasMigrations set in
|
||||
// party.Activate order.
|
||||
// Migrate runs the framework-owned system_files, backend-admin,
|
||||
// deferred-binding (DeferredBindingMigrations, history DeferredHistoryID) and
|
||||
// job-queue (QueueMigrations) sets first, then each plugin's HasMigrations set
|
||||
// in party.Activate order.
|
||||
func Migrate(gdb *gorm.DB, plugins []party.Plugin) error {
|
||||
if gdb == nil {
|
||||
return fmt.Errorf("lagoon: gorm db is nil")
|
||||
@@ -78,6 +79,13 @@ func Migrate(gdb *gorm.DB, plugins []party.Plugin) error {
|
||||
if err := admin.Migrate(); err != nil {
|
||||
return fmt.Errorf("lagoon: migrate backend admin: %w", err)
|
||||
}
|
||||
deferred, err := migrator(gdb, DeferredHistoryID, DeferredBindingMigrations)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := deferred.Migrate(); err != nil {
|
||||
return fmt.Errorf("lagoon: migrate deferred bindings: %w", err)
|
||||
}
|
||||
sqlDB, err := gdb.DB()
|
||||
if err != nil {
|
||||
return fmt.Errorf("lagoon: migrate queue: %w", err)
|
||||
|
||||
253
modules/lagoon/purge.go
Normal file
253
modules/lagoon/purge.go
Normal file
@@ -0,0 +1,253 @@
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user