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
143 changes: 105 additions & 38 deletions apps/sandboxio/internal/processservice/launch.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"io"
"os"
"syscall"
"time"

"github.com/creack/pty"
"golang.org/x/sys/unix"
Expand All @@ -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()
Expand All @@ -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)
Expand Down Expand Up @@ -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},
Expand Down Expand Up @@ -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
Expand All @@ -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))
}
71 changes: 60 additions & 11 deletions apps/sandboxio/internal/processservice/operation.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
Expand All @@ -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 {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down
2 changes: 2 additions & 0 deletions apps/sandboxio/internal/processservice/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading