184 lines
6.7 KiB
Go
184 lines
6.7 KiB
Go
package ingest
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"errors"
|
|
"os"
|
|
"path/filepath"
|
|
"strconv"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// fakeEngine writes a shell script that behaves like tmctl's contract: NDJSON on stdout, human log
|
|
// on stderr, and one of the ratified exit codes. The supervision path is exercised with a real
|
|
// process — a mocked one would prove nothing about pipes, signals or exit codes.
|
|
func fakeEngine(t *testing.T, stdout, stderr string, code int) string {
|
|
t.Helper()
|
|
path := filepath.Join(t.TempDir(), "tmctl")
|
|
script := "#!/bin/sh\nprintf '%s' " + shellQuote(stdout) + "\nprintf '%s' " + shellQuote(stderr) + " >&2\nexit " + strconv.Itoa(code) + "\n"
|
|
if err := os.WriteFile(path, []byte(script), 0o755); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return path
|
|
}
|
|
|
|
func TestRunIngestsStreamAndMapsExitCode(t *testing.T) {
|
|
stream := helloLine + "\n" + `{"seq":2,"type":"bank_stop","data":{"terms_proposed":7}}` + "\n"
|
|
var engineLog bytes.Buffer
|
|
s := &Supervisor{
|
|
Bin: fakeEngine(t, stream, "bank-mining: new terms await owner signature\n", 3),
|
|
Workdir: t.TempDir(),
|
|
EngineLog: &engineLog,
|
|
}
|
|
sink := &recordingSink{}
|
|
|
|
got, err := s.Run(context.Background(), sink, "translate", "--verify-bank")
|
|
if err != nil {
|
|
t.Fatalf("run: %v", err)
|
|
}
|
|
if got != OutcomeBankStop {
|
|
t.Fatalf("outcome = %q, want %q (exit 3 is a deliberate halt, not a failure)", got, OutcomeBankStop)
|
|
}
|
|
if len(sink.applied) != 1 || sink.applied[0].Type != TypeBankStop {
|
|
t.Fatalf("stream not materialized: %+v", sink.applied)
|
|
}
|
|
// The engine's own log must land in the file, not in our structured logger: it carries money.
|
|
if !bytes.Contains(engineLog.Bytes(), []byte("bank-mining")) {
|
|
t.Fatalf("engine stderr not captured: %q", engineLog.String())
|
|
}
|
|
}
|
|
|
|
func TestRunReportsOutcomeEvenWhenStreamBreaks(t *testing.T) {
|
|
s := &Supervisor{Bin: fakeEngine(t, "not a stream\n", "", 1), Workdir: t.TempDir()}
|
|
got, err := s.Run(context.Background(), &recordingSink{})
|
|
if err == nil {
|
|
t.Fatal("a broken stream must be reported")
|
|
}
|
|
if got != OutcomeFailed {
|
|
t.Fatalf("outcome = %q, want %q", got, OutcomeFailed)
|
|
}
|
|
}
|
|
|
|
// PD-12. When the sink stops accepting, the run must END, not continue unwatched: a translation
|
|
// keeps spending for hours, and its ceiling and bank-stop events would go to io.Discard with nobody
|
|
// told. Mutation caught: dropping the stop() after a failed Ingest — the test then waits out the
|
|
// child's sleep and times out.
|
|
func TestFailingSinkStopsTheRun(t *testing.T) {
|
|
stream := helloLine + "\n" + `{"seq":2,"type":"progress","data":{}}` + "\n"
|
|
path := filepath.Join(t.TempDir(), "tmctl")
|
|
// Emits a valid stream, then behaves like a long translation: it stays alive until told to go.
|
|
script := "#!/bin/sh\nprintf '%s' " + shellQuote(stream) + "\nsleep 60\nexit 0\n"
|
|
if err := os.WriteFile(path, []byte(script), 0o755); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
s := &Supervisor{Bin: path, Workdir: t.TempDir()}
|
|
|
|
done := make(chan struct{})
|
|
go func() {
|
|
defer close(done)
|
|
if _, err := s.Run(context.Background(), &failingSink{}, "translate"); err == nil {
|
|
t.Error("a failing sink must be reported")
|
|
}
|
|
}()
|
|
select {
|
|
case <-done:
|
|
case <-time.After(20 * time.Second):
|
|
t.Fatal("the run outlived its sink: the engine was left running while the platform was blind")
|
|
}
|
|
}
|
|
|
|
type failingSink struct{}
|
|
|
|
func (failingSink) Begin(context.Context, Hello) error { return nil }
|
|
|
|
func (failingSink) Apply(context.Context, Envelope) error { return errors.New("database is down") }
|
|
|
|
func TestStatusDecodesTheResyncChannel(t *testing.T) {
|
|
// A real `tmctl status --json` body carries money and snapshot fields; the allowlist ignores
|
|
// them, and this fixture keeps one of each to prove it.
|
|
body := `{"book_id":"gzr","snapshot_id":"a1","committed_usd":1.25,"reserved_usd":0.5,
|
|
"total_units":10,"done":4,"pending":6,"eta_seconds":900,"unsigned_bank_terms":3,
|
|
"chapters":[{"chapter":1,"units_total":2,"units_done":2,"cost_usd":0.4,"verdict":"pass"}]}`
|
|
s := &Supervisor{Bin: fakeEngine(t, body, "", 0), Workdir: t.TempDir()}
|
|
got, err := s.Status(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("status: %v", err)
|
|
}
|
|
if got.BookID != "gzr" || got.Done != 4 || got.ETASeconds != 900 || got.UnsignedBankTerms != 3 {
|
|
t.Fatalf("status = %+v", got)
|
|
}
|
|
if got.Spend == nil || *got.Spend != 1_250_000 {
|
|
t.Fatalf("spend must be metered from status: %v", got.Spend)
|
|
}
|
|
// A status without the figure is not a status reporting zero.
|
|
absent, err := DecodeStatus([]byte(`{"book_id":"gzr"}`))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if absent.Spend != nil {
|
|
t.Fatalf("a missing committed_usd read as %v", *absent.Spend)
|
|
}
|
|
if len(got.Chapters) != 1 || got.Chapters[0].UnitsDone != 2 {
|
|
t.Fatalf("chapters = %+v", got.Chapters)
|
|
}
|
|
}
|
|
|
|
func TestOutcomeOfCoversTheExitContract(t *testing.T) {
|
|
for code, want := range map[int]Outcome{0: OutcomeClean, 1: OutcomeFailed, 2: OutcomeFlagged, 3: OutcomeBankStop, 7: OutcomeFailed} {
|
|
if got := outcomeOf(code); got != want {
|
|
t.Errorf("exit %d = %q, want %q", code, got, want)
|
|
}
|
|
}
|
|
}
|
|
|
|
func shellQuote(s string) string { return "'" + strings.ReplaceAll(s, "'", `'\''`) + "'" }
|
|
|
|
// A run we stopped ourselves is not a failed run. exec reports a cancelled command as
|
|
// context.Canceled instead of an *ExitError, so reading the outcome off the wait error alone marks
|
|
// every gracefully stopped translation as failed — every one in flight on an ordinary SIGTERM.
|
|
// Mutation caught: going back to `errors.As(waitErr, &exitErr)` as the only source of the outcome.
|
|
func TestStoppedRunKeepsTheEnginesOutcome(t *testing.T) {
|
|
path := filepath.Join(t.TempDir(), "tmctl")
|
|
// Behaves like tmctl: stops on the interrupt and exits 0.
|
|
script := "#!/bin/sh\ntrap 'exit 0' INT\nprintf '%s' " + shellQuote(helloLine+"\n") + "\nsleep 30\n"
|
|
if err := os.WriteFile(path, []byte(script), 0o755); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
s := &Supervisor{Bin: path, Workdir: t.TempDir()}
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
type result struct {
|
|
outcome Outcome
|
|
err error
|
|
}
|
|
done := make(chan result, 1)
|
|
go func() {
|
|
o, err := s.Run(ctx, nopSink{}, "translate")
|
|
done <- result{o, err}
|
|
}()
|
|
time.Sleep(200 * time.Millisecond) // let the engine reach its sleep
|
|
cancel()
|
|
|
|
select {
|
|
case got := <-done:
|
|
if got.outcome == OutcomeFailed {
|
|
t.Fatalf("a run stopped by us reported %q; the engine exited 0", got.outcome)
|
|
}
|
|
if !errors.Is(got.err, context.Canceled) {
|
|
t.Fatalf("a stopped run must be distinguishable from a finished one, got err %v", got.err)
|
|
}
|
|
case <-time.After(20 * time.Second):
|
|
t.Fatal("Run did not return")
|
|
}
|
|
}
|
|
|
|
// nopSink accepts everything and records nothing: this test is about the run's OUTCOME, not about
|
|
// what the sink saw.
|
|
type nopSink struct{}
|
|
|
|
func (nopSink) Begin(context.Context, Hello) error { return nil }
|
|
|
|
func (nopSink) Apply(context.Context, Envelope) error { return nil }
|