.
This commit is contained in:
@@ -75,6 +75,12 @@ namespace GitHub.Runner.Common
|
|||||||
await _brokerHttpClient.DeleteSessionAsync(cancellationToken);
|
await _brokerHttpClient.DeleteSessionAsync(cancellationToken);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public async Task DeleteRunnerMessageAsync(Guid sessionId, string jobMessageKey, CancellationToken cancellationToken)
|
||||||
|
{
|
||||||
|
CheckConnection();
|
||||||
|
await _brokerHttpClient.DeleteRunnerMessageAsync(sessionId, jobMessageKey, cancellationToken);
|
||||||
|
}
|
||||||
|
|
||||||
public Task UpdateConnectionIfNeeded(Uri serverUri, VssCredentials credentials)
|
public Task UpdateConnectionIfNeeded(Uri serverUri, VssCredentials credentials)
|
||||||
{
|
{
|
||||||
if (_brokerUri != serverUri || !_hasConnection)
|
if (_brokerUri != serverUri || !_hasConnection)
|
||||||
|
|||||||
@@ -314,7 +314,30 @@ namespace GitHub.Runner.Listener
|
|||||||
|
|
||||||
public async Task DeleteMessageAsync(TaskAgentMessage message)
|
public async Task DeleteMessageAsync(TaskAgentMessage message)
|
||||||
{
|
{
|
||||||
await Task.CompletedTask;
|
Trace.Entering();
|
||||||
|
ArgUtil.NotNull(_session, nameof(_session));
|
||||||
|
|
||||||
|
if (message == null || _session.SessionId == Guid.Empty)
|
||||||
|
{
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
var jobMessageKey = "";
|
||||||
|
if (MessageUtil.IsRunServiceJob(message.MessageType))
|
||||||
|
{
|
||||||
|
var messageRef = StringUtil.ConvertFromJson<RunnerJobRequestRef>(message.Body);
|
||||||
|
jobMessageKey = messageRef.RunnerRequestId;
|
||||||
|
}
|
||||||
|
else
|
||||||
|
{
|
||||||
|
// Broker currently doesn't support delete for other message types
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
using (var cs = new CancellationTokenSource(TimeSpan.FromSeconds(30)))
|
||||||
|
{
|
||||||
|
await _brokerServer.DeleteRunnerMessageAsync(_session.SessionId, jobMessageKey, cs.Token);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private bool IsGetNextMessageExceptionRetriable(Exception ex)
|
private bool IsGetNextMessageExceptionRetriable(Exception ex)
|
||||||
|
|||||||
@@ -563,9 +563,7 @@ namespace GitHub.Runner.Listener
|
|||||||
await runServer.ConnectAsync(new Uri(messageRef.RunServiceUrl), creds);
|
await runServer.ConnectAsync(new Uri(messageRef.RunServiceUrl), creds);
|
||||||
try
|
try
|
||||||
{
|
{
|
||||||
jobRequestMessage =
|
jobRequestMessage = await runServer.GetJobMessageAsync(messageRef.RunnerRequestId, messageQueueLoopTokenSource.Token);
|
||||||
await runServer.GetJobMessageAsync(messageRef.RunnerRequestId,
|
|
||||||
messageQueueLoopTokenSource.Token);
|
|
||||||
}
|
}
|
||||||
catch (TaskOrchestrationJobAlreadyAcquiredException)
|
catch (TaskOrchestrationJobAlreadyAcquiredException)
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -63,8 +63,7 @@ namespace GitHub.Actions.RunService.WebApi
|
|||||||
string os = null,
|
string os = null,
|
||||||
string architecture = null,
|
string architecture = null,
|
||||||
bool? disableUpdate = null,
|
bool? disableUpdate = null,
|
||||||
CancellationToken cancellationToken = default
|
CancellationToken cancellationToken = default)
|
||||||
)
|
|
||||||
{
|
{
|
||||||
var requestUri = new Uri(Client.BaseAddress, "message");
|
var requestUri = new Uri(Client.BaseAddress, "message");
|
||||||
|
|
||||||
@@ -123,8 +122,40 @@ namespace GitHub.Actions.RunService.WebApi
|
|||||||
throw new Exception($"Failed to get job message: {result.Error}");
|
throw new Exception($"Failed to get job message: {result.Error}");
|
||||||
}
|
}
|
||||||
|
|
||||||
public async Task<TaskAgentSession> CreateSessionAsync(
|
public async Task<TaskAgentMessage> GetRunnerMessageAsync(
|
||||||
|
Guid? sessionId,
|
||||||
|
string jobMessageKey,
|
||||||
|
CancellationToken cancellationToken = default)
|
||||||
|
{
|
||||||
|
var requestUri = new Uri(Client.BaseAddress, "message");
|
||||||
|
|
||||||
|
List<KeyValuePair<string, string>> queryParams = new List<KeyValuePair<string, string>>();
|
||||||
|
|
||||||
|
if (sessionId != null)
|
||||||
|
{
|
||||||
|
queryParams.Add("sessionId", sessionId.Value.ToString());
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!string.IsNullOrEmpty(jobMessageKey))
|
||||||
|
{
|
||||||
|
queryParams.Add("jobMessageKey", jobMessageKey);
|
||||||
|
}
|
||||||
|
|
||||||
|
var result = await SendAsync<TaskAgentMessage>(
|
||||||
|
new HttpMethod("DELETE"),
|
||||||
|
requestUri: requestUri,
|
||||||
|
queryParameters: queryParams,
|
||||||
|
cancellationToken: cancellationToken);
|
||||||
|
|
||||||
|
if (result.IsSuccess)
|
||||||
|
{
|
||||||
|
return result.Value;
|
||||||
|
}
|
||||||
|
|
||||||
|
throw new Exception($"Failed to get job message: StatusCode={result.StatusCode} Error={result.Error}");
|
||||||
|
}
|
||||||
|
|
||||||
|
public async Task<TaskAgentSession> CreateSessionAsync(
|
||||||
TaskAgentSession session,
|
TaskAgentSession session,
|
||||||
CancellationToken cancellationToken = default)
|
CancellationToken cancellationToken = default)
|
||||||
{
|
{
|
||||||
|
|||||||
Reference in New Issue
Block a user