package delivery import ( "compress/gzip" "context" "database/sql" "encoding/base64" "encoding/json" "errors" "fmt" "io" "io/fs" "log/slog" "os" "path/filepath" "sync" "time" "unicode/utf8" "gorm.io/driver/sqlite" "gorm.io/gorm" "sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/gormlog" ) // archiveTableQuery counts the archive's table: 0 when the archive // writer has created the file but not yet the table in it. const archiveTableQuery = "SELECT count(*) FROM sqlite_master " + "WHERE type = 'table' AND name = 'archived_events'" // ArchivePath returns where a database target's archive file is: in // the data directory, beside the webhook's event database, under the // name ArchiveFileName gives it. func ArchivePath( dbMgr *database.WebhookDBManager, webhook *database.Webhook, target *database.Target, ) string { return filepath.Join( filepath.Dir(dbMgr.DBPath(webhook.ID)), ArchiveFileName(webhook.Name, target.Name, target.ID), ) } // ArchiveExportFileName returns the name a database target's archive // downloads under: // archive-WEBHOOKNAME-TARGETNAME-YYYYMMDDTHHMMSSZ.json.gz, the names // made safe as in ArchiveFileName and the time in UTC. func ArchiveExportFileName( webhookName, targetName string, at time.Time, ) string { return "archive-" + archiveNamePart(webhookName) + "-" + archiveNamePart(targetName) + "-" + at.UTC().Format("20060102T150405Z") + ".json.gz" } // ArchiveExport is a database target's archive listed for download. It // opens one of the target's files at a time, only when its rows are // about to be written out, and closes it before it opens the next, so // an export holds at most one file open however many the target has. // // Each file is read on its own connection inside one read-only // transaction, so its rows are written out as the file stood when it // was opened. Archives are in WAL mode, where a reader works from a // snapshot and never blocks a writer: archive writes go on while a file // is open, and the export does not see them. SQLite cannot checkpoint // a -wal past an open snapshot, so the open file's -wal grows until the // export has written that file out. type ArchiveExport struct { // periods are the periods of the target's files when the export // was listed, "" for the file without one, in the order // archiveFiles lists them. periods []string // lock is held while currentPath is called and a file is opened, // so that a rename, which holds it too, cannot move the file in // between. lock sync.Locker // currentPath returns the path ArchivePath gives the target under // the names stored for it now, which a rename may have changed since // the export was listed. currentPath func() (string, error) log *slog.Logger } // exportFile is one archive file opened for an export. type exportFile struct { db *sql.DB tx *gorm.DB // period is the period in the file's name, "" for none. period string // empty is true for a file without the archive's table yet. empty bool } // exportedName is how an export names its webhook and its target. type exportedName struct { ID string `json:"id"` Name string `json:"name"` } // NewArchiveExport lists a database target's archive files for export, // given the path ArchivePath gives it (see archiveFiles). It opens none // of them. Its caller holds lock, which every rename of the target's // files runs under, from reading the names path is made of until it // returns, so the files it lists are the ones those names give. // // WriteGzipJSON, called without lock held, finds each file again by its // period under the path currentPath gives, holding lock while it does // and while it opens the file, so a rename during the export loses no // file. A file that is gone by then, emptied by the sweep or moved // away, is skipped. The export never creates a file: with no files, it // has no rows. func NewArchiveExport( path string, lock sync.Locker, currentPath func() (string, error), log *slog.Logger, ) (*ArchiveExport, error) { files, err := archiveFiles(path) if err != nil { return nil, err } x := &ArchiveExport{lock: lock, currentPath: currentPath, log: log} for _, file := range files { x.periods = append(x.periods, file.period) } return x, nil } // openExportFile opens one archive file for an export and takes its // snapshot. The transaction lasts as long as ctx does. func openExportFile( ctx context.Context, file archiveFile, log *slog.Logger, ) (*exportFile, error) { db, err := database.OpenSQLite(file.path, archiveModeExisting) if err != nil { return nil, fmt.Errorf("opening archive %s: %w", file.path, err) } gdb, err := gorm.Open( sqlite.Dialector{Conn: db}, &gorm.Config{ // Never leave this at GORM's default. See // internal/gormlog. Logger: gormlog.New(log), }, ) if err != nil { _ = db.Close() return nil, fmt.Errorf("opening archive %s: %w", file.path, err) } // ReadOnly makes the driver begin a deferred transaction in place // of the BEGIN IMMEDIATE the connection string asks for, so the // export never takes the archive's write lock. tx := gdb.WithContext(ctx).Begin(&sql.TxOptions{ReadOnly: true}) if tx.Error != nil { _ = db.Close() return nil, fmt.Errorf( "reading archive %s: %w", file.path, tx.Error, ) } // The transaction's first read is what takes the snapshot. var tables int err = tx.Raw(archiveTableQuery).Row().Scan(&tables) if err != nil { _ = tx.Rollback() _ = db.Close() return nil, fmt.Errorf("reading archive %s: %w", file.path, err) } return &exportFile{ db: db, tx: tx, period: file.period, empty: tables == 0, }, nil } // WriteGzipJSON writes the export to w as one gzipped JSON object: // webhook and target, each an id and a name; exported_at; and // archived_events, one object per archived row, keyed by column name, // the files in the order archiveFiles lists them. A row from a file // named for a period has "period" beside its columns. A body that is // not valid UTF-8 cannot be a JSON string, so it is written in base64, // with "body_encoding": "base64" beside it. // // Each row is written out before the next is read, so neither the // archive nor its JSON is ever held in memory whole, and each file is // closed once its rows are written, before the next is opened. When it // returns, no file is open. After an error the gzip stream is left // unfinished, so what was written does not decompress as a whole file. func (x *ArchiveExport) WriteGzipJSON( ctx context.Context, w io.Writer, webhook *database.Webhook, target *database.Target, exportedAt time.Time, ) error { head, err := json.Marshal(map[string]any{ "webhook": exportedName{ID: webhook.ID, Name: webhook.Name}, "target": exportedName{ID: target.ID, Name: target.Name}, "exported_at": exportedAt.UTC(), }) if err != nil { return fmt.Errorf("encoding archive export: %w", err) } zw := gzip.NewWriter(w) err = x.writeJSON(ctx, zw, head) if err != nil { return fmt.Errorf("writing archive export: %w", err) } return zw.Close() } // openFile finds the target's archive file for period under the path // currentPath gives now, and opens it for the export, holding x.lock // for both. For a file that is gone, the error wraps fs.ErrNotExist. func (x *ArchiveExport) openFile( ctx context.Context, period string, ) (*exportFile, error) { x.lock.Lock() defer x.lock.Unlock() path, err := x.currentPath() if err != nil { return nil, fmt.Errorf("finding archive file: %w", err) } file := archiveFile{path: archivePeriodPath(path, period), period: period} _, err = os.Stat(file.path) if err != nil { return nil, err } return openExportFile(ctx, file, x.log) } // close ends the file's transaction and closes its connection. func (f *exportFile) close() error { _ = f.tx.Rollback() return f.db.Close() } // writeJSON writes head with archived_events added as its last key, // the rows going into it one at a time. func (x *ArchiveExport) writeJSON( ctx context.Context, w io.Writer, head []byte, ) error { // head goes out without its closing brace, so that // archived_events can follow it. _, err := w.Write(head[:len(head)-1]) if err != nil { return err } _, err = io.WriteString(w, `,"archived_events":[`) if err != nil { return err } err = x.writeRows(ctx, w) if err != nil { return err } _, err = io.WriteString(w, "\n]}\n") return err } // writeRows writes the archived rows of each file to w, one per line, // separated by commas, opening each file in turn and closing it once // its rows are written. func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error { sep := "\n" for _, period := range x.periods { f, err := x.openFile(ctx, period) if errors.Is(err, fs.ErrNotExist) { continue } if err != nil { return err } sep, err = f.writeRows(ctx, w, sep) err = errors.Join(err, f.close()) if err != nil { return err } } return nil } // writeRows writes the file's archived rows to w, oldest first, the // first after sep and each other after ",\n". It returns what goes // before the next row: sep again when the file had no rows. func (f *exportFile) writeRows( ctx context.Context, w io.Writer, sep string, ) (string, error) { if f.empty { return sep, nil } rows, err := f.tx.WithContext(ctx). Model(&archivedEvent{}).Order("id").Rows() if err != nil { return "", err } defer func() { _ = rows.Close() }() for ; rows.Next(); sep = ",\n" { var ev archivedEvent err = f.tx.ScanRows(rows, &ev) if err != nil { return "", err } _, err = io.WriteString(w, sep) if err != nil { return "", err } err = writeRow(w, &ev, f.period) if err != nil { return "", err } } return sep, rows.Err() } // writeRow writes an archived row to w as a JSON object keyed by // column name, its body in base64 when it is not valid UTF-8, with // the period of its file beside them unless that is "". func writeRow(w io.Writer, ev *archivedEvent, period string) error { row := map[string]any{ "id": ev.ID, "event_id": ev.EventID, "webhook_id": ev.WebhookID, "entrypoint_id": ev.EntrypointID, "method": ev.Method, "headers": ev.Headers, "body": ev.Body, "content_type": ev.ContentType, "archived_at": ev.ArchivedAt.UTC(), } if !utf8.ValidString(ev.Body) { row["body"] = base64.StdEncoding.EncodeToString([]byte(ev.Body)) row["body_encoding"] = "base64" } if period != "" { row["period"] = period } line, err := json.Marshal(row) if err != nil { return err } _, err = w.Write(line) return err }