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
128 changes: 128 additions & 0 deletions packages/api/internal/edge/metrics.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
package edge

import (
"context"
"fmt"
"net/http"

"github.com/google/uuid"

"github.com/e2b-dev/infra/packages/api/internal/api"
apiedge "github.com/e2b-dev/infra/packages/shared/pkg/http/edge"
)

func GetClusterSandboxMetrics(ctx context.Context, pool *Pool, sandboxID string, teamID string, clusterID uuid.UUID, qStart *int64, qEnd *int64) ([]api.SandboxMetric, *api.APIError) {
cluster, ok := pool.GetClusterById(clusterID)
if !ok {
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: fmt.Sprintf("Error getting cluster '%s'", clusterID),
Err: fmt.Errorf("cluster with ID '%s' not found", clusterID),
}
}

res, err := cluster.GetHttpClient().V1SandboxMetricsWithResponse(
ctx, sandboxID, &apiedge.V1SandboxMetricsParams{
TeamID: teamID,
Start: qStart,
End: qEnd,
},
)
if err != nil {
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: fmt.Sprintf("Error getting metrics for sandbox '%s'", sandboxID),
Err: fmt.Errorf("error getting metrics for sandbox '%s': %w", sandboxID, err),
}
}

if res.StatusCode() != http.StatusOK {
return nil, &api.APIError{
Code: res.StatusCode(),
ClientMsg: fmt.Sprintf("Error getting metrics for sandbox '%s'", sandboxID),
Err: fmt.Errorf("unexpected response for sandbox - HTTP status '%d'", res.StatusCode()),
}
}

if res.JSON200 == nil {
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: fmt.Sprintf("Error getting metrics for sandbox '%s'", sandboxID),
Err: fmt.Errorf("no metrics returned for sandbox '%s'", sandboxID),
}
}

// Transform edge types (snake_case) to API types (camelCase)
apiMetrics := make([]api.SandboxMetric, len(*res.JSON200))
for i, m := range *res.JSON200 {
apiMetrics[i] = api.SandboxMetric{
Timestamp: m.Timestamp,
TimestampUnix: m.TimestampUnix,
CpuUsedPct: m.CpuUsedPct,
CpuCount: m.CpuCount,
MemTotal: m.MemTotal,
MemUsed: m.MemUsed,
DiskTotal: m.DiskTotal,
DiskUsed: m.DiskUsed,
}
}

return apiMetrics, nil
}

func GetClusterSandboxListMetrics(ctx context.Context, pool *Pool, teamID string, clusterID uuid.UUID, sandboxIDs []string) (map[string]api.SandboxMetric, *api.APIError) {
cluster, ok := pool.GetClusterById(clusterID)
if !ok {
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: fmt.Sprintf("Error getting cluster '%s'", clusterID),
Err: fmt.Errorf("cluster with ID '%s' not found", clusterID),
}
}

res, err := cluster.GetHttpClient().V1SandboxesMetricsWithResponse(
ctx, &apiedge.V1SandboxesMetricsParams{
TeamID: teamID,
SandboxIds: sandboxIDs,
},
)
if err != nil {
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: "Error getting metrics for sandbox list",
Err: fmt.Errorf("error getting metrics for sandbox list: %w", err),
}
}

if res.StatusCode() != http.StatusOK {
return nil, &api.APIError{
Code: res.StatusCode(),
ClientMsg: "Error getting metrics for sandbox list",
Err: fmt.Errorf("unexpected response for sandbox list - HTTP status '%d'", res.StatusCode()),
}
}

if res.JSON200 == nil {
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: "Error getting metrics for sandbox list",
Err: fmt.Errorf("no metrics returned for sandbox list"),
}
}

apiMetrics := make(map[string]api.SandboxMetric)
for key, m := range res.JSON200.Sandboxes {
Comment thread
jakubno marked this conversation as resolved.
apiMetrics[key] = api.SandboxMetric{
Timestamp: m.Timestamp,
TimestampUnix: m.TimestampUnix,
CpuUsedPct: m.CpuUsedPct,
CpuCount: m.CpuCount,
MemTotal: m.MemTotal,
MemUsed: m.MemUsed,
DiskTotal: m.DiskTotal,
DiskUsed: m.DiskUsed,
}
}

return apiMetrics, nil
}
83 changes: 63 additions & 20 deletions packages/api/internal/handlers/sandbox_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,17 +7,18 @@ import (
"time"

"github.com/gin-gonic/gin"
"github.com/launchdarkly/go-sdk-common/v3/ldcontext"
"go.uber.org/zap"

"github.com/e2b-dev/infra/packages/api/internal/api"
"github.com/e2b-dev/infra/packages/api/internal/auth"
"github.com/e2b-dev/infra/packages/api/internal/db/types"
"github.com/e2b-dev/infra/packages/api/internal/edge"
"github.com/e2b-dev/infra/packages/api/internal/utils"
clickhouse "github.com/e2b-dev/infra/packages/clickhouse/pkg"
clickhouseUtils "github.com/e2b-dev/infra/packages/clickhouse/pkg/utils"
featureflags "github.com/e2b-dev/infra/packages/shared/pkg/feature-flags"
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
"github.com/e2b-dev/infra/packages/shared/pkg/telemetry"
)

func (a *APIStore) GetSandboxesSandboxIDMetrics(c *gin.Context, sandboxID string, params api.GetSandboxesSandboxIDMetricsParams) {
Expand All @@ -28,10 +29,20 @@ func (a *APIStore) GetSandboxesSandboxIDMetrics(c *gin.Context, sandboxID string

team := c.Value(auth.TeamContextKey).(*types.Team)

metricsReadFlag, err := a.featureFlags.BoolFlag(ctx, featureflags.MetricsReadFlagName,
featureflags.SandboxContext(sandboxID))
// Build the context for feature flags
ctx = featureflags.SetContext(
ctx,
ldcontext.NewBuilder(sandboxID).
Kind(featureflags.SandboxKind).
Build(),
ldcontext.NewBuilder(team.ID.String()).
Kind(featureflags.TeamKind).
Build(),
)

metricsReadFlag, err := a.featureFlags.BoolFlag(ctx, featureflags.MetricsReadFlagName)
if err != nil {
logger.L().Error(ctx, "error getting metrics read feature flag, soft failing", zap.Error(err))
logger.L().Warn(ctx, "error getting metrics read feature flag, soft failing", zap.Error(err))
}

if !metricsReadFlag {
Expand All @@ -44,35 +55,67 @@ func (a *APIStore) GetSandboxesSandboxIDMetrics(c *gin.Context, sandboxID string
return
}

start, end, err := getSandboxStartEndTime(ctx, a.clickhouseStore, team.ID.String(), sandboxID, params)
// TODO: Remove in [ENG-3377], once edge is migrated
edgeProvidedMetrics, err := a.featureFlags.BoolFlag(ctx, featureflags.EdgeProvidedSandboxMetricsFlagName)
if err != nil {
a.sendAPIStoreError(c, http.StatusInternalServerError, fmt.Sprintf("error when getting metrics time range: %s", err))
logger.L().Warn(ctx, "error getting edge provided metrics feature flag, soft failing", zap.Error(err))
}

var metrics []api.SandboxMetric
var apiErr *api.APIError
if edgeProvidedMetrics {
metrics, apiErr = edge.GetClusterSandboxMetrics(
Comment thread
jakubno marked this conversation as resolved.
ctx,
a.clustersPool,
sandboxID,
team.ID.String(),
utils.WithClusterFallback(team.ClusterID),
params.Start,
params.End,
)
} else {
metrics, apiErr = a.getApiProvidedMetrics(ctx, team, sandboxID, params)
}
if apiErr != nil {
logger.L().Error(ctx, "error getting sandbox metrics", zap.Error(apiErr.Err))
a.sendAPIStoreError(c, apiErr.Code, apiErr.ClientMsg)

return
}

start, end, err = clickhouseUtils.ValidateRange(start, end)
c.JSON(http.StatusOK, metrics)
}

// TODO: Remove in [ENG-3377], once edge is migrated
func (a *APIStore) getApiProvidedMetrics(ctx context.Context, team *types.Team, sandboxID string, params api.GetSandboxesSandboxIDMetricsParams) ([]api.SandboxMetric, *api.APIError) {
start, end, err := getSandboxStartEndTime(ctx, a.clickhouseStore, team.ID.String(), sandboxID, params)
if err != nil {
telemetry.ReportError(ctx, "error validating dates", err, telemetry.WithTeamID(team.ID.String()))
a.sendAPIStoreError(c, http.StatusBadRequest, err.Error())
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: "error getting metrics time range",
Err: fmt.Errorf("error getting metrics time range: %w", err),
}
}

return
start, end, err = clickhouseUtils.ValidateRange(start, end)
if err != nil {
return nil, &api.APIError{
Code: http.StatusBadRequest,
ClientMsg: fmt.Sprintf("error validating time range: %s", err),
Err: err,
}
}

// Calculate the step size
step := clickhouseUtils.CalculateStep(start, end)

metrics, err := a.clickhouseStore.QuerySandboxMetrics(ctx, sandboxID, team.ID.String(), start, end, step)
if err != nil {
logger.L().Error(ctx, "Error fetching sandbox metrics from ClickHouse",
logger.WithSandboxID(sandboxID),
logger.WithTeamID(team.ID.String()),
zap.Error(err),
)

a.sendAPIStoreError(c, http.StatusInternalServerError, fmt.Sprintf("error querying sandbox metrics: %s", err))

return
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: "error querying sandbox metrics",
Err: fmt.Errorf("error querying sandbox metrics: %w", err),
}
}

apiMetrics := make([]api.SandboxMetric, len(metrics))
Expand All @@ -89,7 +132,7 @@ func (a *APIStore) GetSandboxesSandboxIDMetrics(c *gin.Context, sandboxID string
}
}

c.JSON(http.StatusOK, apiMetrics)
return apiMetrics, nil
}

func getSandboxStartEndTime(ctx context.Context, clickhouseStore clickhouse.Clickhouse, teamID, sandboxID string, params api.GetSandboxesSandboxIDMetricsParams) (time.Time, time.Time, error) {
Expand Down
59 changes: 51 additions & 8 deletions packages/api/internal/handlers/sandboxes_list_metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,14 @@ import (

"github.com/gin-gonic/gin"
"github.com/google/uuid"
"github.com/launchdarkly/go-sdk-common/v3/ldcontext"
"go.opentelemetry.io/otel/attribute"
"go.uber.org/zap"

"github.com/e2b-dev/infra/packages/api/internal/api"
"github.com/e2b-dev/infra/packages/api/internal/auth"
"github.com/e2b-dev/infra/packages/api/internal/db/types"
"github.com/e2b-dev/infra/packages/api/internal/edge"
"github.com/e2b-dev/infra/packages/api/internal/utils"
featureflags "github.com/e2b-dev/infra/packages/shared/pkg/feature-flags"
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
Expand All @@ -24,8 +26,9 @@ const maxSandboxMetricsCount = 100
func (a *APIStore) getSandboxesMetrics(
ctx context.Context,
teamID uuid.UUID,
clusterID *uuid.UUID,
sandboxIDs []string,
) (map[string]api.SandboxMetric, error) {
) (map[string]api.SandboxMetric, *api.APIError) {
ctx, span := tracer.Start(ctx, "fetch-sandboxes-metrics")
defer span.End()

Expand All @@ -51,14 +54,46 @@ func (a *APIStore) getSandboxesMetrics(
return make(map[string]api.SandboxMetric), nil
}

metrics, err := a.clickhouseStore.QueryLatestMetrics(ctx, sandboxIDs, teamID.String())
// TODO: Remove in [ENG-3377], once edge is migrated
edgeProvidedMetrics, err := a.featureFlags.BoolFlag(ctx, featureflags.EdgeProvidedSandboxMetricsFlagName)
if err != nil {
logger.L().Warn(ctx, "error getting edge provided metrics feature flag, soft failing", zap.Error(err))
}

var metrics map[string]api.SandboxMetric
var apiErr *api.APIError
if edgeProvidedMetrics {
metrics, apiErr = edge.GetClusterSandboxListMetrics(
ctx,
a.clustersPool,
teamID.String(),
utils.WithClusterFallback(clusterID),
sandboxIDs,
)
} else {
metrics, apiErr = a.getApiProvidedSandboxListMetrics(ctx, teamID.String(), sandboxIDs)
}
if apiErr != nil {
return nil, apiErr
}
Comment thread
jakubno marked this conversation as resolved.

return metrics, nil
}

// TODO: Remove in [ENG-3377], once edge is migrated
func (a *APIStore) getApiProvidedSandboxListMetrics(ctx context.Context, teamID string, sandboxIDs []string) (map[string]api.SandboxMetric, *api.APIError) {
metrics, err := a.clickhouseStore.QueryLatestMetrics(ctx, sandboxIDs, teamID)
if err != nil {
logger.L().Error(ctx, "Error fetching sandbox metrics from ClickHouse",
logger.WithTeamID(teamID.String()),
logger.WithTeamID(teamID),
zap.Error(err),
)

return nil, fmt.Errorf("error querying metrics: %w", err)
return nil, &api.APIError{
Code: http.StatusInternalServerError,
ClientMsg: "Error fetching sandbox metrics",
Err: err,
}
}

apiMetrics := make(map[string]api.SandboxMetric)
Expand Down Expand Up @@ -96,10 +131,18 @@ func (a *APIStore) GetSandboxesMetrics(c *gin.Context, params api.GetSandboxesMe
properties := a.posthog.GetPackageToPosthogProperties(&c.Request.Header)
a.posthog.CreateAnalyticsTeamEvent(ctx, team.ID.String(), "listed running instances with metrics", properties)

sandboxesWithMetrics, err := a.getSandboxesMetrics(ctx, team.ID, params.SandboxIds)
if err != nil {
telemetry.ReportCriticalError(ctx, "error fetching metrics for sandboxes", err)
a.sendAPIStoreError(c, http.StatusInternalServerError, fmt.Sprintf("Error returning metrics for sandboxes for team '%s'", team.ID))
// Build the context for feature flags
ctx = featureflags.SetContext(
ctx,
ldcontext.NewBuilder(team.ID.String()).
Kind(featureflags.TeamKind).
Build(),
)

sandboxesWithMetrics, apiErr := a.getSandboxesMetrics(ctx, team.ID, team.ClusterID, params.SandboxIds)
if apiErr != nil {
logger.L().Error(ctx, "error getting sandbox metrics", zap.Error(apiErr.Err))
a.sendAPIStoreError(c, apiErr.Code, apiErr.ClientMsg)

return
}
Expand Down
13 changes: 7 additions & 6 deletions packages/shared/pkg/feature-flags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,12 +47,13 @@ func newBoolFlag(name string, fallback bool) BoolFlag {
}

var (
MetricsWriteFlagName = newBoolFlag("sandbox-metrics-write", env.IsDevelopment())
MetricsReadFlagName = newBoolFlag("sandbox-metrics-read", env.IsDevelopment())
SnapshotFeatureFlagName = newBoolFlag("use-nfs-for-snapshots", env.IsDevelopment())
TemplateFeatureFlagName = newBoolFlag("use-nfs-for-templates", env.IsDevelopment())
BestOfKCanFit = newBoolFlag("best-of-k-can-fit", true)
BestOfKTooManyStarting = newBoolFlag("best-of-k-too-many-starting", false)
MetricsWriteFlagName = newBoolFlag("sandbox-metrics-write", env.IsDevelopment())
MetricsReadFlagName = newBoolFlag("sandbox-metrics-read", env.IsDevelopment())
SnapshotFeatureFlagName = newBoolFlag("use-nfs-for-snapshots", env.IsDevelopment())
TemplateFeatureFlagName = newBoolFlag("use-nfs-for-templates", env.IsDevelopment())
BestOfKCanFit = newBoolFlag("best-of-k-can-fit", true)
BestOfKTooManyStarting = newBoolFlag("best-of-k-too-many-starting", false)
EdgeProvidedSandboxMetricsFlagName = newBoolFlag("edge-provided-sandbox-metrics", false)
)

type IntFlag struct {
Expand Down