diff --git a/erpnext/accounts/doctype/process_period_closing_voucher/process_period_closing_voucher.py b/erpnext/accounts/doctype/process_period_closing_voucher/process_period_closing_voucher.py index 3239e6f4a00..be6f8ddfbbb 100644 --- a/erpnext/accounts/doctype/process_period_closing_voucher/process_period_closing_voucher.py +++ b/erpnext/accounts/doctype/process_period_closing_voucher/process_period_closing_voucher.py @@ -89,50 +89,55 @@ class ProcessPeriodClosingVoucher(Document): cancel_pcv_processing(self.name) +def initialize_parallel_threads(docname: str): + threads = 4 + timeout = frappe.db.get_single_value("Accounts Settings", "pcv_job_timeout") or 3600 + ppcvd = qb.DocType("Process Period Closing Voucher Detail") + + frappe.db.set_value("Process Period Closing Voucher", docname, "status", "Running") + + if normal_balances := ( + qb.from_(ppcvd) + .select(ppcvd.name, ppcvd.processing_date, ppcvd.report_type, ppcvd.parentfield) + .where(ppcvd.parent.eq(docname) & ppcvd.status.eq("Queued")) + .orderby(ppcvd.parentfield, ppcvd.idx, ppcvd.processing_date) + .limit(threads) + .for_update(skip_locked=True) + .run(as_dict=True) + ): + if not is_scheduler_inactive(): + for x in normal_balances: + frappe.db.set_value( + "Process Period Closing Voucher Detail", + x.name, + "status", + "Running", + ) + frappe.enqueue( + method="erpnext.accounts.doctype.process_period_closing_voucher.process_period_closing_voucher.process_individual_date", + queue="long", + timeout=timeout, + is_async=True, + enqueue_after_commit=True, + docname=docname, + row_name=x.name, + date=x.processing_date, + report_type=x.report_type, + parentfield=x.parentfield, + ) + # keep transaction on PPCV and PPCVD short + # prevents concurrency errors - REPEATABLE READ + if not frappe.in_test: + frappe.db.commit() # nosemgrep + else: + frappe.db.set_value("Process Period Closing Voucher", docname, "status", "Completed") + + @frappe.whitelist() def start_pcv_processing(docname: str): if frappe.db.get_value("Process Period Closing Voucher", docname, "status") in ["Queued", "Running"]: frappe.has_permission("Process Period Closing Voucher", "write", doc=docname, throw=True) - frappe.db.set_value("Process Period Closing Voucher", docname, "status", "Running") - - timeout = frappe.db.get_single_value("Accounts Settings", "pcv_job_timeout") or 3600 - - ppcvd = qb.DocType("Process Period Closing Voucher Detail") - if normal_balances := ( - qb.from_(ppcvd) - .select(ppcvd.processing_date, ppcvd.report_type, ppcvd.parentfield) - .where(ppcvd.parent.eq(docname) & ppcvd.status.eq("Queued")) - .orderby(ppcvd.parentfield, ppcvd.idx, ppcvd.processing_date) - .limit(4) - .for_update(skip_locked=True) - .run(as_dict=True) - ): - if not is_scheduler_inactive(): - for x in normal_balances: - frappe.db.set_value( - "Process Period Closing Voucher Detail", - { - "processing_date": x.processing_date, - "parent": docname, - "report_type": x.report_type, - "parentfield": x.parentfield, - }, - "status", - "Running", - ) - frappe.enqueue( - method="erpnext.accounts.doctype.process_period_closing_voucher.process_period_closing_voucher.process_individual_date", - queue="long", - timeout=timeout, - is_async=True, - enqueue_after_commit=True, - docname=docname, - date=x.processing_date, - report_type=x.report_type, - parentfield=x.parentfield, - ) - else: - frappe.db.set_value("Process Period Closing Voucher", docname, "status", "Completed") + initialize_parallel_threads(docname) @frappe.whitelist() @@ -250,11 +255,11 @@ def get_gle_for_closing_account(pcv, dimension_balance, dimensions): @frappe.whitelist() def schedule_next_date(docname: str): timeout = frappe.db.get_single_value("Accounts Settings", "pcv_job_timeout") or 3600 - ppcvd = qb.DocType("Process Period Closing Voucher Detail") + if to_process := ( qb.from_(ppcvd) - .select(ppcvd.processing_date, ppcvd.report_type, ppcvd.parentfield) + .select(ppcvd.name, ppcvd.processing_date, ppcvd.report_type, ppcvd.parentfield) .where(ppcvd.parent.eq(docname) & ppcvd.status.eq("Queued")) .orderby(ppcvd.parentfield, ppcvd.idx, ppcvd.processing_date) .limit(1) @@ -264,15 +269,15 @@ def schedule_next_date(docname: str): if not is_scheduler_inactive(): frappe.db.set_value( "Process Period Closing Voucher Detail", - { - "processing_date": to_process[0].processing_date, - "parent": docname, - "report_type": to_process[0].report_type, - "parentfield": to_process[0].parentfield, - }, + to_process[0].name, "status", "Running", ) + # keep transaction on PPCV and PPCVD short + # prevents concurrency errors - REPEATABLE READ + if not frappe.in_test: + frappe.db.commit() # nosemgrep + frappe.enqueue( method="erpnext.accounts.doctype.process_period_closing_voucher.process_period_closing_voucher.process_individual_date", queue="long", @@ -280,6 +285,7 @@ def schedule_next_date(docname: str): is_async=True, enqueue_after_commit=True, docname=docname, + row_name=to_process[0].name, date=to_process[0].processing_date, report_type=to_process[0].report_type, parentfield=to_process[0].parentfield, @@ -444,6 +450,11 @@ def summarize_and_post_ledger_entries(docname): make_closing_entries(closing_entries, pcv.name, pcv.company, pcv.period_end_date) + # keep transaction on PPCV and PPCVD short + # prevents concurrency errors - REPEATABLE READ + if not frappe.in_test: + frappe.db.commit() # nosemgrep + frappe.db.set_value("Period Closing Voucher", pcv.name, "gle_processing_status", "Completed") frappe.db.set_value("Process Period Closing Voucher", docname, "status", "Completed") @@ -529,10 +540,10 @@ def build_dimension_wise_balance_dict(gl_entries): return dimension_balances -def process_individual_date(docname: str, date, report_type, parentfield): +def process_individual_date(docname: str, row_name, date, report_type, parentfield): current_date_status = frappe.db.get_value( "Process Period Closing Voucher Detail", - {"processing_date": date, "report_type": report_type, "parentfield": parentfield}, + row_name, "status", ) if current_date_status != "Running": @@ -579,17 +590,20 @@ def process_individual_date(docname: str, date, report_type, parentfield): # save results frappe.db.set_value( "Process Period Closing Voucher Detail", - {"processing_date": date, "parent": docname, "report_type": report_type, "parentfield": parentfield}, + row_name, "closing_balance", frappe.json.dumps(res), ) frappe.db.set_value( "Process Period Closing Voucher Detail", - {"processing_date": date, "parent": docname, "report_type": report_type, "parentfield": parentfield}, + row_name, "status", "Completed", ) + # commit heavy computation before touching PPCV or PPCVD + if not frappe.in_test: + frappe.db.commit() # nosemgrep # chain call schedule_next_date(docname) diff --git a/erpnext/accounts/doctype/process_period_closing_voucher/test_process_period_closing_voucher.py b/erpnext/accounts/doctype/process_period_closing_voucher/test_process_period_closing_voucher.py index f34c1dbedfe..5de93ef1bdd 100644 --- a/erpnext/accounts/doctype/process_period_closing_voucher/test_process_period_closing_voucher.py +++ b/erpnext/accounts/doctype/process_period_closing_voucher/test_process_period_closing_voucher.py @@ -48,18 +48,27 @@ class TestProcessPeriodClosingVoucher(ERPNextTestSuite): ppcv.save() return ppcv - def set_processing_date_status(self, date, ppcv, rpt_type, parentfield, status): + def set_processing_date_status(self, row_name, status): frappe.db.set_value( "Process Period Closing Voucher Detail", - {"processing_date": date, "parent": ppcv, "report_type": rpt_type, "parentfield": parentfield}, + row_name, "status", status, ) - def get_processing_date_closing_balance(self, date, ppcv, rpt_type, parentfield): + def get_row_name(self, ppcv_name, rpt_type, parentfield): + return frappe.db.get_all( + "Process Period Closing Voucher Detail", + filters={"parent": ppcv_name, "report_type": rpt_type, "parentfield": parentfield}, + order_by="report_type, idx", + pluck="name", + limit=1, + )[0] + + def get_processing_date_closing_balance(self, row_name): return frappe.db.get_value( "Process Period Closing Voucher Detail", - {"processing_date": date, "parent": ppcv, "report_type": rpt_type, "parentfield": parentfield}, + row_name, "closing_balance", ) @@ -97,11 +106,10 @@ class TestProcessPeriodClosingVoucher(ERPNextTestSuite): parentfield = "normal_balances" rpt_type = "Profit and Loss" # status has to be set to 'Running' for logic to run - self.set_processing_date_status(today(), ppcv.name, rpt_type, parentfield, "Running") - process_individual_date(ppcv.name, today(), rpt_type, parentfield) - bal = frappe.parse_json( - self.get_processing_date_closing_balance(today(), ppcv.name, rpt_type, parentfield) - ) + row_name = self.get_row_name(ppcv.name, rpt_type, parentfield) + self.set_processing_date_status(row_name, "Running") + process_individual_date(ppcv.name, row_name, today(), rpt_type, parentfield) + bal = frappe.parse_json(self.get_processing_date_closing_balance(row_name)) self.assertEqual(len(bal), 1) expected_pl = { "account": "Sales - _TC", @@ -117,11 +125,10 @@ class TestProcessPeriodClosingVoucher(ERPNextTestSuite): # Balance sheet balance rpt_type = "Balance Sheet" - self.set_processing_date_status(today(), ppcv.name, rpt_type, parentfield, "Running") - process_individual_date(ppcv.name, today(), rpt_type, parentfield) - bal = frappe.parse_json( - self.get_processing_date_closing_balance(today(), ppcv.name, rpt_type, parentfield) - ) + row_name = self.get_row_name(ppcv.name, rpt_type, parentfield) + self.set_processing_date_status(row_name, "Running") + process_individual_date(ppcv.name, row_name, today(), rpt_type, parentfield) + bal = frappe.parse_json(self.get_processing_date_closing_balance(row_name)) self.assertEqual(len(bal), 1) expected_bs = { "account": "Debtors - _TC", @@ -138,11 +145,10 @@ class TestProcessPeriodClosingVoucher(ERPNextTestSuite): # Opening balance parentfield = "z_opening_balances" rpt_type = "Balance Sheet" - self.set_processing_date_status(today(), ppcv.name, rpt_type, parentfield, "Running") - process_individual_date(ppcv.name, today(), rpt_type, parentfield) - bal = frappe.parse_json( - self.get_processing_date_closing_balance(today(), ppcv.name, rpt_type, parentfield) - ) + row_name = self.get_row_name(ppcv.name, rpt_type, parentfield) + self.set_processing_date_status(row_name, "Running") + process_individual_date(ppcv.name, row_name, today(), rpt_type, parentfield) + bal = frappe.parse_json(self.get_processing_date_closing_balance(row_name)) self.assertEqual(len(bal), 2) opening_cash = next(x for x in bal if x["account"] == "Cash - _TC") expected_opening_cash = { diff --git a/erpnext/accounts/doctype/process_period_closing_voucher_detail/process_period_closing_voucher_detail.py b/erpnext/accounts/doctype/process_period_closing_voucher_detail/process_period_closing_voucher_detail.py index f3a8302ac5b..0e0b905c96a 100644 --- a/erpnext/accounts/doctype/process_period_closing_voucher_detail/process_period_closing_voucher_detail.py +++ b/erpnext/accounts/doctype/process_period_closing_voucher_detail/process_period_closing_voucher_detail.py @@ -1,7 +1,7 @@ # Copyright (c) 2025, Frappe Technologies Pvt. Ltd. and contributors # For license information, please see license.txt -# import frappe +import frappe from frappe.model.document import Document @@ -24,3 +24,10 @@ class ProcessPeriodClosingVoucherDetail(Document): # end: auto-generated types pass + + +def on_doctype_update(): + frappe.db.add_index( + "Process Period Closing Voucher Detail", + ["parent", "status", "parentfield", "idx", "processing_date"], + )