package com.example.slflo import com.fasterxml.jackson.databind.JsonNode import java.math.BigDecimal /** What one run did, for the job log. */ data class RunReport( val pulled: Int = 0, val created: Int = 0, val found: Int = 0, val accepted: Int = 0, val rejected: Int = 0, /** The 401 code that stopped the job, when one did. */ val stoppedBy: String? = null, ) { operator fun plus(other: RunReport) = RunReport( pulled + other.pulled, created + other.created, found + other.found, accepted + other.accepted, rejected + other.rejected, stoppedBy ?: other.stoppedBy, ) } /** * Slflo orders into ERP sales orders. * * 1. Pull `GET /orders` by the version watermark, only orders waiting for the ERP * (`filter[status]=ERP_PENDING`), and stage every line under `(id, version, line)`; the * watermark is saved once a page is staged. * 2. Book each staged order: **find before create** by the Slflo order id in the sales * order's customer reference, so a run that died after the create never makes a second one. * 3. Ack: ACCEPTED with the sales order number, or REJECTED (account on hold). * * [recover] works from `filter[ack]=pending` instead, for when the watermark can't be trusted * (staging purged, cursor moved past unbooked orders). */ class OrderImportJob( private val api: SlfloClient, private val erp: ErpBooks, private val staging: StagingTable = StagingTable(), private val watermarks: Watermarks = Watermarks(), ) { fun run(): RunReport = stoppable { pull(mapOf("filter[status]" to "ERP_PENDING"), advance = true) + book() } fun recover(): RunReport = stoppable { pull(mapOf("filter[ack]" to "pending"), advance = false) + book() } private fun stoppable(block: () -> RunReport): RunReport = try { block() } catch (refused: TokenRefused) { RunReport(stoppedBy = refused.code) } private fun pull(filter: Map, advance: Boolean): RunReport { var after = if (advance) watermarks[STREAM] else 0L var pulled = 0 do { val page = api.get("/orders", filter + mapOf("filter[version_after]" to after, "per_page" to api.perPage)) page.path("data").forEach { staging.stage(it) pulled++ } after = page.path("meta").path("next_version_after").asLong(after) if (advance) watermarks[STREAM] = after } while (page.path("meta").path("has_more").asBoolean()) return RunReport(pulled = pulled) } private fun book(): RunReport = staging.unbooked().fold(RunReport()) { report, order -> report + bookOne(order).also { staging.markBooked(order.path("id").asText(), order.path("version").asLong()) } } private fun bookOne(order: JsonNode): RunReport { val id = order.path("id").asText() val account = order.path("customer").path("external_ref").asText(order.path("customer").path("code").asText()) if (erp.onHold(account)) { ack(order, mapOf("outcome" to "REJECTED", "reason_code" to "CUSTOMER_ON_HOLD", "message" to "Customer $account is on hold in the ERP")) return RunReport(rejected = 1) } val existing = erp.findSalesOrder(id) val number = existing ?: erp.createSalesOrder(id, account, linesOf(order)) ack(order, mapOf("outcome" to "ACCEPTED", "external_ref" to number, "external_number" to number)) return RunReport(created = if (existing == null) 1 else 0, found = if (existing == null) 0 else 1, accepted = 1) } /** The key names the order and the version it answers, so "send again" gets a fresh one. */ private fun ack(order: JsonNode, body: Map) { val id = order.path("id").asText() try { api.post("/orders/$id/ack", body, idempotencyKey = "$id@${order.path("version").asLong()}") } catch (refused: ApiRefused) { // Answered elsewhere, or no longer waiting for the ERP: nothing to do. if (refused.code !in setOf("ALREADY_ACKNOWLEDGED", "ORDER_NOT_AWAITING_ACK")) throw refused } } private fun linesOf(order: JsonNode) = order.path("lines").map { line -> ErpBooks.SalesLine( lineNumber = line.path("line_number").asInt(), item = line.path("product").path("external_ref").asText(line.path("product").path("code").asText()), quantity = BigDecimal(line.path("quantity").asText()), amount = BigDecimal(line.path("line_total").asText()), ) } private companion object { const val STREAM = "orders" } }