Skip to content
Open
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
9 changes: 8 additions & 1 deletion runner/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,15 @@ package runner

import (
"context"
"errors"
"fmt"
"log/slog"
"time"

runnerErrors "github.com/cloudbase/garm-provider-common/errors"
commonParams "github.com/cloudbase/garm-provider-common/params"
"github.com/cloudbase/garm/auth"
garmErrors "github.com/cloudbase/garm/internal/errors"
"github.com/cloudbase/garm/params"
)

Expand Down Expand Up @@ -60,7 +62,12 @@ func (r *Runner) SetInstanceToPendingDelete(ctx context.Context) error {
Status: commonParams.InstancePendingDelete,
}

if _, err := r.store.ForceUpdateInstance(r.ctx, instance.ID, updateParams); err != nil {
if _, err := r.store.UpdateInstance(r.ctx, instance.ID, updateParams); err != nil {
var te *runnerErrors.InstanceTransitionError
if errors.As(err, &te) && garmErrors.InstanceIsBeingDeleted(te.From) {
// Already on the deletion lane; treat the refusal as success.
return nil
}
return fmt.Errorf("failed to set instance to pending_delete: %w", err)
}
return nil
Expand Down
144 changes: 144 additions & 0 deletions runner/agent_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
// Copyright 2025 Cloudbase Solutions SRL
//
// Licensed under the Apache License, Version 2.0 (the "License"); you may
// not use this file except in compliance with the License. You may obtain
// a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
// WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
// License for the specific language governing permissions and limitations
// under the License.

package runner

import (
"context"
"fmt"
"testing"

"github.com/stretchr/testify/suite"

commonParams "github.com/cloudbase/garm-provider-common/params"
"github.com/cloudbase/garm/auth"
"github.com/cloudbase/garm/database"
dbCommon "github.com/cloudbase/garm/database/common"
garmTesting "github.com/cloudbase/garm/internal/testing"
"github.com/cloudbase/garm/params"
)

type AgentTestSuite struct {
suite.Suite

adminCtx context.Context
store dbCommon.Store
runner *Runner
instance params.Instance
}

func (s *AgentTestSuite) SetupTest() {
dbCfg := garmTesting.GetTestSqliteDBConfig(s.T())
db, err := database.NewDatabase(context.Background(), dbCfg)
if err != nil {
s.FailNow(fmt.Sprintf("failed to create db connection: %s", err))
}
s.store = db

s.adminCtx = garmTesting.ImpersonateAdminContext(context.Background(), db, s.T())

githubEndpoint := garmTesting.CreateDefaultGithubEndpoint(s.adminCtx, db, s.T())
testCreds := garmTesting.CreateTestGithubCredentials(s.adminCtx, "test-creds", db, s.T(), githubEndpoint)

org, err := db.CreateOrganization(s.adminCtx, "test-org", testCreds, "test-webhook-secret", params.PoolBalancerTypeRoundRobin, false)
if err != nil {
s.FailNow(fmt.Sprintf("failed to create test org: %s", err))
}

entity, err := org.GetEntity()
if err != nil {
s.FailNow(fmt.Sprintf("failed to get entity: %s", err))
}

pool, err := db.CreateEntityPool(s.adminCtx, entity, params.CreatePoolParams{
ProviderName: "test-provider",
MaxRunners: 2,
MinIdleRunners: 0,
Image: "ubuntu:22.04",
Flavor: "medium",
OSType: commonParams.Linux,
OSArch: commonParams.Amd64,
Tags: []string{"linux", "amd64"},
RunnerBootstrapTimeout: 10,
})
if err != nil {
s.FailNow(fmt.Sprintf("failed to create test pool: %s", err))
}

instance, err := db.CreateInstance(s.adminCtx, pool.ID, params.CreateInstanceParams{
Name: "test-agent-instance",
OSType: commonParams.Linux,
OSArch: commonParams.Amd64,
})
if err != nil {
s.FailNow(fmt.Sprintf("failed to create test instance: %s", err))
}
s.instance = instance

s.runner = &Runner{
ctx: s.adminCtx,
store: db,
}
}

func (s *AgentTestSuite) seedStatus(status commonParams.InstanceStatus) {
_, err := s.store.ForceUpdateInstance(s.adminCtx, s.instance.Name, params.UpdateInstanceParams{
Status: status,
})
s.Require().NoError(err)
}

func (s *AgentTestSuite) instanceCtx() context.Context {
return auth.SetInstanceParams(context.Background(), s.instance)
}

// A terminated report on a running instance marks it pending_delete.
func (s *AgentTestSuite) TestSetInstanceToPendingDeleteFromRunning() {
s.seedStatus(commonParams.InstanceRunning)

err := s.runner.SetInstanceToPendingDelete(s.instanceCtx())
s.Require().NoError(err)

updated, err := s.store.GetInstance(s.adminCtx, s.instance.Name)
s.Require().NoError(err)
s.Require().Equal(commonParams.InstancePendingDelete, updated.Status)
}

// A terminated report must not move a deleting instance backwards.
func (s *AgentTestSuite) TestSetInstanceToPendingDeleteDoesNotRegressDeleting() {
s.seedStatus(commonParams.InstanceDeleting)

err := s.runner.SetInstanceToPendingDelete(s.instanceCtx())
s.Require().NoError(err)

updated, err := s.store.GetInstance(s.adminCtx, s.instance.Name)
s.Require().NoError(err)
s.Require().Equal(commonParams.InstanceDeleting, updated.Status)
}

// A terminated report on an already deleted record is a no-op.
func (s *AgentTestSuite) TestSetInstanceToPendingDeleteToleratesDeleted() {
s.seedStatus(commonParams.InstanceDeleted)

err := s.runner.SetInstanceToPendingDelete(s.instanceCtx())
s.Require().NoError(err)

updated, err := s.store.GetInstance(s.adminCtx, s.instance.Name)
s.Require().NoError(err)
s.Require().Equal(commonParams.InstanceDeleted, updated.Status)
}

func TestAgentTestSuite(t *testing.T) {
suite.Run(t, new(AgentTestSuite))
}
22 changes: 18 additions & 4 deletions workers/provider/instance_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import (
commonParams "github.com/cloudbase/garm-provider-common/params"
"github.com/cloudbase/garm/cache"
dbCommon "github.com/cloudbase/garm/database/common"
garmErrors "github.com/cloudbase/garm/internal/errors"
"github.com/cloudbase/garm/params"
"github.com/cloudbase/garm/runner/common"
garmUtil "github.com/cloudbase/garm/util"
Expand Down Expand Up @@ -339,9 +340,10 @@ func (i *instanceManager) consolidateState() error {
}
case commonParams.InstanceRunning:
// Nothing to do. The provider finished creating the instance.
case commonParams.InstancePendingDelete, commonParams.InstancePendingForceDelete:
case commonParams.InstancePendingDelete, commonParams.InstancePendingForceDelete, commonParams.InstanceDeleting:
// Remove or force remove the runner. When force remove is specified, we ignore
// IaaS errors.
// InstanceDeleting resumes an unfinished delete; deletes are idempotent.
if i.instance.Status == commonParams.InstancePendingDelete {
// invoke backoff sleep. We only do this for non forced removals,
// as force delete will always return, regardless of whether or not
Expand All @@ -362,8 +364,8 @@ func (i *instanceManager) consolidateState() error {
}

if err := i.handleDeleteInstanceInProvider(i.instance); err != nil {
slog.ErrorContext(i.ctx, "deleting instance in provider", "error", err, "forced", i.instance.Status == commonParams.InstancePendingForceDelete)
if prevStatus == commonParams.InstancePendingDelete {
slog.ErrorContext(i.ctx, "deleting instance in provider", "error", err, "forced", prevStatus == commonParams.InstancePendingForceDelete)
if prevStatus != commonParams.InstancePendingForceDelete {
i.incrementBackOff()
if err := i.helper.SetInstanceStatus(i.instance.Name, commonParams.InstancePendingDelete, []byte(err.Error()), true); err != nil {
return fmt.Errorf("setting instance status to error: %w", err)
Expand All @@ -373,7 +375,19 @@ func (i *instanceManager) consolidateState() error {
}
}
if err := i.helper.SetInstanceStatus(i.instance.Name, commonParams.InstanceDeleted, nil, false); err != nil {
if !errors.Is(err, runnerErrors.ErrNotFound) {
var transitionErr *runnerErrors.InstanceTransitionError
switch {
case errors.Is(err, runnerErrors.ErrNotFound):
// The record is already gone; nothing left to update.
case errors.As(err, &transitionErr) && garmErrors.InstanceIsBeingDeleted(transitionErr.From):
// The row moved back onto the deletion lane; the provider
// resource is gone, so force the final transition.
if err := i.helper.SetInstanceStatus(i.instance.Name, commonParams.InstanceDeleted, nil, true); err != nil {
if !errors.Is(err, runnerErrors.ErrNotFound) {
return fmt.Errorf("setting instance status to deleted: %w", err)
}
}
default:
return fmt.Errorf("setting instance status to deleted: %w", err)
}
}
Expand Down
163 changes: 163 additions & 0 deletions workers/provider/instance_manager_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,163 @@
// Copyright 2025 Cloudbase Solutions SRL
//
// Licensed under the Apache License, Version 2.0 (the "License"); you may
// not use this file except in compliance with the License. You may obtain
// a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
// WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
// License for the specific language governing permissions and limitations
// under the License.
package provider

import (
"context"
"fmt"
"testing"

"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"

runnerErrors "github.com/cloudbase/garm-provider-common/errors"
commonParams "github.com/cloudbase/garm-provider-common/params"
"github.com/cloudbase/garm/auth"
"github.com/cloudbase/garm/params"
"github.com/cloudbase/garm/runner/common"
runnerCommonMocks "github.com/cloudbase/garm/runner/common/mocks"
)

type statusWrite struct {
status commonParams.InstanceStatus
force bool
}

// fakeProviderHelper records status writes and lets a test script their
// outcome. Everything else returns zero values; consolidateState's deletion
// path only needs SetInstanceStatus and GetControllerInfo.
type fakeProviderHelper struct {
writes []statusWrite
onSetStatus func(w statusWrite) error
}

func (f *fakeProviderHelper) SetInstanceStatus(_ string, status commonParams.InstanceStatus, _ []byte, force bool) error {
w := statusWrite{status: status, force: force}
f.writes = append(f.writes, w)
if f.onSetStatus != nil {
return f.onSetStatus(w)
}
return nil
}

func (f *fakeProviderHelper) InstanceTokenGetter() auth.InstanceTokenGetter { return nil }

func (f *fakeProviderHelper) updateArgsFromProviderInstance(_ string, _ commonParams.ProviderInstance) (params.Instance, error) {
return params.Instance{}, nil
}

func (f *fakeProviderHelper) GetControllerInfo() (params.ControllerInfo, error) {
return params.ControllerInfo{}, nil
}

func (f *fakeProviderHelper) GetGithubEntity(e params.ForgeEntity) (params.ForgeEntity, error) {
return e, nil
}

func newTestManager(status commonParams.InstanceStatus, provider common.Provider, helper providerHelper) *instanceManager {
m := &instanceManager{
ctx: context.Background(),
instance: params.Instance{
Name: "test-instance",
ProviderID: "test-instance",
Status: status,
},
provider: provider,
helper: helper,
}
m.running.Store(true)
return m
}

// A manager whose cached state is deleting must resume the delete.
func TestConsolidateStateResumesDeleting(t *testing.T) {
providerMock := runnerCommonMocks.NewProvider(t)
providerMock.On("DeleteInstance", mock.Anything, "test-instance", mock.Anything).Return(nil)
helper := &fakeProviderHelper{}
m := newTestManager(commonParams.InstanceDeleting, providerMock, helper)

err := m.consolidateState()
require.ErrorIs(t, err, ErrInstanceDeleted)
require.Equal(t, []statusWrite{
{status: commonParams.InstanceDeleting, force: true},
{status: commonParams.InstanceDeleted, force: false},
}, helper.writes)
}

// A resumed delete that fails in the provider requeues.
func TestConsolidateStateResumedDeletingRequeuesOnProviderError(t *testing.T) {
providerMock := runnerCommonMocks.NewProvider(t)
providerMock.On("DeleteInstance", mock.Anything, "test-instance", mock.Anything).Return(fmt.Errorf("provider exploded"))
helper := &fakeProviderHelper{}
m := newTestManager(commonParams.InstanceDeleting, providerMock, helper)

err := m.consolidateState()
require.Error(t, err)
require.NotErrorIs(t, err, ErrInstanceDeleted)
require.Equal(t, []statusWrite{
{status: commonParams.InstanceDeleting, force: true},
{status: commonParams.InstancePendingDelete, force: true},
}, helper.writes)
require.Greater(t, m.deleteBackoff.Nanoseconds(), int64(0))
}

// Force delete keeps its semantics: provider errors are ignored and the
// instance is marked deleted.
func TestConsolidateStateForceDeleteIgnoresProviderError(t *testing.T) {
providerMock := runnerCommonMocks.NewProvider(t)
providerMock.On("DeleteInstance", mock.Anything, "test-instance", mock.Anything).Return(fmt.Errorf("provider exploded"))
helper := &fakeProviderHelper{}
m := newTestManager(commonParams.InstancePendingForceDelete, providerMock, helper)

err := m.consolidateState()
require.ErrorIs(t, err, ErrInstanceDeleted)
require.Equal(t, []statusWrite{
{status: commonParams.InstanceDeleting, force: true},
{status: commonParams.InstanceDeleted, force: false},
}, helper.writes)
}

// A deleted write refused because the row regressed onto the deletion
// lane is retried with force.
func TestConsolidateStateForcesDeletedWhenRowRegressed(t *testing.T) {
providerMock := runnerCommonMocks.NewProvider(t)
providerMock.On("DeleteInstance", mock.Anything, "test-instance", mock.Anything).Return(nil)
helper := &fakeProviderHelper{}
helper.onSetStatus = func(w statusWrite) error {
if w.status == commonParams.InstanceDeleted && !w.force {
return fmt.Errorf("updating instance: %w", runnerErrors.NewInstanceTransitionError(commonParams.InstancePendingDelete, commonParams.InstanceDeleted))
}
return nil
}
m := newTestManager(commonParams.InstancePendingDelete, providerMock, helper)

err := m.consolidateState()
require.ErrorIs(t, err, ErrInstanceDeleted)
require.Equal(t, []statusWrite{
{status: commonParams.InstanceDeleting, force: true},
{status: commonParams.InstanceDeleted, force: false},
{status: commonParams.InstanceDeleted, force: true},
}, helper.writes)
}

// Running remains a no-op.
func TestConsolidateStateRunningIsNoop(t *testing.T) {
providerMock := runnerCommonMocks.NewProvider(t)
helper := &fakeProviderHelper{}
m := newTestManager(commonParams.InstanceRunning, providerMock, helper)

err := m.consolidateState()
require.NoError(t, err)
require.Empty(t, helper.writes)
}
Loading