Files
netbird/management/server/job/manager_test.go
T

79 lines
2.7 KiB
Go

package job
import (
"context"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"go.uber.org/mock/gomock"
"github.com/netbirdio/netbird/management/server/store"
)
func TestCloseChannel_OwnSession(t *testing.T) {
ctx := context.Background()
mockStore := store.NewMockStore(gomock.NewController(t))
jm := NewJobManager(nil, mockStore, nil)
session := NewChannel()
jm.jobChannels["peer-1"] = session
jm.pending["job-1"] = &Event{PeerID: "peer-1"}
jm.pending["job-2"] = &Event{PeerID: "peer-2"}
mockStore.EXPECT().MarkPendingJobsAsFailed(gomock.Any(), "account-1", "peer-1", "job-1", gomock.Any()).Return(nil)
jm.CloseChannel(ctx, "account-1", "peer-1", session)
assert.False(t, jm.IsPeerConnected("peer-1"))
_, err := session.Event(ctx)
assert.ErrorIs(t, err, ErrJobChannelClosed)
assert.NotContains(t, jm.pending, "job-1")
assert.Contains(t, jm.pending, "job-2", "jobs of other peers must be kept")
}
func TestCloseChannel_NewerSessionOwnsPeer(t *testing.T) {
ctx := context.Background()
jm := NewJobManager(nil, store.NewMockStore(gomock.NewController(t)), nil)
stale := NewChannel()
current := NewChannel()
jm.jobChannels["peer-1"] = current
jm.pending["job-1"] = &Event{PeerID: "peer-1"}
jm.CloseChannel(ctx, "account-1", "peer-1", stale)
require.True(t, jm.IsPeerConnected("peer-1"))
assert.Contains(t, jm.pending, "job-1", "jobs served by the newer stream must not be failed")
event := &Event{PeerID: "peer-1"}
require.NoError(t, current.AddEvent(ctx, time.Second, event))
got, err := current.Event(ctx)
require.NoError(t, err, "newer session channel must stay open")
assert.Equal(t, event, got)
}
func TestCreateJobChannel_ReplacesSessionAndDropsItsPendingJobs(t *testing.T) {
ctx := context.Background()
mockStore := store.NewMockStore(gomock.NewController(t))
jm := NewJobManager(nil, mockStore, nil)
stale := NewChannel()
jm.jobChannels["peer-1"] = stale
jm.pending["job-1"] = &Event{PeerID: "peer-1"}
jm.pending["job-2"] = &Event{PeerID: "peer-2"}
mockStore.EXPECT().MarkAllPendingJobsAsFailed(gomock.Any(), "account-1", "peer-1", gomock.Any()).Return(nil)
mockStore.EXPECT().MarkPendingJobsAsFailed(gomock.Any(), "account-1", "peer-1", "job-1", "Pending job cleanup: job stream replaced by a newer stream").Return(nil)
current := jm.CreateJobChannel(ctx, "account-1", "peer-1")
require.NotSame(t, stale, current)
_, err := stale.Event(ctx)
assert.ErrorIs(t, err, ErrJobChannelClosed, "replaced channel must be closed")
assert.True(t, jm.IsPeerConnected("peer-1"))
assert.False(t, jm.IsPeerHasPendingJobs("peer-1"), "jobs of the replaced stream must be dropped")
assert.Contains(t, jm.pending, "job-2", "jobs of other peers must be kept")
}