Author SHA1 Message Date
sneak 8a45947bab Bring the repo up to the standard layout (closes #1)
check / check (push) Successful in 8m9s
Vendors .gitignore, .dockerignore, .editorconfig, REPO_POLICIES.md,
script/cibuild, script/docker, script/lint, script/test and the CI
workflow byte-identical from sneak/prompts commit cc440118c876; the
two ignore files then add this repo's own build output, root-anchored
so cmd/bsdaily/ stays in. The Dockerfile gains a lint phase on the
golangci-lint image already pinned and a test phase on the Debian Go
image with sqlite3 and zstd; the build stage copies a file from each,
so a plain docker build cannot skip them. Nothing installs
golangci-lint on the host. On apt, script/bootstrap refreshes the
package lists once before its first install.

.golangci.yml and the lint cleanup: #6

Model: opus-5-5
2026-10-06 05:09:54 +00:00
22 changed files with 416 additions and 1077 deletions
-99
View File
@@ -1,99 +0,0 @@
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
View File
@@ -1,2 +0,0 @@
node_modules/
yarn.lock
-4
View File
@@ -1,4 +0,0 @@
{
"tabWidth": 4,
"proseWrap": "always"
}
+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.14.0 (Debian-based), 2026-10-06 # golangci/golangci-lint:v2.12.2-alpine, 2026-06-28
FROM golangci/golangci-lint@sha256:ad862ba6b3798cbe0fd9fd7408d498fd74fbd2623a92406b2fd3898faf0bf98f AS lint FROM golangci/golangci-lint:v2.12.2-alpine@sha256:91b27804074a0bacea298707f016911e60cf0cdbc6c7bf5ccacb5f0606d18d60 AS lint
WORKDIR /src WORKDIR /src
+122 -121
View File
@@ -1,26 +1,28 @@
# bsdaily # bsdaily
[bsdaily](https://git.eeqj.de/sneak/bsdaily) is a command-line utility written [bsdaily](https://git.eeqj.de/sneak/bsdaily) is a command-line utility
in [Go](https://golang.org) that carves a single day (or a range of days) of written in [Go](https://golang.org) that carves a single day (or a range of
[Bluesky](https://bsky.app) firehose data out of a large, continuously-growing days) of [Bluesky](https://bsky.app) firehose data out of a large,
SQLite database and writes it out as a self-contained, continuously-growing SQLite database and writes it out as a self-contained,
[zstd](https://facebook.github.io/zstd/)-compressed SQL dump. The dumps are [zstd](https://facebook.github.io/zstd/)-compressed SQL dump. The dumps are
named by date (e.g. `2026-06-27.sql.zst`), organized into per-month directories, named by date (e.g. `2026-06-27.sql.zst`), organized into per-month
and are designed to be published, archived, mirrored, and later re-merged back directories, and are designed to be published, archived, mirrored, and later
into a single database. re-merged back into a single database.
The source database is read from a read-only [ZFS](https://openzfs.org) The source database is read from a read-only [ZFS](https://openzfs.org)
snapshot, so extraction never contends with the live firehose ingester that is snapshot, so extraction never contends with the live firehose ingester that
writing to the original database. The tool is operationally conservative: it is writing to the original database. The tool is operationally
checks free disk space before starting, copies the snapshot to fast scratch conservative: it checks free disk space before starting, copies the snapshot
storage, processes one day at a time to avoid SQLite lock contention, verifies to fast scratch storage, processes one day at a time to avoid SQLite lock
every compressed output before publishing it, and writes output atomically via a contention, verifies every compressed output before publishing it, and
temp-file-and-rename so a partial run never leaves a corrupt `.sql.zst` behind. writes output atomically via a temp-file-and-rename so a partial run never
leaves a corrupt `.sql.zst` behind.
This project was written by [@sneak](https://sneak.berlin) to produce a daily, This project was written by [@sneak](https://sneak.berlin) to produce a
mergeable, publicly-mirrorable archive of the Bluesky firehose. It is currently daily, mergeable, publicly-mirrorable archive of the Bluesky firehose. It is
a one-person effort. The current version is pre-1.0 and there has not yet been a currently a one-person effort. The current version is pre-1.0 and there has
versioned release; [SemVer](https://semver.org) will be used for releases. not yet been a versioned release; [SemVer](https://semver.org) will be used
for releases.
# Build Status # Build Status
@@ -31,14 +33,14 @@ branch must always be green.
Primary development happens on a privately-run Gitea instance at Primary development happens on a privately-run Gitea instance at
[https://git.eeqj.de/sneak/bsdaily](https://git.eeqj.de/sneak/bsdaily) and [https://git.eeqj.de/sneak/bsdaily](https://git.eeqj.de/sneak/bsdaily) and
issues are [tracked there](https://git.eeqj.de/sneak/bsdaily/issues). 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 a standard `go fmt`, syntactically
Markdown), syntactically valid, and must pass the linting defined in the valid, and must pass the linting defined in the repository (presently the
repository (the shared `.golangci.yml` from `sneak/prompts`), which can be run `golangci-lint` defaults), which can be run with a `make lint`. The `main`
with a `make lint`. The `main` branch is protected and all changes must be made branch is protected and all changes must be made via [pull
via [pull requests](https://git.eeqj.de/sneak/bsdaily/pulls) and pass CI to be 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,
tooling requirements, and workflow conventions. tooling requirements, and workflow conventions.
@@ -48,55 +50,55 @@ tooling requirements, and workflow conventions.
This repository adheres to the This repository adheres to the
[Scripts to Rule Them All](https://github.com/github/scripts-to-rule-them-all) [Scripts to Rule Them All](https://github.com/github/scripts-to-rule-them-all)
standard: normalized scripts in `script/` are the entrypoints for the standard: normalized scripts in `script/` are the entrypoints for the
development workflow, and the Makefile targets are thin shims that call them. We development workflow, and the Makefile targets are thin shims that call
provide: them. We provide:
- `script/bootstrap` — install all development dependencies (go, Go module - `script/bootstrap` — install all development dependencies (go, Go
download, node and yarn at pinned versions when absent, and the prettier module download); the linter is not installed on the host
pinned in `package.json` and `yarn.lock`); a node or yarn already installed is
used whatever its version; the linter is not installed on the host
- `script/setup` — make a fresh clone ready for development: runs - `script/setup` — make a fresh clone ready for development: runs
`script/bootstrap`, then `script/install-precommit` `script/bootstrap`, then `script/install-precommit`
- `script/projectname` — print the project name (used for the Docker image tags) - `script/projectname` — print the project name (used for the Docker
- `script/test` — build the Dockerfile's `test` phase, which runs the test suite image tags)
with `-race` (verbose rerun on failure) - `script/test` — build the Dockerfile's `test` phase, which runs the
test suite with `-race` (verbose rerun on failure)
- `script/lint` — build the Dockerfile's `lint` phase, which runs - `script/lint` — build the Dockerfile's `lint` phase, which runs
`golangci-lint run ./...` `golangci-lint run ./...`
- `script/fmt` — format the Go code with `go fmt` and every Markdown file with - `script/fmt` — format the Go code with `go fmt` (writes)
prettier (writes) - `script/fmt-check` — check the Go formatting with `gofmt` (read-only)
- `script/fmt-check` — check the Go formatting with `gofmt` and the Markdown - `script/check` — run `script/test`, `script/lint`, and
formatting with prettier (read-only) `script/fmt-check`
- `script/check` — run `script/test`, `script/lint`, and `script/fmt-check` - `script/docker` — build the Docker image tagged via
- `script/docker` — build the Docker image tagged via `script/projectname`; the `script/projectname`; the build runs the `lint` and `test` phases
build runs the `lint` and `test` phases first first
- `script/cibuild` — CI entrypoint: runs `script/bootstrap` and `script/check`, - `script/cibuild` — CI entrypoint: runs `script/bootstrap` and
then builds the Docker image tagged via `script/projectname` `script/check`, then builds the Docker image tagged via
- `script/precommit` — pre-commit gate: `go mod tidy` (must not change `go.mod` `script/projectname`
or `go.sum`) and `go fmt`, then `script/check` - `script/precommit` — pre-commit gate: `go mod tidy` (must not change
- `script/install-precommit` — install the git pre-commit hook that runs `go.mod` or `go.sum`) and `go fmt`, then `script/check`
`script/precommit` - `script/install-precommit` — install the git pre-commit hook that
runs `script/precommit`
Every Docker build in `script/` is uncached, so the `lint` and `test` phases Every Docker build in `script/` is uncached, so the `lint` and `test`
always run rather than being served from the build cache. phases always run rather than being served from the build cache.
# Problem Statement # Problem Statement
A Bluesky firehose ingester writes every observed post (and associated users, A Bluesky firehose ingester writes every observed post (and associated
hashtags, URLs, and media references) into a single ever-growing SQLite users, hashtags, URLs, and media references) into a single ever-growing
database, `firehose.db`. This database has several properties that make it SQLite database, `firehose.db`. This database has several properties that
awkward to publish or archive directly: make it awkward to publish or archive directly:
- It is **large and always growing**, so re-publishing the whole thing every day - It is **large and always growing**, so re-publishing the whole thing every
is wasteful. day is wasteful.
- It is **continuously written**, so reading from it directly risks lock - It is **continuously written**, so reading from it directly risks lock
contention with the live ingester and inconsistent reads. contention with the live ingester and inconsistent reads.
- It is **monolithic**, so there is no natural unit at which to mirror, share, - It is **monolithic**, so there is no natural unit at which to mirror,
or distribute "just yesterday's posts". share, or distribute "just yesterday's posts".
What is wanted instead is a stable, immutable, per-day artifact: a small file What is wanted instead is a stable, immutable, per-day artifact: a small
containing exactly one calendar day of firehose data, cheap to publish, cheap to file containing exactly one calendar day of firehose data, cheap to publish,
mirror, and trivially re-mergeable into a full database by anyone who collects a cheap to mirror, and trivially re-mergeable into a full database by anyone
set of them. who collects a set of them.
# Proposed Solution # Proposed Solution
@@ -106,81 +108,81 @@ A tool, `bsdaily`, that:
filesystem, so it reads from a consistent point-in-time copy that the live filesystem, so it reads from a consistent point-in-time copy that the live
ingester cannot be writing to; ingester cannot be writing to;
- copies the snapshot's database files to fast scratch storage; - copies the snapshot's database files to fast scratch storage;
- **extracts** a single day's `posts` (and all rows reachable from them) into a - **extracts** a single day's `posts` (and all rows reachable from them) into
fresh, minimal per-day SQLite database; a fresh, minimal per-day SQLite database;
- **dumps** that per-day database to SQL and pipes it through multithreaded zstd - **dumps** that per-day database to SQL and pipes it through multithreaded
compression; zstd compression;
- **verifies** the compressed output (zstd integrity check plus a sanity check - **verifies** the compressed output (zstd integrity check plus a sanity
that the decompressed stream actually looks like SQL); check that the decompressed stream actually looks like SQL);
- **publishes** the result atomically as - **publishes** the result atomically as
`DailiesBase/YYYY-MM/YYYY-MM-DD.sql.zst`. `DailiesBase/YYYY-MM/YYYY-MM-DD.sql.zst`.
Each daily dump is emitted with `INSERT` statements over the full schema Each daily dump is emitted with `INSERT` statements over the full schema
(including the deduplicated `users`, `hashtags`, and `urls` lookup tables), so (including the deduplicated `users`, `hashtags`, and `urls` lookup tables),
any collection of daily dumps can be merged into a single database by rewriting so any collection of daily dumps can be merged into a single database by
`INSERT INTO` to `INSERT OR IGNORE INTO` and replaying them in sequence. Two rewriting `INSERT INTO` to `INSERT OR IGNORE INTO` and replaying them in
helper scripts ([`merge_daily_dumps.sh`](merge_daily_dumps.sh) and sequence. Two helper scripts ([`merge_daily_dumps.sh`](merge_daily_dumps.sh)
[`regenerate_auxiliary_tables.sql`](regenerate_auxiliary_tables.sql)) are and [`regenerate_auxiliary_tables.sql`](regenerate_auxiliary_tables.sql)) are
included to do exactly this and to rebuild the aggregate statistics included to do exactly this and to rebuild the aggregate statistics
(`use_count`, `first_seen`, user `resolved_at`/`updated_at`) afterward. (`use_count`, `first_seen`, user `resolved_at`/`updated_at`) afterward.
# Design Goals # Design Goals
- **Never disturb the live ingester.** All reads come from a ZFS snapshot, never - **Never disturb the live ingester.** All reads come from a ZFS snapshot,
the live database. never the live database.
- **Crash-safe, idempotent runs.** Output is written to a temp file and - **Crash-safe, idempotent runs.** Output is written to a temp file and
atomically renamed; a day whose final output already exists is skipped, so atomically renamed; a day whose final output already exists is skipped, so
re-running a range is safe and resumable. re-running a range is safe and resumable.
- **Mergeable output.** Daily dumps re-combine losslessly into a full database - **Mergeable output.** Daily dumps re-combine losslessly into a full
via `INSERT OR IGNORE`. database via `INSERT OR IGNORE`.
- **Operationally cautious.** Free-space preflight checks on both scratch and - **Operationally cautious.** Free-space preflight checks on both scratch and
output filesystems; explicit verification of every artifact before it is output filesystems; explicit verification of every artifact before it is
published. published.
- **Fast where it's free.** Large snapshot copies use a 256MiB buffer, - **Fast where it's free.** Large snapshot copies use a 256MiB buffer,
pre-allocate the destination, and (on Linux) issue `posix_fadvise` pre-allocate the destination, and (on Linux) issue `posix_fadvise`
sequential/willneed hints; extraction uses aggressive, crash-unsafe-by-design sequential/willneed hints; extraction uses aggressive,
SQLite pragmas because the working data lives only in disposable scratch crash-unsafe-by-design SQLite pragmas because the working data lives only
space. in disposable scratch space.
# Non-Goals # Non-Goals
- **Real-time export.** `bsdaily` operates on daily snapshots; the freshest day - **Real-time export.** `bsdaily` operates on daily snapshots; the freshest
it can produce is the snapshot date minus one. day it can produce is the snapshot date minus one.
- **Schema ownership.** The schema is defined by the upstream firehose ingester; - **Schema ownership.** The schema is defined by the upstream firehose
[`schema.sql`](schema.sql) is included for reference only. `bsdaily` copies ingester; [`schema.sql`](schema.sql) is included for reference only.
whatever table and index DDL it finds in the source. `bsdaily` copies whatever table and index DDL it finds in the source.
- **Cross-platform deployment.** It is built and run on Linux (the free-space - **Cross-platform deployment.** It is built and run on Linux (the
check and fadvise hints use `golang.org/x/sys/unix`; a non-Linux build free-space check and fadvise hints use `golang.org/x/sys/unix`; a non-Linux
compiles but is a no-op for the fadvise hints). The hard-coded paths assume build compiles but is a no-op for the fadvise hints). The hard-coded paths
the production host's ZFS layout. assume the production host's ZFS layout.
# How It Works # How It Works
A single run proceeds as follows: A single run proceeds as follows:
1. **Find the snapshot.** Scan `SnapshotBase` for directories matching 1. **Find the snapshot.** Scan `SnapshotBase` for directories matching
`zfs-auto-snap_daily-YYYY-MM-DD-NNNN`, pick the most recent, and confirm it `zfs-auto-snap_daily-YYYY-MM-DD-NNNN`, pick the most recent, and confirm
contains `firehose.db`. it contains `firehose.db`.
2. **Determine target days.** Default to the snapshot date minus one day; or use 2. **Determine target days.** Default to the snapshot date minus one day;
`--date`, or every day in the inclusive `--from`/`--to` range. or use `--date`, or every day in the inclusive `--from`/`--to` range.
3. **Preflight disk space.** Require at least 500GiB free on the scratch 3. **Preflight disk space.** Require at least 500GiB free on the scratch
filesystem and 20GiB free on the output filesystem. filesystem and 20GiB free on the output filesystem.
4. **Copy the database to scratch.** Copy `firehose.db`, its `-wal`, and (if 4. **Copy the database to scratch.** Copy `firehose.db`, its `-wal`, and (if
present) its `-shm` from the snapshot into a fresh temp directory under present) its `-shm` from the snapshot into a fresh temp directory under
`TmpBase`. `TmpBase`.
5. **Per day**, processed strictly one at a time to avoid SQLite contention: 5. **Per day**, processed strictly one at a time to avoid SQLite contention:
- skip the day if its final output file already exists; - skip the day if its final output file already exists;
- `ATTACH` the copied source DB to a new empty per-day DB, recreate the - `ATTACH` the copied source DB to a new empty per-day DB, recreate the
table DDL, and `INSERT ... SELECT` the target day's `posts` plus all rows table DDL, and `INSERT ... SELECT` the target day's `posts` plus all
reachable from them (`posts_hashtags`, `posts_urls`, `hashtags`, `urls`, rows reachable from them (`posts_hashtags`, `posts_urls`, `hashtags`,
`users`, and `media` if that table exists); `urls`, `users`, and `media` if that table exists);
- abort the day cleanly if there are zero posts (`ErrNoPosts`), rather than - abort the day cleanly if there are zero posts (`ErrNoPosts`), rather
emitting an empty dump; than emitting an empty dump;
- recreate indexes, detach the source, and verify the inserted row count; - recreate indexes, detach the source, and verify the inserted row count;
- `sqlite3 .dump | zstdmt` into a hidden temp file; - `sqlite3 .dump | zstdmt` into a hidden temp file;
- run a zstd integrity check and confirm the decompressed head looks like - run a zstd integrity check and confirm the decompressed head looks like
SQL; SQL;
- atomically rename into place and delete the per-day scratch DB. - atomically rename into place and delete the per-day scratch DB.
6. **Clean up** the temp directory and log a processed/skipped/total summary. 6. **Clean up** the temp directory and log a processed/skipped/total summary.
# Usage # Usage
@@ -217,9 +219,9 @@ sqlite3 merged.db < regenerate_auxiliary_tables.sql
- **Go** (see [`go.mod`](go.mod) for the toolchain version) to build. - **Go** (see [`go.mod`](go.mod) for the toolchain version) to build.
- **Linux** for production use (ZFS snapshots, `statfs` free-space checks, - **Linux** for production use (ZFS snapshots, `statfs` free-space checks,
`posix_fadvise` hints). `posix_fadvise` hints).
- The **`sqlite3`** and **`zstdmt`** (multithreaded zstd) binaries on `PATH`; - The **`sqlite3`** and **`zstdmt`** (multithreaded zstd) binaries on
`zstdcat` is used for verification. SQLite reads/writes during extraction use `PATH`; `zstdcat` is used for verification. SQLite reads/writes during
the pure-Go [`modernc.org/sqlite`](https://pkg.go.dev/modernc.org/sqlite) extraction use the pure-Go [`modernc.org/sqlite`](https://pkg.go.dev/modernc.org/sqlite)
driver, so no cgo is required for that part. driver, so no cgo is required for that part.
# Configuration # Configuration
@@ -239,21 +241,21 @@ elsewhere.
# Data Model # Data Model
The firehose schema (reference copy in [`schema.sql`](schema.sql)) centers on a The firehose schema (reference copy in [`schema.sql`](schema.sql)) centers
`posts` table, with `users` keyed by DID and many-to-many junction tables on a `posts` table, with `users` keyed by DID and many-to-many junction
linking posts to deduplicated `hashtags` and `urls`. An optional `media` table tables linking posts to deduplicated `hashtags` and `urls`. An optional
tracks downloaded blobs by content hash. `bsdaily` does not own this schema; it `media` table tracks downloaded blobs by content hash. `bsdaily` does not
reflects whatever DDL exists in the source snapshot and selects forward from own this schema; it reflects whatever DDL exists in the source snapshot and
`posts` along the foreign-key relationships to produce a referentially-complete selects forward from `posts` along the foreign-key relationships to produce a
per-day slice. referentially-complete per-day slice.
# Use Cases # Use Cases
## Daily public archive ## Daily public archive
Publish one small, immutable file per day to static HTTP (or IPFS, or a mirror Publish one small, immutable file per day to static HTTP (or IPFS, or a
network) so that anyone can fetch exactly the day(s) they want and re-merge them mirror network) so that anyone can fetch exactly the day(s) they want and
locally. re-merge them locally.
## Backfilling a range ## Backfilling a range
@@ -263,17 +265,16 @@ safe to re-run.
## Reconstituting a full database ## Reconstituting a full database
Collect any set of daily dumps and merge them with `INSERT OR IGNORE` to rebuild Collect any set of daily dumps and merge them with `INSERT OR IGNORE` to
a complete, queryable SQLite database, then regenerate the aggregate statistics rebuild a complete, queryable SQLite database, then regenerate the aggregate
tables. statistics tables.
# See Also # See Also
## Links ## Links
- Repo: [https://git.eeqj.de/sneak/bsdaily](https://git.eeqj.de/sneak/bsdaily) - Repo: [https://git.eeqj.de/sneak/bsdaily](https://git.eeqj.de/sneak/bsdaily)
- Issues: - Issues: [https://git.eeqj.de/sneak/bsdaily/issues](https://git.eeqj.de/sneak/bsdaily/issues)
[https://git.eeqj.de/sneak/bsdaily/issues](https://git.eeqj.de/sneak/bsdaily/issues)
- Bluesky: [https://bsky.app](https://bsky.app) - Bluesky: [https://bsky.app](https://bsky.app)
- zstd: [https://facebook.github.io/zstd/](https://facebook.github.io/zstd/) - zstd: [https://facebook.github.io/zstd/](https://facebook.github.io/zstd/)
+28 -31
View File
@@ -1,12 +1,12 @@
# Workflow # Workflow
- branch (from `main`) * branch (from `main`)
- do the work in Next Step * do the work in Next Step
- move Next Step to the top of Completed Steps * move Next Step to the top of Completed Steps
- move the top item of Future Steps into Next Step * move the top item of Future Steps into Next Step
- commit (`TODO.md` changes in the same commit as the work) * commit (`TODO.md` changes in the same commit as the work)
- merge to `main` if the branch is not protected, otherwise open a PR * merge to `main` if the branch is not protected, otherwise open a PR
- push * push
# Status # Status
@@ -14,38 +14,35 @@ pre-1.0
# Next Step # Next Step
Expand tests beyond the compilation smoke test: unit tests for the extraction, Add the canonical `.golangci.yml`, move the lint phase to golangci-lint
verification, and atomic-publish paths. v2.14.0 in the same commit, and fix the findings it surfaces
(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
`script/fmt-check` checks every Markdown file; prettier pinned in
`package.json` and `yarn.lock`, installed by `script/bootstrap`, which
installs node and yarn at pinned versions when absent and otherwise uses the
ones already installed; existing Markdown reformatted.
- 2026-10-05: Brought the repo up to the standard layout: canonical - 2026-10-05: Brought the repo up to the standard layout: canonical
`.gitignore`, `.dockerignore` and `.editorconfig`; `lint` and `test` phases in `.gitignore`, `.dockerignore` and `.editorconfig`; `lint` and `test`
the `Dockerfile`, built by `script/lint` and `script/test`; canonical phases in the `Dockerfile`, built by `script/lint` and `script/test`;
`script/cibuild`, `script/docker` and CI workflow; no linter installed on the canonical `script/cibuild`, `script/docker` and CI workflow; no linter
host; re-vendored `REPO_POLICIES.md`. installed on the host; re-vendored `REPO_POLICIES.md`.
- 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints, Makefile - 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints,
shims, README Entrypoints section Makefile shims, README Entrypoints section
- 2026-06-28: Fixed errcheck lint failures; added compilation smoke test; tidied - 2026-06-28: Fixed errcheck lint failures; added compilation smoke test;
go.mod. tidied go.mod.
- 2026-06-28: Added repo scaffolding: README, LICENSE, Makefile, Dockerfile, - 2026-06-28: Added repo scaffolding: README, LICENSE, Makefile,
REPO_POLICIES.md, and Gitea CI. Dockerfile, REPO_POLICIES.md, and Gitea CI.
- 2026-02-12: Fixed SQLite database locking by removing parallel processing; - 2026-02-12: Fixed SQLite database locking by removing parallel
fixed Linux build via golang.org/x/sys/unix Fadvise. processing; fixed Linux build via golang.org/x/sys/unix Fadvise.
- 2026-02-12: Optimized file copy for large databases; moved temp directory to - 2026-02-12: Optimized file copy for large databases; moved temp
NVMe scratch storage. directory to NVMe scratch storage.
- 2026-02-11: Added date range support. - 2026-02-11: Added date range support.
- 2026-02-09: Initial implementation: single-day extraction, specific-date - 2026-02-09: Initial implementation: single-day extraction, specific-date
targeting, faster pruning of throwaway database copies. targeting, faster pruning of throwaway database copies.
# Future Steps # Future Steps
- Format Markdown with prettier in `script/fmt` and `script/fmt-check`
(https://git.eeqj.de/sneak/bsdaily/issues/7).
- 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.
+48 -88
View File
@@ -1,9 +1,6 @@
// 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"
@@ -13,112 +10,75 @@ 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, fromFlag, toFlag string var dateFlag 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(_ *cobra.Command, _ []string) error { RunE: func(cmd *cobra.Command, args []string) error {
targetDates, err := parseTargetDates(dateFlag, fromFlag, toFlag) hasDate := dateFlag != ""
if err != nil { hasFrom := fromFlag != ""
return err 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")
} }
err = bsdaily.Run(targetDates) var targetDates []time.Time
if err != nil {
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
} }
slog.Info("completed successfully") slog.Info("completed successfully")
return nil return nil
}, },
} }
rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "", rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "", "target date to extract (YYYY-MM-DD); defaults to snapshot date minus one day")
"target date to extract (YYYY-MM-DD); "+ rootCmd.Flags().StringVar(&fromFlag, "from", "", "start of date range to extract (YYYY-MM-DD, inclusive); use with --to")
"defaults to snapshot date minus one day") rootCmd.Flags().StringVar(&toFlag, "to", "", "end of date range to extract (YYYY-MM-DD, inclusive); use with --from")
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")
err := rootCmd.Execute() if err := rootCmd.Execute(); err != nil {
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
}
+6 -17
View File
@@ -1,32 +1,21 @@
package bsdaily_test package bsdaily
import ( import "testing"
"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) {
t.Parallel() if DBFilename == "" || WALFilename == "" || SHMFilename == "" {
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")
} }
} }
+1 -6
View File
@@ -1,11 +1,7 @@
// 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"
@@ -31,5 +27,4 @@ const (
verificationHeadLines = 20 verificationHeadLines = 20
) )
var snapshotPattern = regexp.MustCompile( var snapshotPattern = regexp.MustCompile(`^zfs-auto-snap_daily-(\d{4}-\d{2}-\d{2})-\d{4}$`)
`^zfs-auto-snap_daily-(\d{4}-\d{2}-\d{2})-\d{4}$`)
+8 -28
View File
@@ -1,7 +1,6 @@
package bsdaily package bsdaily
import ( import (
"errors"
"fmt" "fmt"
"io" "io"
"log/slog" "log/slog"
@@ -10,30 +9,20 @@ import (
) )
const ( const (
// 256MB buffer for large file copies from fast storage copyBufferSize = 256 * 1024 * 1024 // 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() {
cerr := srcFile.Close() if cerr := srcFile.Close(); cerr != nil {
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)
} }
}() }()
@@ -48,49 +37,40 @@ 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() {
cerr := dstFile.Close() if cerr := dstFile.Close(); cerr != nil && err == nil {
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
truncErr := dstFile.Truncate(srcInfo.Size()) if err := dstFile.Truncate(srcInfo.Size()); err != nil {
if truncErr != nil { slog.Warn("failed to pre-allocate destination file", "error", err)
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("%w: wrote %d bytes, expected %d", return fmt.Errorf("short copy: wrote %d bytes, expected %d", written, srcInfo.Size())
errShortCopy, written, srcInfo.Size())
} }
err = dstFile.Sync() if err := dstFile.Sync(); err != nil {
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() / oneMB throughputMBps := float64(written) / elapsed.Seconds() / (1024 * 1024)
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
} }
+4 -3
View File
@@ -9,7 +9,8 @@ import (
func applyFileAdvice(file *os.File, size int64) { func applyFileAdvice(file *os.File, size int64) {
fd := int(file.Fd()) fd := int(file.Fd())
_ = unix.Fadvise(fd, 0, size, unix.FADV_SEQUENTIAL) // POSIX_FADV_SEQUENTIAL = 2
// Prefetch the file into the page cache _ = unix.Fadvise(fd, 0, size, 2)
_ = unix.Fadvise(fd, 0, size, unix.FADV_WILLNEED) // POSIX_FADV_WILLNEED = 3 - prefetch file into cache
_ = unix.Fadvise(fd, 0, size, 3)
} }
+3 -16
View File
@@ -1,39 +1,26 @@
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("%w on %s (%s): %.1f GB free, need %.1f GB", return fmt.Errorf("insufficient disk space on %s (%s): %.1f GB free, need %.1f GB",
errInsufficientSpace, path, label, freeGB, minGB) path, label, freeGB, minGB)
} }
return nil return nil
} }
+35 -74
View File
@@ -1,8 +1,6 @@
package bsdaily package bsdaily
import ( import (
"context"
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"os" "os"
@@ -11,44 +9,61 @@ 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"} {
_, err = exec.LookPath(tool) if _, err := exec.LookPath(tool); err != nil {
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)
} }
} }
err = CheckFreeSpace(filepath.Dir(outputPath), MinDailiesFreeBytes, if err := CheckFreeSpace(filepath.Dir(outputPath), MinDailiesFreeBytes, "dailiesBase (pre-dump)"); err != nil {
"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() {
cerr := outFile.Close() if cerr := outFile.Close(); cerr != nil && err == nil {
if cerr != nil && err == nil {
err = fmt.Errorf("closing output: %w", cerr) err = fmt.Errorf("closing output: %w", cerr)
} }
}() }()
err = runDumpPipeline(context.Background(), dbPath, outFile) // Dump all tables but use INSERT OR IGNORE for mergeable imports
// This preserves all data while allowing multiple dumps to be merged
// Users should import with: zstdcat *.sql.zst | sed 's/INSERT INTO/INSERT OR IGNORE INTO/g' | sqlite3 merged.db
dumpCmd := exec.Command("sqlite3", dbPath, ".dump")
zstdCmd := exec.Command("zstdmt", fmt.Sprintf("-%d", zstdCompressionLevel))
pipe, err := dumpCmd.StdoutPipe()
if err != nil { if err != nil {
return err return fmt.Errorf("creating dump stdout pipe: %w", err)
}
zstdCmd.Stdin = pipe
zstdCmd.Stdout = outFile
var dumpStderr, zstdStderr strings.Builder
dumpCmd.Stderr = &dumpStderr
zstdCmd.Stderr = &zstdStderr
slog.Info("starting sqlite3 dump and zstdmt compression")
if err := zstdCmd.Start(); err != nil {
return fmt.Errorf("starting zstdmt: %w", err)
}
if err := dumpCmd.Start(); err != nil {
return fmt.Errorf("starting sqlite3 dump: %w", err)
} }
err = outFile.Sync() if err := dumpCmd.Wait(); err != nil {
if err != nil { return fmt.Errorf("sqlite3 dump failed: %w; stderr: %s", err, dumpStderr.String())
}
if err := zstdCmd.Wait(); err != nil {
return fmt.Errorf("zstdmt failed: %w; stderr: %s", err, zstdStderr.String())
}
if err := outFile.Sync(); err != nil {
return fmt.Errorf("syncing output: %w", err) return fmt.Errorf("syncing output: %w", err)
} }
@@ -56,66 +71,12 @@ 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 errEmptyOutput return fmt.Errorf("compressed output is empty")
}
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
+52 -179
View File
@@ -1,29 +1,22 @@
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"
@@ -31,282 +24,162 @@ 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)"+ 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)
"&_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() {
cerr := db.Close() if cerr := db.Close(); cerr != nil {
if cerr != nil { slog.Warn("failed to close destination database", "path", dstDBPath, "error", cerr)
slog.Warn("failed to close destination database",
"path", dstDBPath, "error", cerr)
} }
}() }()
// Attach source database // Attach source database
_, err = db.ExecContext(ctx, "ATTACH DATABASE ? AS src", srcDBPath) if _, err := db.Exec("ATTACH DATABASE ? AS src", srcDBPath); err != nil {
if err != nil {
return fmt.Errorf("attaching source database: %w", err) return fmt.Errorf("attaching source database: %w", err)
} }
err = createTables(ctx, db)
if err != nil {
return err
}
postCount, err := insertDayRows(ctx, db, targetDay, dayStart, dayEnd)
if err != nil {
return err
}
err = createIndexes(ctx, db)
if err != nil {
return err
}
// Detach source
_, err = db.ExecContext(ctx, "DETACH DATABASE src")
if err != nil {
return fmt.Errorf("detaching source database: %w", err)
}
// Verify post count
var verifyCount int64
err = db.QueryRowContext(ctx, "SELECT COUNT(*) FROM posts").Scan(&verifyCount)
if err != nil {
return fmt.Errorf("verifying post count: %w", err)
}
if verifyCount != postCount {
return fmt.Errorf("%w: inserted %d but found %d",
errPostCountMismatch, postCount, verifyCount)
}
slog.Info("extraction complete", "posts", verifyCount)
return nil
}
// createTables creates every table of the attached source database, empty,
// in the destination database.
func createTables(ctx context.Context, db *sql.DB) error {
// Copy table DDL from source // Copy table DDL from source
slog.Info("copying table DDL from source") slog.Info("copying table DDL from source")
rows, err := db.Query("SELECT sql FROM src.sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%' ORDER BY name")
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 { if err != nil {
return fmt.Errorf("reading source schema: %w", err) return fmt.Errorf("reading source schema: %w", err)
} }
defer func() { defer func() {
cerr := rows.Close() if cerr := rows.Close(); cerr != nil {
if cerr != nil {
slog.Warn("failed to close schema rows", "error", cerr) slog.Warn("failed to close schema rows", "error", cerr)
} }
}() }()
var ddlStatements []string var ddlStatements []string
for rows.Next() { for rows.Next() {
var ddl string var ddl string
if err := rows.Scan(&ddl); err != nil {
err = rows.Scan(&ddl)
if err != nil {
return fmt.Errorf("scanning DDL: %w", err) return fmt.Errorf("scanning DDL: %w", err)
} }
ddlStatements = append(ddlStatements, ddl) ddlStatements = append(ddlStatements, ddl)
} }
if err := rows.Err(); err != nil {
err = rows.Err()
if err != nil {
return fmt.Errorf("iterating DDL rows: %w", err) return fmt.Errorf("iterating DDL rows: %w", err)
} }
for _, ddl := range ddlStatements { for _, ddl := range ddlStatements {
_, err = db.ExecContext(ctx, ddl) if _, err := db.Exec(ddl); err != nil {
if err != nil {
return fmt.Errorf("creating table: %w\nDDL: %s", err, ddl) 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 // Begin transaction for bulk inserts
tx, err := db.BeginTx(ctx, nil) tx, err := db.Begin()
if err != nil { if err != nil {
return 0, fmt.Errorf("beginning transaction: %w", err) return fmt.Errorf("beginning transaction: %w", err)
} }
// After a successful Commit this does nothing and returns ErrTxDone.
defer func() { defer func() {
rerr := tx.Rollback() if err != nil {
if rerr != nil && !errors.Is(rerr, sql.ErrTxDone) { if rerr := tx.Rollback(); rerr != nil && !errors.Is(rerr, sql.ErrTxDone) {
slog.Warn("failed to roll back transaction", "error", rerr) slog.Warn("failed to roll back transaction", "error", rerr)
}
} }
}() }()
// Insert target day's data // Insert target day's data
slog.Info("inserting posts for target day") slog.Info("inserting posts for target day")
result, err := tx.Exec("INSERT INTO posts SELECT * FROM src.posts WHERE timestamp >= ? AND timestamp < ?", dayStart, dayEnd)
//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 { if err != nil {
return 0, fmt.Errorf("inserting posts: %w", err) return fmt.Errorf("inserting posts: %w", err)
} }
postCount, _ := result.RowsAffected() postCount, _ := result.RowsAffected()
slog.Info("inserted posts", "count", postCount) slog.Info("inserted posts", "count", postCount)
if postCount == 0 { if postCount == 0 {
return 0, fmt.Errorf("%w %s - aborting to avoid producing empty output", return fmt.Errorf("%w %s - aborting to avoid producing empty output",
ErrNoPosts, targetDay.Format("2006-01-02")) ErrNoPosts, targetDay.Format("2006-01-02"))
} }
slog.Info("inserting junction and lookup tables") 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 {
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) return fmt.Errorf("inserting posts_hashtags: %w", err)
} }
_, err = tx.ExecContext(ctx, "INSERT INTO posts_urls "+ if _, err := tx.Exec("INSERT INTO posts_urls SELECT * FROM src.posts_urls WHERE post_id IN (SELECT id FROM posts)"); err != nil {
"SELECT * FROM src.posts_urls WHERE post_id IN (SELECT id FROM posts)")
if err != nil {
return fmt.Errorf("inserting posts_urls: %w", err) return fmt.Errorf("inserting posts_urls: %w", err)
} }
_, err = tx.ExecContext(ctx, "INSERT INTO hashtags "+ if _, err := tx.Exec("INSERT INTO hashtags SELECT * FROM src.hashtags WHERE id IN (SELECT hashtag_id FROM posts_hashtags)"); err != nil {
"SELECT * FROM src.hashtags "+
"WHERE id IN (SELECT hashtag_id FROM posts_hashtags)")
if err != nil {
return fmt.Errorf("inserting hashtags: %w", err) return fmt.Errorf("inserting hashtags: %w", err)
} }
_, err = tx.ExecContext(ctx, "INSERT INTO urls "+ if _, err := tx.Exec("INSERT INTO urls SELECT * FROM src.urls WHERE id IN (SELECT url_id FROM posts_urls)"); err != nil {
"SELECT * FROM src.urls WHERE id IN (SELECT url_id FROM posts_urls)")
if err != nil {
return fmt.Errorf("inserting urls: %w", err) return fmt.Errorf("inserting urls: %w", err)
} }
_, err = tx.ExecContext(ctx, "INSERT INTO users "+ if _, err := tx.Exec("INSERT INTO users SELECT * FROM src.users WHERE did IN (SELECT user_did FROM posts)"); err != nil {
"SELECT * FROM src.users WHERE did IN (SELECT user_did FROM posts)")
if err != nil {
return fmt.Errorf("inserting users: %w", err) 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 // Check if media table exists in source and copy if present
var mediaTableExists int var mediaTableExists int
if err := tx.QueryRow("SELECT COUNT(*) FROM src.sqlite_master WHERE type='table' AND name='media'").Scan(&mediaTableExists); err != nil {
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) slog.Warn("checking for media table", "error", err)
} else if mediaTableExists > 0 { } else if mediaTableExists > 0 {
slog.Info("inserting media entries") slog.Info("inserting media entries")
// Get post blob_cids for this day's posts // Get post blob_cids for this day's posts
//nolint:unqueryvet // copies whole rows; the source defines the columns 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 {
_, err = tx.ExecContext(ctx, "INSERT INTO media SELECT * FROM src.media "+ slog.Warn("inserting media (may not have matching entries)", "error", err)
"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 // Commit the transaction before any further database operations
// the destination database. Doing this after the bulk insert is faster if err := tx.Commit(); err != nil {
// than inserting into indexed tables. return fmt.Errorf("committing transaction: %w", err)
func createIndexes(ctx context.Context, db *sql.DB) error { }
tx = nil // Clear tx to ensure defer doesn't try to rollback
// Create indexes after bulk insert for speed // Create indexes after bulk insert for speed
slog.Info("creating indexes") 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")
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 { if err != nil {
return fmt.Errorf("reading source indexes: %w", err) return fmt.Errorf("reading source indexes: %w", err)
} }
defer func() { defer func() {
cerr := idxRows.Close() if cerr := idxRows.Close(); cerr != nil {
if cerr != nil {
slog.Warn("failed to close index rows", "error", cerr) slog.Warn("failed to close index rows", "error", cerr)
} }
}() }()
var idxStatements []string var idxStatements []string
for idxRows.Next() { for idxRows.Next() {
var idxSQL string var idxSQL string
if err := idxRows.Scan(&idxSQL); err != nil {
err = idxRows.Scan(&idxSQL)
if err != nil {
return fmt.Errorf("scanning index DDL: %w", err) return fmt.Errorf("scanning index DDL: %w", err)
} }
idxStatements = append(idxStatements, idxSQL) idxStatements = append(idxStatements, idxSQL)
} }
if err := idxRows.Err(); err != nil {
err = idxRows.Err()
if err != nil {
return fmt.Errorf("iterating index rows: %w", err) return fmt.Errorf("iterating index rows: %w", err)
} }
for _, idxSQL := range idxStatements { for _, idxSQL := range idxStatements {
_, err = db.ExecContext(ctx, idxSQL) if _, err := db.Exec(idxSQL); err != nil {
if err != nil {
return fmt.Errorf("creating index: %w\nDDL: %s", err, idxSQL) return fmt.Errorf("creating index: %w\nDDL: %s", err, idxSQL)
} }
} }
// Detach source
if _, err := db.Exec("DETACH DATABASE src"); err != nil {
return fmt.Errorf("detaching source database: %w", err)
}
// Verify post count
var verifyCount int64
if err := db.QueryRow("SELECT COUNT(*) FROM posts").Scan(&verifyCount); err != nil {
return fmt.Errorf("verifying post count: %w", err)
}
if verifyCount != postCount {
return fmt.Errorf("post count mismatch: inserted %d but found %d", postCount, verifyCount)
}
slog.Info("extraction complete", "posts", verifyCount)
return nil return nil
} }
+78 -165
View File
@@ -9,29 +9,21 @@ 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) {
err := os.Remove(path) if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
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)}
@@ -42,13 +34,10 @@ 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
err = CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase") if err := CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase"); err != nil {
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
} }
@@ -57,53 +46,14 @@ 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 {
rerr := os.RemoveAll(tmpDir) slog.Error("failed to remove temp directory", "path", tmpDir, "error", err)
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)
@@ -115,137 +65,100 @@ func copySnapshotFiles(snapshotDir, tmpDir string) (string, 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("%w: %s", errEmptySource, f) return fmt.Errorf("source file is empty: %s", f)
} }
slog.Info("source file", "path", f, "size_bytes", info.Size()) slog.Info("source file", "path", f, "size_bytes", info.Size())
} }
err := CopyFile(srcDB, dstDB) if err := CopyFile(srcDB, dstDB); err != nil {
if err != nil { return fmt.Errorf("copying database: %w", err)
return "", fmt.Errorf("copying database: %w", err)
} }
if err := CopyFile(srcWAL, dstWAL); err != nil {
err = CopyFile(srcWAL, dstWAL) return fmt.Errorf("copying WAL: %w", err)
if err != nil {
return "", fmt.Errorf("copying WAL: %w", err)
} }
if _, err := os.Stat(srcSHM); err == nil {
_, err = os.Stat(srcSHM) if err := CopyFile(srcSHM, dstSHM); err != nil {
if err == nil { return fmt.Errorf("copying SHM: %w", err)
err = CopyFile(srcSHM, dstSHM)
if err != nil {
return "", fmt.Errorf("copying SHM: %w", err)
} }
} }
return dstDB, nil // 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
// processDay extracts one day from the copied database at dstDB and for _, targetDay := range targetDates {
// publishes its compressed dump. It returns false, with no error, for a dayStr := targetDay.Format("2006-01-02")
// day it skips: one whose output already exists or that has no posts. slog.Info("processing day", "date", dayStr)
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")
if _, err := os.Stat(outputFinal); err == nil {
slog.Info("output already exists, skipping", "path", outputFinal)
skipped++
continue
}
// Check if output already exists // Extract target day into a per-day database
outputDir := filepath.Join(DailiesBase, targetDay.Format("2006-01")) extractedDB := filepath.Join(tmpDir, "extracted-"+dayStr+".db")
outputFinal := filepath.Join(outputDir, dayStr+".sql.zst") 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)
}
_, err := os.Stat(outputFinal) // Dump to SQL and compress
if err == nil { if err := os.MkdirAll(outputDir, 0755); 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) cleanup(extractedDB)
return fmt.Errorf("creating output directory %s: %w", outputDir, err)
return false, nil
} }
return false, fmt.Errorf("extracting day %s: %w", dayStr, err) outputTmp := filepath.Join(outputDir, "."+dayStr+".sql.zst.tmp")
}
// Dump to SQL and compress slog.Info("dumping and compressing", "tmp_output", outputTmp)
//nolint:gosec,mnd // world-readable on purpose: the dailies tree is published if err := DumpAndCompress(extractedDB, outputTmp); err != nil {
err = os.MkdirAll(outputDir, 0755) cleanup(outputTmp)
if err != nil { 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 {
cleanup(extractedDB)
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) cleanup(extractedDB)
processed++
return false, fmt.Errorf("creating output directory %s: %w", outputDir, err)
} }
outputTmp := filepath.Join(outputDir, "."+dayStr+".sql.zst.tmp") slog.Info("run summary", "processed", processed, "skipped", skipped, "total", len(targetDates))
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
} }
+7 -23
View File
@@ -1,7 +1,6 @@
package bsdaily package bsdaily
import ( import (
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"os" "os"
@@ -10,16 +9,10 @@ import (
"time" "time"
) )
var errNoSnapshots = errors.New("no daily snapshots found") func FindLatestDailySnapshot() (dir string, snapshotDate time.Time, err error) {
// 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", return "", time.Time{}, fmt.Errorf("reading snapshot directory %s: %w", SnapshotBase, err)
SnapshotBase, err)
} }
type snapshot struct { type snapshot struct {
@@ -28,30 +21,24 @@ func FindLatestDailySnapshot() (string, time.Time, 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", slog.Warn("skipping snapshot with unparseable date", "name", e.Name(), "error", err)
"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("%w in %s", errNoSnapshots, SnapshotBase) return "", time.Time{}, fmt.Errorf("no daily snapshots found in %s", SnapshotBase)
} }
sort.Slice(snapshots, func(i, j int) bool { sort.Slice(snapshots, func(i, j int) bool {
@@ -59,14 +46,11 @@ func FindLatestDailySnapshot() (string, time.Time, 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 {
_, err = os.Stat(dbPath) return "", time.Time{}, fmt.Errorf("database not found in snapshot %s: %w", dir, err)
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
+19 -81
View File
@@ -1,7 +1,6 @@
package bsdaily package bsdaily
import ( import (
"context"
"errors" "errors"
"fmt" "fmt"
"log/slog" "log/slog"
@@ -10,11 +9,6 @@ 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.
@@ -22,126 +16,70 @@ 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 {
err := testCmd.Run() return fmt.Errorf("zstdmt --test failed: %w; stderr: %s", err, testStderr.String())
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)
content, err := readDecompressedHead(ctx, path) headCmd := exec.Command("head", fmt.Sprintf("-%d", verificationHeadLines))
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
err = catCmd.Start() if err := catCmd.Start(); err != nil {
if err != nil { return fmt.Errorf("starting zstdcat: %w", err)
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)
err = headCmd.Wait() if err := headCmd.Wait(); err != nil {
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
return headOut.String(), nil content := headOut.String()
}
// 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 errEmptyDecompressed return fmt.Errorf("decompressed content is empty")
} }
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("%w; first %d bytes: %s", errNotSQL, return fmt.Errorf("decompressed content does not look like SQL; first %d bytes: %s",
verificationSampleBytes, verificationSampleBytes, content[:min(verificationSampleBytes, len(content))])
content[:min(verificationSampleBytes, len(content))])
} }
slog.Info("SQL content verification passed")
return nil return nil
} }
-5
View File
@@ -1,5 +0,0 @@
{
"devDependencies": {
"prettier": "3.8.1"
}
}
+1 -79
View File
@@ -3,20 +3,11 @@
# this repo. Idempotent: every install is guarded by a check so already # this repo. Idempotent: every install is guarded by a check so already
# installed tools are skipped. Base tooling comes from nix, apt, brew, # installed tools are skipped. Base tooling comes from nix, apt, brew,
# or apk (detected in that order); assumes NOTHING is present (not git, # or apk (detected in that order); assumes NOTHING is present (not git,
# make, go, or node). Node is used directly if installed; otherwise it # make, or go).
# is installed at a pinned version via nvm (installing nvm itself first,
# from a hash-verified release archive, never curl | sh).
set -eu set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)" ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
# Pinned versions, 2026-07-06
NODE_VERSION="22.17.0"
NVM_VERSION="0.40.3"
# sha256 of https://github.com/nvm-sh/nvm/archive/refs/tags/v0.40.3.tar.gz
NVM_SHA256="5f4d6aaa04a177dc93c985e31dbc411ab6b8c6e1e21d8015dbc1372625fcd1d0"
YARN_VERSION="1.22.22"
PKGMGR="" PKGMGR=""
SUDO="" SUDO=""
APT_UPDATED="" APT_UPDATED=""
@@ -65,69 +56,6 @@ missing() {
! command -v "$1" >/dev/null 2>&1 ! command -v "$1" >/dev/null 2>&1
} }
# verify_sha256 <file> <expected-hash>
verify_sha256() {
if command -v sha256sum >/dev/null 2>&1; then
actual="$(sha256sum "$1" | cut -d' ' -f1)"
else
actual="$(shasum -a 256 "$1" | cut -d' ' -f1)"
fi
if [ "$actual" != "$2" ]; then
echo "bootstrap: sha256 mismatch for $1" >&2
echo " expected: $2" >&2
echo " actual: $actual" >&2
exit 1
fi
}
# nvm is a bash script; run a command in a bash with nvm loaded
nvm_sh() {
bash -c ". \"\$HOME/.nvm/nvm.sh\" && $*"
}
ensure_nvm() {
[ -s "$HOME/.nvm/nvm.sh" ] && return 0
# nvm prerequisites; nvm itself requires bash
if missing bash; then pkg_install bash bash bash bash; fi
if missing curl; then pkg_install curl curl curl curl; fi
if missing git; then pkg_install git git git git; fi
tmp="$(mktemp -d)"
curl -fsSL -o "$tmp/nvm.tar.gz" \
"https://github.com/nvm-sh/nvm/archive/refs/tags/v${NVM_VERSION}.tar.gz"
verify_sha256 "$tmp/nvm.tar.gz" "$NVM_SHA256"
mkdir -p "$HOME/.nvm"
tar -xzf "$tmp/nvm.tar.gz" -C "$HOME/.nvm" --strip-components=1
rm -rf "$tmp"
}
ensure_node() {
if ! missing node; then return 0; fi
ensure_nvm
nvm_sh "nvm install $NODE_VERSION"
}
ensure_yarn() {
if ! missing yarn; then return 0; fi
if ! missing corepack; then
corepack enable
corepack prepare "yarn@$YARN_VERSION" --activate
elif [ -s "$HOME/.nvm/nvm.sh" ]; then
nvm_sh "nvm use $NODE_VERSION >/dev/null && corepack enable && \
corepack prepare yarn@$YARN_VERSION --activate"
else
npm install -g "yarn@$YARN_VERSION"
fi
}
install_js_deps() {
if missing yarn && [ -s "$HOME/.nvm/nvm.sh" ]; then
nvm_sh "nvm use $NODE_VERSION >/dev/null && cd \"$ROOT\" && \
yarn install --frozen-lockfile"
else
yarn install --frozen-lockfile
fi
}
main() { main() {
cd "$ROOT" cd "$ROOT"
@@ -140,12 +68,6 @@ main() {
go mod download go mod download
# Node and yarn, then the prettier pinned in package.json and
# yarn.lock, which script/fmt and script/fmt-check run
ensure_node
ensure_yarn
install_js_deps
echo "bootstrap complete" echo "bootstrap complete"
} }
+1 -23
View File
@@ -1,34 +1,12 @@
#!/bin/sh #!/bin/sh
# script/fmt: format all files (writes): the Go code with go fmt and # script/fmt: format all files (writes).
# every Markdown file with prettier.
set -eu set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)" ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
# Must match the pin in script/bootstrap.
NODE_VERSION="22.17.0"
# script/bootstrap installs node and yarn under nvm and leaves neither
# on the PATH of the shell that called it, so resolve the pinned
# toolchain here the way bootstrap's own install step does. nvm is a
# bash script, hence the subshell.
run_yarn() {
if command -v yarn >/dev/null 2>&1; then
exec yarn "$@"
fi
if [ ! -s "$HOME/.nvm/nvm.sh" ]; then
echo "fmt: no yarn; run script/bootstrap first" >&2
exit 1
fi
exec bash -c '. "$HOME/.nvm/nvm.sh" && nvm use "$1" >/dev/null &&
shift && exec yarn "$@"' bash "$NODE_VERSION" "$@"
}
main() { main() {
cd "$ROOT" cd "$ROOT"
go fmt ./... go fmt ./...
# run_yarn replaces this shell, so it stays the last step.
run_yarn run prettier --write '**/*.md' --tab-width 4 --prose-wrap always
} }
main "$@" main "$@"
+1 -23
View File
@@ -1,30 +1,10 @@
#!/bin/sh #!/bin/sh
# script/fmt-check: check formatting (read-only). Same scope as # script/fmt-check: check formatting (read-only). Same scope as
# script/fmt, the Go code and every Markdown file, but fails instead of # script/fmt, but fails instead of writing.
# writing.
set -eu set -eu
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)" ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
# Must match the pin in script/bootstrap.
NODE_VERSION="22.17.0"
# script/bootstrap installs node and yarn under nvm and leaves neither
# on the PATH of the shell that called it, so resolve the pinned
# toolchain here the way bootstrap's own install step does. nvm is a
# bash script, hence the subshell.
run_yarn() {
if command -v yarn >/dev/null 2>&1; then
exec yarn "$@"
fi
if [ ! -s "$HOME/.nvm/nvm.sh" ]; then
echo "fmt-check: no yarn; run script/bootstrap first" >&2
exit 1
fi
exec bash -c '. "$HOME/.nvm/nvm.sh" && nvm use "$1" >/dev/null &&
shift && exec yarn "$@"' bash "$NODE_VERSION" "$@"
}
main() { main() {
cd "$ROOT" cd "$ROOT"
unformatted="$(gofmt -l .)" unformatted="$(gofmt -l .)"
@@ -33,8 +13,6 @@ main() {
echo "$unformatted" >&2 echo "$unformatted" >&2
exit 1 exit 1
fi fi
# run_yarn replaces this shell, so it stays the last step.
run_yarn run prettier --check '**/*.md' --tab-width 4 --prose-wrap always
} }
main "$@" main "$@"
-8
View File
@@ -1,8 +0,0 @@
# THIS IS AN AUTOGENERATED FILE. DO NOT EDIT THIS FILE DIRECTLY.
# yarn lockfile v1
prettier@3.8.1:
version "3.8.1"
resolved "https://registry.yarnpkg.com/prettier/-/prettier-3.8.1.tgz#edf48977cf991558f4fcbd8a3ba6015ba2a3a173"
integrity sha512-UOnG6LftzbdaHZcKoPFtOcCKztrQ57WkHDeRD9t/PTQtmT0NHSeWWepj6pS0z/N7+08BHFDQVUrfmfMRcZwbMg==