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
110 changes: 33 additions & 77 deletions apps/daemon/internal/agenthost/launch_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ type runningView interface {
Wait() (sessionview.Exit, error)
Close() error
Relay() *os.File
RelayLost() <-chan struct{}
Spawn(ctx context.Context, path string, args, env []string, dir string, stdin bool) (*sessionview.Spawned, error)
}

Expand Down Expand Up @@ -127,12 +128,8 @@ func (s *session) spawn(opts clirunner.StartOptions) (*clirunner.Process, error)
}
return nil, fmt.Errorf("%w: spawn: %w", ErrLaunch, err)
}
var stdin io.WriteCloser
if p.Stdin != nil {
stdin = p.Stdin
}
// FromHandle fails only without stdout and stderr, which a spawned process always has.
process, _ := clirunner.FromHandle(p, clirunner.HandleOptions{Parent: opts.Parent, Stdin: stdin, Stdout: p.Stdout, Stderr: p.Stderr, KillTimeout: opts.KillTimeout})
process, _ := clirunner.FromHandle(p, clirunner.HandleOptions{Parent: opts.Parent, Stdin: writer(p.Stdin), Stdout: p.Stdout, Stderr: p.Stderr, KillTimeout: opts.KillTimeout})
return process, nil
}

Expand Down Expand Up @@ -169,26 +166,26 @@ func (s *session) start(lv *liveView, opts clirunner.StartOptions) (*clirunner.P
s.release(lv)
return nil, fmt.Errorf("%w: home: %w", ErrLaunch, err)
}
ends, err := newStdio(opts.NeedStdin)
child, ends, err := sessionview.Stdio([3]*os.File{}, opts.NeedStdin, s.uid, s.uid)
if err != nil {
s.release(lv)
return nil, fmt.Errorf("%w: stdio: %w", ErrLaunch, err)
}
world := worldfs.New(worldExport, s.openFile)
spec := s.spec(world, opts, ends)
spec := s.spec(world, opts, child)
// The gateway serves from the view's network hook until the view has ended.
var stopGateway func()
spec.Network.Setup = func(netns *os.File) (err error) {
stopGateway, err = gateway.Start(netns, s.plan.gateway)
return err
}
v, err := sessionview.Start(startCtx, spec)
ends.closeChild()
closeFiles(child[:])
if err != nil {
if stopGateway != nil {
stopGateway()
}
ends.closeParent()
closeFiles(ends[:])
defer s.release(lv)
// sessionview stops a world that served; Stop reports how that went.
serr := world.Stop()
Expand Down Expand Up @@ -241,7 +238,7 @@ func (s *session) brokerFailed(op string, err error) error {
// own hands a started view to the clirunner.Process the adapter receives. A
// view with a relay gets its own process broker, which serves the view's
// shims in scope until the view ends.
func (s *session) own(lv *liveView, v runningView, world viewWorld, stopGateway func(), opts clirunner.StartOptions, ends *stdio, scope sandboxprocess.Scope) (*clirunner.Process, error) {
func (s *session) own(lv *liveView, v runningView, world viewWorld, stopGateway func(), opts clirunner.StartOptions, ends [3]*os.File, scope sandboxprocess.Scope) (*clirunner.Process, error) {
h := &ownedView{s: s, lv: lv, v: v, world: world, stopGateway: stopGateway, ended: make(chan struct{}), watched: make(chan struct{})}
var brokerErr error
if relay := v.Relay(); relay != nil {
Expand Down Expand Up @@ -269,16 +266,16 @@ func (s *session) own(lv *liveView, v runningView, world viewWorld, stopGateway
}
var process *clirunner.Process
if err == nil {
process, err = clirunner.FromHandle(h, clirunner.HandleOptions{Parent: opts.Parent, Stdin: ends.stdin(),
Stdout: ends.parent[1], Stderr: ends.parent[2], KillTimeout: opts.KillTimeout})
process, err = clirunner.FromHandle(h, clirunner.HandleOptions{Parent: opts.Parent, Stdin: writer(ends[0]),
Stdout: ends[1], Stderr: ends[2], KillTimeout: opts.KillTimeout})
if err != nil {
err = fmt.Errorf("%w: %w", ErrLaunch, err)
}
}
if err != nil {
v.Close()
h.Wait()
ends.closeParent()
closeFiles(ends[:])
return nil, err
}
return process, nil
Expand Down Expand Up @@ -306,33 +303,27 @@ type ownedView struct {
// while the view runs.
func (h *ownedView) watch() {
defer close(h.watched)
var relayEnded <-chan struct{}
var relayLost <-chan struct{}
if h.broker != nil {
relayEnded = h.broker.Done()
}
for {
relayLost = h.v.RelayLost()
}
select {
case <-h.world.Lost():
h.s.fail(fmt.Errorf("%w: world: %w", ErrWorld, h.world.Err()))
return
case <-relayLost:
case <-h.ended:
// A lost relay is reported before the view ends.
select {
case <-h.world.Lost():
h.s.fail(fmt.Errorf("%w: world: %w", ErrWorld, h.world.Err()))
return
case <-relayEnded:
// The broker stops serving on Close, which end calls only after
// watch returns, or when its relay connection ends. sessionview
// ends that connection itself only in its teardown, which starts
// once the launcher has stopped answering, and from then on
// Signal fails with ErrExited, ErrClosed or a lost launcher. So a
// delivered signal means that the view still runs its process
// and the relay was lost while the Harness ran. A failed one
// means that the view is ending, and Wait reports how.
if h.v.Signal(0) == nil {
h.s.brokerFailed("relay", h.broker.Err())
return
}
relayEnded = nil
case <-h.ended:
case <-relayLost:
default:
return
}
}
// The relay's end ends its connection, so the broker stops.
<-h.broker.Done()
h.s.log.Error("process relay lost", "error", h.broker.Err())
h.s.brokerFailed("relay", h.broker.Err())
}

func (h *ownedView) Signal(sig syscall.Signal) error { return h.v.Signal(sig) }
Expand Down Expand Up @@ -378,7 +369,7 @@ func (h *ownedView) end(waitErr error) {
// spec builds the view: the closure and home directories, the agent
// host's /etc files and CA directory, the adapter's overlays and masks, and
// the shim.
func (s *session) spec(world *worldfs.World, opts clirunner.StartOptions, ends *stdio) sessionview.Spec {
func (s *session) spec(world *worldfs.World, opts clirunner.StartOptions, stdio [3]*os.File) sessionview.Spec {
view := s.plan.view
var private []sessionview.PrivateDir
for _, m := range view.Closure {
Expand Down Expand Up @@ -406,57 +397,22 @@ func (s *session) spec(world *worldfs.World, opts clirunner.StartOptions, ends *
Overlays: overlays,
Shim: sessionview.Shim{Binary: s.cfg.Shim, Names: view.Shims, Paths: view.ShimPaths},
Process: sessionview.Process{Path: opts.Binary, Args: append([]string{opts.Binary}, opts.Args...), Env: opts.Env,
Dir: opts.Dir, UID: s.uid, GID: s.uid, Stdin: ends.child[0], Stdout: ends.child[1], Stderr: ends.child[2],
Dir: opts.Dir, UID: s.uid, GID: s.uid, Stdin: stdio[0], Stdout: stdio[1], Stderr: stdio[2],
Grace: opts.KillTimeout},
StagingParent: s.dir.entry(stagingEntry),
CgroupParent: s.cfg.ViewCgroups,
}
}

// stdio holds the view process's stdio: the child ends sessionview passes to
// the process and the parent ends the clirunner.Process owns. Without a stdin
// pipe the child's stdin is /dev/null and there is no parent end.
type stdio struct {
child, parent [3]*os.File
}

func newStdio(needStdin bool) (*stdio, error) {
e := &stdio{}
if needStdin {
r, w, err := os.Pipe()
if err != nil {
return nil, err
}
e.child[0], e.parent[0] = r, w
} else {
null, err := os.Open(os.DevNull)
if err != nil {
return nil, err
}
e.child[0] = null
}
for i := 1; i < 3; i++ {
r, w, err := os.Pipe()
if err != nil {
e.closeChild()
e.closeParent()
return nil, err
}
e.child[i], e.parent[i] = w, r
}
return e, nil
}

func (e *stdio) stdin() io.WriteCloser {
if e.parent[0] == nil {
// writer returns f, or nil without f: a nil *os.File would be a writer
// that fails.
func writer(f *os.File) io.WriteCloser {
if f == nil {
return nil
}
return e.parent[0]
return f
}

func (e *stdio) closeChild() { closeFiles(e.child[:]) }
func (e *stdio) closeParent() { closeFiles(e.parent[:]) }

func closeFiles(files []*os.File) {
for _, f := range files {
if f != nil {
Expand Down
10 changes: 6 additions & 4 deletions apps/daemon/internal/agenthost/session_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,11 +54,11 @@ func TestViewEndReleasesTheSlotBeforeTheProcessEnds(t *testing.T) {
lv := &liveView{}
s.live = lv
s.views.Add(1)
ends, err := newStdio(false)
child, ends, err := sessionview.Stdio([3]*os.File{}, false, uint32(os.Getuid()), uint32(os.Getgid()))
if err != nil {
t.Fatal(err)
}
ends.closeChild()
closeFiles(child[:])
v := &fakeView{exit: make(chan struct{})}
p, err := s.own(lv, v, fakeWorld{stop: stop}, func() {}, clirunner.StartOptions{Parent: context.Background(), KillTimeout: time.Second}, ends, 0)
if err != nil {
Expand Down Expand Up @@ -108,11 +108,11 @@ func TestViewEndStopsTheGateway(t *testing.T) {
lv := &liveView{}
s.live = lv
s.views.Add(1)
ends, err := newStdio(false)
child, ends, err := sessionview.Stdio([3]*os.File{}, false, uint32(os.Getuid()), uint32(os.Getgid()))
if err != nil {
t.Fatal(err)
}
ends.closeChild()
closeFiles(child[:])
v := &fakeView{exit: make(chan struct{})}
p, err := s.own(lv, v, fakeWorld{}, stop, clirunner.StartOptions{Parent: context.Background(), KillTimeout: time.Second}, ends, 0)
if err != nil {
Expand Down Expand Up @@ -417,6 +417,8 @@ func (v *fakeView) Signal(syscall.Signal) error { return nil }

func (v *fakeView) Relay() *os.File { return nil }

func (v *fakeView) RelayLost() <-chan struct{} { return nil }

func (v *fakeView) Spawn(context.Context, string, []string, []string, string, bool) (*sessionview.Spawned, error) {
return nil, v.spawnErr
}
Expand Down
7 changes: 2 additions & 5 deletions apps/daemon/internal/processbroker/broker_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -112,15 +112,12 @@ func (b *Broker) read() {
}
}
b.mu.Lock()
lost := !b.closing
if lost {
if !b.closing {
b.err = fmt.Errorf("%w: %w", ErrRelayLost, err)
}
b.mu.Unlock()
if lost {
b.log.Error("process relay lost", "error", err)
}
b.cancel()
b.conn.CloseWrite() // a relay the broker stops serving ends too
b.conn.Close()
close(b.done)
}
Expand Down
Loading