Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
699d980c92 |
@@ -10,9 +10,3 @@ insert_final_newline = true
|
||||
|
||||
[Makefile]
|
||||
indent_style = tab
|
||||
|
||||
[*.go]
|
||||
indent_style = tab
|
||||
|
||||
# This repository's own sections, such as one for another language it
|
||||
# uses, go below this comment, and a re-vendor keeps them.
|
||||
|
||||
+3
-6
@@ -27,7 +27,7 @@ node_modules/
|
||||
# Environment files. `*.env` covers bare `.env` and the `prod.env`
|
||||
# convention. Only the templates `example.env` and `sample.env` are
|
||||
# re-included below. A repository that commits any other template adds
|
||||
# its own negation at the end of this file, for example `!.env.example`.
|
||||
# its own negation after these lines, for example `!.env.example`.
|
||||
*.[eE][nN][vV]
|
||||
.[eE][nN][vV].*
|
||||
.[eE][nN][vV][rR][cC]
|
||||
@@ -46,11 +46,8 @@ node_modules/
|
||||
[iI][dD]_[eE][dD]25519
|
||||
[iI][dD]_[eE][dD]25519_[sS][kK]
|
||||
|
||||
# This repository's own entries, such as its build outputs, go below
|
||||
# this comment, and a re-vendor keeps them. Anchor a binary built at the
|
||||
# root: `/myapp`, never `myapp`, which also ignores `cmd/myapp/`.
|
||||
|
||||
# Go: logs, test binaries, coverage output, and the binary `make` builds.
|
||||
# Go: logs, test binaries, coverage output, and the binary `make` writes
|
||||
# at the repo root, anchored so it does not also match `cmd/bsdaily/`.
|
||||
*.log
|
||||
*.test
|
||||
*.out
|
||||
|
||||
@@ -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
|
||||
@@ -1,2 +0,0 @@
|
||||
node_modules/
|
||||
yarn.lock
|
||||
@@ -1,4 +0,0 @@
|
||||
{
|
||||
"tabWidth": 4,
|
||||
"proseWrap": "always"
|
||||
}
|
||||
+2
-2
@@ -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
|
||||
|
||||
|
||||
@@ -1,26 +1,28 @@
|
||||
# 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
|
||||
|
||||
@@ -31,14 +33,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 `make fmt` (`go fmt` for Go, prettier for
|
||||
Markdown), 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 a standard `go fmt`, 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.
|
||||
@@ -48,55 +50,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, 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/bootstrap` — install all development dependencies (go, Go
|
||||
module download); 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` 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`
|
||||
- `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`
|
||||
|
||||
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
|
||||
|
||||
@@ -106,63 +108,63 @@ 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
|
||||
@@ -171,11 +173,11 @@ A single run proceeds as follows:
|
||||
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;
|
||||
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
|
||||
@@ -217,9 +219,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
|
||||
@@ -239,21 +241,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
|
||||
|
||||
@@ -263,17 +265,16 @@ 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/)
|
||||
|
||||
|
||||
+11
-23
@@ -1,6 +1,6 @@
|
||||
---
|
||||
title: Repository Policies
|
||||
last_modified: 2026-10-06
|
||||
last_modified: 2026-10-04
|
||||
---
|
||||
|
||||
This document covers repository structure, tooling, and workflow standards. Code
|
||||
@@ -118,9 +118,8 @@ style conventions are in separate documents:
|
||||
and nothing else:
|
||||
|
||||
```sh
|
||||
tag="$(script/projectname)"
|
||||
docker build --no-cache --target lint -t "$tag-lint" .
|
||||
docker build --no-cache --target test -t "$tag-test" .
|
||||
docker build --no-cache --target lint -t "$(script/projectname)-lint" .
|
||||
docker build --no-cache --target test -t "$(script/projectname)-test" .
|
||||
```
|
||||
|
||||
**A stage that is not the last one in the file is built only when the final
|
||||
@@ -133,9 +132,7 @@ style conventions are in separate documents:
|
||||
**Every `docker build` in `script/` is tagged**, here and in
|
||||
`script/cibuild` and `script/docker`. An untagged build leaves a dangling
|
||||
image behind on every invocation, on every developer host and every CI
|
||||
runner; a tagged one replaces the previous image. Each script assigns the
|
||||
tag on its own line before the build, so `set -e` stops it where
|
||||
`script/projectname` fails.
|
||||
runner; a tagged one replaces the previous image.
|
||||
|
||||
Inside a phase the tool is invoked directly — `golangci-lint`, `go test`,
|
||||
`eslint`, `prettier` — never through `make lint` or `script/test`, which are
|
||||
@@ -387,12 +384,9 @@ style conventions are in separate documents:
|
||||
|
||||
- `.gitignore` should be comprehensive from the start: OS files (`.DS_Store`),
|
||||
editor files (`.swp`, `*~`), in-repo agent scratch directories (`.claude/`),
|
||||
`node_modules/`, and the repo's own build outputs. Fetch the standard
|
||||
`.gitignore` from
|
||||
`https://git.eeqj.de/sneak/prompts/raw/branch/main/.gitignore` when setting up
|
||||
a new repo. A repo's `.gitignore` is the standard file followed by the repo's
|
||||
own entries, such as its binaries; a re-vendor replaces the standard part and
|
||||
keeps those entries. These patterns are written to `.gitignore`'s own
|
||||
language build artifacts, and `node_modules/`. Fetch the standard `.gitignore`
|
||||
from `https://git.eeqj.de/sneak/prompts/raw/branch/main/.gitignore` when
|
||||
setting up a new repo. These patterns are written to `.gitignore`'s own
|
||||
semantics, in which an unanchored pattern already matches at every depth; they
|
||||
are not a `.dockerignore` and must not be transplanted into one unmodified.
|
||||
|
||||
@@ -440,15 +434,13 @@ style conventions are in separate documents:
|
||||
byte-identically across repos:
|
||||
|
||||
```sh
|
||||
# The version and the tag each get their own line: a failing command
|
||||
# substitution inside an argument does not trip `set -e`, so the inline
|
||||
# form degrades to an empty constant.
|
||||
# Own line: a failing command substitution inside an argument does not
|
||||
# trip `set -e`, so the inline form degrades to an empty constant.
|
||||
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
||||
[ -n "$version" ] || version="unknown"
|
||||
tag="$(script/projectname)"
|
||||
docker build --no-cache \
|
||||
--build-arg VERSION="$version" \
|
||||
-t "$tag" .
|
||||
-t "$(script/projectname)" .
|
||||
```
|
||||
|
||||
`--always` makes an untagged repo yield an abbreviated commit hash rather
|
||||
@@ -644,11 +636,7 @@ style conventions are in separate documents:
|
||||
Never edit existing migrations after release.
|
||||
|
||||
- All repos should have an `.editorconfig` enforcing the project's indentation
|
||||
settings: the standard file from
|
||||
`https://git.eeqj.de/sneak/prompts/raw/branch/main/.editorconfig`, which sets
|
||||
tabs for `Makefile` and Go files, followed by the repo's own sections, such as
|
||||
one for another language it uses. A re-vendor replaces the standard part and
|
||||
keeps those sections.
|
||||
settings.
|
||||
|
||||
- Avoid putting files in the repo root unless necessary. Root should contain
|
||||
only project-level config files (`README.md`, `AGENTS.md`, `Makefile`,
|
||||
|
||||
@@ -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,38 +14,35 @@ pre-1.0
|
||||
|
||||
# Next Step
|
||||
|
||||
Expand tests beyond the compilation smoke test: unit tests for the extraction,
|
||||
verification, and atomic-publish paths.
|
||||
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
|
||||
|
||||
- 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.
|
||||
|
||||
+42
-82
@@ -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,112 +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("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
|
||||
}
|
||||
|
||||
err = bsdaily.Run(targetDates)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
slog.Info("completed successfully")
|
||||
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
rootCmd.Flags().StringVarP(&dateFlag, "date", "d", "",
|
||||
"target date to extract (YYYY-MM-DD); "+
|
||||
"defaults to snapshot date minus one day")
|
||||
rootCmd.Flags().StringVar(&fromFlag, "from", "",
|
||||
"start of date range to extract (YYYY-MM-DD, inclusive); use with --to")
|
||||
rootCmd.Flags().StringVar(&toFlag, "to", "",
|
||||
"end of date range to extract (YYYY-MM-DD, inclusive); use with --from")
|
||||
|
||||
err := rootCmd.Execute()
|
||||
if err != nil {
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
// parseTargetDates turns the --date, --from and --to flags into the days
|
||||
// to extract. It returns nil when none of them is set, which Run takes
|
||||
// to mean the snapshot date minus one day.
|
||||
func parseTargetDates(dateFlag, fromFlag, toFlag string) ([]time.Time, error) {
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
hasDate := dateFlag != ""
|
||||
hasFrom := fromFlag != ""
|
||||
hasTo := toFlag != ""
|
||||
|
||||
// Validate mutual exclusivity
|
||||
if hasDate && (hasFrom || hasTo) {
|
||||
return nil, errDateExclusive
|
||||
return fmt.Errorf("--date and --from/--to are mutually exclusive")
|
||||
}
|
||||
|
||||
if hasFrom != hasTo {
|
||||
if hasFrom {
|
||||
return nil, errFromRequiresTo
|
||||
return fmt.Errorf("--from requires --to")
|
||||
}
|
||||
|
||||
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)
|
||||
return fmt.Errorf("--to requires --from")
|
||||
}
|
||||
|
||||
var targetDates []time.Time
|
||||
|
||||
if hasDate {
|
||||
t, err := time.Parse("2006-01-02", dateFlag)
|
||||
if err != nil {
|
||||
return fmt.Errorf("invalid --date %q (expected YYYY-MM-DD): %w", dateFlag, err)
|
||||
}
|
||||
targetDates = []time.Time{t}
|
||||
} else if hasFrom {
|
||||
from, err := time.Parse("2006-01-02", fromFlag)
|
||||
if err != nil {
|
||||
return fmt.Errorf("invalid --from %q (expected YYYY-MM-DD): %w", fromFlag, err)
|
||||
}
|
||||
to, err := time.Parse("2006-01-02", toFlag)
|
||||
if err != nil {
|
||||
return fmt.Errorf("invalid --to %q (expected YYYY-MM-DD): %w", toFlag, err)
|
||||
}
|
||||
if from.After(to) {
|
||||
return fmt.Errorf("--from %s is after --to %s", fromFlag, toFlag)
|
||||
}
|
||||
for d := from; !d.After(to); d = d.AddDate(0, 0, 1) {
|
||||
targetDates = append(targetDates, d)
|
||||
}
|
||||
}
|
||||
// else: targetDates remains nil → Run() defaults to snapshot date minus one
|
||||
|
||||
return targetDates, nil
|
||||
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")
|
||||
|
||||
if err := rootCmd.Execute(); err != nil {
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,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}$`)
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
package bsdaily
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
@@ -10,30 +9,20 @@ import (
|
||||
)
|
||||
|
||||
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)
|
||||
|
||||
//nolint:gosec // src is a path this package built
|
||||
srcFile, err := os.Open(src)
|
||||
if err != nil {
|
||||
return fmt.Errorf("opening source %s: %w", src, err)
|
||||
}
|
||||
|
||||
defer func() {
|
||||
cerr := srcFile.Close()
|
||||
if cerr != nil {
|
||||
if cerr := srcFile.Close(); cerr != nil {
|
||||
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())
|
||||
}
|
||||
|
||||
//nolint:gosec // dst is a path this package built
|
||||
dstFile, err := os.Create(dst)
|
||||
if err != nil {
|
||||
return fmt.Errorf("creating destination %s: %w", dst, err)
|
||||
}
|
||||
|
||||
defer func() {
|
||||
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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -1,39 +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)
|
||||
}
|
||||
|
||||
//nolint:gosec,unconvert // Bsize is never negative; Bavail is signed on FreeBSD
|
||||
free := uint64(stat.Bavail) * uint64(stat.Bsize)
|
||||
freeGB := float64(free) / float64(bytesPerGB)
|
||||
minGB := float64(minBytes) / float64(bytesPerGB)
|
||||
|
||||
slog.Info("disk space check", "label", label, "path", path,
|
||||
"free_gb", fmt.Sprintf("%.1f", freeGB),
|
||||
"required_gb", fmt.Sprintf("%.1f", minGB))
|
||||
|
||||
if free < minBytes {
|
||||
return fmt.Errorf("%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
|
||||
}
|
||||
|
||||
+35
-74
@@ -1,8 +1,6 @@
|
||||
package bsdaily
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
@@ -11,44 +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
|
||||
}
|
||||
|
||||
//nolint:gosec // outputPath is a path this package built
|
||||
outFile, err := os.Create(outputPath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("creating output file: %w", err)
|
||||
}
|
||||
|
||||
defer func() {
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -56,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
|
||||
|
||||
+51
-178
@@ -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 = 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
|
||||
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) {
|
||||
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
|
||||
}
|
||||
|
||||
+43
-130
@@ -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) {
|
||||
for _, targetDay := range targetDates {
|
||||
dayStr := targetDay.Format("2006-01-02")
|
||||
|
||||
slog.Info("processing day", "date", dayStr)
|
||||
|
||||
// Check if output already exists
|
||||
outputDir := filepath.Join(DailiesBase, targetDay.Format("2006-01"))
|
||||
outputFinal := filepath.Join(outputDir, dayStr+".sql.zst")
|
||||
|
||||
_, err := os.Stat(outputFinal)
|
||||
if err == nil {
|
||||
if _, err := os.Stat(outputFinal); err == nil {
|
||||
slog.Info("output already exists, skipping", "path", outputFinal)
|
||||
|
||||
return false, nil
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
|
||||
// Extract target day into a per-day database
|
||||
extractedDB := filepath.Join(tmpDir, "extracted-"+dayStr+".db")
|
||||
|
||||
slog.Info("extracting target day", "src", dstDB, "dst", extractedDB)
|
||||
|
||||
err = ExtractDay(dstDB, extractedDB, targetDay)
|
||||
if err != nil {
|
||||
if err := ExtractDay(dstDB, extractedDB, targetDay); err != nil {
|
||||
if errors.Is(err, ErrNoPosts) {
|
||||
slog.Warn("no posts found, skipping day", "date", dayStr)
|
||||
cleanup(extractedDB)
|
||||
|
||||
return false, nil
|
||||
skipped++
|
||||
continue
|
||||
}
|
||||
|
||||
return false, fmt.Errorf("extracting day %s: %w", dayStr, err)
|
||||
return fmt.Errorf("extracting day %s: %w", dayStr, err)
|
||||
}
|
||||
|
||||
// Dump to SQL and compress
|
||||
//nolint:gosec,mnd // world-readable on purpose: the dailies tree is published
|
||||
err = os.MkdirAll(outputDir, 0755)
|
||||
if err != nil {
|
||||
if err := os.MkdirAll(outputDir, 0755); err != nil {
|
||||
cleanup(extractedDB)
|
||||
|
||||
return false, fmt.Errorf("creating output directory %s: %w", outputDir, err)
|
||||
return fmt.Errorf("creating output directory %s: %w", outputDir, err)
|
||||
}
|
||||
|
||||
outputTmp := filepath.Join(outputDir, "."+dayStr+".sql.zst.tmp")
|
||||
|
||||
slog.Info("dumping and compressing", "tmp_output", outputTmp)
|
||||
|
||||
err = DumpAndCompress(extractedDB, outputTmp)
|
||||
if err != nil {
|
||||
if err := DumpAndCompress(extractedDB, outputTmp); err != nil {
|
||||
cleanup(outputTmp)
|
||||
cleanup(extractedDB)
|
||||
|
||||
return false, fmt.Errorf("dump and compress for %s: %w", dayStr, err)
|
||||
return fmt.Errorf("dump and compress for %s: %w", dayStr, err)
|
||||
}
|
||||
|
||||
slog.Info("verifying compressed output")
|
||||
|
||||
err = VerifyOutput(outputTmp)
|
||||
if err != nil {
|
||||
if err := VerifyOutput(outputTmp); err != nil {
|
||||
cleanup(outputTmp)
|
||||
cleanup(extractedDB)
|
||||
|
||||
return false, fmt.Errorf("verification failed for %s: %w", dayStr, err)
|
||||
return 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 {
|
||||
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())
|
||||
|
||||
slog.Info("day completed", "date", dayStr, "path", outputFinal,
|
||||
"size_bytes", info.Size())
|
||||
// Remove extracted DB to reclaim space immediately
|
||||
cleanup(extractedDB)
|
||||
processed++
|
||||
}
|
||||
|
||||
slog.Info("run summary", "processed", processed, "skipped", skipped, "total", len(targetDates))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -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
@@ -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
|
||||
}
|
||||
|
||||
@@ -1,5 +0,0 @@
|
||||
{
|
||||
"devDependencies": {
|
||||
"prettier": "3.8.1"
|
||||
}
|
||||
}
|
||||
+2
-88
@@ -3,23 +3,13 @@
|
||||
# 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, 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).
|
||||
# make, or go).
|
||||
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=""
|
||||
|
||||
detect_pkgmgr() {
|
||||
[ -n "$PKGMGR" ] && return 0
|
||||
@@ -48,14 +38,7 @@ pkg_install() {
|
||||
detect_pkgmgr
|
||||
case "$PKGMGR" in
|
||||
nix) nix-env -iA "nixpkgs.$1" ;;
|
||||
apt)
|
||||
# Package lists may be empty (fresh images); refresh once per run.
|
||||
if [ -z "$APT_UPDATED" ]; then
|
||||
$SUDO env DEBIAN_FRONTEND=noninteractive apt-get update
|
||||
APT_UPDATED=1
|
||||
fi
|
||||
$SUDO env DEBIAN_FRONTEND=noninteractive apt-get install -y "$2"
|
||||
;;
|
||||
apt) $SUDO env DEBIAN_FRONTEND=noninteractive apt-get install -y "$2" ;;
|
||||
brew) brew install "$3" ;;
|
||||
apk) apk add --no-cache "$4" ;;
|
||||
esac
|
||||
@@ -65,69 +48,6 @@ 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"
|
||||
|
||||
@@ -140,12 +60,6 @@ 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"
|
||||
}
|
||||
|
||||
|
||||
+5
-7
@@ -14,17 +14,15 @@ main() {
|
||||
cd "$ROOT"
|
||||
"$SCRIPT_DIR/bootstrap"
|
||||
"$SCRIPT_DIR/check"
|
||||
# The version and the tag each get their own line: a failing
|
||||
# command substitution inside an argument does not trip `set -e`,
|
||||
# so the inline form degrades silently to an empty constant. The
|
||||
# VERSION build argument takes precedence over the version a build
|
||||
# stage derives from the .git in the context.
|
||||
# Own line: a failing command substitution inside an argument does
|
||||
# not trip `set -e`, so the inline form degrades silently to an
|
||||
# empty constant. The VERSION build argument takes precedence over
|
||||
# the version a build stage derives from the .git in the context.
|
||||
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
||||
[ -n "$version" ] || version="unknown"
|
||||
tag="$("$SCRIPT_DIR/projectname")"
|
||||
docker build --no-cache \
|
||||
--build-arg VERSION="$version" \
|
||||
-t "$tag" .
|
||||
-t "$("$SCRIPT_DIR/projectname")" .
|
||||
}
|
||||
|
||||
main "$@"
|
||||
|
||||
+5
-7
@@ -10,17 +10,15 @@ ROOT="$(cd "$SCRIPT_DIR/.." && pwd -P)"
|
||||
|
||||
main() {
|
||||
cd "$ROOT"
|
||||
# The version and the tag each get their own line: a failing
|
||||
# command substitution inside an argument does not trip `set -e`,
|
||||
# so the inline form degrades silently to an empty constant. The
|
||||
# VERSION build argument takes precedence over the version a build
|
||||
# stage derives from the .git in the context.
|
||||
# Own line: a failing command substitution inside an argument does
|
||||
# not trip `set -e`, so the inline form degrades silently to an
|
||||
# empty constant. The VERSION build argument takes precedence over
|
||||
# the version a build stage derives from the .git in the context.
|
||||
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
|
||||
[ -n "$version" ] || version="unknown"
|
||||
tag="$("$SCRIPT_DIR/projectname")"
|
||||
docker build --no-cache \
|
||||
--build-arg VERSION="$version" \
|
||||
-t "$tag" .
|
||||
-t "$("$SCRIPT_DIR/projectname")" .
|
||||
}
|
||||
|
||||
main "$@"
|
||||
|
||||
+1
-23
@@ -1,34 +1,12 @@
|
||||
#!/bin/sh
|
||||
# script/fmt: format all files (writes): the Go code with go fmt and
|
||||
# every Markdown file with prettier.
|
||||
# script/fmt: format all files (writes).
|
||||
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 "$@"
|
||||
|
||||
+1
-23
@@ -1,30 +1,10 @@
|
||||
#!/bin/sh
|
||||
# script/fmt-check: check formatting (read-only). Same scope as
|
||||
# script/fmt, the Go code and every Markdown file, but fails instead of
|
||||
# writing.
|
||||
# script/fmt, 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 .)"
|
||||
@@ -33,8 +13,6 @@ 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 "$@"
|
||||
|
||||
+1
-5
@@ -15,13 +15,9 @@ ROOT="$(cd "$SCRIPT_DIR/.." && pwd -P)"
|
||||
|
||||
main() {
|
||||
cd "$ROOT"
|
||||
# The tag gets its own line: a failing command substitution inside
|
||||
# an argument does not trip `set -e`, so the inline form degrades
|
||||
# silently to an empty constant.
|
||||
tag="$("$SCRIPT_DIR/projectname")"
|
||||
docker build --no-cache \
|
||||
--target lint \
|
||||
-t "$tag-lint" .
|
||||
-t "$("$SCRIPT_DIR/projectname")-lint" .
|
||||
}
|
||||
|
||||
main "$@"
|
||||
|
||||
+1
-5
@@ -11,13 +11,9 @@ ROOT="$(cd "$SCRIPT_DIR/.." && pwd -P)"
|
||||
|
||||
main() {
|
||||
cd "$ROOT"
|
||||
# The tag gets its own line: a failing command substitution inside
|
||||
# an argument does not trip `set -e`, so the inline form degrades
|
||||
# silently to an empty constant.
|
||||
tag="$("$SCRIPT_DIR/projectname")"
|
||||
docker build --no-cache \
|
||||
--target test \
|
||||
-t "$tag-test" .
|
||||
-t "$("$SCRIPT_DIR/projectname")-test" .
|
||||
}
|
||||
|
||||
main "$@"
|
||||
|
||||
@@ -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==
|
||||
Reference in New Issue
Block a user