diff --git a/erpnext/stock/doctype/repost_item_valuation/test_repost_item_valuation.py b/erpnext/stock/doctype/repost_item_valuation/test_repost_item_valuation.py index a2847a120b5..0ed02dde407 100644 --- a/erpnext/stock/doctype/repost_item_valuation/test_repost_item_valuation.py +++ b/erpnext/stock/doctype/repost_item_valuation/test_repost_item_valuation.py @@ -2,21 +2,23 @@ # See license.txt -from unittest.mock import MagicMock, call +from unittest.mock import MagicMock, call, patch import frappe from frappe.tests.utils import FrappeTestCase, change_settings -from frappe.utils import add_days, add_to_date, flt, now, nowdate, today +from frappe.utils import add_days, add_to_date, flt, get_datetime, now, nowdate, today from erpnext.accounts.doctype.sales_invoice.test_sales_invoice import create_sales_invoice from erpnext.accounts.utils import repost_gle_for_stock_vouchers from erpnext.controllers.stock_controller import create_item_wise_repost_entries +from erpnext.stock import stock_ledger from erpnext.stock.doctype.item.test_item import make_item from erpnext.stock.doctype.purchase_receipt.test_purchase_receipt import make_purchase_receipt from erpnext.stock.doctype.repost_item_valuation.repost_item_valuation import ( in_configured_timeslot, ) from erpnext.stock.doctype.stock_entry.stock_entry_utils import make_stock_entry +from erpnext.stock.stock_ledger import update_entries_after from erpnext.stock.tests.test_utils import StockTestMixin from erpnext.stock.utils import PendingRepostingError @@ -574,6 +576,228 @@ class TestRepostItemValuation(FrappeTestCase, StockTestMixin): self.assertSLEs(return_pr, expected_sles) + def test_skip_later_repost_covered_by_manufacture_dependant(self): + """A finished good reposted as a dependant of its raw material makes a later + repost queued for the same finished good and warehouse redundant.""" + rm = self.make_item(properties={"valuation_method": "FIFO"}).name + fg = self.make_item(properties={"valuation_method": "FIFO"}).name + warehouse = "_Test Warehouse - _TC" + + make_stock_entry( + item_code=rm, target=warehouse, qty=100, rate=100, posting_date=add_days(today(), -10) + ) + + manufacture = make_stock_entry( + item_code=rm, + source=warehouse, + qty=10, + purpose="Manufacture", + posting_date=add_days(today(), -5), + do_not_save=True, + ) + manufacture.append( + "items", + { + "item_code": fg, + "t_warehouse": warehouse, + "qty": 1, + "transfer_qty": 1, + "uom": "Nos", + "stock_uom": "Nos", + "conversion_factor": 1.0, + "is_finished_item": 1, + }, + ) + manufacture.save() + manufacture.submit() + + # a repost queued for the finished good, dated after the manufacture entry + later_riv = frappe.get_doc( + doctype="Repost Item Valuation", + based_on="Item and Warehouse", + item_code=fg, + warehouse=warehouse, + posting_date=today(), + posting_time="00:00:01", + ) + later_riv.flags.dont_run_in_test = True + later_riv.submit() + self.assertEqual(later_riv.status, "Queued") + + # reposting the raw material walks the finished good forward as a dependant + rm_riv = frappe.get_doc( + doctype="Repost Item Valuation", + based_on="Item and Warehouse", + item_code=rm, + warehouse=warehouse, + posting_date=add_days(today(), -10), + posting_time="00:00:01", + ) + rm_riv.submit() + + later_riv.load_from_db() + self.assertEqual(later_riv.status, "Skipped") + + def test_repost_covering_earlier_date_is_not_skipped(self): + """A repost for the finished good that starts before the manufacture entry still + has work to do, so it must survive.""" + rm = self.make_item(properties={"valuation_method": "FIFO"}).name + fg = self.make_item(properties={"valuation_method": "FIFO"}).name + warehouse = "_Test Warehouse - _TC" + + make_stock_entry( + item_code=rm, target=warehouse, qty=100, rate=100, posting_date=add_days(today(), -10) + ) + make_stock_entry(item_code=fg, target=warehouse, qty=5, rate=50, posting_date=add_days(today(), -9)) + + manufacture = make_stock_entry( + item_code=rm, + source=warehouse, + qty=10, + purpose="Manufacture", + posting_date=add_days(today(), -5), + do_not_save=True, + ) + manufacture.append( + "items", + { + "item_code": fg, + "t_warehouse": warehouse, + "qty": 1, + "transfer_qty": 1, + "uom": "Nos", + "stock_uom": "Nos", + "conversion_factor": 1.0, + "is_finished_item": 1, + }, + ) + manufacture.save() + manufacture.submit() + + earlier_riv = frappe.get_doc( + doctype="Repost Item Valuation", + based_on="Item and Warehouse", + item_code=fg, + warehouse=warehouse, + posting_date=add_days(today(), -9), + posting_time="00:00:01", + ) + earlier_riv.flags.dont_run_in_test = True + earlier_riv.submit() + + rm_riv = frappe.get_doc( + doctype="Repost Item Valuation", + based_on="Item and Warehouse", + item_code=rm, + warehouse=warehouse, + posting_date=add_days(today(), -10), + posting_time="00:00:01", + ) + rm_riv.submit() + + earlier_riv.load_from_db() + self.assertEqual(earlier_riv.status, "Queued") + earlier_riv.set_status("Skipped") + + def test_repost_covers_every_entry_once_across_batches(self): + """Full rows are fetched REPOST_SLE_BATCH_SIZE at a time, and the prefetched + window is dropped whenever a dependant repost re-sorts the queue. Every active + entry must still be reposted exactly once, in posting order.""" + rm = self.make_item(properties={"valuation_method": "FIFO"}).name + fg = self.make_item(properties={"valuation_method": "FIFO"}).name + warehouse = "_Test Warehouse - _TC" + + make_stock_entry( + item_code=rm, target=warehouse, qty=100, rate=100, posting_date=add_days(today(), -10) + ) + for day, rate in ((-9, 110), (-8, 120), (-7, 130)): + make_stock_entry( + item_code=rm, target=warehouse, qty=10, rate=rate, posting_date=add_days(today(), day) + ) + + manufacture = make_stock_entry( + item_code=rm, + source=warehouse, + qty=10, + purpose="Manufacture", + posting_date=add_days(today(), -6), + do_not_save=True, + ) + manufacture.append( + "items", + { + "item_code": fg, + "t_warehouse": warehouse, + "qty": 1, + "transfer_qty": 1, + "uom": "Nos", + "stock_uom": "Nos", + "conversion_factor": 1.0, + "is_finished_item": 1, + }, + ) + manufacture.save() + manufacture.submit() + + # the finished good is pulled in as a dependant while the raw material is being + # reposted, so the queue grows and is re-sorted part way through + for day in (-5, -4, -3): + make_stock_entry( + item_code=fg, target=warehouse, qty=2, rate=200, posting_date=add_days(today(), day) + ) + + reposted = [] + fetched_batches = [] + original_repost = update_entries_after.repost_stock_ledger_entry + original_fetch = stock_ledger.get_sle_entries_by_names + + def record_repost(self, sle): + reposted.append(sle.name) + return original_repost(self, sle) + + def record_fetch(names): + fetched_batches.append(len(names)) + return original_fetch(names) + + batch_size = 2 + with ( + patch.object(stock_ledger, "REPOST_SLE_BATCH_SIZE", batch_size), + patch.object(update_entries_after, "repost_stock_ledger_entry", record_repost), + patch.object(stock_ledger, "get_sle_entries_by_names", record_fetch), + ): + riv = frappe.get_doc( + doctype="Repost Item Valuation", + based_on="Item and Warehouse", + item_code=rm, + warehouse=warehouse, + posting_date=add_days(today(), -10), + posting_time="00:00:00", + ) + riv.submit() + + active_sles = frappe.get_all( + "Stock Ledger Entry", + filters={"item_code": ("in", [rm, fg]), "warehouse": warehouse, "is_cancelled": 0}, + fields=["name", "posting_datetime", "creation"], + ) + self.assertGreater(len(active_sles), batch_size) + + # every active entry was reposted, and none of them twice + self.assertEqual(sorted(reposted), sorted(row.name for row in active_sles)) + self.assertEqual(len(reposted), len(set(reposted))) + + # and they were reposted in posting order, across the batch boundaries and the + # flush that the dependant discovery triggers + posting_order = { + row.name: (get_datetime(row.posting_datetime), get_datetime(row.creation)) for row in active_sles + } + reposted_order = [posting_order[name] for name in reposted] + self.assertEqual(reposted_order, sorted(reposted_order)) + + # the rows really were fetched a batch at a time, never the whole queue at once + self.assertGreater(len(fetched_batches), 1) + self.assertLessEqual(max(fetched_batches), batch_size) + def test_remove_attached_file(self): item_code = make_item("_Test Remove Attached File Item", properties={"is_stock_item": 1}) diff --git a/erpnext/stock/stock_ledger.py b/erpnext/stock/stock_ledger.py index ebe3b20931d..437d96da2da 100644 --- a/erpnext/stock/stock_ledger.py +++ b/erpnext/stock/stock_ledger.py @@ -2,15 +2,17 @@ # License: GNU General Public License v3. See license.txt import copy +import gc import gzip import json from collections import deque +from itertools import islice import frappe from frappe import _, bold, scrub from frappe.model.meta import get_field_precision from frappe.query_builder import Order -from frappe.query_builder.functions import Sum +from frappe.query_builder.functions import CombineDatetime, Sum from frappe.utils import ( cint, cstr, @@ -45,6 +47,28 @@ from erpnext.stock.utils import ( ) from erpnext.stock.valuation import FIFOValuation, LIFOValuation, round_off_if_near_zero +# Number of stock ledger entries whose full row is loaded in memory at a time while +# reposting. The reposting queue itself only holds the identity/sort keys of the +# entries so that a repost spanning millions of entries does not blow up the worker. +REPOST_SLE_BATCH_SIZE = 500 + +# How many of the most recent messages to keep when trimming `frappe.local.message_log` +# during a repost. The failure handler in Repost Item Valuation reads the tail of this +# log to build the error log, so the recent entries have to survive. +REPOST_MESSAGE_LOG_LIMIT = 50 + +# Columns needed to queue and sort an entry for reposting. The remaining columns are +# fetched in batches of REPOST_SLE_BATCH_SIZE just before the entry is processed. +REPOST_SLE_QUEUE_FIELDS = ( + "name", + "item_code", + "warehouse", + "posting_date", + "posting_time", + "posting_datetime", + "creation", +) + class NegativeStockError(frappe.ValidationError): pass @@ -258,9 +282,70 @@ def repost_future_sle( resume_item_wh_wise_last_posted_sle = {} repost_affected_transaction.update(obj.repost_affected_transaction) + skip_reposts_covered_by_dependant_repost(doc, obj.reposted_dependant_item_wh) update_args_in_repost_item_valuation(doc, index, items_to_be_repost, repost_affected_transaction) +def skip_reposts_covered_by_dependant_repost(doc, reposted_dependant_item_wh): + """Skip queued reposts that a Manufacture/Repack dependant repost has already covered. + + While reposting a raw material, the finished goods produced from it are reposted as + dependants, from the posting datetime of the manufacture entry right through to the + end of their ledger. A separate repost queued for the same finished good and + warehouse at a later datetime therefore has nothing left to do, so it is marked as + Skipped instead of walking the same entries again. + + Only `Item and Warehouse` reposts are skipped. A `Transaction` repost covers several + item-warehouse combinations, so covering one of them says nothing about the rest. + """ + if not doc or not reposted_dependant_item_wh: + return + + riv = frappe.qb.DocType("Repost Item Valuation") + + for (item_code, warehouse), posting_datetime in reposted_dependant_item_wh.items(): + if not posting_datetime: + continue + + ( + frappe.qb.update(riv) + .set(riv.status, "Skipped") + .where( + (riv.item_code == item_code) + & (riv.warehouse == warehouse) + & (riv.name != doc.name) + & (riv.docstatus == 1) + & (riv.status == "Queued") + & (riv.based_on == "Item and Warehouse") + & (CombineDatetime(riv.posting_date, riv.posting_time) >= posting_datetime) + ) + ).run() + + +def release_reposting_memory(): + """Drop process local caches that keep growing over a long running repost. + + `frappe.get_cached_doc`/`get_cached_value` mirror every fetched document in + `frappe.local.cache`, which is never evicted within a job. A repost touching + thousands of distinct Stock Entries, Purchase Receipts or Serial and Batch Bundles + therefore retains all of those documents until the worker exits. Everything dropped + here is still in redis, so it is only re-fetched on demand. + """ + local_cache = getattr(frappe.local, "cache", None) + if isinstance(local_cache, dict): + for key in [key for key in local_cache if b"|document_cache::" in frappe.safe_encode(key)]: + local_cache.pop(key, None) + + # msgprint during reposting (eg. negative stock warnings) accumulates here and is + # never trimmed. Keep the most recent messages so that a failure later in the repost + # can still report them, and drop only the older ones. + message_log = getattr(frappe.local, "message_log", None) + if message_log and len(message_log) > REPOST_MESSAGE_LOG_LIMIT: + frappe.local.message_log = message_log[-REPOST_MESSAGE_LOG_LIMIT:] + + gc.collect() + + def update_args_in_repost_item_valuation( doc, index, @@ -491,6 +576,7 @@ class update_entries_after: self.repost_affected_transaction = args.get("repost_affected_transaction") or set() self.new_items_found = False + self.reposted_dependant_item_wh = {} self.reserved_stock = self.get_reserved_stock() self.data = frappe._dict() @@ -596,8 +682,10 @@ class update_entries_after: def initialize_reposting(self): self._sles = [] + self._sle_batch = {} self.distinct_sles = set() self.distinct_dependant_item_wh = set() + self.reposted_dependant_item_wh = {} self.prev_sle_dict = frappe._dict({}) def get_item_wh_wise_last_posted_sle(self): @@ -640,17 +728,27 @@ class update_entries_after: i = 0 while self._sles: - sle = self._sles.popleft() - if (sle.item_code, sle.warehouse) not in self.distinct_dependant_item_wh: - self.distinct_dependant_item_wh.add((sle.item_code, sle.warehouse)) + queued_sle = self._sles.popleft() + if (queued_sle.item_code, queued_sle.warehouse) not in self.distinct_dependant_item_wh: + self.distinct_dependant_item_wh.add((queued_sle.item_code, queued_sle.warehouse)) - if sle.name in self.distinct_sles: + if queued_sle.name in self.distinct_sles: continue i += 1 - item_wh_key = (sle.item_code, sle.warehouse) + item_wh_key = (queued_sle.item_code, queued_sle.warehouse) if item_wh_key not in self.prev_sle_dict: - self.prev_sle_dict[item_wh_key] = get_previous_sle_of_current_voucher(sle) + self.prev_sle_dict[item_wh_key] = get_previous_sle_of_current_voucher(queued_sle) + + sle = self.get_sle_to_repost(queued_sle) + if not sle: + # the entry was cancelled or deleted after it was queued, so it must not be + # reposted. Cancellation queues its own repost, which picks up from there. + frappe.logger("stock_ledger").info( + f"Skipped {queued_sle.name} while reposting {self.item_code}, " + "entry is no longer active" + ) + continue self.repost_stock_ledger_entry(sle) @@ -664,6 +762,24 @@ class update_entries_after: if i % 2000 == 0: self.update_data_in_repost(len(self._sles), i) + def get_sle_to_repost(self, queued_sle): + """Return the full stock ledger entry row for a queued entry. + + Rows are fetched (and locked) REPOST_SLE_BATCH_SIZE at a time so that only a + small window of complete entries is ever held in memory. + """ + if sle := self._sle_batch.pop(queued_sle.name, None): + return sle + + names = [queued_sle.name] + for row in islice(self._sles, 0, REPOST_SLE_BATCH_SIZE - 1): + if row.name not in self.distinct_sles: + names.append(row.name) + + self._sle_batch = {row.name: row for row in get_sle_entries_by_names(names)} + + return self._sle_batch.pop(queued_sle.name, None) + def sort_sles(self, sles): return sorted( sles, @@ -675,27 +791,40 @@ class update_entries_after: def include_dependant_sle_in_reposting(self, sle): repost_dependant_sle = False - if sle.voucher_type == "Stock Entry" and is_repack_entry(sle.voucher_no): - repack_sles = self.get_sles_for_repack(sle) - for repack_sle in repack_sles: - if (repack_sle.item_code, repack_sle.warehouse) in self.distinct_dependant_item_wh: - continue - repost_dependant_sle = True - self.distinct_dependant_item_wh.add((repack_sle.item_code, repack_sle.warehouse)) - self._sles.extend(self.get_future_entries_to_repost(repack_sle)) + # For a Manufacture/Repack entry the consumed row points at the finished good row, + # so the dependants picked up here are the finished goods produced by this entry. + # Reposting them here makes any queued repost for the same item-warehouse at a + # later date redundant. + produced_by_manufacture = sle.voucher_type == "Stock Entry" and is_manufacture_or_repack_entry( + sle.voucher_no + ) + + if sle.voucher_type == "Stock Entry" and is_repack_entry(sle.voucher_no): + dependant_sles = self.get_sles_for_repack(sle) else: dependant_sles = get_sle_by_voucher_detail_no(sle.dependant_sle_voucher_detail_no) - for depend_sle in dependant_sles: - if (depend_sle.item_code, depend_sle.warehouse) in self.distinct_dependant_item_wh: - continue - repost_dependant_sle = True - self.distinct_dependant_item_wh.add((depend_sle.item_code, depend_sle.warehouse)) - self._sles.extend(self.get_future_entries_to_repost(depend_sle)) + for depend_sle in dependant_sles: + item_wh_key = (depend_sle.item_code, depend_sle.warehouse) + if item_wh_key in self.distinct_dependant_item_wh: + continue + + repost_dependant_sle = True + self.distinct_dependant_item_wh.add(item_wh_key) + self._sles.extend(self.get_future_entries_to_repost(depend_sle)) + + if produced_by_manufacture: + self.reposted_dependant_item_wh.setdefault( + item_wh_key, + depend_sle.posting_datetime + or get_combine_datetime(depend_sle.posting_date, depend_sle.posting_time), + ) if repost_dependant_sle: self._sles = deque(self.sort_sles(self._sles)) + # the queue order changed, the prefetched window is no longer the next batch + self._sle_batch = {} def repost_stock_ledger_entry(self, sle): if isinstance(sle, dict): @@ -723,6 +852,7 @@ class update_entries_after: def reset_vouchers_and_idx(self): self.stock_ledgers_to_repost = [] + self._sle_batch = {} self.prev_sle_dict = frappe._dict() self.item_wh_wise_last_posted_sle = frappe._dict() @@ -749,6 +879,8 @@ class update_entries_after: # To maintain the state of the reposting, so if timeout happens, it can be resumed from the last posted voucher frappe.db.commit() # nosemgrep + release_reposting_memory() + self.publish_real_time_progress(total_sles=total_sles, index=index) def publish_real_time_progress(self, total_sles=None, index=None): @@ -764,7 +896,12 @@ class update_entries_after: ) def get_future_entries_to_repost(self, kwargs): - return get_stock_ledger_entries(kwargs, ">=", "asc", for_update=True, check_serial_no=False) + # The queue holds only the identity and sort keys, and is not locked. Rows are + # locked REPOST_SLE_BATCH_SIZE at a time in `get_sle_to_repost`, so a repost + # spanning millions of entries does not hold a lock on all of them. + return get_stock_ledger_entries( + kwargs, ">=", "asc", check_serial_no=False, fields=REPOST_SLE_QUEUE_FIELDS + ) def get_sles_for_repack(self, sle): return ( @@ -1883,6 +2020,7 @@ def get_stock_ledger_entries( check_serial_no=True, extra_cond=None, for_report=False, + fields=None, ): """get stock ledger entries filtered by specific posting datetime conditions""" conditions = f" and posting_datetime {operator} %(posting_datetime)s" @@ -1942,15 +2080,18 @@ def get_stock_ledger_entries( if for_report and previous_sle.get("project"): conditions += " and project = %(project)s" + select_fields = ", ".join(f"`{field}`" for field in fields) if fields else "*" + # nosemgrep return frappe.db.sql( """ - select *, posting_datetime as "timestamp" + select {select_fields}, posting_datetime as "timestamp" from `tabStock Ledger Entry` where is_cancelled = 0 {conditions} order by posting_datetime {order}, creation {order} {limit} {for_update}""".format( + select_fields=select_fields, conditions=conditions, limit=limit or "", for_update=for_update and "for update" or "", @@ -1962,6 +2103,23 @@ def get_stock_ledger_entries( ) +def get_sle_entries_by_names(names): + """Fetch and lock complete stock ledger entry rows for the given names.""" + if not names: + return [] + + # nosemgrep + return frappe.db.sql( + """ + select *, posting_datetime as "timestamp" + from `tabStock Ledger Entry` + where name in %(names)s and is_cancelled = 0 + for update""", + {"names": names}, + as_dict=1, + ) + + def get_sle_by_voucher_detail_no(voucher_detail_no): return frappe.get_all( "Stock Ledger Entry", @@ -2612,6 +2770,10 @@ def is_repack_entry(stock_entry_id): return frappe.get_cached_value("Stock Entry", stock_entry_id, "purpose") == "Repack" +def is_manufacture_or_repack_entry(stock_entry_id): + return frappe.get_cached_value("Stock Entry", stock_entry_id, "purpose") in ("Manufacture", "Repack") + + def has_correct_data(sle): previous_sle = get_previous_sle( {