Author SHA1 Message Date
sneak 9b0186c481 Say script/bootstrap pins node and yarn only when absent (refs #7)
check / check (push) Waiting to run
The README Entrypoints line and the TODO.md Completed Steps entry said
script/bootstrap installs a pinned node and yarn. It installs them at the
pinned versions only when they are missing and uses a node or yarn already
installed whatever its version, as its header comment says.

Model: opus-5-5
2026-10-06 08:25:55 +00:00
clawbot 8125d4aebf Reformat the existing Markdown with prettier (closes #7)
check / check (push) Waiting to run
The output of make fmt after the previous commit: README.md and TODO.md
are rewrapped at 80 columns, and the TODO.md Workflow list uses -
markers. No wording changes. REPO_POLICIES.md was already formatted.

Model: opus-5-5
2026-10-06 06:31:14 +00:00
clawbot 47e48a2529 Format Markdown with prettier in script/fmt and script/fmt-check (refs #7)
check / check (push) Waiting to run
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 pinned node (nvm from a hash-checked
archive) and yarn (corepack) install, then installs the locked
packages. Both fmt scripts find yarn the canonical way, sourcing nvm
when yarn is not on PATH. The new files and the yarn lookup come from
sneak/prompts at cc440118c876. The vendored REPO_POLICIES.md already
passes prettier, so it is not ignored. The Markdown is reformatted in
the next commit.

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

Model: opus-5-5
2026-10-06 06:30:59 +00:00
22 changed files with 563 additions and 936 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
@@ -0,0 +1,2 @@
node_modules/
yarn.lock
+4
View File
@@ -0,0 +1,4 @@
{
"tabWidth": 4,
"proseWrap": "always"
}
+2 -2
View File
@@ -1,8 +1,8 @@
# Lint phase. The linter is invoked directly rather than through `make
# lint` or `script/lint`, which are themselves a docker build and would
# recurse into a daemon that does not exist in a build step.
# golangci/golangci-lint:v2.14.0 (Debian-based), 2026-10-06
FROM golangci/golangci-lint@sha256:ad862ba6b3798cbe0fd9fd7408d498fd74fbd2623a92406b2fd3898faf0bf98f AS lint
# golangci/golangci-lint:v2.12.2-alpine, 2026-06-28
FROM golangci/golangci-lint:v2.12.2-alpine@sha256:91b27804074a0bacea298707f016911e60cf0cdbc6c7bf5ccacb5f0606d18d60 AS lint
WORKDIR /src
+121 -122
View File
@@ -1,28 +1,26 @@
# bsdaily
[bsdaily](https://git.eeqj.de/sneak/bsdaily) is a command-line utility
written in [Go](https://golang.org) that carves a single day (or a range of
days) of [Bluesky](https://bsky.app) firehose data out of a large,
continuously-growing SQLite database and writes it out as a self-contained,
[bsdaily](https://git.eeqj.de/sneak/bsdaily) is a command-line utility written
in [Go](https://golang.org) that carves a single day (or a range of days) of
[Bluesky](https://bsky.app) firehose data out of a large, continuously-growing
SQLite database and writes it out as a self-contained,
[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, and are designed to be published, archived, mirrored, and later
re-merged back into a single database.
named by date (e.g. `2026-06-27.sql.zst`), organized into per-month directories,
and are designed to be published, archived, mirrored, and later re-merged back
into a single database.
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 writing to the original database. The tool is operationally
conservative: it checks free disk space before starting, copies the snapshot
to fast scratch storage, processes one day at a time to avoid SQLite lock
contention, verifies every compressed output before publishing it, and
writes output atomically via a temp-file-and-rename so a partial run never
leaves a corrupt `.sql.zst` behind.
snapshot, so extraction never contends with the live firehose ingester that is
writing to the original database. The tool is operationally conservative: it
checks free disk space before starting, copies the snapshot to fast scratch
storage, processes one day at a time to avoid SQLite lock contention, verifies
every compressed output before publishing it, and 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, mergeable, publicly-mirrorable archive of the Bluesky firehose. It is
currently a one-person effort. The current version is pre-1.0 and there has
not yet been a versioned release; [SemVer](https://semver.org) will be used
for releases.
This project was written by [@sneak](https://sneak.berlin) to produce a daily,
mergeable, publicly-mirrorable archive of the Bluesky firehose. It is currently
a one-person effort. The current version is pre-1.0 and there has not yet been a
versioned release; [SemVer](https://semver.org) will be used for releases.
# Build Status
@@ -33,14 +31,14 @@ branch must always be green.
Primary development happens on a privately-run Gitea instance at
[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 a standard `go fmt`, syntactically
valid, and must pass the linting defined in the repository (the shared
`.golangci.yml` from `sneak/prompts`), which can be run with a `make lint`.
The `main` branch is protected and all changes must be made via [pull
requests](https://git.eeqj.de/sneak/bsdaily/pulls) and pass CI to be merged.
Changes must always be formatted with `make fmt` (`go fmt` for Go, prettier for
Markdown), syntactically valid, and must pass the linting defined in the
repository (presently the `golangci-lint` defaults), which can be run with a
`make lint`. The `main` branch is protected and all changes must be made via
[pull requests](https://git.eeqj.de/sneak/bsdaily/pulls) and pass CI to be
merged.
See [`REPO_POLICIES.md`](REPO_POLICIES.md) for detailed coding standards,
tooling requirements, and workflow conventions.
@@ -50,55 +48,55 @@ tooling requirements, and workflow conventions.
This repository adheres to the
[Scripts to Rule Them All](https://github.com/github/scripts-to-rule-them-all)
standard: normalized scripts in `script/` are the entrypoints for the
development workflow, and the Makefile targets are thin shims that call
them. We provide:
development workflow, and the Makefile targets are thin shims that call them. We
provide:
- `script/bootstrap` — install all development dependencies (go, Go
module download); the linter is not installed on the host
- `script/bootstrap` — install all development dependencies (go, Go module
download, node and yarn at pinned versions when absent, and the prettier
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/bootstrap`, then `script/install-precommit`
- `script/projectname` — print the project name (used for the Docker
image tags)
- `script/test` — build the Dockerfile's `test` phase, which runs the
test suite with `-race` (verbose rerun on failure)
- `script/projectname` — print the project name (used for the Docker image tags)
- `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
`golangci-lint run ./...`
- `script/fmt` — format the Go code with `go fmt` (writes)
- `script/fmt-check` — check the Go formatting with `gofmt` (read-only)
- `script/check` — run `script/test`, `script/lint`, and
`script/fmt-check`
- `script/docker` — build the Docker image tagged via
`script/projectname`; the build runs the `lint` and `test` phases
first
- `script/cibuild` — CI entrypoint: runs `script/bootstrap` and
`script/check`, then builds the Docker image tagged via
`script/projectname`
- `script/precommit` — pre-commit gate: `go mod tidy` (must not change
`go.mod` or `go.sum`) and `go fmt`, then `script/check`
- `script/install-precommit` — install the git pre-commit hook that
runs `script/precommit`
- `script/fmt` — format the Go code with `go fmt` and every Markdown file with
prettier (writes)
- `script/fmt-check` — check the Go formatting with `gofmt` and the Markdown
formatting with prettier (read-only)
- `script/check` — run `script/test`, `script/lint`, and `script/fmt-check`
- `script/docker` — build the Docker image tagged via `script/projectname`; the
build runs the `lint` and `test` phases first
- `script/cibuild` — CI entrypoint: runs `script/bootstrap` and `script/check`,
then builds the Docker image tagged via `script/projectname`
- `script/precommit` — pre-commit gate: `go mod tidy` (must not change `go.mod`
or `go.sum`) and `go fmt`, then `script/check`
- `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 always run rather than being served from the build cache.
Every Docker build in `script/` is uncached, so the `lint` and `test` phases
always run rather than being served from the build cache.
# Problem Statement
A Bluesky firehose ingester writes every observed post (and associated
users, hashtags, URLs, and media references) into a single ever-growing
SQLite database, `firehose.db`. This database has several properties that
make it awkward to publish or archive directly:
A Bluesky firehose ingester writes every observed post (and associated users,
hashtags, URLs, and media references) into a single ever-growing SQLite
database, `firehose.db`. This database has several properties that make it
awkward to publish or archive directly:
- It is **large and always growing**, so re-publishing the whole thing every
day is wasteful.
- It is **large and always growing**, so re-publishing the whole thing every day
is wasteful.
- It is **continuously written**, so reading from it directly risks lock
contention with the live ingester and inconsistent reads.
- It is **monolithic**, so there is no natural unit at which to mirror,
share, or distribute "just yesterday's posts".
- It is **monolithic**, so there is no natural unit at which to mirror, share,
or distribute "just yesterday's posts".
What is wanted instead is a stable, immutable, per-day artifact: a small
file containing exactly one calendar day of firehose data, cheap to publish,
cheap to mirror, and trivially re-mergeable into a full database by anyone
who collects a set of them.
What is wanted instead is a stable, immutable, per-day artifact: a small file
containing exactly one calendar day of firehose data, cheap to publish, cheap to
mirror, and trivially re-mergeable into a full database by anyone who collects a
set of them.
# Proposed Solution
@@ -108,81 +106,81 @@ A tool, `bsdaily`, that:
filesystem, so it reads from a consistent point-in-time copy that the live
ingester cannot be writing to;
- copies the snapshot's database files to fast scratch storage;
- **extracts** a single day's `posts` (and all rows reachable from them) into
a fresh, minimal per-day SQLite database;
- **dumps** that per-day database to SQL and pipes it through multithreaded
zstd compression;
- **verifies** the compressed output (zstd integrity check plus a sanity
check that the decompressed stream actually looks like SQL);
- **extracts** a single day's `posts` (and all rows reachable from them) into a
fresh, minimal per-day SQLite database;
- **dumps** that per-day database to SQL and pipes it through multithreaded zstd
compression;
- **verifies** the compressed output (zstd integrity check plus a sanity check
that the decompressed stream actually looks like SQL);
- **publishes** the result atomically as
`DailiesBase/YYYY-MM/YYYY-MM-DD.sql.zst`.
Each daily dump is emitted with `INSERT` statements over the full schema
(including the deduplicated `users`, `hashtags`, and `urls` lookup tables),
so any collection of daily dumps can be merged into a single database by
rewriting `INSERT INTO` to `INSERT OR IGNORE INTO` and replaying them in
sequence. Two helper scripts ([`merge_daily_dumps.sh`](merge_daily_dumps.sh)
and [`regenerate_auxiliary_tables.sql`](regenerate_auxiliary_tables.sql)) are
(including the deduplicated `users`, `hashtags`, and `urls` lookup tables), so
any collection of daily dumps can be merged into a single database by rewriting
`INSERT INTO` to `INSERT OR IGNORE INTO` and replaying them in sequence. Two
helper scripts ([`merge_daily_dumps.sh`](merge_daily_dumps.sh) and
[`regenerate_auxiliary_tables.sql`](regenerate_auxiliary_tables.sql)) are
included to do exactly this and to rebuild the aggregate statistics
(`use_count`, `first_seen`, user `resolved_at`/`updated_at`) afterward.
# Design Goals
- **Never disturb the live ingester.** All reads come from a ZFS snapshot,
never the live database.
- **Never disturb the live ingester.** All reads come from a ZFS snapshot, never
the live database.
- **Crash-safe, idempotent runs.** Output is written to a temp file and
atomically renamed; a day whose final output already exists is skipped, so
re-running a range is safe and resumable.
- **Mergeable output.** Daily dumps re-combine losslessly into a full
database via `INSERT OR IGNORE`.
- **Mergeable output.** Daily dumps re-combine losslessly into a full database
via `INSERT OR IGNORE`.
- **Operationally cautious.** Free-space preflight checks on both scratch and
output filesystems; explicit verification of every artifact before it is
published.
- **Fast where it's free.** Large snapshot copies use a 256MiB buffer,
pre-allocate the destination, and (on Linux) issue `posix_fadvise`
sequential/willneed hints; extraction uses aggressive,
crash-unsafe-by-design SQLite pragmas because the working data lives only
in disposable scratch space.
sequential/willneed hints; extraction uses aggressive, crash-unsafe-by-design
SQLite pragmas because the working data lives only in disposable scratch
space.
# Non-Goals
- **Real-time export.** `bsdaily` operates on daily snapshots; the freshest
day it can produce is the snapshot date minus one.
- **Schema ownership.** The schema is defined by the upstream firehose
ingester; [`schema.sql`](schema.sql) is included for reference only.
`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 check and fadvise hints use `golang.org/x/sys/unix`; a non-Linux
build compiles but is a no-op for the fadvise hints). The hard-coded paths
assume the production host's ZFS layout.
- **Real-time export.** `bsdaily` operates on daily snapshots; the freshest day
it can produce is the snapshot date minus one.
- **Schema ownership.** The schema is defined by the upstream firehose ingester;
[`schema.sql`](schema.sql) is included for reference only. `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
check and fadvise hints use `golang.org/x/sys/unix`; a non-Linux build
compiles but is a no-op for the fadvise hints). The hard-coded paths assume
the production host's ZFS layout.
# How It Works
A single run proceeds as follows:
1. **Find the snapshot.** Scan `SnapshotBase` for directories matching
`zfs-auto-snap_daily-YYYY-MM-DD-NNNN`, pick the most recent, and confirm
it contains `firehose.db`.
2. **Determine target days.** Default to the snapshot date minus one day;
or use `--date`, or every day in the inclusive `--from`/`--to` range.
`zfs-auto-snap_daily-YYYY-MM-DD-NNNN`, pick the most recent, and confirm it
contains `firehose.db`.
2. **Determine target days.** Default to the snapshot date minus one day; or use
`--date`, or every day in the inclusive `--from`/`--to` range.
3. **Preflight disk space.** Require at least 500GiB free on the scratch
filesystem and 20GiB free on the output filesystem.
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
`TmpBase`.
5. **Per day**, processed strictly one at a time to avoid SQLite contention:
- skip the day if its final output file already exists;
- `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 reachable from them (`posts_hashtags`, `posts_urls`, `hashtags`,
`urls`, `users`, and `media` if that table exists);
- abort the day cleanly if there are zero posts (`ErrNoPosts`), rather
than emitting an empty dump;
- recreate indexes, detach the source, and verify the inserted row count;
- `sqlite3 .dump | zstdmt` into a hidden temp file;
- run a zstd integrity check and confirm the decompressed head looks like
SQL;
- atomically rename into place and delete the per-day scratch DB.
- skip the day if its final output file already exists;
- `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
reachable from them (`posts_hashtags`, `posts_urls`, `hashtags`, `urls`,
`users`, and `media` if that table exists);
- abort the day cleanly if there are zero posts (`ErrNoPosts`), rather than
emitting an empty dump;
- recreate indexes, detach the source, and verify the inserted row count;
- `sqlite3 .dump | zstdmt` into a hidden temp file;
- run a zstd integrity check and confirm the decompressed head looks like
SQL;
- atomically rename into place and delete the per-day scratch DB.
6. **Clean up** the temp directory and log a processed/skipped/total summary.
# Usage
@@ -219,9 +217,9 @@ sqlite3 merged.db < regenerate_auxiliary_tables.sql
- **Go** (see [`go.mod`](go.mod) for the toolchain version) to build.
- **Linux** for production use (ZFS snapshots, `statfs` free-space checks,
`posix_fadvise` hints).
- The **`sqlite3`** and **`zstdmt`** (multithreaded zstd) binaries on
`PATH`; `zstdcat` is used for verification. SQLite reads/writes during
extraction use the pure-Go [`modernc.org/sqlite`](https://pkg.go.dev/modernc.org/sqlite)
- The **`sqlite3`** and **`zstdmt`** (multithreaded zstd) binaries on `PATH`;
`zstdcat` is used for verification. SQLite reads/writes during extraction use
the pure-Go [`modernc.org/sqlite`](https://pkg.go.dev/modernc.org/sqlite)
driver, so no cgo is required for that part.
# Configuration
@@ -241,21 +239,21 @@ elsewhere.
# Data Model
The firehose schema (reference copy in [`schema.sql`](schema.sql)) centers
on a `posts` table, with `users` keyed by DID and many-to-many junction
tables linking posts to deduplicated `hashtags` and `urls`. An optional
`media` table tracks downloaded blobs by content hash. `bsdaily` does not
own this schema; it reflects whatever DDL exists in the source snapshot and
selects forward from `posts` along the foreign-key relationships to produce a
referentially-complete per-day slice.
The firehose schema (reference copy in [`schema.sql`](schema.sql)) centers on a
`posts` table, with `users` keyed by DID and many-to-many junction tables
linking posts to deduplicated `hashtags` and `urls`. An optional `media` table
tracks downloaded blobs by content hash. `bsdaily` does not own this schema; it
reflects whatever DDL exists in the source snapshot and selects forward from
`posts` along the foreign-key relationships to produce a referentially-complete
per-day slice.
# Use Cases
## Daily public archive
Publish one small, immutable file per day to static HTTP (or IPFS, or a
mirror network) so that anyone can fetch exactly the day(s) they want and
re-merge them locally.
Publish one small, immutable file per day to static HTTP (or IPFS, or a mirror
network) so that anyone can fetch exactly the day(s) they want and re-merge them
locally.
## Backfilling a range
@@ -265,16 +263,17 @@ safe to re-run.
## Reconstituting a full database
Collect any set of daily dumps and merge them with `INSERT OR IGNORE` to
rebuild a complete, queryable SQLite database, then regenerate the aggregate
statistics tables.
Collect any set of daily dumps and merge them with `INSERT OR IGNORE` to rebuild
a complete, queryable SQLite database, then regenerate the aggregate statistics
tables.
# See Also
## Links
- Repo: [https://git.eeqj.de/sneak/bsdaily](https://git.eeqj.de/sneak/bsdaily)
- Issues: [https://git.eeqj.de/sneak/bsdaily/issues](https://git.eeqj.de/sneak/bsdaily/issues)
- Issues:
[https://git.eeqj.de/sneak/bsdaily/issues](https://git.eeqj.de/sneak/bsdaily/issues)
- Bluesky: [https://bsky.app](https://bsky.app)
- zstd: [https://facebook.github.io/zstd/](https://facebook.github.io/zstd/)
+31 -28
View File
@@ -1,12 +1,12 @@
# Workflow
* branch (from `main`)
* do the work in Next Step
* move Next Step to the top of Completed Steps
* move the top item of Future Steps into Next Step
* commit (`TODO.md` changes in the same commit as the work)
* merge to `main` if the branch is not protected, otherwise open a PR
* push
- branch (from `main`)
- do the work in Next Step
- move Next Step to the top of Completed Steps
- move the top item of Future Steps into Next Step
- commit (`TODO.md` changes in the same commit as the work)
- merge to `main` if the branch is not protected, otherwise open a PR
- push
# Status
@@ -14,35 +14,38 @@ pre-1.0
# Next Step
Format Markdown with prettier in `script/fmt` and `script/fmt-check`
(https://git.eeqj.de/sneak/bsdaily/issues/7).
Add the canonical `.golangci.yml`, move the lint phase to golangci-lint v2.14.0
in the same commit, and fix the findings it surfaces
(https://git.eeqj.de/sneak/bsdaily/issues/6).
# 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
`.gitignore`, `.dockerignore` and `.editorconfig`; `lint` and `test`
phases in the `Dockerfile`, built by `script/lint` and `script/test`;
canonical `script/cibuild`, `script/docker` and CI workflow; no linter
installed on the host; re-vendored `REPO_POLICIES.md`.
- 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints,
Makefile shims, README Entrypoints section
- 2026-06-28: Fixed errcheck lint failures; added compilation smoke test;
tidied go.mod.
- 2026-06-28: Added repo scaffolding: README, LICENSE, Makefile,
Dockerfile, REPO_POLICIES.md, and Gitea CI.
- 2026-02-12: Fixed SQLite database locking by removing parallel
processing; fixed Linux build via golang.org/x/sys/unix Fadvise.
- 2026-02-12: Optimized file copy for large databases; moved temp
directory to NVMe scratch storage.
`.gitignore`, `.dockerignore` and `.editorconfig`; `lint` and `test` phases in
the `Dockerfile`, built by `script/lint` and `script/test`; canonical
`script/cibuild`, `script/docker` and CI workflow; no linter installed on the
host; re-vendored `REPO_POLICIES.md`.
- 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints, Makefile
shims, README Entrypoints section
- 2026-06-28: Fixed errcheck lint failures; added compilation smoke test; tidied
go.mod.
- 2026-06-28: Added repo scaffolding: README, LICENSE, Makefile, Dockerfile,
REPO_POLICIES.md, and Gitea CI.
- 2026-02-12: Fixed SQLite database locking by removing parallel processing;
fixed Linux build via golang.org/x/sys/unix Fadvise.
- 2026-02-12: Optimized file copy for large databases; moved temp directory to
NVMe scratch storage.
- 2026-02-11: Added date range support.
- 2026-02-09: Initial implementation: single-day extraction, specific-date
targeting, faster pruning of throwaway database copies.
# Future Steps
- Expand tests beyond the compilation smoke test: unit tests for the
extraction, verification, and atomic-publish paths.
- 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.
+48 -89
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
import (
"errors"
"fmt"
"log/slog"
"os"
@@ -13,113 +10,75 @@ import (
"github.com/spf13/cobra"
)
var (
errDateExclusive = errors.New("--date and --from/--to are mutually exclusive")
errFromRequiresTo = errors.New("--from requires --to")
errToRequiresFrom = errors.New("--to requires --from")
errFromAfterTo = errors.New("--from is after --to")
)
func main() {
logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{
Level: slog.LevelInfo,
}))
slog.SetDefault(logger)
var dateFlag, fromFlag, toFlag string
var dateFlag string
var fromFlag string
var toFlag string
rootCmd := &cobra.Command{
Use: "bsdaily",
Short: "Extract a single day's data from the latest daily snapshot",
SilenceUsage: true,
RunE: func(_ *cobra.Command, _ []string) error {
targetDates, err := parseTargetDates(dateFlag, fromFlag, toFlag)
if err != nil {
return err
RunE: func(cmd *cobra.Command, args []string) error {
hasDate := dateFlag != ""
hasFrom := fromFlag != ""
hasTo := toFlag != ""
// Validate mutual exclusivity
if hasDate && (hasFrom || hasTo) {
return fmt.Errorf("--date and --from/--to are mutually exclusive")
}
if hasFrom != hasTo {
if hasFrom {
return fmt.Errorf("--from requires --to")
}
return fmt.Errorf("--to requires --from")
}
err = bsdaily.Run(targetDates)
if err != nil {
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
}
slog.Info("completed successfully")
return nil
},
}
rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "",
"target date to extract (YYYY-MM-DD); "+
"defaults to snapshot date minus one day")
rootCmd.Flags().StringVar(&fromFlag, "from", "",
"start of date range to extract (YYYY-MM-DD, inclusive); use with --to")
rootCmd.Flags().StringVar(&toFlag, "to", "",
"end of date range to extract (YYYY-MM-DD, inclusive); use with --from")
rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "", "target date to extract (YYYY-MM-DD); defaults to snapshot date minus one day")
rootCmd.Flags().StringVar(&fromFlag, "from", "", "start of date range to extract (YYYY-MM-DD, inclusive); use with --to")
rootCmd.Flags().StringVar(&toFlag, "to", "", "end of date range to extract (YYYY-MM-DD, inclusive); use with --from")
err := rootCmd.Execute()
if err != nil {
if err := rootCmd.Execute(); err != nil {
os.Exit(1)
}
}
// parseTargetDates turns the --date, --from and --to flags into the days
// to extract. It returns nil when none of them is set, which Run takes
// to mean the snapshot date minus one day.
func parseTargetDates(dateFlag, fromFlag, toFlag string) ([]time.Time, error) {
hasDate := dateFlag != ""
hasFrom := fromFlag != ""
hasTo := toFlag != ""
// Validate mutual exclusivity
if hasDate && (hasFrom || hasTo) {
return nil, errDateExclusive
}
if hasFrom != hasTo {
if hasFrom {
return nil, errFromRequiresTo
}
return nil, errToRequiresFrom
}
if hasDate {
t, err := time.Parse("2006-01-02", dateFlag)
if err != nil {
return nil, fmt.Errorf(
"invalid --date %q (expected YYYY-MM-DD): %w", dateFlag, err)
}
return []time.Time{t}, nil
}
if !hasFrom {
return nil, nil
}
from, err := time.Parse("2006-01-02", fromFlag)
if err != nil {
return nil, fmt.Errorf(
"invalid --from %q (expected YYYY-MM-DD): %w", fromFlag, err)
}
to, err := time.Parse("2006-01-02", toFlag)
if err != nil {
return nil, fmt.Errorf(
"invalid --to %q (expected YYYY-MM-DD): %w", toFlag, err)
}
if from.After(to) {
return nil, fmt.Errorf("%w (--from %s, --to %s)",
errFromAfterTo, fromFlag, 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 (
"testing"
"git.eeqj.de/sneak/bsdaily/internal/bsdaily"
)
import "testing"
// TestCompiles is a minimal smoke test that references the package's exported
// surface so that `go test` fails if the package stops compiling. It does not
// touch the filesystem or any of the hard-coded production paths.
func TestCompiles(t *testing.T) {
t.Parallel()
if bsdaily.DBFilename == "" || bsdaily.WALFilename == "" ||
bsdaily.SHMFilename == "" {
if DBFilename == "" || WALFilename == "" || SHMFilename == "" {
t.Fatal("expected database filename constants to be set")
}
if bsdaily.SnapshotBase == "" || bsdaily.TmpBase == "" ||
bsdaily.DailiesBase == "" {
if SnapshotBase == "" || TmpBase == "" || DailiesBase == "" {
t.Fatal("expected base path constants to be set")
}
if bsdaily.MinTmpFreeBytes == 0 || bsdaily.MinDailiesFreeBytes == 0 {
if MinTmpFreeBytes == 0 || MinDailiesFreeBytes == 0 {
t.Fatal("expected free-space thresholds to be set")
}
if bsdaily.ErrNoPosts == nil {
if ErrNoPosts == nil {
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
import "regexp"
// Paths, file names and tuning for a run. They are set for one
// production host; see the README.
const (
SnapshotBase = "/srv/berlin.sneak.fs.blueskyarchive/.zfs/snapshot"
TmpBase = "/srv/storage/tmp"
@@ -31,5 +27,4 @@ const (
verificationHeadLines = 20
)
var snapshotPattern = regexp.MustCompile(
`^zfs-auto-snap_daily-(\d{4}-\d{2}-\d{2})-\d{4}$`)
var snapshotPattern = regexp.MustCompile(`^zfs-auto-snap_daily-(\d{4}-\d{2}-\d{2})-\d{4}$`)
+10 -29
View File
@@ -1,39 +1,28 @@
package bsdaily
import (
"errors"
"fmt"
"io"
"log/slog"
"os"
"path/filepath"
"time"
)
const (
// 256MB buffer for large file copies from fast storage
copyBufferSize = 256 * 1024 * 1024
oneMB = 1024 * 1024
copyBufferSize = 256 * 1024 * 1024 // 256MB buffer for large file copies from fast storage
oneGB = 1024 * 1024 * 1024
)
var errShortCopy = errors.New("short copy")
// CopyFile copies src to dst through a large buffer, pre-allocating dst
// and syncing it to disk before returning.
func CopyFile(src, dst string) (err error) {
startTime := time.Now()
slog.Info("copying file", "src", src, "dst", dst)
srcFile, err := os.Open(filepath.Clean(src))
srcFile, err := os.Open(src)
if err != nil {
return fmt.Errorf("opening source %s: %w", src, err)
}
defer func() {
cerr := srcFile.Close()
if cerr != nil {
if cerr := srcFile.Close(); cerr != nil {
slog.Warn("failed to close source file", "src", src, "error", cerr)
}
}()
@@ -48,48 +37,40 @@ func CopyFile(src, dst string) (err error) {
applyFileAdvice(srcFile, srcInfo.Size())
}
dstFile, err := os.Create(filepath.Clean(dst))
dstFile, err := os.Create(dst)
if err != nil {
return fmt.Errorf("creating destination %s: %w", dst, err)
}
defer func() {
cerr := dstFile.Close()
if cerr != nil && err == nil {
if cerr := dstFile.Close(); cerr != nil && err == nil {
err = fmt.Errorf("closing destination %s: %w", dst, cerr)
}
}()
// Pre-allocate space for the destination file to avoid fragmentation
truncErr := dstFile.Truncate(srcInfo.Size())
if truncErr != nil {
slog.Warn("failed to pre-allocate destination file", "error", truncErr)
if err := dstFile.Truncate(srcInfo.Size()); err != nil {
slog.Warn("failed to pre-allocate destination file", "error", err)
}
// Use a much larger buffer for NVMe-speed copies
buf := make([]byte, copyBufferSize)
written, err := io.CopyBuffer(dstFile, srcFile, buf)
if err != nil {
return fmt.Errorf("copying data: %w", err)
}
if written != srcInfo.Size() {
return fmt.Errorf("%w: wrote %d bytes, expected %d",
errShortCopy, written, srcInfo.Size())
return fmt.Errorf("short copy: wrote %d bytes, expected %d", written, srcInfo.Size())
}
err = dstFile.Sync()
if err != nil {
if err := dstFile.Sync(); err != nil {
return fmt.Errorf("syncing destination %s: %w", dst, err)
}
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,
"elapsed", elapsed.Round(time.Millisecond),
"throughput_mbps", fmt.Sprintf("%.1f", throughputMBps))
return nil
}
+4 -3
View File
@@ -9,7 +9,8 @@ import (
func applyFileAdvice(file *os.File, size int64) {
fd := int(file.Fd())
_ = unix.Fadvise(fd, 0, size, unix.FADV_SEQUENTIAL)
// Prefetch the file into the page cache
_ = unix.Fadvise(fd, 0, size, unix.FADV_WILLNEED)
// POSIX_FADV_SEQUENTIAL = 2
_ = unix.Fadvise(fd, 0, size, 2)
// POSIX_FADV_WILLNEED = 3 - prefetch file into cache
_ = unix.Fadvise(fd, 0, size, 3)
}
+4 -16
View File
@@ -1,38 +1,26 @@
package bsdaily
import (
"errors"
"fmt"
"log/slog"
"golang.org/x/sys/unix"
)
var errInsufficientSpace = errors.New("insufficient disk space")
// CheckFreeSpace returns an error when the filesystem holding path has
// fewer than minBytes bytes available. label names the location in the
// log line and the error.
func CheckFreeSpace(path string, minBytes uint64, label string) error {
var stat unix.Statfs_t
err := unix.Statfs(path, &stat)
if err != nil {
if err := unix.Statfs(path, &stat); err != nil {
return fmt.Errorf("statfs %s (%s): %w", path, label, err)
}
free := stat.Bavail * uint64(stat.Bsize) //nolint:gosec // Bsize is never negative
free := uint64(stat.Bavail) * uint64(stat.Bsize)
freeGB := float64(free) / float64(bytesPerGB)
minGB := float64(minBytes) / float64(bytesPerGB)
slog.Info("disk space check", "label", label, "path", path,
"free_gb", fmt.Sprintf("%.1f", freeGB),
"required_gb", fmt.Sprintf("%.1f", minGB))
if free < minBytes {
return fmt.Errorf("%w on %s (%s): %.1f GB free, need %.1f GB",
errInsufficientSpace, path, label, freeGB, minGB)
return fmt.Errorf("insufficient disk space on %s (%s): %.1f GB free, need %.1f GB",
path, label, freeGB, minGB)
}
return nil
}
+36 -74
View File
@@ -1,8 +1,6 @@
package bsdaily
import (
"context"
"errors"
"fmt"
"log/slog"
"os"
@@ -11,43 +9,61 @@ import (
"strings"
)
var errEmptyOutput = errors.New("compressed output is empty")
// DumpAndCompress writes a `sqlite3 .dump` of the database at dbPath,
// compressed by zstdmt, to outputPath.
func DumpAndCompress(dbPath, outputPath string) (err error) {
for _, tool := range []string{"sqlite3", "zstdmt"} {
_, err = exec.LookPath(tool)
if err != nil {
if _, err := exec.LookPath(tool); err != nil {
return fmt.Errorf("required tool %q not found in PATH: %w", tool, err)
}
}
err = CheckFreeSpace(filepath.Dir(outputPath), MinDailiesFreeBytes,
"dailiesBase (pre-dump)")
if err != nil {
if err := CheckFreeSpace(filepath.Dir(outputPath), MinDailiesFreeBytes, "dailiesBase (pre-dump)"); err != nil {
return err
}
outFile, err := os.Create(filepath.Clean(outputPath))
outFile, err := os.Create(outputPath)
if err != nil {
return fmt.Errorf("creating output file: %w", err)
}
defer func() {
cerr := outFile.Close()
if cerr != nil && err == nil {
if cerr := outFile.Close(); cerr != nil && err == nil {
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 {
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 != nil {
if err := dumpCmd.Wait(); err != nil {
return fmt.Errorf("sqlite3 dump failed: %w; stderr: %s", err, dumpStderr.String())
}
if err := zstdCmd.Wait(); err != nil {
return fmt.Errorf("zstdmt failed: %w; stderr: %s", err, zstdStderr.String())
}
if err := outFile.Sync(); err != nil {
return fmt.Errorf("syncing output: %w", err)
}
@@ -55,66 +71,12 @@ func DumpAndCompress(dbPath, outputPath string) (err error) {
if err != nil {
return fmt.Errorf("stat output: %w", err)
}
const bytesPerMB = 1024 * 1024
slog.Info("compressed output written", "path", outputPath,
"size_bytes", info.Size(), "size_mb", info.Size()/bytesPerMB)
if info.Size() == 0 {
return 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 fmt.Errorf("compressed output is empty")
}
return nil
+52 -179
View File
@@ -1,29 +1,22 @@
package bsdaily
import (
"context"
"database/sql"
"errors"
"fmt"
"log/slog"
"time"
// Registers the "sqlite" driver with database/sql.
_ "modernc.org/sqlite"
)
// ErrNoPosts is returned by ExtractDay when the source holds no posts for
// the target day.
var ErrNoPosts = errors.New("no posts found for target day")
var errPostCountMismatch = errors.New("post count mismatch")
// ExtractDay opens a new empty database at dstDBPath, attaches srcDBPath,
// and copies only the target day's data into it. This is much faster than
// pruning a full copy because it only reads/writes the small slice of data
// being kept.
func ExtractDay(srcDBPath, dstDBPath string, targetDay time.Time) error {
ctx := context.Background()
dayStart := targetDay.Format("2006-01-02") + "T00:00:00"
dayEnd := targetDay.AddDate(0, 0, 1).Format("2006-01-02") + "T00:00:00"
@@ -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
// Use WAL mode for the source attachment to avoid locking issues
pragmas := fmt.Sprintf("?_pragma=journal_mode(WAL)"+
"&_pragma=synchronous(OFF)&_pragma=cache_size(%d)"+
"&_pragma=foreign_keys(OFF)&_pragma=temp_store(MEMORY)"+
"&_pragma=busy_timeout(5000)", sqliteCacheSizeKB)
pragmas := fmt.Sprintf("?_pragma=journal_mode(WAL)&_pragma=synchronous(OFF)&_pragma=cache_size(%d)&_pragma=foreign_keys(OFF)&_pragma=temp_store(MEMORY)&_pragma=busy_timeout(5000)", sqliteCacheSizeKB)
db, err := sql.Open("sqlite", dstDBPath+pragmas)
if err != nil {
return fmt.Errorf("opening destination database: %w", err)
}
defer func() {
cerr := db.Close()
if cerr != nil {
slog.Warn("failed to close destination database",
"path", dstDBPath, "error", cerr)
if cerr := db.Close(); cerr != nil {
slog.Warn("failed to close destination database", "path", dstDBPath, "error", cerr)
}
}()
// Attach source database
_, err = db.ExecContext(ctx, "ATTACH DATABASE ? AS src", srcDBPath)
if err != nil {
if _, err := db.Exec("ATTACH DATABASE ? AS src", srcDBPath); err != nil {
return fmt.Errorf("attaching source database: %w", err)
}
err = copyTables(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
}
// copyTables creates every table of the attached source database, empty,
// in the destination database.
func copyTables(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")
rows, err := db.Query("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 {
if cerr := rows.Close(); 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 {
if err := rows.Scan(&ddl); err != nil {
return fmt.Errorf("scanning DDL: %w", err)
}
ddlStatements = append(ddlStatements, ddl)
}
err = rows.Err()
if err != nil {
if err := rows.Err(); err != nil {
return fmt.Errorf("iterating DDL rows: %w", err)
}
for _, ddl := range ddlStatements {
_, err = db.ExecContext(ctx, ddl)
if err != nil {
if _, err := db.Exec(ddl); 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)
tx, err := db.Begin()
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() {
rerr := tx.Rollback()
if rerr != nil && !errors.Is(rerr, sql.ErrTxDone) {
slog.Warn("failed to roll back transaction", "error", rerr)
if err != nil {
if rerr := tx.Rollback(); rerr != nil && !errors.Is(rerr, sql.ErrTxDone) {
slog.Warn("failed to roll back transaction", "error", rerr)
}
}
}()
// Insert target day's data
slog.Info("inserting posts for target day")
//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)
result, err := tx.Exec("INSERT INTO posts SELECT * FROM src.posts WHERE timestamp >= ? AND timestamp < ?", dayStart, dayEnd)
if err != nil {
return 0, fmt.Errorf("inserting posts: %w", err)
return 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",
return 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 {
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)
}
_, err = tx.ExecContext(ctx, "INSERT INTO posts_urls "+
"SELECT * FROM src.posts_urls WHERE post_id IN (SELECT id FROM posts)")
if err != nil {
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)
}
_, err = tx.ExecContext(ctx, "INSERT INTO hashtags "+
"SELECT * FROM src.hashtags "+
"WHERE id IN (SELECT hashtag_id FROM posts_hashtags)")
if err != nil {
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)
}
_, err = tx.ExecContext(ctx, "INSERT INTO urls "+
"SELECT * FROM src.urls WHERE id IN (SELECT url_id FROM posts_urls)")
if err != nil {
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)
}
_, err = tx.ExecContext(ctx, "INSERT INTO users "+
"SELECT * FROM src.users WHERE did IN (SELECT user_did FROM posts)")
if err != nil {
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)
}
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 {
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
//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)
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)
}
}
}
// 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 {
// 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.QueryContext(ctx, "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.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() {
cerr := idxRows.Close()
if cerr != nil {
if cerr := idxRows.Close(); 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 {
if err := idxRows.Scan(&idxSQL); err != nil {
return fmt.Errorf("scanning index DDL: %w", err)
}
idxStatements = append(idxStatements, idxSQL)
}
err = idxRows.Err()
if err != nil {
if err := idxRows.Err(); err != nil {
return fmt.Errorf("iterating index rows: %w", err)
}
for _, idxSQL := range idxStatements {
_, err = db.ExecContext(ctx, idxSQL)
if err != nil {
if _, err := db.Exec(idxSQL); err != nil {
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
}
+78 -165
View File
@@ -9,29 +9,21 @@ import (
"time"
)
var errEmptySource = errors.New("source file is empty")
// cleanup removes a temporary file, logging a warning if removal fails so
// that leaked scratch files are surfaced rather than silently ignored. A
// missing file is not an error.
func cleanup(path string) {
err := os.Remove(path)
if err != nil && !os.IsNotExist(err) {
if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
slog.Warn("failed to remove temporary file", "path", path, "error", err)
}
}
// Run writes a compressed SQL dump of each day in targetDates, taken from
// the latest daily snapshot, skipping days that already have one or have
// no posts. With no dates it does the day before the snapshot date.
func Run(targetDates []time.Time) error {
snapshotDir, snapshotDate, err := FindLatestDailySnapshot()
if err != nil {
return fmt.Errorf("finding latest snapshot: %w", err)
}
slog.Info("found latest daily snapshot", "dir", snapshotDir,
"snapshot_date", snapshotDate.Format("2006-01-02"))
slog.Info("found latest daily snapshot", "dir", snapshotDir, "snapshot_date", snapshotDate.Format("2006-01-02"))
if len(targetDates) == 0 {
targetDates = []time.Time{snapshotDate.AddDate(0, 0, -1)}
@@ -42,13 +34,10 @@ func Run(targetDates []time.Time) error {
"last", targetDates[len(targetDates)-1].Format("2006-01-02"))
// Check disk space
err = CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase")
if err != nil {
if err := CheckFreeSpace(TmpBase, MinTmpFreeBytes, "tmpBase"); err != nil {
return err
}
err = CheckFreeSpace(DailiesBase, MinDailiesFreeBytes, "dailiesBase")
if err != nil {
if err := CheckFreeSpace(DailiesBase, MinDailiesFreeBytes, "dailiesBase"); err != nil {
return err
}
@@ -57,53 +46,14 @@ func Run(targetDates []time.Time) error {
if err != nil {
return fmt.Errorf("creating temp directory in %s: %w", TmpBase, err)
}
slog.Info("created temp directory", "path", tmpDir)
defer func() {
slog.Info("cleaning up temp directory", "path", tmpDir)
rerr := os.RemoveAll(tmpDir)
if rerr != nil {
slog.Error("failed to remove temp directory",
"path", tmpDir, "error", rerr)
if err := os.RemoveAll(tmpDir); err != nil {
slog.Error("failed to remove temp directory", "path", tmpDir, "error", err)
}
}()
dstDB, err := copySnapshotFiles(snapshotDir, tmpDir)
if err != nil {
return err
}
// Process each day completely before moving to the next. This ensures
// we don't have multiple SQLite operations competing for the same
// source database.
processed := 0
skipped := 0
for _, targetDay := range targetDates {
written, err := processDay(tmpDir, dstDB, targetDay)
if err != nil {
return err
}
if written {
processed++
} else {
skipped++
}
}
slog.Info("run summary", "processed", processed, "skipped", skipped,
"total", len(targetDates))
return nil
}
// copySnapshotFiles copies the database, its WAL and, if present, its SHM
// file from snapshotDir into tmpDir, and returns the copied database's
// path.
func copySnapshotFiles(snapshotDir, tmpDir string) (string, error) {
// Copy database files from snapshot to temp
srcDB := filepath.Join(snapshotDir, DBFilename)
srcWAL := filepath.Join(snapshotDir, WALFilename)
@@ -115,137 +65,100 @@ func copySnapshotFiles(snapshotDir, tmpDir string) (string, error) {
for _, f := range []string{srcDB, srcWAL} {
info, err := os.Stat(f)
if err != nil {
return "", fmt.Errorf("source file missing: %s: %w", f, err)
return fmt.Errorf("source file missing: %s: %w", f, err)
}
if info.Size() == 0 {
return "", fmt.Errorf("%w: %s", errEmptySource, f)
return fmt.Errorf("source file is empty: %s", f)
}
slog.Info("source file", "path", f, "size_bytes", info.Size())
}
err := CopyFile(srcDB, dstDB)
if err != nil {
return "", fmt.Errorf("copying database: %w", err)
if err := CopyFile(srcDB, dstDB); err != nil {
return fmt.Errorf("copying database: %w", err)
}
err = CopyFile(srcWAL, dstWAL)
if err != nil {
return "", fmt.Errorf("copying WAL: %w", err)
if err := CopyFile(srcWAL, dstWAL); err != nil {
return fmt.Errorf("copying WAL: %w", err)
}
_, err = os.Stat(srcSHM)
if err == nil {
err = CopyFile(srcSHM, dstSHM)
if err != nil {
return "", fmt.Errorf("copying SHM: %w", err)
if _, err := os.Stat(srcSHM); err == nil {
if err := CopyFile(srcSHM, dstSHM); 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
// 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")
for _, targetDay := range targetDates {
dayStr := targetDay.Format("2006-01-02")
slog.Info("processing day", "date", dayStr)
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
outputDir := filepath.Join(DailiesBase, targetDay.Format("2006-01"))
outputFinal := filepath.Join(outputDir, dayStr+".sql.zst")
// 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)
}
_, 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)
// Dump to SQL and compress
if err := os.MkdirAll(outputDir, 0755); err != nil {
cleanup(extractedDB)
return false, nil
return fmt.Errorf("creating output directory %s: %w", outputDir, err)
}
return false, fmt.Errorf("extracting day %s: %w", dayStr, err)
}
outputTmp := filepath.Join(outputDir, "."+dayStr+".sql.zst.tmp")
// 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 {
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 {
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)
return false, fmt.Errorf("creating output directory %s: %w", outputDir, err)
processed++
}
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())
slog.Info("run summary", "processed", processed, "skipped", skipped, "total", len(targetDates))
return nil
}
+7 -23
View File
@@ -1,7 +1,6 @@
package bsdaily
import (
"errors"
"fmt"
"log/slog"
"os"
@@ -10,16 +9,10 @@ import (
"time"
)
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) {
func FindLatestDailySnapshot() (dir string, snapshotDate time.Time, err error) {
entries, err := os.ReadDir(SnapshotBase)
if err != nil {
return "", time.Time{}, fmt.Errorf("reading snapshot directory %s: %w",
SnapshotBase, err)
return "", time.Time{}, fmt.Errorf("reading snapshot directory %s: %w", SnapshotBase, err)
}
type snapshot struct {
@@ -28,30 +21,24 @@ func FindLatestDailySnapshot() (string, time.Time, error) {
}
var snapshots []snapshot
for _, e := range entries {
if !e.IsDir() {
continue
}
m := snapshotPattern.FindStringSubmatch(e.Name())
if m == nil {
continue
}
d, err := time.Parse("2006-01-02", m[1])
if err != nil {
slog.Warn("skipping snapshot with unparseable date",
"name", e.Name(), "error", err)
slog.Warn("skipping snapshot with unparseable date", "name", e.Name(), "error", err)
continue
}
snapshots = append(snapshots, snapshot{name: e.Name(), date: d})
}
if len(snapshots) == 0 {
return "", time.Time{}, fmt.Errorf("%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 {
@@ -59,14 +46,11 @@ func FindLatestDailySnapshot() (string, time.Time, error) {
})
latest := snapshots[0]
dir := filepath.Join(SnapshotBase, latest.name)
dir = filepath.Join(SnapshotBase, latest.name)
dbPath := filepath.Join(dir, DBFilename)
_, err = os.Stat(dbPath)
if err != nil {
return "", time.Time{}, fmt.Errorf("database not found in snapshot %s: %w",
dir, err)
if _, err := os.Stat(dbPath); err != nil {
return "", time.Time{}, fmt.Errorf("database not found in snapshot %s: %w", dir, err)
}
return dir, latest.date, nil
+19 -81
View File
@@ -1,7 +1,6 @@
package bsdaily
import (
"context"
"errors"
"fmt"
"log/slog"
@@ -10,11 +9,6 @@ import (
"strings"
)
var (
errEmptyDecompressed = errors.New("decompressed content is empty")
errNotSQL = errors.New("decompressed content does not look like SQL")
)
// killCat terminates the zstdcat process, ignoring the benign case where it
// has already exited (e.g. after receiving SIGPIPE when head closed the pipe)
// and logging any other failure.
@@ -22,126 +16,70 @@ func killCat(cmd *exec.Cmd) {
if cmd.Process == nil {
return
}
err := cmd.Process.Kill()
if err != nil && !errors.Is(err, os.ErrProcessDone) {
if err := cmd.Process.Kill(); err != nil && !errors.Is(err, os.ErrProcessDone) {
slog.Warn("failed to kill zstdcat process", "error", err)
}
}
// VerifyOutput checks that the compressed file at path passes zstdmt's
// integrity test and that its first lines look like SQL.
func VerifyOutput(path string) error {
ctx := context.Background()
slog.Info("running zstdmt integrity check")
//nolint:gosec // path is an output file this package named
testCmd := exec.CommandContext(ctx, "zstdmt", "--test", path)
testCmd := exec.Command("zstdmt", "--test", path)
var testStderr strings.Builder
testCmd.Stderr = &testStderr
err := testCmd.Run()
if err != nil {
return fmt.Errorf("zstdmt --test failed: %w; stderr: %s",
err, testStderr.String())
if err := testCmd.Run(); err != nil {
return fmt.Errorf("zstdmt --test failed: %w; stderr: %s", err, testStderr.String())
}
slog.Info("zstdmt integrity check passed")
slog.Info("verifying SQL content")
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))
catCmd := exec.Command("zstdcat", path)
headCmd := exec.Command("head", fmt.Sprintf("-%d", verificationHeadLines))
pipe, err := catCmd.StdoutPipe()
if err != nil {
return "", fmt.Errorf("creating zstdcat pipe: %w", err)
return fmt.Errorf("creating zstdcat pipe: %w", err)
}
headCmd.Stdin = pipe
var headOut strings.Builder
headCmd.Stdout = &headOut
err = catCmd.Start()
if err != nil {
return "", fmt.Errorf("starting zstdcat: %w", err)
if err := catCmd.Start(); err != nil {
return fmt.Errorf("starting zstdcat: %w", err)
}
err = headCmd.Start()
if err != nil {
if err := headCmd.Start(); err != nil {
killCat(catCmd) // Clean up if head fails to start
return "", fmt.Errorf("starting head: %w", err)
return fmt.Errorf("starting head: %w", err)
}
// Wait for head first (it will exit when it has enough lines)
err = headCmd.Wait()
if err != nil {
if err := headCmd.Wait(); err != nil {
killCat(catCmd)
return "", fmt.Errorf("head command failed: %w", err)
return fmt.Errorf("head command failed: %w", err)
}
// Kill zstdcat since head closed the pipe (expected SIGPIPE)
killCat(catCmd)
_ = catCmd.Wait() // Reap the process
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 {
content := headOut.String()
if len(content) == 0 {
return errEmptyDecompressed
return fmt.Errorf("decompressed content is empty")
}
hasSQLMarker := false
for _, marker := range []string{
"BEGIN TRANSACTION", "CREATE TABLE", "INSERT INTO", "PRAGMA",
} {
for _, marker := range []string{"BEGIN TRANSACTION", "CREATE TABLE", "INSERT INTO", "PRAGMA"} {
if strings.Contains(content, marker) {
hasSQLMarker = true
break
}
}
const verificationSampleBytes = 200
if !hasSQLMarker {
return fmt.Errorf("%w; first %d bytes: %s", errNotSQL,
verificationSampleBytes,
content[:min(verificationSampleBytes, len(content))])
return fmt.Errorf("decompressed content does not look like SQL; first %d bytes: %s",
verificationSampleBytes, content[:min(verificationSampleBytes, len(content))])
}
slog.Info("SQL content verification passed")
return nil
}
+5
View File
@@ -0,0 +1,5 @@
{
"devDependencies": {
"prettier": "3.8.1"
}
}
+79 -1
View File
@@ -3,11 +3,20 @@
# this repo. Idempotent: every install is guarded by a check so already
# installed tools are skipped. Base tooling comes from nix, apt, brew,
# or apk (detected in that order); assumes NOTHING is present (not git,
# make, or go).
# make, go, or node). Node is used directly if installed; otherwise it
# is installed at a pinned version via nvm (installing nvm itself first,
# from a hash-verified release archive, never curl | sh).
set -eu
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=""
SUDO=""
APT_UPDATED=""
@@ -56,6 +65,69 @@ missing() {
! 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() {
cd "$ROOT"
@@ -68,6 +140,12 @@ main() {
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"
}
+23 -1
View File
@@ -1,12 +1,34 @@
#!/bin/sh
# script/fmt: format all files (writes).
# script/fmt: format all files (writes): the Go code with go fmt and
# every Markdown file with prettier.
set -eu
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() {
cd "$ROOT"
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 "$@"
+23 -1
View File
@@ -1,10 +1,30 @@
#!/bin/sh
# script/fmt-check: check formatting (read-only). Same scope as
# script/fmt, but fails instead of writing.
# script/fmt, the Go code and every Markdown file, but fails instead of
# writing.
set -eu
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() {
cd "$ROOT"
unformatted="$(gofmt -l .)"
@@ -13,6 +33,8 @@ main() {
echo "$unformatted" >&2
exit 1
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 "$@"
+8
View File
@@ -0,0 +1,8 @@
# 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==