mirror of
https://github.com/netbirdio/gvisor.git
synced 2026-05-22 17:12:49 -07:00
Add Event controls
Add Event controls and implement "stream" commands. PiperOrigin-RevId: 390691702
This commit is contained in:
@@ -7,6 +7,7 @@ go_library(
|
||||
srcs = [
|
||||
"event.go",
|
||||
"event_any.go",
|
||||
"processor.go",
|
||||
"rate.go",
|
||||
],
|
||||
visibility = ["//:sandbox"],
|
||||
|
||||
@@ -26,3 +26,8 @@ import (
|
||||
func newAny(m proto.Message) (*anypb.Any, error) {
|
||||
return anypb.New(m)
|
||||
}
|
||||
|
||||
func emptyAny() *anypb.Any {
|
||||
var any anypb.Any
|
||||
return &any
|
||||
}
|
||||
|
||||
@@ -0,0 +1,130 @@
|
||||
// Copyright 2021 The gVisor Authors.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package eventchannel
|
||||
|
||||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"google.golang.org/protobuf/proto"
|
||||
pb "gvisor.dev/gvisor/pkg/eventchannel/eventchannel_go_proto"
|
||||
)
|
||||
|
||||
// eventProcessor carries display state across multiple events.
|
||||
type eventProcessor struct {
|
||||
filtering bool
|
||||
// filtered is the number of events omitted since printing the last matching
|
||||
// event. Only meaningful when filtering == true.
|
||||
filtered uint64
|
||||
// allowlist is the set of event names to display. If empty, all events are
|
||||
// displayed.
|
||||
allowlist map[string]bool
|
||||
}
|
||||
|
||||
// newEventProcessor creates a new EventProcessor with filters.
|
||||
func newEventProcessor(filters []string) *eventProcessor {
|
||||
e := &eventProcessor{
|
||||
filtering: len(filters) > 0,
|
||||
allowlist: make(map[string]bool),
|
||||
}
|
||||
for _, f := range filters {
|
||||
e.allowlist[f] = true
|
||||
}
|
||||
return e
|
||||
}
|
||||
|
||||
// processOne reads, parses and displays a single event from the event channel.
|
||||
//
|
||||
// The event channel is a stream of (msglen, payload) packets; this function
|
||||
// processes a single such packet. The msglen is a uvarint-encoded length for
|
||||
// the associated payload. The payload is a binary-encoded 'Any' protobuf, which
|
||||
// in turn encodes an arbitrary event protobuf.
|
||||
func (e *eventProcessor) processOne(src io.Reader, out *os.File) error {
|
||||
// Read and parse the msglen.
|
||||
lenbuf := make([]byte, binary.MaxVarintLen64)
|
||||
if _, err := io.ReadFull(src, lenbuf); err != nil {
|
||||
return err
|
||||
}
|
||||
msglen, consumed := binary.Uvarint(lenbuf)
|
||||
if consumed <= 0 {
|
||||
return fmt.Errorf("couldn't parse the message length")
|
||||
}
|
||||
|
||||
// Read the payload.
|
||||
buf := make([]byte, msglen)
|
||||
// Copy any unused bytes from the len buffer into the payload buffer. These
|
||||
// bytes are actually part of the payload.
|
||||
extraBytes := copy(buf, lenbuf[consumed:])
|
||||
if _, err := io.ReadFull(src, buf[extraBytes:]); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Unmarshal the payload into an "Any" protobuf, which encodes the actual
|
||||
// event.
|
||||
encodedEv := emptyAny()
|
||||
if err := proto.Unmarshal(buf, encodedEv); err != nil {
|
||||
return fmt.Errorf("failed to unmarshal 'any' protobuf message: %v", err)
|
||||
}
|
||||
|
||||
var ev pb.DebugEvent
|
||||
if err := (encodedEv).UnmarshalTo(&ev); err != nil {
|
||||
return fmt.Errorf("failed to decode 'any' protobuf message: %v", err)
|
||||
}
|
||||
|
||||
if e.filtering && e.allowlist[ev.Name] {
|
||||
e.filtered++
|
||||
return nil
|
||||
}
|
||||
|
||||
if e.filtering && e.filtered > 0 {
|
||||
if e.filtered == 1 {
|
||||
fmt.Fprintf(out, "... filtered %d event ...\n\n", e.filtered)
|
||||
} else {
|
||||
fmt.Fprintf(out, "... filtered %d events ...\n\n", e.filtered)
|
||||
}
|
||||
e.filtered = 0
|
||||
}
|
||||
|
||||
// Extract the inner event and display it. Example:
|
||||
//
|
||||
// 2017-10-04 14:35:05.316180374 -0700 PDT m=+1.132485846
|
||||
// cloud_gvisor.MemoryUsage {
|
||||
// total: 23822336
|
||||
// }
|
||||
fmt.Fprintf(out, "%v\n%v {\n", time.Now(), ev.Name)
|
||||
fmt.Fprintf(out, "%v", ev.Text)
|
||||
fmt.Fprintf(out, "}\n\n")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// ProcessAll reads, parses and displays all events from src. The events are
|
||||
// displayed to out.
|
||||
func ProcessAll(src io.Reader, filters []string, out *os.File) error {
|
||||
ep := newEventProcessor(filters)
|
||||
for {
|
||||
switch err := ep.processOne(src, out); err {
|
||||
case nil:
|
||||
continue
|
||||
case io.EOF:
|
||||
return nil
|
||||
default:
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ go_library(
|
||||
name = "control",
|
||||
srcs = [
|
||||
"control.go",
|
||||
"events.go",
|
||||
"fs.go",
|
||||
"lifecycle.go",
|
||||
"logging.go",
|
||||
@@ -20,6 +21,7 @@ go_library(
|
||||
deps = [
|
||||
"//pkg/abi/linux",
|
||||
"//pkg/context",
|
||||
"//pkg/eventchannel",
|
||||
"//pkg/fd",
|
||||
"//pkg/log",
|
||||
"//pkg/sentry/fdimport",
|
||||
|
||||
@@ -0,0 +1,65 @@
|
||||
// Copyright 2021 The gVisor Authors.
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package control
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
|
||||
"gvisor.dev/gvisor/pkg/eventchannel"
|
||||
"gvisor.dev/gvisor/pkg/urpc"
|
||||
)
|
||||
|
||||
// EventsOpts are the arguments for eventchannel-related commands.
|
||||
type EventsOpts struct {
|
||||
urpc.FilePayload
|
||||
}
|
||||
|
||||
// Events is the control server state for eventchannel-related commands.
|
||||
type Events struct {
|
||||
emitter eventchannel.Emitter
|
||||
}
|
||||
|
||||
// AttachDebugEmitter receives a connected unix domain socket FD from the client
|
||||
// and establishes it as a new emitter for the sentry eventchannel. Any existing
|
||||
// emitters are replaced on a subsequent attach.
|
||||
func (e *Events) AttachDebugEmitter(o *EventsOpts, _ *struct{}) error {
|
||||
if len(o.FilePayload.Files) < 1 {
|
||||
return errors.New("no output writer provided")
|
||||
}
|
||||
|
||||
sock, err := o.ReleaseFD(0)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
sockFD := sock.Release()
|
||||
|
||||
// SocketEmitter takes ownership of sockFD.
|
||||
emitter, err := eventchannel.SocketEmitter(sockFD)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create SocketEmitter for FD %d: %v", sockFD, err)
|
||||
}
|
||||
|
||||
// If there is already a debug emitter, close the old one.
|
||||
if e.emitter != nil {
|
||||
e.emitter.Close()
|
||||
}
|
||||
|
||||
e.emitter = eventchannel.DebugEmitterFrom(emitter)
|
||||
|
||||
// Register the new stream destination.
|
||||
eventchannel.AddEmitter(e.emitter)
|
||||
return nil
|
||||
}
|
||||
@@ -121,6 +121,11 @@ const (
|
||||
UsageReduce = "Usage.Reduce"
|
||||
)
|
||||
|
||||
// Events related commands (see events.go for more details).
|
||||
const (
|
||||
EventsAttachDebugEmitter = "Events.AttachDebugEmitter"
|
||||
)
|
||||
|
||||
// ControlSocketAddr generates an abstract unix socket name for the given ID.
|
||||
func ControlSocketAddr(id string) string {
|
||||
return fmt.Sprintf("\x00runsc-sandbox.%s", id)
|
||||
@@ -161,6 +166,7 @@ func newController(fd int, l *Loader) (*controller, error) {
|
||||
}
|
||||
|
||||
ctrl.srv.Register(&debug{})
|
||||
ctrl.srv.Register(&control.Events{})
|
||||
ctrl.srv.Register(&control.Logging{})
|
||||
ctrl.srv.Register(&control.Lifecycle{l.k})
|
||||
ctrl.srv.Register(&control.Fs{l.k})
|
||||
|
||||
@@ -35,9 +35,14 @@ func enableStrace(conf *config.Config) error {
|
||||
}
|
||||
strace.LogMaximumSize = max
|
||||
|
||||
sink := strace.SinkTypeLog
|
||||
if conf.StraceEvent {
|
||||
sink = strace.SinkTypeEvent
|
||||
}
|
||||
|
||||
if len(conf.StraceSyscalls) == 0 {
|
||||
strace.EnableAll(strace.SinkTypeLog)
|
||||
strace.EnableAll(sink)
|
||||
return nil
|
||||
}
|
||||
return strace.Enable(strings.Split(conf.StraceSyscalls, ","), strace.SinkTypeLog)
|
||||
return strace.Enable(strings.Split(conf.StraceSyscalls, ","), sink)
|
||||
}
|
||||
|
||||
@@ -33,6 +33,10 @@ type Events struct {
|
||||
intervalSec int
|
||||
// If true, events will print a single group of stats and exit.
|
||||
stats bool
|
||||
// If true, events will dump all filtered events to stdout.
|
||||
stream bool
|
||||
// filters for streamed events.
|
||||
filters stringSlice
|
||||
}
|
||||
|
||||
// Name implements subcommands.Command.Name.
|
||||
@@ -62,6 +66,8 @@ OPTIONS:
|
||||
func (evs *Events) SetFlags(f *flag.FlagSet) {
|
||||
f.IntVar(&evs.intervalSec, "interval", 5, "set the stats collection interval, in seconds")
|
||||
f.BoolVar(&evs.stats, "stats", false, "display the container's stats then exit")
|
||||
f.BoolVar(&evs.stream, "stream", false, "dump all filtered events to stdout")
|
||||
f.Var(&evs.filters, "filters", "only display matching events")
|
||||
}
|
||||
|
||||
// Execute implements subcommands.Command.Execute.
|
||||
@@ -79,6 +85,13 @@ func (evs *Events) Execute(ctx context.Context, f *flag.FlagSet, args ...interfa
|
||||
Fatalf("loading sandbox: %v", err)
|
||||
}
|
||||
|
||||
if evs.stream {
|
||||
if err := c.Stream(evs.filters, os.Stdout); err != nil {
|
||||
Fatalf("Stream failed: %v", err)
|
||||
}
|
||||
return subcommands.ExitSuccess
|
||||
}
|
||||
|
||||
// Repeatedly get stats from the container.
|
||||
for {
|
||||
// Get the event and print it as JSON.
|
||||
|
||||
@@ -117,6 +117,10 @@ type Config struct {
|
||||
// StraceLogSize is the max size of data blobs to display.
|
||||
StraceLogSize uint `flag:"strace-log-size"`
|
||||
|
||||
// StraceEvent indicates sending strace to events if true. Strace is
|
||||
// sent to log if false.
|
||||
StraceEvent bool `flag:"strace-event"`
|
||||
|
||||
// DisableSeccomp indicates whether seccomp syscall filters should be
|
||||
// disabled. Pardon the double negation, but default to enabled is important.
|
||||
DisableSeccomp bool
|
||||
|
||||
@@ -56,6 +56,7 @@ func RegisterFlags() {
|
||||
flag.Bool("strace", false, "enable strace.")
|
||||
flag.String("strace-syscalls", "", "comma-separated list of syscalls to trace. If --strace is true and this list is empty, then all syscalls will be traced.")
|
||||
flag.Uint("strace-log-size", 1024, "default size (in bytes) to log data argument blobs.")
|
||||
flag.Bool("strace-event", false, "send strace to event.")
|
||||
|
||||
// Flags that control sandbox runtime behavior.
|
||||
flag.String("platform", "ptrace", "specifies which platform to use: ptrace (default), kvm.")
|
||||
|
||||
@@ -670,6 +670,12 @@ func (c *Container) Reduce(wait bool) error {
|
||||
return c.Sandbox.Reduce(c.ID, wait)
|
||||
}
|
||||
|
||||
// Stream dumps all events to out.
|
||||
func (c *Container) Stream(filters []string, out *os.File) error {
|
||||
log.Debugf("Stream in container, cid: %s", c.ID)
|
||||
return c.Sandbox.Stream(c.ID, filters, out)
|
||||
}
|
||||
|
||||
// State returns the metadata of the container.
|
||||
func (c *Container) State() specs.State {
|
||||
return specs.State{
|
||||
|
||||
@@ -2779,3 +2779,53 @@ func TestReduce(t *testing.T) {
|
||||
t.Fatalf("error reduce from container: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestStream checks that Stream dumps expected events.
|
||||
func TestStream(t *testing.T) {
|
||||
spec, conf := sleepSpecConf(t)
|
||||
conf.Strace = true
|
||||
conf.StraceEvent = true
|
||||
conf.StraceSyscalls = ""
|
||||
|
||||
_, bundleDir, cleanup, err := testutil.SetupContainer(spec, conf)
|
||||
if err != nil {
|
||||
t.Fatalf("error setting up container: %v", err)
|
||||
}
|
||||
defer cleanup()
|
||||
|
||||
args := Args{
|
||||
ID: testutil.RandomContainerID(),
|
||||
Spec: spec,
|
||||
BundleDir: bundleDir,
|
||||
}
|
||||
|
||||
cont, err := New(conf, args)
|
||||
if err != nil {
|
||||
t.Fatalf("Creating container: %v", err)
|
||||
}
|
||||
defer cont.Destroy()
|
||||
|
||||
if err := cont.Start(conf); err != nil {
|
||||
t.Fatalf("starting container: %v", err)
|
||||
}
|
||||
|
||||
r, w, err := os.Pipe()
|
||||
if err != nil {
|
||||
t.Fatalf("os.Create(): %v", err)
|
||||
}
|
||||
|
||||
// Spawn a new thread to Stream events as it blocks indefinitely.
|
||||
go func() {
|
||||
cont.Stream(nil, w)
|
||||
}()
|
||||
|
||||
buf := make([]byte, 1024)
|
||||
if _, err := r.Read(buf); err != nil {
|
||||
t.Fatalf("Read out: %v", err)
|
||||
}
|
||||
|
||||
// A syscall strace event includes "Strace".
|
||||
if got, want := string(buf), "Strace"; !strings.Contains(got, want) {
|
||||
t.Errorf("out got %s, want include %s", buf, want)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,12 +17,14 @@ go_library(
|
||||
"//pkg/control/client",
|
||||
"//pkg/control/server",
|
||||
"//pkg/coverage",
|
||||
"//pkg/eventchannel",
|
||||
"//pkg/log",
|
||||
"//pkg/sentry/control",
|
||||
"//pkg/sentry/platform",
|
||||
"//pkg/sync",
|
||||
"//pkg/tcpip/header",
|
||||
"//pkg/tcpip/stack",
|
||||
"//pkg/unet",
|
||||
"//pkg/urpc",
|
||||
"//runsc/boot",
|
||||
"//runsc/boot/platforms",
|
||||
|
||||
@@ -35,10 +35,12 @@ import (
|
||||
"gvisor.dev/gvisor/pkg/control/client"
|
||||
"gvisor.dev/gvisor/pkg/control/server"
|
||||
"gvisor.dev/gvisor/pkg/coverage"
|
||||
"gvisor.dev/gvisor/pkg/eventchannel"
|
||||
"gvisor.dev/gvisor/pkg/log"
|
||||
"gvisor.dev/gvisor/pkg/sentry/control"
|
||||
"gvisor.dev/gvisor/pkg/sentry/platform"
|
||||
"gvisor.dev/gvisor/pkg/sync"
|
||||
"gvisor.dev/gvisor/pkg/unet"
|
||||
"gvisor.dev/gvisor/pkg/urpc"
|
||||
"gvisor.dev/gvisor/runsc/boot"
|
||||
"gvisor.dev/gvisor/runsc/boot/platforms"
|
||||
@@ -1073,6 +1075,37 @@ func (s *Sandbox) Reduce(cid string, wait bool) error {
|
||||
}, nil)
|
||||
}
|
||||
|
||||
// Stream sends the AttachDebugEmitter call for a container in the sandbox, and
|
||||
// dumps filtered events to out.
|
||||
func (s *Sandbox) Stream(cid string, filters []string, out *os.File) error {
|
||||
log.Debugf("Stream sandbox %q", s.ID)
|
||||
conn, err := s.sandboxConnect()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
r, w, err := unet.SocketPair(false)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
wfd, err := w.Release()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to release write socket FD: %v", err)
|
||||
}
|
||||
|
||||
if err := conn.Call(boot.EventsAttachDebugEmitter, &control.EventsOpts{
|
||||
FilePayload: urpc.FilePayload{Files: []*os.File{
|
||||
os.NewFile(uintptr(wfd), "event sink"),
|
||||
}},
|
||||
}, nil); err != nil {
|
||||
return fmt.Errorf("AttachDebugEmitter failed: %v", err)
|
||||
}
|
||||
|
||||
return eventchannel.ProcessAll(r, filters, out)
|
||||
}
|
||||
|
||||
// IsRunning returns true if the sandbox or gofer process is running.
|
||||
func (s *Sandbox) IsRunning() bool {
|
||||
if s.Pid != 0 {
|
||||
|
||||
Reference in New Issue
Block a user