package com.example.slflo import com.fasterxml.jackson.databind.JsonNode import java.security.MessageDigest /** * Payments or returns into ERP journals, by the version watermark, so a reversal or a bounced * cheque comes back as a new version and is booked again. Acked in one `POST /acks` per run. */ class DocumentImportJob( private val api: SlfloClient, private val erp: ErpBooks, /** `payments` or `returns`. */ private val resource: String, private val staging: StagingTable = StagingTable(), private val watermarks: Watermarks = Watermarks(), ) { private val entity = resource.removeSuffix("s") fun run(): RunReport = try { pull() + bookAndAck() } catch (refused: TokenRefused) { RunReport(stoppedBy = refused.code) } private fun pull(): RunReport { var pulled = 0 do { val page = api.get("/$resource", mapOf("filter[version_after]" to watermarks[resource], "per_page" to api.perPage)) page.path("data").forEach { staging.stage(it) pulled++ } watermarks[resource] = page.path("meta").path("next_version_after").asLong(watermarks[resource]) } while (page.path("meta").path("has_more").asBoolean()) return RunReport(pulled = pulled) } private fun bookAndAck(): RunReport { val documents = staging.unbooked() if (documents.isEmpty()) return RunReport() val items = documents.map { document -> val amount = document.path("amount").takeUnless { it.isMissingNode }?.asText() ?: document.path("total").asText() val journal = erp.postJournal(entity, document.path("id").asText(), document.path("version").asLong(), amount) mapOf("id" to document.path("id").asText(), "outcome" to "ACCEPTED", "external_ref" to journal, "external_number" to journal) } val answer = api.post("/acks", mapOf("entity" to entity, "items" to items), keyOf(entity, documents)) documents.forEach { staging.markBooked(it.path("id").asText(), it.path("version").asLong()) } return RunReport(accepted = answer.path("meta").path("applied").asInt() + answer.path("meta").path("unchanged").asInt()) } } /** * Invoices by the gap-free issued counter (`filter[seq_after]`): "this is the last one I have, * what are the next ones" — one integer per stream. */ class InvoiceSeqJob(private val api: SlfloClient, private val erp: ErpBooks, var lastSeq: Long = 0) { /** Reads, books and acks what was issued since [lastSeq]; returns how many. */ fun run(): Int { var read = 0 do { val page = api.get("/invoices", mapOf("filter[seq_after]" to lastSeq, "per_page" to api.perPage)) val invoices = page.path("data").toList() invoices.forEach { invoice -> check(invoice.path("seq").asLong() == lastSeq + 1) { "Gap after seq $lastSeq: got ${invoice.path("seq")}" } lastSeq = invoice.path("seq").asLong() } val items = invoices.map { invoice -> val journal = erp.postJournal("invoice", invoice.path("id").asText(), invoice.path("version").asLong(), invoice.path("totals").path("total").asText()) mapOf("id" to invoice.path("id").asText(), "outcome" to "ACCEPTED", "external_ref" to journal) } if (items.isNotEmpty()) api.post("/acks", mapOf("entity" to "invoice", "items" to items), keyOf("invoice", invoices)) read += invoices.size } while (page.path("meta").path("has_more").asBoolean()) return read } } /** One key per set of (id, version): a retry replays it, a later change gets a new one. */ private fun keyOf(entity: String, documents: List): String { val answered = documents.joinToString(",") { "${it.path("id").asText()}@${it.path("version").asLong()}" } val digest = MessageDigest.getInstance("SHA-256").digest(answered.toByteArray()).joinToString("") { "%02x".format(it) } return "$entity-acks:${digest.take(KEY_DIGEST_LENGTH)}" } private const val KEY_DIGEST_LENGTH = 40