From 3c80b244275769addb11a826e6d4a891e045a72c Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 1 Oct 2026 02:03:00 +0000 Subject: [PATCH 1/3] Deliver a reaped leader's buffered output before Exited The process protocol now orders Exited after every output byte buffered in a captured stream when the service reaps the leader, as a native parent sees all of its child's output once waitpid returns. Output written later by processes that still hold the stream can follow Exited. Exited never waits for an acknowledgement: when the replay limit stops reading, it follows the output read until then. With a PTY, output still in the kernel's asynchronous terminal queue at the reap can follow Exited. The Linux service records each open stream's mark at the reap, which is the offset it has delivered plus FIONREAD/TIOCINQ, and pushes Exited once every stream reaches its mark, closes or is abandoned, or once the retained output reaches the replay limit. Readers now read non-blocking and push under the operation lock, so no byte is uncounted between a read and its event. That also makes the replay budget exact and removes the per-stream reservations. An empty read lowers the mark, so a terminal flush cannot leave it unreachable. CloseOutput closes the file outside the lock because Close waits for a read in progress. --- .../internal/processservice/launch.go | 113 ++++++++++++------ .../internal/processservice/operation.go | 66 ++++++++-- .../internal/processservice/service_test.go | 48 +++++++- docs/process-protocol.md | 2 + 4 files changed, 175 insertions(+), 54 deletions(-) diff --git a/apps/sandboxio/internal/processservice/launch.go b/apps/sandboxio/internal/processservice/launch.go index 1374bfa9..ad9507de 100644 --- a/apps/sandboxio/internal/processservice/launch.go +++ b/apps/sandboxio/internal/processservice/launch.go @@ -18,7 +18,8 @@ import ( ) // launch spawns the process and emits Started or StartFailed, then starts -// the output readers. The reaper reports the leader's exit. +// the output readers. The leader's exit follows the output it left buffered; +// see markLocked. func (op *operation) launch(spec sp.ProcessSpec) { l, f := op.spawn(spec) op.mu.Lock() @@ -37,15 +38,12 @@ func (op *operation) launch(spec sp.ProcessSpec) { op.state = sp.StateRunning op.push(sp.StartedEvent{EventHeader: op.header()}) if op.leaderGone { - op.exitedLocked() + op.markLocked() } cancel, grace := op.cancelPending, op.pendingGrace op.mu.Unlock() - // Each stream reads at most its share of the replay limit, so a reader - // blocked on an idle stream never holds the space another one needs. - chunk := min(int(op.s.caps.MaxDataBytes), int(op.s.caps.MaxReplayBytesPerOperation)/len(l.streams)) for _, st := range l.streams { - go op.read(st, chunk) + go op.read(st) } if cancel { op.cancel(grace) @@ -228,51 +226,43 @@ func openPTY(spec sp.PTYSpec) (master, tty *os.File, err error) { } // read captures one stream until end of file, a read error, or CloseOutput. -// Before each read of up to chunk bytes it reserves chunk bytes of the replay -// limit, so the streams together never exceed it, and it pauses until they -// are free. -func (op *operation) read(st *stream, chunk int) { +// It reads only while the retained output is below the replay limit, and at +// most the space left, so the streams together never exceed it. Each read +// and the push of what it read happen under op.mu, so every byte is either +// still in the kernel's buffer or in an event; markLocked relies on that. +func (op *operation) read(st *stream) { limit := int(op.s.caps.MaxReplayBytesPerOperation) - buf := make([]byte, chunk) + buf := make([]byte, op.s.caps.MaxDataBytes) var disp sp.OutputDisposition - for disp == 0 { + rc, err := st.f.SyscallConn() + for err == nil && disp == 0 { op.mu.Lock() - for op.retained+op.reserved+chunk > limit && !st.abandoned { + for op.retained >= limit && !st.abandoned { op.cond.Wait() } - if st.abandoned { - op.mu.Unlock() - disp = sp.OutputAbandoned - break - } - op.reserved += chunk - op.mu.Unlock() - n, err := st.f.Read(buf) - op.mu.Lock() - op.reserved -= chunk - if n > 0 { - op.push(sp.OutputEvent{EventHeader: op.header(), Stream: st.name, Offset: st.offset, Data: bytes.Clone(buf[:n])}) - st.offset += uint64(n) - } else { - op.cond.Broadcast() // the unused reservation is free again - } - switch { - case err == nil: - case st.abandoned: - disp = sp.OutputAbandoned - case err == io.EOF, op.pty != nil && errors.Is(err, syscall.EIO): // EIO: every slave descriptor closed - disp = sp.OutputDrained - default: - disp = sp.OutputLost - } op.mu.Unlock() + // The callback returns false to wait until the stream is readable. + err = rc.Read(func(fd uintptr) bool { + op.mu.Lock() + defer op.mu.Unlock() + var empty bool + disp, empty = op.readLocked(st, int(fd), buf) + return !empty + }) } st.f.Close() op.mu.Lock() defer op.mu.Unlock() + if disp == 0 { // CloseOutput closed the file, or polling it failed + disp = sp.OutputLost + if st.abandoned { + disp = sp.OutputAbandoned + } + } st.closed = true op.push(sp.StreamClosedEvent{EventHeader: op.header(), Stream: st.name, Offset: st.offset, Disposition: disp}) + op.exitWhenDrainedLocked() op.worst = max(op.worst, disp) if op.openStreams--; op.openStreams > 0 { return @@ -285,3 +275,50 @@ func (op *operation) read(st *stream, chunk int) { op.push(sp.OutputClosedEvent{EventHeader: op.header(), Disposition: worst}) op.settleLocked() } + +// readLocked reads the non-blocking fd once, unless CloseOutput or the replay +// limit stops it, and pushes what it read. It returns the disposition once +// the stream ends, and empty when nothing is buffered. +func (op *operation) readLocked(st *stream, fd int, buf []byte) (disp sp.OutputDisposition, empty bool) { + room := int(op.s.caps.MaxReplayBytesPerOperation) - op.retained + switch { + case st.abandoned: + return sp.OutputAbandoned, false + case room <= 0: + return 0, false + } + buf = buf[:min(len(buf), room)] + n, err := unix.Read(fd, buf) + for err == unix.EINTR { + n, err = unix.Read(fd, buf) + } + switch { + case err == unix.EAGAIN: + // Nothing is buffered, so whatever the mark counted was read or, on a + // terminal, discarded by a flush. + if st.mark > st.offset { + st.mark = st.offset + op.exitWhenDrainedLocked() + } + return 0, true + case err == nil && n > 0: + op.push(sp.OutputEvent{EventHeader: op.header(), Stream: st.name, Offset: st.offset, Data: bytes.Clone(buf[:n])}) + st.offset += uint64(n) + op.exitWhenDrainedLocked() + return 0, false + case err == nil, op.pty != nil && err == unix.EIO: // EIO: every slave descriptor closed + return sp.OutputDrained, false + } + return sp.OutputLost, false +} + +// buffered returns the bytes the kernel holds for reading from f: a PTY +// master's input queue, or a pipe's contents (FIONREAD, the same request). It +// returns 0 once f is closed. +func buffered(f *os.File) uint64 { + n := 0 + if rc, err := f.SyscallConn(); err == nil { + rc.Control(func(fd uintptr) { n, _ = unix.IoctlGetInt(int(fd), unix.TIOCINQ) }) + } + return uint64(max(n, 0)) +} diff --git a/apps/sandboxio/internal/processservice/operation.go b/apps/sandboxio/internal/processservice/operation.go index 69cd13ed..f8665c01 100644 --- a/apps/sandboxio/internal/processservice/operation.go +++ b/apps/sandboxio/internal/processservice/operation.go @@ -41,12 +41,10 @@ type operation struct { settled bool // log holds the retained events first..last; retained counts their - // Output bytes, and reserved the bytes readers may add. acked is the - // highest acknowledged sequence. + // Output bytes. acked is the highest acknowledged sequence. log []sp.Event first, last uint64 retained int - reserved int acked uint64 // observer is the generation of the stream receiving events, zero when @@ -67,6 +65,7 @@ type operation struct { worst sp.OutputDisposition // leaderGone is set when the reaper reaps the leader, with its status. + // Exited follows the output buffered then; see markLocked. leaderGone bool status unix.WaitStatus // cu holds the processes proven to be in the session; the spawn sets it. @@ -88,6 +87,8 @@ type stream struct { offset uint64 abandoned bool closed bool + // mark is the offset Exited waits for once the leader is reaped. + mark uint64 } func newOperation(s *Service, key opKey, digest [sha256.Size]byte) *operation { @@ -165,10 +166,44 @@ func (op *operation) reaped(ws unix.WaitStatus) { defer op.mu.Unlock() op.leaderGone, op.status = true, ws if op.state == sp.StateRunning { - op.exitedLocked() + op.markLocked() } } +// markLocked runs once the leader is reaped and the streams are published. It +// sets each open stream's mark to what it has delivered plus what the kernel +// still buffers for it, so Exited follows every byte written before the reap. +// A read and its push happen under op.mu, so no byte is between the two. Linux +// moves PTY output to the master's input queue, which holds 4 KiB, +// asynchronously through the tty flip buffer; output still there is not +// counted and can follow Exited. +func (op *operation) markLocked() { + for _, st := range op.streams { + if !st.closed && !st.abandoned { + st.mark = st.offset + buffered(st.f) + } + } + op.exitWhenDrainedLocked() +} + +// exitWhenDrainedLocked reports the reaped leader's exit once every stream +// has delivered its mark or closed. At the replay limit no stream is read +// until the client acknowledges, so the exit is reported then: it never +// waits for an acknowledgement. +func (op *operation) exitWhenDrainedLocked() { + if !op.leaderGone || op.state != sp.StateRunning { + return + } + if op.retained < int(op.s.caps.MaxReplayBytesPerOperation) { + for _, st := range op.streams { + if !st.closed && !st.abandoned && st.offset < st.mark { + return + } + } + } + op.exitedLocked() +} + // exitedLocked reports the reaped leader's exit and starts watching the scope. func (op *operation) exitedLocked() { ws := op.status @@ -358,26 +393,35 @@ func (op *operation) closeStdin(offset uint64) error { func (op *operation) closeOutput(name sp.Stream) error { op.mu.Lock() - defer op.mu.Unlock() + st, err := op.abandonLocked(name) + op.mu.Unlock() + if err != nil { + return err + } + // Outside op.mu: Close waits for a read in progress, which holds op.mu. + st.f.Close() + return nil +} + +func (op *operation) abandonLocked(name sp.Stream) (*stream, error) { if op.released { - return released() + return nil, released() } if op.state == sp.StateStarting || op.state == sp.StateStartFailed { - return notRunning("the operation has no output") + return nil, notRunning("the operation has no output") } for _, st := range op.streams { if st.name != name { continue } if st.closed || st.abandoned { - return sp.Fail(sp.CodeOutputClosed, sandboxwire.EffectNone, "stream %d is closed", name) + return nil, sp.Fail(sp.CodeOutputClosed, sandboxwire.EffectNone, "stream %d is closed", name) } st.abandoned = true op.cond.Broadcast() - st.f.Close() - return nil + return st, nil } - return invalid("the operation does not capture stream %d", name) + return nil, invalid("the operation does not capture stream %d", name) } func (op *operation) resize(size sp.WindowSize) error { diff --git a/apps/sandboxio/internal/processservice/service_test.go b/apps/sandboxio/internal/processservice/service_test.go index 49f93a20..261baab3 100644 --- a/apps/sandboxio/internal/processservice/service_test.go +++ b/apps/sandboxio/internal/processservice/service_test.go @@ -10,6 +10,7 @@ import ( "net" "os" "path/filepath" + "runtime" "strconv" "strings" "sync/atomic" @@ -228,17 +229,54 @@ func TestStartFailure(t *testing.T) { } } -// A background job holding stderr keeps output open after the leader exits. -func TestExitedBeforeOutputClosed(t *testing.T) { +// Every byte the leader wrote before it exited precedes Exited, even when +// it wrote more than a pipe buffer just before exiting. With one P the reaper +// often runs before the reader has drained the pipe. +func TestOutputPrecedesExited(t *testing.T) { + defer runtime.GOMAXPROCS(runtime.GOMAXPROCS(1)) h := newHarness(t, DefaultConfig()) - evs := events(t, h.start(h.connect(), pipeSpec("sh", "-c", "sleep 1 >/dev/null & echo hi; exec sleep 0")), sp.EventScopeClosed) + c := h.connect() + const size = 200 << 10 + for range 20 { + evs := events(t, h.start(c, pipeSpec("head", "-c", strconv.Itoa(size), "/dev/zero")), sp.EventExited) + if got := len(output(evs, sp.StreamStdout)); got != size { + t.Fatalf("%d of %d bytes before Exited", got, size) + } + } +} + +// Output a background process writes after the leader exits follows Exited, +// and keeps OutputClosed pending until it ends. +func TestLaterOutputFollowsExited(t *testing.T) { + h := newHarness(t, DefaultConfig()) + evs := events(t, h.start(h.connect(), pipeSpec("sh", "-c", "printf a; (sleep 1; printf b) & exit 0")), sp.EventOutputClosed) _, exited := find[sp.ExitedEvent](t, evs) - closed, outputClosed := find[sp.OutputClosedEvent](t, evs) - if exited > outputClosed || closed.Disposition != sp.OutputDrained || output(evs, sp.StreamStdout) != "hi\n" { + closed := evs[len(evs)-1].(sp.OutputClosedEvent) + if output(evs[:exited], sp.StreamStdout) != "a" || output(evs[exited:], sp.StreamStdout) != "b" || closed.Disposition != sp.OutputDrained { t.Fatalf("events %v", evs) } } +// Exited never waits for an acknowledgement: at the replay limit it follows +// the output already read, and the rest arrives once acknowledged. +func TestExitedAtReplayLimit(t *testing.T) { + cfg := DefaultConfig() + cfg.MaxReplayBytesPerOperation = sandboxwire.MaxChunk + h := newHarness(t, cfg) + const size = sandboxwire.MaxChunk + 1000 // the rest fits in the pipe + op := h.start(h.connect(), pipeSpec("head", "-c", strconv.Itoa(size), "/dev/zero")) + evs := events(t, op, sp.EventExited) + if got := len(output(evs, sp.StreamStdout)); got != sandboxwire.MaxChunk { + t.Fatalf("%d bytes before Exited", got) + } + if err := op.Ack(context.Background(), evs[len(evs)-1].Header().Sequence); err != nil { + t.Fatal(err) + } + if got := len(output(events(t, op, sp.EventOutputClosed), sp.StreamStdout)); got != size-sandboxwire.MaxChunk { + t.Fatalf("%d bytes after Exited", got) + } +} + func TestStdinOffsetsAndHalfClose(t *testing.T) { h := newHarness(t, DefaultConfig()) ctx := context.Background() diff --git a/docs/process-protocol.md b/docs/process-protocol.md index 9aad8504..9cb2267b 100644 --- a/docs/process-protocol.md +++ b/docs/process-protocol.md @@ -86,6 +86,8 @@ Ordering: - `Started` or `StartFailed` comes first. - `Exited` and `OutputClosed` are independent. A background process holding a stream open keeps `OutputClosed` pending after `Exited`, and output can close before the leader exits. +- `Exited` follows every output byte buffered in a captured stream when the service reaps the leader, so the client sees everything the leader wrote before it exited, as a native parent does once `waitpid` returns. Output written later, by processes that still hold the stream, can follow `Exited`, as it does natively. +- `Exited` never waits for an acknowledgement: when the [replay limit](#output-replay-and-flow-control) stops the service reading, `Exited` follows the output read until then, and the rest follows `Exited`. With a PTY, Linux passes terminal output to the master through an asynchronous kernel queue, so output still in that queue when the service reaps the leader can follow `Exited`. - `OutputClosed` follows every output byte the service will deliver, and each `StreamClosed`. - Output of different streams has no order relative to each other or to changes in the file system. From 9cacf91da188e85c0909bbaec02473793703b72a Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 1 Oct 2026 02:31:48 +0000 Subject: [PATCH 2/3] Keep output reads non-blocking and packets whole before Exited The launch now fails with EffectNone unless each output descriptor is non-blocking and watched by the runtime poller, so a read under op.mu never blocks. A read always has at least a page of room under the replay limit, so a packet-mode pipe never loses the rest of a packet; with less room the reader waits for acknowledgements and the exit gate treats it as the limit. The ordering rule now names its exceptions: an abandoned or lost stream, output held back by the replay limit, and the PTY residual. The ordering tests hold the readers through a test seam and use a FIFO handshake instead of scheduling and sleeps. --- .../internal/processservice/launch.go | 46 +++++++++-- .../internal/processservice/operation.go | 13 +++- .../internal/processservice/service.go | 2 + .../internal/processservice/service_test.go | 78 ++++++++++++------- docs/process-protocol.md | 4 +- 5 files changed, 101 insertions(+), 42 deletions(-) diff --git a/apps/sandboxio/internal/processservice/launch.go b/apps/sandboxio/internal/processservice/launch.go index ad9507de..951879f0 100644 --- a/apps/sandboxio/internal/processservice/launch.go +++ b/apps/sandboxio/internal/processservice/launch.go @@ -9,6 +9,7 @@ import ( "io" "os" "syscall" + "time" "github.com/creack/pty" "golang.org/x/sys/unix" @@ -114,6 +115,11 @@ func (op *operation) spawn(spec sp.ProcessSpec) (l launched, f *sp.Failure) { l.streams = append(l.streams, &stream{name: name, f: r}) } } + for _, st := range l.streams { + if err := watch(st.f); err != nil { + return l, ioFail("watch output", err) + } + } attr := &os.ProcAttr{ Env: []string{trampolineEnv}, @@ -225,22 +231,46 @@ func openPTY(spec sp.PTYSpec) (master, tty *os.File, err error) { return os.NewFile(uintptr(fd), "/dev/ptmx"), tty, nil } +// watch makes f non-blocking and checks that the runtime poller watches it, +// so a read of f under op.mu never blocks. Registering f fails, for example, +// at the epoll watch limit. +func watch(f *os.File) error { + rc, err := f.SyscallConn() + if err != nil { + return err + } + if cerr := rc.Control(func(fd uintptr) { err = unix.SetNonblock(int(fd), true) }); cerr != nil { + return cerr + } + if err != nil { + return err + } + return f.SetReadDeadline(time.Time{}) // os.ErrNoDeadline when the poller does not watch f +} + +// minRead is the least room a read has under the replay limit. A pipe in +// packet mode (O_DIRECT) holds packets of up to a page, and a shorter read +// discards the rest of the packet. +var minRead = os.Getpagesize() + // read captures one stream until end of file, a read error, or CloseOutput. -// It reads only while the retained output is below the replay limit, and at -// most the space left, so the streams together never exceed it. Each read -// and the push of what it read happen under op.mu, so every byte is either -// still in the kernel's buffer or in an event; markLocked relies on that. +// It reads only while minRead bytes remain under the replay limit, and at +// most the room left, so the streams together never exceed it. Each read and +// the push of what it read happen under op.mu, so every byte is either still +// in the kernel's buffer or in an event; markLocked relies on that. func (op *operation) read(st *stream) { - limit := int(op.s.caps.MaxReplayBytesPerOperation) buf := make([]byte, op.s.caps.MaxDataBytes) var disp sp.OutputDisposition rc, err := st.f.SyscallConn() for err == nil && disp == 0 { op.mu.Lock() - for op.retained >= limit && !st.abandoned { + for op.room() < minRead && !st.abandoned { op.cond.Wait() } op.mu.Unlock() + if op.s.beforeRead != nil { + op.s.beforeRead() + } // The callback returns false to wait until the stream is readable. err = rc.Read(func(fd uintptr) bool { op.mu.Lock() @@ -280,11 +310,11 @@ func (op *operation) read(st *stream) { // limit stops it, and pushes what it read. It returns the disposition once // the stream ends, and empty when nothing is buffered. func (op *operation) readLocked(st *stream, fd int, buf []byte) (disp sp.OutputDisposition, empty bool) { - room := int(op.s.caps.MaxReplayBytesPerOperation) - op.retained + room := op.room() switch { case st.abandoned: return sp.OutputAbandoned, false - case room <= 0: + case room < minRead: return 0, false } buf = buf[:min(len(buf), room)] diff --git a/apps/sandboxio/internal/processservice/operation.go b/apps/sandboxio/internal/processservice/operation.go index f8665c01..c728810c 100644 --- a/apps/sandboxio/internal/processservice/operation.go +++ b/apps/sandboxio/internal/processservice/operation.go @@ -187,14 +187,14 @@ func (op *operation) markLocked() { } // exitWhenDrainedLocked reports the reaped leader's exit once every stream -// has delivered its mark or closed. At the replay limit no stream is read -// until the client acknowledges, so the exit is reported then: it never -// waits for an acknowledgement. +// has delivered its mark or closed. With less than minRead of room under the +// replay limit no stream is read until the client acknowledges, so the exit +// is reported then: it never waits for an acknowledgement. func (op *operation) exitWhenDrainedLocked() { if !op.leaderGone || op.state != sp.StateRunning { return } - if op.retained < int(op.s.caps.MaxReplayBytesPerOperation) { + if op.room() >= minRead { for _, st := range op.streams { if !st.closed && !st.abandoned && st.offset < st.mark { return @@ -204,6 +204,11 @@ func (op *operation) exitWhenDrainedLocked() { op.exitedLocked() } +// room is how much more output the replay limit lets the operation retain. +func (op *operation) room() int { + return int(op.s.caps.MaxReplayBytesPerOperation) - op.retained +} + // exitedLocked reports the reaped leader's exit and starts watching the scope. func (op *operation) exitedLocked() { ws := op.status diff --git a/apps/sandboxio/internal/processservice/service.go b/apps/sandboxio/internal/processservice/service.go index abf83d5e..f0c525b9 100644 --- a/apps/sandboxio/internal/processservice/service.go +++ b/apps/sandboxio/internal/processservice/service.go @@ -57,6 +57,8 @@ type Service struct { // onExpire runs after an expired owner-loss grace has decided the // cleanup, before the operations are cancelled; tests restore there. onExpire func() + // beforeRead runs before each output read; tests hold the readers there. + beforeRead func() mu sync.Mutex ops map[opKey]*operation diff --git a/apps/sandboxio/internal/processservice/service_test.go b/apps/sandboxio/internal/processservice/service_test.go index 261baab3..3ee75eca 100644 --- a/apps/sandboxio/internal/processservice/service_test.go +++ b/apps/sandboxio/internal/processservice/service_test.go @@ -10,7 +10,6 @@ import ( "net" "os" "path/filepath" - "runtime" "strconv" "strings" "sync/atomic" @@ -155,6 +154,23 @@ func (h *harness) leader(op *sp.Operation) int { return h.svc.ops[opKey{h.att, op.Ref().OperationID}].sid() } +// waitReaped waits until the reaper has reaped the operation's leader. +func (h *harness) waitReaped(op *sp.Operation) { + h.t.Helper() + h.svc.mu.Lock() + o := h.svc.ops[opKey{h.att, op.Ref().OperationID}] + h.svc.mu.Unlock() + for deadline := time.Now().Add(10 * time.Second); time.Now().Before(deadline); time.Sleep(time.Millisecond) { + o.mu.Lock() + gone := o.leaderGone + o.mu.Unlock() + if gone { + return + } + } + h.t.Fatal("the leader was never reaped") +} + // waitMembers waits until the operation's session has n live processes named // comm. func (h *harness) waitMembers(op *sp.Operation, comm string, n int) { @@ -229,19 +245,18 @@ func TestStartFailure(t *testing.T) { } } -// Every byte the leader wrote before it exited precedes Exited, even when -// it wrote more than a pipe buffer just before exiting. With one P the reaper -// often runs before the reader has drained the pipe. +// Exited follows every byte the leader left buffered: the readers wait until +// the leader is reaped, so all its output is still in the pipes then. func TestOutputPrecedesExited(t *testing.T) { - defer runtime.GOMAXPROCS(runtime.GOMAXPROCS(1)) h := newHarness(t, DefaultConfig()) - c := h.connect() - const size = 200 << 10 - for range 20 { - evs := events(t, h.start(c, pipeSpec("head", "-c", strconv.Itoa(size), "/dev/zero")), sp.EventExited) - if got := len(output(evs, sp.StreamStdout)); got != size { - t.Fatalf("%d of %d bytes before Exited", got, size) - } + hold := make(chan struct{}) + h.svc.beforeRead = func() { <-hold } + op := h.start(h.connect(), pipeSpec("sh", "-c", "printf out; printf err >&2")) + h.waitReaped(op) + close(hold) + evs := events(t, op, sp.EventExited) + if output(evs, sp.StreamStdout) != "out" || output(evs, sp.StreamStderr) != "err" { + t.Fatalf("events %v", evs) } } @@ -249,31 +264,38 @@ func TestOutputPrecedesExited(t *testing.T) { // and keeps OutputClosed pending until it ends. func TestLaterOutputFollowsExited(t *testing.T) { h := newHarness(t, DefaultConfig()) - evs := events(t, h.start(h.connect(), pipeSpec("sh", "-c", "printf a; (sleep 1; printf b) & exit 0")), sp.EventOutputClosed) - _, exited := find[sp.ExitedEvent](t, evs) - closed := evs[len(evs)-1].(sp.OutputClosedEvent) - if output(evs[:exited], sp.StreamStdout) != "a" || output(evs[exited:], sp.StreamStdout) != "b" || closed.Disposition != sp.OutputDrained { - t.Fatalf("events %v", evs) + fifo := filepath.Join(t.TempDir(), "fifo") + if err := unix.Mkfifo(fifo, 0o600); err != nil { + t.Fatal(err) + } + op := h.start(h.connect(), pipeSpec("sh", "-c", `printf a; (read x <"$0"; printf b) & exit 0`, fifo)) + before := events(t, op, sp.EventExited) + if err := os.WriteFile(fifo, []byte("\n"), 0); err != nil { + t.Fatal(err) + } + after := events(t, op, sp.EventOutputClosed) + if output(before, sp.StreamStdout) != "a" || output(after, sp.StreamStdout) != "b" || after[len(after)-1].(sp.OutputClosedEvent).Disposition != sp.OutputDrained { + t.Fatalf("events %v then %v", before, after) } } -// Exited never waits for an acknowledgement: at the replay limit it follows -// the output already read, and the rest arrives once acknowledged. -func TestExitedAtReplayLimit(t *testing.T) { +// A pipe in packet mode keeps every packet whole at the replay limit, and +// Exited does not wait for the acknowledgement that lets the rest be read. +func TestPacketsAtReplayLimit(t *testing.T) { cfg := DefaultConfig() - cfg.MaxReplayBytesPerOperation = sandboxwire.MaxChunk + cfg.MaxReplayBytesPerOperation = sandboxwire.MaxChunk + 1000 // not a whole number of packets h := newHarness(t, cfg) - const size = sandboxwire.MaxChunk + 1000 // the rest fits in the pipe - op := h.start(h.connect(), pipeSpec("head", "-c", strconv.Itoa(size), "/dev/zero")) + const packet, packets = 4096, 20 + op := h.start(h.connect(), pipeSpec("dd", "if=/dev/zero", "bs="+strconv.Itoa(packet), "count="+strconv.Itoa(packets), "oflag=direct", "status=none")) evs := events(t, op, sp.EventExited) - if got := len(output(evs, sp.StreamStdout)); got != sandboxwire.MaxChunk { - t.Fatalf("%d bytes before Exited", got) - } + before := len(output(evs, sp.StreamStdout)) if err := op.Ack(context.Background(), evs[len(evs)-1].Header().Sequence); err != nil { t.Fatal(err) } - if got := len(output(events(t, op, sp.EventOutputClosed), sp.StreamStdout)); got != size-sandboxwire.MaxChunk { - t.Fatalf("%d bytes after Exited", got) + evs = events(t, op, sp.EventOutputClosed) + after, closed := len(output(evs, sp.StreamStdout)), evs[len(evs)-1].(sp.OutputClosedEvent) + if before != sandboxwire.MaxChunk || before+after != packet*packets || closed.Disposition != sp.OutputDrained { + t.Fatalf("%d bytes before Exited and %d after, output %v", before, after, closed.Disposition) } } diff --git a/docs/process-protocol.md b/docs/process-protocol.md index 9cb2267b..5791eb67 100644 --- a/docs/process-protocol.md +++ b/docs/process-protocol.md @@ -86,8 +86,8 @@ Ordering: - `Started` or `StartFailed` comes first. - `Exited` and `OutputClosed` are independent. A background process holding a stream open keeps `OutputClosed` pending after `Exited`, and output can close before the leader exits. -- `Exited` follows every output byte buffered in a captured stream when the service reaps the leader, so the client sees everything the leader wrote before it exited, as a native parent does once `waitpid` returns. Output written later, by processes that still hold the stream, can follow `Exited`, as it does natively. -- `Exited` never waits for an acknowledgement: when the [replay limit](#output-replay-and-flow-control) stops the service reading, `Exited` follows the output read until then, and the rest follows `Exited`. With a PTY, Linux passes terminal output to the master through an asynchronous kernel queue, so output still in that queue when the service reaps the leader can follow `Exited`. +- `Exited` follows every output byte buffered in a captured stream when the service reaps the leader, and never waits for an acknowledgement. Exceptions: a stream abandoned by `CloseOutput` or lost to a read failure delivers no more output; output that the [replay limit](#output-replay-and-flow-control) keeps the service from reading follows `Exited`; with a PTY, output still in the kernel's asynchronous queue to the master at the reap can follow `Exited`, and a terminal flush discards output as it does natively. +- Output written after the reap, by processes that still hold a stream, can follow `Exited`, as it does natively. - `OutputClosed` follows every output byte the service will deliver, and each `StreamClosed`. - Output of different streams has no order relative to each other or to changes in the file system. From c71f2a5b856abb9b6315cb6af28c795d385406e8 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 1 Oct 2026 02:33:36 +0000 Subject: [PATCH 3/3] State the page of room the replay limit keeps --- docs/process-protocol.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/process-protocol.md b/docs/process-protocol.md index 5791eb67..93009bbf 100644 --- a/docs/process-protocol.md +++ b/docs/process-protocol.md @@ -122,7 +122,7 @@ The stdin offset counts the bytes the service has accepted, from 0. `WriteStdin` ### Output, replay and flow control - The service retains each operation's events until they are acknowledged, and delivers them to the observer in order. -- Unacknowledged `Output` data of all the operation's streams together is limited to `MaxReplayBytesPerOperation`. At the limit the service stops reading the process's output, so the process blocks on its own writes. Nothing is discarded. +- Unacknowledged `Output` data of all the operation's streams together is limited to `MaxReplayBytesPerOperation`. When less than one memory page of the limit remains, the service stops reading the process's output, so the process blocks on its own writes; a read never splits a packet-mode pipe write. Nothing is discarded. - A slow observer holds back its stream: the service writes events only as fast as the peer reads them. - `Attach` resumes after any sequence in the retained range. A request for acknowledged events returns `ReplayGap`, so missing output is never skipped silently. The events an accepted `Attach` promised stay retained until they are sent, even when another stream acknowledges them first.