Files
bsdaily/internal/bsdaily/extract.go
T
clawbot c16f177575
check / check (push) Successful in 5m14s
Adopt the shared .golangci.yml and fix the code to it (closes #6)
Vendor .golangci.yml byte-identical from sneak/prompts at cc440118 and
move the Dockerfile lint phase to golangci-lint v2.14.0 by the digest
REPO_POLICIES.md names. Fix the code to that config with flags, help
text, output files, SQL and the order of steps unchanged; long
functions are split into named steps.

Judgement call: the extraction transaction is now rolled back on every
early return; the old deferred rollback missed most failures and could
dereference a nil transaction.
Thirteen //nolint directives (gosec, unconvert, mnd, unqueryvet), each
with its reason.

Model: opus-5-5
2026-10-06 14:24:48 +02:00

313 lines
8.6 KiB
Go

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"
slog.Info("extracting day", "from", dayStart, "until", dayEnd)
// 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)
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)
}
}()
// Attach source database
_, err = db.ExecContext(ctx, "ATTACH DATABASE ? AS src", srcDBPath)
if 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")
if err != nil {
return fmt.Errorf("reading source schema: %w", err)
}
defer func() {
cerr := rows.Close()
if cerr != nil {
slog.Warn("failed to close schema rows", "error", cerr)
}
}()
var ddlStatements []string
for rows.Next() {
var ddl string
err = rows.Scan(&ddl)
if err != nil {
return fmt.Errorf("scanning DDL: %w", err)
}
ddlStatements = append(ddlStatements, ddl)
}
err = rows.Err()
if err != nil {
return fmt.Errorf("iterating DDL rows: %w", err)
}
for _, ddl := range ddlStatements {
_, err = db.ExecContext(ctx, ddl)
if err != nil {
return fmt.Errorf("creating table: %w\nDDL: %s", err, ddl)
}
}
return nil
}
// insertDayRows copies the target day's posts, and the rows they refer
// to, from the source database in one transaction. It returns the number
// of posts copied, or ErrNoPosts when there are none.
func insertDayRows(ctx context.Context, db *sql.DB, targetDay time.Time,
dayStart, dayEnd string,
) (int64, error) {
// Begin transaction for bulk inserts
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return 0, fmt.Errorf("beginning transaction: %w", err)
}
// After a successful Commit this does nothing and returns ErrTxDone.
defer func() {
rerr := tx.Rollback()
if rerr != nil && !errors.Is(rerr, sql.ErrTxDone) {
slog.Warn("failed to roll back transaction", "error", rerr)
}
}()
// Insert target day's data
slog.Info("inserting posts for target day")
//nolint:unqueryvet // copies whole rows; the source defines the columns
result, err := tx.ExecContext(ctx, "INSERT INTO posts SELECT * FROM src.posts "+
"WHERE timestamp >= ? AND timestamp < ?", dayStart, dayEnd)
if err != nil {
return 0, fmt.Errorf("inserting posts: %w", err)
}
postCount, _ := result.RowsAffected()
slog.Info("inserted posts", "count", postCount)
if postCount == 0 {
return 0, fmt.Errorf("%w %s - aborting to avoid producing empty output",
ErrNoPosts, targetDay.Format("2006-01-02"))
}
slog.Info("inserting junction and lookup tables")
err = insertRelatedRows(ctx, tx)
if err != nil {
return 0, err
}
copyMediaRows(ctx, tx)
// Commit the transaction before any further database operations
err = tx.Commit()
if err != nil {
return 0, fmt.Errorf("committing transaction: %w", err)
}
return postCount, nil
}
// insertRelatedRows copies the hashtag and URL links of the posts already
// inserted, the hashtags and URLs they link to, and the posting users.
//
//nolint:unqueryvet // copies whole rows; the source defines the columns
func insertRelatedRows(ctx context.Context, tx *sql.Tx) error {
_, err := tx.ExecContext(ctx, "INSERT INTO posts_hashtags "+
"SELECT * FROM src.posts_hashtags WHERE post_id IN (SELECT id FROM posts)")
if err != nil {
return fmt.Errorf("inserting posts_hashtags: %w", err)
}
_, err = tx.ExecContext(ctx, "INSERT INTO posts_urls "+
"SELECT * FROM src.posts_urls WHERE post_id IN (SELECT id FROM posts)")
if err != nil {
return fmt.Errorf("inserting posts_urls: %w", err)
}
_, err = tx.ExecContext(ctx, "INSERT INTO hashtags "+
"SELECT * FROM src.hashtags "+
"WHERE id IN (SELECT hashtag_id FROM posts_hashtags)")
if err != nil {
return fmt.Errorf("inserting hashtags: %w", err)
}
_, err = tx.ExecContext(ctx, "INSERT INTO urls "+
"SELECT * FROM src.urls WHERE id IN (SELECT url_id FROM posts_urls)")
if err != nil {
return fmt.Errorf("inserting urls: %w", err)
}
_, err = tx.ExecContext(ctx, "INSERT INTO users "+
"SELECT * FROM src.users WHERE did IN (SELECT user_did FROM posts)")
if err != nil {
return fmt.Errorf("inserting users: %w", err)
}
return nil
}
// copyMediaRows copies the media rows of the posts already inserted, when
// the source has a media table. Failures are logged, not returned.
func copyMediaRows(ctx context.Context, tx *sql.Tx) {
// Check if media table exists in source and copy if present
var mediaTableExists int
err := tx.QueryRowContext(ctx, "SELECT COUNT(*) FROM src.sqlite_master "+
"WHERE type='table' AND name='media'").Scan(&mediaTableExists)
if err != nil {
slog.Warn("checking for media table", "error", err)
} else if mediaTableExists > 0 {
slog.Info("inserting media entries")
// Get post blob_cids for this day's posts
//nolint:unqueryvet // copies whole rows; the source defines the columns
_, err = tx.ExecContext(ctx, "INSERT INTO media SELECT * FROM src.media "+
"WHERE content_hash IN "+
"(SELECT blob_cids FROM posts WHERE blob_cids IS NOT NULL)")
if err != nil {
slog.Warn("inserting media (may not have matching entries)",
"error", err)
}
}
}
// createIndexes creates every index of the attached source database in
// the destination database. Doing this after the bulk insert is faster
// than inserting into indexed tables.
func createIndexes(ctx context.Context, db *sql.DB) error {
// Create indexes after bulk insert for speed
slog.Info("creating indexes")
idxRows, err := db.QueryContext(ctx, "SELECT sql FROM src.sqlite_master "+
"WHERE type='index' AND name NOT LIKE 'sqlite_%' AND sql IS NOT NULL "+
"ORDER BY name")
if err != nil {
return fmt.Errorf("reading source indexes: %w", err)
}
defer func() {
cerr := idxRows.Close()
if cerr != nil {
slog.Warn("failed to close index rows", "error", cerr)
}
}()
var idxStatements []string
for idxRows.Next() {
var idxSQL string
err = idxRows.Scan(&idxSQL)
if err != nil {
return fmt.Errorf("scanning index DDL: %w", err)
}
idxStatements = append(idxStatements, idxSQL)
}
err = idxRows.Err()
if err != nil {
return fmt.Errorf("iterating index rows: %w", err)
}
for _, idxSQL := range idxStatements {
_, err = db.ExecContext(ctx, idxSQL)
if err != nil {
return fmt.Errorf("creating index: %w\nDDL: %s", err, idxSQL)
}
}
return nil
}