Add DynamoDB sync bridge for workorder-ingest pipeline

- Add ExternalWorkOrderId, Customer, SiteCode, Building, Severity,
  DateReported, ScheduledStart, SourceEmailS3Key to WorkOrder model
- Add Commenter, RecordType, ExternalCommentId to Comments model
- Add AWSSDK.DynamoDBv2 package
- Create SyncController with endpoints:
  - POST /api/Sync/WorkOrders — scans DynamoDB WorkOrders table, upserts into SQL
  - POST /api/Sync/Comments — scans DynamoDB WorkOrderComments, creates Comments
  - POST /api/Sync/All — runs both
- Auto-creates Locations from site_code/building/address during sync
- Maps DynamoDB status values to SHOC status strings
- Maps severity to SHOC priority format (Sev N)
- Deduplicates on ExternalWorkOrderId / ExternalCommentId
This commit is contained in:
Adam Moussa 2026-04-16 17:29:52 -04:00
parent 93fa76a5aa
commit 6f8f67735a
7 changed files with 2657 additions and 0 deletions

View file

@ -7,6 +7,7 @@
</PropertyGroup>
<ItemGroup>
<PackageReference Include="AWSSDK.DynamoDBv2" Version="4.0.17.9" />
<PackageReference Include="Microsoft.AspNetCore.Authentication.JwtBearer" Version="8.0.8" />
<PackageReference Include="Microsoft.EntityFrameworkCore.Design" Version="8.0.8">
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>

View file

@ -0,0 +1,273 @@
using Amazon.DynamoDBv2;
using Amazon.DynamoDBv2.Model;
using Data.SeaHavenIndustries;
using Microsoft.AspNetCore.Authorization;
using Microsoft.AspNetCore.Mvc;
using Microsoft.EntityFrameworkCore;
namespace Api.SeaHavenIndustries.Controllers
{
[Authorize]
[ApiController]
[Route("api/[controller]")]
public class SyncController : Controller
{
private readonly ApplicationDbContext _db;
private readonly AmazonDynamoDBClient _dynamo;
private const string WO_TABLE = "WorkOrders";
private const string COMMENTS_TABLE = "WorkOrderComments";
public SyncController(ApplicationDbContext db)
{
_db = db;
_dynamo = new AmazonDynamoDBClient(Amazon.RegionEndpoint.USEast1);
}
[HttpPost("WorkOrders")]
public async Task<IActionResult> SyncWorkOrders()
{
var synced = 0;
var created = 0;
var updated = 0;
var locationsCreated = 0;
Dictionary<string, AttributeValue>? lastKey = null;
do
{
var request = new ScanRequest
{
TableName = WO_TABLE,
Limit = 100,
ExclusiveStartKey = lastKey
};
var response = await _dynamo.ScanAsync(request);
foreach (var item in response.Items)
{
var externalId = GetString(item, "work_order_id");
if (string.IsNullOrEmpty(externalId)) continue;
var existing = await _db.workOrders
.FirstOrDefaultAsync(w => w.ExternalWorkOrderId == externalId);
var locationId = await ResolveLocationId(
GetString(item, "site_code"),
GetString(item, "building"),
GetString(item, "address"));
if (locationId.HasValue && !locationsCreated.Equals(0))
locationsCreated++;
if (existing == null)
{
var wo = new WorkOrder
{
ExternalWorkOrderId = externalId,
WorkerOrderNumber = externalId,
WorkerOrderTitle = GetString(item, "description"),
Description = GetString(item, "description"),
Status = MapStatus(GetString(item, "wo_status")),
Priority = MapSeverityToPriority(GetString(item, "severity")),
Severity = GetString(item, "severity"),
Customer = GetString(item, "customer"),
SiteCode = GetString(item, "site_code"),
Building = GetString(item, "building"),
LocationId = locationId,
DueDate = ParseDate(GetString(item, "due_date")),
DateReported = ParseDate(GetString(item, "date_reported")),
ScheduledStart = ParseDate(GetString(item, "scheduled_start")),
AssignTo = null,
SourceEmailS3Key = GetString(item, "source_email_s3_key"),
CreatedDate = ParseDate(GetString(item, "created_at")) ?? DateTime.UtcNow,
istemplate = false
};
_db.workOrders.Add(wo);
created++;
}
else
{
existing.WorkerOrderTitle = GetString(item, "description") ?? existing.WorkerOrderTitle;
existing.Description = GetString(item, "description") ?? existing.Description;
existing.Status = MapStatus(GetString(item, "wo_status")) ?? existing.Status;
existing.Priority = MapSeverityToPriority(GetString(item, "severity")) ?? existing.Priority;
existing.Severity = GetString(item, "severity") ?? existing.Severity;
existing.SiteCode = GetString(item, "site_code") ?? existing.SiteCode;
existing.Building = GetString(item, "building") ?? existing.Building;
existing.LocationId = locationId ?? existing.LocationId;
existing.DueDate = ParseDate(GetString(item, "due_date")) ?? existing.DueDate;
existing.DateReported = ParseDate(GetString(item, "date_reported")) ?? existing.DateReported;
existing.ScheduledStart = ParseDate(GetString(item, "scheduled_start")) ?? existing.ScheduledStart;
existing.SourceEmailS3Key = GetString(item, "source_email_s3_key") ?? existing.SourceEmailS3Key;
updated++;
}
synced++;
if (synced % 100 == 0)
await _db.SaveChangesAsync();
}
lastKey = response.LastEvaluatedKey;
} while (lastKey != null && lastKey.Count > 0);
await _db.SaveChangesAsync();
return Ok(new
{
synced,
created,
updated,
locationsCreated = await _db.Locations.CountAsync()
});
}
[HttpPost("Comments")]
public async Task<IActionResult> SyncComments()
{
var synced = 0;
var created = 0;
var skipped = 0;
Dictionary<string, AttributeValue>? lastKey = null;
do
{
var request = new ScanRequest
{
TableName = COMMENTS_TABLE,
Limit = 100,
ExclusiveStartKey = lastKey
};
var response = await _dynamo.ScanAsync(request);
foreach (var item in response.Items)
{
var externalCommentId = GetString(item, "comment_id");
var externalWoId = GetString(item, "work_order_id");
if (string.IsNullOrEmpty(externalCommentId) || string.IsNullOrEmpty(externalWoId))
continue;
var alreadyExists = await _db.Comments
.AnyAsync(c => c.ExternalCommentId == externalCommentId);
if (alreadyExists)
{
skipped++;
synced++;
continue;
}
var workOrder = await _db.workOrders
.FirstOrDefaultAsync(w => w.ExternalWorkOrderId == externalWoId);
if (workOrder == null)
{
skipped++;
continue;
}
var comment = new Comments
{
ExternalCommentId = externalCommentId,
WorkerOrderId = workOrder.Id,
Commenttext = GetString(item, "text"),
Commenter = GetString(item, "commenter"),
RecordType = GetString(item, "record_type"),
CreatedDate = ParseDate(GetString(item, "created_at")) ?? DateTime.UtcNow,
};
_db.Comments.Add(comment);
created++;
synced++;
if (synced % 100 == 0)
await _db.SaveChangesAsync();
}
lastKey = response.LastEvaluatedKey;
} while (lastKey != null && lastKey.Count > 0);
await _db.SaveChangesAsync();
return Ok(new { synced, created, skipped });
}
[HttpPost("All")]
public async Task<IActionResult> SyncAll()
{
var woResult = await SyncWorkOrders() as OkObjectResult;
var commentResult = await SyncComments() as OkObjectResult;
return Ok(new
{
workOrders = woResult?.Value,
comments = commentResult?.Value
});
}
private async Task<int?> ResolveLocationId(string? siteCode, string? building, string? address)
{
if (string.IsNullOrWhiteSpace(siteCode) && string.IsNullOrWhiteSpace(building))
return null;
var matchCode = siteCode ?? building;
var existing = await _db.Locations
.FirstOrDefaultAsync(l => l.Name == matchCode || l.Title == matchCode);
if (existing != null)
return existing.Id;
var location = new Locations
{
Name = matchCode,
Title = building,
Address = address,
Status = "Active"
};
_db.Locations.Add(location);
await _db.SaveChangesAsync();
return location.Id;
}
private static string? GetString(Dictionary<string, AttributeValue> item, string key)
{
if (item.TryGetValue(key, out var val) && val.S != null)
return val.S;
return null;
}
private static DateTime? ParseDate(string? value)
{
if (string.IsNullOrWhiteSpace(value)) return null;
if (DateTime.TryParse(value, out var dt)) return dt;
return null;
}
private static string? MapStatus(string? dynamoStatus)
{
return dynamoStatus switch
{
"new" => "Open",
"assigned" => "Open",
"in_progress" => "In Progress",
"on_hold" => "On Hold",
"completed" => "Done",
"cancelled" => "Cancelled",
"unknown" => "Open",
_ => "Open"
};
}
private static string? MapSeverityToPriority(string? severity)
{
if (string.IsNullOrWhiteSpace(severity)) return null;
return $"Sev {severity}";
}
}
}

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,129 @@
using System;
using Microsoft.EntityFrameworkCore.Migrations;
#nullable disable
namespace Data.SeaHavenIndustries.Migrations
{
/// <inheritdoc />
public partial class AddExternalSyncFields : Migration
{
/// <inheritdoc />
protected override void Up(MigrationBuilder migrationBuilder)
{
migrationBuilder.AddColumn<string>(
name: "Building",
table: "workOrders",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "Customer",
table: "workOrders",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<DateTime>(
name: "DateReported",
table: "workOrders",
type: "datetime2",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "ExternalWorkOrderId",
table: "workOrders",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<DateTime>(
name: "ScheduledStart",
table: "workOrders",
type: "datetime2",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "Severity",
table: "workOrders",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "SiteCode",
table: "workOrders",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "SourceEmailS3Key",
table: "workOrders",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "Commenter",
table: "Comments",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "ExternalCommentId",
table: "Comments",
type: "nvarchar(max)",
nullable: true);
migrationBuilder.AddColumn<string>(
name: "RecordType",
table: "Comments",
type: "nvarchar(max)",
nullable: true);
}
/// <inheritdoc />
protected override void Down(MigrationBuilder migrationBuilder)
{
migrationBuilder.DropColumn(
name: "Building",
table: "workOrders");
migrationBuilder.DropColumn(
name: "Customer",
table: "workOrders");
migrationBuilder.DropColumn(
name: "DateReported",
table: "workOrders");
migrationBuilder.DropColumn(
name: "ExternalWorkOrderId",
table: "workOrders");
migrationBuilder.DropColumn(
name: "ScheduledStart",
table: "workOrders");
migrationBuilder.DropColumn(
name: "Severity",
table: "workOrders");
migrationBuilder.DropColumn(
name: "SiteCode",
table: "workOrders");
migrationBuilder.DropColumn(
name: "SourceEmailS3Key",
table: "workOrders");
migrationBuilder.DropColumn(
name: "Commenter",
table: "Comments");
migrationBuilder.DropColumn(
name: "ExternalCommentId",
table: "Comments");
migrationBuilder.DropColumn(
name: "RecordType",
table: "Comments");
}
}
}

View file

@ -376,6 +376,9 @@ namespace Data.SeaHavenIndustries.Migrations
SqlServerPropertyBuilderExtensions.UseIdentityColumn(b.Property<int>("Id"));
b.Property<string>("Commenter")
.HasColumnType("nvarchar(max)");
b.Property<string>("Commenttext")
.HasColumnType("nvarchar(max)");
@ -391,6 +394,9 @@ namespace Data.SeaHavenIndustries.Migrations
b.Property<string>("Documents")
.HasColumnType("nvarchar(max)");
b.Property<string>("ExternalCommentId")
.HasColumnType("nvarchar(max)");
b.Property<bool?>("IsDeleted")
.HasColumnType("bit");
@ -400,6 +406,9 @@ namespace Data.SeaHavenIndustries.Migrations
b.Property<int?>("LastModifierUserId")
.HasColumnType("int");
b.Property<string>("RecordType")
.HasColumnType("nvarchar(max)");
b.Property<string>("UserId")
.HasColumnType("nvarchar(450)");
@ -1393,9 +1402,18 @@ namespace Data.SeaHavenIndustries.Migrations
b.Property<string>("BeforPhotoIssue")
.HasColumnType("nvarchar(max)");
b.Property<string>("Building")
.HasColumnType("nvarchar(max)");
b.Property<DateTime?>("CreatedDate")
.HasColumnType("datetime2");
b.Property<string>("Customer")
.HasColumnType("nvarchar(max)");
b.Property<DateTime?>("DateReported")
.HasColumnType("datetime2");
b.Property<string>("DeleterUserId")
.HasColumnType("nvarchar(max)");
@ -1408,6 +1426,9 @@ namespace Data.SeaHavenIndustries.Migrations
b.Property<DateTime?>("DueDate")
.HasColumnType("datetime2");
b.Property<string>("ExternalWorkOrderId")
.HasColumnType("nvarchar(max)");
b.Property<bool?>("IsDeleted")
.HasColumnType("bit");
@ -1426,6 +1447,12 @@ namespace Data.SeaHavenIndustries.Migrations
b.Property<string>("Priority")
.HasColumnType("nvarchar(max)");
b.Property<DateTime?>("ScheduledStart")
.HasColumnType("datetime2");
b.Property<string>("Severity")
.HasColumnType("nvarchar(max)");
b.Property<string>("SignOffAttachment")
.HasColumnType("nvarchar(max)");
@ -1435,6 +1462,12 @@ namespace Data.SeaHavenIndustries.Migrations
b.Property<string>("SignOffSignature")
.HasColumnType("nvarchar(max)");
b.Property<string>("SiteCode")
.HasColumnType("nvarchar(max)");
b.Property<string>("SourceEmailS3Key")
.HasColumnType("nvarchar(max)");
b.Property<string>("Status")
.HasColumnType("nvarchar(max)");

View file

@ -7,6 +7,9 @@ namespace Data.SeaHavenIndustries
public class Comments : FullAuditEntity
{
public string? Commenttext { get; set; }
public string? Commenter { get; set; }
public string? RecordType { get; set; }
public string? ExternalCommentId { get; set; }
public string? Documents { get; set; }
public string? UserId { get; set; }
[ForeignKey(nameof(UserId))]

View file

@ -6,9 +6,17 @@ namespace Data.SeaHavenIndustries
{
public class WorkOrder : FullAuditEntity
{
public string? ExternalWorkOrderId { get; set; }
public string? WorkerOrderNumber { get; set; }
public string? WorkerOrderTitle { get; set; }
public string? Description { get; set; }
public string? Customer { get; set; }
public string? SiteCode { get; set; }
public string? Building { get; set; }
public string? Severity { get; set; }
public DateTime? DateReported { get; set; }
public DateTime? ScheduledStart { get; set; }
public string? SourceEmailS3Key { get; set; }
public bool? istemplate { get; set; }
public string? PO { get; set; }
public string? TT { get; set; }