Skip to content

Commit c369063

Browse files
arkamare2b-bot[bot]
authored andcommitted
refactor(orchestrator): inject the sandbox logger into the NBD device path
GitOrigin-RevId: 7ae8b67f0030f3e807dd74e6586205d22a96082d
1 parent 3ecebe7 commit c369063

16 files changed

Lines changed: 245 additions & 64 deletions

‎packages/orchestrator/pkg/sandbox/nbd/devicehelper.go‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block"
1111
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/nbd/testutils"
1212
"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
13+
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
1314
)
1415

1516
// GetNBDDevice provisions a one-shot device pool, opens a direct-path mount
@@ -50,7 +51,7 @@ func GetNBDDevice(ctx context.Context, backend block.Device, featureFlags *featu
5051
close(poolClosed)
5152
}()
5253

53-
mnt := NewDirectPathMount(backend, devicePool, featureFlags, mountOpts...)
54+
mnt := NewDirectPathMount(backend, devicePool, featureFlags, logger.L(), mountOpts...)
5455

5556
mntIndex, err := mnt.Open(ctx)
5657
if err != nil {

‎packages/orchestrator/pkg/sandbox/nbd/dispatch.go‎

Lines changed: 11 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -104,10 +104,9 @@ type Dispatch struct {
104104
responseHeader []byte
105105
writeLock sync.Mutex
106106
prov Provider
107-
// provName is the concrete backend type name, cached at construction so
108-
// error logs can identify which storage layer failed without reflection
109-
// on every call.
110-
provName string
107+
// logger is tagged with the backend type name at construction, so the
108+
// reflection stays out of the per-request path.
109+
logger logger.Logger
111110
pendingResponses sync.WaitGroup
112111
shuttingDown bool
113112
shuttingDownLock sync.Mutex
@@ -119,12 +118,12 @@ type Dispatch struct {
119118
asyncWriteZeroes bool
120119
}
121120

122-
func NewDispatch(fp io.ReadWriter, prov Provider, asyncWriteZeroes bool) *Dispatch {
121+
func NewDispatch(fp io.ReadWriter, prov Provider, asyncWriteZeroes bool, lg logger.Logger) *Dispatch {
123122
d := &Dispatch{
124123
responseHeader: make([]byte, 16),
125124
fp: fp,
126125
prov: prov,
127-
provName: fmt.Sprintf("%T", prov),
126+
logger: lg.With(zap.String("nbd_provider", fmt.Sprintf("%T", prov))),
128127
fatal: make(chan error, 1),
129128
asyncWriteZeroes: asyncWriteZeroes,
130129
}
@@ -338,10 +337,9 @@ func (d *Dispatch) cmdRead(ctx context.Context, cmdHandle uint64, cmdFrom uint64
338337
// Per-request backend failure: signal it to the NBD client via the
339338
// response error byte and keep the dispatch loop alive. Only
340339
// writeResponse errors (dead NBD socket) escalate through d.fatal.
341-
logger.L().Error(ctx, "nbd backend read failed",
340+
d.logger.Error(ctx, "nbd backend read failed",
342341
zap.Error(readErr),
343342
zap.String("nbd_op", "read"),
344-
zap.String("nbd_provider", d.provName),
345343
zap.Uint64("nbd_handle", handle),
346344
zap.Uint64("nbd_offset", from),
347345
zap.Uint32("nbd_length", length),
@@ -360,10 +358,9 @@ func (d *Dispatch) cmdRead(ctx context.Context, cmdHandle uint64, cmdFrom uint64
360358
select {
361359
case d.fatal <- err:
362360
default:
363-
logger.L().Error(ctx, "nbd error cmd read",
361+
d.logger.Error(ctx, "nbd error cmd read",
364362
zap.Error(err),
365363
zap.String("nbd_op", "read"),
366-
zap.String("nbd_provider", d.provName),
367364
zap.Uint64("nbd_handle", cmdHandle),
368365
zap.Uint64("nbd_offset", cmdFrom),
369366
zap.Uint32("nbd_length", cmdLength),
@@ -410,10 +407,9 @@ func (d *Dispatch) cmdWrite(ctx context.Context, cmdHandle uint64, cmdFrom uint6
410407
}
411408

412409
if writeErr != nil {
413-
logger.L().Error(ctx, "nbd backend write failed",
410+
d.logger.Error(ctx, "nbd backend write failed",
414411
zap.Error(writeErr),
415412
zap.String("nbd_op", "write"),
416-
zap.String("nbd_provider", d.provName),
417413
zap.Uint64("nbd_handle", handle),
418414
zap.Uint64("nbd_offset", from),
419415
zap.Int("nbd_length", len(data)),
@@ -432,10 +428,9 @@ func (d *Dispatch) cmdWrite(ctx context.Context, cmdHandle uint64, cmdFrom uint6
432428
select {
433429
case d.fatal <- err:
434430
default:
435-
logger.L().Error(ctx, "nbd error cmd write",
431+
d.logger.Error(ctx, "nbd error cmd write",
436432
zap.Error(err),
437433
zap.String("nbd_op", "write"),
438-
zap.String("nbd_provider", d.provName),
439434
zap.Uint64("nbd_handle", cmdHandle),
440435
zap.Uint64("nbd_offset", cmdFrom),
441436
zap.Int("nbd_length", len(cmdData)),
@@ -484,9 +479,8 @@ func (d *Dispatch) cmdWriteZeroes(ctx context.Context, cmdHandle uint64, cmdFrom
484479
var respErr uint32
485480
if zeroErr != nil {
486481
respErr = 1
487-
logger.L().Error(ctx, "nbd backend write-zeroes failed",
482+
d.logger.Error(ctx, "nbd backend write-zeroes failed",
488483
zap.Error(zeroErr),
489-
zap.String("nbd_provider", d.provName),
490484
zap.Uint64("nbd_handle", cmdHandle),
491485
zap.Uint64("nbd_offset", cmdFrom),
492486
zap.Int64("nbd_length", cmdLength),
@@ -515,10 +509,9 @@ func (d *Dispatch) cmdWriteZeroes(ctx context.Context, cmdHandle uint64, cmdFrom
515509
select {
516510
case d.fatal <- err:
517511
default:
518-
logger.L().Error(ctx, "nbd error cmd write-zeroes",
512+
d.logger.Error(ctx, "nbd error cmd write-zeroes",
519513
zap.Error(err),
520514
zap.String("nbd_op", "write-zeroes"),
521-
zap.String("nbd_provider", d.provName),
522515
zap.Uint64("nbd_handle", cmdHandle),
523516
zap.Uint64("nbd_offset", cmdFrom),
524517
zap.Int64("nbd_length", cmdLength),

‎packages/orchestrator/pkg/sandbox/nbd/dispatch_fault_test.go‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import (
1414
"github.com/stretchr/testify/require"
1515

1616
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block"
17+
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
1718
)
1819

1920
// replyConn is a fake NBD socket: Read serves queued request bytes, Write
@@ -93,7 +94,7 @@ func TestDispatch_MmapFault(t *testing.T) {
9394
require.NoError(t, os.Truncate(path, 0))
9495

9596
conn := &replyConn{reqCh: make(chan []byte, 8), replies: make(chan Response, 8)}
96-
d := NewDispatch(conn, &faultProv{cache: cache}, true)
97+
d := NewDispatch(conn, &faultProv{cache: cache}, true, logger.L())
9798

9899
done := make(chan error, 1)
99100
go func() { done <- d.Handle(t.Context()) }()
Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
//go:build linux
2+
3+
package nbd
4+
5+
import (
6+
"context"
7+
"errors"
8+
"io"
9+
"testing"
10+
"time"
11+
12+
"github.com/stretchr/testify/require"
13+
14+
"github.com/e2b-dev/infra/packages/shared/pkg/sandboxtypes"
15+
)
16+
17+
// errProv fails every backend call, so one request is enough to reach the
18+
// dispatcher's error log.
19+
type errProv struct{}
20+
21+
var errBackend = errors.New("backend unavailable")
22+
23+
func (errProv) ReadAt(context.Context, []byte, int64) (int, error) { return 0, errBackend }
24+
25+
func (errProv) Size(context.Context) (int64, error) { return 0, nil }
26+
27+
func (errProv) WriteAt([]byte, int64) (int, error) { return 0, errBackend }
28+
29+
func (errProv) WriteZeroesAt(int64, int64) (int, error) { return 0, errBackend }
30+
31+
// A backend failure is the line an operator chases, so it has to name the
32+
// sandbox it belongs to rather than an offset and a device index they then
33+
// have to join back to one.
34+
func TestDispatch_BackendFailureCarriesSandboxIdentity(t *testing.T) {
35+
t.Parallel()
36+
37+
runtime := sandboxtypes.RuntimeMetadata{
38+
SandboxID: "sbx-log-identity",
39+
TemplateID: "tmpl-log-identity",
40+
TeamID: "team-log-identity",
41+
BuildID: "build-log-identity",
42+
ExecutionID: "exec-log-identity",
43+
}
44+
45+
conn := &replyConn{reqCh: make(chan []byte, 1), replies: make(chan Response, 1)}
46+
d := NewDispatch(conn, errProv{}, true, runtime.Logger())
47+
48+
done := make(chan error, 1)
49+
go func() { done <- d.Handle(t.Context()) }()
50+
51+
conn.reqCh <- nbdRequest(NBDCmdRead, 1, 0, 4096)
52+
select {
53+
case resp := <-conn.replies:
54+
require.NotZero(t, resp.Error, "a failing backend read must reply with an NBD error")
55+
case <-time.After(5 * time.Second):
56+
t.Fatal("no reply for the failing read")
57+
}
58+
59+
close(conn.reqCh) // Handle sees io.EOF and exits
60+
select {
61+
case err := <-done:
62+
require.ErrorIs(t, err, io.EOF)
63+
case <-time.After(5 * time.Second):
64+
t.Fatal("dispatch loop did not exit")
65+
}
66+
67+
var fields map[string]any
68+
for _, entry := range testLogObserver.FilterMessage("nbd backend read failed").All() {
69+
if entry.ContextMap()["sandbox.id"] == runtime.SandboxID {
70+
fields = entry.ContextMap()
71+
72+
break
73+
}
74+
}
75+
76+
require.NotNil(t, fields, "the backend read failure must be logged for this sandbox")
77+
require.Equal(t, runtime.TemplateID, fields["template.id"])
78+
require.Equal(t, runtime.TeamID, fields["team.id"])
79+
require.Equal(t, runtime.BuildID, fields["build.id"])
80+
require.Equal(t, runtime.ExecutionID, fields["execution.id"])
81+
require.Equal(t, "nbd.errProv", fields["nbd_provider"])
82+
}

‎packages/orchestrator/pkg/sandbox/nbd/dispatch_writezeroes_test.go‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,8 @@ import (
99
"sync"
1010
"testing"
1111
"time"
12+
13+
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
1214
)
1315

1416
// ctrlConn is a fake NBD socket driving the real Dispatch.Handle loop.
@@ -136,7 +138,7 @@ func TestDispatchWriteZeroesReadLoopStall(t *testing.T) {
136138
firstWrite: make(chan struct{}),
137139
}
138140
prov := &stallProv{seen: map[int64]bool{}, wz: make(chan struct{})}
139-
d := NewDispatch(conn, prov, tc.asyncWriteZeroes)
141+
d := NewDispatch(conn, prov, tc.asyncWriteZeroes, logger.L())
140142

141143
done := make(chan struct{})
142144
go func() {

‎packages/orchestrator/pkg/sandbox/nbd/path_direct.go‎

Lines changed: 22 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,8 @@ type DirectPathMount struct {
8888
devicePool *DevicePool
8989
featureFlags *featureflags.Client
9090

91+
logger logger.Logger
92+
9193
Backend block.Device
9294
deviceIndex uint32
9395
blockSize uint64
@@ -127,12 +129,13 @@ func WithDeadconnTimeout(d time.Duration) MountOption {
127129
return func(m *DirectPathMount) { m.deadconnTimeout = d }
128130
}
129131

130-
func NewDirectPathMount(b block.Device, devicePool *DevicePool, featureFlags *featureflags.Client, opts ...MountOption) *DirectPathMount {
132+
func NewDirectPathMount(b block.Device, devicePool *DevicePool, featureFlags *featureflags.Client, lg logger.Logger, opts ...MountOption) *DirectPathMount {
131133
m := &DirectPathMount{
132134
Backend: b,
133135
blockSize: 4096,
134136
devicePool: devicePool,
135137
featureFlags: featureFlags,
138+
logger: lg,
136139
socksClient: make([]*os.File, 0),
137140
socksServer: make([]io.Closer, 0),
138141
deviceIndex: math.MaxUint32,
@@ -168,7 +171,7 @@ func (d *DirectPathMount) Open(ctx context.Context) (retDeviceIndex uint32, err
168171
span.SetStatus(codes.Error, err.Error())
169172
}
170173
span.End()
171-
logger.L().Debug(ctx, "opening direct path mount", zap.Uint32("device_index", d.deviceIndex), zap.Error(err))
174+
d.logger.Debug(ctx, "opening direct path mount", zap.Uint32("device_index", d.deviceIndex), zap.Error(err))
172175
}()
173176

174177
telemetry.ReportEvent(ctx, "opening direct path mount")
@@ -221,20 +224,20 @@ func (d *DirectPathMount) Open(ctx context.Context) (retDeviceIndex uint32, err
221224
}
222225
server.Close()
223226

224-
dispatch := NewDispatch(serverc, d.Backend, asyncWriteZeroes)
225-
// Capture deviceIndex for the goroutine closure — it's reassigned on
226-
// each retry iteration of the outer for-loop (not a range loop, so
227-
// Go 1.22+ loop variable fix doesn't apply).
228-
devIdx := deviceIndex
227+
// Capture deviceIndex here — it's reassigned on each retry
228+
// iteration of the outer for-loop (not a range loop, so Go 1.22+
229+
// loop variable fix doesn't apply).
230+
connLogger := d.logger.With(
231+
zap.Uint32("device_index", deviceIndex),
232+
zap.Int("socket_index", i),
233+
)
234+
235+
dispatch := NewDispatch(serverc, d.Backend, asyncWriteZeroes, connLogger)
229236
// Start reading commands on the socket and dispatching them to our provider
230237
d.handlersWg.Go(func() {
231238
handleErr := dispatch.Handle(handlerCtx)
232239
// The error is expected to happen if the nbd (socket connection) is closed
233-
logger.L().Info(handlerCtx, "closing handler for NBD commands",
234-
zap.Error(handleErr),
235-
zap.Uint32("device_index", devIdx),
236-
zap.Int("socket_index", i),
237-
)
240+
connLogger.Info(handlerCtx, "closing handler for NBD commands", zap.Error(handleErr))
238241
})
239242

240243
d.socksServer = append(d.socksServer, serverc)
@@ -264,19 +267,19 @@ func (d *DirectPathMount) Open(ctx context.Context) (retDeviceIndex uint32, err
264267
break
265268
}
266269

267-
logger.L().Error(ctx, "error opening NBD, retrying", zap.Error(connectErr), zap.Uint32("device_index", deviceIndex))
270+
d.logger.Error(ctx, "error opening NBD, retrying", zap.Error(connectErr), zap.Uint32("device_index", deviceIndex))
268271

269272
// Sometimes (rare), there seems to be a BADF error here. Lets just retry for now...
270273
// Close things down and try again...
271274
err := closeSocketPairs(d.socksClient, d.socksServer)
272275
if err != nil {
273-
logger.L().Error(ctx, "error closing socket pairs on error opening NBD", zap.Error(err))
276+
d.logger.Error(ctx, "error closing socket pairs on error opening NBD", zap.Error(err))
274277
}
275278

276279
// Release the device back to the pool
277280
err = d.devicePool.ReleaseDevice(ctx, deviceIndex)
278281
if err != nil {
279-
logger.L().Error(ctx, "error opening NBD, error releasing device", zap.Error(err), zap.Uint32("device_index", deviceIndex))
282+
d.logger.Error(ctx, "error opening NBD, error releasing device", zap.Error(err), zap.Uint32("device_index", deviceIndex))
280283
}
281284

282285
if strings.Contains(connectErr.Error(), "invalid argument") {
@@ -348,7 +351,7 @@ func (d *DirectPathMount) closeConnected(ctx context.Context, deviceIndex uint32
348351
// Warn, not Debug: this is rare, it means a device and a pool slot came within
349352
// one branch of being stranded, and the build failure the caller reports does
350353
// not name either of them.
351-
logger.L().Warn(ctx, "tearing down a connected NBD device Open cannot return",
354+
d.logger.Warn(ctx, "tearing down a connected NBD device Open cannot return",
352355
zap.Uint32("device_index", deviceIndex),
353356
zap.String("stage", stage),
354357
)
@@ -468,7 +471,7 @@ func (d *DirectPathMount) Close(ctx context.Context) error {
468471
stage.Store("sync")
469472
watchdog := time.AfterFunc(deviceCloseWarnThreshold, func() {
470473
nbdSlowCloseCounter.Add(ctx, 1, metric.WithAttributes(attribute.String("stage", stage.Load().(string))))
471-
logger.L().Warn(ctx, "NBD device descriptor close stalled",
474+
d.logger.Warn(ctx, "NBD device descriptor close stalled",
472475
zap.Duration("threshold", deviceCloseWarnThreshold),
473476
zap.Uint32("device_index", idx),
474477
)
@@ -497,7 +500,7 @@ func (d *DirectPathMount) Close(ctx context.Context) error {
497500
}
498501

499502
if !watchdog.Stop() {
500-
logger.L().Warn(ctx, "NBD device descriptor close finished after stalling",
503+
d.logger.Warn(ctx, "NBD device descriptor close finished after stalling",
501504
zap.Duration("duration", time.Since(closeStart)),
502505
zap.Uint32("device_index", idx),
503506
)
@@ -559,7 +562,7 @@ func (d *DirectPathMount) Close(ctx context.Context) error {
559562
// Release the device back to the pool, retry if it is in use
560563
if idx != math.MaxUint32 {
561564
telemetry.ReportEvent(ctx, "releasing device to the pool")
562-
err := d.devicePool.ReleaseDevice(ctx, idx, WithInfiniteRetry())
565+
err := d.devicePool.ReleaseDevice(ctx, idx, WithInfiniteRetry(), WithLogger(d.logger))
563566
if err != nil {
564567
errs = append(errs, fmt.Errorf("error releasing overlay device: %w", err))
565568
}

‎packages/orchestrator/pkg/sandbox/nbd/path_direct_cancel_test.go‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import (
1818
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block"
1919
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/nbd/testutils"
2020
"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
21+
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
2122
"github.com/e2b-dev/infra/packages/shared/pkg/storage/header"
2223
)
2324

@@ -107,7 +108,7 @@ func TestPathDirect_OpenCancelledAfterConnect(t *testing.T) {
107108
// from here on nothing else may take the slot whose release is under test, and nothing
108109
// can, because the device is connected and so reads as in-use.
109110
connected := DeviceSlot(math.MaxUint32)
110-
mnt := NewDirectPathMount(overlay, pool, featureFlags,
111+
mnt := NewDirectPathMount(overlay, pool, featureFlags, logger.L(),
111112
withAfterConnect(func(deviceIndex uint32) {
112113
connected = deviceIndex
113114
stopFeeder()

‎packages/orchestrator/pkg/sandbox/nbd/path_direct_close_events_test.go‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import (
1616
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block"
1717
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/nbd/testutils"
1818
"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
19+
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
1920
"github.com/e2b-dev/infra/packages/shared/pkg/storage/header"
2021
)
2122

@@ -62,7 +63,7 @@ func TestCloseClosesTheDescriptorBeforeTheHandlerTeardown(t *testing.T) {
6263

6364
pool := newPartitionedPool(t)
6465

65-
mnt := NewDirectPathMount(overlay, pool, featureFlags)
66+
mnt := NewDirectPathMount(overlay, pool, featureFlags, logger.L())
6667

6768
_, err = mnt.Open(ctx)
6869
require.NoError(t, err)

0 commit comments

Comments
 (0)