Author SHA1 Message Date
clawbot f5c7614768 Stamp the git tag or short commit into the binary (closes #4)
check / check (push) Waiting to run
A plain `docker build .` now stamps bsdaily's version: the VERSION
build argument when one is given, otherwise `git describe --tags
--always` on the .git in the build context, as the canonical
Dockerfile does. The build fails if .git is there and the version
still comes out empty, dev or unknown. A host `make` build stamps the
same `git describe`, or dev. bsdaily logs the version on the first
line of every run and prints it with --version; a build that stamps
nothing, or an empty value, reports dev.

Judgement call: the build line also takes the canonical -trimpath and
-s -w.
One //nolint (gochecknoglobals): -X can only set a package-level
variable.

Model: opus-5-5
Co-authored-by: clawbot <sneak+clawbot@sneak.cloud>
2026-10-06 18:41:51 +02:00
clawbot c16f177575 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 14:24:48 +02:00
clawbot 925a3896f7 Format Markdown with prettier in make fmt and make fmt-check (closes #7)
check / check (push) Successful in 5m9s
script/fmt now also runs prettier over every Markdown file, and
script/fmt-check checks them without writing. prettier is pinned in
package.json and yarn.lock. script/bootstrap keeps its Go and apt
handling and gains the canonical node and yarn install (nvm from a
hash-checked archive, yarn through corepack), at the pinned versions
when they are absent; a node or yarn already installed is used as is.
Both fmt scripts find yarn the canonical way. The new files come from
sneak/prompts at cc440118c876. README.md and TODO.md are rewrapped by
make fmt, with no wording changes.

Deviation: package.json drops the canonical "license": "MIT" line; this
repo is WTFPL.

Model: opus-5-5
Co-authored-by: clawbot <sneak+clawbot@sneak.cloud>
2026-10-06 12:24:51 +02:00
16 changed files with 930 additions and 360 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
+21 -5
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
@@ -47,7 +47,9 @@ FROM golang:1.26.4-alpine@sha256:3ad57304ad93bbec8548a0437ad9e06a455660655d9af01
COPY --from=lint /src/go.sum /dev/null COPY --from=lint /src/go.sum /dev/null
COPY --from=test /src/go.sum /dev/null COPY --from=test /src/go.sum /dev/null
ARG VERSION=dev RUN apk add --no-cache git
# A tar-stream context keeps the sender's file owners, which git refuses.
RUN git config --system --add safe.directory /src
WORKDIR /src WORKDIR /src
@@ -58,8 +60,22 @@ RUN go mod download
# Copy source code # Copy source code
COPY . . COPY . .
# Build (pure Go, no CGO required since we use modernc.org/sqlite) # Build (pure Go, no CGO required since we use modernc.org/sqlite).
RUN CGO_ENABLED=0 go build -o /bsdaily ./cmd/bsdaily # The VERSION build arg when one is given, otherwise
# `git describe --tags --always` on the .git in the build context. With
# .git present, a version that is still empty, dev or unknown fails the
# build: git is missing or could not read the checkout.
ARG VERSION
RUN VERSION="${VERSION:-$(git describe --tags --always)}"; \
if [ -e .git ]; then \
case "$VERSION" in ""|dev|unknown) \
echo "version is '$VERSION' although .git is present" >&2; \
exit 1 ;; \
esac; \
fi; \
CGO_ENABLED=0 go build -trimpath \
-ldflags="-s -w -X main.Version=${VERSION}" \
-o /bsdaily ./cmd/bsdaily
# Runtime stage # Runtime stage
# alpine:3.21, 2026-06-28 # alpine:3.21, 2026-06-28
+5 -3
View File
@@ -1,7 +1,9 @@
.PHONY: all bootstrap setup check test lint fmt fmt-check build clean deps test-coverage test-integration install release release-snapshot docker hooks .PHONY: all bootstrap setup check test lint fmt fmt-check build clean deps test-coverage test-integration install release release-snapshot docker hooks
# Version number # Stamped into the binary: the same `git describe` a plain `docker build .`
VERSION := 0.1.0-dev # runs, or dev when it prints nothing (outside a git checkout, or where git is
# missing). ?= so that a VERSION already in the environment takes precedence.
VERSION ?= $(or $(shell git describe --tags --always 2>/dev/null),dev)
# Default target # Default target
all: bsdaily all: bsdaily
@@ -36,7 +38,7 @@ lint:
# Build binary (pure Go; no CGO required since we use modernc.org/sqlite). # Build binary (pure Go; no CGO required since we use modernc.org/sqlite).
bsdaily: internal/*/*.go cmd/bsdaily/*.go bsdaily: internal/*/*.go cmd/bsdaily/*.go
CGO_ENABLED=0 go build -o $@ ./cmd/bsdaily CGO_ENABLED=0 go build -ldflags "-X main.Version=$(VERSION)" -o $@ ./cmd/bsdaily
# Clean build artifacts. # Clean build artifacts.
clean: clean:
+20 -4
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,
@@ -189,6 +189,7 @@ A single run proceeds as follows:
bsdaily # extract the snapshot date minus one day bsdaily # extract the snapshot date minus one day
bsdaily --date 2026-06-27 # extract a single specific day bsdaily --date 2026-06-27 # extract a single specific day
bsdaily --from 2026-06-01 --to 2026-06-27 # extract an inclusive range bsdaily --from 2026-06-01 --to 2026-06-27 # extract an inclusive range
bsdaily --version # print the version and exit
``` ```
Flags: Flags:
@@ -197,9 +198,24 @@ Flags:
`--from`/`--to`. `--from`/`--to`.
- `--from YYYY-MM-DD` — start of an inclusive range (requires `--to`). - `--from YYYY-MM-DD` — start of an inclusive range (requires `--to`).
- `--to YYYY-MM-DD` — end of an inclusive range (requires `--from`). - `--to YYYY-MM-DD` — end of an inclusive range (requires `--from`).
- `-v`, `--version` — print the version and exit.
With no flags, the tool extracts the day before the latest snapshot. All With no flags, the tool extracts the day before the latest snapshot. All
progress is logged as structured `slog` text to stderr. progress is logged as structured `slog` text to stderr; the first line of every
run carries the version.
The version is set at link time and depends on how the binary was built:
- `docker build .` takes it from the `VERSION` build argument when one is given,
otherwise from `git describe --tags --always` on the `.git` in the build
context. The build fails if `.git` is there and the version still comes out
empty, `dev` or `unknown`. With neither `.git` nor `VERSION`, the binary
reports `dev`.
- `script/docker`, `script/cibuild` and `make docker` pass the host's
`git describe --tags --always --dirty` as `VERSION`, so on a modified tree the
version ends in `-dirty`. When that prints nothing, they pass `unknown`.
- `make` stamps the host's `git describe --tags --always`, without `-dirty`, or
`dev` when that prints nothing.
## Merging dumps back into a database ## Merging dumps back into a database
+9 -5
View File
@@ -14,12 +14,18 @@ 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: A plain `docker build .` and a host `make` build stamp the git tag
or short commit into the binary, which `bsdaily` logs on the first line of
every run and prints with `--version`
(https://git.eeqj.de/sneak/bsdaily/issues/4).
- 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 +52,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.
+103 -49
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,126 @@ 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")
)
// Version is the git tag or short commit, set at link time with
// -X main.Version=... by the Dockerfile and the Makefile. A build that sets
// nothing, or sets it empty, reports dev.
//
//nolint:gochecknoglobals // -X can only set a package-level variable
var Version string
func main() { func main() {
if Version == "" {
Version = "dev"
}
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",
Version: Version,
SilenceUsage: true, SilenceUsage: true,
RunE: func(cmd *cobra.Command, args []string) error { RunE: func(_ *cobra.Command, _ []string) error {
hasDate := dateFlag != "" slog.Info("starting", "version", Version)
hasFrom := fromFlag != ""
hasTo := toFlag != ""
// Validate mutual exclusivity targetDates, err := parseTargetDates(dateFlag, fromFlag, toFlag)
if hasDate && (hasFrom || hasTo) { if err != nil {
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
} }