From 3ef410b14052466eba45df58a1c26ec7f74e7651 Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 6 Oct 2026 07:26:40 +0000 Subject: [PATCH] Adopt the shared .golangci.yml and fix the code to it (closes #6) Vendor .golangci.yml byte-identical from sneak/prompts at cc440118 and move the Dockerfile lint phase to golangci-lint v2.14.0 by the digest REPO_POLICIES.md names. Fix the code to that config with flags, help text, output files, SQL and the order of steps unchanged; long functions are split into named steps. Judgement call: the extraction transaction is now rolled back on every early return; the old deferred rollback missed most failures and could dereference a nil transaction. Thirteen //nolint directives (gosec, unconvert, mnd, unqueryvet), each with its reason. Model: opus-5-5 --- .golangci.yml | 99 ++++++++ Dockerfile | 4 +- README.md | 6 +- TODO.md | 10 +- cmd/bsdaily/main.go | 140 +++++++----- internal/bsdaily/bsdaily_test.go | 23 +- internal/bsdaily/config.go | 7 +- internal/bsdaily/copy.go | 36 ++- internal/bsdaily/copy_linux.go | 7 +- internal/bsdaily/disk.go | 19 +- internal/bsdaily/dump.go | 109 ++++++--- internal/bsdaily/extract.go | 375 +++++++++++++++++++++---------- internal/bsdaily/run.go | 261 ++++++++++++++------- internal/bsdaily/snapshot.go | 30 ++- internal/bsdaily/verify.go | 100 +++++++-- 15 files changed, 872 insertions(+), 354 deletions(-) create mode 100644 .golangci.yml diff --git a/.golangci.yml b/.golangci.yml new file mode 100644 index 0000000..1b73eb9 --- /dev/null +++ b/.golangci.yml @@ -0,0 +1,99 @@ +version: "2" + +# Config schema uses the golangci-lint v2 layout (settings live under +# linters.settings, not top-level linters-settings) so that the +# thresholds below are actually applied by golangci-lint >= v2. + +run: + timeout: 5m + modules-download-mode: readonly + +linters: + default: all + enable: + # Successor to the deprecated gomodguard. Named explicitly, rather than + # left to `default: all`, because it carries the module policy below. + - gomodguard_v2 + disable: + # Genuinely incompatible with project patterns + - exhaustruct # Requires all struct fields + - exhaustruct_v5 # Requires all struct fields (successor to exhaustruct) + - godot # Requires comments to end with periods + - wrapcheck # Too verbose for internal packages + - varnamelen # Short names like db, id are idiomatic Go + # Deprecated: the warning is attached to the old name, so it is + # silenced by disabling that name, not by enabling the successor. + - wsl # Deprecated, replaced by wsl_v5 + - gomodguard # Deprecated, replaced by gomodguard_v2 + settings: + lll: + line-length: 88 + funlen: + lines: 80 + statements: 50 + cyclop: + max-complexity: 15 + dupl: + threshold: 100 + depguard: + # Test-support code must not be compiled into the shipped binary. A + # test-support package exists to hand a test privileges the program + # itself must never have, so a file that is not a test must not import + # one. Test files, and the files inside a package whose directory name + # ends in `test`, are where that code belongs, and are exempt. + # + # The deny list below is the one part of this file a repository is + # expected to extend, and the only part it may. depguard matches an + # import path against a list of prefixes, so it cannot be told "any path + # whose last segment ends in test"; a repository's own test-support + # packages have to be named here one at a time, by full import path, + # under a module path that differs from repository to repository. Add + # them; change nothing else. + rules: + test-support: + list-mode: lax + files: + - "$all" + - "!$test" + - "!**/*test/**" + deny: + - pkg: net/http/httptest + desc: >- + Test-support code belongs in test files and in packages whose + directory name ends in test, not in the shipped binary. + # Only decisions already recorded in the Go package defaults are + # listed here. Every entry matches the module path exactly. + gomodguard_v2: + blocked: + - module: github.com/rs/zerolog + recommendations: + - log/slog + reason: "Structured logging is stdlib log/slog." + # One entry per pre-fork module path, because the later releases + # are separate paths. A prefix match would be shorter but would + # also reach github.com/go-redis/redismock, the test double for + # the successor these entries recommend. + - module: github.com/go-redis/redis + recommendations: + - github.com/redis/go-redis/v9 + reason: "Pre-fork module; use the maintained go-redis v9." + - module: github.com/go-redis/redis/v7 + recommendations: + - github.com/redis/go-redis/v9 + reason: "Pre-fork module; use the maintained go-redis v9." + - module: github.com/go-redis/redis/v8 + recommendations: + - github.com/redis/go-redis/v9 + reason: "Pre-fork module; use the maintained go-redis v9." + - module: github.com/sergi/go-diff + recommendations: + - github.com/aymanbagabas/go-udiff + reason: "No unified diff output; use go-udiff." + - module: github.com/hexops/gotextdiff + recommendations: + - github.com/aymanbagabas/go-udiff + reason: "Unmaintained fork; use go-udiff." + +issues: + max-issues-per-linter: 0 + max-same-issues: 0 diff --git a/Dockerfile b/Dockerfile index 85f283b..e7c0d8b 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,8 +1,8 @@ # Lint phase. The linter is invoked directly rather than through `make # lint` or `script/lint`, which are themselves a docker build and would # recurse into a daemon that does not exist in a build step. -# golangci/golangci-lint:v2.12.2-alpine, 2026-06-28 -FROM golangci/golangci-lint:v2.12.2-alpine@sha256:91b27804074a0bacea298707f016911e60cf0cdbc6c7bf5ccacb5f0606d18d60 AS lint +# golangci/golangci-lint:v2.14.0 (Debian-based), 2026-10-06 +FROM golangci/golangci-lint@sha256:ad862ba6b3798cbe0fd9fd7408d498fd74fbd2623a92406b2fd3898faf0bf98f AS lint WORKDIR /src diff --git a/README.md b/README.md index b06f8ae..0aaecdb 100644 --- a/README.md +++ b/README.md @@ -35,9 +35,9 @@ issues are [tracked there](https://git.eeqj.de/sneak/bsdaily/issues). Changes must always be formatted with `make fmt` (`go fmt` for Go, prettier for Markdown), syntactically valid, and must pass the linting defined in the -repository (presently the `golangci-lint` defaults), which can be run with a -`make lint`. The `main` branch is protected and all changes must be made via -[pull requests](https://git.eeqj.de/sneak/bsdaily/pulls) and pass CI to be +repository (the shared `.golangci.yml` from `sneak/prompts`), which can be run +with a `make lint`. The `main` branch is protected and all changes must be made +via [pull requests](https://git.eeqj.de/sneak/bsdaily/pulls) and pass CI to be merged. See [`REPO_POLICIES.md`](REPO_POLICIES.md) for detailed coding standards, diff --git a/TODO.md b/TODO.md index 91cac96..96f3ea9 100644 --- a/TODO.md +++ b/TODO.md @@ -14,12 +14,14 @@ pre-1.0 # Next Step -Add the canonical `.golangci.yml`, move the lint phase to golangci-lint v2.14.0 -in the same commit, and fix the findings it surfaces -(https://git.eeqj.de/sneak/bsdaily/issues/6). +Expand tests beyond the compilation smoke test: unit tests for the extraction, +verification, and atomic-publish paths. # Completed Steps +- 2026-10-06: Added the canonical `.golangci.yml`, moved the lint phase to + golangci-lint v2.14.0, and fixed the code to pass it + (https://git.eeqj.de/sneak/bsdaily/issues/6). - 2026-10-06: Formatted Markdown with prettier: `script/fmt` writes and `script/fmt-check` checks every Markdown file; prettier pinned in `package.json` and `yarn.lock`, installed by `script/bootstrap`, which @@ -46,6 +48,4 @@ in the same commit, and fix the findings it surfaces # Future Steps -- Expand tests beyond the compilation smoke test: unit tests for the extraction, - verification, and atomic-publish paths. - Cut a first SemVer release once compliance and test coverage land. diff --git a/cmd/bsdaily/main.go b/cmd/bsdaily/main.go index 7d5409e..8a1bdf5 100644 --- a/cmd/bsdaily/main.go +++ b/cmd/bsdaily/main.go @@ -1,6 +1,9 @@ +// Package main is the bsdaily command. It extracts one day, or a range +// of days, from the latest daily snapshot. package main import ( + "errors" "fmt" "log/slog" "os" @@ -10,75 +13,112 @@ import ( "github.com/spf13/cobra" ) +var ( + errDateExclusive = errors.New("--date and --from/--to are mutually exclusive") + errFromRequiresTo = errors.New("--from requires --to") + errToRequiresFrom = errors.New("--to requires --from") + errFromAfterTo = errors.New("is after --to") +) + func main() { logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{ Level: slog.LevelInfo, })) slog.SetDefault(logger) - var dateFlag string - var fromFlag string - var toFlag string + var dateFlag, fromFlag, toFlag string rootCmd := &cobra.Command{ Use: "bsdaily", Short: "Extract a single day's data from the latest daily snapshot", SilenceUsage: true, - RunE: func(cmd *cobra.Command, args []string) error { - hasDate := dateFlag != "" - hasFrom := fromFlag != "" - hasTo := toFlag != "" - - // Validate mutual exclusivity - if hasDate && (hasFrom || hasTo) { - return fmt.Errorf("--date and --from/--to are mutually exclusive") - } - if hasFrom != hasTo { - if hasFrom { - return fmt.Errorf("--from requires --to") - } - return fmt.Errorf("--to requires --from") - } - - var targetDates []time.Time - - if hasDate { - t, err := time.Parse("2006-01-02", dateFlag) - if err != nil { - return fmt.Errorf("invalid --date %q (expected YYYY-MM-DD): %w", dateFlag, err) - } - targetDates = []time.Time{t} - } else if hasFrom { - from, err := time.Parse("2006-01-02", fromFlag) - if err != nil { - return fmt.Errorf("invalid --from %q (expected YYYY-MM-DD): %w", fromFlag, err) - } - to, err := time.Parse("2006-01-02", toFlag) - if err != nil { - return fmt.Errorf("invalid --to %q (expected YYYY-MM-DD): %w", toFlag, err) - } - if from.After(to) { - return fmt.Errorf("--from %s is after --to %s", fromFlag, toFlag) - } - for d := from; !d.After(to); d = d.AddDate(0, 0, 1) { - targetDates = append(targetDates, d) - } - } - // else: targetDates remains nil → Run() defaults to snapshot date minus one - - if err := bsdaily.Run(targetDates); err != nil { + RunE: func(_ *cobra.Command, _ []string) error { + targetDates, err := parseTargetDates(dateFlag, fromFlag, toFlag) + if err != nil { return err } + + err = bsdaily.Run(targetDates) + if err != nil { + return err + } + slog.Info("completed successfully") + return nil }, } - rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "", "target date to extract (YYYY-MM-DD); defaults to snapshot date minus one day") - rootCmd.Flags().StringVar(&fromFlag, "from", "", "start of date range to extract (YYYY-MM-DD, inclusive); use with --to") - rootCmd.Flags().StringVar(&toFlag, "to", "", "end of date range to extract (YYYY-MM-DD, inclusive); use with --from") + rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "", + "target date to extract (YYYY-MM-DD); "+ + "defaults to snapshot date minus one day") + rootCmd.Flags().StringVar(&fromFlag, "from", "", + "start of date range to extract (YYYY-MM-DD, inclusive); use with --to") + rootCmd.Flags().StringVar(&toFlag, "to", "", + "end of date range to extract (YYYY-MM-DD, inclusive); use with --from") - if err := rootCmd.Execute(); err != nil { + err := rootCmd.Execute() + if err != nil { os.Exit(1) } } + +// parseTargetDates turns the --date, --from and --to flags into the days +// to extract. It returns nil when none of them is set, which Run takes +// to mean the snapshot date minus one day. +func parseTargetDates(dateFlag, fromFlag, toFlag string) ([]time.Time, error) { + hasDate := dateFlag != "" + hasFrom := fromFlag != "" + hasTo := toFlag != "" + + // Validate mutual exclusivity + if hasDate && (hasFrom || hasTo) { + return nil, errDateExclusive + } + + if hasFrom != hasTo { + if hasFrom { + return nil, errFromRequiresTo + } + + return nil, errToRequiresFrom + } + + if hasDate { + t, err := time.Parse("2006-01-02", dateFlag) + if err != nil { + return nil, fmt.Errorf( + "invalid --date %q (expected YYYY-MM-DD): %w", dateFlag, err) + } + + return []time.Time{t}, nil + } + + if !hasFrom { + return nil, nil + } + + from, err := time.Parse("2006-01-02", fromFlag) + if err != nil { + return nil, fmt.Errorf( + "invalid --from %q (expected YYYY-MM-DD): %w", fromFlag, err) + } + + to, err := time.Parse("2006-01-02", toFlag) + if err != nil { + return nil, fmt.Errorf( + "invalid --to %q (expected YYYY-MM-DD): %w", toFlag, err) + } + + if from.After(to) { + return nil, fmt.Errorf("--from %s %w %s", fromFlag, errFromAfterTo, toFlag) + } + + var targetDates []time.Time + + for d := from; !d.After(to); d = d.AddDate(0, 0, 1) { + targetDates = append(targetDates, d) + } + + return targetDates, nil +} diff --git a/internal/bsdaily/bsdaily_test.go b/internal/bsdaily/bsdaily_test.go index 279e81f..4deda3b 100644 --- a/internal/bsdaily/bsdaily_test.go +++ b/internal/bsdaily/bsdaily_test.go @@ -1,21 +1,32 @@ -package bsdaily +package bsdaily_test -import "testing" +import ( + "testing" + + "git.eeqj.de/sneak/bsdaily/internal/bsdaily" +) // TestCompiles is a minimal smoke test that references the package's exported // surface so that `go test` fails if the package stops compiling. It does not // touch the filesystem or any of the hard-coded production paths. func TestCompiles(t *testing.T) { - if DBFilename == "" || WALFilename == "" || SHMFilename == "" { + t.Parallel() + + if bsdaily.DBFilename == "" || bsdaily.WALFilename == "" || + bsdaily.SHMFilename == "" { t.Fatal("expected database filename constants to be set") } - if SnapshotBase == "" || TmpBase == "" || DailiesBase == "" { + + if bsdaily.SnapshotBase == "" || bsdaily.TmpBase == "" || + bsdaily.DailiesBase == "" { t.Fatal("expected base path constants to be set") } - if MinTmpFreeBytes == 0 || MinDailiesFreeBytes == 0 { + + if bsdaily.MinTmpFreeBytes == 0 || bsdaily.MinDailiesFreeBytes == 0 { t.Fatal("expected free-space thresholds to be set") } - if ErrNoPosts == nil { + + if bsdaily.ErrNoPosts == nil { t.Fatal("expected ErrNoPosts sentinel to be set") } } diff --git a/internal/bsdaily/config.go b/internal/bsdaily/config.go index b534ea8..2a856bd 100644 --- a/internal/bsdaily/config.go +++ b/internal/bsdaily/config.go @@ -1,7 +1,11 @@ +// Package bsdaily extracts single days of firehose data from the latest +// daily ZFS snapshot into zstd-compressed SQL dumps. package bsdaily import "regexp" +// Paths, file names and tuning for a run. They are set for one +// production host; see the README. const ( SnapshotBase = "/srv/berlin.sneak.fs.blueskyarchive/.zfs/snapshot" TmpBase = "/srv/storage/tmp" @@ -27,4 +31,5 @@ const ( verificationHeadLines = 20 ) -var snapshotPattern = regexp.MustCompile(`^zfs-auto-snap_daily-(\d{4}-\d{2}-\d{2})-\d{4}$`) +var snapshotPattern = regexp.MustCompile( + `^zfs-auto-snap_daily-(\d{4}-\d{2}-\d{2})-\d{4}$`) diff --git a/internal/bsdaily/copy.go b/internal/bsdaily/copy.go index 901ca0b..9bf93e5 100644 --- a/internal/bsdaily/copy.go +++ b/internal/bsdaily/copy.go @@ -1,6 +1,7 @@ package bsdaily import ( + "errors" "fmt" "io" "log/slog" @@ -9,20 +10,30 @@ import ( ) const ( - copyBufferSize = 256 * 1024 * 1024 // 256MB buffer for large file copies from fast storage + // 256MB buffer for large file copies from fast storage + copyBufferSize = 256 * 1024 * 1024 + oneMB = 1024 * 1024 oneGB = 1024 * 1024 * 1024 ) +var errShortCopy = errors.New("short copy") + +// CopyFile copies src to dst through a large buffer, pre-allocating dst +// and syncing it to disk before returning. func CopyFile(src, dst string) (err error) { startTime := time.Now() + slog.Info("copying file", "src", src, "dst", dst) + //nolint:gosec // src is a path this package built srcFile, err := os.Open(src) if err != nil { return fmt.Errorf("opening source %s: %w", src, err) } + defer func() { - if cerr := srcFile.Close(); cerr != nil { + cerr := srcFile.Close() + if cerr != nil { slog.Warn("failed to close source file", "src", src, "error", cerr) } }() @@ -37,40 +48,49 @@ func CopyFile(src, dst string) (err error) { applyFileAdvice(srcFile, srcInfo.Size()) } + //nolint:gosec // dst is a path this package built dstFile, err := os.Create(dst) if err != nil { return fmt.Errorf("creating destination %s: %w", dst, err) } + defer func() { - if cerr := dstFile.Close(); cerr != nil && err == nil { + cerr := dstFile.Close() + if cerr != nil && err == nil { err = fmt.Errorf("closing destination %s: %w", dst, cerr) } }() // Pre-allocate space for the destination file to avoid fragmentation - if err := dstFile.Truncate(srcInfo.Size()); err != nil { - slog.Warn("failed to pre-allocate destination file", "error", err) + truncErr := dstFile.Truncate(srcInfo.Size()) + if truncErr != nil { + slog.Warn("failed to pre-allocate destination file", "error", truncErr) } // Use a much larger buffer for NVMe-speed copies buf := make([]byte, copyBufferSize) + written, err := io.CopyBuffer(dstFile, srcFile, buf) if err != nil { return fmt.Errorf("copying data: %w", err) } if written != srcInfo.Size() { - return fmt.Errorf("short copy: wrote %d bytes, expected %d", written, srcInfo.Size()) + return fmt.Errorf("%w: wrote %d bytes, expected %d", + errShortCopy, written, srcInfo.Size()) } - if err := dstFile.Sync(); err != nil { + err = dstFile.Sync() + if err != nil { return fmt.Errorf("syncing destination %s: %w", dst, err) } elapsed := time.Since(startTime) - throughputMBps := float64(written) / elapsed.Seconds() / (1024 * 1024) + throughputMBps := float64(written) / elapsed.Seconds() / oneMB + slog.Info("file copied", "dst", dst, "bytes", written, "elapsed", elapsed.Round(time.Millisecond), "throughput_mbps", fmt.Sprintf("%.1f", throughputMBps)) + return nil } diff --git a/internal/bsdaily/copy_linux.go b/internal/bsdaily/copy_linux.go index d513062..5a49c76 100644 --- a/internal/bsdaily/copy_linux.go +++ b/internal/bsdaily/copy_linux.go @@ -9,8 +9,7 @@ import ( func applyFileAdvice(file *os.File, size int64) { fd := int(file.Fd()) - // POSIX_FADV_SEQUENTIAL = 2 - _ = unix.Fadvise(fd, 0, size, 2) - // POSIX_FADV_WILLNEED = 3 - prefetch file into cache - _ = unix.Fadvise(fd, 0, size, 3) + _ = unix.Fadvise(fd, 0, size, unix.FADV_SEQUENTIAL) + // Prefetch the file into the page cache + _ = unix.Fadvise(fd, 0, size, unix.FADV_WILLNEED) } diff --git a/internal/bsdaily/disk.go b/internal/bsdaily/disk.go index 60cd32e..2523654 100644 --- a/internal/bsdaily/disk.go +++ b/internal/bsdaily/disk.go @@ -1,26 +1,39 @@ package bsdaily import ( + "errors" "fmt" "log/slog" "golang.org/x/sys/unix" ) +var errInsufficientSpace = errors.New("insufficient disk space") + +// CheckFreeSpace returns an error when the filesystem holding path has +// fewer than minBytes bytes available. label names the location in the +// log line and the error. func CheckFreeSpace(path string, minBytes uint64, label string) error { var stat unix.Statfs_t - if err := unix.Statfs(path, &stat); err != nil { + + err := unix.Statfs(path, &stat) + if err != nil { return fmt.Errorf("statfs %s (%s): %w", path, label, err) } + + //nolint:gosec,unconvert // Bsize is never negative; Bavail is signed on FreeBSD free := uint64(stat.Bavail) * uint64(stat.Bsize) freeGB := float64(free) / float64(bytesPerGB) minGB := float64(minBytes) / float64(bytesPerGB) + slog.Info("disk space check", "label", label, "path", path, "free_gb", fmt.Sprintf("%.1f", freeGB), "required_gb", fmt.Sprintf("%.1f", minGB)) + if free < minBytes { - return fmt.Errorf("insufficient disk space on %s (%s): %.1f GB free, need %.1f GB", - path, label, freeGB, minGB) + return fmt.Errorf("%w on %s (%s): %.1f GB free, need %.1f GB", + errInsufficientSpace, path, label, freeGB, minGB) } + return nil } diff --git a/internal/bsdaily/dump.go b/internal/bsdaily/dump.go index 8cbe19e..86a36b2 100644 --- a/internal/bsdaily/dump.go +++ b/internal/bsdaily/dump.go @@ -1,6 +1,8 @@ package bsdaily import ( + "context" + "errors" "fmt" "log/slog" "os" @@ -9,61 +11,44 @@ import ( "strings" ) +var errEmptyOutput = errors.New("compressed output is empty") + +// DumpAndCompress writes a `sqlite3 .dump` of the database at dbPath, +// compressed by zstdmt, to outputPath. func DumpAndCompress(dbPath, outputPath string) (err error) { for _, tool := range []string{"sqlite3", "zstdmt"} { - if _, err := exec.LookPath(tool); err != nil { + _, err = exec.LookPath(tool) + if err != nil { return fmt.Errorf("required tool %q not found in PATH: %w", tool, err) } } - if err := CheckFreeSpace(filepath.Dir(outputPath), MinDailiesFreeBytes, "dailiesBase (pre-dump)"); err != nil { + err = CheckFreeSpace(filepath.Dir(outputPath), MinDailiesFreeBytes, + "dailiesBase (pre-dump)") + if err != nil { return err } + //nolint:gosec // outputPath is a path this package built outFile, err := os.Create(outputPath) if err != nil { return fmt.Errorf("creating output file: %w", err) } + defer func() { - if cerr := outFile.Close(); cerr != nil && err == nil { + cerr := outFile.Close() + if cerr != nil && err == nil { err = fmt.Errorf("closing output: %w", cerr) } }() - // Dump all tables but use INSERT OR IGNORE for mergeable imports - // This preserves all data while allowing multiple dumps to be merged - // Users should import with: zstdcat *.sql.zst | sed 's/INSERT INTO/INSERT OR IGNORE INTO/g' | sqlite3 merged.db - dumpCmd := exec.Command("sqlite3", dbPath, ".dump") - zstdCmd := exec.Command("zstdmt", fmt.Sprintf("-%d", zstdCompressionLevel)) - - pipe, err := dumpCmd.StdoutPipe() + err = runDumpPipeline(context.Background(), dbPath, outFile) if err != nil { - return fmt.Errorf("creating dump stdout pipe: %w", err) - } - zstdCmd.Stdin = pipe - zstdCmd.Stdout = outFile - - var dumpStderr, zstdStderr strings.Builder - dumpCmd.Stderr = &dumpStderr - zstdCmd.Stderr = &zstdStderr - - slog.Info("starting sqlite3 dump and zstdmt compression") - - if err := zstdCmd.Start(); err != nil { - return fmt.Errorf("starting zstdmt: %w", err) - } - if err := dumpCmd.Start(); err != nil { - return fmt.Errorf("starting sqlite3 dump: %w", err) + return err } - if err := dumpCmd.Wait(); err != nil { - return fmt.Errorf("sqlite3 dump failed: %w; stderr: %s", err, dumpStderr.String()) - } - if err := zstdCmd.Wait(); err != nil { - return fmt.Errorf("zstdmt failed: %w; stderr: %s", err, zstdStderr.String()) - } - - if err := outFile.Sync(); err != nil { + err = outFile.Sync() + if err != nil { return fmt.Errorf("syncing output: %w", err) } @@ -71,12 +56,66 @@ func DumpAndCompress(dbPath, outputPath string) (err error) { if err != nil { return fmt.Errorf("stat output: %w", err) } + const bytesPerMB = 1024 * 1024 + slog.Info("compressed output written", "path", outputPath, "size_bytes", info.Size(), "size_mb", info.Size()/bytesPerMB) if info.Size() == 0 { - return fmt.Errorf("compressed output is empty") + return errEmptyOutput + } + + return nil +} + +// runDumpPipeline runs `sqlite3 dbPath .dump | zstdmt` with the +// compressed stream going to outFile, and waits for both to finish. +func runDumpPipeline(ctx context.Context, dbPath string, outFile *os.File) error { + // The dump holds plain INSERT INTO statements. merge_daily_dumps.sh + // rewrites them to INSERT OR IGNORE INTO so that several dumps can be + // merged into one database. + //nolint:gosec // dbPath is a scratch file this package created + dumpCmd := exec.CommandContext(ctx, "sqlite3", dbPath, ".dump") + //nolint:gosec // the argument is built from a constant + zstdCmd := exec.CommandContext(ctx, "zstdmt", + fmt.Sprintf("-%d", zstdCompressionLevel)) + + pipe, err := dumpCmd.StdoutPipe() + if err != nil { + return fmt.Errorf("creating dump stdout pipe: %w", err) + } + + zstdCmd.Stdin = pipe + zstdCmd.Stdout = outFile + + var dumpStderr, zstdStderr strings.Builder + + dumpCmd.Stderr = &dumpStderr + zstdCmd.Stderr = &zstdStderr + + slog.Info("starting sqlite3 dump and zstdmt compression") + + err = zstdCmd.Start() + if err != nil { + return fmt.Errorf("starting zstdmt: %w", err) + } + + err = dumpCmd.Start() + if err != nil { + return fmt.Errorf("starting sqlite3 dump: %w", err) + } + + err = dumpCmd.Wait() + if err != nil { + return fmt.Errorf("sqlite3 dump failed: %w; stderr: %s", + err, dumpStderr.String()) + } + + err = zstdCmd.Wait() + if err != nil { + return fmt.Errorf("zstdmt failed: %w; stderr: %s", + err, zstdStderr.String()) } return nil diff --git a/internal/bsdaily/extract.go b/internal/bsdaily/extract.go index 18dafc0..12b74a3 100644 --- a/internal/bsdaily/extract.go +++ b/internal/bsdaily/extract.go @@ -1,22 +1,29 @@ package bsdaily import ( + "context" "database/sql" "errors" "fmt" "log/slog" "time" + // Registers the "sqlite" driver with database/sql. _ "modernc.org/sqlite" ) +// ErrNoPosts is returned by ExtractDay when the source holds no posts for +// the target day. var ErrNoPosts = errors.New("no posts found for target day") +var errPostCountMismatch = errors.New("post count mismatch") + // ExtractDay opens a new empty database at dstDBPath, attaches srcDBPath, // and copies only the target day's data into it. This is much faster than // pruning a full copy because it only reads/writes the small slice of data // being kept. func ExtractDay(srcDBPath, dstDBPath string, targetDay time.Time) error { + ctx := context.Background() dayStart := targetDay.Format("2006-01-02") + "T00:00:00" dayEnd := targetDay.AddDate(0, 0, 1).Format("2006-01-02") + "T00:00:00" @@ -24,162 +31,282 @@ func ExtractDay(srcDBPath, dstDBPath string, targetDay time.Time) error { // Maximum performance pragmas - we don't care about crash safety for temp files // Use WAL mode for the source attachment to avoid locking issues - pragmas := fmt.Sprintf("?_pragma=journal_mode(WAL)&_pragma=synchronous(OFF)&_pragma=cache_size(%d)&_pragma=foreign_keys(OFF)&_pragma=temp_store(MEMORY)&_pragma=busy_timeout(5000)", sqliteCacheSizeKB) + pragmas := fmt.Sprintf("?_pragma=journal_mode(WAL)"+ + "&_pragma=synchronous(OFF)&_pragma=cache_size(%d)"+ + "&_pragma=foreign_keys(OFF)&_pragma=temp_store(MEMORY)"+ + "&_pragma=busy_timeout(5000)", sqliteCacheSizeKB) + db, err := sql.Open("sqlite", dstDBPath+pragmas) if err != nil { return fmt.Errorf("opening destination database: %w", err) } + defer func() { - if cerr := db.Close(); cerr != nil { - slog.Warn("failed to close destination database", "path", dstDBPath, "error", cerr) + cerr := db.Close() + if cerr != nil { + slog.Warn("failed to close destination database", + "path", dstDBPath, "error", cerr) } }() // Attach source database - if _, err := db.Exec("ATTACH DATABASE ? AS src", srcDBPath); err != nil { + _, err = db.ExecContext(ctx, "ATTACH DATABASE ? AS src", srcDBPath) + if err != nil { return fmt.Errorf("attaching source database: %w", err) } - // Copy table DDL from source - slog.Info("copying table DDL from source") - rows, err := db.Query("SELECT sql FROM src.sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%' ORDER BY name") + err = createTables(ctx, db) if err != nil { - return fmt.Errorf("reading source schema: %w", err) - } - defer func() { - if cerr := rows.Close(); cerr != nil { - slog.Warn("failed to close schema rows", "error", cerr) - } - }() - - var ddlStatements []string - for rows.Next() { - var ddl string - if err := rows.Scan(&ddl); err != nil { - return fmt.Errorf("scanning DDL: %w", err) - } - ddlStatements = append(ddlStatements, ddl) - } - if err := rows.Err(); err != nil { - return fmt.Errorf("iterating DDL rows: %w", err) + return err } - for _, ddl := range ddlStatements { - if _, err := db.Exec(ddl); err != nil { - return fmt.Errorf("creating table: %w\nDDL: %s", err, ddl) - } - } - - // Begin transaction for bulk inserts - tx, err := db.Begin() + postCount, err := insertDayRows(ctx, db, targetDay, dayStart, dayEnd) if err != nil { - return fmt.Errorf("beginning transaction: %w", err) + return err } - defer func() { - if err != nil { - if rerr := tx.Rollback(); rerr != nil && !errors.Is(rerr, sql.ErrTxDone) { - slog.Warn("failed to roll back transaction", "error", rerr) - } - } - }() - // Insert target day's data - slog.Info("inserting posts for target day") - result, err := tx.Exec("INSERT INTO posts SELECT * FROM src.posts WHERE timestamp >= ? AND timestamp < ?", dayStart, dayEnd) + err = createIndexes(ctx, db) if err != nil { - return fmt.Errorf("inserting posts: %w", err) - } - postCount, _ := result.RowsAffected() - slog.Info("inserted posts", "count", postCount) - - if postCount == 0 { - return fmt.Errorf("%w %s - aborting to avoid producing empty output", - ErrNoPosts, targetDay.Format("2006-01-02")) - } - - slog.Info("inserting junction and lookup tables") - if _, err := tx.Exec("INSERT INTO posts_hashtags SELECT * FROM src.posts_hashtags WHERE post_id IN (SELECT id FROM posts)"); err != nil { - return fmt.Errorf("inserting posts_hashtags: %w", err) - } - - if _, err := tx.Exec("INSERT INTO posts_urls SELECT * FROM src.posts_urls WHERE post_id IN (SELECT id FROM posts)"); err != nil { - return fmt.Errorf("inserting posts_urls: %w", err) - } - - if _, err := tx.Exec("INSERT INTO hashtags SELECT * FROM src.hashtags WHERE id IN (SELECT hashtag_id FROM posts_hashtags)"); err != nil { - return fmt.Errorf("inserting hashtags: %w", err) - } - - if _, err := tx.Exec("INSERT INTO urls SELECT * FROM src.urls WHERE id IN (SELECT url_id FROM posts_urls)"); err != nil { - return fmt.Errorf("inserting urls: %w", err) - } - - if _, err := tx.Exec("INSERT INTO users SELECT * FROM src.users WHERE did IN (SELECT user_did FROM posts)"); err != nil { - return fmt.Errorf("inserting users: %w", err) - } - - // Check if media table exists in source and copy if present - var mediaTableExists int - if err := tx.QueryRow("SELECT COUNT(*) FROM src.sqlite_master WHERE type='table' AND name='media'").Scan(&mediaTableExists); err != nil { - slog.Warn("checking for media table", "error", err) - } else if mediaTableExists > 0 { - slog.Info("inserting media entries") - // Get post blob_cids for this day's posts - if _, err := tx.Exec("INSERT INTO media SELECT * FROM src.media WHERE content_hash IN (SELECT blob_cids FROM posts WHERE blob_cids IS NOT NULL)"); err != nil { - slog.Warn("inserting media (may not have matching entries)", "error", err) - } - } - - // Commit the transaction before any further database operations - if err := tx.Commit(); err != nil { - return fmt.Errorf("committing transaction: %w", err) - } - tx = nil // Clear tx to ensure defer doesn't try to rollback - - // Create indexes after bulk insert for speed - slog.Info("creating indexes") - idxRows, err := db.Query("SELECT sql FROM src.sqlite_master WHERE type='index' AND name NOT LIKE 'sqlite_%' AND sql IS NOT NULL ORDER BY name") - if err != nil { - return fmt.Errorf("reading source indexes: %w", err) - } - defer func() { - if cerr := idxRows.Close(); cerr != nil { - slog.Warn("failed to close index rows", "error", cerr) - } - }() - - var idxStatements []string - for idxRows.Next() { - var idxSQL string - if err := idxRows.Scan(&idxSQL); err != nil { - return fmt.Errorf("scanning index DDL: %w", err) - } - idxStatements = append(idxStatements, idxSQL) - } - if err := idxRows.Err(); err != nil { - return fmt.Errorf("iterating index rows: %w", err) - } - - for _, idxSQL := range idxStatements { - if _, err := db.Exec(idxSQL); err != nil { - return fmt.Errorf("creating index: %w\nDDL: %s", err, idxSQL) - } + return err } // Detach source - if _, err := db.Exec("DETACH DATABASE src"); err != nil { + _, err = db.ExecContext(ctx, "DETACH DATABASE src") + if err != nil { return fmt.Errorf("detaching source database: %w", err) } // Verify post count var verifyCount int64 - if err := db.QueryRow("SELECT COUNT(*) FROM posts").Scan(&verifyCount); err != nil { + + err = db.QueryRowContext(ctx, "SELECT COUNT(*) FROM posts").Scan(&verifyCount) + if err != nil { return fmt.Errorf("verifying post count: %w", err) } + if verifyCount != postCount { - return fmt.Errorf("post count mismatch: inserted %d but found %d", postCount, verifyCount) + return fmt.Errorf("%w: inserted %d but found %d", + errPostCountMismatch, postCount, verifyCount) } + slog.Info("extraction complete", "posts", verifyCount) return nil } + +// createTables creates every table of the attached source database, empty, +// in the destination database. +func createTables(ctx context.Context, db *sql.DB) error { + // Copy table DDL from source + slog.Info("copying table DDL from source") + + rows, err := db.QueryContext(ctx, "SELECT sql FROM src.sqlite_master "+ + "WHERE type='table' AND name NOT LIKE 'sqlite_%' ORDER BY name") + if err != nil { + return fmt.Errorf("reading source schema: %w", err) + } + + defer func() { + cerr := rows.Close() + if cerr != nil { + slog.Warn("failed to close schema rows", "error", cerr) + } + }() + + var ddlStatements []string + + for rows.Next() { + var ddl string + + err = rows.Scan(&ddl) + if err != nil { + return fmt.Errorf("scanning DDL: %w", err) + } + + ddlStatements = append(ddlStatements, ddl) + } + + err = rows.Err() + if err != nil { + return fmt.Errorf("iterating DDL rows: %w", err) + } + + for _, ddl := range ddlStatements { + _, err = db.ExecContext(ctx, ddl) + if err != nil { + return fmt.Errorf("creating table: %w\nDDL: %s", err, ddl) + } + } + + return nil +} + +// insertDayRows copies the target day's posts, and the rows they refer +// to, from the source database in one transaction. It returns the number +// of posts copied, or ErrNoPosts when there are none. +func insertDayRows(ctx context.Context, db *sql.DB, targetDay time.Time, + dayStart, dayEnd string, +) (int64, error) { + // Begin transaction for bulk inserts + tx, err := db.BeginTx(ctx, nil) + if err != nil { + return 0, fmt.Errorf("beginning transaction: %w", err) + } + + // After a successful Commit this does nothing and returns ErrTxDone. + defer func() { + rerr := tx.Rollback() + if rerr != nil && !errors.Is(rerr, sql.ErrTxDone) { + slog.Warn("failed to roll back transaction", "error", rerr) + } + }() + + // Insert target day's data + slog.Info("inserting posts for target day") + + //nolint:unqueryvet // copies whole rows; the source defines the columns + result, err := tx.ExecContext(ctx, "INSERT INTO posts SELECT * FROM src.posts "+ + "WHERE timestamp >= ? AND timestamp < ?", dayStart, dayEnd) + if err != nil { + return 0, fmt.Errorf("inserting posts: %w", err) + } + + postCount, _ := result.RowsAffected() + + slog.Info("inserted posts", "count", postCount) + + if postCount == 0 { + return 0, fmt.Errorf("%w %s - aborting to avoid producing empty output", + ErrNoPosts, targetDay.Format("2006-01-02")) + } + + slog.Info("inserting junction and lookup tables") + + err = insertRelatedRows(ctx, tx) + if err != nil { + return 0, err + } + + copyMediaRows(ctx, tx) + + // Commit the transaction before any further database operations + err = tx.Commit() + if err != nil { + return 0, fmt.Errorf("committing transaction: %w", err) + } + + return postCount, nil +} + +// insertRelatedRows copies the hashtag and URL links of the posts already +// inserted, the hashtags and URLs they link to, and the posting users. +// +//nolint:unqueryvet // copies whole rows; the source defines the columns +func insertRelatedRows(ctx context.Context, tx *sql.Tx) error { + _, err := tx.ExecContext(ctx, "INSERT INTO posts_hashtags "+ + "SELECT * FROM src.posts_hashtags WHERE post_id IN (SELECT id FROM posts)") + if err != nil { + return fmt.Errorf("inserting posts_hashtags: %w", err) + } + + _, err = tx.ExecContext(ctx, "INSERT INTO posts_urls "+ + "SELECT * FROM src.posts_urls WHERE post_id IN (SELECT id FROM posts)") + if err != nil { + return fmt.Errorf("inserting posts_urls: %w", err) + } + + _, err = tx.ExecContext(ctx, "INSERT INTO hashtags "+ + "SELECT * FROM src.hashtags "+ + "WHERE id IN (SELECT hashtag_id FROM posts_hashtags)") + if err != nil { + return fmt.Errorf("inserting hashtags: %w", err) + } + + _, err = tx.ExecContext(ctx, "INSERT INTO urls "+ + "SELECT * FROM src.urls WHERE id IN (SELECT url_id FROM posts_urls)") + if err != nil { + return fmt.Errorf("inserting urls: %w", err) + } + + _, err = tx.ExecContext(ctx, "INSERT INTO users "+ + "SELECT * FROM src.users WHERE did IN (SELECT user_did FROM posts)") + if err != nil { + return fmt.Errorf("inserting users: %w", err) + } + + return nil +} + +// copyMediaRows copies the media rows of the posts already inserted, when +// the source has a media table. Failures are logged, not returned. +func copyMediaRows(ctx context.Context, tx *sql.Tx) { + // Check if media table exists in source and copy if present + var mediaTableExists int + + err := tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM src.sqlite_master "+ + "WHERE type='table' AND name='media'").Scan(&mediaTableExists) + if err != nil { + slog.Warn("checking for media table", "error", err) + } else if mediaTableExists > 0 { + slog.Info("inserting media entries") + + // Get post blob_cids for this day's posts + //nolint:unqueryvet // copies whole rows; the source defines the columns + _, err = tx.ExecContext(ctx, "INSERT INTO media SELECT * FROM src.media "+ + "WHERE content_hash IN "+ + "(SELECT blob_cids FROM posts WHERE blob_cids IS NOT NULL)") + if err != nil { + slog.Warn("inserting media (may not have matching entries)", + "error", err) + } + } +} + +// createIndexes creates every index of the attached source database in +// the destination database. Doing this after the bulk insert is faster +// than inserting into indexed tables. +func createIndexes(ctx context.Context, db *sql.DB) error { + // Create indexes after bulk insert for speed + slog.Info("creating indexes") + + idxRows, err := db.QueryContext(ctx, "SELECT sql FROM src.sqlite_master "+ + "WHERE type='index' AND name NOT LIKE 'sqlite_%' AND sql IS NOT NULL "+ + "ORDER BY name") + if err != nil { + return fmt.Errorf("reading source indexes: %w", err) + } + + defer func() { + cerr := idxRows.Close() + if cerr != nil { + slog.Warn("failed to close index rows", "error", cerr) + } + }() + + var idxStatements []string + + for idxRows.Next() { + var idxSQL string + + err = idxRows.Scan(&idxSQL) + if err != nil { + return fmt.Errorf("scanning index DDL: %w", err) + } + + idxStatements = append(idxStatements, idxSQL) + } + + err = idxRows.Err() + if err != nil { + return fmt.Errorf("iterating index rows: %w", err) + } + + for _, idxSQL := range idxStatements { + _, err = db.ExecContext(ctx, idxSQL) + if err != nil { + return fmt.Errorf("creating index: %w\nDDL: %s", err, idxSQL) + } + } + + return nil +} diff --git a/internal/bsdaily/run.go b/internal/bsdaily/run.go index 2093d6f..47896d6 100644 --- a/internal/bsdaily/run.go +++ b/internal/bsdaily/run.go @@ -9,21 +9,29 @@ import ( "time" ) +var errEmptySource = errors.New("source file is empty") + // cleanup removes a temporary file, logging a warning if removal fails so // that leaked scratch files are surfaced rather than silently ignored. A // missing file is not an error. func cleanup(path string) { - if err := os.Remove(path); err != nil && !os.IsNotExist(err) { + err := os.Remove(path) + if err != nil && !os.IsNotExist(err) { slog.Warn("failed to remove temporary file", "path", path, "error", err) } } +// Run writes a compressed SQL dump of each day in targetDates, taken from +// the latest daily snapshot, skipping days that already have one or have +// no posts. With no dates it does the day before the snapshot date. func Run(targetDates []time.Time) error { snapshotDir, snapshotDate, err := FindLatestDailySnapshot() if err != nil { return fmt.Errorf("finding latest snapshot: %w", err) } - slog.Info("found latest daily snapshot", "dir", snapshotDir, "snapshot_date", snapshotDate.Format("2006-01-02")) + + slog.Info("found latest daily snapshot", "dir", snapshotDir, + "snapshot_date", snapshotDate.Format("2006-01-02")) if len(targetDates) == 0 { targetDates = []time.Time{snapshotDate.AddDate(0, 0, -1)} @@ -34,10 +42,13 @@ func Run(targetDates []time.Time) error { "last", targetDates[len(targetDates)-1].Format("2006-01-02")) // Check disk space - if err := CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase"); err != nil { + err = CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase") + if err != nil { return err } - if err := CheckFreeSpace(DailiesBase, MinDailiesFreeBytes, "dailiesBase"); err != nil { + + err = CheckFreeSpace(DailiesBase, MinDailiesFreeBytes, "dailiesBase") + if err != nil { return err } @@ -46,14 +57,53 @@ func Run(targetDates []time.Time) error { if err != nil { return fmt.Errorf("creating temp directory in %s: %w", TmpBase, err) } + slog.Info("created temp directory", "path", tmpDir) + defer func() { slog.Info("cleaning up temp directory", "path", tmpDir) - if err := os.RemoveAll(tmpDir); err != nil { - slog.Error("failed to remove temp directory", "path", tmpDir, "error", err) + + rerr := os.RemoveAll(tmpDir) + if rerr != nil { + slog.Error("failed to remove temp directory", + "path", tmpDir, "error", rerr) } }() + dstDB, err := copySnapshotFiles(snapshotDir, tmpDir) + if err != nil { + return err + } + + // Process each day completely before moving to the next. This ensures + // we don't have multiple SQLite operations competing for the same + // source database. + processed := 0 + skipped := 0 + + for _, targetDay := range targetDates { + written, err := processDay(tmpDir, dstDB, targetDay) + if err != nil { + return err + } + + if written { + processed++ + } else { + skipped++ + } + } + + slog.Info("run summary", "processed", processed, "skipped", skipped, + "total", len(targetDates)) + + return nil +} + +// copySnapshotFiles copies the database, its WAL and, if present, its SHM +// file from snapshotDir into tmpDir, and returns the copied database's +// path. +func copySnapshotFiles(snapshotDir, tmpDir string) (string, error) { // Copy database files from snapshot to temp srcDB := filepath.Join(snapshotDir, DBFilename) srcWAL := filepath.Join(snapshotDir, WALFilename) @@ -65,100 +115,137 @@ func Run(targetDates []time.Time) error { for _, f := range []string{srcDB, srcWAL} { info, err := os.Stat(f) if err != nil { - return fmt.Errorf("source file missing: %s: %w", f, err) + return "", fmt.Errorf("source file missing: %s: %w", f, err) } + if info.Size() == 0 { - return fmt.Errorf("source file is empty: %s", f) + return "", fmt.Errorf("%w: %s", errEmptySource, f) } + slog.Info("source file", "path", f, "size_bytes", info.Size()) } - if err := CopyFile(srcDB, dstDB); err != nil { - return fmt.Errorf("copying database: %w", err) - } - if err := CopyFile(srcWAL, dstWAL); err != nil { - return fmt.Errorf("copying WAL: %w", err) - } - if _, err := os.Stat(srcSHM); err == nil { - if err := CopyFile(srcSHM, dstSHM); err != nil { - return fmt.Errorf("copying SHM: %w", err) - } + err := CopyFile(srcDB, dstDB) + if err != nil { + return "", fmt.Errorf("copying database: %w", err) } - // Process each day completely before moving to the next - // This ensures we don't have multiple SQLite operations competing for the same source database - processed := 0 - skipped := 0 + err = CopyFile(srcWAL, dstWAL) + if err != nil { + return "", fmt.Errorf("copying WAL: %w", err) + } - for _, targetDay := range targetDates { - dayStr := targetDay.Format("2006-01-02") - slog.Info("processing day", "date", dayStr) - - // Check if output already exists - outputDir := filepath.Join(DailiesBase, targetDay.Format("2006-01")) - outputFinal := filepath.Join(outputDir, dayStr+".sql.zst") - if _, err := os.Stat(outputFinal); err == nil { - slog.Info("output already exists, skipping", "path", outputFinal) - skipped++ - continue - } - - // Extract target day into a per-day database - extractedDB := filepath.Join(tmpDir, "extracted-"+dayStr+".db") - slog.Info("extracting target day", "src", dstDB, "dst", extractedDB) - if err := ExtractDay(dstDB, extractedDB, targetDay); err != nil { - if errors.Is(err, ErrNoPosts) { - slog.Warn("no posts found, skipping day", "date", dayStr) - cleanup(extractedDB) - skipped++ - continue - } - return fmt.Errorf("extracting day %s: %w", dayStr, err) - } - - // Dump to SQL and compress - if err := os.MkdirAll(outputDir, 0755); err != nil { - cleanup(extractedDB) - return fmt.Errorf("creating output directory %s: %w", outputDir, err) - } - - outputTmp := filepath.Join(outputDir, "."+dayStr+".sql.zst.tmp") - - slog.Info("dumping and compressing", "tmp_output", outputTmp) - if err := DumpAndCompress(extractedDB, outputTmp); err != nil { - cleanup(outputTmp) - cleanup(extractedDB) - return fmt.Errorf("dump and compress for %s: %w", dayStr, err) - } - - slog.Info("verifying compressed output") - if err := VerifyOutput(outputTmp); err != nil { - cleanup(outputTmp) - cleanup(extractedDB) - return fmt.Errorf("verification failed for %s: %w", dayStr, err) - } - - // Atomic rename to final path - slog.Info("renaming to final output", "from", outputTmp, "to", outputFinal) - if err := os.Rename(outputTmp, outputFinal); err != nil { - cleanup(outputTmp) - cleanup(extractedDB) - return fmt.Errorf("atomic rename for %s: %w", dayStr, err) - } - - info, err := os.Stat(outputFinal) + _, err = os.Stat(srcSHM) + if err == nil { + err = CopyFile(srcSHM, dstSHM) if err != nil { - cleanup(extractedDB) - return fmt.Errorf("stat final output: %w", err) + return "", fmt.Errorf("copying SHM: %w", err) } - slog.Info("day completed", "date", dayStr, "path", outputFinal, "size_bytes", info.Size()) - - // Remove extracted DB to reclaim space immediately - cleanup(extractedDB) - processed++ } - slog.Info("run summary", "processed", processed, "skipped", skipped, "total", len(targetDates)) + return dstDB, nil +} + +// processDay extracts one day from the copied database at dstDB and +// publishes its compressed dump. It returns false, with no error, for a +// day it skips: one whose output already exists or that has no posts. +func processDay(tmpDir, dstDB string, targetDay time.Time) (bool, error) { + dayStr := targetDay.Format("2006-01-02") + + slog.Info("processing day", "date", dayStr) + + // Check if output already exists + outputDir := filepath.Join(DailiesBase, targetDay.Format("2006-01")) + outputFinal := filepath.Join(outputDir, dayStr+".sql.zst") + + _, err := os.Stat(outputFinal) + if err == nil { + slog.Info("output already exists, skipping", "path", outputFinal) + + return false, nil + } + + // Extract target day into a per-day database + extractedDB := filepath.Join(tmpDir, "extracted-"+dayStr+".db") + + slog.Info("extracting target day", "src", dstDB, "dst", extractedDB) + + err = ExtractDay(dstDB, extractedDB, targetDay) + if err != nil { + if errors.Is(err, ErrNoPosts) { + slog.Warn("no posts found, skipping day", "date", dayStr) + cleanup(extractedDB) + + return false, nil + } + + return false, fmt.Errorf("extracting day %s: %w", dayStr, err) + } + + // Dump to SQL and compress + //nolint:gosec,mnd // world-readable on purpose: the dailies tree is published + err = os.MkdirAll(outputDir, 0755) + if err != nil { + cleanup(extractedDB) + + return false, fmt.Errorf("creating output directory %s: %w", outputDir, err) + } + + outputTmp := filepath.Join(outputDir, "."+dayStr+".sql.zst.tmp") + + slog.Info("dumping and compressing", "tmp_output", outputTmp) + + err = DumpAndCompress(extractedDB, outputTmp) + if err != nil { + cleanup(outputTmp) + cleanup(extractedDB) + + return false, fmt.Errorf("dump and compress for %s: %w", dayStr, err) + } + + slog.Info("verifying compressed output") + + err = VerifyOutput(outputTmp) + if err != nil { + cleanup(outputTmp) + cleanup(extractedDB) + + return false, fmt.Errorf("verification failed for %s: %w", dayStr, err) + } + + err = publishOutput(outputTmp, outputFinal, dayStr) + + // Remove extracted DB to reclaim space immediately + cleanup(extractedDB) + + if err != nil { + return false, err + } + + return true, nil +} + +// publishOutput renames the verified temporary output to its final path +// and logs the finished day. The rename is atomic, so the final path never +// holds a partial file. +func publishOutput(outputTmp, outputFinal, dayStr string) error { + // Atomic rename to final path + slog.Info("renaming to final output", "from", outputTmp, "to", outputFinal) + + err := os.Rename(outputTmp, outputFinal) + if err != nil { + cleanup(outputTmp) + + return fmt.Errorf("atomic rename for %s: %w", dayStr, err) + } + + info, err := os.Stat(outputFinal) + if err != nil { + return fmt.Errorf("stat final output: %w", err) + } + + slog.Info("day completed", "date", dayStr, "path", outputFinal, + "size_bytes", info.Size()) return nil } diff --git a/internal/bsdaily/snapshot.go b/internal/bsdaily/snapshot.go index a8a3824..53cf843 100644 --- a/internal/bsdaily/snapshot.go +++ b/internal/bsdaily/snapshot.go @@ -1,6 +1,7 @@ package bsdaily import ( + "errors" "fmt" "log/slog" "os" @@ -9,10 +10,16 @@ import ( "time" ) -func FindLatestDailySnapshot() (dir string, snapshotDate time.Time, err error) { +var errNoSnapshots = errors.New("no daily snapshots found") + +// FindLatestDailySnapshot returns the directory and date of the newest +// daily snapshot in SnapshotBase. It fails when there is none or when the +// newest one has no database file. +func FindLatestDailySnapshot() (string, time.Time, error) { entries, err := os.ReadDir(SnapshotBase) if err != nil { - return "", time.Time{}, fmt.Errorf("reading snapshot directory %s: %w", SnapshotBase, err) + return "", time.Time{}, fmt.Errorf("reading snapshot directory %s: %w", + SnapshotBase, err) } type snapshot struct { @@ -21,24 +28,30 @@ func FindLatestDailySnapshot() (dir string, snapshotDate time.Time, err error) { } var snapshots []snapshot + for _, e := range entries { if !e.IsDir() { continue } + m := snapshotPattern.FindStringSubmatch(e.Name()) if m == nil { continue } + d, err := time.Parse("2006-01-02", m[1]) if err != nil { - slog.Warn("skipping snapshot with unparseable date", "name", e.Name(), "error", err) + slog.Warn("skipping snapshot with unparseable date", + "name", e.Name(), "error", err) + continue } + snapshots = append(snapshots, snapshot{name: e.Name(), date: d}) } if len(snapshots) == 0 { - return "", time.Time{}, fmt.Errorf("no daily snapshots found in %s", SnapshotBase) + return "", time.Time{}, fmt.Errorf("%w in %s", errNoSnapshots, SnapshotBase) } sort.Slice(snapshots, func(i, j int) bool { @@ -46,11 +59,14 @@ func FindLatestDailySnapshot() (dir string, snapshotDate time.Time, err error) { }) latest := snapshots[0] - dir = filepath.Join(SnapshotBase, latest.name) + dir := filepath.Join(SnapshotBase, latest.name) dbPath := filepath.Join(dir, DBFilename) - if _, err := os.Stat(dbPath); err != nil { - return "", time.Time{}, fmt.Errorf("database not found in snapshot %s: %w", dir, err) + + _, err = os.Stat(dbPath) + if err != nil { + return "", time.Time{}, fmt.Errorf("database not found in snapshot %s: %w", + dir, err) } return dir, latest.date, nil diff --git a/internal/bsdaily/verify.go b/internal/bsdaily/verify.go index c419c3d..b0c1f09 100644 --- a/internal/bsdaily/verify.go +++ b/internal/bsdaily/verify.go @@ -1,6 +1,7 @@ package bsdaily import ( + "context" "errors" "fmt" "log/slog" @@ -9,6 +10,11 @@ import ( "strings" ) +var ( + errEmptyDecompressed = errors.New("decompressed content is empty") + errNotSQL = errors.New("decompressed content does not look like SQL") +) + // killCat terminates the zstdcat process, ignoring the benign case where it // has already exited (e.g. after receiving SIGPIPE when head closed the pipe) // and logging any other failure. @@ -16,70 +22,126 @@ func killCat(cmd *exec.Cmd) { if cmd.Process == nil { return } - if err := cmd.Process.Kill(); err != nil && !errors.Is(err, os.ErrProcessDone) { + + err := cmd.Process.Kill() + if err != nil && !errors.Is(err, os.ErrProcessDone) { slog.Warn("failed to kill zstdcat process", "error", err) } } +// VerifyOutput checks that the compressed file at path passes zstdmt's +// integrity test and that its first lines look like SQL. func VerifyOutput(path string) error { + ctx := context.Background() + slog.Info("running zstdmt integrity check") - testCmd := exec.Command("zstdmt", "--test", path) + + //nolint:gosec // path is an output file this package named + testCmd := exec.CommandContext(ctx, "zstdmt", "--test", path) + var testStderr strings.Builder + testCmd.Stderr = &testStderr - if err := testCmd.Run(); err != nil { - return fmt.Errorf("zstdmt --test failed: %w; stderr: %s", err, testStderr.String()) + + err := testCmd.Run() + if err != nil { + return fmt.Errorf("zstdmt --test failed: %w; stderr: %s", + err, testStderr.String()) } + slog.Info("zstdmt integrity check passed") slog.Info("verifying SQL content") - catCmd := exec.Command("zstdcat", path) - headCmd := exec.Command("head", fmt.Sprintf("-%d", verificationHeadLines)) + + content, err := readDecompressedHead(ctx, path) + if err != nil { + return err + } + + err = checkLooksLikeSQL(content) + if err != nil { + return err + } + + slog.Info("SQL content verification passed") + + return nil +} + +// readDecompressedHead returns the first verificationHeadLines lines of +// the decompressed file at path, read through `zstdcat path | head`. +func readDecompressedHead(ctx context.Context, path string) (string, error) { + //nolint:gosec // path is an output file this package named + catCmd := exec.CommandContext(ctx, "zstdcat", path) + //nolint:gosec // the argument is built from a constant + headCmd := exec.CommandContext(ctx, "head", + fmt.Sprintf("-%d", verificationHeadLines)) pipe, err := catCmd.StdoutPipe() if err != nil { - return fmt.Errorf("creating zstdcat pipe: %w", err) + return "", fmt.Errorf("creating zstdcat pipe: %w", err) } + headCmd.Stdin = pipe var headOut strings.Builder + headCmd.Stdout = &headOut - if err := catCmd.Start(); err != nil { - return fmt.Errorf("starting zstdcat: %w", err) + err = catCmd.Start() + if err != nil { + return "", fmt.Errorf("starting zstdcat: %w", err) } - if err := headCmd.Start(); err != nil { + + err = headCmd.Start() + if err != nil { killCat(catCmd) // Clean up if head fails to start - return fmt.Errorf("starting head: %w", err) + + return "", fmt.Errorf("starting head: %w", err) } // Wait for head first (it will exit when it has enough lines) - if err := headCmd.Wait(); err != nil { + err = headCmd.Wait() + if err != nil { killCat(catCmd) - return fmt.Errorf("head command failed: %w", err) + + return "", fmt.Errorf("head command failed: %w", err) } // Kill zstdcat since head closed the pipe (expected SIGPIPE) killCat(catCmd) + _ = catCmd.Wait() // Reap the process - content := headOut.String() + return headOut.String(), nil +} + +// checkLooksLikeSQL returns an error when content is empty or contains +// none of the keywords expected near the start of a `sqlite3 .dump`. +func checkLooksLikeSQL(content string) error { if len(content) == 0 { - return fmt.Errorf("decompressed content is empty") + return errEmptyDecompressed } hasSQLMarker := false - for _, marker := range []string{"BEGIN TRANSACTION", "CREATE TABLE", "INSERT INTO", "PRAGMA"} { + + for _, marker := range []string{ + "BEGIN TRANSACTION", "CREATE TABLE", "INSERT INTO", "PRAGMA", + } { if strings.Contains(content, marker) { hasSQLMarker = true + break } } + const verificationSampleBytes = 200 + if !hasSQLMarker { - return fmt.Errorf("decompressed content does not look like SQL; first %d bytes: %s", - verificationSampleBytes, content[:min(verificationSampleBytes, len(content))]) + return fmt.Errorf("%w; first %d bytes: %s", errNotSQL, + verificationSampleBytes, + content[:min(verificationSampleBytes, len(content))]) } - slog.Info("SQL content verification passed") return nil }