feat(domain): scan uploads with clamd over INSTREAM; 422 refused, 503 scanner down (refs #192)

Verified against a real clamd 1.4.6: clean → Clean, EICAR → Infected,
closed port → Unavailable.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
not
2026-10-02 09:30:05 +02:00
co-authored by Claude Opus 5.5
parent ed18f1fc5f
commit f380fe582a
2 changed files with 57 additions and 7 deletions
+16 -4
View File
@@ -37,6 +37,10 @@ builder.Services.AddSingleton(sp => sp.GetRequiredService<IConfiguration>()
builder.Services.AddSingleton(sp => sp.GetRequiredService<IConfiguration>() builder.Services.AddSingleton(sp => sp.GetRequiredService<IConfiguration>()
.GetSection("Acl").Get<AclOptions>() .GetSection("Acl").Get<AclOptions>()
?? throw new InvalidOperationException("Missing configuration section 'Acl'")); ?? throw new InvalidOperationException("Missing configuration section 'Acl'"));
// clamd defaults to the compose/chart service name; ClamAv__* overrides it (S-29, ADR-0036).
builder.Services.AddSingleton(sp => sp.GetRequiredService<IConfiguration>()
.GetSection("ClamAv").Get<ClamAvOptions>() ?? new ClamAvOptions());
builder.Services.AddSingleton<IDocumentScanner, ClamdDocumentScanner>();
// The in-memory registration store is shared between the submit endpoint and the worker (ADR-0009). // The in-memory registration store is shared between the submit endpoint and the worker (ADR-0009).
builder.Services.AddSingleton<IRegistrationStore, InMemoryRegistrationStore>(); builder.Services.AddSingleton<IRegistrationStore, InMemoryRegistrationStore>();
@@ -158,8 +162,8 @@ app.MapPost("/registrations/{id}/withdraw", async (string id, WithdrawRequest bo
// Provide documents (S-10a): the zorgprofessional supplies the documents their registration is parked // Provide documents (S-10a): the zorgprofessional supplies the documents their registration is parked
// waiting for, completing the WachtOpDocumenten task so the process advances to beoordeling (ADR-0017). // waiting for, completing the WachtOpDocumenten task so the process advances to beoordeling (ADR-0017).
// Owner-scoped by the caller's bsn (the BFF forwards it from the DigiD token); unknown or not-the- // Owner-scoped by the caller's bsn (the BFF forwards it from the DigiD token); unknown or not-the-
// caller's is 404 (indistinguishable). Idempotent — completing an already-left wait is a no-op. The // caller's is 404 (indistinguishable). Idempotent — completing an already-left wait is a no-op. Only a
// real file upload + ZGW storage is S-10b; this endpoint is the trigger that unblocks the process. // PDF that clamd scans clean is stored and unblocks the process (S-29).
app.MapPost("/registrations/{id}/documents", async (string id, ProvideDocumentsRequest body, ProvideDocuments provide, CancellationToken ct) => app.MapPost("/registrations/{id}/documents", async (string id, ProvideDocumentsRequest body, ProvideDocuments provide, CancellationToken ct) =>
{ {
if (!Guid.TryParse(id, out var guid)) if (!Guid.TryParse(id, out var guid))
@@ -177,8 +181,16 @@ app.MapPost("/registrations/{id}/documents", async (string id, ProvideDocumentsR
var command = new ProvideDocumentsCommand( var command = new ProvideDocumentsCommand(
new RegistrationId(guid), body.Bsn, content, new RegistrationId(guid), body.Bsn, content,
body.FileName ?? "diploma.pdf", body.ContentType ?? "application/pdf"); body.FileName ?? "diploma.pdf", body.ContentType ?? "application/pdf");
var outcome = await provide.HandleAsync(command, ct); // A refused file is 422 with a machine-readable reason the portal words for the citizen; an
return outcome == ProvideDocumentsOutcome.Accepted ? Results.NoContent() : Results.NotFound(); // unreachable scanner is 503 — retryable, and nothing was stored (S-29, ADR-0036).
return await provide.HandleAsync(command, ct) switch
{
ProvideDocumentsOutcome.Accepted => Results.NoContent(),
ProvideDocumentsOutcome.NotAPdf => Results.UnprocessableEntity(new { reason = "not-a-pdf" }),
ProvideDocumentsOutcome.Infected => Results.UnprocessableEntity(new { reason = "infected" }),
ProvideDocumentsOutcome.ScannerUnavailable => Results.StatusCode(StatusCodes.Status503ServiceUnavailable),
_ => Results.NotFound(),
};
}); });
// The behandelaar's werkbak (S-12): the registrations awaiting beoordeling, read from the open // The behandelaar's werkbak (S-12): the registrations awaiting beoordeling, read from the open
@@ -1,10 +1,48 @@
using System.Buffers.Binary;
using System.Net.Sockets;
using System.Text;
using Big.Application; using Big.Application;
namespace Big.Infrastructure; namespace Big.Infrastructure;
/// <summary>Scans a document with clamd over its INSTREAM protocol (S-29, ADR-0036).</summary> /// <summary>
/// Scans a document with clamd over its INSTREAM protocol (S-29, ADR-0036): <c>zINSTREAM\0</c>, the
/// document as one big-endian length-prefixed chunk, a zero-length terminator, then one reply —
/// <c>stream: OK</c> or <c>stream: &lt;signature&gt; FOUND</c>. Anything else (an ERROR reply, a refused
/// connection, a timeout) is <see cref="ScanVerdict.Unavailable"/>, so the caller fails closed.
/// </summary>
public sealed class ClamdDocumentScanner(ClamAvOptions options) : IDocumentScanner public sealed class ClamdDocumentScanner(ClamAvOptions options) : IDocumentScanner
{ {
public Task<ScanVerdict> ScanAsync(byte[] content, CancellationToken ct = default) public async Task<ScanVerdict> ScanAsync(byte[] content, CancellationToken ct = default)
=> Task.FromResult(ScanVerdict.Clean); {
ArgumentNullException.ThrowIfNull(content);
using var timeout = CancellationTokenSource.CreateLinkedTokenSource(ct);
timeout.CancelAfter(options.Timeout);
string reply;
try
{
using var client = new TcpClient();
await client.ConnectAsync(options.Host, options.Port, timeout.Token);
var stream = client.GetStream();
var length = new byte[4];
BinaryPrimitives.WriteInt32BigEndian(length, content.Length);
await stream.WriteAsync("zINSTREAM\0"u8.ToArray(), timeout.Token);
await stream.WriteAsync(length, timeout.Token);
await stream.WriteAsync(content, timeout.Token);
await stream.WriteAsync(new byte[4], timeout.Token);
using var reader = new StreamReader(stream, Encoding.ASCII);
reply = (await reader.ReadToEndAsync(timeout.Token)).TrimEnd('\0', '\n');
}
catch (Exception e) when (e is SocketException or IOException
|| (e is OperationCanceledException && !ct.IsCancellationRequested))
{
return ScanVerdict.Unavailable;
}
if (reply == "stream: OK") return ScanVerdict.Clean;
return reply.EndsWith(" FOUND", StringComparison.Ordinal) ? ScanVerdict.Infected : ScanVerdict.Unavailable;
}
} }