Skip to content
Merged
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
1 change: 1 addition & 0 deletions eng/ci/templates/jobs/run-linux-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ jobs:
echo "Starting WorkerProxy image: $(localImage)"
containerId="$(docker run --detach \
--publish 127.0.0.1::80 \
--env WORKERPROXY__PODNAME=ci-worker-pod \
$(localImage))"
echo "Started WorkerProxy container: $containerId"
cleanup() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

using System;
using System.Collections.Generic;
using Azure.Functions.WorkerProxy.Rpc;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;

Expand All @@ -12,6 +13,7 @@ namespace Azure.Functions.WorkerProxy.Http;
/// Captures the worker's HTTP destination before advertising WorkerProxy to the runtime.
/// </summary>
internal sealed partial class WorkerHttpCapabilityProvider(IOptions<WorkerProxyOptions> options, ILogger<WorkerHttpCapabilityProvider> logger)
: IWorkerCapabilityFinalizer
{
private const string HttpUriCapability = "HttpUri";

Expand Down
40 changes: 40 additions & 0 deletions src/Functions.WorkerProxy/Management/ManagementApiEndpoints.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

using System.Threading.Tasks;
using Azure.Functions.WorkerProxy.State;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Http;
using Microsoft.AspNetCore.Routing;

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Registers worker lifecycle APIs on the management listener.
/// </summary>
/// <remarks>
/// Platform callers are expected to send UTF-8 JSON. Assignment uses framework JSON binding:
/// malformed or incompatible JSON with a supported charset returns HTTP 400;
/// unsupported media types return HTTP 415. Binding errors do not guarantee our validation envelope.
/// Other binding failures follow framework behavior. Successfully bound requests use our
/// field-validation envelope and lifecycle error codes.
/// </remarks>
internal static class ManagementApiEndpoints
{
public static void Map(IEndpointRouteBuilder endpoints)
{
endpoints.MapGet("/admin/worker/ready", ManagementApiHandlers.GetWorkerReady)
.AddEndpointFilter(DisableResponseCaching).AllowAnonymous();
endpoints.MapPut("/admin/worker/assignment",
(WorkerAssignRequest request, WorkerPodStateManager manager) => ManagementApiHandlers.AssignWorker(request, manager))
.AllowAnonymous();
endpoints.MapGet("/admin/worker/state", ManagementApiHandlers.GetWorkerStateAsync)
.AddEndpointFilter(DisableResponseCaching).AllowAnonymous();
}

private static ValueTask<object?> DisableResponseCaching(EndpointFilterInvocationContext context, EndpointFilterDelegate next)
{
context.HttpContext.Response.Headers.CacheControl = "no-store";
return next(context);
}
}
113 changes: 113 additions & 0 deletions src/Functions.WorkerProxy/Management/ManagementApiHandlers.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

using System;
using System.Collections.Generic;
using System.Globalization;
using System.Threading;
using System.Threading.Tasks;
using Azure.Functions.WorkerProxy.State;
using Microsoft.AspNetCore.Http;

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Maps validated management requests and worker lifecycle outcomes to HTTP results.
/// </summary>
internal static class ManagementApiHandlers
{
public static IResult GetWorkerReady(WorkerPodStateManager manager) =>
manager.State.IsWorkerReady ? TypedResults.Ok() : TypedResults.StatusCode(StatusCodes.Status503ServiceUnavailable);

public static IResult AssignWorker(WorkerAssignRequest request, WorkerPodStateManager manager)
{
if (!WorkerAssignRequestValidator.TryCreateAssignment(
request, out WorkerAssignment? assignment, out IReadOnlyList<RequestValidationError> errors))
{
return ValidationError(errors);
}

return manager.Assign(assignment) switch
{
WorkerAssignmentResult.Created => TypedResults.Created("/admin/worker/assignment"),
WorkerAssignmentResult.AlreadyAssigned => TypedResults.NoContent(),
WorkerAssignmentResult.WorkerNotReady => Error(
StatusCodes.Status503ServiceUnavailable, WorkerApiErrorCodes.WorkerNotReady, "The worker has not established a valid StartStream."),
WorkerAssignmentResult.AssignmentConflict => Error(
StatusCodes.Status409Conflict, WorkerApiErrorCodes.AssignmentConflict, "The pod is already assigned to a different assignment."),
WorkerAssignmentResult.WorkerTerminated => Error(
StatusCodes.Status409Conflict, WorkerApiErrorCodes.WorkerTerminated, "The assigned worker stream has terminated."),
_ => throw new InvalidOperationException("Unexpected worker assignment result.")
};
}

/// <summary>
/// Parses the optional revision query and returns or polls worker state, honoring request cancellation.
/// </summary>
public static async Task<IResult> GetWorkerStateAsync(HttpRequest request, WorkerPodStateManager manager)
{
if (!TryGetLastKnownRevision(request.Query, out long? lastKnownRevision))
{
return InvalidRevision();
}

return await GetInstanceStateAsync(lastKnownRevision, manager, request.HttpContext.RequestAborted);
}

public static async Task<IResult> GetInstanceStateAsync(
long? revision,
WorkerPodStateManager manager,
CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
if (revision is not { } lastKnownRevision)
{
return StateResponse(manager.State);
}

// Revisions never decrease, so a revision valid here remains valid when the manager registers the poll.
if (lastKnownRevision < 0 || lastKnownRevision > manager.State.Revision)
{
return InvalidRevision();
}

WorkerStatePollResult result = await manager.WaitForChangeAsync(lastKnownRevision, cancellationToken);
return result.State is { } state ? StateResponse(state) : TypedResults.NoContent();
}

internal static IResult InvalidRevision() =>
ValidationError([new(WorkerApiErrorCodes.InvalidRevision, "lastKnownRevision")]);

private static bool TryGetLastKnownRevision(IQueryCollection query, out long? lastKnownRevision)
{
lastKnownRevision = null;
if (!query.TryGetValue("lastKnownRevision", out var revisions))
{
return true;
}

if (revisions is not [string value])
{
return false;
}

if (!long.TryParse(value, NumberStyles.AllowLeadingSign, CultureInfo.InvariantCulture, out long revision))
{
return false;
}

lastKnownRevision = revision;
return true;
}

private static IResult ValidationError(IReadOnlyList<RequestValidationError> errors) =>
TypedResults.Json(new RequestValidationResponse(errors),
WorkerProxyJsonContext.Default.RequestValidationResponse, statusCode: StatusCodes.Status400BadRequest);

private static IResult Error(int statusCode, string code, string detail) =>
TypedResults.Json(new WorkerApiErrorResponse(new(code, detail)),
WorkerProxyJsonContext.Default.WorkerApiErrorResponse, statusCode: statusCode);

private static IResult StateResponse(WorkerPodState state) =>
TypedResults.Json(WorkerInstanceState.FromState(state), WorkerProxyJsonContext.Default.WorkerInstanceState);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Identifies an invalid request field, or request for a body-level error, without echoing its value.
/// </summary>
internal sealed record RequestValidationError(string Code, string Target);
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

using System.Collections.Generic;

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Contains detected field errors, or one request-body error, returned with HTTP 400.
/// </summary>
internal sealed record RequestValidationResponse(IReadOnlyList<RequestValidationError> Errors);
17 changes: 17 additions & 0 deletions src/Functions.WorkerProxy/Management/WorkerApiError.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Provides a stable error code and optional diagnostic detail without echoing request values.
/// </summary>
/// <param name="Code">A case-sensitive contract identifier; existing codes must not be renamed or repurposed.</param>
/// <param name="Detail">Diagnostic text that may change and must not be used for client decisions.</param>
/// <remarks>
/// WorkerNotReady (503) permits retry after readiness. WorkerTerminated (409) is terminal for the
/// assigned session; retrying the same assignment on this pod cannot recover it.
/// AssignmentConflict (409) rejects a different assignment; do not retry that request unchanged.
/// Clients must inspect Code to distinguish the two 409 outcomes and handle unknown codes gracefully.
/// </remarks>
internal sealed record WorkerApiError(string Code, string? Detail = null);
19 changes: 19 additions & 0 deletions src/Functions.WorkerProxy/Management/WorkerApiErrorCodes.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Defines the stable error codes returned by the management APIs.
/// </summary>
internal static class WorkerApiErrorCodes
{
// Clients branch on these exact, case-sensitive wire values. Do not change or repurpose them.
// Keep explicit literals rather than nameof so symbol renames cannot change the contract.
public const string Required = "Required";
public const string InvalidValue = "InvalidValue";
public const string InvalidRevision = "InvalidRevision";
public const string WorkerNotReady = "WorkerNotReady";
public const string WorkerTerminated = "WorkerTerminated";
public const string AssignmentConflict = "AssignmentConflict";
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Wraps worker lifecycle failures in the management API error envelope.
/// </summary>
internal sealed record WorkerApiErrorResponse(WorkerApiError Error);
37 changes: 37 additions & 0 deletions src/Functions.WorkerProxy/Management/WorkerAssignRequest.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
// Copyright (c) .NET Foundation. All rights reserved.
// Licensed under the MIT License. See License.txt in the project root for license information.

using System.Collections.Generic;

namespace Azure.Functions.WorkerProxy.Management;

/// <summary>
/// Describes assignment identity and the caller-selected startup mode.
/// </summary>
/// <remarks>
/// All fields are required. The directory may be empty for Preconfigured but must be nonblank
/// for SpecializationRequired. Configuration is currently retained for retry comparison only.
/// </remarks>
internal sealed class WorkerAssignRequest
{
/// <summary>
/// Gets the required, case-sensitive startup mode selected by the caller.
/// </summary>
/// <remarks>
/// <c>Preconfigured</c> means the worker starts with its final application configuration (BYOC).
/// <c>SpecializationRequired</c> means the worker requires configuration through specialization.
/// Assignment currently records either mode without performing specialization.
/// The raw string is retained so invalid values can be reported as field-level validation errors.
/// </remarks>
public string? StartupMode { get; init; }

public string? FunctionAppName { get; init; }

public string? FunctionGroupName { get; init; }

public bool? IsAlwaysReady { get; init; }

public Dictionary<string, string?>? Environment { get; init; }

public string? FunctionAppDirectory { get; init; }
}
Loading
Loading