package delivery import ( "compress/gzip" "context" "database/sql" "encoding/base64" "encoding/json" "errors" "fmt" "io" "log/slog" "path/filepath" "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 opened for download: // every one of its files, each read on its own connection inside one // read-only transaction, so it writes out the archive as it stood when // OpenArchiveExport returned. // // Archives are in WAL mode, where a reader works from a snapshot and // never blocks a writer: archive writes go on while an export is open, // and the export does not see them. SQLite cannot checkpoint a -wal // past an open snapshot, so a file's -wal grows until the export has // written out that file. type ArchiveExport struct { files []*exportFile } // 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"` } // OpenArchiveExport opens every one of a database target's archive // files for export, given the path ArchivePath gives it (see // archiveFiles), and takes the snapshots the export reads. It never // creates a file: with no files, the export has no rows. // // Once it has returned, the files are open, so a rename or a move of // them does not affect the export, which reads the same files under // their new names. // // The transactions last as long as ctx does, so ctx must last for the // whole export. func OpenArchiveExport( ctx context.Context, path string, log *slog.Logger, ) (*ArchiveExport, error) { files, err := archiveFiles(path) if err != nil { return nil, err } x := &ArchiveExport{} for _, file := range files { f, err := openExportFile(ctx, file, log) if err != nil { _ = x.Close() return nil, err } x.files = append(x.files, f) } return x, nil } // openExportFile opens one archive file for an export and takes its // snapshot. 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. 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() } // Close ends the transactions of the files still open and closes // their connections. func (x *ArchiveExport) Close() error { errs := make([]error, 0, len(x.files)) for _, f := range x.files { errs = append(errs, f.close()) } return errors.Join(errs...) } // close ends the file's transaction and closes its connection, once. func (f *exportFile) close() error { if f.db == nil { return nil } _ = f.tx.Rollback() err := f.db.Close() f.db = nil return err } // 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, and closes each file once its rows are written. func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error { sep := "\n" for _, f := range x.files { var err error sep, err = f.writeRows(ctx, w, sep) if err != nil { return err } 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 }