Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
8ce377a
fix(runnerhub): report a token-store fault as Unavailable, not Unauth…
rigel-mintaka Oct 5, 2026
b70eefa
test(runnerhub): assert the fixed Unavailable message (RIG-4529)
rigel-mintaka Oct 5, 2026
064ebc3
fix(stack): confirm cross-process teardown by group exit, not socket …
rigel-mintaka Oct 4, 2026
5ca6d3e
fix(stack): treat another uid's process group as recycled (RIG-3937)
rigel-mintaka Oct 5, 2026
d86a7d9
feat(stack): skip groups recorded in an earlier boot (RIG-4570)
rigel-mintaka Oct 6, 2026
a5fd62e
fix(stack): read a per-boot UUID and refuse a reboot drop while the s…
rigel-mintaka Oct 6, 2026
c1745f9
fix(stack): bound container teardown by ctx and never block Phase A (…
rigel-mintaka Oct 5, 2026
fe93f42
fix(stack): drain consumers before infra; hard-kill without stop-time…
rigel-mintaka Oct 5, 2026
0fc7c01
fix(stack): still stop infra when down is cancelled mid-teardown (RIG…
rigel-mintaka Oct 5, 2026
3062afc
test(stack): drive the rebooted-record gateway through ctx teardown (…
rigel-mintaka Oct 6, 2026
de20f34
fix(stack): default the microVM runner run root to the runtime dir (R…
rigel-mintaka Oct 5, 2026
da5149a
fix(fabric): give the max-deliveries advisory fetch its own deadline …
rigel-mintaka Oct 7, 2026
b142073
test(fabric): pin the advisory fetch deadline near its floor (RIG-4786)
rigel-mintaka Oct 7, 2026
ccacccd
Merging de20f34b171ce91dd50af764f099411a530e2759 into trunk-temp/pr-1…
trunk-io[bot] Oct 8, 2026
0e9e165
Merging b70eefaca6c0c474379ca9beec1155a439d83fa9 into trunk-temp/pr-1…
trunk-io[bot] Oct 8, 2026
be92dc0
Merging 3062afc7f0b9a2ccbaccd3854d697f5d5be33ff5 into trunk-temp/pr-1…
trunk-io[bot] Oct 8, 2026
99b6e65
Merging b14207312fa6907f9ebd4358d8fe7aeeb0a14283 into trunk-temp/pr-1…
trunk-io[bot] Oct 8, 2026
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
12 changes: 12 additions & 0 deletions docs/designs/infra/release/compass-distribution/design.md
Original file line number Diff line number Diff line change
Expand Up @@ -419,6 +419,18 @@ Mechanics, grounded in the current seams:
`readPgidFile` dispatches on the tag; the hard-error-on-malformed
discipline is unchanged ("signaling off a half-understood record is
exactly the blast radius the design forbids", `pgidfile.go:100-103`).
+ **Boot identity (RIG-4570 option A).** The v2 header is
`<version> <writerPid> [<bootid>]`: the writer stamps the current boot
(Linux `/proc/sys/kernel/random/boot_id`, darwin `kern.bootsessionuuid`;
both are per-boot UUIDs that a wall-clock step cannot move). A
`down` that reads a boot id differing from the current boot signals no
process entry and drops the record; container entries keep their name
teardown. If the stack socket still answers under a mismatch, `down`
refuses, keeps the record, and signals nothing. A boot id that is not a
UUID (8-4-4-4-12 hex) is a malformed header. A missing boot id (older
builds, or an unreadable id) is *unknown* and falls back to the
per-entry identity check. A v2 reader that predates the boot id refuses a
three-column header as malformed. A v1 header never carries the column.
+ **Cross-version rule.** A v1-only binary never half-parses a v2 record —
but by the *entry-line grammar*, not a header-version check: shipped v1
`readPgidFile` stores `header[0]` as `Version` and never compares it to
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,8 @@ supervised children** (containers scoped out — Open Question 0).
because `up` always exits after a successful spawn (`main.go:235-238`), the
writer pid is dead in every linger teardown, so it discriminates nothing about
whether the *children* are alive. Plain text, trailing newline, 0600.
(Since DL-262 v2 the header also carries an optional boot id; a record from
an earlier boot has no live group to signal — see the DL-262 record.)
- **Write timing**: rewritten **atomically (temp + rename in the state dir)
after each successful child spawn** in `spawnChain`
(`go/internal/stack/stack.go:171-228`), i.e. the file always reflects the set
Expand Down
5 changes: 5 additions & 0 deletions go/cmd/compass-app/embedded.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import (
"path/filepath"
"runtime"
"strings"
"syscall"

"connectrpc.com/connect"

Expand Down Expand Up @@ -245,6 +246,10 @@ func runStackDown(bin string) func(ctx context.Context, args []string) error {
//nolint:gosec // G204: bin is operator/PATH-resolved (resolveStackBin) and
// the argv is pipeline-assembled (stackDownArgs), not user input.
cmd := exec.CommandContext(ctx, bin, args...)
// down has already consumed its teardown record, so a timeout must SIGTERM
// it: SIGKILL would skip the survivor rewrite and leak the stack.
cmd.Cancel = func() error { return cmd.Process.Signal(syscall.SIGTERM) }
cmd.WaitDelay = stackDownCancelGrace
cmd.Env = prependExecDirToPath(os.Environ(), filepath.Dir(bin))
stderr, cleanup, capErr := captureStderr(cmd)
if capErr != nil {
Expand Down
33 changes: 33 additions & 0 deletions go/cmd/compass-app/embedded_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,11 @@ import (
"errors"
"net"
"net/http"
"os"
"path/filepath"
"slices"
"strings"
"syscall"
"testing"
"time"

Expand Down Expand Up @@ -383,6 +385,37 @@ func TestRunStackDownZeroExitSucceeds(t *testing.T) {
}
}

// TestRunStackDownCancelSendsSIGTERM: a timed-out down has already consumed its
// teardown record, so cancel must SIGTERM it (letting it rewrite survivors), not
// SIGKILL it. The child traps TERM and leaves a marker only a SIGTERM can write.
func TestRunStackDownCancelSendsSIGTERM(t *testing.T) {
dir := t.TempDir()
ready := filepath.Join(dir, "ready")
marker := filepath.Join(dir, "rewrote")
if err := syscall.Mkfifo(ready, 0o600); err != nil {
t.Fatalf("mkfifo: %v", err)
}
ctx, cancel := context.WithTimeout(context.Background(), embeddedTestTimeout)
defer cancel()
go func() {
// Opening the FIFO blocks until the child has armed its trap.
f, err := os.Open(ready)
if err == nil {
_ = f.Close() // read end of a gate FIFO; nothing to flush
}
cancel()
}()

script := "trap 'echo ok > \"$1\"; exit 3' TERM; echo > \"$2\"; while :; do sleep 1 & wait; done"
err := runStackDown("/bin/sh")(ctx, []string{"-c", script, "sh", marker, ready})
if err == nil {
t.Fatal("stackDown err = nil, want the cancelled child's exit error")
}
if _, statErr := os.Stat(marker); statErr != nil {
t.Fatalf("child did not run its SIGTERM handler (marker: %v); err = %v", statErr, err)
}
}

// classifyDeps returns a preflight.Deps whose every injected effect passes (host
// GOOS linux) — mirroring preflight_test.go's okDeps. Tests override individual
// fields to drive one failing check at a time through classifyPreflight.
Expand Down
4 changes: 4 additions & 0 deletions go/cmd/compass-app/lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,10 @@ import (
// reusing it.
const stackDownTimeout = 60 * time.Second

// stackDownCancelGrace is how long a timed-out down gets after SIGTERM to
// rewrite its survivor record before os/exec escalates to SIGKILL.
const stackDownCancelGrace = 10 * time.Second

// quitController is the explicit "Quit and stop stack" orchestration over its
// injected effects. It holds the teardown seam (stackDown), the argv inputs
// (params, resolved once in run()), the app-quit indirection (quit, wired to
Expand Down
35 changes: 12 additions & 23 deletions go/cmd/compass-stack/cross_process_podman_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -384,21 +384,19 @@ func waitServerAnswering(t *testing.T, deps stack.Deps, socketPath string) {
// its process-group id and the group leader's start-time token. The start time
// is what closes the pid-recycling window — the same (Pgid, StartTime) identity
// the production teardown checks (internal/stack/pgidfile.go pgidEntry,
// adapters/groupsignal.go Alive).
// adapters/groupsignal.go Liveness).
type recordedGroup struct {
pgid int
startTime uint64
}

// waitGroupsGone polls every recorded process group until it is gone or the
// budget elapses — the authoritative "the children are actually dead" proof
// after a cross-process down. "Gone" is identity-checked, mirroring production's
// teardown gate (adapters/groupsignal.go Alive): a group is gone when it is
// ESRCH, OR it still exists but its leader's start time no longer equals the
// recorded token (the kernel recycled the pid to an unrelated leader). Without
// the identity check the probe is a false-FAILURE risk on the shared box — a
// dead child's pgid reused by another process would read as "still alive" and
// fail the test at budget even though teardown worked.
// after a cross-process down. "Gone" mirrors production's teardown gate
// (adapters/groupsignal.go Liveness): a group is gone when it is ESRCH, OR its
// leader's start time no longer equals the recorded token (a recycled pid).
// Without the identity check a dead child's pgid reused by another process
// would read as "still alive" and fail the test even though teardown worked.
//
// It signals nothing: kill(-pgid, 0) is the existence probe (signal 0), run
// only on pgids read from the stack's OWN stack.pgids record, never a scan
Expand Down Expand Up @@ -429,13 +427,10 @@ func waitGroupsGone(t *testing.T, groups []recordedGroup, budget time.Duration)
}
}

// groupAlive reports whether the process group named by pgid still exists AND
// its leader's start time equals the recorded token — the same existence-then-
// identity gate production teardown uses (adapters/groupsignal.go Alive). A
// group that is ESRCH, or whose leader start time no longer matches (a recycled
// pid), or whose /proc entry cannot be read is reported not-alive: for a
// post-down liveness probe the safe verdict is "gone", never a false "alive"
// off a pid the kernel reused. It signals nothing — kill(-pgid, 0) is signal 0.
// groupAlive reports whether the recorded process group still needs teardown,
// mirroring adapters/groupsignal.go Liveness: ESRCH or a recycled leader is
// gone; a matching leader or an unreadable one (members outliving a reaped
// leader keep the pgid) is alive. It signals nothing — kill(-pgid, 0) is signal 0.
// An unexpected kill errno (not ESRCH/EPERM) is surfaced as an error.
func groupAlive(pgid int, startTime uint64) (bool, error) {
err := syscall.Kill(-pgid, 0)
Expand All @@ -445,16 +440,10 @@ func groupAlive(pgid int, startTime uint64) (bool, error) {
case err != nil && !errors.Is(err, syscall.EPERM):
return false, err // unexpected errno
}
// Exists (nil or EPERM). Confirm identity via the leader's start time; a read
// failure means the leader vanished or /proc is unreadable — treat as gone.
got, rerr := readLeaderStartTime(pgid)
if rerr != nil {
// Deliberate: a /proc read failure on an existing pgid means the leader
// vanished between the two syscalls (or /proc is unreadable) — the safe
// post-down verdict is "gone", never a false "alive". Mirrors production
// adapters/groupsignal.go Alive, which also treats a read failure as
// not-alive.
return false, nil //nolint:nilerr // read failure => leader gone => not-alive (see comment)
// The group exists without a readable leader: orphaned members remain.
return true, nil //nolint:nilerr // read failure on an existing group => orphaned => alive
}
return got == startTime, nil
}
Expand Down
1 change: 1 addition & 0 deletions go/internal/auth/interceptor.go
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,7 @@ func resolveBearer(ctx context.Context, st *store.Store, header string) (store.A
// Oracle-safe: every resolution failure — unknown, revoked, or a cross-door
// (Runner) token — is one indistinguishable CodeUnauthenticated to the client.
// The distinct sentinel is logged (debug) as a server-side audit signal only.
// A store fault (ErrTokenLookupFailed) also stays Unauthenticated on this door.
slog.DebugContext(ctx, "network door rejected bearer token", "reason", err)
return store.AccountID(""), "", connect.NewError(connect.CodeUnauthenticated, errInvalidToken)
}
Expand Down
42 changes: 26 additions & 16 deletions go/internal/auth/token.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,21 +64,25 @@ func IssueAccountToken(ctx context.Context, st *store.Store, account store.Accou
return "", errors.New("persisting issued token: hash collision on two successive mints")
}

// Sentinel resolution failures returned by ResolveToken. They exist so the server
// can LOG which case fired (audit), NOT so a door tells them apart to the client:
// a door MUST map all three to the same bare CodeUnauthenticated, or the response
// becomes an oracle for whether a token is unknown, revoked, or for the other door.
// Sentinel resolution failures returned by ResolveToken. The three credential
// verdicts exist so the server can LOG which case fired (audit), NOT so a door
// tells them apart to the client: a door MUST map them to the same bare
// CodeUnauthenticated, or the response becomes an oracle for whether a token is
// unknown, revoked, or for the other door. ErrTokenLookupFailed is not a verdict.
var (
// ErrTokenNotFound: the presented token was never issued (or the store has
// no live record of it). Any unexpected store error folds here too, so
// resolution fails closed.
// no live record of it).
ErrTokenNotFound = errors.New("auth: token not found")
// ErrTokenRevoked: the token was issued but has since been withdrawn.
ErrTokenRevoked = errors.New("auth: token revoked")
// ErrWrongKind: the token resolves, but to the other subject kind — a Runner
// token presented to the account door, or an account token to the Runner
// door. The OQ7 cross-door rejection (design.md:1308-1314).
ErrWrongKind = errors.New("auth: token subject kind mismatch")
// ErrTokenLookupFailed: the store could not answer, so no verdict exists. It
// still fails closed; a door MAY report it as Unavailable since it says nothing
// about the token.
ErrTokenLookupFailed = errors.New("auth: token lookup failed")
)

// ResolveToken authenticates a presented bearer to a subject of the required
Expand All @@ -88,32 +92,38 @@ var (
// comparison never touches a stored plaintext — there is none), then verifies
// the resolved subject's kind against want.
//
// It returns a distinct sentinel per failure — ErrTokenNotFound (never issued or
// an unexpected store error, folded here to fail closed), ErrTokenRevoked
// (withdrawn), ErrWrongKind (issued for the other door) — so the server can log
// which fired. Every caller MUST map all three to the same bare
// It returns a distinct sentinel per failure — ErrTokenNotFound (never issued),
// ErrTokenRevoked (withdrawn), ErrWrongKind (issued for the other door) — so the
// server can log which fired. Every caller MUST map those three to the same bare
// CodeUnauthenticated: the distinction is a server-side audit signal, never a
// client-visible one (a distinguishable response is a token-existence oracle).
// Any other store error, including ctx cancellation, wraps ErrTokenLookupFailed.
// Both the account door (want=SubjectAccount) and the Runner door
// (want=SubjectRunner) share this one resolver, so the security-critical
// resolve+kind-gate lives and is tested in exactly one place; each door adds only
// its own trivial typed wrap on the returned Subject.
func ResolveToken(ctx context.Context, st *store.Store, presented string, want store.SubjectKind) (store.Subject, error) {
subj, err := st.ResolveTokenHash(ctx, hashToken(presented))
if err != nil {
if errors.Is(err, store.ErrTokenRevoked) {
return store.Subject{}, ErrTokenRevoked
}
// store.ErrNotFound — and any other store error — is not a live
// credential; fail closed as not-found.
return store.Subject{}, ErrTokenNotFound
return store.Subject{}, tokenResolutionError(err)
}
if subj.Kind != want {
return store.Subject{}, ErrWrongKind
}
return subj, nil
}

// tokenResolutionError distinguishes credential verdicts from store failures.
func tokenResolutionError(err error) error {
if errors.Is(err, store.ErrTokenRevoked) {
return ErrTokenRevoked
}
if errors.Is(err, store.ErrNotFound) {
return ErrTokenNotFound
}
return fmt.Errorf("%w: %w", ErrTokenLookupFailed, err)
}

// RevokeToken withdraws a bearer token by its presented plaintext, marking the
// stored hash revoked so ResolveToken thereafter fails it as ErrTokenRevoked.
// Hashing lives here (hashToken), the one place issuance and resolution agree on
Expand Down
52 changes: 45 additions & 7 deletions go/internal/auth/token_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,9 +38,8 @@ func TestIssueThenResolveRoundTripsToTheIssuedAccount(t *testing.T) {
}
}

// unknown_token_resolves_to_none: a token the store never issued must not
// resolve — even with an unrelated live token present — and the failure is the
// distinct ErrTokenNotFound sentinel.
// Unknown tokens return ErrTokenNotFound, while operational lookup errors use
// ErrTokenLookupFailed and retain their cause.
func TestUnknownTokenResolvesToNotFound(t *testing.T) {
ctx := context.Background()
st, admin, _ := openTestStore(t)
Expand All @@ -56,10 +55,49 @@ func TestUnknownTokenResolvesToNotFound(t *testing.T) {
}
}

// a revoked token stops resolving and surfaces the distinct ErrTokenRevoked
// sentinel — separate from ErrTokenNotFound so the server can tell a withdrawn
// credential from an unknown one (the distinction is audit-only; the door still
// maps both to one CodeUnauthenticated).
func TestResolveTokenLookupFailuresAreNotNotFound(t *testing.T) {
ctx := context.Background()
st, _, _ := openTestStore(t)

_, err := ResolveToken(ctx, st, "never-issued", store.SubjectAccount)
if !errors.Is(err, ErrTokenNotFound) {
t.Fatalf("not-found lookup error = %v, want ErrTokenNotFound", err)
}
canceledCtx, cancel := context.WithCancel(ctx)
cancel()
_, err = ResolveToken(canceledCtx, st, "lookup-failure", store.SubjectAccount)
if !errors.Is(err, ErrTokenLookupFailed) {
t.Fatalf("canceled lookup error = %v, want ErrTokenLookupFailed", err)
}
if errors.Is(err, ErrTokenNotFound) {
t.Fatalf("canceled lookup error = %v, must not be ErrTokenNotFound", err)
}
}

func TestResolveTokenStoreLookupFailureIsNotNotFound(t *testing.T) {
ctx := context.Background()
st, _, _ := openTestStore(t)
st.Close()

_, err := ResolveToken(ctx, st, "lookup-failure", store.SubjectAccount)
if !errors.Is(err, ErrTokenLookupFailed) {
t.Fatalf("closed-store lookup error = %v, want ErrTokenLookupFailed", err)
}
if errors.Is(err, ErrTokenNotFound) {
t.Fatalf("closed-store lookup error = %v, must not be ErrTokenNotFound", err)
}
}

func TestTokenResolutionErrorPreservesCause(t *testing.T) {
cause := errors.New("conn refused")
err := tokenResolutionError(cause)
if !errors.Is(err, ErrTokenLookupFailed) || !errors.Is(err, cause) {
t.Fatalf("lookup error = %v, want lookup sentinel and original cause", err)
}
}

// a revoked token stops resolving and surfaces ErrTokenRevoked, separate from
// ErrTokenNotFound; both are credential verdicts that doors map to Unauthenticated.
func TestRevokedTokenResolvesToRevoked(t *testing.T) {
ctx := context.Background()
st, admin, _ := openTestStore(t)
Expand Down
2 changes: 2 additions & 0 deletions go/internal/fabric/SUBJECTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -277,6 +277,8 @@ instance handles each advisory.
- The reason reads `fabric: dropped by the server after N delivery attempts
(callback outlived ack_wait)`.
- A message that aged out of the stream before the fetch is logged, not parked.
- The fetch has its own bounded deadline. This is needed because the advisory arrives
after `AckWait`.
- The fabric's NATS user needs subscribe permission on
`$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>`.

Expand Down
8 changes: 6 additions & 2 deletions go/internal/fabric/event_fabric.go
Original file line number Diff line number Diff line change
Expand Up @@ -579,6 +579,9 @@ func (f *Fabric) publishDLQ(ctx context.Context, subject string, data []byte, ca
return reason, nil
}

// minAdvisoryGetTimeout floors the advisory's stream fetch under a short AckWait.
const minAdvisoryGetTimeout = 5 * time.Second

// maxDeliveriesAdvisory is the server's notice that a consumer gave up on a
// message after MaxDeliver attempts. Only the fields the park needs.
type maxDeliveriesAdvisory struct {
Expand Down Expand Up @@ -615,8 +618,9 @@ func (f *Fabric) parkOnMaxDeliveries(ctx context.Context, subject string) (*nats
f.notifyParkDecided(path, false)
return
}
// Bounded so a stalled fetch cannot hold the claim in flight and park its waiters.
getCtx, cancelGet := context.WithTimeout(context.WithoutCancel(ctx), f.cfg.ackWait())
// Bounded so a stalled fetch cannot hold the claim; floored because AckWait
// has already expired by the advisory and says nothing about fetch latency.
getCtx, cancelGet := context.WithTimeout(context.WithoutCancel(ctx), max(f.cfg.ackWait(), minAdvisoryGetTimeout))
getParkedMsg := f.getParkedMsg
if getParkedMsg == nil {
getParkedMsg = func(ctx context.Context, seq uint64) (*jetstream.RawStreamMsg, error) {
Expand Down
Loading
Loading