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
18 changes: 17 additions & 1 deletion coderd/coderd.go
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,7 @@ import (
"github.com/coder/coder/v2/coderd/workspaceconnwatcher"
"github.com/coder/coder/v2/coderd/workspacestats"
"github.com/coder/coder/v2/coderd/wsbuilder"
"github.com/coder/coder/v2/coderd/wsbuildorchestrator"
"github.com/coder/coder/v2/coderd/x/chatd"
"github.com/coder/coder/v2/coderd/x/chatd/chatprovider"
"github.com/coder/coder/v2/coderd/x/chatd/mcpclient"
Expand Down Expand Up @@ -950,6 +951,19 @@ func New(options *Options) *API {

api.workspaceAgentConnWatcher = workspaceconnwatcher.New(api.ctx, options.Logger, options.Pubsub, options.Database)

api.workspaceBuildOrchestrator = wsbuildorchestrator.New(wsbuildorchestrator.Options{
Logger: options.Logger,
Database: options.Database,
Pubsub: options.Pubsub,
FileCache: api.FileCache,
BuildUsageChecker: api.BuildUsageChecker,
DeploymentValues: options.DeploymentValues,
Experiments: api.Experiments,
BuilderMetrics: options.WorkspaceBuilderMetrics,
Clock: api.Clock,
})
api.workspaceBuildOrchestrator.Start(api.ctx)

apiKeyMiddleware := httpmw.ExtractAPIKeyMW(httpmw.ExtractAPIKeyConfig{
DB: options.Database,
ActivateDormantUser: ActivateDormantUser(options.Logger, &api.Auditor, options.Database),
Expand Down Expand Up @@ -2293,7 +2307,8 @@ type API struct {
// profiler is process-global, so concurrent collections would fail.
ProfileCollecting atomic.Bool

workspaceAgentConnWatcher *workspaceconnwatcher.Watcher
workspaceAgentConnWatcher *workspaceconnwatcher.Watcher
workspaceBuildOrchestrator *wsbuildorchestrator.Orchestrator
}

// Close waits for all WebSocket connections to drain before returning.
Expand Down Expand Up @@ -2358,6 +2373,7 @@ func (api *API) Close() error {
_ = api.AppEncryptionKeyCache.Close()
_ = api.UpdatesProvider.Close()
api.workspaceAgentConnWatcher.Close()
api.workspaceBuildOrchestrator.Close()

if current := api.PrebuildsReconciler.Load(); current != nil {
ctx, giveUp := context.WithTimeoutCause(context.Background(), time.Second*30, xerrors.New("gave up waiting for reconciler to stop before shutdown"))
Expand Down
16 changes: 12 additions & 4 deletions coderd/database/dbauthz/dbauthz.go
Original file line number Diff line number Diff line change
Expand Up @@ -671,10 +671,11 @@ var (
Identifier: rbac.RoleIdentifier{Name: "dbpurge"},
DisplayName: "DB Purge Daemon",
Site: rbac.Permissions(map[string][]policy.Action{
rbac.ResourceSystem.Type: {policy.ActionDelete},
rbac.ResourceNotificationMessage.Type: {policy.ActionDelete},
rbac.ResourceApiKey.Type: {policy.ActionDelete},
rbac.ResourceAibridgeInterception.Type: {policy.ActionDelete},
rbac.ResourceSystem.Type: {policy.ActionDelete},
rbac.ResourceNotificationMessage.Type: {policy.ActionDelete},
rbac.ResourceApiKey.Type: {policy.ActionDelete},
rbac.ResourceAibridgeInterception.Type: {policy.ActionDelete},
rbac.ResourceWorkspaceBuildOrchestration.Type: {policy.ActionDelete},
// Chat auto-archive sets archived=true on inactive chats.
rbac.ResourceChat.Type: {policy.ActionRead, policy.ActionUpdate},
// Purge old boundary logs past the retention period.
Expand Down Expand Up @@ -2389,6 +2390,13 @@ func (q *querier) DeleteOldWorkspaceAgentStats(ctx context.Context) error {
return q.db.DeleteOldWorkspaceAgentStats(ctx)
}

func (q *querier) DeleteOldWorkspaceBuildOrchestrations(ctx context.Context, arg database.DeleteOldWorkspaceBuildOrchestrationsParams) (int64, error) {
if err := q.authorizeContext(ctx, policy.ActionDelete, rbac.ResourceWorkspaceBuildOrchestration.AnyOrganization()); err != nil {
return 0, err
}
return q.db.DeleteOldWorkspaceBuildOrchestrations(ctx, arg)
}

func (q *querier) DeleteOrganizationMember(ctx context.Context, arg database.DeleteOrganizationMemberParams) error {
return deleteQ[database.OrganizationMember](q.log, q.auth, func(ctx context.Context, arg database.DeleteOrganizationMemberParams) (database.OrganizationMember, error) {
member, err := database.ExpectOne(q.OrganizationMembers(ctx, database.OrganizationMembersParams{
Expand Down
8 changes: 8 additions & 0 deletions coderd/database/dbauthz/dbauthz_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4156,6 +4156,14 @@ func (s *MethodTestSuite) TestWorkspace() {
Asserts(rbac.ResourceWorkspaceBuildOrchestration.AnyOrganization(), policy.ActionUpdate).
Returns(orchestration)
}))
s.Run("DeleteOldWorkspaceBuildOrchestrations", s.Mocked(func(dbm *dbmock.MockStore, _ *gofakeit.Faker, check *expects) {
arg := database.DeleteOldWorkspaceBuildOrchestrationsParams{
BeforeTime: dbtime.Now(),
LimitCount: 100,
}
dbm.EXPECT().DeleteOldWorkspaceBuildOrchestrations(gomock.Any(), arg).Return(int64(0), nil).AnyTimes()
check.Args(arg).Asserts(rbac.ResourceWorkspaceBuildOrchestration.AnyOrganization(), policy.ActionDelete)
}))
s.Run("Start/InsertWorkspaceBuildParameters", s.Mocked(func(dbm *dbmock.MockStore, faker *gofakeit.Faker, check *expects) {
w := testutil.Fake(s.T(), faker, database.Workspace{})
b := testutil.Fake(s.T(), faker, database.WorkspaceBuild{
Expand Down
8 changes: 8 additions & 0 deletions coderd/database/dbmetrics/querymetrics.go

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

15 changes: 15 additions & 0 deletions coderd/database/dbmock/dbmock.go

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

16 changes: 16 additions & 0 deletions coderd/database/dbpurge/dbpurge.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,11 @@ const (
// long enough to cover the maximum interval of a heartbeat event (currently
// 1 hour) plus some buffer.
maxTelemetryHeartbeatAge = 24 * time.Hour
// Operational handoff state; terminal rows are kept for debugging, then
// purged.
workspaceBuildOrchestrationTerminalRetention = 24 * time.Hour
// Batch size for workspace build orchestration deletion.
workspaceBuildOrchestrationsBatchSize = 10000
// Chat and chat file batch sizes stay smaller than audit/connection
// log batches because chat_files rows carry bytea blobs.
chatsBatchSize = 1000
Expand Down Expand Up @@ -275,6 +280,15 @@ func (i *instance) purgeTick(ctx context.Context, db database.Store, start time.
}
}

deleteOldWorkspaceBuildOrchestrationsBefore := start.Add(-workspaceBuildOrchestrationTerminalRetention)
purgedWorkspaceBuildOrchestrations, err := tx.DeleteOldWorkspaceBuildOrchestrations(ctx, database.DeleteOldWorkspaceBuildOrchestrationsParams{
BeforeTime: deleteOldWorkspaceBuildOrchestrationsBefore,
LimitCount: workspaceBuildOrchestrationsBatchSize,
})
if err != nil {
return xerrors.Errorf("failed to delete old workspace build orchestrations: %w", err)
}

var purgedChats, purgedChatFiles, purgedChatDebugRuns int64
if purgeChats {
purgedChats, purgedChatFiles, err = i.purgeChatsInTx(ctx, tx, start, chatRetentionDays)
Expand Down Expand Up @@ -304,6 +318,7 @@ func (i *instance) purgeTick(ctx context.Context, db database.Store, start time.
slog.F("audit_logs", purgedAuditLogs),
slog.F("boundary_logs", purgedBoundaryLogs),
slog.F("boundary_sessions", purgedBoundarySessions),
slog.F("workspace_build_orchestrations", purgedWorkspaceBuildOrchestrations),
slog.F("chats", purgedChats),
slog.F("chat_files", purgedChatFiles),
slog.F("chat_debug_runs", purgedChatDebugRuns),
Expand All @@ -318,6 +333,7 @@ func (i *instance) purgeTick(ctx context.Context, db database.Store, start time.
i.recordsPurged.WithLabelValues("audit_logs").Add(float64(purgedAuditLogs))
i.recordsPurged.WithLabelValues("boundary_logs").Add(float64(purgedBoundaryLogs))
i.recordsPurged.WithLabelValues("boundary_sessions").Add(float64(purgedBoundarySessions))
i.recordsPurged.WithLabelValues("workspace_build_orchestrations").Add(float64(purgedWorkspaceBuildOrchestrations))
i.recordsPurged.WithLabelValues("chats").Add(float64(purgedChats))
i.recordsPurged.WithLabelValues("chat_debug_runs").Add(float64(purgedChatDebugRuns))
i.recordsPurged.WithLabelValues("chat_files").Add(float64(purgedChatFiles))
Expand Down
147 changes: 147 additions & 0 deletions coderd/database/dbpurge/dbpurge_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,11 @@ func TestMetrics(t *testing.T) {
})
require.GreaterOrEqual(t, auditLogs, 0)

workspaceBuildOrchestrations := promhelp.CounterValue(t, reg, "coderd_dbpurge_records_purged_total", prometheus.Labels{
"record_type": "workspace_build_orchestrations",
})
require.GreaterOrEqual(t, workspaceBuildOrchestrations, 0)

chats := promhelp.CounterValue(t, reg, "coderd_dbpurge_records_purged_total", prometheus.Labels{
"record_type": "chats",
})
Expand Down Expand Up @@ -247,6 +252,7 @@ func TestMetrics(t *testing.T) {
mDB.EXPECT().DeleteOldNotificationMessages(gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().ExpirePrebuildsAPIKeys(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().DeleteOldTelemetryLocks(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().DeleteOldWorkspaceBuildOrchestrations(gomock.Any(), gomock.Any()).Return(int64(0), nil).AnyTimes()
mDB.EXPECT().DeleteOldAuditLogConnectionEvents(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().DeleteOldChatDebugRuns(gomock.Any(), gomock.AssignableToTypeOf(database.DeleteOldChatDebugRunsParams{})).Return(int64(0), nil).MinTimes(1)
mDB.EXPECT().InTx(gomock.Any(), database.DefaultTXOptions().WithID("db_purge")).
Expand Down Expand Up @@ -297,6 +303,7 @@ func TestMetrics(t *testing.T) {
mDB.EXPECT().DeleteOldNotificationMessages(gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().ExpirePrebuildsAPIKeys(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().DeleteOldTelemetryLocks(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().DeleteOldWorkspaceBuildOrchestrations(gomock.Any(), gomock.Any()).Return(int64(0), nil).AnyTimes()
mDB.EXPECT().DeleteOldAuditLogConnectionEvents(gomock.Any(), gomock.Any()).Return(nil).AnyTimes()
mDB.EXPECT().DeleteOldChats(gomock.Any(), gomock.AssignableToTypeOf(database.DeleteOldChatsParams{})).Return(int64(0), nil).MinTimes(1)
mDB.EXPECT().DeleteOldChatFiles(gomock.Any(), gomock.AssignableToTypeOf(database.DeleteOldChatFilesParams{})).Return(int64(0), nil).MinTimes(1)
Expand Down Expand Up @@ -1117,6 +1124,146 @@ func TestDeleteOldTelemetryHeartbeats(t *testing.T) {
}, testutil.WaitShort, testutil.IntervalFast, "it should delete old telemetry heartbeats")
}

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

ctx := testutil.Context(t, testutil.WaitShort)
db, _, rawDB := dbtestutil.NewDBWithSQLDB(t)

org := dbgen.Organization(t, db, database.Organization{})
user := dbgen.User(t, db, database.User{})
versionJob := dbgen.ProvisionerJob(t, db, nil, database.ProvisionerJob{
OrganizationID: org.ID,
Type: database.ProvisionerJobTypeTemplateVersionImport,
})
version := dbgen.TemplateVersion(t, db, database.TemplateVersion{
OrganizationID: org.ID,
JobID: versionJob.ID,
CreatedBy: user.ID,
})
template := dbgen.Template(t, db, database.Template{
OrganizationID: org.ID,
ActiveVersionID: version.ID,
CreatedBy: user.ID,
})
workspace := dbgen.Workspace(t, db, database.WorkspaceTable{
OwnerID: user.ID,
OrganizationID: org.ID,
TemplateID: template.ID,
})

now := dbtime.Now()
cutoff := now.Add(-24 * time.Hour)
buildTime := cutoff.Add(-time.Hour)
oldCompletedTime := cutoff.Add(-3 * time.Minute)
oldFailedTime := cutoff.Add(-2 * time.Minute)
oldCanceledTime := cutoff.Add(-time.Minute)
oldPendingTime := cutoff.Add(-time.Minute)
recentTime := cutoff.Add(time.Minute)

createBuild := func(buildNumber int32, createdAt time.Time) database.WorkspaceBuild {
return mustCreateWorkspaceBuild(t, db, org, version, workspace.ID, createdAt, buildNumber)
}
insertOrchestration := func(parentBuild database.WorkspaceBuild, updatedAt time.Time) database.WorkspaceBuildOrchestration {
orchestration, err := db.InsertWorkspaceBuildOrchestration(ctx, database.InsertWorkspaceBuildOrchestrationParams{
ID: uuid.New(),
CreatedAt: updatedAt,
UpdatedAt: updatedAt,
ParentBuildID: parentBuild.ID,
ChildTransition: database.WorkspaceTransitionStart,
ChildRichParameterValues: json.RawMessage("[]"),
})
require.NoError(t, err)
return orchestration
}

// Given: old terminal orchestration rows (completed, failed,
// canceled), an old pending row, and a recent terminal row.
oldCompletedParent := createBuild(1, buildTime)
oldCompletedChild := createBuild(2, buildTime)
oldCompleted := insertOrchestration(oldCompletedParent, oldCompletedTime)
_, err := db.UpdateWorkspaceBuildOrchestrationCompletedByID(ctx, database.UpdateWorkspaceBuildOrchestrationCompletedByIDParams{
ID: oldCompleted.ID,
ChildBuildID: uuid.NullUUID{UUID: oldCompletedChild.ID, Valid: true},
UpdatedAt: oldCompletedTime,
})
require.NoError(t, err)

oldFailedParent := createBuild(3, buildTime)
oldFailed := insertOrchestration(oldFailedParent, oldFailedTime)
_, err = db.UpdateWorkspaceBuildOrchestrationFailedByID(ctx, database.UpdateWorkspaceBuildOrchestrationFailedByIDParams{
ID: oldFailed.ID,
Error: sql.NullString{String: "failed", Valid: true},
UpdatedAt: oldFailedTime,
})
require.NoError(t, err)

oldCanceledParent := createBuild(4, buildTime)
oldCanceled := insertOrchestration(oldCanceledParent, oldCanceledTime)
_, err = db.UpdateWorkspaceBuildOrchestrationCanceledByID(ctx, database.UpdateWorkspaceBuildOrchestrationCanceledByIDParams{
ID: oldCanceled.ID,
UpdatedAt: oldCanceledTime,
})
require.NoError(t, err)

oldPendingParent := createBuild(5, buildTime)
oldPending := insertOrchestration(oldPendingParent, oldPendingTime)

recentCompletedParent := createBuild(6, buildTime)
recentCompletedChild := createBuild(7, buildTime)
recentCompleted := insertOrchestration(recentCompletedParent, recentTime)
_, err = db.UpdateWorkspaceBuildOrchestrationCompletedByID(ctx, database.UpdateWorkspaceBuildOrchestrationCompletedByIDParams{
ID: recentCompleted.ID,
ChildBuildID: uuid.NullUUID{UUID: recentCompletedChild.ID, Valid: true},
UpdatedAt: recentTime,
})
require.NoError(t, err)

// When: old workspace build orchestrations are deleted with LimitCount 1
deleted, err := db.DeleteOldWorkspaceBuildOrchestrations(ctx, database.DeleteOldWorkspaceBuildOrchestrationsParams{
BeforeTime: cutoff,
LimitCount: 1,
})
require.NoError(t, err)
require.EqualValues(t, 1, deleted)

// Then: only the oldest terminal row is deleted.
assertOrchestrationDeleted(ctx, t, rawDB, oldCompletedParent.ID)
assertOrchestrationExists(ctx, t, rawDB, oldFailedParent.ID, oldFailed.ID)
assertOrchestrationExists(ctx, t, rawDB, oldCanceledParent.ID, oldCanceled.ID)
assertOrchestrationExists(ctx, t, rawDB, oldPendingParent.ID, oldPending.ID)
assertOrchestrationExists(ctx, t, rawDB, recentCompletedParent.ID, recentCompleted.ID)

// When: old workspace build orchestrations are deleted again.
deleted, err = db.DeleteOldWorkspaceBuildOrchestrations(ctx, database.DeleteOldWorkspaceBuildOrchestrationsParams{
BeforeTime: cutoff,
LimitCount: 10,
})
require.NoError(t, err)
require.EqualValues(t, 2, deleted)

// Then: the remaining old terminal rows are deleted.
assertOrchestrationDeleted(ctx, t, rawDB, oldFailedParent.ID)
assertOrchestrationDeleted(ctx, t, rawDB, oldCanceledParent.ID)
assertOrchestrationExists(ctx, t, rawDB, oldPendingParent.ID, oldPending.ID)
assertOrchestrationExists(ctx, t, rawDB, recentCompletedParent.ID, recentCompleted.ID)
}

func assertOrchestrationDeleted(ctx context.Context, t *testing.T, rawDB *sql.DB, parentBuildID uuid.UUID) {
t.Helper()

_, err := dbtestutil.GetWorkspaceBuildOrchestrationByParentBuildID(ctx, rawDB, parentBuildID)
require.ErrorIs(t, err, sql.ErrNoRows)
}

func assertOrchestrationExists(ctx context.Context, t *testing.T, rawDB *sql.DB, parentBuildID uuid.UUID, orchestrationID uuid.UUID) {
t.Helper()

orchestration, err := dbtestutil.GetWorkspaceBuildOrchestrationByParentBuildID(ctx, rawDB, parentBuildID)
require.NoError(t, err)
require.Equal(t, orchestrationID, orchestration.ID)
}

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

Expand Down
1 change: 1 addition & 0 deletions coderd/database/querier.go

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

31 changes: 31 additions & 0 deletions coderd/database/queries.sql.go

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

Loading
Loading