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
41 changes: 41 additions & 0 deletions docs/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@ flowchart TB

LB["Load balancer<br/>api.* | *.domain wildcard | docker.*"]

VC["volume-content API (belt)<br/>api.&lt;domain&gt;"]

subgraph controlplane["Control plane (API node pool)"]
API["API<br/>REST :80, gRPC :5009/:5109"]
DashAPI["dashboard-api :3010"]
Expand Down Expand Up @@ -68,6 +70,8 @@ flowchart TB
SDK -->|REST| LB --> API
Browser -->|"port-sandboxid.domain"| LB --> CP
Docker --> LB --> DRP
API -.->|"mint content token + domain"| SDK
SDK -->|"volume content (token-authed)<br/>api.&lt;BYOC or default domain&gt;"| VC
API -->|"gRPC Create/Delete/Pause"| ORCH
API -->|"gRPC TemplateCreate"| TM
CP -->|"lookup sandbox → node"| RD
Expand Down Expand Up @@ -306,6 +310,43 @@ sequenceDiagram
E-->>U: response
```

### Volume content

Persistent volumes (`packages/orchestrator/pkg/volumes/`) are managed through the control-plane
API (`POST/GET /volumes`), but their **content** — reading and writing files — is served by a
separate volume-content API (belt, `e2b-dev/belt`) that the SDK talks to directly, not through the
control-plane API. The API's role is to mint the credential and tell the SDK where to send content
traffic.

```mermaid
sequenceDiagram
autonumber
participant U as SDK
participant API as API
participant PG as PostgreSQL
participant VC as volume-content API (belt)

U->>API: POST /volumes (create) or GET /volumes/{id}
API->>PG: persist / load volume row
API->>API: mint JWT (aud = https://api.&lt;domain&gt;)<br/>resolve domain
API-->>U: { volumeID, name, token, domain? }
Note over U: domain is returned only for BYOC teams;<br/>SDK stores it and falls back to api.&lt;E2B_DOMAIN&gt; otherwise
U->>VC: /volumecontent/{id}/... at api.&lt;domain&gt;<br/>Authorization: Bearer token
VC->>VC: verify token (audience must match its own origin)
VC-->>U: file content
```

- **Domain selection.** The token's audience and the content host are the same origin,
`https://api.<domain>`. For teams on a **custom (BYOC) cluster** (`team.ClusterID` set), the API
returns that cluster's domain (`cluster.SandboxDomain`, resolved in
`handlers.volumeContentDomain`) so content traffic goes to the BYOC cluster's edge instead of the
control-plane host. For teams on the default cluster the response omits `domain` and the SDK uses
its configured default (`api.<E2B_DOMAIN>`); the audience then uses the deployment's `DOMAIN_NAME`.
- **Token.** A short-lived JWT (`handlers.generateVolumeContentToken`, config in
`cfg.VolumesTokenConfig`) signed by the API, scoped to the team and volume, presented as a bearer
token on every content request. Its `aud` claim is `https://api.<domain>`, so a token minted for
one cluster's origin is not accepted by another.

### Pause and resume

- **Pause**: API records a snapshot row in Postgres, then gRPC `Pause` to the node. The
Expand Down
416 changes: 212 additions & 204 deletions packages/api/internal/api/api.gen.go

Large diffs are not rendered by default.

29 changes: 29 additions & 0 deletions packages/api/internal/clusters/mock.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
package clusters

import (
"github.com/google/uuid"

"github.com/e2b-dev/infra/packages/shared/pkg/smap"
)

// NewTestPool builds a Pool pre-populated with the given clusters, for use in
// tests that need cluster lookups (e.g. GetClusterById) without spinning up a
// full synchronization loop.
func NewTestPool(clusters ...*Cluster) *Pool {
m := smap.New[*Cluster]()
for _, c := range clusters {
m.Insert(c.ID.String(), c)
}

return &Pool{clusters: m}
}

// NewTestCluster builds a minimal Cluster carrying just an ID and sandbox
// domain, for tests. Other collaborators (instances, synchronization,
// resources) are left nil.
func NewTestCluster(id uuid.UUID, sandboxDomain *string) *Cluster {
return &Cluster{
ID: id,
SandboxDomain: sandboxDomain,
}
}
13 changes: 12 additions & 1 deletion packages/api/internal/handlers/volume_create.go
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,16 @@ func (a *APIStore) PostVolumes(c *gin.Context) {

clusterID := clustershared.WithClusterFallback(team.ClusterID)

// Resolve the BYOC domain up front so we fail before allocating any
// resources if the team's cluster can't be found.
domain, err := a.volumeContentDomain(team)
if err != nil {
a.sendAPIStoreError(c, http.StatusServiceUnavailable, "Cluster not found")
telemetry.ReportError(ctx, "cluster not found", err)

return
}

// The volume identity we intend to create. We generate the ID up front so we
// can hand it to the orchestrator, which is given the chance to adjust these
// values (e.g. resolve a placeholder volume type) and returns the definitive
Expand Down Expand Up @@ -164,7 +174,7 @@ func (a *APIStore) PostVolumes(c *gin.Context) {
Set("volume_type", volume.VolumeType),
)

token, apiErr := generateVolumeContentToken(a.config.VolumesToken, volume, team)
token, apiErr := generateVolumeContentToken(a.config.VolumesToken, volume, team, a.volumeTokenAudience(domain))
if apiErr != nil {
a.sendAPIStoreError(c, apiErr.Code, apiErr.ClientMsg)
telemetry.ReportCriticalError(ctx, apiErr.ClientMsg, apiErr.Err)
Expand All @@ -176,6 +186,7 @@ func (a *APIStore) PostVolumes(c *gin.Context) {
VolumeID: volume.ID.String(),
Name: volume.Name,
Token: token,
Domain: domain,
}

c.JSON(http.StatusCreated, result)
Expand Down
11 changes: 10 additions & 1 deletion packages/api/internal/handlers/volume_get.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,15 @@ func (a *APIStore) GetVolumesVolumeID(c *gin.Context, volumeID api.VolumeID) {
return
}

token, err := generateVolumeContentToken(a.config.VolumesToken, volume, team)
domain, domainErr := a.volumeContentDomain(team)
if domainErr != nil {
a.sendAPIStoreError(c, http.StatusServiceUnavailable, "Cluster not found")
telemetry.ReportError(c.Request.Context(), "cluster not found", domainErr)

return
}

token, err := generateVolumeContentToken(a.config.VolumesToken, volume, team, a.volumeTokenAudience(domain))
if err != nil {
a.sendAPIStoreError(c, http.StatusInternalServerError, "failed to sign token")
telemetry.ReportCriticalError(c.Request.Context(), "failed to sign token", err)
Expand All @@ -27,6 +35,7 @@ func (a *APIStore) GetVolumesVolumeID(c *gin.Context, volumeID api.VolumeID) {
VolumeID: volume.ID.String(),
Name: volume.Name,
Token: token,
Domain: domain,
}

c.JSON(http.StatusOK, result)
Expand Down
10 changes: 9 additions & 1 deletion packages/api/internal/handlers/volume_get_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,8 @@ func TestGenerateVolumeContentToken_SetsTokidHeader(t *testing.T) {
}

// Act: generate token
tokenStr, apiErr := generateVolumeContentToken(config, volume, team)
audience := "https://api.custom.example.com"
tokenStr, apiErr := generateVolumeContentToken(config, volume, team, audience)
require.Nil(t, apiErr)
require.NotEmpty(t, tokenStr)

Expand All @@ -64,4 +65,11 @@ func TestGenerateVolumeContentToken_SetsTokidHeader(t *testing.T) {
tokidVal, ok := parsed.Header["tokid"].(string)
require.True(t, ok, "expected custom header 'tokid' to be present and a string")
require.Equal(t, config.SigningKeyName, tokidVal)

// Assert: the audience claim is the origin we passed, not the cluster ID.
claims, ok := parsed.Claims.(jwt.MapClaims)
require.True(t, ok)
aud, err := claims.GetAudience()
require.NoError(t, err)
require.Equal(t, jwt.ClaimStrings{audience}, aud)
}
10 changes: 5 additions & 5 deletions packages/api/internal/handlers/volume_token.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,14 @@ import (
"github.com/e2b-dev/infra/packages/api/internal/cfg"
"github.com/e2b-dev/infra/packages/auth/pkg/types"
"github.com/e2b-dev/infra/packages/db/queries"
"github.com/e2b-dev/infra/packages/shared/pkg/clusters"
)

var ErrVolumesTokenNotConfigured = errors.New("volumes content token signing is not supported by this deployment")

func generateVolumeContentToken(config cfg.VolumesTokenConfig, volume queries.Volume, team *types.Team) (string, *api.APIError) {
// generateVolumeContentToken mints the JWT the SDK presents to the volume
// content API. audience is the origin the token is valid for
// (`https://api.<domain>`, see volumeTokenAudience).
func generateVolumeContentToken(config cfg.VolumesTokenConfig, volume queries.Volume, team *types.Team, audience string) (string, *api.APIError) {
if !config.IsConfigured() {
return "", &api.APIError{
Err: ErrVolumesTokenNotConfigured,
Expand All @@ -27,14 +29,12 @@ func generateVolumeContentToken(config cfg.VolumesTokenConfig, volume queries.Vo
}
}

clusterID := clusters.WithClusterFallback(team.ClusterID)

now := time.Now()
expiration := now.Add(config.Duration)

claims := jwt.MapClaims{
// registered
"aud": clusterID.String(),
"aud": audience,
"exp": jwt.NewNumericDate(expiration),
"iat": jwt.NewNumericDate(now),
"iss": config.Issuer,
Expand Down
32 changes: 32 additions & 0 deletions packages/api/internal/handlers/volume_util.go
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,38 @@ func (a *APIStore) getVolume(c *gin.Context, volumeID string) (queries.Volume, *

var ErrClusterNotFound = errors.New("cluster not found")

// volumeContentDomain resolves the domain the SDK should use as the destination
// for volume content requests. Teams connected to a custom (BYOC) cluster get
// that cluster's domain, which the SDK uses in place of the default
// api.<E2B_DOMAIN> host. Teams on the default cluster get nil, leaving the SDK
// on its configured default.
func (a *APIStore) volumeContentDomain(team *types.Team) (*string, error) {
if team.ClusterID == nil {
return nil, nil
}

cluster, ok := a.clusters.GetClusterById(*team.ClusterID)
if !ok {
return nil, fmt.Errorf("%w: %s", ErrClusterNotFound, team.ClusterID.String())
}

return cluster.SandboxDomain, nil

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Reject BYOC clusters that have no volume domain

When a team is assigned to a cluster whose sandbox_proxy_domain is null—which remains valid under AdminClusterCreateRequest in spec/openapi-dashboard.yml:422-424—this returns (nil, nil), so omitempty removes domain and the SDK falls back to the default control-plane host. Volume-content requests are then routed to the wrong cluster; return a configuration error here or make the domain mandatory for BYOC clusters.

AGENTS.md reference: AGENTS.md:L22-L26

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is acceptable, it's deep into undefined behavior. If there's no client proxy domain, the whole cluster doesn't work, so this is the least of that cluster's problem.

}

// volumeTokenAudience builds the audience claim for a volume content token: the
// origin the SDK targets for content requests, i.e. `https://api.<domain>`. The
// domain is the team's BYOC cluster domain when set (as resolved by
// volumeContentDomain) and the deployment default otherwise, matching the host
// the token is actually presented to.
func (a *APIStore) volumeTokenAudience(domain *string) string {
host := a.config.DomainName
if domain != nil && *domain != "" {
host = *domain
}

return fmt.Sprintf("https://api.%s", host)
}

var ErrNoHealthyOrchestratorFound = errors.New("no healthy orchestrator found")

var ErrUnknownVolumeType = errors.New("unknown volume type")
Expand Down
66 changes: 66 additions & 0 deletions packages/api/internal/handlers/volume_util_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,78 @@ package handlers
import (
"testing"

"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"github.com/e2b-dev/infra/packages/api/internal/api"
"github.com/e2b-dev/infra/packages/api/internal/cfg"
"github.com/e2b-dev/infra/packages/api/internal/clusters"
"github.com/e2b-dev/infra/packages/api/internal/orchestrator/nodemanager"
"github.com/e2b-dev/infra/packages/auth/pkg/types"
authqueries "github.com/e2b-dev/infra/packages/db/pkg/auth/queries"
)

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

store := &APIStore{config: cfg.Config{DomainName: "e2b.app"}}

t.Run("falls back to the deployment domain", func(t *testing.T) {
t.Parallel()

assert.Equal(t, "https://api.e2b.app", store.volumeTokenAudience(nil))
})

t.Run("uses the BYOC domain when set", func(t *testing.T) {
t.Parallel()

domain := "custom.example.com"
assert.Equal(t, "https://api.custom.example.com", store.volumeTokenAudience(&domain))
})
}

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

clusterID := uuid.New()
domain := "custom.example.com"

store := &APIStore{
clusters: clusters.NewTestPool(clusters.NewTestCluster(clusterID, &domain)),
}

teamWith := func(id *uuid.UUID) *types.Team {
return &types.Team{Team: &authqueries.Team{ID: uuid.New(), ClusterID: id}}
}

t.Run("no cluster returns nil domain", func(t *testing.T) {
t.Parallel()

got, err := store.volumeContentDomain(teamWith(nil))
require.NoError(t, err)
assert.Nil(t, got)
})

t.Run("BYOC cluster returns its domain", func(t *testing.T) {
t.Parallel()

got, err := store.volumeContentDomain(teamWith(&clusterID))
require.NoError(t, err)
require.NotNil(t, got)
assert.Equal(t, domain, *got)
})

t.Run("unknown cluster returns error", func(t *testing.T) {
t.Parallel()

unknown := uuid.New()
got, err := store.volumeContentDomain(teamWith(&unknown))
require.ErrorIs(t, err, ErrClusterNotFound)
assert.Nil(t, got)
})
}

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

Expand Down
7 changes: 7 additions & 0 deletions spec/openapi.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2058,6 +2058,13 @@ components:
token:
type: string
description: Auth token to use for interacting with volume content
domain:
type: string
description: |
Domain to use as the destination for volume content requests,
replacing the default `api.<E2B_DOMAIN>`. Only returned when the
team is connected to a custom (BYOC) cluster; absent otherwise, in
which case the default domain is used.
Comment on lines +2064 to +2067

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Document the BYOC volume-content route

This contract redirects SDK volume-content traffic from the configured control-plane API host to the BYOC cluster domain, but docs/ARCHITECTURE.md:33-46 still models SDK traffic only through the main load balancer and API and contains no volume-content path. Add this cross-service routing flow to the architecture document as required for routing and topology changes.

AGENTS.md reference: AGENTS.md:L7-L9

Useful? React with 👍 / 👎.

required:
- volumeID
- name
Expand Down
6 changes: 6 additions & 0 deletions tests/integration/internal/api/generated.go

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

Loading