Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions packages/clickhouse/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ require (
github.com/ClickHouse/clickhouse-go/v2 v2.40.1
github.com/e2b-dev/infra/packages/shared v0.0.0
github.com/google/uuid v1.6.0
github.com/stretchr/testify v1.11.1
go.opentelemetry.io/otel v1.43.0
go.opentelemetry.io/otel/metric v1.43.0
go.opentelemetry.io/otel/trace v1.43.0
Expand All @@ -24,6 +25,7 @@ require (
github.com/cenkalti/backoff/v5 v5.0.3 // indirect
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/coder/websocket v1.8.13 // indirect
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect
github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/elastic/go-sysinfo v1.15.4 // indirect
Expand Down Expand Up @@ -65,6 +67,7 @@ require (
github.com/patrickmn/go-cache v2.1.0+incompatible // indirect
github.com/paulmach/orb v0.11.1 // indirect
github.com/pierrec/lz4/v4 v4.1.22 // indirect
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
github.com/pressly/goose/v3 v3.26.0 // indirect
github.com/prometheus/procfs v0.17.0 // indirect
github.com/redis/go-redis/v9 v9.17.3 // indirect
Expand Down
18 changes: 17 additions & 1 deletion packages/clickhouse/pkg/clickhouse.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package clickhouse

import (
"context"
"errors"
"fmt"
"time"

Expand Down Expand Up @@ -38,7 +39,6 @@ func NewDriver(connectionString string) (driver.Conn, error) {

options.MaxOpenConns = 10
options.MaxIdleConns = 3
options.TLS = nil

conn, err := clickhouse.Open(options)
if err != nil {
Expand All @@ -57,6 +57,22 @@ func New(connectionString string) (*Client, error) {
return &Client{conn: conn}, nil
}

// EndpointFromDSN returns the credential-stripped host:port for use in logs
// and metric attributes. Never use the raw DSN there — it contains the password.
// On parse failure returns a fixed sentinel so the DSN (which clickhouse-go's
// url.Error embeds verbatim) never reaches a log line.
func EndpointFromDSN(dsn string) (string, error) {
options, err := clickhouse.ParseDSN(dsn)
if err != nil {
return "", errors.New("parse DSN")
}
if len(options.Addr) == 0 {
return "", errors.New("DSN has no addresses")
}

return options.Addr[0], nil
Comment thread
rguliyev marked this conversation as resolved.
}

// Close drains the queue and flushes remaining items
func (c *Client) Close(context.Context) error {
return c.conn.Close()
Expand Down
10 changes: 7 additions & 3 deletions packages/clickhouse/pkg/events/delivery.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,9 @@ type ClickhouseDelivery struct {

var tracer = otel.Tracer("github.com/e2b-dev/infra/packages/clickhouse/pkg/events")

func NewDefaultClickhouseSandboxEventsDelivery(ctx context.Context, conn driver.Conn, featureFlags *featureflags.Client) (*ClickhouseDelivery, error) {
const DefaultBatcherName = "sandbox-events"

func NewDefaultClickhouseSandboxEventsDelivery(ctx context.Context, conn driver.Conn, featureFlags *featureflags.Client, batcherName string) (*ClickhouseDelivery, error) {
maxBatchSize := featureFlags.IntFlag(ctx, featureflags.ClickhouseBatcherMaxBatchSize)

maxDelay := time.Duration(featureFlags.IntFlag(ctx, featureflags.ClickhouseBatcherMaxDelay)) * time.Millisecond
Expand All @@ -62,7 +64,7 @@ func NewDefaultClickhouseSandboxEventsDelivery(ctx context.Context, conn driver.

return NewClickhouseSandboxEventsDelivery(
ctx, conn, batcher.BatcherOptions{
Name: "sandbox-events",
Name: batcherName,
MaxBatchSize: maxBatchSize,
MaxDelay: maxDelay,
QueueSize: batcherQueueSize,
Expand Down Expand Up @@ -112,7 +114,8 @@ func (c *ClickhouseDelivery) Publish(_ context.Context, _ string, event events.S
})
}

func (c *ClickhouseDelivery) Close(context.Context) error {
// Close waits for queued items to flush.
func (c *ClickhouseDelivery) Close(_ context.Context) error {
return c.batcher.Stop()
}

Expand All @@ -128,6 +131,7 @@ func (c *ClickhouseDelivery) batchInserter(ctx context.Context, events []Sandbox

return fmt.Errorf("error preparing batch: %w", err)
}
defer batch.Close()

for _, event := range events {
Comment thread
rguliyev marked this conversation as resolved.
err := batch.Append(
Expand Down
11 changes: 9 additions & 2 deletions packages/clickhouse/pkg/hoststats/delivery.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,18 +47,21 @@ type ClickhouseDelivery struct {

var tracer = otel.Tracer("github.com/e2b-dev/infra/packages/clickhouse/pkg/hoststats")

const DefaultBatcherName = "sandbox-host-stats"

func NewDefaultClickhouseHostStatsDelivery(
ctx context.Context,
conn driver.Conn,
featureFlags *featureflags.Client,
batcherName string,
) (*ClickhouseDelivery, error) {
maxBatchSize := featureFlags.IntFlag(ctx, featureflags.ClickhouseBatcherMaxBatchSize)
maxDelay := time.Duration(featureFlags.IntFlag(ctx, featureflags.ClickhouseBatcherMaxDelay)) * time.Millisecond
batcherQueueSize := featureFlags.IntFlag(ctx, featureflags.ClickhouseBatcherQueueSize)

return NewClickhouseHostStatsDelivery(
ctx, conn, batcher.BatcherOptions{
Name: "sandbox-host-stats",
Name: batcherName,
MaxBatchSize: maxBatchSize,
MaxDelay: maxDelay,
QueueSize: batcherQueueSize,
Expand Down Expand Up @@ -93,7 +96,10 @@ func (c *ClickhouseDelivery) Push(stat SandboxHostStat) error {
return c.batcher.Push(stat)
}

func (c *ClickhouseDelivery) Close(context.Context) error {
// Close waits for queued items to flush. ctx is ignored: honoring it would
// leak the flush goroutine onto a connection the caller is about to tear
// down. Real deadlines require plumbing ctx through batcher.Stop.
func (c *ClickhouseDelivery) Close(_ context.Context) error {
return c.batcher.Stop()
}

Expand All @@ -109,6 +115,7 @@ func (c *ClickhouseDelivery) batchInserter(ctx context.Context, stats []SandboxH

return fmt.Errorf("error preparing batch: %w", err)
}
defer batch.Close()

for _, stat := range stats {
Comment thread
rguliyev marked this conversation as resolved.
err := batch.Append(
Expand Down
55 changes: 55 additions & 0 deletions packages/clickhouse/pkg/hoststats/hoststats.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ package hoststats

import (
"context"
"errors"
"sync"
"time"

"github.com/google/uuid"
Expand Down Expand Up @@ -56,3 +58,56 @@ func NewNoopDelivery() Delivery {

func (d *noopDelivery) Push(_ SandboxHostStat) error { return nil }
func (d *noopDelivery) Close(_ context.Context) error { return nil }

// multiDelivery fans out to every target. Push is serial (each target's Push
// is a non-blocking batcher send); Close is parallel so a stalled target
// can't block the others from draining.
type multiDelivery struct {
targets []Delivery
}

var _ Delivery = (*multiDelivery)(nil)

// NewMultiDelivery returns noop for 0 targets and the target directly for 1,
// so callers can wrap unconditionally.
func NewMultiDelivery(targets ...Delivery) Delivery {
switch len(targets) {
case 0:
return NewNoopDelivery()
case 1:
return targets[0]
default:
return &multiDelivery{targets: targets}
}
}

func (m *multiDelivery) Push(stat SandboxHostStat) error {
var err error
for _, t := range m.targets {
if e := t.Push(stat); e != nil {
err = errors.Join(err, e)
}
}

return err
Comment thread
rguliyev marked this conversation as resolved.
}

func (m *multiDelivery) Close(ctx context.Context) error {
var (
wg sync.WaitGroup
mu sync.Mutex
errs error
)
for _, t := range m.targets {
wg.Go(func() {
if e := t.Close(ctx); e != nil {
mu.Lock()
errs = errors.Join(errs, e)
mu.Unlock()
}
})
}
wg.Wait()

return errs
}
Comment thread
rguliyev marked this conversation as resolved.
136 changes: 136 additions & 0 deletions packages/clickhouse/pkg/hoststats/hoststats_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
package hoststats

import (
"context"
"errors"
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

type fakeDelivery struct {
pushed []SandboxHostStat
pushErr error
closed bool
closeErr error
}

func (f *fakeDelivery) Push(stat SandboxHostStat) error {
f.pushed = append(f.pushed, stat)

return f.pushErr
}

func (f *fakeDelivery) Close(_ context.Context) error {
f.closed = true

return f.closeErr
}

func TestMultiDelivery_PushFansOutToAllTargets(t *testing.T) {
t.Parallel()

a, b := &fakeDelivery{}, &fakeDelivery{}
md := NewMultiDelivery(a, b)

stat := SandboxHostStat{SandboxID: "sbx-1"}
require.NoError(t, md.Push(stat))

assert.Equal(t, []SandboxHostStat{stat}, a.pushed)
assert.Equal(t, []SandboxHostStat{stat}, b.pushed)
}

func TestMultiDelivery_PushContinuesAfterTargetError(t *testing.T) {
t.Parallel()

bad := &fakeDelivery{pushErr: errors.New("boom")}
good := &fakeDelivery{}
md := NewMultiDelivery(bad, good)

stat := SandboxHostStat{SandboxID: "sbx-1"}
err := md.Push(stat)

require.Error(t, err)
require.ErrorContains(t, err, "boom")
// Good target still received the stat — best-effort, no early return.
assert.Equal(t, []SandboxHostStat{stat}, good.pushed)
assert.Equal(t, []SandboxHostStat{stat}, bad.pushed)
}

func TestMultiDelivery_CloseAllTargets(t *testing.T) {
t.Parallel()

a := &fakeDelivery{closeErr: errors.New("a closed badly")}
b := &fakeDelivery{}
md := NewMultiDelivery(a, b)

err := md.Close(context.Background())

require.Error(t, err)
require.ErrorContains(t, err, "a closed badly")
assert.True(t, a.closed)
assert.True(t, b.closed)
}

func TestMultiDelivery_PushJoinsAllErrors(t *testing.T) {
t.Parallel()

errA := errors.New("push failed A")
errB := errors.New("push failed B")
a := &fakeDelivery{pushErr: errA}
b := &fakeDelivery{pushErr: errB}
md := NewMultiDelivery(a, b)

stat := SandboxHostStat{SandboxID: "sbx-1"}
err := md.Push(stat)

require.Error(t, err)
require.ErrorIs(t, err, errA)
require.ErrorIs(t, err, errB)
assert.Equal(t, []SandboxHostStat{stat}, a.pushed)
assert.Equal(t, []SandboxHostStat{stat}, b.pushed)
}

func TestMultiDelivery_CloseJoinsAllErrors(t *testing.T) {
t.Parallel()

errA := errors.New("close failed A")
errB := errors.New("close failed B")
a := &fakeDelivery{closeErr: errA}
b := &fakeDelivery{closeErr: errB}
md := NewMultiDelivery(a, b)

err := md.Close(context.Background())

require.Error(t, err)
require.ErrorIs(t, err, errA)
require.ErrorIs(t, err, errB)
assert.True(t, a.closed)
assert.True(t, b.closed)
}

func TestMultiDelivery_EmptyTargets(t *testing.T) {
t.Parallel()

md := NewMultiDelivery()
assert.NoError(t, md.Push(SandboxHostStat{}))
assert.NoError(t, md.Close(context.Background()))
}

func TestMultiDelivery_SingleTargetBypassesWrapper(t *testing.T) {
t.Parallel()

single := &fakeDelivery{}
md := NewMultiDelivery(single)

// One-target case returns the underlying delivery directly — no wrapper.
assert.Same(t, single, md)

stat := SandboxHostStat{SandboxID: "sbx-1"}
require.NoError(t, md.Push(stat))
assert.Equal(t, []SandboxHostStat{stat}, single.pushed)

require.NoError(t, md.Close(context.Background()))
assert.True(t, single.closed)
}
8 changes: 8 additions & 0 deletions packages/clickhouse/pkg/team.go
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,10 @@ func (c *Client) QueryMaxStartRateTeamMetrics(ctx context.Context, teamID string

// No data -> return 0
if !rows.Next() {
if err := rows.Err(); err != nil {
return MaxTeamMetric{}, fmt.Errorf("error iterating max start rate team metrics: %w", err)
}

return MaxTeamMetric{
Value: 0,
Timestamp: time.Now(),
Expand Down Expand Up @@ -158,6 +162,10 @@ func (c *Client) QueryMaxConcurrentTeamMetrics(ctx context.Context, teamID strin

// No data -> return 0
if !rows.Next() {
if err := rows.Err(); err != nil {
return MaxTeamMetric{}, fmt.Errorf("error iterating max concurrent team metrics: %w", err)
}

return MaxTeamMetric{
Value: 0,
Timestamp: time.Now(),
Expand Down
2 changes: 2 additions & 0 deletions packages/orchestrator/cmd/resume-build/fph_bench.go
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
//go:build linux

package main

import (
Expand Down
Loading
Loading