diff --git a/apps/sandboxio/internal/processservice/launch.go b/apps/sandboxio/internal/processservice/launch.go index 1374bfa9..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" @@ -18,7 +19,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 +39,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) @@ -116,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}, @@ -227,52 +231,68 @@ 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. -// 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) { - limit := int(op.s.caps.MaxReplayBytesPerOperation) - buf := make([]byte, chunk) +// 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) { + 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.room() < minRead && !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 + if op.s.beforeRead != nil { + op.s.beforeRead() } - 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 +305,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 := op.room() + switch { + case st.abandoned: + return sp.OutputAbandoned, false + case room < minRead: + 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..c728810c 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,49 @@ 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. 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.room() >= minRead { + for _, st := range op.streams { + if !st.closed && !st.abandoned && st.offset < st.mark { + return + } + } + } + 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 @@ -358,26 +398,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.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 49f93a20..3ee75eca 100644 --- a/apps/sandboxio/internal/processservice/service_test.go +++ b/apps/sandboxio/internal/processservice/service_test.go @@ -154,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) { @@ -228,17 +245,60 @@ func TestStartFailure(t *testing.T) { } } -// A background job holding stderr keeps output open after the leader exits. -func TestExitedBeforeOutputClosed(t *testing.T) { +// 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) { 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) - _, 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" { + 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) } } +// 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()) + 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) + } +} + +// 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 + 1000 // not a whole number of packets + h := newHarness(t, cfg) + 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) + before := len(output(evs, sp.StreamStdout)) + if err := op.Ack(context.Background(), evs[len(evs)-1].Header().Sequence); err != nil { + t.Fatal(err) + } + 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) + } +} + 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..93009bbf 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, 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. @@ -120,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.