// Pull connector: every Slflo order waiting for the ERP becomes one ERP sales order, then an ack. using System.Text.Json.Nodes; namespace Slflo.Samples; public sealed class Staging { public long Watermark { get; set; } // "(id, version, line)" -> the order as pulled, and whether it has been booked public Dictionary Lines { get; set; } = new(); } public sealed class StagedLine { public string OrderId { get; set; } = ""; public long Version { get; set; } public string Document { get; set; } = ""; public bool Booked { get; set; } } public static class PullOrders { public static async Task> RunAsync(SlfloClient api) { var staging = StateFiles.Load("staging-cs.json"); var erp = FakeErp.Open(); var result = new Dictionary { ["pulled"] = 0, ["created"] = 0, ["found"] = 0, ["accepted"] = 0, ["rejected"] = 0 }; // 1. Pull from the watermark, only orders waiting for the ERP, and stage every line. while (true) { var page = await api.GetAsync("/orders", new Dictionary { ["filter[status]"] = "ERP_PENDING", ["filter[version_after]"] = staging.Watermark, ["per_page"] = 100, }); foreach (var order in page["data"]!.AsArray()) { var id = order!["id"]!.GetValue(); var version = order["version"]!.GetValue(); var lines = order["lines"]!.AsArray().Select(l => l!["line_number"]!.GetValue()).DefaultIfEmpty(0); foreach (var line in lines) staging.Lines.TryAdd($"{id}|{version}|{line}", new StagedLine { OrderId = id, Version = version, Document = order.ToJsonString() }); result["pulled"]++; } staging.Watermark = Math.Max(staging.Watermark, page["meta"]!["next_version_after"]!.GetValue()); StateFiles.Save("staging-cs.json", staging); // the watermark moves only once its page is staged if (!page["meta"]!["has_more"]!.GetValue()) break; } // 2. Book each staged order once: find before create, then ack. var unbooked = staging.Lines.Values.Where(l => !l.Booked).GroupBy(l => (l.OrderId, l.Version)).OrderBy(g => g.Key.Version); foreach (var group in unbooked) { var order = JsonNode.Parse(group.First().Document)!; var account = order["customer"]!["external_ref"]?.GetValue() ?? order["customer"]!["code"]!.GetValue(); object answer; if (FakeErp.OnHold(account)) { answer = new { outcome = "REJECTED", reason_code = "CUSTOMER_ON_HOLD", message = $"Customer {account} is on hold in the ERP" }; result["rejected"]++; } else { var number = erp.FindSalesOrder(group.Key.OrderId); if (number is null) { number = erp.CreateSalesOrder(group.Key.OrderId); result["created"]++; } else result["found"]++; answer = new { outcome = "ACCEPTED", external_ref = number, external_number = number }; result["accepted"]++; } try { await api.PostAsync($"/orders/{group.Key.OrderId}/ack", answer, $"{group.Key.OrderId}@{group.Key.Version}"); } catch (ApiErrorException refused) when (refused.Code is "ALREADY_ACKNOWLEDGED" or "ORDER_NOT_AWAITING_ACK") { // Answered elsewhere, or no longer waiting: nothing to do. } foreach (var line in group) line.Booked = true; StateFiles.Save("staging-cs.json", staging); } return result; } }