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: 2 additions & 1 deletion packages/api/internal/cfg/model.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,8 @@ type Config struct {
AnalyticsCollectorAPIToken string `env:"ANALYTICS_COLLECTOR_API_TOKEN"`
AnalyticsCollectorHost string `env:"ANALYTICS_COLLECTOR_HOST"`

ClickhouseConnectionString string `env:"CLICKHOUSE_CONNECTION_STRING"`
ClickhouseConnectionString string `env:"CLICKHOUSE_CONNECTION_STRING"`
ClickhouseConnectionStrings []string `env:"CLICKHOUSE_CONNECTION_STRINGS" envSeparator:";"`

LokiPassword string `env:"LOKI_PASSWORD"`
LokiURL string `env:"LOKI_URL,required"`
Expand Down
29 changes: 19 additions & 10 deletions packages/api/internal/handlers/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -109,16 +109,19 @@ func NewAPIStore(ctx context.Context, tel *telemetry.Client, redisClient redis.U

logger.L().Info(ctx, "Created database client")

var clickhouseStore clickhouse.Clickhouse

clickhouseConnectionString := config.ClickhouseConnectionString
if clickhouseConnectionString == "" {
clickhouseStore = clickhouse.NewNoopClient()
} else {
clickhouseStore, err = clickhouse.New(clickhouseConnectionString)
if err != nil {
logger.L().Fatal(ctx, "initializing ClickHouse store", zap.Error(err))
}
// LD-gated switcher: empty flag → singular DSN (self-managed); "0", "1", …
// → alternates from CLICKHOUSE_CONNECTION_STRINGS. Lets reads shift between
// clusters per-query without restarts. Empty singular DSN falls back to a
// noop client.
clickhouseStore, err := clickhouse.NewSwitchingClient(
ctx,
featureFlags,
config.ClickhouseConnectionString,
config.ClickhouseConnectionStrings,
clickhouse.WithAllowNoopDefault(true),
)
if err != nil {
logger.L().Fatal(ctx, "initializing ClickHouse switching client", zap.Error(err))
}

posthogClient, posthogErr := analyticscollector.NewPosthogClient(ctx, config.PosthogAPIKey)
Expand Down Expand Up @@ -318,6 +321,12 @@ func (a *APIStore) Close(ctx context.Context) error {
errs = append(errs, fmt.Errorf("closing snapshot cache: %w", err))
}

if a.clickhouseStore != nil {
if err := a.clickhouseStore.Close(ctx); err != nil {
errs = append(errs, fmt.Errorf("closing ClickHouse store: %w", err))
}
}

return errors.Join(errs...)
}

Expand Down
93 changes: 93 additions & 0 deletions packages/clickhouse/pkg/switcher.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
package clickhouse

import (
"context"
"time"

"github.com/e2b-dev/infra/packages/shared/pkg/featureflags"
"github.com/e2b-dev/infra/packages/shared/pkg/utils/switching"
)

// SwitchingClient is a Clickhouse client that routes each read to one of
// several DSNs based on the clickhouse-read-endpoint LaunchDarkly flag.
// Each call delegates to switcher.Resolve(ctx), so the active endpoint can
// change between calls without restarting. See
// packages/shared/pkg/utils/switching for the underlying mechanism.
type SwitchingClient struct {
switcher *switching.Switcher[Clickhouse]
}

var _ Clickhouse = (*SwitchingClient)(nil)

// NewSwitchingClient builds N+1 clients (one for defaultDSN, one per alternate)
// and selects between them per-call using ClickhouseReadEndpointFlag. An empty
// flag value (the LD default) selects the default client; "0", "1", … select
// alternateDSNs[i]. Invalid values fall back to default + rate-limited warning.
func NewSwitchingClient(
ctx context.Context,
ff *featureflags.Client,
defaultDSN string,
alternateDSNs []string,
opts ...Option,
) (*SwitchingClient, error) {
var sOpts []switching.Option[Clickhouse]
for _, opt := range opts {
if opt != nil {
sOpts = append(sOpts, switching.Option[Clickhouse](opt))
}
}

s, err := switching.New[Clickhouse](
ctx,
ff,
featureflags.ClickhouseReadEndpointFlag,
defaultDSN,
alternateDSNs,
func(dsn string) (Clickhouse, error) { return New(dsn) },
append(sOpts, switching.WithNoopFactory(func() (Clickhouse, error) {
return NewNoopClient(), nil
}))...,
)
if err != nil {
return nil, err
}

return &SwitchingClient{switcher: s}, nil
}

// Option mirrors switching.Option for caller convenience.
type Option switching.Option[Clickhouse]

// WithAllowNoopDefault enables falling back to a noop client when the
// default DSN is empty.
func WithAllowNoopDefault(allow bool) Option {
return Option(switching.WithAllowNoopDefault[Clickhouse](allow))
}

func (s *SwitchingClient) Close(ctx context.Context) error {
return s.switcher.Close(ctx)
}

func (s *SwitchingClient) QuerySandboxTimeRange(ctx context.Context, sandboxID, teamID string) (time.Time, time.Time, error) {
return s.switcher.Resolve(ctx).QuerySandboxTimeRange(ctx, sandboxID, teamID)
}

func (s *SwitchingClient) QuerySandboxMetrics(ctx context.Context, sandboxID, teamID string, start, end time.Time, step time.Duration) ([]Metrics, error) {
return s.switcher.Resolve(ctx).QuerySandboxMetrics(ctx, sandboxID, teamID, start, end, step)
}

func (s *SwitchingClient) QueryLatestMetrics(ctx context.Context, sandboxIDs []string, teamID string) ([]Metrics, error) {
return s.switcher.Resolve(ctx).QueryLatestMetrics(ctx, sandboxIDs, teamID)
}

func (s *SwitchingClient) QueryTeamMetrics(ctx context.Context, teamID string, start, end time.Time, step time.Duration) ([]TeamMetrics, error) {
return s.switcher.Resolve(ctx).QueryTeamMetrics(ctx, teamID, start, end, step)
}

func (s *SwitchingClient) QueryMaxStartRateTeamMetrics(ctx context.Context, teamID string, start, end time.Time, step time.Duration) (MaxTeamMetric, error) {
return s.switcher.Resolve(ctx).QueryMaxStartRateTeamMetrics(ctx, teamID, start, end, step)
}

func (s *SwitchingClient) QueryMaxConcurrentTeamMetrics(ctx context.Context, teamID string, start, end time.Time) (MaxTeamMetric, error) {
return s.switcher.Resolve(ctx).QueryMaxConcurrentTeamMetrics(ctx, teamID, start, end)
}
11 changes: 11 additions & 0 deletions packages/dashboard-api/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,7 @@ require (
github.com/goccy/go-yaml v1.19.2 // indirect
github.com/golang-jwt/jwt/v5 v5.3.1 // indirect
github.com/gorilla/mux v1.8.1 // indirect
github.com/gregjones/httpcache v0.0.0-20190611155906-901d90724c79 // indirect
github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.2 // indirect
github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect
github.com/hashicorp/go-cleanhttp v0.5.2 // indirect
Expand All @@ -92,6 +93,14 @@ require (
github.com/json-iterator/go v1.1.12 // indirect
github.com/klauspost/compress v1.18.5 // indirect
github.com/klauspost/cpuid/v2 v2.3.0 // indirect
github.com/launchdarkly/ccache v1.1.0 // indirect
github.com/launchdarkly/eventsource v1.10.0 // indirect
github.com/launchdarkly/go-jsonstream/v3 v3.1.0 // indirect
github.com/launchdarkly/go-sdk-common/v3 v3.3.0 // indirect
github.com/launchdarkly/go-sdk-events/v3 v3.5.0 // indirect
github.com/launchdarkly/go-semver v1.0.3 // indirect
github.com/launchdarkly/go-server-sdk-evaluation/v3 v3.0.1 // indirect
github.com/launchdarkly/go-server-sdk/v7 v7.13.0 // indirect
github.com/leodido/go-urn v1.4.0 // indirect
github.com/lib/pq v1.11.2 // indirect
github.com/lufia/plan9stats v0.0.0-20240909124753-873cd0166683 // indirect
Expand All @@ -115,6 +124,7 @@ require (
github.com/oasdiff/yaml3 v0.0.12 // indirect
github.com/opencontainers/go-digest v1.0.0 // indirect
github.com/opencontainers/image-spec v1.1.1 // indirect
github.com/patrickmn/go-cache v2.1.0+incompatible // indirect
github.com/paulmach/orb v0.11.1 // indirect
github.com/pelletier/go-toml/v2 v2.3.1 // indirect
github.com/perimeterx/marshmallow v1.1.5 // indirect
Expand Down Expand Up @@ -160,6 +170,7 @@ require (
go.uber.org/multierr v1.11.0 // indirect
golang.org/x/arch v0.25.0 // indirect
golang.org/x/crypto v0.51.0 // indirect
golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f // indirect
golang.org/x/mod v0.36.0 // indirect
golang.org/x/net v0.55.0 // indirect
golang.org/x/oauth2 v0.36.0 // indirect
Expand Down
26 changes: 26 additions & 0 deletions packages/dashboard-api/go.sum

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 5 additions & 0 deletions packages/shared/pkg/featureflags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -433,6 +433,11 @@ var (
DefaultPersistentVolumeType = NewStringFlag("default-persistent-volume-type", "")
BuildNodeInfo = NewJSONFlag("preferred-build-node", ldvalue.Null())
FirecrackerVersions = NewJSONFlag("firecracker-versions", ldvalue.FromJSONMarshal(FirecrackerVersionMap))

// ClickhouseReadEndpointFlag selects which ClickHouse DSN to use for reads.
// "" (empty) → singular CLICKHOUSE_CONNECTION_STRING (self-managed default).
// "0", "1", ... → index into CLICKHOUSE_CONNECTION_STRINGS
ClickhouseReadEndpointFlag = NewStringFlag("clickhouse-read-endpoint", "")
)

// ResolveFirecrackerVersion resolves the firecracker version using the FirecrackerVersions feature flag.
Expand Down
Loading
Loading