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
17 changes: 17 additions & 0 deletions src/Runner.Listener/Runner.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
152 changes: 152 additions & 0 deletions src/Test/L0/Listener/RunnerL0.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<TaskAgentMessage> Messages, TaskCompletionSource<bool> Drained) ArrangeRunServiceRunner(TestHostContext hc, Runner.Listener.Runner runner, bool ephemeral)
{
hc.SetSingleton<IConfigurationManager>(_configurationManager.Object);
hc.SetSingleton<IJobNotification>(_jobNotification.Object);
hc.SetSingleton<IMessageListener>(_messageListener.Object);
hc.SetSingleton<IPromptManager>(_promptManager.Object);
hc.SetSingleton<IRunnerServer>(_runnerServer.Object);
hc.SetSingleton<IConfigurationStore>(_configStore.Object);
hc.SetSingleton<ISelfUpdater>(_updater.Object);
hc.SetSingleton<ICredentialManager>(_credentialManager.Object);
hc.EnqueueInstance<IErrorThrottler>(_acquireJobThrottler.Object);
hc.EnqueueInstance<IRunServer>(_runServer.Object);
hc.EnqueueInstance<IJobDispatcher>(_jobDispatcher.Object);

runner.Initialize(hc);
var settings = new RunnerSettings
{
PoolId = 43242,
AgentId = 5678,
Ephemeral = ephemeral,
ServerUrl = "https://github.com",
};

var messages = new Queue<TaskAgentMessage>();
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<bool>(TaskCreationOptions.RunContinuationsAsynchronously);

_configurationManager.Setup(x => x.LoadSettings())
.Returns(settings);
_configurationManager.Setup(x => x.IsConfigured())
.Returns(true);
_messageListener.Setup(x => x.CreateSessionAsync(It.IsAny<CancellationToken>()))
.Returns(Task.FromResult<CreateSessionResult>(CreateSessionResult.Success));
_messageListener.Setup(x => x.GetNextMessageAsync(It.IsAny<CancellationToken>()))
.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<TaskAgentMessage>()))
.Returns(Task.CompletedTask);
_jobNotification.Setup(x => x.StartClient(It.IsAny<String>()))
.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<CancellationToken>()))
.ThrowsAsync((Exception)Activator.CreateInstance(acquireException, "Job assignment is invalid: MissingKey"));

//Act
var command = new CommandSettings(hc, new string[] { "run" });
Task<int> 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<CancellationToken>()), Times.Once());
_jobDispatcher.Verify(x => x.Run(It.IsAny<Pipelines.AgentJobRequestMessage>(), It.IsAny<bool>()), Times.Never());
_acquireJobThrottler.Verify(x => x.IncrementAndWaitAsync(It.IsAny<CancellationToken>()), 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<TaskAgentMessage>()), 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<CancellationToken>()))
.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<int> 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<CancellationToken>()), Times.Once());
_acquireJobThrottler.Verify(x => x.IncrementAndWaitAsync(It.IsAny<CancellationToken>()), Times.Once());
_jobDispatcher.Verify(x => x.Run(It.IsAny<Pipelines.AgentJobRequestMessage>(), It.IsAny<bool>()), Times.Never());
// Skipping must still drain the message, or Broker redelivers it forever.
_messageListener.Verify(x => x.DeleteMessageAsync(It.IsAny<TaskAgentMessage>()), Times.Once());
_configurationManager.Verify(x => x.DeleteLocalRunnerConfig(), Times.Never());
}
}

[Fact]
[Trait("Level", "L0")]
[Trait("Category", "Runner")]
Expand Down