feat: add WatchAllWorkspaceBuilds endpoint for autostart scaletests (#22057)

This PR adds a `WatchAllWorkspaces` function with `watch-all-workspaces`
endpoint, which can be used to listen on a single global pubsub channel
for _all_ workspace build updates, and makes use of it in the autostart
scaletest.

This negates the need to use a workspace watch pubsub channel _per_
workspace, which has auth overhead associated with each call. This is
especially relevant in situations such as the autostart scaletest, where
we need to start/stop a set of workspaces before we can configure their
autostart config. The overhead associated with all the watch requests
skews the scaletest results and makes it harder to reason about the
performance of the autostart feature itself.

The autostart scaletest also no longer generates its own metrics nor
does it wait for all the workspaces to actually start via autostart. We
should update the scaletest dashboard after both PRs are merged to
measure autostart performance via the new metrics.



The new function/endpoint and its usage in the autostart scaletest are
gated behind an experiment feature flag, this is something we should
discuss whether we want to enable the endpoint in prod by default or
not. If so, we can remove the experiment.

---------

Signed-off-by: Callum Styan <callumstyan@gmail.com>
Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
Co-authored-by: Callum Styan <callum@coder.com>
This commit is contained in:
Callum Styan
2026-03-13 20:37:41 -07:00
committed by GitHub
co-authored by Claude Opus 4.5 Callum Styan
parent b492c42624
commit 36665e17b2
23 changed files with 1296 additions and 249 deletions
+32 -3
View File
@@ -1134,6 +1134,31 @@ const docTemplate = `{
}
}
},
"/experimental/watch-all-workspacebuilds": {
"get": {
"security": [
{
"CoderSessionToken": []
}
],
"produces": [
"application/json"
],
"tags": [
"Workspaces"
],
"summary": "Watch all workspace builds",
"operationId": "watch-all-workspace-builds",
"responses": {
"101": {
"description": "Switching Protocols"
}
},
"x-apidocgen": {
"skip": true
}
}
},
"/experiments": {
"get": {
"security": [
@@ -15226,7 +15251,8 @@ const docTemplate = `{
"web-push",
"oauth2",
"agents",
"mcp-server-http"
"mcp-server-http",
"workspace-build-updates"
],
"x-enum-comments": {
"ExperimentAgents": "Enables agent-powered chat functionality.",
@@ -15236,6 +15262,7 @@ const docTemplate = `{
"ExperimentNotifications": "Sends notifications via SMTP and webhooks following certain events.",
"ExperimentOAuth2": "Enables OAuth2 provider functionality.",
"ExperimentWebPush": "Enables web push notifications through the browser.",
"ExperimentWorkspaceBuildUpdates": "Enables publishing workspace build updates to the all builds pubsub channel.",
"ExperimentWorkspaceUsage": "Enables the new workspace usage tracking."
},
"x-enum-descriptions": [
@@ -15246,7 +15273,8 @@ const docTemplate = `{
"Enables web push notifications through the browser.",
"Enables OAuth2 provider functionality.",
"Enables agent-powered chat functionality.",
"Enables the MCP HTTP server functionality."
"Enables the MCP HTTP server functionality.",
"Enables publishing workspace build updates to the all builds pubsub channel."
],
"x-enum-varnames": [
"ExperimentExample",
@@ -15256,7 +15284,8 @@ const docTemplate = `{
"ExperimentWebPush",
"ExperimentOAuth2",
"ExperimentAgents",
"ExperimentMCPServerHTTP"
"ExperimentMCPServerHTTP",
"ExperimentWorkspaceBuildUpdates"
]
},
"codersdk.ExternalAPIKeyScopes": {
+28 -3
View File
@@ -983,6 +983,27 @@
}
}
},
"/experimental/watch-all-workspacebuilds": {
"get": {
"security": [
{
"CoderSessionToken": []
}
],
"produces": ["application/json"],
"tags": ["Workspaces"],
"summary": "Watch all workspace builds",
"operationId": "watch-all-workspace-builds",
"responses": {
"101": {
"description": "Switching Protocols"
}
},
"x-apidocgen": {
"skip": true
}
}
},
"/experiments": {
"get": {
"security": [
@@ -13744,7 +13765,8 @@
"web-push",
"oauth2",
"agents",
"mcp-server-http"
"mcp-server-http",
"workspace-build-updates"
],
"x-enum-comments": {
"ExperimentAgents": "Enables agent-powered chat functionality.",
@@ -13754,6 +13776,7 @@
"ExperimentNotifications": "Sends notifications via SMTP and webhooks following certain events.",
"ExperimentOAuth2": "Enables OAuth2 provider functionality.",
"ExperimentWebPush": "Enables web push notifications through the browser.",
"ExperimentWorkspaceBuildUpdates": "Enables publishing workspace build updates to the all builds pubsub channel.",
"ExperimentWorkspaceUsage": "Enables the new workspace usage tracking."
},
"x-enum-descriptions": [
@@ -13764,7 +13787,8 @@
"Enables web push notifications through the browser.",
"Enables OAuth2 provider functionality.",
"Enables agent-powered chat functionality.",
"Enables the MCP HTTP server functionality."
"Enables the MCP HTTP server functionality.",
"Enables publishing workspace build updates to the all builds pubsub channel."
],
"x-enum-varnames": [
"ExperimentExample",
@@ -13774,7 +13798,8 @@
"ExperimentWebPush",
"ExperimentOAuth2",
"ExperimentAgents",
"ExperimentMCPServerHTTP"
"ExperimentMCPServerHTTP",
"ExperimentWorkspaceBuildUpdates"
]
},
"codersdk.ExternalAPIKeyScopes": {
+7
View File
@@ -1204,6 +1204,13 @@ func New(options *Options) *API {
// MCP HTTP transport endpoint with mandatory authentication
r.Mount("/http", api.mcpHTTPHandler())
})
r.Route("/watch-all-workspacebuilds", func(r chi.Router) {
r.Use(
apiKeyMiddleware,
httpmw.RequireExperiment(api.Experiments, codersdk.ExperimentWorkspaceBuildUpdates),
)
r.Get("/", api.watchAllWorkspaceBuilds)
})
})
r.Route("/api/v2", func(r chi.Router) {
@@ -1289,6 +1289,21 @@ func (s *server) FailJob(ctx context.Context, failJob *proto.FailedJob) (*proto.
if err != nil {
return nil, xerrors.Errorf("publish workspace update: %w", err)
}
// Publish workspace build update to the all builds channel if the experiment is enabled.
if s.Experiments.Enabled(codersdk.ExperimentWorkspaceBuildUpdates) {
err = wspubsub.PublishWorkspaceBuildUpdate(ctx, s.Pubsub, codersdk.WorkspaceBuildUpdate{
WorkspaceID: workspace.ID,
WorkspaceName: workspace.Name,
BuildID: build.ID,
Transition: string(build.Transition),
JobStatus: string(database.ProvisionerJobStatusFailed),
BuildNumber: build.BuildNumber,
})
if err != nil {
s.Logger.Warn(ctx, "failed to publish workspace build update", slog.Error(err))
}
}
case *proto.FailedJob_TemplateImport_:
}
@@ -2489,6 +2504,21 @@ func (s *server) completeWorkspaceBuildJob(ctx context.Context, job database.Pro
return xerrors.Errorf("update workspace: %w", err)
}
// Publish workspace build update to the all builds channel if the experiment is enabled.
if s.Experiments.Enabled(codersdk.ExperimentWorkspaceBuildUpdates) {
err = wspubsub.PublishWorkspaceBuildUpdate(ctx, s.Pubsub, codersdk.WorkspaceBuildUpdate{
WorkspaceID: workspace.ID,
WorkspaceName: workspace.Name,
BuildID: workspaceBuild.ID,
Transition: string(workspaceBuild.Transition),
JobStatus: string(database.ProvisionerJobStatusSucceeded),
BuildNumber: workspaceBuild.BuildNumber,
})
if err != nil {
s.Logger.Warn(ctx, "failed to publish workspace build update", slog.Error(err))
}
}
if input.PrebuiltWorkspaceBuildStage == sdkproto.PrebuiltWorkspaceBuildStage_CLAIM {
s.Logger.Info(ctx, "workspace prebuild successfully claimed by user",
slog.F("workspace_id", workspace.ID))
+15
View File
@@ -766,6 +766,21 @@ func (api *API) patchCancelWorkspaceBuild(rw http.ResponseWriter, r *http.Reques
WorkspaceID: workspace.ID,
})
// Publish workspace build update to the all builds channel if the experiment is enabled.
if api.Experiments.Enabled(codersdk.ExperimentWorkspaceBuildUpdates) {
err = wspubsub.PublishWorkspaceBuildUpdate(ctx, api.Pubsub, codersdk.WorkspaceBuildUpdate{
WorkspaceID: workspace.ID,
WorkspaceName: workspace.Name,
BuildID: workspaceBuild.ID,
Transition: string(workspaceBuild.Transition),
JobStatus: string(database.ProvisionerJobStatusCanceled),
BuildNumber: workspaceBuild.BuildNumber,
})
if err != nil {
api.Logger.Warn(ctx, "failed to publish workspace build update", slog.Error(err))
}
}
httpapi.Write(ctx, rw, http.StatusOK, codersdk.Response{
Message: "Job has been marked as canceled...",
})
+74
View File
@@ -43,6 +43,8 @@ import (
"github.com/coder/coder/v2/coderd/wspubsub"
"github.com/coder/coder/v2/codersdk"
"github.com/coder/coder/v2/codersdk/agentsdk"
"github.com/coder/coder/v2/codersdk/wsjson"
"github.com/coder/websocket"
)
var (
@@ -2175,6 +2177,78 @@ func (api *API) watchWorkspace(
}
}
// @Summary Watch all workspace builds
// @ID watch-all-workspace-builds
// @Security CoderSessionToken
// @Produce json
// @Tags Workspaces
// @Success 101
// @Router /experimental/watch-all-workspacebuilds [get]
// @x-apidocgen {"skip": true}
func (api *API) watchAllWorkspaceBuilds(rw http.ResponseWriter, r *http.Request) {
ctx := r.Context()
// Buffer enough updates to avoid blocking the pubsub callback while we're
// accepting the WebSocket connection. Accepting the connection signals to
// the client that the server is subscribed and ready to forward events.
updates := make(chan codersdk.WorkspaceBuildUpdate, 256)
cancelSubscribe, err := api.Pubsub.SubscribeWithErr(wspubsub.AllWorkspaceEventChannel,
wspubsub.HandleWorkspaceBuildUpdate(
func(_ context.Context, update codersdk.WorkspaceBuildUpdate, err error) {
if err != nil {
api.Logger.Warn(ctx, "workspace build update subscription error", slog.Error(err))
return
}
select {
case updates <- update:
default:
api.Logger.Warn(ctx, "workspace build update dropped, client too slow")
}
}))
if err != nil {
httpapi.Write(ctx, rw, http.StatusInternalServerError, codersdk.Response{
Message: "Internal error subscribing to workspace build events.",
Detail: err.Error(),
})
return
}
defer cancelSubscribe()
conn, err := websocket.Accept(rw, r, nil)
if err != nil {
httpapi.Write(ctx, rw, http.StatusBadRequest, codersdk.Response{
Message: "Failed to accept WebSocket.",
Detail: err.Error(),
})
return
}
defer conn.Close(websocket.StatusNormalClosure, "done")
// CloseRead starts a goroutine to read and discard messages from the client,
// including Pong messages sent in response to our Ping heartbeats.
_ = conn.CloseRead(context.Background())
ctx, cancel := context.WithCancel(ctx)
go httpapi.HeartbeatClose(ctx, api.Logger, cancel, conn)
defer cancel()
enc := wsjson.NewEncoder[codersdk.WorkspaceBuildUpdate](conn, websocket.MessageText)
for {
select {
case <-ctx.Done():
return
case update, ok := <-updates:
if !ok {
return
}
if err := enc.Encode(update); err != nil {
return
}
}
}
}
// @Summary Get workspace timings by ID
// @ID get-workspace-timings-by-id
// @Security CoderSessionToken
+107
View File
@@ -3674,6 +3674,113 @@ func TestWorkspaceWatcher(t *testing.T) {
wait("second is for the build cancel", nil)
}
func TestWatchAllWorkspaceBuilds(t *testing.T) {
t.Parallel()
// Enable the workspace build updates experiment.
client, closer := coderdtest.NewWithProvisionerCloser(t, &coderdtest.Options{
IncludeProvisionerDaemon: true,
DeploymentValues: coderdtest.DeploymentValues(t, func(dv *codersdk.DeploymentValues) {
dv.Experiments = []string{string(codersdk.ExperimentWorkspaceBuildUpdates)}
}),
})
defer closer.Close()
user := coderdtest.CreateFirstUser(t, client)
// Create a simple template version.
version := coderdtest.CreateTemplateVersion(t, client, user.OrganizationID, &echo.Responses{
Parse: echo.ParseComplete,
ProvisionPlan: echo.PlanComplete,
ProvisionGraph: []*proto.Response{{
Type: &proto.Response_Graph{
Graph: &proto.GraphComplete{
Resources: []*proto.Resource{{
Name: "example",
Type: "aws_instance",
}},
},
},
}},
})
coderdtest.AwaitTemplateVersionJobCompleted(t, client, version.ID)
template := coderdtest.CreateTemplate(t, client, user.OrganizationID, version.ID)
ctx, cancel := context.WithTimeout(context.Background(), testutil.WaitLong)
defer cancel()
// Subscribe to all workspace build updates via SSE BEFORE creating workspaces
// so we can use it to wait for the initial builds.
decoder, err := client.WatchAllWorkspaceBuilds(ctx)
require.NoError(t, err)
defer decoder.Close()
updates := decoder.Chan()
logger := testutil.Logger(t).Named(t.Name())
// Helper to wait for a specific update.
waitForUpdate := func(event string, workspaceID uuid.UUID, expectedTransition, expectedStatus string) codersdk.WorkspaceBuildUpdate {
t.Helper()
for {
select {
case <-ctx.Done():
require.FailNow(t, "timed out waiting for event", event)
return codersdk.WorkspaceBuildUpdate{}
case update, ok := <-updates:
if !ok {
require.FailNow(t, "updates channel closed", event)
return codersdk.WorkspaceBuildUpdate{}
}
logger.Info(ctx, "received workspace build update",
slog.F("event", event),
slog.F("workspace_id", update.WorkspaceID),
slog.F("build_id", update.BuildID),
slog.F("transition", update.Transition),
slog.F("job_status", update.JobStatus),
slog.F("build_number", update.BuildNumber))
if update.WorkspaceID == workspaceID && update.Transition == expectedTransition && update.JobStatus == expectedStatus {
return update
}
// Keep waiting if this isn't the update we're looking for.
logger.Info(ctx, "skipping update, not matching expected",
slog.F("expected_workspace_id", workspaceID),
slog.F("expected_transition", expectedTransition),
slog.F("expected_status", expectedStatus))
}
}
}
// Create two workspaces and wait for their initial builds via the SSE channel.
workspace1 := coderdtest.CreateWorkspace(t, client, template.ID)
update := waitForUpdate("workspace1 initial build", workspace1.ID, "start", "succeeded")
require.Equal(t, workspace1.ID, update.WorkspaceID)
require.Equal(t, int32(1), update.BuildNumber)
workspace2 := coderdtest.CreateWorkspace(t, client, template.ID)
update = waitForUpdate("workspace2 initial build", workspace2.ID, "start", "succeeded")
require.Equal(t, workspace2.ID, update.WorkspaceID)
require.Equal(t, int32(1), update.BuildNumber)
// Stop workspace 1.
_ = coderdtest.CreateWorkspaceBuild(t, client, workspace1, database.WorkspaceTransitionStop)
update = waitForUpdate("workspace1 stop", workspace1.ID, "stop", "succeeded")
require.Equal(t, workspace1.ID, update.WorkspaceID)
// Stop workspace 2.
_ = coderdtest.CreateWorkspaceBuild(t, client, workspace2, database.WorkspaceTransitionStop)
update = waitForUpdate("workspace2 stop", workspace2.ID, "stop", "succeeded")
require.Equal(t, workspace2.ID, update.WorkspaceID)
// Start workspace 1 again.
_ = coderdtest.CreateWorkspaceBuild(t, client, workspace1, database.WorkspaceTransitionStart)
update = waitForUpdate("workspace1 start", workspace1.ID, "start", "succeeded")
require.Equal(t, workspace1.ID, update.WorkspaceID)
// Start workspace 2 again.
_ = coderdtest.CreateWorkspaceBuild(t, client, workspace2, database.WorkspaceTransitionStart)
update = waitForUpdate("workspace2 start", workspace2.ID, "start", "succeeded")
require.Equal(t, workspace2.ID, update.WorkspaceID)
}
func mustLocation(t *testing.T, location string) *time.Location {
t.Helper()
loc, err := time.LoadLocation(location)
+39
View File
@@ -7,8 +7,47 @@ import (
"github.com/google/uuid"
"golang.org/x/xerrors"
"github.com/coder/coder/v2/coderd/database/pubsub"
"github.com/coder/coder/v2/codersdk"
)
// AllWorkspaceEventChannel is a global channel that receives events for all
// workspaces. This is useful when you need to watch N workspaces without
// creating N separate subscriptions.
const AllWorkspaceEventChannel = "workspace_updates:all"
// HandleWorkspaceBuildUpdate wraps a callback to parse WorkspaceBuildUpdate
// messages from the pubsub.
func HandleWorkspaceBuildUpdate(cb func(ctx context.Context, payload codersdk.WorkspaceBuildUpdate, err error)) func(ctx context.Context, message []byte, err error) {
return func(ctx context.Context, message []byte, err error) {
if err != nil {
cb(ctx, codersdk.WorkspaceBuildUpdate{}, xerrors.Errorf("workspace build update pubsub: %w", err))
return
}
var payload codersdk.WorkspaceBuildUpdate
if err := json.Unmarshal(message, &payload); err != nil {
cb(ctx, codersdk.WorkspaceBuildUpdate{}, xerrors.Errorf("unmarshal workspace build update: %w", err))
return
}
cb(ctx, payload, nil)
}
}
// PublishWorkspaceBuildUpdate is a helper to publish a workspace build update
// to the AllWorkspaceEventChannel. This should be called when a build
// completes (succeeds, fails, or is canceled).
func PublishWorkspaceBuildUpdate(_ context.Context, ps pubsub.Pubsub, update codersdk.WorkspaceBuildUpdate) error {
msg, err := json.Marshal(update)
if err != nil {
return xerrors.Errorf("marshal workspace build update: %w", err)
}
if err := ps.Publish(AllWorkspaceEventChannel, msg); err != nil {
return xerrors.Errorf("publish workspace build update: %w", err)
}
return nil
}
// WorkspaceEventChannel can be used to subscribe to events for
// workspaces owned by the provided user ID.
func WorkspaceEventChannel(ownerID uuid.UUID) string {