diff --git a/docs/roadmap.md b/docs/roadmap.md index c3275da..7750855 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -113,7 +113,7 @@ - [x] создание, чтение, подтверждение и отмена через Operator API; - [ ] выполнение, retry, verification и rollback в полной state machine; - [x] обязательные idempotency key, audit reason и job TTL; -- [ ] Agent job channel только для собственного `node_id`; +- [x] атомарный Agent job channel только для собственного `node_id`; - [x] начальные строгие JSON schemas v1 для `preflight` и `system.base-install`; - [ ] versioned JSON schemas для остальных action type; - [ ] restricted root helper через Unix socket; diff --git a/src/ServerMonitorManager.Control/Program.cs b/src/ServerMonitorManager.Control/Program.cs index fcd7fc8..e84c0fe 100644 --- a/src/ServerMonitorManager.Control/Program.cs +++ b/src/ServerMonitorManager.Control/Program.cs @@ -377,6 +377,20 @@ agents.MapPost("/heartbeat", async ( }); } }); +agents.MapGet("/provisioning/jobs/next", async ( + HttpContext context, + ControlStore controlStore, + CancellationToken cancellationToken) => +{ + var nodeId = context.User.FindFirstValue(ClaimTypes.NameIdentifier); + if (string.IsNullOrWhiteSpace(nodeId) || !NodeIdValidator.IsValid(nodeId)) + { + return Results.Forbid(); + } + + var job = await controlStore.ClaimNextProvisioningJobAsync(nodeId, cancellationToken); + return job is null ? Results.NoContent() : Results.Ok(job); +}); var control = app.MapGroup("/api/v1/control").RequireAuthorization("Operator"); control.MapGet("/agents", async (ControlStore controlStore, CancellationToken cancellationToken) => diff --git a/src/ServerMonitorManager.Control/ProvisioningStore.cs b/src/ServerMonitorManager.Control/ProvisioningStore.cs index 35f0eed..0c20ad2 100644 --- a/src/ServerMonitorManager.Control/ProvisioningStore.cs +++ b/src/ServerMonitorManager.Control/ProvisioningStore.cs @@ -96,6 +96,56 @@ public sealed partial class ControlStore return await reader.ReadAsync(cancellationToken) ? ReadProvisioningJob(reader) : null; } + public async Task ClaimNextProvisioningJobAsync( + string nodeId, + CancellationToken cancellationToken = default) + { + await using var connection = await OpenAsync(cancellationToken); + await using var transaction = (SqliteTransaction)await connection.BeginTransactionAsync(cancellationToken); + var now = DateTimeOffset.UtcNow; + var command = connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = """ + UPDATE provisioning_jobs SET + state = $preflight, + updated_at = $now, + version = version + 1 + WHERE id = ( + SELECT id FROM provisioning_jobs + WHERE node_id = $node + AND state = $queued + AND expires_at > $now + ORDER BY created_at, id + LIMIT 1) + AND state = $queued + RETURNING *; + """; + command.Parameters.AddWithValue("$node", nodeId); + command.Parameters.AddWithValue("$queued", ProvisioningJobStates.Queued); + command.Parameters.AddWithValue("$preflight", ProvisioningJobStates.Preflight); + command.Parameters.AddWithValue("$now", now.ToString("O")); + ProvisioningJob? job; + await using (var reader = await command.ExecuteReaderAsync(cancellationToken)) + { + job = await reader.ReadAsync(cancellationToken) ? ReadProvisioningJob(reader) : null; + } + if (job is null) + { + await transaction.CommitAsync(cancellationToken); + return null; + } + + await WriteProvisioningEventAsync( + connection, transaction, job.Id, "job.claimed", job.State, + "Provisioning job claimed by its assigned Node.", now, cancellationToken); + await WriteAuditAsync( + connection, transaction, nodeId, "provisioning.job.claimed", job.Id, + JsonSerializer.Serialize(new { job.NodeId, job.ActionType, job.SchemaVersion }), + cancellationToken); + await transaction.CommitAsync(cancellationToken); + return job; + } + public Task ConfirmProvisioningJobAsync( string id, ProvisioningJobCommandRequest request, @@ -148,7 +198,9 @@ public sealed partial class ControlStore } var allowed = requiredState is not null ? current.State == requiredState - : current.State is ProvisioningJobStates.Queued or ProvisioningJobStates.AwaitingConfirmation; + : current.State is ProvisioningJobStates.Queued + or ProvisioningJobStates.Preflight + or ProvisioningJobStates.AwaitingConfirmation; if (!allowed || current.ExpiresAt <= DateTimeOffset.UtcNow) { throw new ProvisioningTransitionException(current.State, targetState); diff --git a/src/ServerMonitorManager.Core/Contracts.cs b/src/ServerMonitorManager.Core/Contracts.cs index ad4d79b..8622801 100644 --- a/src/ServerMonitorManager.Core/Contracts.cs +++ b/src/ServerMonitorManager.Core/Contracts.cs @@ -176,6 +176,7 @@ public sealed record ProvisioningJob( public static class ProvisioningJobStates { public const string Queued = "Queued"; + public const string Preflight = "Preflight"; public const string AwaitingConfirmation = "AwaitingConfirmation"; public const string Cancelled = "Cancelled"; } diff --git a/tests/ServerMonitorManager.Control.Tests/ControlApiTests.cs b/tests/ServerMonitorManager.Control.Tests/ControlApiTests.cs index 0b11e8c..17a0749 100644 --- a/tests/ServerMonitorManager.Control.Tests/ControlApiTests.cs +++ b/tests/ServerMonitorManager.Control.Tests/ControlApiTests.cs @@ -118,6 +118,50 @@ public sealed class ControlApiTests : IAsyncDisposable $"/api/v1/control/provisioning/jobs/{job.Id}", cancellationToken)).StatusCode); } + [Fact] + public async Task AgentReceivesOnlyItsOwnProvisioningJob() + { + var nodeId = $"node-{Guid.NewGuid():N}"[..13]; + var store = _factory.Services.GetRequiredService(); + var cancellationToken = TestContext.Current.CancellationToken; + var token = await store.CreateEnrollmentTokenAsync(nodeId, TimeSpan.FromMinutes(10), cancellationToken); + await store.EnrollAsync( + new ServerMonitorManager.Core.EnrollmentRequest( + nodeId, token, "csr", Guid.NewGuid().ToString()), + () => new ServerMonitorManager.Control.IssuedCertificate( + "certificate", "ca", Guid.NewGuid().ToString("N"), DateTimeOffset.UtcNow.AddDays(1)), + cancellationToken); + using var parameters = System.Text.Json.JsonDocument.Parse("{}"); + var created = await store.CreateProvisioningJobAsync( + nodeId, + new ServerMonitorManager.Core.ProvisioningJobCreateRequest( + "preflight", 1, parameters.RootElement.Clone(), 60, + "API Agent isolation test", Guid.NewGuid().ToString()), + "windows-pc", + cancellationToken); + + using var other = _factory.CreateClient(); + other.DefaultRequestHeaders.Add("X-Test-Identity", "other-node"); + other.DefaultRequestHeaders.Add("X-Test-Role", "Agent"); + Assert.Equal(HttpStatusCode.NoContent, (await other.GetAsync( + "/api/v1/agents/provisioning/jobs/next", cancellationToken)).StatusCode); + + using var assigned = _factory.CreateClient(); + assigned.DefaultRequestHeaders.Add("X-Test-Identity", nodeId); + assigned.DefaultRequestHeaders.Add("X-Test-Role", "Agent"); + var response = await assigned.GetAsync( + "/api/v1/agents/provisioning/jobs/next", cancellationToken); + Assert.Equal(HttpStatusCode.OK, response.StatusCode); + var claimed = await response.Content.ReadFromJsonAsync( + cancellationToken); + Assert.NotNull(claimed); + Assert.Equal(created.Id, claimed.Id); + Assert.Equal(nodeId, claimed.NodeId); + Assert.Equal(ServerMonitorManager.Core.ProvisioningJobStates.Preflight, claimed.State); + Assert.Equal(HttpStatusCode.NoContent, (await assigned.GetAsync( + "/api/v1/agents/provisioning/jobs/next", cancellationToken)).StatusCode); + } + public async ValueTask DisposeAsync() => await _factory.DisposeAsync(); private sealed class ControlApiFactory : WebApplicationFactory diff --git a/tests/ServerMonitorManager.Control.Tests/ControlStoreTests.cs b/tests/ServerMonitorManager.Control.Tests/ControlStoreTests.cs index 5036f1c..87bba8c 100644 --- a/tests/ServerMonitorManager.Control.Tests/ControlStoreTests.cs +++ b/tests/ServerMonitorManager.Control.Tests/ControlStoreTests.cs @@ -83,6 +83,35 @@ public sealed class ControlStoreTests : IAsyncDisposable Assert.Equal(ProvisioningJobStates.AwaitingConfirmation, next.State); } + [Fact] + public async Task ProvisioningJobCanBeClaimedOnlyOnceByAssignedNode() + { + var cancellationToken = TestContext.Current.CancellationToken; + var store = CreateStore(); + await store.InitializeAsync(cancellationToken); + await EnrollAgentAsync(store, "home", "C3D4", cancellationToken); + await EnrollAgentAsync(store, "other", "E5F6", cancellationToken); + using var parameters = JsonDocument.Parse("{}"); + var created = await store.CreateProvisioningJobAsync( + "home", + new ProvisioningJobCreateRequest( + "preflight", 1, parameters.RootElement.Clone(), 60, + "Inspect server", Guid.NewGuid().ToString()), + "operator", + cancellationToken); + + Assert.Null(await store.ClaimNextProvisioningJobAsync("other", cancellationToken)); + var claims = await Task.WhenAll( + store.ClaimNextProvisioningJobAsync("home", cancellationToken), + store.ClaimNextProvisioningJobAsync("home", cancellationToken)); + + var claimed = Assert.Single(claims, job => job is not null)!; + Assert.Equal(created.Id, claimed.Id); + Assert.Equal("home", claimed.NodeId); + Assert.Equal(ProvisioningJobStates.Preflight, claimed.State); + Assert.Null(Assert.Single(claims, job => job is null)); + } + [Fact] public void DiagnosticsExportOmitsRawIdentitiesAndNormalizesStates() {