diff --git a/bpf/stepbpf.c b/bpf/stepbpf.c index 46bd2db..5ac35c3 100644 --- a/bpf/stepbpf.c +++ b/bpf/stepbpf.c @@ -14,21 +14,41 @@ //go:build ignore -// Step-attribution programs that require kernel BTF (tp_btf tracepoints and -// the cgroup/sock_create hook). Kept in a separate collection from tcbpf.c -// so a kernel without BTF degrades attribution only — the firewall proper -// loads and enforces unaffected. The three shared maps below are declared -// with identical shapes in tcbpf.c, which owns them; this collection is -// loaded with MapReplacements pointing at the tcbpf map fds. +// Step-attribution programs that require kernel BTF (tp_btf tracepoints, +// the cgroup/sock_create hook and the task iterator). Kept in a separate +// collection from tcbpf.c so a kernel without BTF degrades attribution +// only — the firewall proper loads and enforces unaffected. The shared maps +// below are declared with identical shapes in tcbpf.c, which owns them; +// this collection is loaded with MapReplacements pointing at the tcbpf map +// fds. +// +// Pid numbering. Everything the kernel compares or keys on in this file is +// a global (init-namespace) id: that is what the tracepoint arguments and +// bpf_get_current_pid_tgid() hand out. Userspace, however, only ever sees +// pids as numbered in the daemon's own pid namespace, and the two differ +// whenever cargowall runs inside a container (ARC runners are Kubernetes +// pods). map_task_nspid is the bridge — global tid → daemon-namespace +// tgid — seeded for pre-existing tasks by step_task_iter and kept current +// by step_fork/step_exit. It is what lets the reconciler read +// /proc//cmdline for a boundary event, lets userspace resolve the +// process behind a socket, and lets pkg/steps translate the /proc-numbered +// pids it discovers into the global ids the kernel needs. #include "vmlinux.h" #include #include +#include + +// Inode of the daemon's pid namespace (stat /proc/self/ns/pid), written by +// userspace before load. Picks the entry of a task's upid chain that is +// "the" pid userspace can act on. Zero disables translation outright +// rather than matching garbage. +const volatile __u32 pidns_ino = 0; // ---- Shared with tcbpf.c (replaced at load time; shapes must match) ---- struct step_state { - __u32 worker_tgid; // Runner.Worker tgid (0 = not discovered) + __u32 worker_tgid; // Runner.Worker global tgid (0 = not discovered) __u32 enabled; // 0 = feature off, all step hooks no-op __u64 next_ordinal; // next step ordinal, atomically incremented }; @@ -54,10 +74,19 @@ struct { __uint(max_entries, 65536); } map_sock_step SEC(".maps"); +struct { + __uint(type, BPF_MAP_TYPE_HASH); + __type(key, __u32); + __type(value, __u32); + __uint(max_entries, 32768); +} map_task_nspid SEC(".maps"); + // ---- Owned by this collection ---- // New-step notifications for the userspace reconciler: emitted once per new // direct child process of Runner.Worker, before the child runs user code. +// tgid is the child's pid as numbered in the daemon's namespace, so the +// reconciler can read its /proc entry directly (0 if untranslatable). struct step_child_event { __u32 tgid; __u32 ordinal; @@ -71,16 +100,68 @@ struct { __uint(max_entries, 64 * 1024); } map_step_events SEC(".maps"); +// One record per task in the daemon's pid-namespace subtree, streamed by +// step_task_iter. tid/tgid/ppid are global; ns_tgid is the process id +// userspace sees. pkg/steps builds its process-tree views from these +// instead of scanning /proc, which cannot reveal global ids from inside a +// namespace. +struct task_iter_rec { + __u32 tid; + __u32 tgid; + __u32 ppid; + __u32 ns_tgid; +}; + +const struct task_iter_rec *btf_anchor_task_iter_rec __attribute__((unused)); + +// MAX_PID_NS_LEVEL: a task's upid chain is at most this deep. +#define PIDNS_MAX_LEVEL 32 + static __always_inline struct step_state *step_state(void) { __u32 key = 0; return bpf_map_lookup_elem(&map_step_state, &key); } +// ns_tgid_of returns t's process id as numbered in the daemon's pid +// namespace, or 0 when t lives outside that namespace's subtree (another +// pod, the host) or translation is disabled. Walks the group leader's upid +// chain — numbers[0] is the init namespace, numbers[level] the task's own — +// looking for the namespace whose inode userspace passed in. The daemon's +// namespace appears at most once in a chain, so order is irrelevant; a +// plain bounded index keeps the flexible-array access simple. +static __always_inline __u32 ns_tgid_of(struct task_struct *t) +{ + if (!pidns_ino) + return 0; + struct pid *p = BPF_CORE_READ(t, group_leader, thread_pid); + if (!p) + return 0; + unsigned int level = BPF_CORE_READ(p, level); + if (level >= PIDNS_MAX_LEVEL) + level = PIDNS_MAX_LEVEL - 1; + for (unsigned int i = 0; i < PIDNS_MAX_LEVEL; i++) { + if (i > level) + break; + struct upid up; + if (bpf_core_read(&up, sizeof(up), &p->numbers[i])) + break; + if (BPF_CORE_READ(up.ns, ns.inum) == pidns_ino) + return up.nr; + } + return 0; +} + // Tag propagation. Runs synchronously inside fork(), in the parent's // context, so the child's tag exists before its first instruction — step // attribution has no detection race for host processes. // -// Two cases: +// The namespace bridge is maintained first and unconditionally (not gated +// on enabled): userspace relies on map_task_nspid being complete from +// attach time onward, so the seeding iterator only has to cover tasks that +// predate the attach. A new thread copies its process's entry (one lookup); +// a new process walks its upid chain. +// +// Then two tagging cases: // 1. New process (child tgid == child tid) whose parent process is // Runner.Worker → allocate the next step ordinal and notify userspace. // Worker *thread* creation deliberately falls through to case 2: a new @@ -92,22 +173,33 @@ static __always_inline struct step_state *step_state(void) { SEC("tp_btf/sched_process_fork") int BPF_PROG(step_fork, struct task_struct *parent, struct task_struct *child) { - struct step_state *st = step_state(); - if (!st || !st->enabled) - return 0; - __u32 parent_tid = parent->pid; __u32 parent_tgid = parent->tgid; __u32 child_tid = child->pid; __u32 child_tgid = child->tgid; + __u32 ns_tgid = 0; + if (child_tgid == parent_tgid) { + __u32 *pns = bpf_map_lookup_elem(&map_task_nspid, &parent_tid); + if (pns) + ns_tgid = *pns; + } + if (!ns_tgid) + ns_tgid = ns_tgid_of(child); + if (ns_tgid) + bpf_map_update_elem(&map_task_nspid, &child_tid, &ns_tgid, BPF_ANY); + + struct step_state *st = step_state(); + if (!st || !st->enabled) + return 0; + if (parent_tgid == st->worker_tgid && child_tgid == child_tid) { __u32 ord = (__u32)__sync_fetch_and_add(&st->next_ordinal, 1); bpf_map_update_elem(&map_task_step, &child_tid, &ord, BPF_ANY); struct step_child_event *ev = bpf_ringbuf_reserve(&map_step_events, sizeof(*ev), 0); if (ev) { - ev->tgid = child_tgid; + ev->tgid = ns_tgid; ev->ordinal = ord; bpf_ringbuf_submit(ev, 0); } @@ -120,14 +212,15 @@ int BPF_PROG(step_fork, struct task_struct *parent, struct task_struct *child) return 0; } -// PID-reuse hygiene: drop the tag when the task exits. Unconditional — a -// delete on an absent key is cheap, and gating on step_state would cost the -// same lookup. +// PID-reuse hygiene: drop the tag and the namespace entry when the task +// exits. Unconditional — a delete on an absent key is cheap, and gating on +// step_state would cost the same lookup. SEC("tp_btf/sched_process_exit") int BPF_PROG(step_exit, struct task_struct *task) { __u32 tid = task->pid; bpf_map_delete_elem(&map_task_step, &tid); + bpf_map_delete_elem(&map_task_nspid, &tid); return 0; } @@ -150,4 +243,32 @@ int cg_sock_create(struct bpf_sock *ctx) return 1; } +// Task walk for userspace. Two jobs in one pass: seed map_task_nspid for +// every task that predates the tracepoint attach, and stream the process +// tree (global ids plus the daemon-namespace tgid) so pkg/steps can seed +// step tags and tag container subtrees without /proc. The kernel already +// scopes the iteration to the opener's pid namespace — only the daemon's +// subtree is visited — and ns_tgid_of filters, defensively, anything it +// cannot map. Re-running it is idempotent, so container tagging takes a +// fresh pass per call. +SEC("iter/task") +int step_task_iter(struct bpf_iter__task *ctx) +{ + struct task_struct *t = ctx->task; + if (!t) + return 0; + __u32 ns_tgid = ns_tgid_of(t); + if (!ns_tgid) + return 0; + struct task_iter_rec rec = { + .tid = t->pid, + .tgid = t->tgid, + .ppid = BPF_CORE_READ(t, real_parent, tgid), + .ns_tgid = ns_tgid, + }; + bpf_map_update_elem(&map_task_nspid, &rec.tid, &ns_tgid, BPF_ANY); + bpf_seq_write(ctx->meta->seq, &rec, sizeof(rec)); + return 0; +} + char _license[] SEC("license") = "GPL"; diff --git a/bpf/stepbpf_bpfeb.go b/bpf/stepbpf_bpfeb.go index 085527e..d0c96f5 100644 --- a/bpf/stepbpf_bpfeb.go +++ b/bpf/stepbpf_bpfeb.go @@ -20,6 +20,14 @@ type StepBpfStepState struct { NextOrdinal uint64 } +type StepBpfTaskIterRec struct { + _ structs.HostLayout + Tid uint32 + Tgid uint32 + Ppid uint32 + NsTgid uint32 +} + // Names of all BPF objects in the ELF. // // Used for safe lookups in a Collection or CollectionSpec. @@ -27,11 +35,15 @@ const ( StepBpfMapMapSockStep = "map_sock_step" StepBpfMapMapStepEvents = "map_step_events" StepBpfMapMapStepState = "map_step_state" + StepBpfMapMapTaskNspid = "map_task_nspid" StepBpfMapMapTaskStep = "map_task_step" StepBpfProgCgSockCreate = "cg_sock_create" StepBpfProgStepExit = "step_exit" StepBpfProgStepFork = "step_fork" + StepBpfProgStepTaskIter = "step_task_iter" StepBpfVarBtfAnchorStepChildEvent = "btf_anchor_step_child_event" + StepBpfVarBtfAnchorTaskIterRec = "btf_anchor_task_iter_rec" + StepBpfVarPidnsIno = "pidns_ino" ) // LoadStepBpf returns the embedded CollectionSpec for StepBpf. @@ -79,6 +91,7 @@ type StepBpfProgramSpecs struct { CgSockCreate *ebpf.ProgramSpec `ebpf:"cg_sock_create"` StepExit *ebpf.ProgramSpec `ebpf:"step_exit"` StepFork *ebpf.ProgramSpec `ebpf:"step_fork"` + StepTaskIter *ebpf.ProgramSpec `ebpf:"step_task_iter"` } // StepBpfMapSpecs contains maps before they are loaded into the kernel. @@ -88,6 +101,7 @@ type StepBpfMapSpecs struct { MapSockStep *ebpf.MapSpec `ebpf:"map_sock_step"` MapStepEvents *ebpf.MapSpec `ebpf:"map_step_events"` MapStepState *ebpf.MapSpec `ebpf:"map_step_state"` + MapTaskNspid *ebpf.MapSpec `ebpf:"map_task_nspid"` MapTaskStep *ebpf.MapSpec `ebpf:"map_task_step"` } @@ -96,6 +110,8 @@ type StepBpfMapSpecs struct { // It can be passed ebpf.CollectionSpec.Assign. type StepBpfVariableSpecs struct { BtfAnchorStepChildEvent *ebpf.VariableSpec `ebpf:"btf_anchor_step_child_event"` + BtfAnchorTaskIterRec *ebpf.VariableSpec `ebpf:"btf_anchor_task_iter_rec"` + PidnsIno *ebpf.VariableSpec `ebpf:"pidns_ino"` } // StepBpfObjects contains all objects after they have been loaded into the kernel. @@ -121,6 +137,7 @@ type StepBpfMaps struct { MapSockStep *ebpf.Map `ebpf:"map_sock_step"` MapStepEvents *ebpf.Map `ebpf:"map_step_events"` MapStepState *ebpf.Map `ebpf:"map_step_state"` + MapTaskNspid *ebpf.Map `ebpf:"map_task_nspid"` MapTaskStep *ebpf.Map `ebpf:"map_task_step"` } @@ -129,6 +146,7 @@ func (m *StepBpfMaps) Close() error { m.MapSockStep, m.MapStepEvents, m.MapStepState, + m.MapTaskNspid, m.MapTaskStep, ) } @@ -138,6 +156,8 @@ func (m *StepBpfMaps) Close() error { // It can be passed to LoadStepBpfObjects or ebpf.CollectionSpec.LoadAndAssign. type StepBpfVariables struct { BtfAnchorStepChildEvent *ebpf.Variable `ebpf:"btf_anchor_step_child_event"` + BtfAnchorTaskIterRec *ebpf.Variable `ebpf:"btf_anchor_task_iter_rec"` + PidnsIno *ebpf.Variable `ebpf:"pidns_ino"` } // StepBpfPrograms contains all programs after they have been loaded into the kernel. @@ -147,6 +167,7 @@ type StepBpfPrograms struct { CgSockCreate *ebpf.Program `ebpf:"cg_sock_create"` StepExit *ebpf.Program `ebpf:"step_exit"` StepFork *ebpf.Program `ebpf:"step_fork"` + StepTaskIter *ebpf.Program `ebpf:"step_task_iter"` } func (p *StepBpfPrograms) Close() error { @@ -154,6 +175,7 @@ func (p *StepBpfPrograms) Close() error { p.CgSockCreate, p.StepExit, p.StepFork, + p.StepTaskIter, ) } diff --git a/bpf/stepbpf_bpfeb.o b/bpf/stepbpf_bpfeb.o index ea08050..2fe0741 100644 Binary files a/bpf/stepbpf_bpfeb.o and b/bpf/stepbpf_bpfeb.o differ diff --git a/bpf/stepbpf_bpfel.go b/bpf/stepbpf_bpfel.go index f195dc7..61cf7b9 100644 --- a/bpf/stepbpf_bpfel.go +++ b/bpf/stepbpf_bpfel.go @@ -20,6 +20,14 @@ type StepBpfStepState struct { NextOrdinal uint64 } +type StepBpfTaskIterRec struct { + _ structs.HostLayout + Tid uint32 + Tgid uint32 + Ppid uint32 + NsTgid uint32 +} + // Names of all BPF objects in the ELF. // // Used for safe lookups in a Collection or CollectionSpec. @@ -27,11 +35,15 @@ const ( StepBpfMapMapSockStep = "map_sock_step" StepBpfMapMapStepEvents = "map_step_events" StepBpfMapMapStepState = "map_step_state" + StepBpfMapMapTaskNspid = "map_task_nspid" StepBpfMapMapTaskStep = "map_task_step" StepBpfProgCgSockCreate = "cg_sock_create" StepBpfProgStepExit = "step_exit" StepBpfProgStepFork = "step_fork" + StepBpfProgStepTaskIter = "step_task_iter" StepBpfVarBtfAnchorStepChildEvent = "btf_anchor_step_child_event" + StepBpfVarBtfAnchorTaskIterRec = "btf_anchor_task_iter_rec" + StepBpfVarPidnsIno = "pidns_ino" ) // LoadStepBpf returns the embedded CollectionSpec for StepBpf. @@ -79,6 +91,7 @@ type StepBpfProgramSpecs struct { CgSockCreate *ebpf.ProgramSpec `ebpf:"cg_sock_create"` StepExit *ebpf.ProgramSpec `ebpf:"step_exit"` StepFork *ebpf.ProgramSpec `ebpf:"step_fork"` + StepTaskIter *ebpf.ProgramSpec `ebpf:"step_task_iter"` } // StepBpfMapSpecs contains maps before they are loaded into the kernel. @@ -88,6 +101,7 @@ type StepBpfMapSpecs struct { MapSockStep *ebpf.MapSpec `ebpf:"map_sock_step"` MapStepEvents *ebpf.MapSpec `ebpf:"map_step_events"` MapStepState *ebpf.MapSpec `ebpf:"map_step_state"` + MapTaskNspid *ebpf.MapSpec `ebpf:"map_task_nspid"` MapTaskStep *ebpf.MapSpec `ebpf:"map_task_step"` } @@ -96,6 +110,8 @@ type StepBpfMapSpecs struct { // It can be passed ebpf.CollectionSpec.Assign. type StepBpfVariableSpecs struct { BtfAnchorStepChildEvent *ebpf.VariableSpec `ebpf:"btf_anchor_step_child_event"` + BtfAnchorTaskIterRec *ebpf.VariableSpec `ebpf:"btf_anchor_task_iter_rec"` + PidnsIno *ebpf.VariableSpec `ebpf:"pidns_ino"` } // StepBpfObjects contains all objects after they have been loaded into the kernel. @@ -121,6 +137,7 @@ type StepBpfMaps struct { MapSockStep *ebpf.Map `ebpf:"map_sock_step"` MapStepEvents *ebpf.Map `ebpf:"map_step_events"` MapStepState *ebpf.Map `ebpf:"map_step_state"` + MapTaskNspid *ebpf.Map `ebpf:"map_task_nspid"` MapTaskStep *ebpf.Map `ebpf:"map_task_step"` } @@ -129,6 +146,7 @@ func (m *StepBpfMaps) Close() error { m.MapSockStep, m.MapStepEvents, m.MapStepState, + m.MapTaskNspid, m.MapTaskStep, ) } @@ -138,6 +156,8 @@ func (m *StepBpfMaps) Close() error { // It can be passed to LoadStepBpfObjects or ebpf.CollectionSpec.LoadAndAssign. type StepBpfVariables struct { BtfAnchorStepChildEvent *ebpf.Variable `ebpf:"btf_anchor_step_child_event"` + BtfAnchorTaskIterRec *ebpf.Variable `ebpf:"btf_anchor_task_iter_rec"` + PidnsIno *ebpf.Variable `ebpf:"pidns_ino"` } // StepBpfPrograms contains all programs after they have been loaded into the kernel. @@ -147,6 +167,7 @@ type StepBpfPrograms struct { CgSockCreate *ebpf.Program `ebpf:"cg_sock_create"` StepExit *ebpf.Program `ebpf:"step_exit"` StepFork *ebpf.Program `ebpf:"step_fork"` + StepTaskIter *ebpf.Program `ebpf:"step_task_iter"` } func (p *StepBpfPrograms) Close() error { @@ -154,6 +175,7 @@ func (p *StepBpfPrograms) Close() error { p.CgSockCreate, p.StepExit, p.StepFork, + p.StepTaskIter, ) } diff --git a/bpf/stepbpf_bpfel.o b/bpf/stepbpf_bpfel.o index cf24ea8..149e724 100644 Binary files a/bpf/stepbpf_bpfel.o and b/bpf/stepbpf_bpfel.o differ diff --git a/bpf/stepbpf_test.go b/bpf/stepbpf_test.go index c188acd..3b979b7 100644 --- a/bpf/stepbpf_test.go +++ b/bpf/stepbpf_test.go @@ -17,9 +17,12 @@ package bpf import ( + "io" "os" "os/exec" "strconv" + "strings" + "syscall" "testing" "time" "unsafe" @@ -37,17 +40,51 @@ type stepChildEventT struct { Ordinal uint32 } +// pidnsInode is the daemon-side half of the namespace bridge: the nsfs +// inode of this process's pid namespace, as pkg/steps passes it in. +func pidnsInode(t *testing.T) uint32 { + t.Helper() + fi, err := os.Stat("/proc/self/ns/pid") + require.NoError(t, err) + return uint32(fi.Sys().(*syscall.Stat_t).Ino) +} + +// globalTgid runs the task iterator and returns the global tgid of the +// process numbered nsPid in this namespace — the translation pkg/steps +// performs for the worker pid and for seeding. +func globalTgid(t *testing.T, iter *link.Iter, nsPid int) (uint32, bool) { + t.Helper() + rd, err := iter.Open() + require.NoError(t, err) + defer rd.Close() + data, err := io.ReadAll(rd) + require.NoError(t, err) + const recSize = int(unsafe.Sizeof(StepBpfTaskIterRec{})) + for off := 0; off+recSize <= len(data); off += recSize { + rec := (*StepBpfTaskIterRec)(unsafe.Pointer(&data[off])) + if rec.Tid == rec.Tgid && int(rec.NsTgid) == nsPid { + return rec.Tgid, true + } + } + return 0, false +} + // TestStepFork verifies the step-attribution collection end to end on a real // kernel: it loads StepBpf against TcBpf's shared maps, attaches the tp_btf -// tracepoints, declares the test process itself as Runner.Worker, forks a -// child, and asserts the kernel assigned it the configured ordinal and -// emitted the new-step ringbuf notification — the exact contract -// pkg/steps.Start builds on. Requires root and kernel BTF. +// tracepoints and the task iterator, declares the test process itself as +// Runner.Worker (by the global tgid the iterator reports for our /proc +// pid), forks a child, and asserts the kernel assigned it the configured +// ordinal under its global tid, recorded its namespace pid in +// map_task_nspid, and emitted the new-step ringbuf notification carrying +// that namespace pid — the exact contract pkg/steps.Start builds on. +// Requires root and kernel BTF. TestStepFork_InPidNamespace re-runs it +// where global and namespace pids differ. func TestStepFork(t *testing.T) { tcObjs := loadBPFObjects(t) spec, err := LoadStepBpf() require.NoError(t, err) + require.NoError(t, spec.Variables[StepBpfVarPidnsIno].Set(pidnsInode(t))) var objs StepBpfObjects if err := spec.LoadAndAssign(&objs, &ebpf.CollectionOptions{ @@ -55,6 +92,7 @@ func TestStepFork(t *testing.T) { "map_task_step": tcObjs.MapTaskStep, "map_sock_step": tcObjs.MapSockStep, "map_step_state": tcObjs.MapStepState, + "map_task_nspid": tcObjs.MapTaskNspid, }, }); err != nil { t.Skipf("step BPF objects not loadable (kernel BTF required): %v", err) @@ -75,16 +113,23 @@ func TestStepFork(t *testing.T) { require.NoError(t, err, "attach sched_process_exit") t.Cleanup(func() { exitLink.Close() }) + iter, err := link.AttachIter(link.IterOptions{Program: objs.StepTaskIter}) + require.NoError(t, err, "attach task iterator") + t.Cleanup(func() { iter.Close() }) + rd, err := ringbuf.NewReader(objs.MapStepEvents) require.NoError(t, err) t.Cleanup(func() { rd.Close() }) - // Declare this test process as the worker: its next direct child gets - // ordinal 5. Uses the bpf2go-generated state struct so kernel and Go - // layouts cannot drift. + // Declare this test process as the worker — by global tgid, which is + // what step_fork compares parent->tgid against — so its next direct + // child gets ordinal 5. Uses the bpf2go-generated state struct so kernel + // and Go layouts cannot drift. + selfTgid, ok := globalTgid(t, iter, os.Getpid()) + require.True(t, ok, "the iterator must list this process") const wantOrdinal = 5 require.NoError(t, tcObjs.MapStepState.Put(uint32(0), TcBpfStepState{ - WorkerTgid: uint32(os.Getpid()), + WorkerTgid: selfTgid, Enabled: 1, NextOrdinal: wantOrdinal, })) @@ -94,11 +139,17 @@ func TestStepFork(t *testing.T) { childPID := uint32(child.Process.Pid) defer func() { _ = child.Wait() }() - // The fork hook runs synchronously inside fork(), so the tag must be - // visible as soon as Start returns. + // The fork hook runs synchronously inside fork(), so both the tag (under + // the child's global tid) and its namespace-pid entry must be visible as + // soon as Start returns. + childTgid, ok := globalTgid(t, iter, int(childPID)) + require.True(t, ok, "the iterator must list the child") var got uint32 - require.NoError(t, tcObjs.MapTaskStep.Lookup(childPID, &got), + require.NoError(t, tcObjs.MapTaskStep.Lookup(childTgid, &got), "child tid must be tagged immediately after fork") + var nsPid uint32 + require.NoError(t, tcObjs.MapTaskNspid.Lookup(childTgid, &nsPid)) + require.Equal(t, childPID, nsPid, "map_task_nspid must carry the child's pid as we number it") // The new-step notification must arrive for the reconciler. Other test // machinery may fork too, so scan until the child's event or deadline, @@ -119,17 +170,41 @@ func TestStepFork(t *testing.T) { break } - // Exit hygiene: once the child is reaped, its tag must be gone. + // Exit hygiene: once the child is reaped, its tag and namespace entry + // must be gone. require.NoError(t, child.Wait()) require.Eventually(t, func() bool { var v uint32 - return tcObjs.MapTaskStep.Lookup(childPID, &v) != nil - }, 2*time.Second, 50*time.Millisecond, "exit tracepoint must delete the child's tag") + return tcObjs.MapTaskStep.Lookup(childTgid, &v) != nil && + tcObjs.MapTaskNspid.Lookup(childTgid, &v) != nil + }, 2*time.Second, 50*time.Millisecond, "exit tracepoint must delete the child's entries") // Disable so later tests in this package see inert step hooks. require.NoError(t, tcObjs.MapStepState.Put(uint32(0), TcBpfStepState{})) } +// TestStepFork_InPidNamespace re-executes TestStepFork inside a fresh pid +// (and mount, for /proc) namespace — the shape of an ARC runner pod, where +// /proc pids and the kernel's global pids disagree and the pre-iterator +// contract failed on the very first lookup. +func TestStepFork_InPidNamespace(t *testing.T) { + if os.Geteuid() != 0 { + t.Skip("needs root") + } + unshare, err := exec.LookPath("unshare") + if err != nil { + t.Skip("unshare not available") + } + out, err := exec.Command(unshare, "--pid", "--fork", "--mount-proc", + os.Args[0], "-test.run", "^TestStepFork$", "-test.v").CombinedOutput() + t.Logf("inside pid namespace:\n%s", out) + if strings.Contains(string(out), "--- SKIP") { + t.Skip("inner test skipped") + } + require.NoError(t, err) + require.Contains(t, string(out), "--- PASS: TestStepFork") +} + // TestBlockedEventLayoutMatchesBTF pins the kernel↔userspace event layout: // every member of the C struct blocked_event (as compiled, via BTF) must sit // at the same offset as its field in the Go mirror, and the sizes must @@ -196,6 +271,37 @@ func TestStepChildEventLayoutMatchesBTF(t *testing.T) { require.Equal(t, uintptr(st.Members[1].Offset.Bytes()), unsafe.Offsetof(ev.Ordinal)) } +// TestTaskIterRecLayoutMatchesBTF pins the iterator record layout against +// the bpf2go mirror pkg/steps decodes with: the C struct is four packed +// u32s, and a regenerated mirror that drifted (a field reordered or +// widened in C without `go generate`) would break the unsafe cast. +func TestTaskIterRecLayoutMatchesBTF(t *testing.T) { + spec, err := LoadStepBpf() + require.NoError(t, err) + + typ, err := spec.Types.AnyTypeByName("task_iter_rec") + require.NoError(t, err) + st, ok := typ.(*btf.Struct) + require.True(t, ok, "task_iter_rec must be a struct, got %T", typ) + + var rec StepBpfTaskIterRec + require.Equal(t, uint32(unsafe.Sizeof(rec)), st.Size, "struct size") + require.Len(t, st.Members, 4) + want := []struct { + name string + off uintptr + }{ + {"tid", unsafe.Offsetof(rec.Tid)}, + {"tgid", unsafe.Offsetof(rec.Tgid)}, + {"ppid", unsafe.Offsetof(rec.Ppid)}, + {"ns_tgid", unsafe.Offsetof(rec.NsTgid)}, + } + for i, w := range want { + require.Equal(t, w.name, st.Members[i].Name) + require.Equal(t, w.off, uintptr(st.Members[i].Offset.Bytes())) + } +} + // TestStepMapsShrunkLoad pins the memory optimization in cmd/start.go: // with step attribution off, the step maps are resized down before load, // and the collection must still pass the verifier at the small sizes. @@ -205,6 +311,7 @@ func TestStepMapsShrunkLoad(t *testing.T) { require.NoError(t, err) spec.Maps["map_task_step"].MaxEntries = 64 spec.Maps["map_sock_step"].MaxEntries = 64 + spec.Maps["map_task_nspid"].MaxEntries = 64 var objs TcBpfObjects require.NoError(t, spec.LoadAndAssign(&objs, nil)) diff --git a/bpf/tcbpf.c b/bpf/tcbpf.c index 72d9945..ffb1058 100644 --- a/bpf/tcbpf.c +++ b/bpf/tcbpf.c @@ -101,7 +101,9 @@ struct { __uint(max_entries, 256 * 1024); } map_events SEC(".maps"); -// LRU hash map: socket cookie → PID (populated by cgroup programs, read by TC) +// LRU hash map: socket cookie → PID (populated by cgroup programs, read by +// TC). The PID is numbered in the daemon's pid namespace when the step +// tracker is running — see map_task_nspid — so userspace can resolve it. struct { __uint(type, BPF_MAP_TYPE_LRU_HASH); __type(key, __u64); @@ -153,7 +155,7 @@ struct { #define STEP_ORD_PREDAEMON 0xFFFFFFFEu // worker children that predate cargowall struct step_state { - __u32 worker_tgid; // Runner.Worker tgid (0 = not discovered) + __u32 worker_tgid; // Runner.Worker global tgid (0 = not discovered) __u32 enabled; // 0 = feature off, all step hooks no-op __u64 next_ordinal; // next step ordinal, atomically incremented }; @@ -183,6 +185,31 @@ struct { __uint(max_entries, 65536); } map_sock_step SEC(".maps"); +// global tid → tgid as numbered in the daemon's pid namespace. Maintained +// by stepbpf.c (fork/exit tracepoints plus the task iterator that seeds +// pre-existing tasks); empty when step attribution is off. Kernel-side +// identity stays global everywhere — this table exists so the pid handed +// to userspace is one it can look up in its own /proc, which differs from +// the global number whenever cargowall runs inside a container (ARC +// runners). Keyed by tid so any thread resolves its process. +struct { + __uint(type, BPF_MAP_TYPE_HASH); + __type(key, __u32); + __type(value, __u32); + __uint(max_entries, 32768); +} map_task_nspid SEC(".maps"); + +// Helper: the calling thread's process id as userspace should see it — the +// daemon-namespace tgid when the step tracker has translated this task, +// the global tgid otherwise (identical outside a pid namespace, and the +// best available answer while step attribution is off). +static __always_inline __u32 sock_owner_pid(void) { + __u64 pid_tgid = bpf_get_current_pid_tgid(); + __u32 tid = (__u32)pid_tgid; + __u32 *ns = bpf_map_lookup_elem(&map_task_nspid, &tid); + return ns ? *ns : (__u32)(pid_tgid >> 32); +} + // Helper: copy the calling thread's step tag onto a socket cookie. // Shared by the connect/sendmsg hooks below (fallback for sockets created // before attach) and cg_sock_create in stepbpf.c (primary path). @@ -543,7 +570,7 @@ static __always_inline int handle_ipv6(struct __sk_buff *skb, __u32 l3_offset) { SEC("cgroup/connect4") int cg_connect4(struct bpf_sock_addr *ctx) { __u64 cookie = bpf_get_socket_cookie(ctx); - __u32 pid = bpf_get_current_pid_tgid() >> 32; + __u32 pid = sock_owner_pid(); bpf_map_update_elem(&map_sock_pid, &cookie, &pid, BPF_ANY); step_tag_socket(cookie); return 1; @@ -552,7 +579,7 @@ int cg_connect4(struct bpf_sock_addr *ctx) { SEC("cgroup/connect6") int cg_connect6(struct bpf_sock_addr *ctx) { __u64 cookie = bpf_get_socket_cookie(ctx); - __u32 pid = bpf_get_current_pid_tgid() >> 32; + __u32 pid = sock_owner_pid(); bpf_map_update_elem(&map_sock_pid, &cookie, &pid, BPF_ANY); step_tag_socket(cookie); return 1; @@ -563,7 +590,7 @@ int cg_connect6(struct bpf_sock_addr *ctx) { SEC("cgroup/sendmsg4") int cg_sendmsg4(struct bpf_sock_addr *ctx) { __u64 cookie = bpf_get_socket_cookie(ctx); - __u32 pid = bpf_get_current_pid_tgid() >> 32; + __u32 pid = sock_owner_pid(); bpf_map_update_elem(&map_sock_pid, &cookie, &pid, BPF_ANY); step_tag_socket(cookie); return 1; @@ -572,7 +599,7 @@ int cg_sendmsg4(struct bpf_sock_addr *ctx) { SEC("cgroup/sendmsg6") int cg_sendmsg6(struct bpf_sock_addr *ctx) { __u64 cookie = bpf_get_socket_cookie(ctx); - __u32 pid = bpf_get_current_pid_tgid() >> 32; + __u32 pid = sock_owner_pid(); bpf_map_update_elem(&map_sock_pid, &cookie, &pid, BPF_ANY); step_tag_socket(cookie); return 1; diff --git a/bpf/tcbpf_bpfeb.go b/bpf/tcbpf_bpfeb.go index 5dd01ab..6934e0f 100644 --- a/bpf/tcbpf_bpfeb.go +++ b/bpf/tcbpf_bpfeb.go @@ -77,6 +77,7 @@ const ( TcBpfMapMapSockPid = "map_sock_pid" TcBpfMapMapSockStep = "map_sock_step" TcBpfMapMapStepState = "map_step_state" + TcBpfMapMapTaskNspid = "map_task_nspid" TcBpfMapMapTaskStep = "map_task_step" TcBpfProgCgConnect4 = "cg_connect4" TcBpfProgCgConnect6 = "cg_connect6" @@ -153,6 +154,7 @@ type TcBpfMapSpecs struct { MapSockPid *ebpf.MapSpec `ebpf:"map_sock_pid"` MapSockStep *ebpf.MapSpec `ebpf:"map_sock_step"` MapStepState *ebpf.MapSpec `ebpf:"map_step_state"` + MapTaskNspid *ebpf.MapSpec `ebpf:"map_task_nspid"` MapTaskStep *ebpf.MapSpec `ebpf:"map_task_step"` } @@ -195,6 +197,7 @@ type TcBpfMaps struct { MapSockPid *ebpf.Map `ebpf:"map_sock_pid"` MapSockStep *ebpf.Map `ebpf:"map_sock_step"` MapStepState *ebpf.Map `ebpf:"map_step_state"` + MapTaskNspid *ebpf.Map `ebpf:"map_task_nspid"` MapTaskStep *ebpf.Map `ebpf:"map_task_step"` } @@ -212,6 +215,7 @@ func (m *TcBpfMaps) Close() error { m.MapSockPid, m.MapSockStep, m.MapStepState, + m.MapTaskNspid, m.MapTaskStep, ) } diff --git a/bpf/tcbpf_bpfeb.o b/bpf/tcbpf_bpfeb.o index 91cc607..3503beb 100644 Binary files a/bpf/tcbpf_bpfeb.o and b/bpf/tcbpf_bpfeb.o differ diff --git a/bpf/tcbpf_bpfel.go b/bpf/tcbpf_bpfel.go index 43c3529..4434a5c 100644 --- a/bpf/tcbpf_bpfel.go +++ b/bpf/tcbpf_bpfel.go @@ -77,6 +77,7 @@ const ( TcBpfMapMapSockPid = "map_sock_pid" TcBpfMapMapSockStep = "map_sock_step" TcBpfMapMapStepState = "map_step_state" + TcBpfMapMapTaskNspid = "map_task_nspid" TcBpfMapMapTaskStep = "map_task_step" TcBpfProgCgConnect4 = "cg_connect4" TcBpfProgCgConnect6 = "cg_connect6" @@ -153,6 +154,7 @@ type TcBpfMapSpecs struct { MapSockPid *ebpf.MapSpec `ebpf:"map_sock_pid"` MapSockStep *ebpf.MapSpec `ebpf:"map_sock_step"` MapStepState *ebpf.MapSpec `ebpf:"map_step_state"` + MapTaskNspid *ebpf.MapSpec `ebpf:"map_task_nspid"` MapTaskStep *ebpf.MapSpec `ebpf:"map_task_step"` } @@ -195,6 +197,7 @@ type TcBpfMaps struct { MapSockPid *ebpf.Map `ebpf:"map_sock_pid"` MapSockStep *ebpf.Map `ebpf:"map_sock_step"` MapStepState *ebpf.Map `ebpf:"map_step_state"` + MapTaskNspid *ebpf.Map `ebpf:"map_task_nspid"` MapTaskStep *ebpf.Map `ebpf:"map_task_step"` } @@ -212,6 +215,7 @@ func (m *TcBpfMaps) Close() error { m.MapSockPid, m.MapSockStep, m.MapStepState, + m.MapTaskNspid, m.MapTaskStep, ) } diff --git a/bpf/tcbpf_bpfel.o b/bpf/tcbpf_bpfel.o index 6670b13..ea37333 100644 Binary files a/bpf/tcbpf_bpfel.o and b/bpf/tcbpf_bpfel.o differ diff --git a/cargowall.go b/cargowall.go index 9712012..a49ba4d 100644 --- a/cargowall.go +++ b/cargowall.go @@ -23,7 +23,7 @@ import ( ) //go:generate go run github.com/cilium/ebpf/cmd/bpf2go -tags linux -go-package bpf -cc clang -output-dir bpf TcBpf bpf/tcbpf.c -//go:generate go run github.com/cilium/ebpf/cmd/bpf2go -tags linux -go-package bpf -cc clang -output-dir bpf StepBpf bpf/stepbpf.c +//go:generate go run github.com/cilium/ebpf/cmd/bpf2go -tags linux -go-package bpf -cc clang -output-dir bpf -type task_iter_rec StepBpf bpf/stepbpf.c //go:generate go run github.com/cilium/ebpf/cmd/bpf2go -tags linux -go-package bpf -cc clang -output-dir bpf OriginBpf bpf/originbpf.c var version = "dev" diff --git a/cmd/start.go b/cmd/start.go index 463edca..d3773ec 100644 --- a/cmd/start.go +++ b/cmd/start.go @@ -503,14 +503,15 @@ func startCargoWall(cmd *StartCmd, hooks *StartHooks, teardowns *teardownList) e if err != nil { return fmt.Errorf("failed to load TC eBPF spec: %w", err) } - // The step-attribution maps preallocate ~7MB of kernel memory at their - // full size (LRU hash always preallocates). Shrink them when the feature - // is off — every non-GitHub run — but never when it is on: stepbpf.c - // declares the full sizes and MapReplacements rejects mismatched specs. + // The step-attribution maps preallocate ~8MB of kernel memory at their + // full size (hash maps preallocate by default). Shrink them when the + // feature is off — every non-GitHub run — but never when it is on: + // stepbpf.c declares the full sizes and MapReplacements rejects + // mismatched specs. if !cmd.StepAttribution { // Present unless the generated spec is stale (go generate not run); // verify-bpf-generated-code guards that, but don't panic if it slips. - for _, name := range []string{"map_task_step", "map_sock_step"} { + for _, name := range []string{"map_task_step", "map_sock_step", "map_task_nspid"} { if m := spec.Maps[name]; m != nil { m.MaxEntries = 64 } @@ -536,6 +537,12 @@ func startCargoWall(cmd *StartCmd, hooks *StartHooks, teardowns *teardownList) e // Attach cgroup programs for PID tracking via socket cookie. // Best-effort: if attachment fails, TC filtering still works but PID will be 0. + // Attached before steps.Start on purpose: a socket connecting before + // the step tracker has seeded map_task_nspid records the global tgid + // (unresolvable from inside a pid namespace), but attaching later would + // leave those same sockets with no pid at all. The window is daemon + // startup — the job's step is still waiting on readiness — and only the + // process name of a connection made in it is affected. cgroupProgs := []struct { prog *ebpf.Program attach ebpf.AttachType diff --git a/design.md b/design.md index c5911a7..bba556b 100644 --- a/design.md +++ b/design.md @@ -398,10 +398,35 @@ sequenceDiagram | `cg_connect6` | cgroup/connect6 | Maps IPv6 TCP socket cookie → PID | | `cg_sendmsg4` | cgroup/sendmsg4 | Maps IPv4 UDP socket cookie → PID | | `cg_sendmsg6` | cgroup/sendmsg6 | Maps IPv6 UDP socket cookie → PID | -| `step_fork` / `step_exit` | tp_btf tracepoints (stepbpf.c, separate collection; needs kernel BTF) | Step attribution: mint/inherit per-step process tags, drop at exit | +| `step_fork` / `step_exit` | tp_btf tracepoints (stepbpf.c, separate collection; needs kernel BTF) | Step attribution: mint/inherit per-step process tags, drop at exit. Also maintains `map_task_nspid` (global tid → tgid as numbered in the daemon's pid namespace) for every fork | +| `step_task_iter` | iter/task (stepbpf.c) | Task walk run by `pkg/steps` at start and per container tag: seeds `map_task_nspid` for tasks that predate the attach and streams the process tree in global ids, so userspace never has to translate pids through `/proc` | | `cg_sock_create` | cgroup/sock_create (stepbpf.c) | Copies the creating thread's step tag onto each socket cookie | | `cg_origin_egress` | cgroup_skb/egress (originbpf.c, separate collection) | Flow-origin recorder and — in enforce mode — the primary egress verdict (issue #106). Runs in socket context, pre-NAT, inside the originating netns, so it can both attribute and enforce container traffic. Behavior is set by a runtime mode: observe (pass, record only), shadow (compute the verdict, report would-blocks, pass), enforce (drop denied traffic). Shares the rule maps with `tc_egress` via `bpf/verdict.h` | +### Pid namespaces (issue #135) + +The kernel numbers tasks globally; everything the daemon reads from `/proc` +— the `Runner.Worker` pid, its own pid, container leaders reported by +dockerd — is numbered in the daemon's pid namespace. On a GitHub-hosted VM +the two coincide. Inside a container (ARC runners are Kubernetes pods) they +do not, and before #135 the worker's `/proc` pid was compared against +`parent->tgid` in `step_fork`, seeding wrote `/proc` tids as map keys, and +`map_sock_pid` handed userspace global tgids it could not resolve — every +event came out unattributed with an empty process name. + +The rule now is: kernel-side identity is global everywhere; translation +happens at the boundary, in the kernel, through `map_task_nspid`. +`step_task_iter` (scoped by the kernel to the opener's namespace) seeds it +and gives `pkg/steps` the process tree in global ids; `step_fork` keeps it +current (a thread copies its process's entry, a process walks its upid +chain to the daemon's namespace, identified by the nsfs inode passed in as +`pidns_ino`); the cgroup hooks write the namespace tgid into `map_sock_pid`; +boundary events carry the namespace pid so the reconciler's `/proc` read +works. Kernel floor is unchanged: the `iter/task` target, `bpf_seq_write` and +the `BPF_ITER_CREATE` command that opens an iterator instance all shipped in +5.8 (uapi `bpf.h` at v5.8 lists `BPF_ITER_CREATE`; v5.7 does not), like the +ring buffer. + ## Audit Mode Audit mode allows CargoWall to run in a log-only configuration — all traffic decisions are recorded but no packets are dropped. diff --git a/pkg/steps/steps.go b/pkg/steps/steps.go index 9633c72..02e4042 100644 --- a/pkg/steps/steps.go +++ b/pkg/steps/steps.go @@ -21,17 +21,30 @@ // thread's tag, and the cgroup hooks copy the tag onto each socket cookie — // so every TC event carries the step that (transitively) created its socket. // Attribution only: no verdict consults these maps. +// +// Pid numbering: the kernel side keys everything by global (init-namespace) +// ids; everything this package reads from /proc — the worker pid, the +// daemon's own pid, container leaders reported by dockerd — is numbered in +// the daemon's pid namespace, a different space whenever cargowall runs +// inside a container (ARC runners are Kubernetes pods). The step_task_iter +// BPF iterator bridges the two: it walks the daemon's namespace subtree, +// seeds the kernel's global→namespace table (map_task_nspid, which is what +// makes boundary events and socket pids resolvable) and streams the +// process tree in global ids, so seeding and container tagging never scan +// /proc. On a plain VM the translation is the identity. package steps import ( "errors" "fmt" + "io" "log/slog" "math" "os" "strconv" "strings" "sync" + "syscall" "time" "unicode/utf8" "unsafe" @@ -78,13 +91,16 @@ const maxOrdinalBase = uint64(events.MaxRealStepOrdinal) // Tracker owns the step-attribution BPF programs and the reconciler // goroutine that turns new-step ringbuf notifications into audit events. type Tracker struct { - workerPID int + workerPID int // as numbered in the daemon's pid namespace (/proc) + workerTgid uint32 // the same process, global — what step_fork compares workerCmdline string objs bpf.StepBpfObjects links []link.Link + iter *link.Iter // step_task_iter; also in links for Close reader *ringbuf.Reader stateMap *ebpf.Map taskMap *ebpf.Map + nspidMap *ebpf.Map // map_task_nspid; cleared at Close, see there auditLogger *events.AuditLogger logger *slog.Logger done chan struct{} @@ -96,6 +112,10 @@ type Tracker struct { // DNS-path client attribution (sock_diag lookups, caching, map joins) — // see clientResolver. Constructed in Start, closed in Close. resolver *clientResolver + + // snapshotFn produces the process tree the seeding and container-tagging + // paths work from; iterSnapshot in production, injectable for tests. + snapshotFn func() (*procSnapshot, error) } // Start loads the step-attribution collection against tcObjs' shared maps, @@ -129,10 +149,24 @@ func Start(tcObjs *bpf.TcBpfObjects, opts Options, auditLogger *events.AuditLogg return nil, fmt.Errorf("worker pid %d out of range", workerPID) } + // Which pid namespace "our" pids belong to. Every /proc-derived id in + // this package (the worker pid above, seeding, container leaders) is + // numbered in it; the kernel numbers tasks globally. Identical on a + // plain VM, different inside a container (ARC runners). + pidnsIno, err := pidNamespaceInode() + if err != nil { + return nil, err + } + spec, err := bpf.LoadStepBpf() if err != nil { return nil, fmt.Errorf("failed to load step BPF spec: %w", err) } + if v := spec.Variables[bpf.StepBpfVarPidnsIno]; v == nil { + return nil, errors.New("step BPF spec has no pidns_ino variable (stale generated code)") + } else if err := v.Set(pidnsIno); err != nil { + return nil, fmt.Errorf("failed to set pidns_ino: %w", err) + } t := &Tracker{ workerPID: workerPID, @@ -141,19 +175,22 @@ func Start(tcObjs *bpf.TcBpfObjects, opts Options, auditLogger *events.AuditLogg workerCmdline: readCmdline(workerPID), stateMap: tcObjs.MapStepState, taskMap: tcObjs.MapTaskStep, + nspidMap: tcObjs.MapTaskNspid, auditLogger: auditLogger, logger: logger, done: make(chan struct{}), resolver: newClientResolver(tcObjs.MapSockStep, tcObjs.MapSockPid, logger), } + t.snapshotFn = t.iterSnapshot - // The three shared maps are owned by the tcbpf collection; replacing them - // here makes both collections operate on the same kernel maps. + // The shared maps are owned by the tcbpf collection; replacing them here + // makes both collections operate on the same kernel maps. if err := spec.LoadAndAssign(&t.objs, &ebpf.CollectionOptions{ MapReplacements: map[string]*ebpf.Map{ "map_task_step": tcObjs.MapTaskStep, "map_sock_step": tcObjs.MapSockStep, "map_step_state": tcObjs.MapStepState, + "map_task_nspid": tcObjs.MapTaskNspid, }, }); err != nil { return nil, fmt.Errorf("failed to load step BPF objects (kernel BTF required): %w", err) @@ -164,22 +201,47 @@ func Start(tcObjs *bpf.TcBpfObjects, opts Options, auditLogger *events.AuditLogg return nil, err } + // First task walk: seeds map_task_nspid for everything that predates the + // attach (the tracepoints keep it current from here on) and hands back + // the process tree in global ids — which is how the worker's /proc pid + // becomes the global tgid step_fork compares against. + snap, err := t.snapshotFn() + if err != nil { + t.Close() + return nil, err + } + workerTgid, ok := snap.global[workerPID] + if !ok { + t.Close() + return nil, fmt.Errorf("%s pid %d is not visible to the task iterator", workerComm, workerPID) + } + t.workerTgid = workerTgid + if pidnsIno != initPidNamespaceInode { + logger.Info("Running inside a pid namespace; translating task ids", + "pidns_ino", pidnsIno, "worker_pid", workerPID, "worker_tgid", workerTgid) + } + // Seed while the programs are attached but still disabled (enabled=0 makes - // them no-ops), then enable, then seed once more: the second pass catches - // processes forked during the first scan, so the only unattributed window - // is forks that both start and create sockets between the enable flip and - // the rescan. Seeding is strictly create-only (BPF_NOEXIST), so the second - // pass can never clobber an ordinal the now-live fork tracepoint assigned. - t.seedExisting() + // the tagging hooks no-ops), then enable, then seed once more from a fresh + // walk: the second pass catches processes forked during the first, so the + // only unattributed window is forks that both start and create sockets + // between the enable flip and the rescan. Seeding is strictly create-only + // (BPF_NOEXIST), so the second pass can never clobber an ordinal the + // now-live fork tracepoint assigned. + t.seedExisting(snap) if err := t.stateMap.Put(uint32(0), bpf.TcBpfStepState{ - WorkerTgid: uint32(workerPID), + WorkerTgid: workerTgid, Enabled: 1, NextOrdinal: max(opts.OrdinalBase, 1), }); err != nil { t.Close() return nil, fmt.Errorf("failed to seed step state: %w", err) } - t.seedExisting() + if snap, err := t.snapshotFn(); err == nil { + t.seedExisting(snap) + } else { + logger.Warn("Second seeding pass skipped", "error", err) + } rd, err := ringbuf.NewReader(t.objs.MapStepEvents) if err != nil { @@ -213,9 +275,34 @@ func (t *Tracker) Close() { for _, l := range t.links { _ = l.Close() } + // The tcbpf connect/sendmsg hooks outlive this tracker (they are + // attached by cmd/start.go, and stay up when Start fails after the + // iterator has already seeded the table). With the tracepoints gone + // nothing maintains map_task_nspid, so a recycled tid would resolve to + // a dead task's process. Empty it: the hooks then take their documented + // global-tgid fallback, exactly as when attribution was never on. + clearMap(t.nspidMap) t.objs.Close() } +// clearMap deletes every entry of a u32→u32 hash map. +func clearMap(m *ebpf.Map) { + if m == nil { + return + } + var ( + keys []uint32 + k, v uint32 + ) + it := m.Iterate() + for it.Next(&k, &v) { + keys = append(keys, k) + } + for _, k := range keys { + _ = m.Delete(k) + } +} + func (t *Tracker) attach() error { forkLink, err := link.AttachTracing(link.TracingOptions{ Program: t.objs.StepFork, @@ -246,73 +333,67 @@ func (t *Tracker) attach() error { return fmt.Errorf("failed to attach cgroup sock_create: %w", err) } t.links = append(t.links, sockLink) + + iter, err := link.AttachIter(link.IterOptions{Program: t.objs.StepTaskIter}) + if err != nil { + return fmt.Errorf("failed to attach task iterator: %w", err) + } + t.links = append(t.links, iter) + t.iter = iter return nil } -// seedExisting tags processes that predate the daemon. Worker threads get -// StepOrdinalRunner so infra traffic (log upload, action downloads) is -// labeled as runner overhead; existing worker child subtrees — steps that -// ran or started before cargowall attached — get StepOrdinalPreDaemon, -// except the daemon's own subtree, which gets StepOrdinalRunner (cargowall -// is infrastructure, not a workflow step). Every write is create-only: -// once the fork tracepoint is live it is the sole authority on new tags, -// and the post-enable rescan must never replace a kernel-assigned ordinal. -// Best-effort by nature: /proc scans race process creation, which is why -// Start runs it twice around the enable flip. -func (t *Tracker) seedExisting() { - children := buildChildrenMap() - +// seedExisting tags processes that predate the daemon, from one task-walk +// snapshot. Worker threads get StepOrdinalRunner so infra traffic (log +// upload, action downloads) is labeled as runner overhead; existing worker +// child subtrees — steps that ran or started before cargowall attached — +// get StepOrdinalPreDaemon, except the daemon's own subtree, which gets +// StepOrdinalRunner (cargowall is infrastructure, not a workflow step). +// Every write is create-only: once the fork tracepoint is live it is the +// sole authority on new tags, and the post-enable rescan must never replace +// a kernel-assigned ordinal. Best-effort by nature: a walk races process +// creation, which is why Start runs it twice around the enable flip. +func (t *Tracker) seedExisting(snap *procSnapshot) { // Decide each pid's intended tag before writing anything, so a pid in // the daemon's subtree is written exactly once with the right value // (create-only writes mean there is no second chance). - self := make(map[int]bool) - for _, pid := range subtreePids(os.Getpid(), children) { - self[pid] = true + self := make(map[uint32]bool) + selfTgid, selfVisible := snap.global[os.Getpid()] + if selfVisible { + for _, pid := range snap.subtree(selfTgid) { + self[pid] = true + } } - t.tagTasks(t.workerPID, events.StepOrdinalRunner, ebpf.UpdateNoExist) - for _, root := range children[t.workerPID] { - for _, pid := range subtreePids(root, children) { + t.tagTasks(snap, t.workerTgid, events.StepOrdinalRunner, ebpf.UpdateNoExist) + for _, root := range snap.children[t.workerTgid] { + for _, pid := range snap.subtree(root) { ordinal := events.StepOrdinalPreDaemon if self[pid] { ordinal = events.StepOrdinalRunner } - t.tagTasks(pid, ordinal, ebpf.UpdateNoExist) + t.tagTasks(snap, pid, ordinal, ebpf.UpdateNoExist) } } // Standalone runs (daemon not under the worker): still label our own // traffic as infrastructure. Create-only, so a no-op when the loop // above already covered us. - for _, pid := range subtreePids(os.Getpid(), children) { - t.tagTasks(pid, events.StepOrdinalRunner, ebpf.UpdateNoExist) - } -} - -// subtreePids returns pid plus all its descendant process ids. -func subtreePids(pid int, children map[int][]int) []int { - pids := []int{pid} - for i := 0; i < len(pids); i++ { - pids = append(pids, children[pids[i]]...) + if selfVisible { + for _, pid := range snap.subtree(selfTgid) { + t.tagTasks(snap, pid, events.StepOrdinalRunner, ebpf.UpdateNoExist) + } } - return pids } -// tagTasks writes ordinal for every thread of pid. Seeding passes -// UpdateNoExist so a tid the fork tracepoint already tagged keeps its -// kernel-assigned ordinal; container leaders pass UpdateAny (see -// TagContainerProcess for why overwrite is safe there and only there). -func (t *Tracker) tagTasks(pid int, ordinal uint32, flags ebpf.MapUpdateFlags) { - tids, err := os.ReadDir("/proc/" + strconv.Itoa(pid) + "/task") - if err != nil { - return // process exited mid-scan - } - for _, tid := range tids { - n, err := strconv.ParseUint(tid.Name(), 10, 32) - if err != nil { - continue - } - _ = t.taskMap.Update(uint32(n), ordinal, flags) +// tagTasks writes ordinal for every thread of the (global) tgid as listed +// in snap. Seeding passes UpdateNoExist so a tid the fork tracepoint +// already tagged keeps its kernel-assigned ordinal; container leaders pass +// UpdateAny (see TagContainerProcess for why overwrite is safe there and +// only there). +func (t *Tracker) tagTasks(snap *procSnapshot, tgid, ordinal uint32, flags ebpf.MapUpdateFlags) { + for _, tid := range snap.tids[tgid] { + _ = t.taskMap.Update(tid, ordinal, flags) } } @@ -322,6 +403,12 @@ func (t *Tracker) tagTasks(pid int, ordinal uint32, flags ebpf.MapUpdateFlags) { // userspace bridge that puts the launching step's ordinal on them, after // which sock_create tags their sockets like any other tagged task. // +// pid is as dockerd reported it, i.e. numbered in the daemon's namespace +// when dockerd shares it (the caller's cgroup identity check has already +// confirmed that pid is the expected container in our /proc). A pid the +// task walk cannot see — exited, or a dockerd outside our namespace — is +// skipped. +// // The leader's threads are written with UpdateAny: a freshly shim-forked // leader is untagged in the normal case, an exec re-tag into a long-lived // container must replace the container's older ordinal, and overwrite @@ -329,10 +416,13 @@ func (t *Tracker) tagTasks(pid int, ordinal uint32, flags ebpf.MapUpdateFlags) { // create-only — children forked after an earlier tag already carry correct // kernel-inherited ordinals that must never be clobbered. func (t *Tracker) TagContainerProcess(pid int, ordinal uint32) { - t.tagTasks(pid, ordinal, ebpf.UpdateAny) - children := buildChildrenMap() - for _, p := range subtreePids(pid, children)[1:] { - t.tagTasks(p, ordinal, ebpf.UpdateNoExist) + snap, leader, ok := t.resolve(pid) + if !ok { + return + } + t.tagTasks(snap, leader, ordinal, ebpf.UpdateAny) + for _, p := range snap.subtree(leader)[1:] { + t.tagTasks(snap, p, ordinal, ebpf.UpdateNoExist) } } @@ -346,11 +436,115 @@ func (t *Tracker) TagContainerProcess(pid int, ordinal uint32) { // adoption is not an exec re-tag, and a recycled-tid stale entry is the // rarer wrong to optimize for than demoting a live container. func (t *Tracker) AdoptContainerProcess(pid int, ordinal uint32) { - for _, p := range subtreePids(pid, buildChildrenMap()) { - t.tagTasks(p, ordinal, ebpf.UpdateNoExist) + snap, leader, ok := t.resolve(pid) + if !ok { + return + } + for _, p := range snap.subtree(leader) { + t.tagTasks(snap, p, ordinal, ebpf.UpdateNoExist) } } +// resolve takes a fresh task walk and translates a daemon-namespace pid to +// the global tgid the step maps are keyed by. +func (t *Tracker) resolve(pid int) (*procSnapshot, uint32, bool) { + snap, err := t.snapshotFn() + if err != nil { + t.logger.Warn("Task walk failed", "error", err) + return nil, 0, false + } + leader, ok := snap.global[pid] + if !ok { + return nil, 0, false // exited mid-flight, or not numbered in our namespace + } + return snap, leader, true +} + +// initPidNamespaceInode is PROC_PID_INIT_INO, the fixed nsfs inode of the +// initial pid namespace — what a plain VM sees; anything else means the +// daemon is inside a container. +const initPidNamespaceInode uint32 = 0xEFFFFFFC + +// pidNamespaceInode identifies the daemon's pid namespace the way the +// kernel does (ns_common.inum), for the step_task_iter/step_fork walk. +func pidNamespaceInode() (uint32, error) { + fi, err := os.Stat("/proc/self/ns/pid") + if err != nil { + return 0, fmt.Errorf("failed to stat pid namespace: %w", err) + } + st, ok := fi.Sys().(*syscall.Stat_t) + if !ok { + return 0, errors.New("pid namespace stat has no inode") + } + return uint32(st.Ino), nil +} + +// procSnapshot is one pass of the step_task_iter walk: every task in the +// daemon's pid-namespace subtree, keyed by the global ids the step maps +// use, plus the translation from the daemon-namespace pids userspace +// discovers. +type procSnapshot struct { + tids map[uint32][]uint32 // global tgid → its threads' global tids + children map[uint32][]uint32 // global parent tgid → child global tgids + global map[int]uint32 // daemon-namespace tgid → global tgid +} + +// iterSnapshot runs the task iterator once and parses its output. +func (t *Tracker) iterSnapshot() (*procSnapshot, error) { + rd, err := t.iter.Open() + if err != nil { + return nil, fmt.Errorf("failed to open task iterator: %w", err) + } + defer rd.Close() + data, err := io.ReadAll(rd) + if err != nil { + return nil, fmt.Errorf("failed to read task iterator: %w", err) + } + return parseTaskRecords(data), nil +} + +// parseTaskRecords decodes the iterator's stream of task_iter_rec (the +// bpf2go-generated mirror of the C struct, so the layout has one source). +// Threads contribute their tid to the leader's list and nothing else: they +// share the leader's tree position and namespace pid. A trailing partial +// record (a truncated read) is ignored. +func parseTaskRecords(data []byte) *procSnapshot { + snap := &procSnapshot{ + tids: make(map[uint32][]uint32), + children: make(map[uint32][]uint32), + global: make(map[int]uint32), + } + const recSize = int(unsafe.Sizeof(bpf.StepBpfTaskIterRec{})) + for off := 0; off+recSize <= len(data); off += recSize { + rec := (*bpf.StepBpfTaskIterRec)(unsafe.Pointer(&data[off])) + snap.tids[rec.Tgid] = append(snap.tids[rec.Tgid], rec.Tid) + if rec.Tid != rec.Tgid { + continue + } + snap.global[int(rec.NsTgid)] = rec.Tgid + if rec.Ppid != rec.Tgid { + snap.children[rec.Ppid] = append(snap.children[rec.Ppid], rec.Tgid) + } + } + return snap +} + +// subtree returns root plus all its descendant tgids, breadth-first. The +// seen set guards against a cyclic snapshot; real trees terminate. +func (s *procSnapshot) subtree(root uint32) []uint32 { + pids := []uint32{root} + seen := map[uint32]bool{root: true} + for i := 0; i < len(pids); i++ { + for _, c := range s.children[pids[i]] { + if !seen[c] { + seen[c] = true + pids = append(pids, c) + } + } + } + return pids +} + // boundary records one step_child_event as observed by run(), stamped with // wall-clock receive time so external event streams that carry their own // timestamps (Docker's timeNano) can be resolved against the step that was @@ -414,6 +608,8 @@ func (t *Tracker) run() { // resolve against it promptly. t.recordBoundary(ev.Ordinal, time.Now()) + // ev.Tgid is numbered in our pid namespace (see map_task_nspid), so + // /proc reads work from inside a container too. cmdline := sanitizeCmdline(t.stepCmdline(int(ev.Tgid))) t.logger.Info("Workflow step process started", "step_ordinal", ev.Ordinal, @@ -571,27 +767,6 @@ func readComm(pid int) string { return strings.TrimSpace(string(comm)) } -// buildChildrenMap snapshots the process tree as parent pid → child pids. -func buildChildrenMap() map[int][]int { - children := make(map[int][]int) - procs, err := os.ReadDir("/proc") - if err != nil { - return children - } - for _, p := range procs { - pid, err := strconv.Atoi(p.Name()) - if err != nil { - continue - } - ppid, ok := readPPid(pid) - if !ok { - continue - } - children[ppid] = append(children[ppid], pid) - } - return children -} - // readPPid extracts the parent pid from /proc//stat. The comm field // (2) can contain spaces and parentheses, so parse from the last ')' — // state is the field after it, ppid the one after that. diff --git a/pkg/steps/steps_test.go b/pkg/steps/steps_test.go index 7cc60d7..a747548 100644 --- a/pkg/steps/steps_test.go +++ b/pkg/steps/steps_test.go @@ -17,18 +17,24 @@ package steps import ( + "bytes" + "encoding/binary" + "errors" "log/slog" "os" "os/exec" - "runtime" + "strconv" + "strings" "testing" "time" "github.com/cilium/ebpf" + "github.com/cilium/ebpf/link" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "golang.org/x/sys/unix" + "github.com/code-cargo/cargowall/bpf" "github.com/code-cargo/cargowall/pkg/events" ) @@ -67,15 +73,60 @@ func TestReadComm_Self(t *testing.T) { assert.LessOrEqual(t, len(comm), 15) } -func TestBuildChildrenMap_ContainsSelf(t *testing.T) { - children := buildChildrenMap() - assert.Contains(t, children[os.Getppid()], os.Getpid()) +// rec encodes one task_iter_rec as the step_task_iter iterator streams it. +func rec(tid, tgid, ppid, nsTgid uint32) []byte { + var buf bytes.Buffer + _ = binary.Write(&buf, binary.NativeEndian, bpf.StepBpfTaskIterRec{Tid: tid, Tgid: tgid, Ppid: ppid, NsTgid: nsTgid}) + return buf.Bytes() } -func TestSubtreePids_IncludesSelfAndDescendants(t *testing.T) { - children := map[int][]int{100: {200, 300}, 300: {400}} - assert.ElementsMatch(t, []int{100, 200, 300, 400}, subtreePids(100, children)) - assert.Equal(t, []int{400}, subtreePids(400, children)) +func concat(recs ...[]byte) []byte { + var out []byte + for _, r := range recs { + out = append(out, r...) + } + return out +} + +func TestParseTaskRecords_BuildsTreeAndTranslation(t *testing.T) { + // Global ids on the left, namespace pids on the right — deliberately + // different so a test that conflates them fails. + snap := parseTaskRecords(concat( + rec(1, 1, 0, 1), + rec(100, 100, 1, 7), + rec(101, 100, 1, 7), // thread of 100: contributes a tid, nothing else + rec(200, 200, 100, 8), + rec(300, 300, 200, 9), + rec(400, 400, 1, 10), + rec(5, 5, 5, 5), // self-parented: no child edge to itself + []byte{1, 2, 3}, // truncated trailing record: ignored + )) + + assert.Equal(t, []uint32{100, 101}, snap.tids[100]) + assert.Equal(t, uint32(100), snap.global[7]) + assert.Equal(t, uint32(200), snap.global[8]) + _, threadListed := snap.global[0] + assert.False(t, threadListed) + assert.Equal(t, []uint32{100, 400}, snap.children[1]) + assert.Equal(t, []uint32{200}, snap.children[100]) + assert.Empty(t, snap.children[5]) + assert.Equal(t, []uint32{100, 200, 300}, snap.subtree(100)) + assert.Equal(t, []uint32{400}, snap.subtree(400)) +} + +func TestProcSnapshot_SubtreeTerminatesOnCycle(t *testing.T) { + snap := &procSnapshot{children: map[uint32][]uint32{1: {2}, 2: {1}}} + assert.Equal(t, []uint32{1, 2}, snap.subtree(1)) +} + +func TestPidNamespaceInode_MatchesProcLink(t *testing.T) { + got, err := pidNamespaceInode() + require.NoError(t, err) + link, err := os.Readlink("/proc/self/ns/pid") + require.NoError(t, err) + want, err := strconv.ParseUint(strings.TrimSuffix(strings.TrimPrefix(link, "pid:["), "]"), 10, 32) + require.NoError(t, err) + assert.Equal(t, uint32(want), got) } func TestFindAncestorByComm_FindsParent(t *testing.T) { @@ -232,50 +283,217 @@ func newTaskMap(t *testing.T) *ebpf.Map { return m } -func TestTagContainerProcess_LeaderOverwriteDescendantCreateOnly(t *testing.T) { +// containerSnapshot is a synthetic task walk: leader global tgid 1000 (pid +// 50 in our namespace) with two threads, a pre-tagged child 2000 and an +// untagged child 3000, each with a grandchild. +func containerSnapshot() *procSnapshot { + return &procSnapshot{ + tids: map[uint32][]uint32{1000: {1000, 1001}, 2000: {2000}, 3000: {3000}, 3100: {3100}}, + children: map[uint32][]uint32{1000: {2000, 3000}, 3000: {3100}}, + global: map[int]uint32{50: 1000, 51: 2000, 52: 3000, 53: 3100}, + } +} + +func newSnapshotTracker(t *testing.T, snap *procSnapshot, err error) (*Tracker, *ebpf.Map) { + t.Helper() m := newTaskMap(t) - tr := &Tracker{taskMap: m} - - // Pin this goroutine to one OS thread so its tid is guaranteed to exist - // in /proc//task both when pre-seeding and when asserting — the Go - // runtime creates and parks threads at will, so any other tid could race. - runtime.LockOSThread() - defer runtime.UnlockOSThread() - selfTid := uint32(unix.Gettid()) - - // Two descendants of the leader (this test process): one already carries - // an ordinal — standing in for a child the fork tracepoint tagged — and - // one is untagged. - startSleeper := func() int { - cmd := exec.Command("sleep", "30") - require.NoError(t, cmd.Start()) - t.Cleanup(func() { - _ = cmd.Process.Kill() - _ = cmd.Wait() - }) - return cmd.Process.Pid + tr := &Tracker{ + taskMap: m, + logger: slog.Default(), + snapshotFn: func() (*procSnapshot, error) { return snap, err }, } - preTaggedPid := startSleeper() - untaggedPid := startSleeper() + return tr, m +} + +func TestTagContainerProcess_LeaderOverwriteDescendantCreateOnly(t *testing.T) { + tr, m := newSnapshotTracker(t, containerSnapshot(), nil) - // Pre-seed: our own tid simulates a stale leader entry (exec re-tag into - // a long-lived container / recycled tid), the first child a + // Pre-seed: a leader thread simulates a stale entry (exec re-tag into a + // long-lived container / recycled tid), the first child a // kernel-inherited descendant tag. - require.NoError(t, m.Put(selfTid, uint32(99))) - require.NoError(t, m.Put(uint32(preTaggedPid), uint32(7))) + require.NoError(t, m.Put(uint32(1001), uint32(99))) + require.NoError(t, m.Put(uint32(2000), uint32(7))) - tr.TagContainerProcess(os.Getpid(), 42) + // The caller speaks namespace pids; the map is keyed by global tids. + tr.TagContainerProcess(50, 42) - // Leader threads are written UpdateAny: the stale ordinal must be replaced. var got uint32 - require.NoError(t, m.Lookup(selfTid, &got)) - assert.Equal(t, uint32(42), got, "leader tid must be overwritten (UpdateAny)") - - // Descendants are create-only: a kernel-inherited ordinal survives... - require.NoError(t, m.Lookup(uint32(preTaggedPid), &got)) + require.NoError(t, m.Lookup(uint32(1000), &got)) + assert.Equal(t, uint32(42), got, "leader must be tagged under its global tid") + require.NoError(t, m.Lookup(uint32(1001), &got)) + assert.Equal(t, uint32(42), got, "leader thread must be overwritten (UpdateAny)") + require.NoError(t, m.Lookup(uint32(2000), &got)) assert.Equal(t, uint32(7), got, "pre-tagged descendant must keep its ordinal (UpdateNoExist)") - - // ...while an untagged descendant picks up the container's ordinal. - require.NoError(t, m.Lookup(uint32(untaggedPid), &got)) + require.NoError(t, m.Lookup(uint32(3000), &got)) assert.Equal(t, uint32(42), got, "untagged descendant must be tagged") + require.NoError(t, m.Lookup(uint32(3100), &got)) + assert.Equal(t, uint32(42), got, "grandchild must be tagged") + assert.Error(t, m.Lookup(uint32(50), &got), "the namespace pid itself must never be used as a key") +} + +func TestAdoptContainerProcess_CreateOnlyForLeaderToo(t *testing.T) { + tr, m := newSnapshotTracker(t, containerSnapshot(), nil) + require.NoError(t, m.Put(uint32(1000), uint32(7))) + + tr.AdoptContainerProcess(50, 42) + + var got uint32 + require.NoError(t, m.Lookup(uint32(1000), &got)) + assert.Equal(t, uint32(7), got, "adoption must not demote a live leader") + require.NoError(t, m.Lookup(uint32(1001), &got)) + assert.Equal(t, uint32(42), got) + require.NoError(t, m.Lookup(uint32(3000), &got)) + assert.Equal(t, uint32(42), got) +} + +func TestTagContainerProcess_UnknownPidAndWalkFailureAreNoOps(t *testing.T) { + tr, m := newSnapshotTracker(t, containerSnapshot(), nil) + tr.TagContainerProcess(999, 42) // exited, or numbered in another namespace + var got uint32 + assert.Error(t, m.Lookup(uint32(999), &got)) + + failing, m2 := newSnapshotTracker(t, nil, errors.New("iterator closed")) + failing.TagContainerProcess(50, 42) + assert.Error(t, m2.Lookup(uint32(1000), &got)) +} + +// loadTcObjects loads the tcbpf collection that owns the tracker's shared +// maps, exactly as cmd/start.go does before steps.Start. +func loadTcObjects(t *testing.T) *bpf.TcBpfObjects { + t.Helper() + spec, err := bpf.LoadTcBpf() + require.NoError(t, err) + var objs bpf.TcBpfObjects + if err := spec.LoadAndAssign(&objs, nil); err != nil { + t.Skipf("tcbpf objects not loadable (needs root): %v", err) + } + t.Cleanup(func() { objs.Close() }) + return &objs +} + +// connectedUDPCookie opens a UDP socket, connects it (which runs the +// cgroup/connect4 hook in this thread's context) and returns its socket +// cookie — the key the hook wrote map_sock_pid and map_sock_step under. +func connectedUDPCookie(t *testing.T) uint64 { + t.Helper() + fd, err := unix.Socket(unix.AF_INET, unix.SOCK_DGRAM|unix.SOCK_CLOEXEC, 0) + require.NoError(t, err) + t.Cleanup(func() { _ = unix.Close(fd) }) + require.NoError(t, unix.Connect(fd, &unix.SockaddrInet4{Port: 9, Addr: [4]byte{127, 0, 0, 1}})) + cookie, err := unix.GetsockoptUint64(fd, unix.SOL_SOCKET, unix.SO_COOKIE) + require.NoError(t, err) + return cookie +} + +// TestStart_TagsWorkerChild is the production path end to end: Start with +// this test process declared as Runner.Worker, then fork a child and check +// every id crossing — seeding keyed by global tids, the boundary event and +// map_task_nspid carrying our namespace's numbering, and the tcbpf connect +// hook recording a pid we can resolve. Run on its own it covers a plain +// VM; TestStart_InPidNamespace re-runs it inside a fresh pid namespace, +// where every one of those crossings used to be wrong. +func TestStart_TagsWorkerChild(t *testing.T) { + objs := loadTcObjects(t) + // The connect hook is attached by cmd/start.go, not by Start; mirror + // that so the socket-owner path is exercised too. + connectLink, err := link.AttachCgroup(link.CgroupOptions{ + Path: "/sys/fs/cgroup", Attach: ebpf.AttachCGroupInet4Connect, Program: objs.CgConnect4, + }) + require.NoError(t, err, "attach cgroup/connect4") + t.Cleanup(func() { _ = connectLink.Close() }) + + tr, err := Start(objs, Options{WorkerPID: os.Getpid(), OrdinalBase: 100}, nil, slog.Default()) + if err != nil && strings.Contains(err.Error(), "BTF") { + t.Skipf("step tracker needs kernel BTF: %v", err) + } + require.NoError(t, err) + closed := false + t.Cleanup(func() { + if !closed { + tr.Close() + } + }) + + snap, err := tr.snapshotFn() + require.NoError(t, err) + selfTgid, ok := snap.global[os.Getpid()] + require.True(t, ok, "the daemon must see itself in the task walk") + assert.Equal(t, selfTgid, tr.workerTgid) + + // Seeding: the worker's (our) threads are runner infrastructure, keyed + // by global tid, and the namespace table names us by our own pid. + var got uint32 + require.NoError(t, objs.MapTaskStep.Lookup(selfTgid, &got)) + assert.Equal(t, uint32(events.StepOrdinalRunner), got) + require.NoError(t, objs.MapTaskNspid.Lookup(selfTgid, &got)) + assert.Equal(t, uint32(os.Getpid()), got) + + child := exec.Command("sleep", "5") + require.NoError(t, child.Start()) + t.Cleanup(func() { _ = child.Process.Kill(); _ = child.Wait() }) + childPID := child.Process.Pid + + snap, err = tr.snapshotFn() + require.NoError(t, err) + childTgid, ok := snap.global[childPID] + require.True(t, ok, "child must be visible to the task walk") + var ordinal uint32 + require.NoError(t, objs.MapTaskStep.Lookup(childTgid, &ordinal), + "worker child must be tagged under its global tid") + // Ordinals are opaque group ids, not positions: transient runtime forks + // (Go's one-time pidfd probe) consume them too, so only the base is a + // floor — see Options.OrdinalBase. + assert.GreaterOrEqual(t, ordinal, uint32(100)) + require.NoError(t, objs.MapTaskNspid.Lookup(childTgid, &got)) + assert.Equal(t, uint32(childPID), got) + + // The reconciler saw the boundary event (its tgid is our numbering, so + // the cmdline read behind it works too — visible in -v output). + assert.Eventually(t, func() bool { return tr.OrdinalAt(time.Now()) >= ordinal }, + 2*time.Second, 20*time.Millisecond, "boundary must reach the reconciler") + + // Socket owner: the connect hook stores the pid as we number it (what + // lookupProcessName reads), and the step tag we were seeded with. + cookie := connectedUDPCookie(t) + require.NoError(t, objs.MapSockPid.Lookup(cookie, &got)) + assert.Equal(t, uint32(os.Getpid()), got, "map_sock_pid must carry our namespace's pid") + require.NoError(t, objs.MapSockStep.Lookup(cookie, &got)) + assert.Equal(t, uint32(events.StepOrdinalRunner), got) + + // After Close the connect hook is still attached but nothing maintains + // the translation table, so it must be empty: the hook then falls back + // to the global tgid rather than resolving a recycled tid to a dead task. + tr.Close() + closed = true + var k, v uint32 + assert.False(t, objs.MapTaskNspid.Iterate().Next(&k, &v), "map_task_nspid must be empty after Close") +} + +// TestStart_InPidNamespace re-executes TestStart_TagsWorkerChild inside a +// fresh pid (and mount, for /proc) namespace — the shape of an ARC runner +// pod, where /proc pids and kernel pids disagree. +func TestStart_InPidNamespace(t *testing.T) { + runInPidNamespace(t, "^TestStart_TagsWorkerChild$") +} + +// runInPidNamespace runs the named test of this binary under +// unshare --pid --fork --mount-proc and requires it to pass there. Needs +// root (like every BPF test here) and util-linux. +func runInPidNamespace(t *testing.T, run string) { + t.Helper() + if os.Geteuid() != 0 { + t.Skip("needs root") + } + unshare, err := exec.LookPath("unshare") + if err != nil { + t.Skip("unshare not available") + } + out, err := exec.Command(unshare, "--pid", "--fork", "--mount-proc", + os.Args[0], "-test.run", run, "-test.v").CombinedOutput() + t.Logf("inside pid namespace:\n%s", out) + if strings.Contains(string(out), "--- SKIP") { + t.Skip("inner test skipped") + } + require.NoError(t, err) + require.Contains(t, string(out), "--- PASS") }