package delivery import ( "compress/gzip" "context" "database/sql" "encoding/base64" "encoding/json" "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. // It reads the file 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 the -wal // past an open snapshot, so the -wal grows until the export is closed. type ArchiveExport struct { db *sql.DB tx *gorm.DB // empty is true when there is nothing to read: no file, or 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 the archive file at path for export and // takes the snapshot the export reads. It never creates the file: with // no file at path, the export has no rows. // // Once it has returned, the file is open, so a rename or a move of it // does not affect the export, which reads the same file under its new // name. // // The transaction lasts 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) { if !fileExists(path) { return &ArchiveExport{empty: true}, nil } db, err := database.OpenSQLite(path, archiveModeExisting) if err != nil { return nil, fmt.Errorf("opening archive %s: %w", 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", 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", 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", path, err) } return &ArchiveExport{db: db, tx: tx, 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. // 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. 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 export's transaction and closes its connection. func (x *ArchiveExport) Close() error { if x.db == nil { return nil } _ = x.tx.Rollback() return x.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 each archived row to w, oldest first, one per line, // separated by commas. func (x *ArchiveExport) writeRows(ctx context.Context, w io.Writer) error { if x.empty { return nil } rows, err := x.tx.WithContext(ctx). Model(&archivedEvent{}).Order("id").Rows() if err != nil { return err } defer func() { _ = rows.Close() }() for sep := "\n"; rows.Next(); sep = ",\n" { var ev archivedEvent err = x.tx.ScanRows(rows, &ev) if err != nil { return err } _, err = io.WriteString(w, sep) if err != nil { return err } err = writeRow(w, &ev) if err != nil { return err } } return 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. func writeRow(w io.Writer, ev *archivedEvent) 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" } line, err := json.Marshal(row) if err != nil { return err } _, err = w.Write(line) return err }