chore: modify task status scaletest to use Agent API dRPC (#22356)

relates to #21335

Modifies our taskstatus scaletest load generator to use the dRPC connection to mimic what an actual running Task would do via the MCP server (c.f. PRs below this one in the stack).

Disclosure: I used AI to generate large portions of this PR, but hand-reviewed and tweaked.
This commit is contained in:
Spike Curtis
2026-03-04 22:12:35 +04:00
committed by GitHub
parent 8327e1f65f
commit fda181bb26
3 changed files with 111 additions and 79 deletions
+49 -28
View File
@@ -9,6 +9,7 @@ import (
"golang.org/x/xerrors" "golang.org/x/xerrors"
"cdr.dev/slog/v3" "cdr.dev/slog/v3"
agentproto "github.com/coder/coder/v2/agent/proto"
"github.com/coder/coder/v2/codersdk" "github.com/coder/coder/v2/codersdk"
"github.com/coder/coder/v2/codersdk/agentsdk" "github.com/coder/coder/v2/codersdk/agentsdk"
"github.com/coder/quartz" "github.com/coder/quartz"
@@ -41,15 +42,20 @@ type client interface {
initialize(logger slog.Logger) initialize(logger slog.Logger)
} }
// appStatusPatcher abstracts the details of using agentsdk.Client for updating app status. // appStatusUpdater abstracts the details of updating app status via the
// This interface is separate from client because it requires an agent token which is only // Agent dRPC API. This interface is separate from client because it
// available after creating an external workspace. // requires an agent token which is only available after creating an
type appStatusPatcher interface { // external workspace.
// patchAppStatus updates the status of a workspace app. type appStatusUpdater interface {
patchAppStatus(ctx context.Context, req agentsdk.PatchAppStatus) error // updateAppStatus sends a status update for a workspace app.
updateAppStatus(ctx context.Context, req *agentproto.UpdateAppStatusRequest) error
// initialize sets up the patcher with the provided logger and agent token. // initialize establishes the dRPC connection using the provided
initialize(logger slog.Logger, agentToken string) // agent token. Must be called before updateAppStatus.
initialize(ctx context.Context, logger slog.Logger, agentToken string) error
// close cleanly shuts down the underlying dRPC connection.
close() error
} }
// sdkClient is the concrete implementation of the client interface using // sdkClient is the concrete implementation of the client interface using
@@ -103,42 +109,57 @@ func (c *sdkClient) initialize(logger slog.Logger) {
c.coderClient.SetLogBodies(true) c.coderClient.SetLogBodies(true)
} }
// sdkAppStatusPatcher is the concrete implementation of the appStatusPatcher interface // sdkAppStatusUpdater is the concrete implementation of the
// using agentsdk.Client. // appStatusUpdater interface. It dials the Agent dRPC endpoint once
type sdkAppStatusPatcher struct { // during initialize and reuses the connection for all subsequent
agentClient *agentsdk.Client // UpdateAppStatus calls.
url *url.URL type sdkAppStatusUpdater struct {
httpClient *http.Client drpcClient agentproto.DRPCAgentClient28
url *url.URL
httpClient *http.Client
} }
// newAppStatusPatcher creates a new appStatusPatcher implementation. // newAppStatusUpdater creates a new appStatusUpdater implementation.
func newAppStatusPatcher(client *codersdk.Client) appStatusPatcher { func newAppStatusUpdater(client *codersdk.Client) appStatusUpdater {
return &sdkAppStatusPatcher{ return &sdkAppStatusUpdater{
url: client.URL, url: client.URL,
httpClient: client.HTTPClient, httpClient: client.HTTPClient,
} }
} }
func (p *sdkAppStatusPatcher) patchAppStatus(ctx context.Context, req agentsdk.PatchAppStatus) error { func (u *sdkAppStatusUpdater) updateAppStatus(ctx context.Context, req *agentproto.UpdateAppStatusRequest) error {
if p.agentClient == nil { if u.drpcClient == nil {
panic("agentClient not initialized - call initialize first") return xerrors.New("dRPC client not initialized - call initialize first")
} }
return p.agentClient.PatchAppStatus(ctx, req) _, err := u.drpcClient.UpdateAppStatus(ctx, req)
return err
} }
func (p *sdkAppStatusPatcher) initialize(logger slog.Logger, agentToken string) { func (u *sdkAppStatusUpdater) close() error {
// Create and configure the agent client with the provided token if u.drpcClient == nil {
p.agentClient = agentsdk.New( return nil
p.url, }
return u.drpcClient.DRPCConn().Close()
}
func (u *sdkAppStatusUpdater) initialize(ctx context.Context, logger slog.Logger, agentToken string) error {
agentClient := agentsdk.New(
u.url,
agentsdk.WithFixedToken(agentToken), agentsdk.WithFixedToken(agentToken),
codersdk.WithHTTPClient(p.httpClient), codersdk.WithHTTPClient(u.httpClient),
codersdk.WithLogger(logger), codersdk.WithLogger(logger),
codersdk.WithLogBodies(), codersdk.WithLogBodies(),
) )
drpcClient, _, err := agentClient.ConnectRPC28WithRole(ctx, "")
if err != nil {
return xerrors.Errorf("connect to agent dRPC endpoint: %w", err)
}
u.drpcClient = drpcClient
return nil
} }
// Ensure sdkClient implements the client interface. // Ensure sdkClient implements the client interface.
var _ client = (*sdkClient)(nil) var _ client = (*sdkClient)(nil)
// Ensure sdkAppStatusPatcher implements the appStatusPatcher interface. // Ensure sdkAppStatusUpdater implements the appStatusUpdater interface.
var _ appStatusPatcher = (*sdkAppStatusPatcher)(nil) var _ appStatusUpdater = (*sdkAppStatusUpdater)(nil)
+18 -10
View File
@@ -13,8 +13,8 @@ import (
"cdr.dev/slog/v3" "cdr.dev/slog/v3"
"cdr.dev/slog/v3/sloggers/sloghuman" "cdr.dev/slog/v3/sloggers/sloghuman"
agentproto "github.com/coder/coder/v2/agent/proto"
"github.com/coder/coder/v2/codersdk" "github.com/coder/coder/v2/codersdk"
"github.com/coder/coder/v2/codersdk/agentsdk"
"github.com/coder/coder/v2/scaletest/harness" "github.com/coder/coder/v2/scaletest/harness"
"github.com/coder/coder/v2/scaletest/loadtestutil" "github.com/coder/coder/v2/scaletest/loadtestutil"
"github.com/coder/quartz" "github.com/coder/quartz"
@@ -30,7 +30,7 @@ type createExternalWorkspaceResult struct {
type Runner struct { type Runner struct {
client client client client
patcher appStatusPatcher updater appStatusUpdater
cfg Config cfg Config
logger slog.Logger logger slog.Logger
@@ -55,7 +55,7 @@ var (
func NewRunner(coderClient *codersdk.Client, cfg Config) *Runner { func NewRunner(coderClient *codersdk.Client, cfg Config) *Runner {
return &Runner{ return &Runner{
client: newClient(coderClient), client: newClient(coderClient),
patcher: newAppStatusPatcher(coderClient), updater: newAppStatusUpdater(coderClient),
cfg: cfg, cfg: cfg,
clock: quartz.NewReal(), clock: quartz.NewReal(),
reportTimes: make(map[int]time.Time), reportTimes: make(map[int]time.Time),
@@ -96,9 +96,17 @@ func (r *Runner) Run(ctx context.Context, name string, logs io.Writer) error {
r.workspaceID = result.workspaceID r.workspaceID = result.workspaceID
r.logger.Info(ctx, "created external workspace", slog.F("workspace_id", r.workspaceID)) r.logger.Info(ctx, "created external workspace", slog.F("workspace_id", r.workspaceID))
// Initialize the patcher with the agent token // Establish the dRPC connection using the agent token.
r.patcher.initialize(r.logger, result.agentToken) if err := r.updater.initialize(ctx, r.logger, result.agentToken); err != nil {
r.logger.Info(ctx, "initialized app status patcher with agent token") r.cfg.Metrics.ReportTaskStatusErrorsTotal.WithLabelValues(r.cfg.MetricLabelValues...).Inc()
return xerrors.Errorf("initialize app status updater: %w", err)
}
defer func() {
if err := r.updater.close(); err != nil {
r.logger.Error(ctx, "failed to close app status updater", slog.Error(err))
}
}()
r.logger.Info(ctx, "initialized app status updater with agent token")
workspaceUpdatesCtx, cancelWorkspaceUpdates := context.WithCancel(ctx) workspaceUpdatesCtx, cancelWorkspaceUpdates := context.WithCancel(ctx)
defer cancelWorkspaceUpdates() defer cancelWorkspaceUpdates()
@@ -227,11 +235,11 @@ func (r *Runner) reportTaskStatus(ctx context.Context) error {
} }
r.mu.Unlock() r.mu.Unlock()
err := r.patcher.patchAppStatus(ctx, agentsdk.PatchAppStatus{ err := r.updater.updateAppStatus(ctx, &agentproto.UpdateAppStatusRequest{
AppSlug: r.cfg.AppSlug, Slug: r.cfg.AppSlug,
Message: statusUpdatePrefix + strconv.Itoa(msgNo), Message: statusUpdatePrefix + strconv.Itoa(msgNo),
State: codersdk.WorkspaceAppStatusStateWorking, State: agentproto.UpdateAppStatusRequest_WORKING,
URI: "https://example.com/example-status/", Uri: "https://example.com/example-status/",
}) })
if err != nil { if err != nil {
r.logger.Error(ctx, "failed to report task status", slog.Error(err)) r.logger.Error(ctx, "failed to report task status", slog.Error(err))
+44 -41
View File
@@ -15,8 +15,8 @@ import (
"cdr.dev/slog/v3" "cdr.dev/slog/v3"
"cdr.dev/slog/v3/sloggers/sloghuman" "cdr.dev/slog/v3/sloggers/sloghuman"
agentproto "github.com/coder/coder/v2/agent/proto"
"github.com/coder/coder/v2/codersdk" "github.com/coder/coder/v2/codersdk"
"github.com/coder/coder/v2/codersdk/agentsdk"
"github.com/coder/coder/v2/testutil" "github.com/coder/coder/v2/testutil"
"github.com/coder/quartz" "github.com/coder/quartz"
) )
@@ -115,43 +115,46 @@ func (m *fakeClient) deleteWorkspace(ctx context.Context, workspaceID uuid.UUID)
return nil return nil
} }
// fakeAppStatusPatcher implements the appStatusPatcher interface for testing // fakeAppStatusUpdater implements the appStatusUpdater interface for testing.
type fakeAppStatusPatcher struct { type fakeAppStatusUpdater struct {
t *testing.T t *testing.T
logger slog.Logger logger slog.Logger
agentToken string agentToken string
// Channels for controlling the behavior // Channels for controlling the behavior
patchStatusCalls chan agentsdk.PatchAppStatus updateStatusCalls chan *agentproto.UpdateAppStatusRequest
patchStatusErrors chan error updateStatusErrors chan error
} }
func newFakeAppStatusPatcher(t *testing.T) *fakeAppStatusPatcher { func newFakeAppStatusUpdater(t *testing.T) *fakeAppStatusUpdater {
return &fakeAppStatusPatcher{ return &fakeAppStatusUpdater{
t: t, t: t,
patchStatusCalls: make(chan agentsdk.PatchAppStatus), updateStatusCalls: make(chan *agentproto.UpdateAppStatusRequest),
patchStatusErrors: make(chan error, 1), updateStatusErrors: make(chan error, 1),
} }
} }
func (p *fakeAppStatusPatcher) initialize(logger slog.Logger, agentToken string) { func (u *fakeAppStatusUpdater) initialize(_ context.Context, logger slog.Logger, agentToken string) error {
p.logger = logger u.logger = logger
p.agentToken = agentToken u.agentToken = agentToken
return nil
} }
func (p *fakeAppStatusPatcher) patchAppStatus(ctx context.Context, req agentsdk.PatchAppStatus) error { func (*fakeAppStatusUpdater) close() error {
assert.NotEmpty(p.t, p.agentToken) return nil
p.logger.Debug(ctx, "called fake PatchAppStatus", slog.F("req", req)) }
// Send the request to the channel so tests can verify it
func (u *fakeAppStatusUpdater) updateAppStatus(ctx context.Context, req *agentproto.UpdateAppStatusRequest) error {
assert.NotEmpty(u.t, u.agentToken)
u.logger.Debug(ctx, "called fake UpdateAppStatus", slog.F("req", req))
select { select {
case p.patchStatusCalls <- req: case u.updateStatusCalls <- req:
case <-ctx.Done(): case <-ctx.Done():
return ctx.Err() return ctx.Err()
} }
// Check if there's an error to return
select { select {
case err := <-p.patchStatusErrors: case err := <-u.updateStatusErrors:
return err return err
default: default:
return nil return nil
@@ -165,7 +168,7 @@ func TestRunner_Run(t *testing.T) {
mClock := quartz.NewMock(t) mClock := quartz.NewMock(t)
fClient := newFakeClient(t) fClient := newFakeClient(t)
fPatcher := newFakeAppStatusPatcher(t) fUpdater := newFakeAppStatusUpdater(t)
templateID := uuid.UUID{5, 6, 7, 8} templateID := uuid.UUID{5, 6, 7, 8}
workspaceName := "test-workspace" workspaceName := "test-workspace"
appSlug := "test-app" appSlug := "test-app"
@@ -190,7 +193,7 @@ func TestRunner_Run(t *testing.T) {
} }
runner := &Runner{ runner := &Runner{
client: fClient, client: fClient,
patcher: fPatcher, updater: fUpdater,
cfg: cfg, cfg: cfg,
clock: mClock, clock: mClock,
reportTimes: make(map[int]time.Time), reportTimes: make(map[int]time.Time),
@@ -224,17 +227,17 @@ func TestRunner_Run(t *testing.T) {
// Wait for the initial TickerFunc call before advancing time, otherwise our ticks will be off. // Wait for the initial TickerFunc call before advancing time, otherwise our ticks will be off.
reportTickerTrap.MustWait(ctx).MustRelease(ctx) reportTickerTrap.MustWait(ctx).MustRelease(ctx)
// at this point, the patcher must be initialized // at this point, the updater must be initialized
require.Equal(t, testAgentToken, fPatcher.agentToken) require.Equal(t, testAgentToken, fUpdater.agentToken)
updateDelay := time.Duration(0) updateDelay := time.Duration(0)
for i := 0; i < 4; i++ { for i := 0; i < 4; i++ {
tickWaiter := mClock.Advance((10 * time.Second) - updateDelay) tickWaiter := mClock.Advance((10 * time.Second) - updateDelay)
patchCall := testutil.RequireReceive(ctx, t, fPatcher.patchStatusCalls) updateCall := testutil.RequireReceive(ctx, t, fUpdater.updateStatusCalls)
require.Equal(t, appSlug, patchCall.AppSlug) require.Equal(t, appSlug, updateCall.Slug)
require.Equal(t, fmt.Sprintf("scaletest status update:%d", i), patchCall.Message) require.Equal(t, fmt.Sprintf("scaletest status update:%d", i), updateCall.Message)
require.Equal(t, codersdk.WorkspaceAppStatusStateWorking, patchCall.State) require.Equal(t, agentproto.UpdateAppStatusRequest_WORKING, updateCall.State)
tickWaiter.MustWait(ctx) tickWaiter.MustWait(ctx)
// Send workspace update 1, 2, 3, or 4 seconds after the report // Send workspace update 1, 2, 3, or 4 seconds after the report
@@ -287,7 +290,7 @@ func TestRunner_RunMissedUpdate(t *testing.T) {
mClock := quartz.NewMock(t) mClock := quartz.NewMock(t)
fClient := newFakeClient(t) fClient := newFakeClient(t)
fPatcher := newFakeAppStatusPatcher(t) fUpdater := newFakeAppStatusUpdater(t)
templateID := uuid.UUID{5, 6, 7, 8} templateID := uuid.UUID{5, 6, 7, 8}
workspaceName := "test-workspace" workspaceName := "test-workspace"
appSlug := "test-app" appSlug := "test-app"
@@ -312,7 +315,7 @@ func TestRunner_RunMissedUpdate(t *testing.T) {
} }
runner := &Runner{ runner := &Runner{
client: fClient, client: fClient,
patcher: fPatcher, updater: fUpdater,
cfg: cfg, cfg: cfg,
clock: mClock, clock: mClock,
reportTimes: make(map[int]time.Time), reportTimes: make(map[int]time.Time),
@@ -349,10 +352,10 @@ func TestRunner_RunMissedUpdate(t *testing.T) {
updateDelay := time.Duration(0) updateDelay := time.Duration(0)
for i := 0; i < 4; i++ { for i := 0; i < 4; i++ {
tickWaiter := mClock.Advance((10 * time.Second) - updateDelay) tickWaiter := mClock.Advance((10 * time.Second) - updateDelay)
patchCall := testutil.RequireReceive(testCtx, t, fPatcher.patchStatusCalls) updateCall := testutil.RequireReceive(testCtx, t, fUpdater.updateStatusCalls)
require.Equal(t, appSlug, patchCall.AppSlug) require.Equal(t, appSlug, updateCall.Slug)
require.Equal(t, fmt.Sprintf("scaletest status update:%d", i), patchCall.Message) require.Equal(t, fmt.Sprintf("scaletest status update:%d", i), updateCall.Message)
require.Equal(t, codersdk.WorkspaceAppStatusStateWorking, patchCall.State) require.Equal(t, agentproto.UpdateAppStatusRequest_WORKING, updateCall.State)
tickWaiter.MustWait(testCtx) tickWaiter.MustWait(testCtx)
// Send workspace update 1, 2, 3, or 4 seconds after the report // Send workspace update 1, 2, 3, or 4 seconds after the report
@@ -412,7 +415,7 @@ func TestRunner_Run_WithErrors(t *testing.T) {
mClock := quartz.NewMock(t) mClock := quartz.NewMock(t)
fClient := newFakeClient(t) fClient := newFakeClient(t)
fPatcher := newFakeAppStatusPatcher(t) fUpdater := newFakeAppStatusUpdater(t)
templateID := uuid.UUID{5, 6, 7, 8} templateID := uuid.UUID{5, 6, 7, 8}
workspaceName := "test-workspace" workspaceName := "test-workspace"
appSlug := "test-app" appSlug := "test-app"
@@ -437,7 +440,7 @@ func TestRunner_Run_WithErrors(t *testing.T) {
} }
runner := &Runner{ runner := &Runner{
client: fClient, client: fClient,
patcher: fPatcher, updater: fUpdater,
cfg: cfg, cfg: cfg,
clock: mClock, clock: mClock,
reportTimes: make(map[int]time.Time), reportTimes: make(map[int]time.Time),
@@ -467,8 +470,8 @@ func TestRunner_Run_WithErrors(t *testing.T) {
for i := 0; i < 4; i++ { for i := 0; i < 4; i++ {
tickWaiter := mClock.Advance(10 * time.Second) tickWaiter := mClock.Advance(10 * time.Second)
testutil.RequireSend(testCtx, t, fPatcher.patchStatusErrors, xerrors.New("a bad thing happened")) testutil.RequireSend(testCtx, t, fUpdater.updateStatusErrors, xerrors.New("a bad thing happened"))
_ = testutil.RequireReceive(testCtx, t, fPatcher.patchStatusCalls) _ = testutil.RequireReceive(testCtx, t, fUpdater.updateStatusCalls)
tickWaiter.MustWait(testCtx) tickWaiter.MustWait(testCtx)
} }
@@ -513,7 +516,7 @@ func TestRunner_Run_BuildFailed(t *testing.T) {
mClock := quartz.NewMock(t) mClock := quartz.NewMock(t)
fClient := newFakeClient(t) fClient := newFakeClient(t)
fPatcher := newFakeAppStatusPatcher(t) fUpdater := newFakeAppStatusUpdater(t)
templateID := uuid.UUID{5, 6, 7, 8} templateID := uuid.UUID{5, 6, 7, 8}
workspaceName := "test-workspace" workspaceName := "test-workspace"
appSlug := "test-app" appSlug := "test-app"
@@ -538,7 +541,7 @@ func TestRunner_Run_BuildFailed(t *testing.T) {
} }
runner := &Runner{ runner := &Runner{
client: fClient, client: fClient,
patcher: fPatcher, updater: fUpdater,
cfg: cfg, cfg: cfg,
clock: mClock, clock: mClock,
reportTimes: make(map[int]time.Time), reportTimes: make(map[int]time.Time),
@@ -671,7 +674,7 @@ func TestRunner_Cleanup(t *testing.T) {
runner := &Runner{ runner := &Runner{
client: fakeClient, client: fakeClient,
patcher: newFakeAppStatusPatcher(t), updater: newFakeAppStatusUpdater(t),
cfg: cfg, cfg: cfg,
clock: quartz.NewMock(t), clock: quartz.NewMock(t),
} }