shoc-backend/SeaHaven.Services/Implementation/WorkOrderWebhookService.cs
2026-07-27 16:33:06 -03:00

431 lines
18 KiB
C#

using System.Diagnostics.Metrics;
using System.Globalization;
using System.Security.Cryptography;
using System.Text;
using System.Text.Json;
using System.Text.Json.Serialization;
using Data.SeaHavenIndustries.Enums;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using SeaHaven.DataServices.Interfaces;
using SeaHaven.Services.Configuration;
using SeaHaven.Services.Constants;
using SeaHaven.Services.Helpers;
using SeaHaven.Services.Interfaces;
namespace SeaHaven.Services.Implementation
{
public sealed class WorkOrderWebhookService : IWorkOrderWebhookService
{
private static readonly Meter Meter = new("SeaHaven.WorkOrderWebhook", "1.0");
private static readonly Counter<long> Accepted = Meter.CreateCounter<long>("work_order_webhook.accepted");
private static readonly Counter<long> Duplicates = Meter.CreateCounter<long>("work_order_webhook.duplicate");
private static readonly Counter<long> Rejected = Meter.CreateCounter<long>("work_order_webhook.rejected");
private static readonly Counter<long> Invalid = Meter.CreateCounter<long>("work_order_webhook.invalid");
private static readonly Counter<long> Failed = Meter.CreateCounter<long>("work_order_webhook.failed");
private static readonly JsonSerializerOptions JsonOptions = new()
{
PropertyNameCaseInsensitive = false,
PropertyNamingPolicy = JsonNamingPolicy.SnakeCaseLower
};
private readonly IWorkOrderWebhookSecretProvider _secretProvider;
private readonly IWorkOrderWebhookDataService _dataService;
private readonly IOptionsMonitor<WorkOrderWebhookOptions> _options;
private readonly TimeProvider _timeProvider;
private readonly ILogger<WorkOrderWebhookService> _logger;
public WorkOrderWebhookService(
IWorkOrderWebhookSecretProvider secretProvider,
IWorkOrderWebhookDataService dataService,
IOptionsMonitor<WorkOrderWebhookOptions> options,
TimeProvider timeProvider,
ILogger<WorkOrderWebhookService> logger)
{
_secretProvider = secretProvider;
_dataService = dataService;
_options = options;
_timeProvider = timeProvider;
_logger = logger;
}
public int MaximumBodyBytes =>
Math.Clamp(_options.CurrentValue.MaxBodyBytes, 1, 1_048_576);
public async Task<WorkOrderWebhookResult> ProcessAsync(
WorkOrderWebhookRequest request,
CancellationToken cancellationToken)
{
var options = _options.CurrentValue;
if (!options.Enabled)
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.Disabled);
if (!TryReadTimestamp(request.Timestamp, options.AllowedClockSkewSeconds, out var timestamp)
|| string.IsNullOrWhiteSpace(request.KeyId)
|| request.KeyId.Length > 128
|| string.IsNullOrWhiteSpace(request.Signature))
{
Rejected.Add(1);
return Unauthorized();
}
WorkOrderWebhookSecretResult secretResult;
try
{
secretResult = await _secretProvider.GetSecretAsync(request.KeyId, cancellationToken);
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
throw;
}
catch (Exception)
{
Failed.Add(1);
_logger.LogError("Work-order webhook signing secret could not be loaded.");
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.Unavailable);
}
if (secretResult.Status == WorkOrderWebhookSecretStatus.Unavailable)
{
Failed.Add(1);
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.Unavailable);
}
if (secretResult.Status != WorkOrderWebhookSecretStatus.Found
|| secretResult.Secret is not { Length: > 0 })
{
Rejected.Add(1);
return Unauthorized();
}
var signatureValid = false;
try
{
signatureValid = VerifySignature(
request.Timestamp!,
request.Body,
request.Signature!,
secretResult.Secret);
}
finally
{
CryptographicOperations.ZeroMemory(secretResult.Secret);
}
if (!signatureValid)
{
Rejected.Add(1);
return Unauthorized();
}
WorkOrderWebhookEnvelope? envelope;
try
{
envelope = JsonSerializer.Deserialize<WorkOrderWebhookEnvelope>(request.Body, JsonOptions);
}
catch (JsonException)
{
Invalid.Add(1);
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.InvalidEnvelope);
}
if (!TryCreateMutation(envelope, request.Body, out var mutation))
{
Invalid.Add(1);
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.InvalidEnvelope);
}
try
{
var persisted = await _dataService.ApplyAsync(mutation!, cancellationToken);
switch (persisted.Status)
{
case WorkOrderWebhookPersistenceStatus.Applied:
Accepted.Add(1);
_logger.LogInformation(
"Work-order webhook delivery {DeliveryId} was accepted. State mutation skipped: {StateMutationSkipped}.",
mutation!.DeliveryId,
persisted.StateMutationSkipped);
return new WorkOrderWebhookResult(
WorkOrderWebhookStatus.Applied,
persisted.StateMutationSkipped);
case WorkOrderWebhookPersistenceStatus.Duplicate:
Duplicates.Add(1);
_logger.LogInformation(
"Work-order webhook delivery {DeliveryId} was already processed.",
mutation!.DeliveryId);
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.Duplicate);
case WorkOrderWebhookPersistenceStatus.HashConflict:
Rejected.Add(1);
_logger.LogWarning(
"Work-order webhook delivery {DeliveryId} conflicts with an existing delivery.",
mutation!.DeliveryId);
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.HashConflict);
default:
throw new InvalidOperationException("Unknown persistence result.");
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
throw;
}
catch (Exception ex)
{
Failed.Add(1);
_logger.LogError(
ex,
"Work-order webhook delivery {DeliveryId} could not be persisted.",
mutation!.DeliveryId);
return new WorkOrderWebhookResult(WorkOrderWebhookStatus.Unavailable);
}
}
private bool TryReadTimestamp(
string? value,
int allowedClockSkewSeconds,
out DateTimeOffset timestamp)
{
timestamp = default;
if (string.IsNullOrEmpty(value)
|| value.Any(c => c is < '0' or > '9')
|| !long.TryParse(value, NumberStyles.None, CultureInfo.InvariantCulture, out var seconds))
return false;
try
{
timestamp = DateTimeOffset.FromUnixTimeSeconds(seconds);
}
catch (ArgumentOutOfRangeException)
{
return false;
}
var skew = (_timeProvider.GetUtcNow() - timestamp).Duration();
return skew <= TimeSpan.FromSeconds(Math.Clamp(allowedClockSkewSeconds, 0, 3600));
}
private static bool VerifySignature(
string timestamp,
byte[] body,
string signature,
byte[] secret)
{
const string prefix = "v1=";
if (!signature.StartsWith(prefix, StringComparison.Ordinal)
|| signature.Length != prefix.Length + 64)
return false;
byte[] supplied;
try
{
supplied = Convert.FromHexString(signature[prefix.Length..]);
}
catch (FormatException)
{
return false;
}
var timestampBytes = Encoding.UTF8.GetBytes(timestamp);
var signed = new byte[timestampBytes.Length + 1 + body.Length];
timestampBytes.CopyTo(signed, 0);
signed[timestampBytes.Length] = (byte)'.';
body.CopyTo(signed, timestampBytes.Length + 1);
var expected = HMACSHA256.HashData(secret, signed);
return CryptographicOperations.FixedTimeEquals(expected, supplied);
}
private bool TryCreateMutation(
WorkOrderWebhookEnvelope? envelope,
byte[] body,
out WorkOrderWebhookMutation? mutation)
{
mutation = null;
if (envelope == null
|| envelope.SchemaVersion != 1
|| string.IsNullOrWhiteSpace(envelope.DeliveryId)
|| envelope.DeliveryId.Length > 128
|| string.IsNullOrWhiteSpace(envelope.EventType)
|| envelope.EventType.Length > 64
|| envelope.OccurredAt == null
|| !string.Equals(envelope.Source, WorkOrderSourceIdentity.WireSource, StringComparison.Ordinal)
|| envelope.Replay == null
|| envelope.Data == null
|| string.IsNullOrWhiteSpace(envelope.Data.WorkOrderId)
|| envelope.Data.WorkOrderId.Length > 450
|| !KnownEvents.Contains(envelope.EventType))
return false;
var isComment = envelope.EventType == "work_order.comment_added";
if (isComment
&& (string.IsNullOrWhiteSpace(envelope.Data.CommentId)
|| envelope.Data.CommentId.Length > 450
|| string.IsNullOrWhiteSpace(envelope.Data.Text)))
return false;
var updatedAt = envelope.Data.UpdatedAt
?? envelope.Data.IngestedAt
?? envelope.OccurredAt;
if (updatedAt == null)
return false;
var bodyHash = Convert.ToHexString(SHA256.HashData(body)).ToLowerInvariant();
var versionHash = isComment
? WorkOrderExternalVersion.ComputeComment(new WorkOrderExternalVersion.Comment(
envelope.Data.WorkOrderId,
envelope.Data.CommentId!,
envelope.Data.RecordType,
envelope.Data.Commenter,
envelope.Data.Text,
envelope.Data.CreatedAt,
envelope.Data.IngestedAt,
envelope.Data.SourceEmailS3Key))
: WorkOrderExternalVersion.ComputeWorkOrder(new WorkOrderExternalVersion.WorkOrder(
envelope.Data.WorkOrderId,
envelope.Data.WoStatus ?? envelope.Data.Status,
envelope.Data.Description,
envelope.Data.Customer,
envelope.Data.SiteCode,
envelope.Data.Building,
envelope.Data.Address,
envelope.Data.Severity,
envelope.Data.Priority,
envelope.Data.AssignedTo,
envelope.Data.DateReported,
envelope.Data.ScheduledStart,
envelope.Data.DueDate,
envelope.Data.RecordType,
envelope.Data.CreatedAt,
envelope.Data.SourceEmailS3Key));
var statusSource = envelope.Data.WoStatus ?? envelope.Data.Status;
var status = WorkOrderIngestFieldMapper.MapStatus(statusSource);
mutation = new WorkOrderWebhookMutation
{
DeliveryId = envelope.DeliveryId,
EventType = envelope.EventType,
OccurredAt = envelope.OccurredAt.Value,
ProcessedAt = _timeProvider.GetUtcNow(),
BodySha256 = bodyHash,
ExternalWorkOrderId = envelope.Data.WorkOrderId,
WorkerOrderNumber = envelope.Data.WorkOrderId,
Source = WorkOrderSourceIdentity.CanonicalSource,
UpdatedAt = updatedAt.Value,
VersionHash = versionHash,
IsStateEvent = !isComment,
IsCancelled = envelope.EventType == "work_order.cancelled",
Title = envelope.Data.Title,
Description = envelope.Data.Description,
Status = status,
LifecycleStatus = envelope.EventType == "work_order.cancelled"
? LifecycleStatus.Canceled
: LifecycleStatusMapper.FromLegacyStatus(status),
Severity = envelope.Data.Severity,
Priority = envelope.Data.Priority
?? WorkOrderIngestFieldMapper.MapSeverityToPriority(envelope.Data.Severity),
AssignedTo = envelope.Data.AssignedTo,
RecordType = envelope.Data.RecordType,
SourceEmailS3Key = envelope.Data.SourceEmailS3Key,
Customer = envelope.Data.Customer,
SiteCode = envelope.Data.SiteCode,
Building = envelope.Data.Building,
Address = envelope.Data.Address,
DueDate = envelope.Data.DueDate?.UtcDateTime,
DateReported = envelope.Data.DateReported?.UtcDateTime,
ScheduledStart = envelope.Data.ScheduledStart?.UtcDateTime,
CreatedAt = envelope.Data.CreatedAt?.UtcDateTime,
CommentId = isComment ? envelope.Data.CommentId : null,
CommentText = isComment ? envelope.Data.Text : null,
Commenter = isComment ? envelope.Data.Commenter : null,
CommentType = isComment ? envelope.Data.CommentType : null
};
return true;
}
private static WorkOrderWebhookResult Unauthorized() =>
new(WorkOrderWebhookStatus.Unauthorized);
private static readonly HashSet<string> KnownEvents = new(StringComparer.Ordinal)
{
"work_order.created",
"work_order.updated",
"work_order.cancelled",
"work_order.comment_added"
};
private sealed class WorkOrderWebhookEnvelope
{
[JsonPropertyName("schema_version")]
public int? SchemaVersion { get; set; }
[JsonPropertyName("delivery_id")]
public string? DeliveryId { get; set; }
[JsonPropertyName("event_type")]
public string? EventType { get; set; }
[JsonPropertyName("occurred_at")]
public DateTimeOffset? OccurredAt { get; set; }
public string? Source { get; set; }
public bool? Replay { get; set; }
public WorkOrderWebhookData? Data { get; set; }
}
private sealed class WorkOrderWebhookData
{
[JsonPropertyName("work_order_id")]
public string? WorkOrderId { get; set; }
public string? Description { get; set; }
public string? Title { get; set; }
[JsonPropertyName("wo_status")]
public string? WoStatus { get; set; }
public string? Status { get; set; }
public string? Severity { get; set; }
public string? Priority { get; set; }
public string? Customer { get; set; }
[JsonPropertyName("assigned_to")]
public string? AssignedTo { get; set; }
[JsonPropertyName("site_code")]
public string? SiteCode { get; set; }
public string? Building { get; set; }
public string? Address { get; set; }
[JsonPropertyName("due_date")]
public DateTimeOffset? DueDate { get; set; }
[JsonPropertyName("date_reported")]
public DateTimeOffset? DateReported { get; set; }
[JsonPropertyName("scheduled_start")]
public DateTimeOffset? ScheduledStart { get; set; }
[JsonPropertyName("created_at")]
public DateTimeOffset? CreatedAt { get; set; }
[JsonPropertyName("updated_at")]
public DateTimeOffset? UpdatedAt { get; set; }
[JsonPropertyName("ingested_at")]
public DateTimeOffset? IngestedAt { get; set; }
[JsonPropertyName("record_type")]
public string? RecordType { get; set; }
[JsonPropertyName("source_email_s3_key")]
public string? SourceEmailS3Key { get; set; }
[JsonPropertyName("comment_id")]
public string? CommentId { get; set; }
public string? Text { get; set; }
public string? Commenter { get; set; }
[JsonPropertyName("comment_type")]
public string? CommentType { get; set; }
}
}
}