Author SHA1 Message Date
clawbot 3ef410b140 Adopt the shared .golangci.yml and fix the code to it (closes #6)
check / check (push) Waiting to run
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
2026-10-06 10:28:23 +00:00
15 changed files with 872 additions and 354 deletions
+99
View File
@@ -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
+2 -2
View File
@@ -1,8 +1,8 @@
# Lint phase. The linter is invoked directly rather than through `make # Lint phase. The linter is invoked directly rather than through `make
# lint` or `script/lint`, which are themselves a docker build and would # lint` or `script/lint`, which are themselves a docker build and would
# recurse into a daemon that does not exist in a build step. # recurse into a daemon that does not exist in a build step.
# golangci/golangci-lint:v2.12.2-alpine, 2026-06-28 # golangci/golangci-lint:v2.14.0 (Debian-based), 2026-10-06
FROM golangci/golangci-lint:v2.12.2-alpine@sha256:91b27804074a0bacea298707f016911e60cf0cdbc6c7bf5ccacb5f0606d18d60 AS lint FROM golangci/golangci-lint@sha256:ad862ba6b3798cbe0fd9fd7408d498fd74fbd2623a92406b2fd3898faf0bf98f AS lint
WORKDIR /src WORKDIR /src
+3 -3
View File
@@ -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 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 Markdown), syntactically valid, and must pass the linting defined in the
repository (presently the `golangci-lint` defaults), which can be run with a repository (the shared `.golangci.yml` from `sneak/prompts`), which can be run
`make lint`. The `main` branch is protected and all changes must be made via with a `make lint`. The `main` branch is protected and all changes must be made
[pull requests](https://git.eeqj.de/sneak/bsdaily/pulls) and pass CI to be via [pull requests](https://git.eeqj.de/sneak/bsdaily/pulls) and pass CI to be
merged. merged.
See [`REPO_POLICIES.md`](REPO_POLICIES.md) for detailed coding standards, See [`REPO_POLICIES.md`](REPO_POLICIES.md) for detailed coding standards,
+5 -5
View File
@@ -14,12 +14,14 @@ pre-1.0
# Next Step # Next Step
Add the canonical `.golangci.yml`, move the lint phase to golangci-lint v2.14.0 Expand tests beyond the compilation smoke test: unit tests for the extraction,
in the same commit, and fix the findings it surfaces verification, and atomic-publish paths.
(https://git.eeqj.de/sneak/bsdaily/issues/6).
# Completed Steps # 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 - 2026-10-06: Formatted Markdown with prettier: `script/fmt` writes and
`script/fmt-check` checks every Markdown file; prettier pinned in `script/fmt-check` checks every Markdown file; prettier pinned in
`package.json` and `yarn.lock`, installed by `script/bootstrap`, which `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 # 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. - Cut a first SemVer release once compliance and test coverage land.
+90 -50
View File
@@ -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 package main
import ( import (
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"os" "os"
@@ -10,75 +13,112 @@ import (
"github.com/spf13/cobra" "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() { func main() {
logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{ logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{
Level: slog.LevelInfo, Level: slog.LevelInfo,
})) }))
slog.SetDefault(logger) slog.SetDefault(logger)
var dateFlag string var dateFlag, fromFlag, toFlag string
var fromFlag string
var toFlag string
rootCmd := &cobra.Command{ rootCmd := &cobra.Command{
Use: "bsdaily", Use: "bsdaily",
Short: "Extract a single day's data from the latest daily snapshot", Short: "Extract a single day's data from the latest daily snapshot",
SilenceUsage: true, SilenceUsage: true,
RunE: func(cmd *cobra.Command, args []string) error { RunE: func(_ *cobra.Command, _ []string) error {
hasDate := dateFlag != "" targetDates, err := parseTargetDates(dateFlag, fromFlag, toFlag)
hasFrom := fromFlag != "" if err != nil {
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 {
return err return err
} }
err = bsdaily.Run(targetDates)
if err != nil {
return err
}
slog.Info("completed successfully") slog.Info("completed successfully")
return nil return nil
}, },
} }
rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "", "target date to extract (YYYY-MM-DD); defaults to snapshot date minus one day") rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "",
rootCmd.Flags().StringVar(&fromFlag, "from", "", "start of date range to extract (YYYY-MM-DD, inclusive); use with --to") "target date to extract (YYYY-MM-DD); "+
rootCmd.Flags().StringVar(&toFlag, "to", "", "end of date range to extract (YYYY-MM-DD, inclusive); use with --from") "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) 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
}
+17 -6
View File
@@ -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 // 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 // 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. // touch the filesystem or any of the hard-coded production paths.
func TestCompiles(t *testing.T) { 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") 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") 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") t.Fatal("expected free-space thresholds to be set")
} }
if ErrNoPosts == nil {
if bsdaily.ErrNoPosts == nil {
t.Fatal("expected ErrNoPosts sentinel to be set") t.Fatal("expected ErrNoPosts sentinel to be set")
} }
} }
+6 -1
View File
@@ -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 package bsdaily
import "regexp" import "regexp"
// Paths, file names and tuning for a run. They are set for one
// production host; see the README.
const ( const (
SnapshotBase = "/srv/berlin.sneak.fs.blueskyarchive/.zfs/snapshot" SnapshotBase = "/srv/berlin.sneak.fs.blueskyarchive/.zfs/snapshot"
TmpBase = "/srv/storage/tmp" TmpBase = "/srv/storage/tmp"
@@ -27,4 +31,5 @@ const (
verificationHeadLines = 20 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}$`)
+28 -8
View File
@@ -1,6 +1,7 @@
package bsdaily package bsdaily
import ( import (
"errors"
"fmt" "fmt"
"io" "io"
"log/slog" "log/slog"
@@ -9,20 +10,30 @@ import (
) )
const ( 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 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) { func CopyFile(src, dst string) (err error) {
startTime := time.Now() startTime := time.Now()
slog.Info("copying file", "src", src, "dst", dst) slog.Info("copying file", "src", src, "dst", dst)
//nolint:gosec // src is a path this package built
srcFile, err := os.Open(src) srcFile, err := os.Open(src)
if err != nil { if err != nil {
return fmt.Errorf("opening source %s: %w", src, err) return fmt.Errorf("opening source %s: %w", src, err)
} }
defer func() { 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) 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()) applyFileAdvice(srcFile, srcInfo.Size())
} }
//nolint:gosec // dst is a path this package built
dstFile, err := os.Create(dst) dstFile, err := os.Create(dst)
if err != nil { if err != nil {
return fmt.Errorf("creating destination %s: %w", dst, err) return fmt.Errorf("creating destination %s: %w", dst, err)
} }
defer func() { 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) err = fmt.Errorf("closing destination %s: %w", dst, cerr)
} }
}() }()
// Pre-allocate space for the destination file to avoid fragmentation // Pre-allocate space for the destination file to avoid fragmentation
if err := dstFile.Truncate(srcInfo.Size()); err != nil { truncErr := dstFile.Truncate(srcInfo.Size())
slog.Warn("failed to pre-allocate destination file", "error", err) if truncErr != nil {
slog.Warn("failed to pre-allocate destination file", "error", truncErr)
} }
// Use a much larger buffer for NVMe-speed copies // Use a much larger buffer for NVMe-speed copies
buf := make([]byte, copyBufferSize) buf := make([]byte, copyBufferSize)
written, err := io.CopyBuffer(dstFile, srcFile, buf) written, err := io.CopyBuffer(dstFile, srcFile, buf)
if err != nil { if err != nil {
return fmt.Errorf("copying data: %w", err) return fmt.Errorf("copying data: %w", err)
} }
if written != srcInfo.Size() { 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) return fmt.Errorf("syncing destination %s: %w", dst, err)
} }
elapsed := time.Since(startTime) 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, slog.Info("file copied", "dst", dst, "bytes", written,
"elapsed", elapsed.Round(time.Millisecond), "elapsed", elapsed.Round(time.Millisecond),
"throughput_mbps", fmt.Sprintf("%.1f", throughputMBps)) "throughput_mbps", fmt.Sprintf("%.1f", throughputMBps))
return nil return nil
} }
+3 -4
View File
@@ -9,8 +9,7 @@ import (
func applyFileAdvice(file *os.File, size int64) { func applyFileAdvice(file *os.File, size int64) {
fd := int(file.Fd()) fd := int(file.Fd())
// POSIX_FADV_SEQUENTIAL = 2 _ = unix.Fadvise(fd, 0, size, unix.FADV_SEQUENTIAL)
_ = unix.Fadvise(fd, 0, size, 2) // Prefetch the file into the page cache
// POSIX_FADV_WILLNEED = 3 - prefetch file into cache _ = unix.Fadvise(fd, 0, size, unix.FADV_WILLNEED)
_ = unix.Fadvise(fd, 0, size, 3)
} }
+16 -3
View File
@@ -1,26 +1,39 @@
package bsdaily package bsdaily
import ( import (
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"golang.org/x/sys/unix" "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 { func CheckFreeSpace(path string, minBytes uint64, label string) error {
var stat unix.Statfs_t 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) 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) free := uint64(stat.Bavail) * uint64(stat.Bsize)
freeGB := float64(free) / float64(bytesPerGB) freeGB := float64(free) / float64(bytesPerGB)
minGB := float64(minBytes) / float64(bytesPerGB) minGB := float64(minBytes) / float64(bytesPerGB)
slog.Info("disk space check", "label", label, "path", path, slog.Info("disk space check", "label", label, "path", path,
"free_gb", fmt.Sprintf("%.1f", freeGB), "free_gb", fmt.Sprintf("%.1f", freeGB),
"required_gb", fmt.Sprintf("%.1f", minGB)) "required_gb", fmt.Sprintf("%.1f", minGB))
if free < minBytes { if free < minBytes {
return fmt.Errorf("insufficient disk space on %s (%s): %.1f GB free, need %.1f GB", return fmt.Errorf("%w on %s (%s): %.1f GB free, need %.1f GB",
path, label, freeGB, minGB) errInsufficientSpace, path, label, freeGB, minGB)
} }
return nil return nil
} }
+74 -35
View File
@@ -1,6 +1,8 @@
package bsdaily package bsdaily
import ( import (
"context"
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"os" "os"
@@ -9,61 +11,44 @@ import (
"strings" "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) { func DumpAndCompress(dbPath, outputPath string) (err error) {
for _, tool := range []string{"sqlite3", "zstdmt"} { 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) 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 return err
} }
//nolint:gosec // outputPath is a path this package built
outFile, err := os.Create(outputPath) outFile, err := os.Create(outputPath)
if err != nil { if err != nil {
return fmt.Errorf("creating output file: %w", err) return fmt.Errorf("creating output file: %w", err)
} }
defer func() { 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) err = fmt.Errorf("closing output: %w", cerr)
} }
}() }()
// Dump all tables but use INSERT OR IGNORE for mergeable imports err = runDumpPipeline(context.Background(), dbPath, outFile)
// 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()
if err != nil { if err != nil {
return fmt.Errorf("creating dump stdout pipe: %w", err) return 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)
} }
if err := dumpCmd.Wait(); err != nil { err = outFile.Sync()
return fmt.Errorf("sqlite3 dump failed: %w; stderr: %s", err, dumpStderr.String()) if err != nil {
}
if err := zstdCmd.Wait(); err != nil {
return fmt.Errorf("zstdmt failed: %w; stderr: %s", err, zstdStderr.String())
}
if err := outFile.Sync(); err != nil {
return fmt.Errorf("syncing output: %w", err) return fmt.Errorf("syncing output: %w", err)
} }
@@ -71,12 +56,66 @@ func DumpAndCompress(dbPath, outputPath string) (err error) {
if err != nil { if err != nil {
return fmt.Errorf("stat output: %w", err) return fmt.Errorf("stat output: %w", err)
} }
const bytesPerMB = 1024 * 1024 const bytesPerMB = 1024 * 1024
slog.Info("compressed output written", "path", outputPath, slog.Info("compressed output written", "path", outputPath,
"size_bytes", info.Size(), "size_mb", info.Size()/bytesPerMB) "size_bytes", info.Size(), "size_mb", info.Size()/bytesPerMB)
if info.Size() == 0 { 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 return nil
+251 -124
View File
@@ -1,22 +1,29 @@
package bsdaily package bsdaily
import ( import (
"context"
"database/sql" "database/sql"
"errors" "errors"
"fmt" "fmt"
"log/slog" "log/slog"
"time" "time"
// Registers the "sqlite" driver with database/sql.
_ "modernc.org/sqlite" _ "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 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, // 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 // 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 // pruning a full copy because it only reads/writes the small slice of data
// being kept. // being kept.
func ExtractDay(srcDBPath, dstDBPath string, targetDay time.Time) error { func ExtractDay(srcDBPath, dstDBPath string, targetDay time.Time) error {
ctx := context.Background()
dayStart := targetDay.Format("2006-01-02") + "T00:00:00" dayStart := targetDay.Format("2006-01-02") + "T00:00:00"
dayEnd := targetDay.AddDate(0, 0, 1).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 // Maximum performance pragmas - we don't care about crash safety for temp files
// Use WAL mode for the source attachment to avoid locking issues // 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) db, err := sql.Open("sqlite", dstDBPath+pragmas)
if err != nil { if err != nil {
return fmt.Errorf("opening destination database: %w", err) return fmt.Errorf("opening destination database: %w", err)
} }
defer func() { defer func() {
if cerr := db.Close(); cerr != nil { cerr := db.Close()
slog.Warn("failed to close destination database", "path", dstDBPath, "error", cerr) if cerr != nil {
slog.Warn("failed to close destination database",
"path", dstDBPath, "error", cerr)
} }
}() }()
// Attach source database // 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) return fmt.Errorf("attaching source database: %w", err)
} }
// Copy table DDL from source err = createTables(ctx, db)
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")
if err != nil { if err != nil {
return fmt.Errorf("reading source schema: %w", err) return 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)
} }
for _, ddl := range ddlStatements { postCount, err := insertDayRows(ctx, db, targetDay, dayStart, dayEnd)
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()
if err != nil { 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 err = createIndexes(ctx, db)
slog.Info("inserting posts for target day")
result, err := tx.Exec("INSERT INTO posts SELECT * FROM src.posts WHERE timestamp >= ? AND timestamp < ?", dayStart, dayEnd)
if err != nil { if err != nil {
return fmt.Errorf("inserting posts: %w", err) return 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)
}
} }
// Detach source // 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) return fmt.Errorf("detaching source database: %w", err)
} }
// Verify post count // Verify post count
var verifyCount int64 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) return fmt.Errorf("verifying post count: %w", err)
} }
if verifyCount != postCount { 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) slog.Info("extraction complete", "posts", verifyCount)
return nil 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
}
+174 -87
View File
@@ -9,21 +9,29 @@ import (
"time" "time"
) )
var errEmptySource = errors.New("source file is empty")
// cleanup removes a temporary file, logging a warning if removal fails so // cleanup removes a temporary file, logging a warning if removal fails so
// that leaked scratch files are surfaced rather than silently ignored. A // that leaked scratch files are surfaced rather than silently ignored. A
// missing file is not an error. // missing file is not an error.
func cleanup(path string) { 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) 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 { func Run(targetDates []time.Time) error {
snapshotDir, snapshotDate, err := FindLatestDailySnapshot() snapshotDir, snapshotDate, err := FindLatestDailySnapshot()
if err != nil { if err != nil {
return fmt.Errorf("finding latest snapshot: %w", err) 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 { if len(targetDates) == 0 {
targetDates = []time.Time{snapshotDate.AddDate(0, 0, -1)} 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")) "last", targetDates[len(targetDates)-1].Format("2006-01-02"))
// Check disk space // Check disk space
if err := CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase"); err != nil { err = CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase")
if err != nil {
return err return err
} }
if err := CheckFreeSpace(DailiesBase, MinDailiesFreeBytes, "dailiesBase"); err != nil {
err = CheckFreeSpace(DailiesBase, MinDailiesFreeBytes, "dailiesBase")
if err != nil {
return err return err
} }
@@ -46,14 +57,53 @@ func Run(targetDates []time.Time) error {
if err != nil { if err != nil {
return fmt.Errorf("creating temp directory in %s: %w", TmpBase, err) return fmt.Errorf("creating temp directory in %s: %w", TmpBase, err)
} }
slog.Info("created temp directory", "path", tmpDir) slog.Info("created temp directory", "path", tmpDir)
defer func() { defer func() {
slog.Info("cleaning up temp directory", "path", tmpDir) 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 // Copy database files from snapshot to temp
srcDB := filepath.Join(snapshotDir, DBFilename) srcDB := filepath.Join(snapshotDir, DBFilename)
srcWAL := filepath.Join(snapshotDir, WALFilename) srcWAL := filepath.Join(snapshotDir, WALFilename)
@@ -65,100 +115,137 @@ func Run(targetDates []time.Time) error {
for _, f := range []string{srcDB, srcWAL} { for _, f := range []string{srcDB, srcWAL} {
info, err := os.Stat(f) info, err := os.Stat(f)
if err != nil { 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 { 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()) slog.Info("source file", "path", f, "size_bytes", info.Size())
} }
if err := CopyFile(srcDB, dstDB); err != nil { err := CopyFile(srcDB, dstDB)
return fmt.Errorf("copying database: %w", err) if 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)
}
} }
// Process each day completely before moving to the next err = CopyFile(srcWAL, dstWAL)
// This ensures we don't have multiple SQLite operations competing for the same source database if err != nil {
processed := 0 return "", fmt.Errorf("copying WAL: %w", err)
skipped := 0 }
for _, targetDay := range targetDates { _, err = os.Stat(srcSHM)
dayStr := targetDay.Format("2006-01-02") if err == nil {
slog.Info("processing day", "date", dayStr) err = CopyFile(srcSHM, dstSHM)
// 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)
if err != nil { if err != nil {
cleanup(extractedDB) return "", fmt.Errorf("copying SHM: %w", err)
return fmt.Errorf("stat final output: %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 return nil
} }
+23 -7
View File
@@ -1,6 +1,7 @@
package bsdaily package bsdaily
import ( import (
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"os" "os"
@@ -9,10 +10,16 @@ import (
"time" "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) entries, err := os.ReadDir(SnapshotBase)
if err != nil { 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 { type snapshot struct {
@@ -21,24 +28,30 @@ func FindLatestDailySnapshot() (dir string, snapshotDate time.Time, err error) {
} }
var snapshots []snapshot var snapshots []snapshot
for _, e := range entries { for _, e := range entries {
if !e.IsDir() { if !e.IsDir() {
continue continue
} }
m := snapshotPattern.FindStringSubmatch(e.Name()) m := snapshotPattern.FindStringSubmatch(e.Name())
if m == nil { if m == nil {
continue continue
} }
d, err := time.Parse("2006-01-02", m[1]) d, err := time.Parse("2006-01-02", m[1])
if err != nil { 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 continue
} }
snapshots = append(snapshots, snapshot{name: e.Name(), date: d}) snapshots = append(snapshots, snapshot{name: e.Name(), date: d})
} }
if len(snapshots) == 0 { 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 { sort.Slice(snapshots, func(i, j int) bool {
@@ -46,11 +59,14 @@ func FindLatestDailySnapshot() (dir string, snapshotDate time.Time, err error) {
}) })
latest := snapshots[0] latest := snapshots[0]
dir = filepath.Join(SnapshotBase, latest.name) dir := filepath.Join(SnapshotBase, latest.name)
dbPath := filepath.Join(dir, DBFilename) 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 return dir, latest.date, nil
+81 -19
View File
@@ -1,6 +1,7 @@
package bsdaily package bsdaily
import ( import (
"context"
"errors" "errors"
"fmt" "fmt"
"log/slog" "log/slog"
@@ -9,6 +10,11 @@ import (
"strings" "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 // killCat terminates the zstdcat process, ignoring the benign case where it
// has already exited (e.g. after receiving SIGPIPE when head closed the pipe) // has already exited (e.g. after receiving SIGPIPE when head closed the pipe)
// and logging any other failure. // and logging any other failure.
@@ -16,70 +22,126 @@ func killCat(cmd *exec.Cmd) {
if cmd.Process == nil { if cmd.Process == nil {
return 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) 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 { func VerifyOutput(path string) error {
ctx := context.Background()
slog.Info("running zstdmt integrity check") 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 var testStderr strings.Builder
testCmd.Stderr = &testStderr 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("zstdmt integrity check passed")
slog.Info("verifying SQL content") 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() pipe, err := catCmd.StdoutPipe()
if err != nil { if err != nil {
return fmt.Errorf("creating zstdcat pipe: %w", err) return "", fmt.Errorf("creating zstdcat pipe: %w", err)
} }
headCmd.Stdin = pipe headCmd.Stdin = pipe
var headOut strings.Builder var headOut strings.Builder
headCmd.Stdout = &headOut headCmd.Stdout = &headOut
if err := catCmd.Start(); err != nil { err = catCmd.Start()
return fmt.Errorf("starting zstdcat: %w", err) 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 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) // 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) 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) // Kill zstdcat since head closed the pipe (expected SIGPIPE)
killCat(catCmd) killCat(catCmd)
_ = catCmd.Wait() // Reap the process _ = 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 { if len(content) == 0 {
return fmt.Errorf("decompressed content is empty") return errEmptyDecompressed
} }
hasSQLMarker := false 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) { if strings.Contains(content, marker) {
hasSQLMarker = true hasSQLMarker = true
break break
} }
} }
const verificationSampleBytes = 200 const verificationSampleBytes = 200
if !hasSQLMarker { if !hasSQLMarker {
return fmt.Errorf("decompressed content does not look like SQL; first %d bytes: %s", return fmt.Errorf("%w; first %d bytes: %s", errNotSQL,
verificationSampleBytes, content[:min(verificationSampleBytes, len(content))]) verificationSampleBytes,
content[:min(verificationSampleBytes, len(content))])
} }
slog.Info("SQL content verification passed")
return nil return nil
} }