diff --git a/src/compute-plane-services/nvca/pkg/nvca/queue_manager.go b/src/compute-plane-services/nvca/pkg/nvca/queue_manager.go index fe1f39a9c..bfc7d970c 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/queue_manager.go +++ b/src/compute-plane-services/nvca/pkg/nvca/queue_manager.go @@ -449,9 +449,10 @@ func (qm *QueueManager) SyncQueues(ctx context.Context) error { // Hydrate the creation messages with the new request message IDs in matrics // based on round-robin starting index. - creationMessagesByGPU, anyQueuePullError := createCreationMessageMatricesByGPU( + creationMessagesByGPU, hasCreationErrors := createCreationMessageMatricesByGPU( log, qwInputs, qwOutputs, currQueueRingMetadata, existingMsgIDs, ) + anyQueuePullError = anyQueuePullError || hasCreationErrors // Create worker pool to the number of termination messages plus the unique number of GPUs // Termination requests can act in parallel while creation requests must be locked to a single GPU type per goroutine diff --git a/src/compute-plane-services/nvca/pkg/nvca/queue_manager_test.go b/src/compute-plane-services/nvca/pkg/nvca/queue_manager_test.go index b547c8665..9c88c160c 100644 --- a/src/compute-plane-services/nvca/pkg/nvca/queue_manager_test.go +++ b/src/compute-plane-services/nvca/pkg/nvca/queue_manager_test.go @@ -644,6 +644,47 @@ func TestSyncQueuesWithBk8s(t *testing.T) { require.EventuallyWithT(t, verifyPodsDeleted, 120*time.Second, 100*time.Millisecond) } +// TestSyncQueuesTerminationErrorNotMasked verifies that a termination queue +// pull error keeps the manager status not OK, even when the creation queue +// pull succeeds. +func TestSyncQueuesTerminationErrorNotMasked(t *testing.T) { + ctx, cancel := context.WithCancel(newTestContext()) + t.Cleanup(cancel) + + clients := mockKubeClients() + b := NewBackendk8sCacheBuilder(). + WithNamespaceLabels(labels.Set{"foo": "bar"}). + WithClients(clients). + WithStaticGPUCapacity(10) + bc, _, err := b.Start(ctx) + require.NoError(t, err) + + // Use distinct queue URLs (mod=true) so the termination queue can fail + // independently of the creation queue. + queueCreds := getTestQueueCreds(true) + createQueue := queueCreds.CreationQueues[testGPUNameDefault] + + qc := &mockqueue.Client{ + Use10MillisForWaits: true, + FailReceiveQueueURL: queueCreds.TerminationQueue.QueueURL, + FailReceiveErr: fmt.Errorf("simulated termination queue pull failure"), + } + qc.AddMessage(createQueue.QueueURL, queue.ReceiveMessageOutput{ + MessageID: creationMessageId, + ReceiptHandle: "randomIdTermFail", + Body: []byte(goodCM), + }) + + bsc := newMockBackendStatusCacheFromK8s(t, bc) + metrics := nvcametrics.FromContext(ctx) + qm := NewQueueManager(bc, bsc, qc, queueCreds, featureflag.DefaultFetcher, types.MaintenanceModeNone, metrics) + assert.True(t, qm.StatusOK()) + + err = qm.SyncQueues(ctx) + assert.NoError(t, err) + assert.False(t, qm.StatusOK(), "status must not be OK when termination queue pull fails") +} + func TestSyncQueuesWithBk8sFailedCaching(t *testing.T) { origUUID := GetUseUUIDForRequestObjName() SetUseUUIDForRequestObjName(false) diff --git a/src/compute-plane-services/nvca/pkg/queue/mock/mock.go b/src/compute-plane-services/nvca/pkg/queue/mock/mock.go index e030997fe..ee708af91 100644 --- a/src/compute-plane-services/nvca/pkg/queue/mock/mock.go +++ b/src/compute-plane-services/nvca/pkg/queue/mock/mock.go @@ -38,6 +38,10 @@ type Client struct { // Speed up tests using 10's of milliseconds instead of seconds. Use10MillisForWaits bool + // If set, ReceiveMessage returns FailReceiveErr for FailReceiveQueueURL. + FailReceiveQueueURL string + FailReceiveErr error + queues map[string][]QueueMessage mu sync.RWMutex } @@ -107,6 +111,10 @@ done: func (c *Client) ReceiveMessage(ctx context.Context, input queue.ReceiveMessageInput) ([]queue.ReceiveMessageOutput, error) { qi := input.QueueInfo + if c.FailReceiveErr != nil && qi.QueueURL == c.FailReceiveQueueURL { + return nil, c.FailReceiveErr + } + c.mu.Lock() if c.queues == nil { c.queues = map[string][]QueueMessage{}