diff --git a/src/Runner.Listener/Runner.cs b/src/Runner.Listener/Runner.cs index 31c87130257..2cd9e9aea47 100644 --- a/src/Runner.Listener/Runner.cs +++ b/src/Runner.Listener/Runner.cs @@ -740,6 +740,23 @@ ex is TaskOrchestrationJobNotFoundException || // HTTP status 404 ex is TaskOrchestrationJobAlreadyAcquiredException || // HTTP status 409 ex is TaskOrchestrationJobUnprocessableException) // HTTP status 422 { + // 404 and 409 mean the assignment is gone, and the service consumes an + // ephemeral runner's registration when it assigns the job, so there is no + // session left to listen on: skipping would leave the runner alive but + // deregistered, never assigned work again. 422 says the job itself is + // unprocessable, not that the assignment moved, so it keeps skipping. + if (settings.Ephemeral && + (ex is TaskOrchestrationJobNotFoundException || ex is TaskOrchestrationJobAlreadyAcquiredException)) + { + _term.WriteLine("The job assigned to this ephemeral runner is no longer available. Cleaning up local configuration."); + Trace.Info($"Ephemeral runner lost its job assignment. Exiting runner. {ex.Message}"); + // The registration is already gone, so deleting the session would only + // stall on a call the service is bound to reject. + skipSessionDeletion = true; + runOnceJobCompleted = true; + return Constants.Runner.ReturnCode.Success; + } + Trace.Info($"Skipping message Job. {ex.Message}"); await _acquireJobThrottler.IncrementAndWaitAsync(messageQueueLoopTokenSource.Token); continue; diff --git a/src/Test/L0/Listener/RunnerL0.cs b/src/Test/L0/Listener/RunnerL0.cs index fcf442244f3..8c7cede6f51 100644 --- a/src/Test/L0/Listener/RunnerL0.cs +++ b/src/Test/L0/Listener/RunnerL0.cs @@ -1056,6 +1056,158 @@ public async Task TestEphemeralRunnerJobRequestMessageFromRunServiceExitsOnAckno } } + // Shared arrange for the RunService acquire-path tests. The message loop parks instead of + // dequeuing an empty queue: a bare Dequeue surfaces a regression as "Queue empty" thrown + // from the mock, which hides the assertion that was meant to catch it. + private (Queue Messages, TaskCompletionSource Drained) ArrangeRunServiceRunner(TestHostContext hc, Runner.Listener.Runner runner, bool ephemeral) + { + hc.SetSingleton(_configurationManager.Object); + hc.SetSingleton(_jobNotification.Object); + hc.SetSingleton(_messageListener.Object); + hc.SetSingleton(_promptManager.Object); + hc.SetSingleton(_runnerServer.Object); + hc.SetSingleton(_configStore.Object); + hc.SetSingleton(_updater.Object); + hc.SetSingleton(_credentialManager.Object); + hc.EnqueueInstance(_acquireJobThrottler.Object); + hc.EnqueueInstance(_runServer.Object); + hc.EnqueueInstance(_jobDispatcher.Object); + + runner.Initialize(hc); + var settings = new RunnerSettings + { + PoolId = 43242, + AgentId = 5678, + Ephemeral = ephemeral, + ServerUrl = "https://github.com", + }; + + var messages = new Queue(); + messages.Enqueue(new TaskAgentMessage() + { + Body = JsonUtility.ToString(new RunnerJobRequestRef() { BillingOwnerId = "github", RunnerRequestId = "999", RunServiceUrl = "https://run-service.com" }), + MessageId = 4234, + MessageType = JobRequestMessageTypes.RunnerJobRequest + }); + + // Completed once the runner is back asking for work, which means the message loop has + // already re-checked its cancellation token. Signalling any earlier races the shutdown. + var drained = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + _configurationManager.Setup(x => x.LoadSettings()) + .Returns(settings); + _configurationManager.Setup(x => x.IsConfigured()) + .Returns(true); + _messageListener.Setup(x => x.CreateSessionAsync(It.IsAny())) + .Returns(Task.FromResult(CreateSessionResult.Success)); + _messageListener.Setup(x => x.GetNextMessageAsync(It.IsAny())) + .Returns(async (CancellationToken token) => + { + if (0 == messages.Count) + { + drained.TrySetResult(true); + await Task.Delay(Timeout.Infinite, token); + } + + return messages.Dequeue(); + }); + _messageListener.Setup(x => x.DeleteSessionAsync()) + .Returns(Task.CompletedTask); + _messageListener.Setup(x => x.DeleteMessageAsync(It.IsAny())) + .Returns(Task.CompletedTask); + _jobNotification.Setup(x => x.StartClient(It.IsAny())) + .Callback(() => + { + + }); + _credentialManager.Setup(x => x.LoadCredentials(true)).Returns(new VssCredentials()); + _configStore.Setup(x => x.IsServiceConfigured()).Returns(false); + + return (messages, drained); + } + + [Theory] + [InlineData(typeof(TaskOrchestrationJobNotFoundException))] // HTTP status 404 + [InlineData(typeof(TaskOrchestrationJobAlreadyAcquiredException))] // HTTP status 409 + [Trait("Level", "L0")] + [Trait("Category", "Runner")] + public async Task TestEphemeralRunnerJobRequestMessageFromRunServiceExitsOnLostJobAssignment(Type acquireException) + { + // Name the context per row: TestHostContext derives its trace file from + // [CallerMemberName] and deletes any existing one, so rows would overwrite each other. + using (var hc = new TestHostContext(this, $"{nameof(TestEphemeralRunnerJobRequestMessageFromRunServiceExitsOnLostJobAssignment)}_{acquireException.Name}")) + { + //Arrange + var runner = new Runner.Listener.Runner(); + ArrangeRunServiceRunner(hc, runner, ephemeral: true); + _runServer.Setup(x => x.GetJobMessageAsync("999", "github", It.IsAny())) + .ThrowsAsync((Exception)Activator.CreateInstance(acquireException, "Job assignment is invalid: MissingKey")); + + //Act + var command = new CommandSettings(hc, new string[] { "run" }); + Task runnerTask = runner.ExecuteCommand(command); + + //Assert + await Task.WhenAny(runnerTask, Task.Delay(30000)); + + Assert.True(runnerTask.IsCompleted, $"{nameof(runner.ExecuteCommand)} timed out."); + Assert.True(!runnerTask.IsFaulted, runnerTask.Exception?.ToString()); + Assert.Equal(Constants.Runner.ReturnCode.Success, await runnerTask); + + _runServer.Verify(x => x.GetJobMessageAsync("999", "github", It.IsAny()), Times.Once()); + _jobDispatcher.Verify(x => x.Run(It.IsAny(), It.IsAny()), Times.Never()); + _acquireJobThrottler.Verify(x => x.IncrementAndWaitAsync(It.IsAny()), Times.Never()); + // The registration is already consumed, so deleting the session would only stall. + _messageListener.Verify(x => x.DeleteSessionAsync(), Times.Never()); + _messageListener.Verify(x => x.DeleteMessageAsync(It.IsAny()), Times.Once()); + _configurationManager.Verify(x => x.DeleteLocalRunnerConfig(), Times.Once()); + } + } + + [Theory] + // 422 says the job is unprocessable, not that the assignment moved: an ephemeral runner keeps listening. + [InlineData(true, false, typeof(TaskOrchestrationJobUnprocessableException))] + // A persistent runner keeps skipping whatever the service answers. + [InlineData(false, false, typeof(TaskOrchestrationJobNotFoundException))] + [InlineData(false, false, typeof(TaskOrchestrationJobAlreadyAcquiredException))] + // --once is a client-side flag, so the registration is not consumed and the runner must + // keep listening. This pins settings.Ephemeral as the predicate over the in-scope runOnce. + [InlineData(false, true, typeof(TaskOrchestrationJobAlreadyAcquiredException))] + [Trait("Level", "L0")] + [Trait("Category", "Runner")] + public async Task TestRunnerJobRequestMessageFromRunServiceContinuesOnLostJobAssignment(bool ephemeral, bool runOnce, Type acquireException) + { + using (var hc = new TestHostContext(this, $"{nameof(TestRunnerJobRequestMessageFromRunServiceContinuesOnLostJobAssignment)}_{ephemeral}_{runOnce}_{acquireException.Name}")) + { + //Arrange + var runner = new Runner.Listener.Runner(); + var arrange = ArrangeRunServiceRunner(hc, runner, ephemeral: ephemeral); + _runServer.Setup(x => x.GetJobMessageAsync("999", "github", It.IsAny())) + .ThrowsAsync((Exception)Activator.CreateInstance(acquireException, "Job assignment is invalid: MissingKey")); + + //Act + var command = new CommandSettings(hc, runOnce ? new string[] { "run", "--once" } : new string[] { "run" }); + Task runnerTask = runner.ExecuteCommand(command); + + //Assert + //the runner skips the job and goes back to listening, so it stops only once we shut it down + await Task.WhenAny(arrange.Drained.Task, runnerTask, Task.Delay(30000)); + Assert.True(arrange.Drained.Task.IsCompletedSuccessfully, $"the runner did not go back to listening. {runnerTask.Exception?.ToString()}"); + + hc.ShutdownRunner(ShutdownReason.UserCancelled); + await Task.WhenAny(runnerTask, Task.Delay(30000)); + + Assert.True(runnerTask.IsCompleted, $"{nameof(runner.ExecuteCommand)} timed out."); + Assert.True(runnerTask.IsCanceled); + _runServer.Verify(x => x.GetJobMessageAsync("999", "github", It.IsAny()), Times.Once()); + _acquireJobThrottler.Verify(x => x.IncrementAndWaitAsync(It.IsAny()), Times.Once()); + _jobDispatcher.Verify(x => x.Run(It.IsAny(), It.IsAny()), Times.Never()); + // Skipping must still drain the message, or Broker redelivers it forever. + _messageListener.Verify(x => x.DeleteMessageAsync(It.IsAny()), Times.Once()); + _configurationManager.Verify(x => x.DeleteLocalRunnerConfig(), Times.Never()); + } + } + [Fact] [Trait("Level", "L0")] [Trait("Category", "Runner")]