Skip to content

Commit 6e2cbad

Browse files
pmakanju-e2be2b-bot[bot]
authored andcommitted
feat(api): add the outbox worker that tears down a deleted team's resources
GitOrigin-RevId: 4675d05f809ca00f11b6dc503fe06dc97d904f5c
1 parent 871e0d1 commit 6e2cbad

23 files changed

Lines changed: 920 additions & 173 deletions

‎docs/ARCHITECTURE.md‎

Lines changed: 23 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -197,6 +197,25 @@ The control-plane entry point (Gin, OpenAPI-generated from `spec/openapi.yml`, p
197197
operation.
198198
- **Extra listeners**: internal gRPC :5009 and edge gRPC :5109 expose `ResumeSandbox` so
199199
client-proxy can wake paused sandboxes on incoming traffic.
200+
- **Outbox** (`internal/outbox`): a River client on the shared database's `river` schema
201+
works jobs that other services enqueue there for work needing the orchestrator
202+
connections. Its job arguments live in `packages/db/pkg/outbox`. The one kind today,
203+
`teardown_team_resources`, stops a deleted team's workloads. A guard runs first on every
204+
attempt: a team that exists and is not blocked fails the attempt, and the job retries rather
205+
than kill a live team's sandboxes; a missing team row does not stop it. Then
206+
`kill_sandboxes` kills the team's running sandboxes and succeeds only once the sandbox store
207+
holds none of them in any state, so a sandbox still pausing or being killed makes the job
208+
retry. A sandbox that is not running cannot be paused, so no new snapshot of the team
209+
appears afterwards. Killing through the orchestrator is what removes the sandboxes' store
210+
records, so the job deletes no Redis state itself, and every attempt is safe to run again.
211+
Builds are not cancelled: they end on their own build timeout. An attempt is bounded at
212+
30 minutes; retries back off from 30 seconds to a 15-minute cap, so the 100 attempts last
213+
about a day before River discards the job and keeps it. Volumes, templates, snapshots and
214+
the team row are left alone. Steps report
215+
`api.outbox.steps.finished` and `api.outbox.step.duration`; finished jobs report
216+
`outbox.jobs.finished`, and the backlog reports `outbox.jobs` and `outbox.oldest_*_age`.
217+
`OUTBOX_MAX_WORKERS` (default 10) and `OUTBOX_BACKLOG_INTERVAL` (default 30s) configure it;
218+
shutdown stops it after the sandbox-work drain, before the database and Redis clients close.
200219
- Reads ClickHouse for sandbox/team metrics endpoints. Sandbox and template-build logs default to
201220
Loki, with a LaunchDarkly-gated ClickHouse read path (`logs-read-config`) for local-cluster logs
202221
during the log storage migration. `LOKI_URL` is optional: without it the api has no Loki client
@@ -423,8 +442,9 @@ uncoded or unknown reasons as generic errors.
423442
`DELETE /v1/management/projects/{teamID}` is declared and answers 501. `envs`, `snapshots` and
424443
`volumes` reference `teams` with `ON DELETE NO ACTION` and templates are only soft-deleted, so a
425444
project that ever built one pins its team row — and releasing it needs the API service's
426-
orchestrator connections, which this service does not have. Projects are not deleted from control
427-
planes today.
445+
orchestrator connections, which this service does not have. The API's `teardown_team_resources`
446+
outbox job (see the API section) is the worker for that release; nothing enqueues it yet, and
447+
projects are not deleted from control planes today.
428448

429449
## Data stores
430450

@@ -447,7 +467,7 @@ The same database holds River's job tables in the `river` schema, which a goose
447467
creates. The db-migrator applies the goose migrations and then River's own, under one advisory
448468
lock; `make migrate` runs it, because the goose CLI cannot apply River's. At startup the API and
449469
dashboard-api refuse a database whose goose version is older than the one they were built against,
450-
or whose River migrations are not current.
470+
or whose River migrations are not current. The API works jobs from that schema (see its Outbox).
451471

452472
A template and a paused-sandbox snapshot have the **same artifact shape** — a snapshot is just a
453473
new build whose memfile/rootfs are stored as diffs against the template it came from (diff chains

‎packages/api/go.mod‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,9 @@ require (
4848
github.com/oapi-codegen/runtime v1.6.0
4949
github.com/posthog/posthog-go v0.0.0-20230801140217-d607812dee69
5050
github.com/redis/go-redis/v9 v9.21.0
51+
github.com/riverqueue/river v0.47.0
52+
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.47.0
53+
github.com/riverqueue/river/rivertype v0.47.0
5154
github.com/stretchr/testify v1.12.1
5255
go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.71.0
5356
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.71.0
@@ -329,12 +332,10 @@ require (
329332
github.com/quic-go/quic-go v0.61.0 // indirect
330333
github.com/redis/go-redis/extra/rediscmd/v9 v9.17.3 // indirect
331334
github.com/redis/go-redis/extra/redisotel/v9 v9.17.3 // indirect
332-
github.com/riverqueue/river v0.47.0 // indirect
333335
github.com/riverqueue/river/riverdriver v0.47.0 // indirect
334336
github.com/riverqueue/river/riverdriver/riverdatabasesql v0.47.0 // indirect
335-
github.com/riverqueue/river/riverdriver/riverpgxv5 v0.47.0 // indirect
336337
github.com/riverqueue/river/rivershared v0.47.0 // indirect
337-
github.com/riverqueue/river/rivertype v0.47.0 // indirect
338+
github.com/riverqueue/rivercontrib/otelriver v0.12.0 // indirect
338339
github.com/rs/zerolog v1.34.0 // indirect
339340
github.com/samber/lo v1.53.0 // indirect
340341
github.com/santhosh-tekuri/jsonschema/v6 v6.0.2 // indirect
@@ -357,6 +358,10 @@ require (
357358
github.com/tdewolff/test v1.0.12 // indirect
358359
github.com/testcontainers/testcontainers-go v0.44.0 // indirect
359360
github.com/testcontainers/testcontainers-go/modules/postgres v0.44.0 // indirect
361+
github.com/tidwall/gjson v1.19.0 // indirect
362+
github.com/tidwall/match v1.2.0 // indirect
363+
github.com/tidwall/pretty v1.2.1 // indirect
364+
github.com/tidwall/sjson v1.2.5 // indirect
360365
github.com/tjhop/slog-gokit v0.1.4 // indirect
361366
github.com/tklauser/go-sysconf v0.4.0 // indirect
362367
github.com/tklauser/numcpus v0.12.0 // indirect

‎packages/api/go.sum‎

Lines changed: 7 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎packages/api/internal/cfg/model.go‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -167,6 +167,13 @@ type Config struct {
167167
// placement (max of CPU and pool). A node that reports no pool scores
168168
// 0.5, so clusters without a pool should set this false and rank on CPU.
169169
BestOfKHugepageMemory bool `env:"BEST_OF_K_HUGEPAGE_MEMORY" envDefault:"true"`
170+
171+
// OutboxMaxWorkers is how many River outbox jobs one API replica works at
172+
// once.
173+
OutboxMaxWorkers int `env:"OUTBOX_MAX_WORKERS" envDefault:"10"`
174+
// OutboxBacklogInterval is how often the outbox backlog gauges are read
175+
// from the database.
176+
OutboxBacklogInterval time.Duration `env:"OUTBOX_BACKLOG_INTERVAL" envDefault:"30s"`
170177
}
171178

172179
type FailureCondition string
Lines changed: 7 additions & 101 deletions
Original file line numberDiff line numberDiff line change
@@ -1,132 +1,38 @@
11
package handlers
22

33
import (
4-
"context"
5-
"fmt"
64
"net/http"
7-
"sync/atomic"
8-
"time"
95

106
"github.com/gin-gonic/gin"
117
"github.com/google/uuid"
128
"go.uber.org/zap"
13-
"golang.org/x/sync/errgroup"
149

1510
"github.com/e2b-dev/infra/packages/api/internal/api"
16-
dbtypes "github.com/e2b-dev/infra/packages/db/pkg/types"
17-
"github.com/e2b-dev/infra/packages/db/queries"
18-
"github.com/e2b-dev/infra/packages/shared/pkg/clusters"
19-
templatemanagergrpc "github.com/e2b-dev/infra/packages/shared/pkg/grpc/template-manager"
2011
"github.com/e2b-dev/infra/packages/shared/pkg/logger"
2112
)
2213

23-
// buildDeleteTimeout bounds the node-side delete once the build has been
24-
// recorded as failed.
25-
const buildDeleteTimeout = 30 * time.Second
26-
27-
type buildCanceller interface {
28-
SetTerminalStatus(ctx context.Context, buildID uuid.UUID, statusGroup dbtypes.BuildStatusGroup, reason *templatemanagergrpc.TemplateBuildStatusReason) (bool, error)
29-
DeleteBuild(ctx context.Context, buildID uuid.UUID, templateID string, clusterID uuid.UUID, nodeID string) error
30-
}
31-
32-
// cancelBuild ends a build and stops it on its node. The node-side delete takes
33-
// the build's artifacts with it, so it only follows a write that ended the
34-
// build: one that finished on its own between the listing and here keeps both
35-
// its outcome and the artifacts that outcome refers to.
36-
func cancelBuild(ctx context.Context, tm buildCanceller, build queries.GetCancellableTemplateBuildsByTeamRow) error {
37-
recorded, err := tm.SetTerminalStatus(ctx, build.BuildID, dbtypes.BuildStatusGroupFailed, &templatemanagergrpc.TemplateBuildStatusReason{
38-
Message: "cancelled by admin",
39-
})
40-
if err != nil {
41-
return fmt.Errorf("failed to set build status to failed: %w", err)
42-
}
43-
44-
if !recorded || build.ClusterNodeID == nil {
45-
return nil
46-
}
47-
48-
// The write above dropped the build from every listing, so nothing retries
49-
// this stop, and the request context dies with the caller or its deadline.
50-
deleteCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), buildDeleteTimeout)
51-
defer cancel()
52-
53-
err = tm.DeleteBuild(deleteCtx, build.BuildID, build.TemplateID, clusters.WithClusterFallback(build.ClusterID), *build.ClusterNodeID)
54-
if err != nil {
55-
return fmt.Errorf("failed to delete build on node: %w", err)
56-
}
57-
58-
return nil
59-
}
60-
6114
func (a *APIStore) PostAdminTeamsTeamIDBuildsCancel(c *gin.Context, teamID uuid.UUID) {
6215
ctx := c.Request.Context()
6316
ctx, span := tracer.Start(ctx, "cancel admin-team-builds")
6417
defer span.End()
6518

6619
logger.L().Info(ctx, "Admin cancelling all builds for team", logger.WithTeamID(teamID.String()))
6720

68-
builds, err := a.sqlcDB.GetCancellableTemplateBuildsByTeam(ctx, teamID)
21+
cancelledCount, failedCount, err := a.templateManager.CancelTeamBuilds(ctx, teamID, "cancelled by admin")
6922
if err != nil {
7023
a.sendAPIStoreError(c, http.StatusInternalServerError, "Failed to get builds")
7124

7225
return
7326
}
7427

75-
logger.L().Info(ctx, "Found builds to cancel",
76-
logger.WithTeamID(teamID.String()),
77-
zap.Int("count", len(builds)),
78-
)
79-
80-
cancelledCount := atomic.Int64{}
81-
failedCount := atomic.Int64{}
82-
83-
wg := errgroup.Group{}
84-
wg.SetLimit(10)
85-
86-
for _, b := range builds {
87-
wg.Go(func() error {
88-
buildID := b.BuildID
89-
templateID := b.TemplateID
90-
91-
err := cancelBuild(ctx, a.templateManager, b)
92-
if err != nil {
93-
logger.L().Error(ctx, "Failed to cancel build",
94-
zap.String("buildID", buildID.String()),
95-
zap.String("templateID", templateID),
96-
logger.WithTeamID(teamID.String()),
97-
zap.Error(err))
98-
failedCount.Add(1)
99-
100-
return nil
101-
}
102-
103-
logger.L().Debug(ctx, "Successfully cancelled build",
104-
zap.String("buildID", buildID.String()),
105-
zap.String("templateID", templateID),
106-
logger.WithTeamID(teamID.String()))
107-
cancelledCount.Add(1)
108-
109-
return nil
110-
})
111-
}
112-
113-
err = wg.Wait()
114-
if err != nil {
115-
a.sendAPIStoreError(c, http.StatusInternalServerError, "Failed to cancel builds")
116-
117-
return
118-
}
119-
12028
logger.L().Info(ctx, "Completed cancelling team builds",
12129
logger.WithTeamID(teamID.String()),
122-
zap.Int64("cancelled", cancelledCount.Load()),
123-
zap.Int64("failed", failedCount.Load()),
30+
zap.Int("cancelled", cancelledCount),
31+
zap.Int("failed", failedCount),
12432
)
12533

126-
result := api.AdminBuildCancelResult{
127-
CancelledCount: int(cancelledCount.Load()),
128-
FailedCount: int(failedCount.Load()),
129-
}
130-
131-
c.JSON(http.StatusOK, result)
34+
c.JSON(http.StatusOK, api.AdminBuildCancelResult{
35+
CancelledCount: cancelledCount,
36+
FailedCount: failedCount,
37+
})
13238
}

‎packages/api/internal/handlers/admin_kill_team_sandboxes.go‎

Lines changed: 7 additions & 57 deletions
Original file line numberDiff line numberDiff line change
@@ -2,12 +2,10 @@ package handlers
22

33
import (
44
"net/http"
5-
"sync/atomic"
65

76
"github.com/gin-gonic/gin"
87
"github.com/google/uuid"
98
"go.uber.org/zap"
10-
"golang.org/x/sync/errgroup"
119

1210
"github.com/e2b-dev/infra/packages/api/internal/api"
1311
"github.com/e2b-dev/infra/packages/api/internal/sandbox"
@@ -28,58 +26,13 @@ func (a *APIStore) PostAdminTeamsTeamIDSandboxesKill(c *gin.Context, teamID uuid
2826

2927
logger.L().Info(ctx, "Admin killing all sandboxes for team", logger.WithTeamID(teamID.String()))
3028

31-
// Get all running sandboxes for the team
32-
sandboxes, err := a.orchestrator.GetSandboxes(ctx, teamID, []sandbox.State{sandbox.StateRunning})
29+
killedCount, failedCount, err := a.orchestrator.KillTeamSandboxes(ctx, teamID, sandbox.KillReasonAdmin)
3330
if err != nil {
3431
a.sendAPIStoreError(c, http.StatusInternalServerError, "Failed to get sandboxes")
3532

3633
return
3734
}
3835

39-
logger.L().Info(ctx, "Found sandboxes to kill",
40-
logger.WithTeamID(teamID.String()),
41-
zap.Int("count", len(sandboxes)),
42-
)
43-
44-
killedCount := atomic.Int64{}
45-
failedCount := atomic.Int64{}
46-
47-
wg := errgroup.Group{}
48-
wg.SetLimit(10)
49-
50-
// Kill each sandbox
51-
for _, sbx := range sandboxes {
52-
wg.Go(func() error {
53-
err := a.orchestrator.RemoveSandbox(ctx, sbx.TeamID, sbx.SandboxID, sandbox.RemoveOpts{
54-
Action: sandbox.StateActionKill,
55-
Reason: sandbox.KillReasonAdmin,
56-
})
57-
if err != nil {
58-
logger.L().Error(ctx, "Failed to kill sandbox",
59-
logger.WithSandboxID(sbx.SandboxID),
60-
logger.WithTeamID(teamID.String()),
61-
zap.String("kill_reason", sandbox.KillReasonAdmin.String()),
62-
zap.Error(err))
63-
failedCount.Add(1)
64-
} else {
65-
logger.L().Debug(ctx, "Successfully killed sandbox",
66-
logger.WithSandboxID(sbx.SandboxID),
67-
logger.WithTeamID(teamID.String()),
68-
zap.String("kill_reason", sandbox.KillReasonAdmin.String()))
69-
killedCount.Add(1)
70-
}
71-
72-
return nil
73-
})
74-
}
75-
76-
err = wg.Wait()
77-
if err != nil {
78-
a.sendAPIStoreError(c, http.StatusInternalServerError, "Failed to kill sandboxes")
79-
80-
return
81-
}
82-
8336
// Invalidate auth cache for this team so subsequent requests re-check against DB
8437
if err := a.authService.InvalidateTeamCache(ctx, teamID); err != nil {
8538
logger.L().Error(ctx, "Failed to invalidate auth cache for team",
@@ -89,15 +42,12 @@ func (a *APIStore) PostAdminTeamsTeamIDSandboxesKill(c *gin.Context, teamID uuid
8942

9043
logger.L().Info(ctx, "Completed killing team sandboxes",
9144
zap.String("teamID", teamID.String()),
92-
zap.Int64("killed", killedCount.Load()),
93-
zap.Int64("failed", failedCount.Load()),
45+
zap.Int("killed", killedCount),
46+
zap.Int("failed", failedCount),
9447
)
9548

96-
// Return result
97-
result := api.AdminSandboxKillResult{
98-
KilledCount: int(killedCount.Load()),
99-
FailedCount: int(failedCount.Load()),
100-
}
101-
102-
c.JSON(http.StatusOK, result)
49+
c.JSON(http.StatusOK, api.AdminSandboxKillResult{
50+
KilledCount: killedCount,
51+
FailedCount: failedCount,
52+
})
10353
}

‎packages/api/internal/handlers/store.go‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import (
2626
"github.com/e2b-dev/infra/packages/api/internal/cfg"
2727
"github.com/e2b-dev/infra/packages/api/internal/clusters"
2828
"github.com/e2b-dev/infra/packages/api/internal/orchestrator"
29+
"github.com/e2b-dev/infra/packages/api/internal/outbox"
2930
"github.com/e2b-dev/infra/packages/api/internal/sandbox"
3031
managementv1 "github.com/e2b-dev/infra/packages/api/internal/secretsstore/management/v1"
3132
template_manager "github.com/e2b-dev/infra/packages/api/internal/template-manager"
@@ -473,6 +474,22 @@ func NewAPIStore(ctx context.Context, tel *telemetry.Client, redisClient redis.U
473474
return a
474475
}
475476

477+
// NewOutbox builds the River client that works the API's outbox jobs on the
478+
// store's clients. Stop it before Close.
479+
func (a *APIStore) NewOutbox(l logger.Logger) (*outbox.River, error) {
480+
return outbox.New(outbox.Dependencies{
481+
Pool: a.sqlcDB.Pool(),
482+
Sandboxes: a.orchestrator,
483+
Teams: a.authDB,
484+
Logger: l,
485+
Telemetry: a.Telemetry,
486+
Config: outbox.Config{
487+
MaxWorkers: a.config.OutboxMaxWorkers,
488+
BacklogInterval: a.config.OutboxBacklogInterval,
489+
},
490+
})
491+
}
492+
476493
// Drain stops admitting sandbox work that outlives its request and waits for
477494
// what is in flight. It runs before Close, which tears down the clients that
478495
// work uses.

0 commit comments

Comments
 (0)