Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 27 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ happened, not just whether pg_restore was noisy.

> [!NOTE]
> Pre-1.0. The pipeline works end to end (live wrapper, error classification,
> exit-code normalization); an honest ETA is still to come.
> exit-code normalization, stdin restores); an honest ETA is still to come.

## Install

Expand Down Expand Up @@ -83,7 +83,8 @@ ele --replay <plan> <log> # replay a captured stderr log through the live view
```

`--plan` and `--replay` are offline: they never connect to a database. For both,
`<plan>` is either a dump or a saved `pg_restore -l` listing (see below).
`<plan>` is either a dump or a saved `pg_restore -l` listing (see below);
`--replay` also takes `-`, meaning no plan at all.

`--replay` is the safe dry-run. Feed it a captured stderr log and a plan, and it
reruns the whole pipeline: on a terminal it drives the real live block (progress,
Expand All @@ -96,11 +97,26 @@ Saving the listing once lets you dry-run (or re-`--plan`) without the dump on ha
```sh
pg_restore -l latest.dump > latest.toc # capture the plan once (no database)
ele --replay latest.toc ele-20260719.log # replay any captured log against it
ele --replay - ele-20260719.log # replay with no plan: the countless view
```

Tune the animation length with `ELE_REPLAY_SECONDS` (default 12; `0` feeds
instantly and just prints the summary).

### Restoring from stdin

```sh
cat latest.dump | ele -d myapp_dev --clean --no-owner
```

Preflight needs a file it can open, so a dump arriving on a pipe has no table of
contents and no per-phase totals. `ele` still runs the whole view: spinner,
current-object line, grouped errors, summary and exit code, with a per-phase
object count in place of each bar. Percentages and the skipped-object count are
gone, because both need a total.

`ele --replay - <log>` shows the same view offline, against a captured log.

## How It Works

- **Preflight** runs `pg_restore -l` (under `LC_ALL=C`) and parses the table of
Expand All @@ -115,7 +131,10 @@ instantly and just prints the summary).
through untouched.
- **Parse & aggregate** turn that stream into progress and grouped errors,
surviving the way parallel (`-j`) workers interleave their output. Errors are
fingerprinted (quoted identifiers blanked) and classified benign or real.
fingerprinted (quoted identifiers blanked) and classified benign or real. The
table of message wordings the parser matches on is generated from PostgreSQL's
own `pg_dump` message catalogs (`just gen-messages`), across release tags 13
to 18, so a new major is a regeneration rather than a rewrite.
- **Exit code** is normalized: if every error was benign, `ele` exits `0` with a
note; any real error exits nonzero. Set `ELE_STRICT_EXIT=1` to keep
pg_restore's raw code.
Expand All @@ -127,12 +146,16 @@ collision with pg_restore's own flag surface.

| Variable | Effect |
|---|---|
| `ELE_LOG=path` | Write the raw log here instead of `./ele-<timestamp>.log`. `/dev/null` discards it. |
| `ELE_STRICT_EXIT=1` | Don't normalize the exit code; return pg_restore's own. |
| `NO_COLOR=1` | Disable color. |
| `ELE_PLAIN=1` | Disable the repaint block; emit periodic one-line progress instead. |
| `ELE_PASSTHROUGH=1` | Don't wrap at all: raw `pg_restore` output and its own exit code. |

`ele` also drops to plain output automatically when stderr isn't a terminal, or
under `CI` / `CLAUDECODE`.
under `CI` / `CLAUDECODE`. Plain output is still the aggregated view - a summary
and periodic progress, never the firehose. `ELE_PASSTHROUGH=1` is the only way
back to raw `pg_restore`.

## Feedback

Expand Down
3 changes: 2 additions & 1 deletion cmd/ele/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,8 @@ func printUsage(w io.Writer) {
" ele --plan <plan> print the parsed restore plan and exit\n"+
" ele --replay <plan> <log> dry-run: replay a captured stderr log through the\n"+
" live view\n"+
" (<plan> is a dump or a saved pg_restore -l listing)\n")
" (<plan> is a dump or a saved pg_restore -l listing;\n"+
" '-' replays with no plan, as a stdin restore runs)\n")
}

func preflightOnly(out io.Writer, dumpPath string) error {
Expand Down
85 changes: 55 additions & 30 deletions internal/aggregator/aggregator.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,11 @@ type Aggregator struct {
cfg Config
now func() time.Time // injectable clock for item timing

// countless: no plan, so no denominators. Phases are still counted - the
// event's own description gives the section - they just have no total to
// count against. This is the stdin restore, where preflight is impossible.
countless bool

// phase denominators, from the plan
total map[toc.Section]int
// phase progress
Expand Down Expand Up @@ -70,39 +75,43 @@ type inflightItem struct {
start time.Time
}

// New builds an Aggregator for a restore plan and config.
// New builds an Aggregator for a restore plan and config. A nil plan is the
// countless restore: everything still aggregates, only the denominators are
// missing.
func New(plan *toc.RestorePlan, cfg Config) *Aggregator {
pre, data, post, _ := plan.PhaseCounts()
bytesTotal, byteSized := plan.DataBytes(), false
// DataBytes is only meaningful when preflight sized every data file; the
// caller signals that via a nonzero total together with ByteSized. Here we
// treat any positive total as byte-capable and let the renderer decide.
if bytesTotal > 0 {
byteSized = true
}
return &Aggregator{
plan: plan,
cfg: cfg,
now: time.Now,
total: map[toc.Section]int{
toc.PreData: pre, toc.Data: data, toc.PostData: post,
},
a := &Aggregator{
plan: plan,
cfg: cfg,
now: time.Now,
countless: plan == nil,
total: map[toc.Section]int{},
done: map[toc.Section]int{},
sectionDone: map[toc.Section]bool{},
doneIDs: map[int]bool{},
bytesTotal: bytesTotal,
byteSized: byteSized,
inflight: map[int]*inflightItem{},
groups: map[string]*errorGroup{},
}
if plan == nil {
return a
}

pre, data, post, _ := plan.PhaseCounts()
a.total = map[toc.Section]int{toc.PreData: pre, toc.Data: data, toc.PostData: post}
// DataBytes is only meaningful when preflight sized every data file; the
// caller signals that via a nonzero total together with ByteSized. Here we
// treat any positive total as byte-capable and let the renderer decide.
if bytesTotal := plan.DataBytes(); bytesTotal > 0 {
a.bytesTotal, a.byteSized = bytesTotal, true
}
return a
}

// Feed folds one event into the running state.
func (a *Aggregator) Feed(ev parser.Event) {
switch ev.Kind {
case parser.KindProcessingItem:
a.parallel = true
a.completeID(ev.DumpID)
a.completeID(ev.DumpID, ev.Desc)

case parser.KindLaunchItem:
a.parallel = true
Expand All @@ -116,7 +125,7 @@ func (a *Aggregator) Feed(ev parser.Event) {
})
delete(a.inflight, ev.DumpID)
}
a.completeID(ev.DumpID)
a.completeID(ev.DumpID, ev.Desc)

case parser.KindCreating:
if !a.parallel {
Expand Down Expand Up @@ -152,20 +161,35 @@ func (a *Aggregator) Feed(ev parser.Event) {
}

// completeID marks the entry with the given dump id done, once. Its phase and
// byte size come from the plan. Ids not in the plan are ignored.
func (a *Aggregator) completeID(id int) {
// byte size come from the plan; ids the plan doesn't know are ignored, since a
// real plan is the authority on what this restore contains.
//
// Countless mode has no plan to consult, so the event's own description carries
// the phase instead - toc.SectionOf maps it exactly the way the plan would have.
// Sizes stay unknown either way: only preflight can stat a data file.
func (a *Aggregator) completeID(id int, desc string) {
if a.doneIDs[id] {
return
}
e, ok := a.plan.Get(id)
if !ok {
if a.plan != nil {
e, ok := a.plan.Get(id)
if !ok {
return
}
a.doneIDs[id] = true
a.incPhase(e.Section)
if e.Section == toc.Data && e.HasBytes {
a.bytesDone += e.Bytes
}
return
}

s := toc.SectionOf(desc)
if s == toc.SectionUnknown {
return
}
a.doneIDs[id] = true
a.incPhase(e.Section)
if e.Section == toc.Data && e.HasBytes {
a.bytesDone += e.Bytes
}
a.incPhase(s)
}

// completeSerial advances a phase in serial mode, where events carry no dump id
Expand Down Expand Up @@ -204,10 +228,11 @@ func (a *Aggregator) setCurrent(desc, name string) {

// incPhase increments a phase's done count, capped at its total so a stray
// double-count can never push a bar past 100%. The first real completion also
// ends the DROP wave.
// ends the DROP wave. Countless mode has no total to cap against - and nothing
// to overflow, since it draws counts rather than bars.
func (a *Aggregator) incPhase(s toc.Section) {
a.dropWaveOver = true
if a.done[s] < a.total[s] {
if a.countless || a.done[s] < a.total[s] {
a.done[s]++
}
}
Expand Down
70 changes: 69 additions & 1 deletion internal/aggregator/aggregator_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,14 +29,25 @@ func loadPlan(t *testing.T) *toc.RestorePlan {
// replay drives a captured stderr fixture (from the parser package's testdata)
// through parser then aggregator, exactly as the runner will wire them.
func replay(t *testing.T, fixture string, cfg Config) *Aggregator {
t.Helper()
return feed(t, New(loadPlan(t), cfg), fixture)
}

// replayCountless drives the same fixture with no plan - the stdin restore,
// where preflight is impossible and only numerators exist.
func replayCountless(t *testing.T, fixture string, cfg Config) *Aggregator {
t.Helper()
return feed(t, New(nil, cfg), fixture)
}

func feed(t *testing.T, agg *Aggregator, fixture string) *Aggregator {
t.Helper()
f, err := os.Open(filepath.Join("..", "parser", "testdata", fixture))
if err != nil {
t.Fatalf("open fixture: %v", err)
}
defer f.Close()

agg := New(loadPlan(t), cfg)
p := parser.New()
sc := bufio.NewScanner(f)
sc.Buffer(make([]byte, 0, 64*1024), 4*1024*1024)
Expand Down Expand Up @@ -321,3 +332,60 @@ func TestClassifyBenign(t *testing.T) {
t.Error("unique violation must be real")
}
}

// TestCountlessSerial is the stdin restore: no plan, so no denominators, but
// every phase still counts and the error groups are unchanged. Counts come from
// each event's own object description rather than a TOC lookup.
func TestCountlessSerial(t *testing.T) {
cfg := Config{Clean: true, NoOwner: true}
s := replayCountless(t, "clean-serial.stderr", cfg).Snapshot()
planned := replay(t, "clean-serial.stderr", cfg).Snapshot()

if !s.Countless {
t.Error("Countless = false, want true for a nil plan")
}
for _, p := range []PhaseProgress{s.Pre, s.Data, s.Post} {
if p.Total != 0 {
t.Errorf("%s total = %d, want 0 (no plan to count against)", p.Section, p.Total)
}
if p.Done == 0 {
t.Errorf("%s made no progress", p.Section)
}
}
if s.ByteSized {
t.Error("ByteSized = true; only preflight can size data files")
}

// The planned run caps each phase at its denominator, so countless counting
// can match it but never fall behind it.
if s.Pre.Done < planned.Pre.Done || s.Data.Done < planned.Data.Done || s.Post.Done < planned.Post.Done {
t.Errorf("countless %d/%d/%d behind planned %d/%d/%d",
s.Pre.Done, s.Data.Done, s.Post.Done,
planned.Pre.Done, planned.Data.Done, planned.Post.Done)
}

// Error classification never depended on the plan.
if s.ErrTotal != planned.ErrTotal || s.ErrBenign != planned.ErrBenign || s.ErrReal != planned.ErrReal {
t.Errorf("errors = %d/%d/%d, want the planned run's %d/%d/%d",
s.ErrTotal, s.ErrBenign, s.ErrReal,
planned.ErrTotal, planned.ErrBenign, planned.ErrReal)
}
}

// TestCountlessParallel confirms the -j path counts too: item events carry the
// object description alongside the dump id, so a missing plan costs only the
// denominators - in-flight tracking and timings are unaffected.
func TestCountlessParallel(t *testing.T) {
agg := replayCountless(t, "clean-j4.stderr", Config{Clean: true, NoOwner: true})
s := agg.Snapshot()

if s.Data.Done == 0 || s.Post.Done == 0 {
t.Errorf("data %d, post %d: want both counted", s.Data.Done, s.Post.Done)
}
if len(s.InFlight) != 0 {
t.Errorf("in-flight not drained: %+v", s.InFlight)
}
if len(s.Slowest) == 0 {
t.Error("no item timings; launching/finished pairs should still time")
}
}
6 changes: 6 additions & 0 deletions internal/aggregator/snapshot.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,11 @@ type Snapshot struct {
Data PhaseProgress
Post PhaseProgress

// Countless: the restore ran without a plan (a dump on stdin), so every
// Total is 0 and the phases carry counts only. The renderer drops the bars
// rather than drawing three empty tracks.
Countless bool

ByteSized bool
BytesDone int64
BytesTotal int64
Expand Down Expand Up @@ -101,6 +106,7 @@ func (a *Aggregator) Snapshot() Snapshot {
Pre: a.phase(toc.PreData),
Data: a.phase(toc.Data),
Post: a.phase(toc.PostData),
Countless: a.countless,
ByteSized: a.byteSized,
BytesDone: a.bytesDone,
BytesTotal: a.bytesTotal,
Expand Down
Loading
Loading