//! Pushing issued invoices and their payments to ERPNext. //! //! The database mutex is never held across an `await`: each step takes the lock, reads or writes, and //! releases it before the next network call. Every outcome after the pre-flight checks is recorded in //! `erpnext_sync` (also failures), and an existing remote document is never overwritten. use super::client::{ErpClient, Upload}; use super::config::{self, ErpnextConfig, NamingMode}; use super::errors::{ErpError, ErrorKind}; use super::mapping::{ self, build_address, build_customer, build_sales_invoice, paise_to_decimal, remarks_marker, InvoiceContext, Vendor, SALES_INVOICE_V2, }; use crate::commands::archive::read_archive_impl; use crate::commands::invoice::get_invoice_impl; use crate::gst; use crate::integrations::{ InvoiceSink, PaymentRequest, PushRequest, PushedInvoice, PushedPayment, SyncStatus, }; use crate::models::{Client, Invoice}; use percent_encoding::{utf8_percent_encode, AsciiSet, NON_ALPHANUMERIC}; use rusqlite::{params, Connection, OptionalExtension}; use serde::Serialize; use serde_json::{json, Map, Value}; use sha2::{Digest, Sha256}; use std::collections::HashMap; use std::path::Path; use std::sync::Mutex; type Db = Mutex; const DOCTYPE_INVOICE: &str = "Sales Invoice"; const GET_PAYMENT_ENTRY: &str = "erpnext.accounts.doctype.payment_entry.payment_entry.get_payment_entry"; fn pre(message: impl Into) -> ErpError { ErpError::new(ErrorKind::Precondition, message) } fn with_db(db: &Db, f: impl FnOnce(&mut Connection) -> Result) -> Result { let mut conn = db.lock().map_err(|e| pre(format!("The database is busy: {e}")))?; f(&mut conn).map_err(|e| pre(format!("Could not read or save the sync state: {e}"))) } fn sha256_hex(bytes: &[u8]) -> String { format!("{:x}", Sha256::digest(bytes)) } // ---- results ---- /// One invoice's push outcome. A failure is a normal result (`ok: false`), so a bulk push can carry on. #[derive(Debug, Clone, Serialize, PartialEq)] #[serde(rename_all = "camelCase")] pub struct PushResult { pub invoice_id: i64, pub number: String, pub ok: bool, /// `synced`, `error` or `conflict`; `refused` when nothing was sent and nothing was recorded. pub status: String, pub remote_name: String, /// 0 draft, 1 submitted. pub remote_docstatus: i64, /// The document was created by this call (not found already there). pub created: bool, /// Nothing needed doing: already synced with the same payload. pub no_op: bool, /// The archived PDF is attached on the remote document. pub attached: bool, pub error: Option, pub error_kind: Option, pub warnings: Vec, } impl PushResult { fn refused(invoice_id: i64, number: &str, e: ErpError) -> Self { PushResult { invoice_id, number: number.to_string(), ok: false, status: "refused".into(), remote_name: String::new(), remote_docstatus: 0, created: false, no_op: false, attached: false, error: Some(e.to_string()), error_kind: Some(e.kind), warnings: Vec::new(), } } } #[derive(Debug, Clone, Serialize, PartialEq)] #[serde(rename_all = "camelCase")] pub struct PaymentPushResult { pub payment_id: i64, pub invoice_id: i64, pub ok: bool, pub entry_name: Option, /// The payment had already been sent; nothing was posted. pub already_synced: bool, pub error: Option, pub error_kind: Option, } // ---- local state ---- #[derive(Debug, Clone, Default)] struct SyncRow { remote_name: String, remote_docstatus: i64, status: String, last_error: String, payload_hash: String, synced_at: Option, attachment_sha256: String, } fn load_sync(conn: &Connection, invoice_id: i64) -> Result, String> { conn.query_row( "SELECT remote_name, remote_docstatus, status, last_error, payload_hash, synced_at, attachment_sha256 FROM erpnext_sync WHERE invoice_id = ?1", params![invoice_id], |r| { Ok(SyncRow { remote_name: r.get(0)?, remote_docstatus: r.get(1)?, status: r.get(2)?, last_error: r.get(3)?, payload_hash: r.get(4)?, synced_at: r.get(5)?, attachment_sha256: r.get(6)?, }) }, ) .optional() .map_err(|e| e.to_string()) } /// One transaction per write, so the row is never half updated. fn write_sync(conn: &mut Connection, invoice_id: i64, row: &SyncRow) -> Result<(), String> { let tx = conn.transaction().map_err(|e| e.to_string())?; tx.execute( "INSERT INTO erpnext_sync (invoice_id, remote_name, remote_docstatus, status, last_error, payload_hash, synced_at, attachment_sha256) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) ON CONFLICT(invoice_id) DO UPDATE SET remote_name = excluded.remote_name, remote_docstatus = excluded.remote_docstatus, status = excluded.status, last_error = excluded.last_error, payload_hash = excluded.payload_hash, synced_at = excluded.synced_at, attachment_sha256 = excluded.attachment_sha256", params![ invoice_id, row.remote_name, row.remote_docstatus, row.status, row.last_error, row.payload_hash, row.synced_at, row.attachment_sha256 ], ) .map_err(|e| e.to_string())?; tx.commit().map_err(|e| e.to_string()) } fn status_of(invoice_id: i64, row: Option) -> SyncStatus { match row { None => SyncStatus::none(invoice_id), Some(r) => SyncStatus { invoice_id, status: r.status, remote_name: r.remote_name, remote_docstatus: r.remote_docstatus, last_error: r.last_error, synced_at: r.synced_at, attached: !r.attachment_sha256.is_empty(), }, } } pub fn sync_status(conn: &Connection, invoice_id: i64) -> Result { Ok(status_of(invoice_id, load_sync(conn, invoice_id)?)) } /// Every invoice that has a sync row, for the History list. pub fn sync_statuses(conn: &Connection) -> Result, String> { let mut stmt = conn .prepare("SELECT invoice_id FROM erpnext_sync ORDER BY invoice_id") .map_err(|e| e.to_string())?; let ids = stmt .query_map([], |r| r.get::<_, i64>(0)) .map_err(|e| e.to_string())? .collect::>>() .map_err(|e| e.to_string())?; ids.into_iter().map(|id| sync_status(conn, id)).collect() } /// `/app/sales-invoice/`, with the name URL-encoded (a mirrored number contains a slash). pub fn open_url(conn: &Connection, invoice_id: i64) -> Result { const KEEP: &AsciiSet = &NON_ALPHANUMERIC.remove(b'-').remove(b'_').remove(b'.').remove(b'~'); let row = load_sync(conn, invoice_id)?; let name = row.map(|r| r.remote_name).unwrap_or_default(); if name.is_empty() { return Err("This invoice has not been sent to ERPNext yet.".into()); } let cfg = config::load(conn)?; let base = super::client::normalize_base_url(&cfg.base_url).map_err(|e| e.to_string())?; Ok(format!("{base}/app/sales-invoice/{}", utf8_percent_encode(&name, KEEP))) } // ---- loading ---- struct ClientRow { /// `None` when the invoice has no saved client (the details come from the invoice itself). id: Option, client: Client, customer: Option, address: Option, } struct Pdf { sha256: String, bytes: Vec, } struct Loaded { cfg: ErpnextConfig, invoice: Invoice, vendor: Vendor, client: ClientRow, item_codes: Vec>, sync: Option, pdf: Option, /// Why there is no PDF to attach although attaching is switched on. pdf_warning: Option, india_compliance: bool, } fn load_client(conn: &Connection, inv: &Invoice) -> Result { if let Some(id) = inv.client_id { let row = conn .query_row( "SELECT name, address, gstin, state_code, po_number, created_at, address_line1, address_line2, city, pincode, gst_category, default_notes, payment_terms_days, erpnext_customer, erpnext_address FROM clients WHERE id = ?1", params![id], |r| { let customer: Option = r.get(13)?; let address: Option = r.get(14)?; Ok(ClientRow { id: Some(id), client: Client { id: Some(id), name: r.get(0)?, address: r.get(1)?, gstin: r.get(2)?, state_code: r.get(3)?, po_number: r.get(4)?, created_at: r.get(5)?, address_line1: r.get(6)?, address_line2: r.get(7)?, city: r.get(8)?, pincode: r.get(9)?, gst_category: r.get(10)?, default_notes: r.get(11)?, payment_terms_days: r.get(12)?, invoice_count: 0, }, customer: customer.filter(|c| !c.trim().is_empty()), address: address.filter(|a| !a.trim().is_empty()), }) }, ) .optional() .map_err(|e| e.to_string())?; if let Some(row) = row { return Ok(row); } } // No saved client: use what the invoice froze. let gstin = inv.client_gstin.trim(); let has_gstin = !gstin.is_empty() && !gstin.eq_ignore_ascii_case("NA"); Ok(ClientRow { id: None, client: Client { id: None, name: inv.client_name.clone(), address: inv.client_address.clone(), gstin: inv.client_gstin.clone(), state_code: String::new(), po_number: String::new(), created_at: String::new(), address_line1: String::new(), address_line2: String::new(), city: String::new(), pincode: String::new(), gst_category: if has_gstin { "registered_regular" } else { "unregistered" }.into(), default_notes: String::new(), payment_terms_days: None, invoice_count: 0, }, customer: None, address: None, }) } /// Invoice rows are matched to item presets by description (case-insensitive) to find their ERPNext item code. fn load_item_codes(conn: &Connection, inv: &Invoice) -> Result>, String> { let mut stmt = conn .prepare( "SELECT description, erpnext_item_code FROM item_presets WHERE erpnext_item_code IS NOT NULL AND trim(erpnext_item_code) <> '' ORDER BY id DESC", ) .map_err(|e| e.to_string())?; let mut by_desc: HashMap = HashMap::new(); let rows = stmt .query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))) .map_err(|e| e.to_string())?; for row in rows { let (desc, code) = row.map_err(|e| e.to_string())?; // ORDER BY id DESC: the oldest preset wins on a repeated description. by_desc.insert(desc.trim().to_lowercase(), code.trim().to_string()); } Ok(inv.items.iter().map(|i| by_desc.get(&i.description.trim().to_lowercase()).cloned()).collect()) } fn load_for_push(db: &Db, local_dir: &Path, invoice_id: i64) -> Result { let conn = db .lock() .map_err(|e| (String::new(), pre(format!("The database is busy: {e}"))))?; let invoice = get_invoice_impl(&conn, invoice_id).map_err(|e| (String::new(), pre(e)))?; let number = invoice.number.clone(); let fail = |e: ErpError| (number.clone(), e); match invoice.status.as_str() { "issued" => {} "cancelled" => { return Err(fail(pre(format!( "Invoice {number} is cancelled, so it is not sent to ERPNext. Cancelled invoices are never pushed." )))) } other => { return Err(fail(pre(format!( "Invoice {number} is a {other}; only issued invoices are sent to ERPNext." )))) } } let cfg = config::load(&conn).map_err(|e| fail(pre(e)))?; let vendor = Vendor::from_snapshot(&invoice.vendor_snapshot) .ok_or_else(|| fail(pre(format!("Invoice {number} has no supplier details recorded, so it cannot be sent."))))?; let client = load_client(&conn, &invoice).map_err(|e| fail(pre(e)))?; let item_codes = load_item_codes(&conn, &invoice).map_err(|e| fail(pre(e)))?; let sync = load_sync(&conn, invoice_id).map_err(|e| fail(pre(e)))?; let (mut pdf, mut pdf_warning) = (None, None); if cfg.attach_pdf { match invoice.archived_pdf_sha256.as_deref() { None => { pdf_warning = Some(format!( "Invoice {number} is not archived, so it was sent without its PDF. Export it once, then push again to attach it." )) } Some(sha) => match read_archive_impl(&conn, local_dir, invoice_id) { Ok(bytes) => pdf = Some(Pdf { sha256: sha.to_string(), bytes }), Err(e) => pdf_warning = Some(format!("The PDF could not be attached: {e}.")), }, } } let india_compliance = serde_json::from_str::(&cfg.last_detect_result) .ok() .and_then(|v| v.get("indiaCompliance").and_then(Value::as_bool)) .unwrap_or(false); Ok(Loaded { cfg, invoice, vendor, client, item_codes, sync, pdf, pdf_warning, india_compliance }) } // ---- address checks (local, before anything is sent) ---- /// First two PIN digits that belong to each GST state code. Deliberately a little generous at the borders /// (postal circles and states do not line up exactly); it only catches obvious mismatches. fn pin_prefixes(state_code: &str) -> Option<&'static [&'static str]> { Some(match state_code { "01" => &["18", "19"], "02" => &["17"], "03" => &["14", "15", "16"], "04" => &["16"], "05" => &["24", "25", "26"], "06" => &["12", "13"], "07" => &["11"], "08" => &["30", "31", "32", "33", "34"], "09" => &["20", "21", "22", "23", "24", "25", "26", "27", "28"], "10" => &["80", "81", "82", "83", "84", "85"], "11" => &["73", "75"], "12" | "13" | "14" | "15" | "16" | "17" => &["79"], "18" => &["78"], "19" => &["70", "71", "72", "73", "74"], "20" => &["81", "82", "83"], "21" => &["75", "76", "77"], "22" => &["49"], "23" => &["45", "46", "47", "48"], "24" => &["36", "37", "38", "39"], "26" => &["39"], "27" => &["40", "41", "42", "43", "44"], "29" => &["56", "57", "58", "59"], "30" => &["40"], "31" => &["68"], "32" => &["67", "68", "69"], "33" => &["60", "61", "62", "63", "64"], "34" => &["53", "60", "67"], "35" => &["74"], "36" => &["50", "51", "52"], "37" => &["50", "51", "52", "53"], "38" => &["19"], _ => return None, }) } fn gstin_applies(client: &Client) -> bool { let g = client.gstin.trim(); !g.is_empty() && !g.eq_ignore_ascii_case("NA") && matches!(client.gst_category.as_str(), "registered_regular" | "composition" | "sez") } /// The checks India Compliance would make on an Address, done here so the error is readable and local. pub fn validate_address(client: &Client) -> Result<(), String> { let state_code = client.state_code.trim(); let state_name = mapping::state_name(state_code).unwrap_or("that state"); let pin = client.pincode.trim(); if !pin.is_empty() { if pin.len() != 6 || !pin.bytes().all(|b| b.is_ascii_digit()) || pin.starts_with('0') { return Err(format!("{}: the PIN code \"{pin}\" is not a valid 6-digit PIN.", client.name.trim())); } if let Some(allowed) = pin_prefixes(state_code) { if !allowed.iter().any(|p| pin.starts_with(p)) { return Err(format!( "{}: the PIN code {pin} does not belong to {state_name}. Fix the client's address before sending.", client.name.trim() )); } } } if gstin_applies(client) && !state_code.is_empty() { let gstin = client.gstin.trim().to_ascii_uppercase(); if gstin.len() >= 2 && &gstin[0..2] != state_code { return Err(format!( "{}: the GSTIN {gstin} starts with state code {}, but the address state is {state_name} ({state_code}).", client.name.trim(), &gstin[0..2] )); } } Ok(()) } // ---- remote steps ---- fn doc_name(response: &Value) -> Option { response .get("data") .and_then(|d| d.get("name")) .and_then(Value::as_str) .map(str::trim) .filter(|n| !n.is_empty()) .map(str::to_string) } fn doc_docstatus(doc: &Value) -> i64 { doc.get("docstatus").and_then(Value::as_i64).unwrap_or(0) } fn doc_total_paise(doc: &Value) -> Option { let v = doc.get("grand_total")?; let n = v.as_f64().or_else(|| v.as_str().and_then(|s| s.trim().parse().ok()))?; Some(gst::rupees_to_paise(n)) } async fn ensure_customer(db: &Db, http: &ErpClient, l: &Loaded) -> Result { if let Some(c) = &l.client.customer { return Ok(c.clone()); } let client = &l.client.client; let name = client.name.trim(); let mut found: Option = None; if l.india_compliance && gstin_applies(client) { let rows = http .list_resource( "Customer", &["name"], json!([["gstin", "=", client.gstin.trim().to_ascii_uppercase()]]), "creation asc", ) .await?; found = rows.first().and_then(|r| r.get("name")).and_then(Value::as_str).map(str::to_string); } if found.is_none() { let rows = http .list_resource("Customer", &["name"], json!([["customer_name", "=", name]]), "creation asc") .await?; found = rows.first().and_then(|r| r.get("name")).and_then(Value::as_str).map(str::to_string); } let customer = match found { Some(c) => c, None => { if !l.cfg.create_missing_customers { return Err(pre(format!( "The customer \"{name}\" does not exist in ERPNext, and creating customers is switched off in the ERPNext settings." ))); } let req = build_customer(client, &l.cfg, l.india_compliance).map_err(pre)?; let resp = http.post(req.path, &req.body, req.idempotent).await?; // A duplicate name comes back as "X - 1": always use what the server returned. doc_name(&resp).ok_or_else(|| ErpError::protocol("ERPNext did not return the new customer's name."))? } }; if let Some(id) = l.client.id { with_db(db, |c| { c.execute("UPDATE clients SET erpnext_customer = ?1 WHERE id = ?2", params![customer, id]) .map(|_| ()) .map_err(|e| e.to_string()) })?; } Ok(customer) } /// Returns the address name and, when none could be made, a warning. async fn ensure_address( db: &Db, http: &ErpClient, l: &Loaded, customer: &str, ) -> Result<(Option, Option), ErpError> { if let Some(a) = &l.client.address { return Ok((Some(a.clone()), None)); } let client = &l.client.client; let Some(client_id) = l.client.id else { return Ok(( None, Some("The client is not saved, so the invoice was sent without a customer address.".into()), )); }; let has_address = [&client.address_line1, &client.address, &client.city, &client.pincode] .iter() .any(|s| !s.trim().is_empty()); if !has_address { return Ok(( None, Some(format!("{} has no address saved, so the invoice was sent without a customer address.", client.name.trim())), )); } validate_address(client).map_err(pre)?; let req = build_address(client, customer, l.india_compliance).map_err(pre)?; let resp = http.post(req.path, &req.body, req.idempotent).await?; let name = doc_name(&resp).ok_or_else(|| ErpError::protocol("ERPNext did not return the new address's name."))?; with_db(db, |c| { c.execute("UPDATE clients SET erpnext_address = ?1 WHERE id = ?2", params![name, client_id]) .map(|_| ()) .map_err(|e| e.to_string()) })?; Ok((Some(name), None)) } struct RemoteDoc { name: String, docstatus: i64, created: bool, } /// An existing remote document is only accepted when its total equals Voiced's; anything else is a /// conflict that is reported and never overwritten. fn accept_existing(name: &str, doc: &Value, inv: &Invoice) -> Result { let docstatus = doc_docstatus(doc); if docstatus == 2 { return Err(ErpError::new( ErrorKind::Conflict, format!("ERPNext already has {name} for invoice {}, but it was cancelled there. Resolve it in ERPNext first.", inv.number), )); } let ours = gst::rupees_to_paise(inv.total); match doc_total_paise(doc) { Some(theirs) if theirs == ours => Ok(RemoteDoc { name: name.to_string(), docstatus, created: false }), Some(theirs) => Err(ErpError::new( ErrorKind::Conflict, format!( "ERPNext already has {name} for invoice {}, but its total is {} and Voiced's is {}. Voiced does not overwrite it: fix or delete the ERPNext document, then push again.", inv.number, paise_to_decimal(theirs), paise_to_decimal(ours) ), )), None => Err(ErpError::new( ErrorKind::Conflict, format!("ERPNext already has {name} for invoice {}, but its total could not be read to compare.", inv.number), )), } } async fn find_by_remarks(http: &ErpClient, l: &Loaded) -> Result, ErpError> { let marker = remarks_marker(&l.invoice.number); let rows = http .list_resource( DOCTYPE_INVOICE, &["name", "docstatus", "grand_total", "remarks"], json!([["remarks", "like", format!("{marker}%")], ["docstatus", "!=", 2]]), "creation asc", ) .await?; // "like" is a prefix match: "INV/2026-001" must not pick up "INV/2026-0010". let hit = rows.iter().find(|r| { let remarks = r.get("remarks").and_then(Value::as_str).unwrap_or(""); remarks == marker || remarks.starts_with(&format!("{marker}\n")) }); match hit { Some(doc) => { let name = doc.get("name").and_then(Value::as_str).unwrap_or_default(); accept_existing(name, doc, &l.invoice).map(Some) } None => Ok(None), } } async fn create_or_find(http: &ErpClient, l: &Loaded, body: &mapping::BuiltRequest) -> Result { let inv = &l.invoice; if l.cfg.naming_mode == NamingMode::Series { if let Some(found) = find_by_remarks(http, l).await? { return Ok(found); } } match http.post(body.path, &body.body, body.idempotent).await { Ok(resp) => { let doc = resp.get("data").cloned().unwrap_or(Value::Null); let name = match (doc_name(&resp), l.cfg.naming_mode) { (Some(n), _) => n, (None, NamingMode::Mirror) => inv.number.clone(), (None, NamingMode::Series) => { return Err(ErpError::protocol("ERPNext did not return the new Sales Invoice's name.")) } }; Ok(RemoteDoc { name, docstatus: doc_docstatus(&doc), created: true }) } Err(e) if e.kind == ErrorKind::Duplicate && l.cfg.naming_mode == NamingMode::Mirror => { // The mirrored name is taken: either a repeat of an earlier push or someone else's document. let mut path: Vec<&str> = SALES_INVOICE_V2.to_vec(); path.push(&inv.number); match http.get(&path, &[]).await { Ok(resp) => { let doc = resp.get("data").cloned().unwrap_or(Value::Null); accept_existing(&inv.number, &doc, inv) } Err(_) => Err(e), } } Err(e) => Err(e), } } /// `POST .../method/submit` on API v2 (mirror mode needs v2 anyway); series mode uses the v1 `run_method` form. /// Never `frappe.client.submit`, which overwrites the whole document. async fn submit_remote(http: &ErpClient, cfg: &ErpnextConfig, name: &str) -> Result { let resp = match cfg.naming_mode { NamingMode::Mirror => { let mut path: Vec<&str> = SALES_INVOICE_V2.to_vec(); path.extend([name, "method", "submit"]); http.post(&path, &json!({}), false).await? } NamingMode::Series => { http.post(&["api", "resource", DOCTYPE_INVOICE, name], &json!({ "run_method": "submit" }), false).await? } }; // A 2xx without a docstatus is taken as submitted; one that says otherwise is not. match resp.get("data").and_then(|d| d.get("docstatus")).and_then(Value::as_i64) { Some(1) | None => Ok(1), Some(other) => Err(ErpError::protocol(format!( "ERPNext accepted the submit request but the document is still at docstatus {other}." ))), } } fn attachment_file_name(number: &str) -> String { let cleaned: String = number .chars() .map(|c| if c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.') { c } else { '-' }) .collect(); let cleaned = cleaned.trim_matches('-'); format!("{}.pdf", if cleaned.is_empty() { "invoice" } else { cleaned }) } async fn attach_pdf(http: &ErpClient, l: &Loaded, remote_name: &str, pdf: &Pdf) -> Result<(), ErpError> { let file_name = attachment_file_name(&l.invoice.number); let fields = [ ("doctype", DOCTYPE_INVOICE.to_string()), ("docname", remote_name.to_string()), ("is_private", "1".to_string()), ]; let resp = http .post_file( &["api", "method", "upload_file"], &Upload { file_name: &file_name, mime: "application/pdf", bytes: &pdf.bytes, fields: &fields }, ) .await?; let ok = resp.get("message").map(|m| m.get("name").is_some() || m.get("file_url").is_some()).unwrap_or(false); if ok { Ok(()) } else { Err(ErpError::protocol("ERPNext did not confirm the PDF upload.")) } } // ---- the invoice push ---- #[derive(Default)] struct Progress { step: &'static str, remote_name: String, remote_docstatus: i64, payload_hash: String, attachment_sha256: String, created: bool, no_op: bool, warnings: Vec, } fn persist(db: &Db, invoice_id: i64, prev: Option<&SyncRow>, st: &Progress, status: &str, error: &str) -> Result<(), ErpError> { let synced_at = if status == "synced" { Some(chrono::Utc::now().to_rfc3339()) } else { prev.and_then(|p| p.synced_at.clone()) }; let row = SyncRow { remote_name: st.remote_name.clone(), remote_docstatus: st.remote_docstatus, status: status.to_string(), last_error: error.to_string(), payload_hash: st.payload_hash.clone(), synced_at, attachment_sha256: st.attachment_sha256.clone(), }; with_db(db, |c| write_sync(c, invoice_id, &row)) } /// The message shown to the user: which step failed, then the error's own readable text. Local refusals and /// conflicts already say what is wrong and get no prefix. fn failure_text(step: &str, e: &ErpError) -> String { if step.is_empty() || matches!(e.kind, ErrorKind::Config | ErrorKind::Precondition | ErrorKind::Conflict) { e.to_string() } else { format!("Could not {step}: {e}") } } async fn run_push(db: &Db, http: &ErpClient, l: &Loaded, want_submit: bool, st: &mut Progress) -> Result<(), ErpError> { let inv = &l.invoice; let cfg = &l.cfg; let prev = l.sync.as_ref(); // A row with a remote name (synced, or failed after the document was created) is the same remote document; // a conflict row is re-checked from scratch. let existing = prev.filter(|s| !s.remote_name.is_empty() && s.status != "conflict"); if let Some(p) = prev { st.remote_name = p.remote_name.clone(); st.remote_docstatus = p.remote_docstatus; st.payload_hash = p.payload_hash.clone(); st.attachment_sha256 = p.attachment_sha256.clone(); } st.step = "look up or create the Customer"; let customer = ensure_customer(db, http, l).await?; st.step = "create the Address"; let (address, address_warning) = ensure_address(db, http, l, &customer).await?; st.warnings.extend(address_warning); st.step = ""; let ctx = InvoiceContext { invoice: inv, config: cfg, vendor: &l.vendor, customer: &customer, customer_address: address.as_deref(), item_codes: &l.item_codes, india_compliance: l.india_compliance, submit: false, }; let built = build_sales_invoice(&ctx).map_err(pre)?; // serde_json keeps object keys sorted, so the serialisation (and the hash) is stable. let hash = sha256_hex(built.body.to_string().as_bytes()); match existing { Some(p) => { if p.payload_hash.is_empty() { st.payload_hash = hash; } else if p.payload_hash != hash { st.warnings.push( "The settings or client details changed since this invoice was sent. The ERPNext document was left as it is." .into(), ); } } None => { st.payload_hash = hash; st.step = "create the Sales Invoice"; let doc = create_or_find(http, l, &built).await?; st.remote_name = doc.name; st.remote_docstatus = doc.docstatus; st.created = doc.created; // Keep the remote name even if the next steps fail. persist(db, inv.id, prev, st, "synced", "")?; } } let need_submit = want_submit && st.remote_docstatus == 0; let pdf = l.pdf.as_ref().filter(|p| p.sha256 != st.attachment_sha256); if let (Some(w), None) = (&l.pdf_warning, &l.pdf) { st.warnings.push(w.clone()); } if existing.is_some() && prev.is_some_and(|p| p.status == "synced") && !need_submit && pdf.is_none() { st.no_op = true; return Ok(()); } if need_submit { st.step = "submit the Sales Invoice"; st.remote_docstatus = submit_remote(http, cfg, &st.remote_name.clone()).await?; persist(db, inv.id, prev, st, "synced", "")?; } if let Some(pdf) = pdf { st.step = "attach the PDF"; match attach_pdf(http, l, &st.remote_name.clone(), pdf).await { Ok(()) => st.attachment_sha256 = pdf.sha256.clone(), // The invoice itself is in ERPNext; a failed upload is a warning, retried by the next push. Err(e) => st.warnings.push(format!("The PDF was not attached: {e}")), } } Ok(()) } pub async fn push_invoice(db: &Db, local_dir: &Path, http: &ErpClient, invoice_id: i64, submit: Option) -> PushResult { let loaded = match load_for_push(db, local_dir, invoice_id) { Ok(l) => l, Err((number, e)) => return PushResult::refused(invoice_id, &number, e), }; let want_submit = submit.unwrap_or(loaded.cfg.submit_on_push); let mut st = Progress::default(); let outcome = run_push(db, http, &loaded, want_submit, &mut st).await; let number = loaded.invoice.number.clone(); let prev = loaded.sync.as_ref(); let (error, status): (Option<(ErpError, String)>, &str) = match outcome { Ok(()) => { let write = if st.no_op { Ok(()) } else { persist(db, invoice_id, prev, &st, "synced", "") }; match write { Ok(()) => (None, "synced"), Err(e) => { let text = e.to_string(); (Some((e, text)), "error") } } } Err(e) => { let text = failure_text(st.step, &e); let status = if e.kind == ErrorKind::Conflict { "conflict" } else { "error" }; // Best effort: the original error is what the caller needs to see. let _ = persist(db, invoice_id, prev, &st, status, &text); (Some((e, text)), status) } }; PushResult { invoice_id, number, ok: error.is_none(), status: status.to_string(), remote_name: st.remote_name, remote_docstatus: st.remote_docstatus, created: st.created, no_op: st.no_op, attached: !st.attachment_sha256.is_empty(), error: error.as_ref().map(|(_, t)| t.clone()), error_kind: error.as_ref().map(|(e, _)| e.kind), warnings: st.warnings, } } /// Pushes one invoice after another; a failing row never stops the rest. pub async fn push_invoices( db: &Db, local_dir: &Path, http: &ErpClient, ids: &[i64], submit: Option, ) -> Vec { let mut seen = std::collections::HashSet::new(); let mut out = Vec::with_capacity(ids.len()); for &id in ids { if seen.insert(id) { out.push(push_invoice(db, local_dir, http, id, submit).await); } } out } // ---- payments ---- pub struct PaymentEntryInput<'a> { pub payment_id: i64, pub invoice_number: &'a str, pub remote_invoice: &'a str, pub paid_on: &'a str, pub reference: &'a str, /// Cash received. pub cash_paise: i64, pub tds_paise: i64, pub tds_account: &'a str, pub cost_center: &'a str, } /// Turns the unsaved dict from `get_payment_entry` into the Payment Entry to insert and submit. /// /// UNVERIFIED against a live ERPNext (check in F4): the deduction row fields (`account`, `cost_center`, /// `amount`), the sign ERPNext expects for a TDS deduction on a receipt, and whether `allocated_amount` must be /// cash plus TDS for the difference amount to come out zero. Everything that depends on those guesses is here. pub fn build_payment_entry(draft: &Value, p: &PaymentEntryInput) -> Result { let mut doc: Map = draft.as_object().cloned().ok_or("ERPNext returned an unexpected payment draft.")?; doc.retain(|k, _| !k.starts_with("__")); doc.insert("doctype".into(), json!("Payment Entry")); doc.insert("posting_date".into(), json!(p.paid_on)); let reference = if p.reference.trim().is_empty() { format!("Voiced payment {}", p.payment_id) } else { p.reference.trim().to_string() }; doc.insert("reference_no".into(), json!(reference)); doc.insert("reference_date".into(), json!(p.paid_on)); doc.insert("paid_amount".into(), mapping::money(p.cash_paise)); doc.insert("received_amount".into(), mapping::money(p.cash_paise)); doc.insert("remarks".into(), json!(format!("Voiced payment {} for invoice {}", p.payment_id, p.invoice_number))); let allocated = mapping::money(p.cash_paise + p.tds_paise); let refs = doc.get_mut("references").and_then(Value::as_array_mut).ok_or("ERPNext returned no invoice reference for this payment.")?; let target = refs .iter_mut() .find(|r| r.get("reference_name").and_then(Value::as_str) == Some(p.remote_invoice)) .ok_or_else(|| format!("ERPNext's payment draft does not reference {}.", p.remote_invoice))?; target["allocated_amount"] = allocated; let deductions = if p.tds_paise > 0 { let mut row = Map::new(); row.insert("account".into(), json!(p.tds_account)); if !p.cost_center.trim().is_empty() { row.insert("cost_center".into(), json!(p.cost_center.trim())); } row.insert("amount".into(), mapping::money(p.tds_paise)); vec![Value::Object(row)] } else { Vec::new() }; doc.insert("deductions".into(), Value::Array(deductions)); doc.insert("docstatus".into(), json!(1)); Ok(Value::Object(doc)) } struct PaymentRow { invoice_id: i64, paid_on: String, amount_paise: i64, tds_paise: i64, reference: String, entry: Option, } fn payment_failure(payment_id: i64, invoice_id: i64, step: &str, e: ErpError) -> PaymentPushResult { PaymentPushResult { payment_id, invoice_id, ok: false, entry_name: None, already_synced: false, error: Some(failure_text(step, &e)), error_kind: Some(e.kind), } } pub async fn push_payment(db: &Db, http: &ErpClient, payment_id: i64) -> PaymentPushResult { let loaded = with_db(db, |conn| { let row = conn .query_row( "SELECT invoice_id, paid_on, amount_paise, tds_paise, reference, erpnext_payment_entry FROM payments WHERE id = ?1", params![payment_id], |r| { Ok(PaymentRow { invoice_id: r.get(0)?, paid_on: r.get(1)?, amount_paise: r.get(2)?, tds_paise: r.get(3)?, reference: r.get(4)?, entry: r.get::<_, Option>(5)?.filter(|e| !e.trim().is_empty()), }) }, ) .optional() .map_err(|e| e.to_string())? .ok_or_else(|| "Payment not found".to_string())?; let number: String = conn .query_row("SELECT number FROM invoices WHERE id = ?1", params![row.invoice_id], |r| r.get(0)) .map_err(|e| e.to_string())?; let sync = load_sync(conn, row.invoice_id)?; let cfg = config::load(conn)?; Ok((row, number, sync, cfg)) }); let (row, number, sync, cfg) = match loaded { Ok(v) => v, Err(e) => return payment_failure(payment_id, 0, "", e), }; let fail = |e: ErpError| payment_failure(payment_id, row.invoice_id, "", e); if let Some(entry) = &row.entry { return PaymentPushResult { payment_id, invoice_id: row.invoice_id, ok: true, entry_name: Some(entry.clone()), already_synced: true, error: None, error_kind: None, }; } let remote_invoice = match sync { // A later unrelated failure may have flipped the row to `error`; the remote document is still submitted. Some(s) if !s.remote_name.is_empty() && s.status != "conflict" && s.remote_docstatus == 1 => s.remote_name, Some(s) if !s.remote_name.is_empty() => { return fail(pre(format!( "Invoice {number} is not submitted in ERPNext yet. Submit the invoice in ERPNext first (or push it again with \"submit\" on), then send the payment." ))) } _ => { return fail(pre(format!( "Invoice {number} has not been sent to ERPNext. Send it first, and submit it, then send the payment." ))) } }; if cfg.payment_bank_account.trim().is_empty() { return fail(pre("Set the payment bank account in the ERPNext settings first.")); } if row.tds_paise > 0 && cfg.tds_account.trim().is_empty() { return fail(pre("This payment has TDS: set the TDS account in the ERPNext settings first.")); } let query = [ ("dt", DOCTYPE_INVOICE.to_string()), ("dn", remote_invoice.clone()), ("bank_account", cfg.payment_bank_account.trim().to_string()), ("party_amount", paise_to_decimal(row.amount_paise + row.tds_paise)), ]; let draft = match http.get(&["api", "method", GET_PAYMENT_ENTRY], &query).await { Ok(v) => v.get("message").cloned().unwrap_or(Value::Null), Err(e) => return payment_failure(payment_id, row.invoice_id, "prepare the Payment Entry", e), }; let body = match build_payment_entry( &draft, &PaymentEntryInput { payment_id, invoice_number: &number, remote_invoice: &remote_invoice, paid_on: &row.paid_on, reference: &row.reference, cash_paise: row.amount_paise, tds_paise: row.tds_paise, tds_account: cfg.tds_account.trim(), cost_center: &cfg.cost_center, }, ) { Ok(b) => b, Err(e) => return fail(ErpError::protocol(e)), }; // Not retried after a 5xx or timeout: it may have been created, and a second entry would double-count. let resp = match http.post(&["api", "resource", "Payment Entry"], &body, false).await { Ok(v) => v, Err(e) => return payment_failure(payment_id, row.invoice_id, "create the Payment Entry", e), }; let Some(name) = doc_name(&resp) else { return fail(ErpError::protocol("ERPNext did not return the new Payment Entry's name.")); }; if let Err(e) = with_db(db, |c| { c.execute("UPDATE payments SET erpnext_payment_entry = ?1 WHERE id = ?2", params![name, payment_id]) .map(|_| ()) .map_err(|e| e.to_string()) }) { return fail(e); } PaymentPushResult { payment_id, invoice_id: row.invoice_id, ok: true, entry_name: Some(name), already_synced: false, error: None, error_kind: None, } } // ---- the sink ---- /// `InvoiceSink` for ERPNext. The commands call `push_invoice`/`push_payment` above for the richer results; /// the trait is the seam other targets implement. pub struct ErpnextSink<'a> { pub db: &'a Db, pub local_dir: &'a Path, pub http: &'a ErpClient, } impl InvoiceSink for ErpnextSink<'_> { fn id(&self) -> &'static str { "erpnext" } fn push_invoice<'a>( &'a self, request: PushRequest, ) -> impl std::future::Future> + Send + 'a { async move { let r = push_invoice(self.db, self.local_dir, self.http, request.invoice_id, request.submit).await; if r.ok { Ok(PushedInvoice { remote_name: r.remote_name, remote_docstatus: r.remote_docstatus, created: r.created, warnings: r.warnings, }) } else { Err(r.error.unwrap_or_else(|| "The push failed.".into())) } } } fn push_payment<'a>( &'a self, request: PaymentRequest, ) -> impl std::future::Future> + Send + 'a { async move { let r = push_payment(self.db, self.http, request.payment_id).await; match (r.ok, r.entry_name) { (true, Some(name)) => Ok(PushedPayment { remote_name: name, created: !r.already_synced }), _ => Err(r.error.unwrap_or_else(|| "The payment push failed.".into())), } } } fn status(&self, invoice_id: i64) -> Result { let conn = self.db.lock().map_err(|e| e.to_string())?; sync_status(&conn, invoice_id) } } #[cfg(test)] mod tests;