textmachine/platform/internal/pgstore/exports.go

272 lines
10 KiB
Go

package pgstore
import (
"context"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgtype"
)
// exports.go: the storage of the export door (canon §createExport/§getExport). The SQL itself is
// generated from `queries/exports.sql`; what lives here is the domain shape and the two rules that
// are not expressible in one statement.
// ExportState is the canon's four-state ladder, and it is a state rather than a boolean for the
// reason the canon gives: a boolean merges three situations into "not ready" and a poll on it never
// ends.
type ExportState string
const (
// ExportPending — the build is running or queued.
ExportPending ExportState = "pending"
// ExportReady — the artifact exists and the link works.
ExportReady ExportState = "ready"
// ExportFailed — the build ended in an error; `FailureCode` says which.
ExportFailed ExportState = "failed"
// ExportExpired — it was built and the link has lapsed. The file is gone with it.
ExportExpired ExportState = "expired"
)
// Export is one requested copy of a book.
type Export struct {
ID string
BookID string
Format string
State ExportState
// Path is where the artifact lives while it exists. SERVER TOPOLOGY: it never crosses the wire,
// the same rule that keeps the engine's own `BuildReport.files` off it.
Path string
// SizeBytes and Complete are nil until the build answers. Not zero-and-false: "the file has a
// hole" and "nobody has looked yet" are different facts, and the second is what `pending` means.
SizeBytes *int64
Complete *bool
// FailureCode is set only in ExportFailed.
FailureCode string
RequestedAt time.Time
FinishedAt *time.Time
ExpiresAt *time.Time
// Revision is the BOOK's revision at the moment this row was read — the number every
// book-scoped response carries. Read in the same statement as the row, so the two cannot be a
// transaction apart.
Revision int64
}
// ErrNoExport is an export that does not exist, or belongs to another account. ONE error for both,
// deliberately: telling them apart on the wire would answer "whose book is this" to a caller who is
// not its owner.
var ErrNoExport = errors.New("pgstore: no such export")
// ExportBuild is what the build worker needs: the row's own identity and the book's project
// directory. No owner — the worker is not a caller, and a job re-driven after a restart has no
// session behind it.
type ExportBuild struct {
ID string
BookID string
Format string
State ExportState
Workdir string
// HasTree says the book has been cut into chapters. A book without one has nothing to build from
// and must not reach the engine: see exports.Service.Build for why the question is answered here
// rather than at the door.
HasTree bool
}
// CreateExport records a requested export AND the queue entry that builds it in ONE transaction, so
// the two commit together or not at all — the same rule the admission of a run holds, for a weaker
// but real reason: a row with no job is an export that stays `pending` until the stale sweep fails
// it, which is a poll that ends on the wrong answer, and a job with no row is a worker that reads
// nothing.
//
// The book's revision is read inside the same transaction: the canon stamps every book-scoped
// response with it, and a number fetched afterwards would be from a different instant than the row.
func (s *Store) CreateExport(ctx context.Context, bookID, format string, now time.Time,
enqueue func(context.Context, Tx, string) error) (Export, error) {
out := Export{BookID: bookID, Format: format, State: ExportPending, RequestedAt: now}
err := s.inTx(ctx, func(tx pgx.Tx) error {
out.ID = newID("exp")
if err := s.q.WithTx(tx).InsertExport(ctx, InsertExportParams{
ID: out.ID, BookID: bookID, Format: format, RequestedAt: now,
}); err != nil {
return err
}
if err := tx.QueryRow(ctx, `select revision from books where id = $1`, bookID).
Scan(&out.Revision); err != nil {
return err
}
return enqueue(ctx, tx, out.ID)
})
if err != nil {
return Export{}, fmt.Errorf("pgstore: create export: %w", err)
}
return out, nil
}
// ReadExport is the poll, and the ownership check is inside the statement rather than beside it.
func (s *Store) ReadExport(ctx context.Context, userID, bookID, exportID string) (Export, error) {
row, err := s.q.ReadExportForOwner(ctx, ReadExportForOwnerParams{
ID: exportID, BookID: bookID, OwnerID: userID,
})
if errors.Is(err, pgx.ErrNoRows) {
return Export{}, ErrNoExport
}
if err != nil {
return Export{}, fmt.Errorf("pgstore: read export: %w", err)
}
out := Export{
ID: row.ID, BookID: row.BookID, Format: row.Format, State: ExportState(row.State),
RequestedAt: row.RequestedAt, FinishedAt: row.FinishedAt, ExpiresAt: row.ExpiresAt,
Revision: row.Revision,
}
if row.Path != nil {
out.Path = *row.Path
}
if row.FailureCode != nil {
out.FailureCode = *row.FailureCode
}
if row.SizeBytes.Valid {
n := row.SizeBytes.Int64
out.SizeBytes = &n
}
if row.Complete.Valid {
b := row.Complete.Bool
out.Complete = &b
}
return out, nil
}
// ClaimExportBuild stamps the moment a worker picked a build up, and says whether it may build.
//
// The stamp is what lets the GC tell a build whose process is gone from one that is merely waiting
// its turn behind three parses. The guards are what make the claim a claim: `state = 'pending'`
// stops a worker building a row the stale sweep has already answered for, and `started_at is null`
// stops a SECOND worker joining one that is already building — both would write to the same path,
// and the loser's cleanup would delete what the winner published.
func (s *Store) ClaimExportBuild(ctx context.Context, exportID string, now time.Time) error {
n, err := s.q.StampExportStart(ctx, StampExportStartParams{ID: exportID, Now: &now})
if err != nil {
return fmt.Errorf("pgstore: claim export build: %w", err)
}
if n == 0 {
return ErrExportSettled
}
return nil
}
// ReadExportForBuild is the worker's read of the job it was handed.
func (s *Store) ReadExportForBuild(ctx context.Context, exportID string) (ExportBuild, error) {
row, err := s.q.ReadExportForBuild(ctx, exportID)
if errors.Is(err, pgx.ErrNoRows) {
return ExportBuild{}, ErrNoExport
}
if err != nil {
return ExportBuild{}, fmt.Errorf("pgstore: read export for build: %w", err)
}
return ExportBuild{ID: row.ID, BookID: row.BookID, Format: row.Format,
State: ExportState(row.State), Workdir: row.Workdir, HasTree: row.HasTree}, nil
}
// ErrExportSettled is a write against an export that is no longer `pending`: the stale sweep failed
// it, or a duplicate worker got there first. Returned rather than swallowed because the CALLER is
// then holding a built file nobody will serve, and deleting it is its business, not this layer's.
var ErrExportSettled = errors.New("pgstore: the export is no longer pending")
// FinishExport publishes a built artifact.
func (s *Store) FinishExport(ctx context.Context, exportID, path string, size int64, complete bool,
now, expiresAt time.Time) error {
n, err := s.q.FinishExport(ctx, FinishExportParams{
ID: exportID, Path: &path,
SizeBytes: pgtype.Int8{Int64: size, Valid: true},
Complete: pgtype.Bool{Bool: complete, Valid: true},
FinishedAt: &now, ExpiresAt: &expiresAt,
})
if err != nil {
return fmt.Errorf("pgstore: finish export: %w", err)
}
if n == 0 {
return ErrExportSettled
}
return nil
}
// FailExport records why a build did not produce a file. Idempotent against the stale sweep by the
// same `pending` guard as FinishExport: the first verdict stands.
func (s *Store) FailExport(ctx context.Context, exportID, failureCode string, now time.Time) error {
n, err := s.q.FailExport(ctx, FailExportParams{
ID: exportID, FailureCode: &failureCode, FinishedAt: &now,
})
if err != nil {
return fmt.Errorf("pgstore: fail export: %w", err)
}
if n == 0 {
return ErrExportSettled
}
return nil
}
// ExportSweepBatch bounds one pass of either GC question. Small because the work behind each row is
// a file deletion and the sweep shares its tick with every other reconciler.
const ExportSweepBatch = 200
// ExpireExports moves every lapsed artifact to `expired`. The bytes are removed afterwards, by
// Unlinked/ForgetExportPath — the row's move and the unlink are two acts, and the row keeps naming
// the file until the file is actually gone.
func (s *Store) ExpireExports(ctx context.Context, now time.Time) (int, error) {
rows, err := s.q.ExpireReadyExports(ctx, ExpireReadyExportsParams{Now: &now, Lim: ExportSweepBatch})
if err != nil {
return 0, fmt.Errorf("pgstore: expire exports: %w", err)
}
return len(rows), nil
}
// Artifact is one file the GC still owes an unlink.
type Artifact struct {
ID string
Path string
}
// Unlinked lists the artifacts whose row has lapsed and whose bytes are still on disk.
//
// It is the retry the one-shot version did not have: an unlink can fail — a full disk, a read-only
// mount, a crash between the two acts — and a path cleared in the same statement as the state left
// a file nothing could ever name again, because every other read of this table selects `ready`.
func (s *Store) Unlinked(ctx context.Context) ([]Artifact, error) {
rows, err := s.q.UnlinkedExports(ctx, ExportSweepBatch)
if err != nil {
return nil, fmt.Errorf("pgstore: list unlinked exports: %w", err)
}
out := make([]Artifact, 0, len(rows))
for _, r := range rows {
if r.Path != nil && *r.Path != "" {
out = append(out, Artifact{ID: r.ID, Path: *r.Path})
}
}
return out, nil
}
// ForgetExportPath stops a row naming a file, and is called ONLY after the file is gone.
func (s *Store) ForgetExportPath(ctx context.Context, exportID string) error {
if err := s.q.ForgetExportPath(ctx, exportID); err != nil {
return fmt.Errorf("pgstore: forget export path: %w", err)
}
return nil
}
// FailStaleExports ends the polls of builds nobody is coming back for, and returns the ids it
// ended so the caller can say so once rather than per row.
//
// TWO cutoffs: `startedCutoff` judges a build a worker HAS picked up, `queuedCutoff` one still
// waiting its turn. See the statement for why they cannot be one number.
func (s *Store) FailStaleExports(ctx context.Context, failureCode string,
startedCutoff, queuedCutoff, now time.Time) ([]string, error) {
ids, err := s.q.FailStalePendingExports(ctx, FailStalePendingExportsParams{
FailureCode: &failureCode, Now: &now,
StartedCutoff: &startedCutoff, QueuedCutoff: queuedCutoff, Lim: ExportSweepBatch,
})
if err != nil {
return nil, fmt.Errorf("pgstore: fail stale exports: %w", err)
}
return ids, nil
}