From 34a4294c1221bf83a2565fbd050475aff3add779 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 9 Oct 2026 08:24:34 +0000 Subject: [PATCH 1/2] Contain sandbox process descendants in delegated cgroups --- .github/workflows/check.yml | 5 + .../internal/processservice/cgroup.go | 211 ++++++++ .../internal/processservice/cgroup_test.go | 450 ++++++++++++++++++ .../internal/processservice/launch.go | 31 +- .../internal/processservice/operation.go | 11 +- .../internal/processservice/scope.go | 45 +- .../internal/processservice/service.go | 20 +- docs/process-protocol.md | 10 +- docs/zh/process-protocol.md | 12 +- 9 files changed, 771 insertions(+), 24 deletions(-) create mode 100644 apps/sandboxio/internal/processservice/cgroup.go create mode 100644 apps/sandboxio/internal/processservice/cgroup_test.go diff --git a/.github/workflows/check.yml b/.github/workflows/check.yml index b44cd4c1c..e87bb0a32 100644 --- a/.github/workflows/check.yml +++ b/.github/workflows/check.yml @@ -129,6 +129,11 @@ jobs: - name: Test Runtime and shared Go contracts if: matrix.part == 'runtime' run: make check-go + - name: Test delegated process scopes + if: matrix.part == 'runtime' + run: | + go test -race -c -o "$RUNNER_TEMP/process-scope.test" ./apps/sandboxio/internal/processservice + sudo systemd-run --quiet --collect --wait --pipe --unit="oac-process-scope-${GITHUB_RUN_ID}-${GITHUB_RUN_ATTEMPT}" --property=Delegate=yes --uid="$(id -u)" --gid="$(id -g)" /usr/bin/env OAC_TEST_PROCESS_CGROUP=1 "$RUNNER_TEMP/process-scope.test" -test.run '^TestCgroup' -test.v - name: Test microsandbox provider and Linux helper if: matrix.part == 'runtime' run: make check-microsandbox-provider diff --git a/apps/sandboxio/internal/processservice/cgroup.go b/apps/sandboxio/internal/processservice/cgroup.go new file mode 100644 index 000000000..218e25c18 --- /dev/null +++ b/apps/sandboxio/internal/processservice/cgroup.go @@ -0,0 +1,211 @@ +//go:build linux + +package processservice + +import ( + "errors" + "fmt" + "io/fs" + "os" + "path/filepath" + "slices" + "strconv" + "strings" + "sync" + + "golang.org/x/sys/unix" + + sp "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxprocess" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +// delegatedCgroup advertises the stronger scope only after exercising the +// Provider's delegation. A missing or read-only delegation leaves POSIX scope +// available; callers still select the scope explicitly from Describe. +func delegatedCgroup() string { + raw, err := os.ReadFile("/proc/self/cgroup") + if err != nil { + return "" + } + for line := range strings.Lines(string(raw)) { + path, ok := strings.CutPrefix(strings.TrimSpace(line), "0::") + if !ok || !filepath.IsAbs(path) || filepath.Clean(path) != path { + continue + } + parent := filepath.Join("/sys/fs/cgroup", path) + var stat unix.Statfs_t + if unix.Statfs(parent, &stat) != nil || stat.Type != unix.CGROUP2_SUPER_MAGIC { + return "" + } + // This also rejects a mount whose root does not match the cgroup namespace. + procs, err := os.ReadFile(filepath.Join(parent, "cgroup.procs")) + if err != nil || !slices.Contains(strings.Fields(string(procs)), strconv.Itoa(os.Getpid())) { + return "" + } + g, err := newProcessCgroup(parent) + if err != nil { + return "" + } + err = g.signal(unix.SIGKILL) // the exclusive probe group is empty + _, observation := g.observe() + // No process was ever placed in this exclusive probe group. + g.closed = true + cleanup := g.cleanup() + if err == nil && observation == nil && cleanup == nil { + return parent + } + return "" + } + return "" +} + +// processCgroup owns one operation's cgroup. It is containment for process +// lifetime, not isolation from other processes running as the same account. +// Its descriptor pins the group while a signal races scope removal. +type processCgroup struct { + mu sync.Mutex + root *os.Root + path string + closed bool +} + +func newProcessCgroup(parent string) (*processCgroup, error) { + path, err := os.MkdirTemp(parent, "oac-process-") + if err != nil { + return nil, err + } + root, err := os.OpenRoot(path) + if err != nil { + _ = os.Remove(path) + return nil, err + } + return &processCgroup{root: root, path: path}, nil +} + +func (g *processCgroup) place(pid int) error { + g.mu.Lock() + defer g.mu.Unlock() + return g.root.WriteFile("cgroup.procs", []byte(strconv.Itoa(pid)), 0) +} + +func (g *processCgroup) signal(sig unix.Signal) error { + g.mu.Lock() + defer g.mu.Unlock() + if g.closed { + return notRunning("the cgroup is empty") + } + if sig == unix.SIGKILL { + if err := g.root.WriteFile("cgroup.kill", []byte("1"), 0); err != nil { + return sp.Fail(sp.CodeIO, sandboxwire.EffectPossible, "kill cgroup: %v", err) + } + return nil + } + // Membership comes from this operation's cgroup, never a process-tree scan. + // A pidfd plus a second membership read excludes a PID reused outside it. + var errs []error + sent := 0 + err := fs.WalkDir(g.root.FS(), ".", func(path string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if !d.IsDir() { + return nil + } + procs, err := g.root.ReadFile(filepath.Join(path, "cgroup.procs")) + if err != nil { + return err + } + for _, word := range strings.Fields(string(procs)) { + pid, err := strconv.Atoi(word) + if err != nil { + return err + } + fd, err := unix.PidfdOpen(pid, 0) + if errors.Is(err, unix.ESRCH) { + continue + } + if err != nil { + return err + } + current, err := g.root.ReadFile(filepath.Join(path, "cgroup.procs")) + if err == nil && slices.Contains(strings.Fields(string(current)), word) { + err = unix.PidfdSendSignal(fd, sig, nil, 0) + if err == nil { + sent++ + } + } + unix.Close(fd) + if err != nil && !errors.Is(err, unix.ESRCH) { + errs = append(errs, err) + } + } + return nil + }) + return signalResult(sent, err, errors.Join(errs...), "the cgroup is empty") +} + +// observe reports closed only after the kernel reports the whole subtree +// empty. An observation failure keeps the operation unknown for the next poll. +func (g *processCgroup) observe() (int, error) { + g.mu.Lock() + defer g.mu.Unlock() + if g.closed { + return 0, nil + } + raw, err := g.root.ReadFile("cgroup.events") + if err != nil { + return 0, err + } + empty := false + for line := range strings.Lines(string(raw)) { + switch strings.TrimSpace(line) { + case "populated 1": + return 1, nil + case "populated 0": + empty = true + } + } + if !empty { + return 0, fmt.Errorf("cgroup.events has no populated state") + } + g.closed = true + return 0, nil +} + +// cleanup removes only an already-empty owned subtree. A removal failure is +// distinct from an unknown process scope; Shutdown makes one more attempt. +func (g *processCgroup) cleanup() error { + g.mu.Lock() + defer g.mu.Unlock() + if g.root == nil { + return nil + } + if !g.closed { + return errors.New("the process cgroup is not confirmed empty") + } + + var dirs []string + err := fs.WalkDir(g.root.FS(), ".", func(path string, d fs.DirEntry, err error) error { + if err != nil { + return err + } + if d.IsDir() && path != "." { + dirs = append(dirs, path) + } + return nil + }) + if err != nil { + return err + } + for _, path := range slices.Backward(dirs) { + if err := g.root.Remove(path); err != nil { + return err + } + } + if err := os.Remove(g.path); err != nil { + return err + } + err = g.root.Close() + g.root = nil + return err +} diff --git a/apps/sandboxio/internal/processservice/cgroup_test.go b/apps/sandboxio/internal/processservice/cgroup_test.go new file mode 100644 index 000000000..4b87461bd --- /dev/null +++ b/apps/sandboxio/internal/processservice/cgroup_test.go @@ -0,0 +1,450 @@ +//go:build linux + +package processservice + +import ( + "context" + "errors" + "io/fs" + "os" + "path/filepath" + "slices" + "strconv" + "strings" + "testing" + "time" + + "golang.org/x/sys/unix" + + sp "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxprocess" + "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxwire" +) + +// Run the test binary in a delegated cgroup, for example with systemd-run +// --user --property=Delegate=yes. The gate makes missing delegation a failure +// in the qualification job rather than silently skipping that coverage. +func cgroupHarness(t *testing.T) *harness { + t.Helper() + h := newHarness(t, DefaultConfig()) + if !slices.Contains(h.svc.caps.Scopes, sp.ScopeCgroupV2) { + if os.Getenv("OAC_TEST_PROCESS_CGROUP") == "1" { + t.Fatal("a writable cgroup v2 delegation with cgroup.kill is required") + } + t.Skip("run in a delegated cgroup with OAC_TEST_PROCESS_CGROUP=1") + } + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + h.svc.Shutdown(ctx) + }) + return h +} + +func cgroupSpec(argv ...string) sp.ProcessSpec { + s := pipeSpec(argv...) + s.Scope = sp.ScopeCgroupV2 + return s +} + +func awaitFile(t *testing.T, path string, want string) { + t.Helper() + for deadline := time.Now().Add(10 * time.Second); ; time.Sleep(time.Millisecond) { + b, err := os.ReadFile(path) + if err == nil && string(b) == want { + return + } + if time.Now().After(deadline) { + t.Fatalf("file %s did not reach %q: %q %v", path, want, b, err) + } + } +} + +func cgroupRecord(t *testing.T, h *harness, op *sp.Operation) *operation { + t.Helper() + h.svc.mu.Lock() + defer h.svc.mu.Unlock() + return h.svc.ops[opKey{h.att, op.Ref().OperationID}] +} + +// The descendant writes only after the observer opens the gate. Cancelling +// its cgroup before that gate must make any later effect impossible, even +// after both its original process group and parent have gone. +func TestCgroupCancelDetachedDescendants(t *testing.T) { + for _, kind := range []string{"fork", "setsid", "double-fork"} { + for _, leaderExit := range []bool{false, true} { + t.Run(kind+"/leader-exited="+strconv.FormatBool(leaderExit), func(t *testing.T) { + h := cgroupHarness(t) + dir := t.TempDir() + writer := filepath.Join(dir, "writer") + script := "#!/bin/sh\ntrap '' TERM HUP\necho $$ > \"$1/pid\"\nprintf a > \"$1/ticks\"\nwhile [ ! -e \"$1/gate\" ]; do sleep .01; done\nprintf b >> \"$1/ticks\"\n" + if err := os.WriteFile(writer, []byte(script), 0700); err != nil { + t.Fatal(err) + } + launch := `"$1" "$2" &` + if kind == "setsid" { + launch = `setsid "$1" "$2" &` + } + if kind == "double-fork" { + launch = `(setsid "$1" "$2" &) ;` + } + if leaderExit { + launch += " exit 0" + } else { + launch += " while :; do sleep 1; done" + } + c := h.connect() + op := h.start(c, cgroupSpec("sh", "-c", launch, "leader", writer, dir)) + awaitFile(t, filepath.Join(dir, "ticks"), "a") + raw, err := os.ReadFile(filepath.Join(dir, "pid")) + if err != nil { + t.Fatal(err) + } + pid, err := strconv.Atoi(strings.TrimSpace(string(raw))) + if err != nil { + t.Fatal(err) + } + fd, err := unix.PidfdOpen(pid, 0) + if err != nil { + t.Fatal(err) + } + defer unix.Close(fd) + if leaderExit { + events(t, op, sp.EventExited) + } else { + events(t, op, sp.EventStarted) + } + if st, err := op.Inspect(t.Context()); err != nil || st.Scope != sp.ScopeStateActive { + t.Fatalf("background scope was not retained: %+v %v", st, err) + } + record := cgroupRecord(t, h, op) + group := record.cgroup.path + // Losing and replacing the observer does not cancel background work. + c.Close() + c = h.connect() + op, _, err = c.Attach(t.Context(), h.svc.instance, op.Ref().OperationID, 0) + if err != nil { + t.Fatal(err) + } + if err := op.Cancel(t.Context(), 0); err != nil { + t.Fatal(err) + } + events(t, op, sp.EventScopeClosed) + if _, err := os.Stat(group); !errors.Is(err, fs.ErrNotExist) { + t.Fatalf("closed group remains: %v", err) + } + if st, err := readStat(pid); err == nil && st.live() { + t.Fatalf("detached writer remains live: %+v", st) + } + if err := os.WriteFile(filepath.Join(dir, "gate"), []byte("go"), 0600); err != nil { + t.Fatal(err) + } + awaitFile(t, filepath.Join(dir, "ticks"), "a") + }) + } + } +} + +func TestCgroupBackgroundCompletesAfterLeader(t *testing.T) { + h := cgroupHarness(t) + dir := t.TempDir() + op := h.start(h.connect(), cgroupSpec("sh", "-c", `(while [ ! -e "$0/gate" ]; do sleep .01; done; printf completed > "$0/result") & exit 0`, dir)) + events(t, op, sp.EventExited) + wantCode(t, op.Release(t.Context()), sp.CodeBusy) + if st, err := op.Inspect(t.Context()); err != nil || st.Scope != sp.ScopeStateActive { + t.Fatalf("scope %+v %v", st, err) + } + if err := os.WriteFile(filepath.Join(dir, "gate"), []byte("go"), 0600); err != nil { + t.Fatal(err) + } + events(t, op, sp.EventScopeClosed) + awaitFile(t, filepath.Join(dir, "result"), "completed") +} + +func TestCgroupFailedLaunchLeavesNoScope(t *testing.T) { + h := cgroupHarness(t) + op := h.start(h.connect(), cgroupSpec("no-such-command")) + ev := next(t, op) + if f, ok := ev.(sp.StartFailedEvent); !ok || f.Failure.Code != sp.CodeNotFound { + t.Fatalf("first event: %+v", ev) + } + record := cgroupRecord(t, h, op) + if _, err := os.Stat(record.cgroup.path); !errors.Is(err, fs.ErrNotExist) { + t.Fatalf("failed launch group remains: %v", err) + } + if err := op.Release(t.Context()); err != nil { + t.Fatal(err) + } +} + +func TestCgroupCancelWhileStarting(t *testing.T) { + h := cgroupHarness(t) + spec := cgroupSpec("sh", "-c", "sleep 30 & wait") + key := opKey{h.att, sandboxwire.NewID()} + op := newOperation(h.svc, key, spec.Digest()) + h.svc.mu.Lock() + h.svc.ops[key] = op + h.svc.active.Add(1) + h.svc.mu.Unlock() + if err := op.cancel(0); err != nil { + t.Fatal(err) + } + op.launch(sp.StartRequest{OperationRef: sp.OperationRef{ServerInstanceID: h.svc.instance, OperationID: key.operation}, Spec: spec}) + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + op.awaitScope(ctx) + if ctx.Err() != nil || op.inspect().Scope != sp.ScopeStateClosed { + t.Fatalf("pending cancellation did not settle: %+v %v", op.inspect(), ctx.Err()) + } + if _, err := os.Stat(op.cgroup.path); !errors.Is(err, fs.ErrNotExist) { + t.Fatalf("cancelled group remains: %v", err) + } +} + +func TestCgroupCancelDoesNotReachOtherOperation(t *testing.T) { + h := cgroupHarness(t) + c := h.connect() + first := h.start(c, cgroupSpec("sleep", "30")) + other := h.start(c, cgroupSpec("sleep", "30")) + events(t, first, sp.EventStarted) + events(t, other, sp.EventStarted) + if err := first.Cancel(t.Context(), 0); err != nil { + t.Fatal(err) + } + events(t, first, sp.EventScopeClosed) + if st, err := other.Inspect(t.Context()); err != nil || st.State != sp.StateRunning || st.Scope != sp.ScopeStateActive { + t.Fatalf("other operation changed: %+v %v", st, err) + } + if err := other.Cancel(t.Context(), 0); err != nil { + t.Fatal(err) + } + events(t, other, sp.EventScopeClosed) +} + +func TestCgroupPlacementAndGracefulSignal(t *testing.T) { + h := cgroupHarness(t) + op := h.start(h.connect(), cgroupSpec("sh", "-c", `trap 'printf stopped; exit 0' TERM; cat /proc/self/cgroup; printf ready; while :; do sleep 1; done`)) + var before []sp.Event + for !strings.Contains(output(before, sp.StreamStdout), "ready") { + before = append(before, next(t, op)) + } + group := filepath.Base(cgroupRecord(t, h, op).cgroup.path) + if !strings.Contains(output(before, sp.StreamStdout), "/"+group+"\n") { + t.Fatalf("target started outside its cgroup: %q", output(before, sp.StreamStdout)) + } + if err := op.Signal(t.Context(), 15, sp.TargetScope); err != nil { + t.Fatal(err) + } + evs := events(t, op, sp.EventScopeClosed) + exit, _ := find[sp.ExitedEvent](t, evs) + if exit.Status.Kind != sp.ExitCode || exit.Status.Code != 0 { + t.Fatalf("TERM handler did not finish: %+v", exit.Status) + } + if !strings.Contains(output(evs, sp.StreamStdout), "stopped") { + t.Fatalf("TERM handler output missing: %q", output(evs, sp.StreamStdout)) + } +} + +func TestCgroupObservationFailureDoesNotSettle(t *testing.T) { + h := cgroupHarness(t) + dir := t.TempDir() + op := h.start(h.connect(), cgroupSpec("sh", "-c", `(while [ ! -e "$0/gate" ]; do sleep .01; done) & exit 0`, dir)) + group := cgroupRecord(t, h, op).cgroup.path + eventsFile := filepath.Join(group, "cgroup.events") + if err := os.Chmod(eventsFile, 0); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(eventsFile, 0444) }) + evs := events(t, op, sp.EventObservationLost) + lost, _ := find[sp.ObservationLostEvent](t, evs) + if lost.Observation != sp.ObservationScope { + t.Fatalf("wrong lost observation: %+v", lost) + } + if st, err := op.Inspect(t.Context()); err != nil || st.Scope != sp.ScopeStateUnknown { + t.Fatalf("scope %+v %v", st, err) + } + wantCode(t, op.Release(t.Context()), sp.CodeBusy) + if err := os.Chmod(eventsFile, 0444); err != nil { + t.Fatal(err) + } + if err := op.Cancel(t.Context(), 0); err != nil { + t.Fatal(err) + } + events(t, op, sp.EventScopeClosed) +} + +func TestCgroupOutputDrainsAfterLeader(t *testing.T) { + h := cgroupHarness(t) + dir := t.TempDir() + op := h.start(h.connect(), cgroupSpec("sh", "-c", `printf leader; (while [ ! -e "$0/gate" ]; do sleep .01; done; printf child) & exit 0`, dir)) + evs := events(t, op, sp.EventExited) + if output(evs, sp.StreamStdout) != "leader" { + t.Fatalf("leader output: %q", output(evs, sp.StreamStdout)) + } + if err := os.WriteFile(filepath.Join(dir, "gate"), []byte("go"), 0600); err != nil { + t.Fatal(err) + } + scope, drained := false, false + for !scope || !drained { + ev := next(t, op) + evs = append(evs, ev) + switch ev.(type) { + case sp.ScopeClosedEvent: + scope = true + case sp.OutputClosedEvent: + drained = true + } + } + if output(evs, sp.StreamStdout) != "leaderchild" { + t.Fatalf("lost descendant output: %q", output(evs, sp.StreamStdout)) + } + if err := op.Release(t.Context()); err != nil { + t.Fatal(err) + } +} + +// A launch that has already failed still owns its scope until observation +// confirms it empty. A filesystem error must not expose an invalid status, +// claim settlement, or prevent other requests and bounded Shutdown. +func TestCgroupFailedLaunchObservationAndShutdown(t *testing.T) { + h := cgroupHarness(t) + spec := cgroupSpec("no-such-command") + key := opKey{h.att, sandboxwire.NewID()} + record := newOperation(h.svc, key, spec.Digest()) + h.svc.mu.Lock() + h.svc.ops[key] = record + h.svc.active.Add(1) + h.svc.mu.Unlock() + _, failure := record.spawn(sp.StartRequest{OperationRef: sp.OperationRef{ServerInstanceID: h.svc.instance, OperationID: key.operation}, Spec: spec}) + if failure == nil { + t.Fatal("launch unexpectedly succeeded") + } + eventsFile := filepath.Join(record.cgroup.path, "cgroup.events") + if err := os.Chmod(eventsFile, 0); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(eventsFile, 0444) }) + record.mu.Lock() + record.startFailure = failure + record.killing = true + record.mu.Unlock() + go record.watchScope() + c := h.connect() + op, _, err := c.Attach(t.Context(), h.svc.instance, key.operation, 0) + if err != nil { + t.Fatal(err) + } + for deadline := time.Now().Add(5 * time.Second); ; time.Sleep(time.Millisecond) { + st, err := op.Inspect(t.Context()) + if err != nil { + t.Fatal(err) + } + if err := st.Validate(); err != nil { + t.Fatal(err) + } + if st.Scope == sp.ScopeStateUnknown { + if st.State != sp.StateStarting || st.StartFailure != nil { + t.Fatalf("invalid pending failure: %+v", st) + } + break + } + if time.Now().After(deadline) { + t.Fatal("observation failure was not visible") + } + } + select { + case ev := <-op.Events(): + t.Fatalf("event before confirmed failure: %+v", ev) + default: + } + other := h.start(c, cgroupSpec("true")) + events(t, other, sp.EventScopeClosed) + ctx, cancel := context.WithTimeout(t.Context(), 20*time.Millisecond) + done := make(chan struct{}) + go func() { h.svc.Shutdown(ctx); close(done) }() + select { + case <-done: + case <-time.After(2 * time.Second): + _ = os.Chmod(eventsFile, 0444) + cancel() + t.Fatal("shutdown ignored its deadline") + } + err = ctx.Err() + cancel() + if err != context.DeadlineExceeded { + t.Fatalf("shutdown returned without its unresolved scope: %v", err) + } + if err := os.Chmod(eventsFile, 0444); err != nil { + t.Fatal(err) + } + ev := next(t, op) + if _, ok := ev.(sp.StartFailedEvent); !ok { + t.Fatalf("wrong first event: %+v", ev) + } + if st, err := op.Inspect(t.Context()); err != nil || st.State != sp.StateStartFailed || st.Scope != sp.ScopeStateClosed { + t.Fatalf("failure not settled: %+v %v", st, err) + } +} + +func TestCgroupEmptyDirectoryFailureDoesNotKeepScopeActive(t *testing.T) { + h := cgroupHarness(t) + g, err := newProcessCgroup(h.svc.cgroupParent) + if err != nil { + t.Fatal(err) + } + // Inject a directory removal failure after the kernel confirms emptiness. + // Cleanup failure must not turn an empty process scope into Unknown. + original := g.path + g.path = filepath.Join(original, "missing-parent", "group") + if live, err := g.observe(); live != 0 || err != nil { + t.Fatalf("empty group: %d %v", live, err) + } + if err := g.cleanup(); err == nil { + t.Fatal("injected directory removal failure was hidden") + } + if live, err := g.observe(); live != 0 || err != nil { + t.Fatalf("cleanup error changed empty scope: %d %v", live, err) + } + g.path = original + if err := g.cleanup(); err != nil { + t.Fatal(err) + } +} + +func TestCgroupKillFailureKeepsNestedScopeActive(t *testing.T) { + h := cgroupHarness(t) + dir := t.TempDir() + script := `group=/sys/fs/cgroup$(sed -n 's/^0:://p' /proc/self/cgroup) +mkdir "$group/nested" +echo $$ > "$group/nested/cgroup.procs" +(trap '' TERM HUP; printf a > "$0/ticks"; while [ ! -e "$0/gate" ]; do sleep .01; done; printf b >> "$0/ticks"; while :; do sleep 1; done) & +exit 0` + op := h.start(h.connect(), cgroupSpec("sh", "-c", script, dir)) + awaitFile(t, filepath.Join(dir, "ticks"), "a") + events(t, op, sp.EventExited) + group := cgroupRecord(t, h, op).cgroup.path + killFile := filepath.Join(group, "cgroup.kill") + if err := os.Chmod(killFile, 0); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(killFile, 0200) }) + wantCode(t, op.Signal(t.Context(), 9, sp.TargetScope), sp.CodeIO) + if err := op.Cancel(t.Context(), 0); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, "gate"), nil, 0600); err != nil { + t.Fatal(err) + } + awaitFile(t, filepath.Join(dir, "ticks"), "ab") + if st, err := op.Inspect(t.Context()); err != nil || st.Scope != sp.ScopeStateActive { + t.Fatalf("failed kill falsely closed scope: %+v %v", st, err) + } + wantCode(t, op.Release(t.Context()), sp.CodeBusy) + if err := os.Chmod(killFile, 0200); err != nil { + t.Fatal(err) + } + // The existing cancellation poll retries the failed KILL without a new Cancel. + events(t, op, sp.EventScopeClosed) + if _, err := os.Stat(group); !errors.Is(err, fs.ErrNotExist) { + t.Fatalf("nested scope was not removed: %v", err) + } +} diff --git a/apps/sandboxio/internal/processservice/launch.go b/apps/sandboxio/internal/processservice/launch.go index 4e9af85d7..54ce14f76 100644 --- a/apps/sandboxio/internal/processservice/launch.go +++ b/apps/sandboxio/internal/processservice/launch.go @@ -25,11 +25,19 @@ func (op *operation) launch(req sp.StartRequest) { l, f := op.spawn(req) op.mu.Lock() if f != nil { - op.state, op.startFailure = sp.StateStartFailed, f - op.scope, op.stdinClosed = sp.ScopeStateClosed, true - op.push(sp.StartFailedEvent{EventHeader: op.header(), Failure: *f}) - op.settleLocked() + op.startFailure = f + op.cond.Broadcast() + if op.cgroup != nil { + // A possibly-executed target must be gone before StartFailed + // establishes settlement. Observe asynchronously, so a damaged + // scope cannot block the stream or the service's bounded shutdown. + op.killing = true + op.mu.Unlock() + go op.watchScope() + return + } op.mu.Unlock() + op.scopeClosed() if op.cu != nil { op.cu.close() } @@ -65,6 +73,13 @@ func (op *operation) spawn(req sp.StartRequest) (l launched, f *sp.Failure) { ioFail := func(what string, err error) *sp.Failure { return sp.Fail(sp.CodeIO, sandboxwire.EffectNone, "%s: %v", what, err) } + if req.Spec.Scope == sp.ScopeCgroupV2 { + var err error + op.cgroup, err = newProcessCgroup(op.s.cgroupParent) + if err != nil { + return l, ioFail("create process cgroup", err) + } + } // Descriptors the child inherits close after the start; the parent's // close only when the launch fails. Each is listed as soon as it exists. var child, parent []*os.File @@ -157,6 +172,14 @@ func (op *operation) spawn(req sp.StartRequest) (l launched, f *sp.Failure) { closeAll(child) child = nil + if op.cgroup != nil { + // The trusted trampoline is blocked reading launchFD. Place it before + // releasing the request, so the target never runs outside its scope. + if err := op.cgroup.place(pid); err != nil { + op.killPinned(true, unix.SIGKILL) + return l, ioFail("place process in cgroup", err) + } + } _, werr := launchW.Write(sp.Encode(req)) launchW.Close() status, rerr := io.ReadAll(statusR) diff --git a/apps/sandboxio/internal/processservice/operation.go b/apps/sandboxio/internal/processservice/operation.go index e6894af41..e87fb96a7 100644 --- a/apps/sandboxio/internal/processservice/operation.go +++ b/apps/sandboxio/internal/processservice/operation.go @@ -69,7 +69,8 @@ type operation struct { leaderGone bool status unix.WaitStatus // cu holds the processes proven to be in the session; the spawn sets it. - cu *custody + cu *custody + cgroup *processCgroup // A Cancel that arrives while the operation is starting waits for the // launch. @@ -135,10 +136,14 @@ func (op *operation) settleLocked() { } func (op *operation) statusLocked() sp.OperationStatus { + var failure *sp.Failure + if op.state == sp.StateStartFailed { + failure = op.startFailure + } return sp.OperationStatus{ State: op.state, Exit: op.exit, - StartFailure: op.startFailure, + StartFailure: failure, StdinOffset: op.stdinOffset, StdinClosed: op.stdinClosed, Output: op.output, @@ -413,7 +418,7 @@ func (op *operation) closeOutput(name sp.Stream) error { // has published them. func (op *operation) abandonOutput() { op.mu.Lock() - for op.state == sp.StateStarting { + for op.state == sp.StateStarting && op.startFailure == nil { op.cond.Wait() } streams := op.streams diff --git a/apps/sandboxio/internal/processservice/scope.go b/apps/sandboxio/internal/processservice/scope.go index 0972702f2..91e3d92f5 100644 --- a/apps/sandboxio/internal/processservice/scope.go +++ b/apps/sandboxio/internal/processservice/scope.go @@ -13,6 +13,7 @@ import ( "sync" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" "golang.org/x/sys/unix" sp "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxprocess" @@ -374,9 +375,12 @@ func (op *operation) signalGroup(pgrp int, sig unix.Signal) error { return signalResult(sent, lost, failed, "the process group is empty") } -// signalScope signals the initial process group by ID while its leader pins -// it, then every other member in custody. +// signalScope uses the selected scope. A POSIX scope signals the initial +// process group while its leader pins it, then every other member in custody. func (op *operation) signalScope(sig unix.Signal) error { + if op.cgroup != nil { + return op.cgroup.signal(sig) + } sid := op.sid() grouped, gerr := op.killPinned(true, sig) sent, _, lost, failed := op.cu.sweep(func(st procStat) bool { return !grouped || st.pgrp != sid }, sig) @@ -386,11 +390,11 @@ func (op *operation) signalScope(sig unix.Signal) error { return signalResult(sent, lost, errors.Join(gerr, failed), "the scope is empty") } -// watchScope polls until the session is confirmed empty, repeating KILL once +// watchScope polls until the scope is confirmed empty, repeating KILL once // Cancel's grace has passed so members forked meanwhile die too. A failed // poll makes the scope Unknown and reports ObservationLost; polling and the // KILL escalation go on, and the operation settles only once a poll confirms -// the session empty. +// the scope empty. func (op *operation) watchScope() { delay := scopePollFirst for { @@ -400,16 +404,36 @@ func (op *operation) watchScope() { sig = unix.SIGKILL } op.mu.Unlock() - _, live, lost, _ := op.cu.sweep(func(procStat) bool { return true }, sig) + var live int + var lost error + if op.cgroup != nil { + if sig != 0 { + op.cgroup.signal(sig) + } + live, lost = op.cgroup.observe() + } else { + _, live, lost, _ = op.cu.sweep(func(procStat) bool { return true }, sig) + } if lost == nil && live == 0 { - op.cu.close() + if op.cgroup != nil { + if err := op.cgroup.cleanup(); err != nil { + log.Bg().Warn("empty process cgroup cleanup failed", "operation", op.key.operation, "error", err) + } + } + if op.cu != nil { + op.cu.close() + } op.scopeClosed() return } op.mu.Lock() if lost != nil && op.scope == sp.ScopeStateActive { op.scope = sp.ScopeStateUnknown - op.push(sp.ObservationLostEvent{EventHeader: op.header(), Observation: sp.ObservationScope, Failure: *sp.Fail(sp.CodeIO, sandboxwire.EffectPossible, "observe session: %v", lost)}) + // Started or StartFailed must be the first event. A failed launch + // being drained is still inspectable as Starting with Unknown scope. + if op.state != sp.StateStarting { + op.push(sp.ObservationLostEvent{EventHeader: op.header(), Observation: sp.ObservationScope, Failure: *sp.Fail(sp.CodeIO, sandboxwire.EffectPossible, "observe process scope: %v", lost)}) + } } op.mu.Unlock() time.Sleep(delay) @@ -440,6 +464,11 @@ func (op *operation) scopeClosed() { op.killTimer.Stop() } op.scope = sp.ScopeStateClosed - op.push(sp.ScopeClosedEvent{EventHeader: op.header()}) + if op.state == sp.StateStarting && op.startFailure != nil { + op.state, op.stdinClosed = sp.StateStartFailed, true + op.push(sp.StartFailedEvent{EventHeader: op.header(), Failure: *op.startFailure}) + } else { + op.push(sp.ScopeClosedEvent{EventHeader: op.header()}) + } op.settleLocked() } diff --git a/apps/sandboxio/internal/processservice/service.go b/apps/sandboxio/internal/processservice/service.go index 5ec94e6ca..dbd02baf9 100644 --- a/apps/sandboxio/internal/processservice/service.go +++ b/apps/sandboxio/internal/processservice/service.go @@ -17,6 +17,7 @@ import ( "sync/atomic" "time" + "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log" "golang.org/x/sys/unix" sp "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxprocess" @@ -49,9 +50,10 @@ var signals = []sp.Signal{1, 2, 3, 9, 10, 12, 14, 15, 18, 19, 20, 21, 22, 28} // // Service is one incarnation of the process service. Its operation records // live in memory; a new Service has a new ServerInstanceID. type Service struct { - instance sandboxwire.ID - caps sp.Capabilities - cfg Config + instance sandboxwire.ID + caps sp.Capabilities + cfg Config + cgroupParent string // the service's verified writable cgroup, or empty // stat reads a process's /proc stat; tests replace it. stat func(pid int) (procStat, error) // onExpire runs after an expired owner-loss grace has decided the @@ -104,8 +106,13 @@ func New(cfg Config) (*Service, error) { if err := probePidfd(); err != nil { return nil, fmt.Errorf("%w: %v", ErrPidfdUnsupported, err) } + parent := delegatedCgroup() + if parent != "" { + caps.Scopes = append(caps.Scopes, sp.ScopeCgroupV2) + } return &Service{ - instance: sandboxwire.NewID(), caps: caps, cfg: cfg, stat: readStat, + cgroupParent: parent, + instance: sandboxwire.NewID(), caps: caps, cfg: cfg, stat: readStat, ops: map[opKey]*operation{}, owners: map[sandboxwire.ID]*time.Timer{}, stale: map[sandboxwire.ID]bool{}, }, nil } @@ -361,6 +368,11 @@ func (s *Service) Shutdown(ctx context.Context) { s.endAll(ops) for _, op := range ops { op.awaitScope(ctx) + if ctx.Err() == nil && op.cgroup != nil { + if err := op.cgroup.cleanup(); err != nil { + log.Warn(ctx, "empty process cgroup cleanup failed at shutdown", "operation", op.key.operation, "error", err) + } + } } } diff --git a/docs/process-protocol.md b/docs/process-protocol.md index 33a589cac..d37f051c7 100644 --- a/docs/process-protocol.md +++ b/docs/process-protocol.md @@ -36,7 +36,9 @@ A service must: - never drop an event it has not been told was delivered, as described in [Output, replay and flow control](#output-replay-and-flow-control); - treat the loss of a stream as nothing more than the loss of an observer, as described in [Ownership](#ownership). -The Linux service calls `processservice.Init()` first in the binary's `main`. Go cannot set a child's umask, so each launch re-executes the service binary as a trampoline that reads the launch from an inherited descriptor, marks every inherited descriptor above 2 close-on-exec, applies the umask and working directory, and execs the target. `Init` runs that trampoline and returns at once in a normal start. The Linux service launches every operation in a new session with `setsid`, observes the session through `/proc`, and advertises `ScopePOSIXSession` only. It requires `pidfd_open` and `pidfd_send_signal` (Linux 5.3 or later): without them `processservice.New` fails with `ErrPidfdUnsupported`. +The Linux service calls `processservice.Init()` first in the binary's `main`. Go cannot set a child's umask, so each launch re-executes the service binary as a trampoline that reads the launch from an inherited descriptor, marks every inherited descriptor above 2 close-on-exec, applies the umask and working directory, and execs the target. `Init` runs that trampoline and returns at once in a normal start. The Linux service launches every operation in a new session with `setsid` and advertises `ScopePOSIXSession`. It requires `pidfd_open` and `pidfd_send_signal` (Linux 5.3 or later): without them `processservice.New` fails with `ErrPidfdUnsupported`. + +The Linux implementation also advertises `ScopeCgroupV2` when its current cgroup, read from `/proc/self/cgroup` under `/sys/fs/cgroup`, permits creating and removing an operation subgroup, writing `cgroup.kill`, and observing `cgroup.events`. The Provider must delegate that current cgroup to the service account before starting SandboxIO; this requires a writable cgroup v2 mount and `cgroup.kill` (Linux 5.14 or later). No resource controller needs to be enabled. Each selected cgroup scope owns a new subgroup. The service places the blocked trampoline there before sending its launch spec, so the target and its descendants start inside the scope. A failed placement never executes the target. Delegation supplies process lifetime control, not isolation from other processes with the same account; the Provider owns that isolation. Before serving, `main` makes the process a child subreaper (`prctl(PR_SET_CHILD_SUBREAPER)`) and runs `processservice.Reap(ctx)` for the life of the process. `Reap` is the process's only `wait`: it reaps every child, delivers each leader's exit to its operation, and reaps the orphaned descendants the subreaper inherits. Nothing else in the binary may wait for children, and no operation observes an exit while `Reap` is not running. When the binary stops, it calls `Shutdown(ctx)` after its streams have ended: `Shutdown` cancels every live operation and abandons its output, as revoking [ownership](#ownership) does, and returns once each scope has closed or `ctx` ends. @@ -139,6 +141,8 @@ With a PTY, output arrives as the `Terminal` stream, stdin writes are terminal i Both scopes start the process in a new POSIX session. `ScopeCgroupV2` also places it in a new cgroup and is advertised only when the service enforces that. A descendant can leave a `ScopePOSIXSession` scope by calling `setsid`; that is the scope's limit. The Linux service observes a session scope by polling `/proc` for live members, so `ScopeClosed` is best-effort within that limit. +For the Linux implementation's `ScopeCgroupV2`, `TargetScope` signals use the operation's cgroup subtree: KILL writes `cgroup.kill`; other signals use pidfds after rechecking membership in that subgroup. Forking, double-forking and `setsid` do not leave this scope. `ScopeClosed` requires `cgroup.events` to report the entire subtree empty, independently of the leader's exit and output draining. A failed observation keeps the scope `Unknown` and is retried. Background processes may continue after a normal leader exit until they finish or the scope is cancelled. + A request never names a process ID. `Signal` takes a signal number from `Capabilities.Signals` and one of these targets: | Target | Receives the signal | @@ -150,7 +154,7 @@ A request never names a process ID. `Signal` takes a signal number from `Capabil A target with no process returns `NotRunning`, and so does every target after `ScopeClosed` and a `TargetPTYForegroundGroup` whose terminal has no foreground group. -A signal reaches only a process the service can prove is in the operation's session, never one that reused a PID or a session ID. While the leader is unreaped, its PID pins the session and initial process group IDs, and the Linux service signals the leader or that group by ID. Otherwise it signals through the pidfds it holds for the session's processes, and only those a refresh in the same request proved still in the session; when that refresh fails, it signals nothing more and the request fails with `IO`. It takes a pidfd for a process showing the session ID only when a process it already holds stayed in the session across that read, and it updates the held set before reaping any held process. When it holds no process but a live process still shows the session ID, it cannot tell the session from a new one with the same ID: it signals nothing, the target returns `NotRunning`, and the scope becomes `Unknown` with `ObservationLost`. +For `ScopePOSIXSession` and the session-based signal targets, a signal reaches only a process the service can prove is in the operation's session, never one that reused a PID or a session ID. While the leader is unreaped, its PID pins the session and initial process group IDs, and the Linux service signals the leader or that group by ID. Otherwise it signals through the pidfds it holds for the session's processes, and only those a refresh in the same request proved still in the session; when that refresh fails, it signals nothing more and the request fails with `IO`. It takes a pidfd for a process showing the session ID only when a process it already holds stayed in the session across that read, and it updates the held set before reaping any held process. When it holds no process but a live process still shows the session ID, it cannot tell the session from a new one with the same ID: it signals nothing, the target returns `NotRunning`, and the scope becomes `Unknown` with `ObservationLost`. ### Ownership @@ -200,3 +204,5 @@ A failure carries a `Code`, an `Effect` and a message. `EffectNone` means the re ## Verification `go test ./internal/sandboxprocess` checks the golden frames in `internal/sandboxprocess/testdata`, and `go test -fuzz FuzzDecode ./internal/sandboxprocess` fuzzes the decoder. `go test ./apps/sandboxio/internal/processservice` runs the Linux service over in-memory streams with real processes. + +The backend runtime CI job runs `TestCgroup` with `OAC_TEST_PROCESS_CGROUP=1` in a systemd delegated unit. That gate fails when delegation is missing. Locally, compile the test binary with `go test -race -c -o "$HOME/.oac/processservice.test" ./apps/sandboxio/internal/processservice` and run it with `systemd-run --user --collect --wait --pipe --property=Delegate=yes --setenv=OAC_TEST_PROCESS_CGROUP=1 "$HOME/.oac/processservice.test" -test.run TestCgroup`. These tests exercise detached descendants, cancellation, reconnect, output drain, scope observation failures and bounded shutdown with real kernel cgroups. diff --git a/docs/zh/process-protocol.md b/docs/zh/process-protocol.md index 5f9248bf7..f5f94ade4 100644 --- a/docs/zh/process-protocol.md +++ b/docs/zh/process-protocol.md @@ -1,7 +1,7 @@ --- title: "进程协议" source: docs/process-protocol.md -source_hash: c2658d8462e8699b9ab5c5122c2f36ee19e30b5890bed69d54c545d60d596ff9 +source_hash: 5111fa281a69270ee6301648134a3266003c055256d164574b9b9a5941a73564 --- 进程协议定义 agent host 如何在沙箱中启动和控制进程。沙箱内的 Sandbox I/O 服务提供该协议,agent host 的 broker 是其客户端。协议依据明确的 spec 启动进程,以有序事件流式传输其输出,在精确 offset 处接受 stdin,并将 leader 退出、输出结束和进程 scope 结束作为独立事实报告。 @@ -38,7 +38,9 @@ Go 客户端为 `sandboxprocess.NewClient(stream)`。`Start` 和 `Attach` 返回 - 绝不丢弃未被告知已交付的事件,见[输出、重放与流量控制](#output-replay-and-flow-control); - 将 stream 丢失仅视为 observer 丢失,见[所有权](#ownership)。 -Linux 服务在二进制的 `main` 中首先调用 `processservice.Init()`。Go 无法设置子进程的 umask,因此每次启动都会将服务二进制作为 trampoline 重新执行:它从继承的描述符读取启动信息,将所有大于 2 的继承描述符标记为 close-on-exec,应用 umask 和工作目录,然后 exec 目标。`Init` 负责运行该 trampoline,正常启动时立即返回。Linux 服务用 `setsid` 在新会话中启动每个 operation,通过 `/proc` 观测会话,并且只声明 `ScopePOSIXSession`。它要求 `pidfd_open` 和 `pidfd_send_signal`(Linux 5.3 或更高版本):缺少它们时 `processservice.New` 以 `ErrPidfdUnsupported` 失败。 +Linux 服务在二进制的 `main` 中首先调用 `processservice.Init()`。Go 无法设置子进程的 umask,因此每次启动都会将服务二进制作为 trampoline 重新执行:它从继承的描述符读取启动信息,将所有大于 2 的继承描述符标记为 close-on-exec,应用 umask 和工作目录,然后 exec 目标。`Init` 负责运行该 trampoline,正常启动时立即返回。Linux 服务用 `setsid` 在新会话中启动每个 operation,并声明 `ScopePOSIXSession`。它要求 `pidfd_open` 和 `pidfd_send_signal`(Linux 5.3 或更高版本):缺少它们时 `processservice.New` 以 `ErrPidfdUnsupported` 失败。 + +Linux 实现从 `/proc/self/cgroup` 读取当前 cgroup,并映射到 `/sys/fs/cgroup`;当该组允许创建和删除 operation 子组、写入 `cgroup.kill` 及观测 `cgroup.events` 时,还会声明 `ScopeCgroupV2`。Provider 必须在启动 SandboxIO 前将该当前 cgroup 委派给服务账号;这需要可写的 cgroup v2 挂载和 `cgroup.kill`(Linux 5.14 或更高版本)。无需启用资源 controller。每个选择 cgroup scope 的 operation 独占一个新子组。服务先将阻塞中的 trampoline 放入该组,再发送 launch spec,保证目标及其后代从启动起就在 scope 内。放置失败时绝不执行目标。委派提供进程生命周期控制,不提供针对同账号其他进程的隔离;隔离由 Provider 负责。 开始服务前,`main` 将进程设为 child subreaper(`prctl(PR_SET_CHILD_SUBREAPER)`),并在进程整个生命周期内运行 `processservice.Reap(ctx)`。`Reap` 是进程中唯一的 `wait`:它回收每个子进程,将每个 leader 的退出交付给对应 operation,并回收 subreaper 继承的孤儿后代进程。二进制中其他任何代码都不得等待子进程;`Reap` 未运行时,任何 operation 都观测不到退出。二进制停止时,在其 stream 结束后调用 `Shutdown(ctx)`:`Shutdown` 像撤销[所有权](#ownership)时一样,取消每个存活的 operation 并放弃其输出,在每个 scope 关闭或 `ctx` 结束后返回。 @@ -141,6 +143,8 @@ stdin offset 从 0 开始,计算服务已接受的字节数。`WriteStdin` 和 两种 scope 都在新的 POSIX 会话中启动进程。`ScopeCgroupV2` 还会将进程放入新的 cgroup,且仅在服务强制执行这一点时才声明。后代进程可以调用 `setsid` 离开 `ScopePOSIXSession` scope;这是该 scope 的局限。Linux 服务通过轮询 `/proc` 中的存活成员来观测会话 scope,因此 `ScopeClosed` 在此局限内尽力而为。 +对于 Linux 实现的 `ScopeCgroupV2`,`TargetScope` 信号作用于 operation 的 cgroup 子树:KILL 写入 `cgroup.kill`;其他信号在再次核对子组成员后通过 pidfd 发送。fork、double-fork 和 `setsid` 都不会离开该 scope。`ScopeClosed` 要求 `cgroup.events` 确认整棵子树为空,与 leader 退出和输出排空相互独立。观测失败时 scope 保持 `Unknown`,服务继续重试。后台进程可在 leader 正常退出后继续运行,直到它们完成或 scope 被取消。 + 请求从不指定进程 ID。`Signal` 接受 `Capabilities.Signals` 中的信号编号和以下目标之一: | 目标 | 接收信号的进程 | @@ -152,7 +156,7 @@ stdin offset 从 0 开始,计算服务已接受的字节数。`WriteStdin` 和 没有进程的目标返回 `NotRunning`;`ScopeClosed` 之后的任何目标,以及终端没有前台进程组的 `TargetPTYForegroundGroup`,同样返回 `NotRunning`。 -信号只送达服务能证明属于该 operation 会话的进程,绝不送达复用了 PID 或会话 ID 的进程。leader 尚未被回收时,其 PID 固定住会话 ID 和初始进程组 ID,Linux 服务按 ID 向 leader 或该进程组发送信号。否则,它通过为会话进程持有的 pidfd 发送信号,并且只发给同一请求内刷新证明仍在会话中的进程;刷新失败时,不再发送任何信号,请求以 `IO` 失败。只有当它已持有的某个进程在该次读取期间始终留在会话中时,它才为显示该会话 ID 的进程获取 pidfd;回收任何已持有进程之前,它会先更新持有集合。当它未持有任何进程,但仍有存活进程显示该会话 ID 时,它无法区分该会话与具有相同 ID 的新会话:它不发送任何信号,目标返回 `NotRunning`,scope 随 `ObservationLost` 变为 `Unknown`。 +对于 `ScopePOSIXSession` 及基于会话的信号目标,信号只送达服务能证明属于该 operation 会话的进程,绝不送达复用了 PID 或会话 ID 的进程。leader 尚未被回收时,其 PID 固定住会话 ID 和初始进程组 ID,Linux 服务按 ID 向 leader 或该进程组发送信号。否则,它通过为会话进程持有的 pidfd 发送信号,并且只发给同一请求内刷新证明仍在会话中的进程;刷新失败时,不再发送任何信号,请求以 `IO` 失败。只有当它已持有的某个进程在该次读取期间始终留在会话中时,它才为显示该会话 ID 的进程获取 pidfd;回收任何已持有进程之前,它会先更新持有集合。当它未持有任何进程,但仍有存活进程显示该会话 ID 时,它无法区分该会话与具有相同 ID 的新会话:它不发送任何信号,目标返回 `NotRunning`,scope 随 `ObservationLost` 变为 `Unknown`。 ### 所有权 {#ownership} @@ -202,3 +206,5 @@ stdin offset 从 0 开始,计算服务已接受的字节数。`WriteStdin` 和 ## 验证 {#verification} `go test ./internal/sandboxprocess` 检查 `internal/sandboxprocess/testdata` 中的 golden frame,`go test -fuzz FuzzDecode ./internal/sandboxprocess` 对解码器进行 fuzz 测试。`go test ./apps/sandboxio/internal/processservice` 使用真实进程,在内存 stream 上运行 Linux 服务。 + +backend runtime CI job 在 systemd 委派单元内以 `OAC_TEST_PROCESS_CGROUP=1` 运行 `TestCgroup`;缺少委派时该门禁失败。本地可用 `go test -race -c -o "$HOME/.oac/processservice.test" ./apps/sandboxio/internal/processservice` 编译测试二进制,再用 `systemd-run --user --collect --wait --pipe --property=Delegate=yes --setenv=OAC_TEST_PROCESS_CGROUP=1 "$HOME/.oac/processservice.test" -test.run TestCgroup` 运行。这些测试使用真实内核 cgroup 验证脱离会话的后代、取消、重连、输出排空、scope 观测失败及有界关闭。 From c2d4a2b635625a15e8097bf037d6a7cd8aa0104e Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Fri, 9 Oct 2026 08:38:34 +0000 Subject: [PATCH 2/2] Fix cgroup delegation and cancellation races --- .../internal/processservice/cgroup.go | 33 ++++++++- .../internal/processservice/cgroup_test.go | 68 +++++++++++++++++-- .../internal/processservice/launch.go | 21 +++--- .../internal/processservice/operation.go | 2 + .../internal/processservice/scope.go | 7 ++ .../internal/processservice/service_test.go | 12 +++- docs/process-protocol.md | 2 +- docs/zh/process-protocol.md | 4 +- 8 files changed, 131 insertions(+), 18 deletions(-) diff --git a/apps/sandboxio/internal/processservice/cgroup.go b/apps/sandboxio/internal/processservice/cgroup.go index 218e25c18..19925d40a 100644 --- a/apps/sandboxio/internal/processservice/cgroup.go +++ b/apps/sandboxio/internal/processservice/cgroup.go @@ -42,11 +42,31 @@ func delegatedCgroup() string { if err != nil || !slices.Contains(strings.Fields(string(procs)), strconv.Itoa(os.Getpid())) { return "" } + // Migrating a child requires write access to the common ancestor's + // cgroup.procs as well as the destination. Opening changes no membership. + migration, err := os.OpenFile(filepath.Join(parent, "cgroup.procs"), os.O_WRONLY, 0) + if err != nil { + return "" + } + migration.Close() g, err := newProcessCgroup(parent) if err != nil { return "" } - err = g.signal(unix.SIGKILL) // the exclusive probe group is empty + kind, err := g.root.ReadFile("cgroup.type") + if err == nil && strings.TrimSpace(string(kind)) != "domain" { + err = errors.New("operation cgroup is not a domain") + } + if err == nil { + var destination *os.File + destination, err = g.root.OpenFile("cgroup.procs", os.O_WRONLY, 0) + if err == nil { + destination.Close() + // Probe the control file directly: Signal on an empty scope + // must instead report NotRunning. + err = g.root.WriteFile("cgroup.kill", []byte("1"), 0) + } + } _, observation := g.observe() // No process was ever placed in this exclusive probe group. g.closed = true @@ -95,9 +115,16 @@ func (g *processCgroup) signal(sig unix.Signal) error { return notRunning("the cgroup is empty") } if sig == unix.SIGKILL { + live, observation := g.observeLocked() + if observation == nil && live == 0 { + return notRunning("the cgroup is empty") + } if err := g.root.WriteFile("cgroup.kill", []byte("1"), 0); err != nil { return sp.Fail(sp.CodeIO, sandboxwire.EffectPossible, "kill cgroup: %v", err) } + if observation != nil { + return sp.Fail(sp.CodeIO, sandboxwire.EffectPossible, "kill cgroup with unknown membership: %v", observation) + } return nil } // Membership comes from this operation's cgroup, never a process-tree scan. @@ -149,6 +176,10 @@ func (g *processCgroup) signal(sig unix.Signal) error { func (g *processCgroup) observe() (int, error) { g.mu.Lock() defer g.mu.Unlock() + return g.observeLocked() +} + +func (g *processCgroup) observeLocked() (int, error) { if g.closed { return 0, nil } diff --git a/apps/sandboxio/internal/processservice/cgroup_test.go b/apps/sandboxio/internal/processservice/cgroup_test.go index 4b87461bd..e663b7170 100644 --- a/apps/sandboxio/internal/processservice/cgroup_test.go +++ b/apps/sandboxio/internal/processservice/cgroup_test.go @@ -413,20 +413,24 @@ func TestCgroupEmptyDirectoryFailureDoesNotKeepScopeActive(t *testing.T) { func TestCgroupKillFailureKeepsNestedScopeActive(t *testing.T) { h := cgroupHarness(t) dir := t.TempDir() - script := `group=/sys/fs/cgroup$(sed -n 's/^0:://p' /proc/self/cgroup) + script := `trap '' TERM HUP +group=/sys/fs/cgroup$(sed -n 's/^0:://p' /proc/self/cgroup) mkdir "$group/nested" echo $$ > "$group/nested/cgroup.procs" (trap '' TERM HUP; printf a > "$0/ticks"; while [ ! -e "$0/gate" ]; do sleep .01; done; printf b >> "$0/ticks"; while :; do sleep 1; done) & -exit 0` +while :; do sleep 1; done` op := h.start(h.connect(), cgroupSpec("sh", "-c", script, dir)) awaitFile(t, filepath.Join(dir, "ticks"), "a") - events(t, op, sp.EventExited) + events(t, op, sp.EventStarted) group := cgroupRecord(t, h, op).cgroup.path killFile := filepath.Join(group, "cgroup.kill") if err := os.Chmod(killFile, 0); err != nil { t.Fatal(err) } - t.Cleanup(func() { _ = os.Chmod(killFile, 0200) }) + t.Cleanup(func() { + _ = os.Chmod(killFile, 0200) + _ = cgroupRecord(t, h, op).cgroup.signal(unix.SIGKILL) + }) wantCode(t, op.Signal(t.Context(), 9, sp.TargetScope), sp.CodeIO) if err := op.Cancel(t.Context(), 0); err != nil { t.Fatal(err) @@ -448,3 +452,59 @@ exit 0` t.Fatalf("nested scope was not removed: %v", err) } } + +func TestCgroupDelegationRequiresParentMigrationPermission(t *testing.T) { + h := cgroupHarness(t) + procs := filepath.Join(h.svc.cgroupParent, "cgroup.procs") + info, err := os.Stat(procs) + if err != nil { + t.Fatal(err) + } + if err := os.Chmod(procs, 0400); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(procs, info.Mode().Perm()) }) + service, err := New(DefaultConfig()) + if err != nil { + t.Fatal(err) + } + if slices.Contains(service.caps.Scopes, sp.ScopeCgroupV2) { + t.Fatal("advertised cgroup scope without parent migration permission") + } +} + +func TestCgroupKillEmptyReturnsNotRunning(t *testing.T) { + h := cgroupHarness(t) + group, err := newProcessCgroup(h.svc.cgroupParent) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _, _ = group.observe(); _ = group.cleanup() }) + wantCode(t, group.signal(unix.SIGKILL), sp.CodeNotRunning) +} + +func TestCgroupKillWithUnknownMembershipStillStopsProcesses(t *testing.T) { + h := cgroupHarness(t) + dir := t.TempDir() + op := h.start(h.connect(), cgroupSpec("sh", "-c", `trap '' TERM HUP; printf ready > "$0/ready"; while :; do sleep 1; done`, dir)) + awaitFile(t, filepath.Join(dir, "ready"), "ready") + group := cgroupRecord(t, h, op).cgroup.path + eventsFile := filepath.Join(group, "cgroup.events") + if err := os.Chmod(eventsFile, 0); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = os.Chmod(eventsFile, 0444) }) + // Unknown membership cannot prove a successful Signal, but must not + // prevent the cgroup kill from stopping processes whose observation failed. + wantCode(t, op.Signal(t.Context(), 9, sp.TargetScope), sp.CodeIO) + awaitFile(t, filepath.Join(group, "cgroup.procs"), "") + events(t, op, sp.EventObservationLost) + if st, err := op.Inspect(t.Context()); err != nil || st.Scope != sp.ScopeStateUnknown { + t.Fatalf("unknown scope was falsely settled: %+v %v", st, err) + } + wantCode(t, op.Release(t.Context()), sp.CodeBusy) + if err := os.Chmod(eventsFile, 0444); err != nil { + t.Fatal(err) + } + events(t, op, sp.EventScopeClosed) +} diff --git a/apps/sandboxio/internal/processservice/launch.go b/apps/sandboxio/internal/processservice/launch.go index 54ce14f76..29106983c 100644 --- a/apps/sandboxio/internal/processservice/launch.go +++ b/apps/sandboxio/internal/processservice/launch.go @@ -144,13 +144,21 @@ func (op *operation) spawn(req sp.StartRequest) (l launched, f *sp.Failure) { } reaping.RLock() p, err := os.StartProcess("/proc/self/exe", []string{trampolineArg0}, attr) - var ferr error + var ferr, placementErr error if err == nil { // The reaper cannot reap the child while reaping is held, so its PID // still names it. var fd int if fd, ferr = unix.PidfdOpen(p.Pid, 0); ferr == nil { op.cu = newCustody(p.Pid, fd, op.s.stat) + if op.cgroup != nil { + // The trusted trampoline is blocked on launchFD. Keep reaping + // excluded until placement: a pidfd alone does not reserve its PID. + placementErr = op.cgroup.place(p.Pid) + if placementErr != nil { + unix.Kill(p.Pid, unix.SIGKILL) + } + } } else { unix.Kill(p.Pid, unix.SIGKILL) // it has not read the launch, so it never execs } @@ -165,6 +173,9 @@ func (op *operation) spawn(req sp.StartRequest) (l launched, f *sp.Failure) { if ferr != nil { return l, ioFail("open process descriptor", ferr) } + if placementErr != nil { + return l, ioFail("place process in cgroup", placementErr) + } l.pid = pid op.mu.Lock() op.pid = l.pid // for killPinned; the rest is published with Started @@ -172,14 +183,6 @@ func (op *operation) spawn(req sp.StartRequest) (l launched, f *sp.Failure) { closeAll(child) child = nil - if op.cgroup != nil { - // The trusted trampoline is blocked reading launchFD. Place it before - // releasing the request, so the target never runs outside its scope. - if err := op.cgroup.place(pid); err != nil { - op.killPinned(true, unix.SIGKILL) - return l, ioFail("place process in cgroup", err) - } - } _, werr := launchW.Write(sp.Encode(req)) launchW.Close() status, rerr := io.ReadAll(statusR) diff --git a/apps/sandboxio/internal/processservice/operation.go b/apps/sandboxio/internal/processservice/operation.go index e87fb96a7..6a96eed33 100644 --- a/apps/sandboxio/internal/processservice/operation.go +++ b/apps/sandboxio/internal/processservice/operation.go @@ -80,6 +80,7 @@ type operation struct { killTimer *time.Timer killAt time.Time killing bool + watching bool } type stream struct { @@ -580,6 +581,7 @@ func (op *operation) kill() { op.killing = true op.mu.Unlock() op.signalScope(unix.SIGKILL) + go op.watchScope() } func (op *operation) release() error { diff --git a/apps/sandboxio/internal/processservice/scope.go b/apps/sandboxio/internal/processservice/scope.go index 91e3d92f5..95abeee6c 100644 --- a/apps/sandboxio/internal/processservice/scope.go +++ b/apps/sandboxio/internal/processservice/scope.go @@ -396,6 +396,13 @@ func (op *operation) signalScope(sig unix.Signal) error { // KILL escalation go on, and the operation settles only once a poll confirms // the scope empty. func (op *operation) watchScope() { + op.mu.Lock() + if op.watching { + op.mu.Unlock() + return + } + op.watching = true + op.mu.Unlock() delay := scopePollFirst for { op.mu.Lock() diff --git a/apps/sandboxio/internal/processservice/service_test.go b/apps/sandboxio/internal/processservice/service_test.go index 72146d091..ddb56857c 100644 --- a/apps/sandboxio/internal/processservice/service_test.go +++ b/apps/sandboxio/internal/processservice/service_test.go @@ -414,7 +414,17 @@ func TestCancelEscalatesToKill(t *testing.T) { if err := op.Cancel(context.Background(), 200); err != nil { t.Fatal(err) } - evs := events(t, op, sp.EventScopeClosed) + var evs []sp.Event + for exited, closed := false, false; !exited || !closed; { + ev := next(t, op) + evs = append(evs, ev) + switch ev.(type) { + case sp.ExitedEvent: + exited = true + case sp.ScopeClosedEvent: + closed = true + } + } if e, _ := find[sp.ExitedEvent](t, evs); e.Status.Signal != 9 { t.Fatalf("exit %+v", e.Status) } diff --git a/docs/process-protocol.md b/docs/process-protocol.md index d37f051c7..ad2eebf48 100644 --- a/docs/process-protocol.md +++ b/docs/process-protocol.md @@ -38,7 +38,7 @@ A service must: The Linux service calls `processservice.Init()` first in the binary's `main`. Go cannot set a child's umask, so each launch re-executes the service binary as a trampoline that reads the launch from an inherited descriptor, marks every inherited descriptor above 2 close-on-exec, applies the umask and working directory, and execs the target. `Init` runs that trampoline and returns at once in a normal start. The Linux service launches every operation in a new session with `setsid` and advertises `ScopePOSIXSession`. It requires `pidfd_open` and `pidfd_send_signal` (Linux 5.3 or later): without them `processservice.New` fails with `ErrPidfdUnsupported`. -The Linux implementation also advertises `ScopeCgroupV2` when its current cgroup, read from `/proc/self/cgroup` under `/sys/fs/cgroup`, permits creating and removing an operation subgroup, writing `cgroup.kill`, and observing `cgroup.events`. The Provider must delegate that current cgroup to the service account before starting SandboxIO; this requires a writable cgroup v2 mount and `cgroup.kill` (Linux 5.14 or later). No resource controller needs to be enabled. Each selected cgroup scope owns a new subgroup. The service places the blocked trampoline there before sending its launch spec, so the target and its descendants start inside the scope. A failed placement never executes the target. Delegation supplies process lifetime control, not isolation from other processes with the same account; the Provider owns that isolation. +The Linux implementation also advertises `ScopeCgroupV2` when its current cgroup, read from `/proc/self/cgroup` under `/sys/fs/cgroup`, permits creating and removing a domain operation subgroup, opening both parent and destination `cgroup.procs` for migration, writing `cgroup.kill`, and observing `cgroup.events`. The Provider must delegate that current cgroup to the service account before starting SandboxIO; this requires a writable cgroup v2 mount and `cgroup.kill` (Linux 5.14 or later). No resource controller needs to be enabled. Each selected cgroup scope owns a new subgroup. The service places the blocked trampoline there before sending its launch spec, so the target and its descendants start inside the scope. A failed placement never executes the target. Delegation supplies process lifetime control, not isolation from other processes with the same account; the Provider owns that isolation. Before serving, `main` makes the process a child subreaper (`prctl(PR_SET_CHILD_SUBREAPER)`) and runs `processservice.Reap(ctx)` for the life of the process. `Reap` is the process's only `wait`: it reaps every child, delivers each leader's exit to its operation, and reaps the orphaned descendants the subreaper inherits. Nothing else in the binary may wait for children, and no operation observes an exit while `Reap` is not running. When the binary stops, it calls `Shutdown(ctx)` after its streams have ended: `Shutdown` cancels every live operation and abandons its output, as revoking [ownership](#ownership) does, and returns once each scope has closed or `ctx` ends. diff --git a/docs/zh/process-protocol.md b/docs/zh/process-protocol.md index f5f94ade4..b2dedcf1d 100644 --- a/docs/zh/process-protocol.md +++ b/docs/zh/process-protocol.md @@ -1,7 +1,7 @@ --- title: "进程协议" source: docs/process-protocol.md -source_hash: 5111fa281a69270ee6301648134a3266003c055256d164574b9b9a5941a73564 +source_hash: ce053cea2ee275fd425025d1cf60635d68f01d3497515f7ae2dd97f84c135489 --- 进程协议定义 agent host 如何在沙箱中启动和控制进程。沙箱内的 Sandbox I/O 服务提供该协议,agent host 的 broker 是其客户端。协议依据明确的 spec 启动进程,以有序事件流式传输其输出,在精确 offset 处接受 stdin,并将 leader 退出、输出结束和进程 scope 结束作为独立事实报告。 @@ -40,7 +40,7 @@ Go 客户端为 `sandboxprocess.NewClient(stream)`。`Start` 和 `Attach` 返回 Linux 服务在二进制的 `main` 中首先调用 `processservice.Init()`。Go 无法设置子进程的 umask,因此每次启动都会将服务二进制作为 trampoline 重新执行:它从继承的描述符读取启动信息,将所有大于 2 的继承描述符标记为 close-on-exec,应用 umask 和工作目录,然后 exec 目标。`Init` 负责运行该 trampoline,正常启动时立即返回。Linux 服务用 `setsid` 在新会话中启动每个 operation,并声明 `ScopePOSIXSession`。它要求 `pidfd_open` 和 `pidfd_send_signal`(Linux 5.3 或更高版本):缺少它们时 `processservice.New` 以 `ErrPidfdUnsupported` 失败。 -Linux 实现从 `/proc/self/cgroup` 读取当前 cgroup,并映射到 `/sys/fs/cgroup`;当该组允许创建和删除 operation 子组、写入 `cgroup.kill` 及观测 `cgroup.events` 时,还会声明 `ScopeCgroupV2`。Provider 必须在启动 SandboxIO 前将该当前 cgroup 委派给服务账号;这需要可写的 cgroup v2 挂载和 `cgroup.kill`(Linux 5.14 或更高版本)。无需启用资源 controller。每个选择 cgroup scope 的 operation 独占一个新子组。服务先将阻塞中的 trampoline 放入该组,再发送 launch spec,保证目标及其后代从启动起就在 scope 内。放置失败时绝不执行目标。委派提供进程生命周期控制,不提供针对同账号其他进程的隔离;隔离由 Provider 负责。 +Linux 实现从 `/proc/self/cgroup` 读取当前 cgroup,并映射到 `/sys/fs/cgroup`;当该组允许创建和删除 domain operation 子组、以写入模式打开父组及目标组的 `cgroup.procs` 以迁移进程、写入 `cgroup.kill` 及观测 `cgroup.events` 时,还会声明 `ScopeCgroupV2`。Provider 必须在启动 SandboxIO 前将该当前 cgroup 委派给服务账号;这需要可写的 cgroup v2 挂载和 `cgroup.kill`(Linux 5.14 或更高版本)。无需启用资源 controller。每个选择 cgroup scope 的 operation 独占一个新子组。服务先将阻塞中的 trampoline 放入该组,再发送 launch spec,保证目标及其后代从启动起就在 scope 内。放置失败时绝不执行目标。委派提供进程生命周期控制,不提供针对同账号其他进程的隔离;隔离由 Provider 负责。 开始服务前,`main` 将进程设为 child subreaper(`prctl(PR_SET_CHILD_SUBREAPER)`),并在进程整个生命周期内运行 `processservice.Reap(ctx)`。`Reap` 是进程中唯一的 `wait`:它回收每个子进程,将每个 leader 的退出交付给对应 operation,并回收 subreaper 继承的孤儿后代进程。二进制中其他任何代码都不得等待子进程;`Reap` 未运行时,任何 operation 都观测不到退出。二进制停止时,在其 stream 结束后调用 `Shutdown(ctx)`:`Shutdown` 像撤销[所有权](#ownership)时一样,取消每个存活的 operation 并放弃其输出,在每个 scope 关闭或 `ctx` 结束后返回。