1818import os
1919import traceback
2020import warnings
21+ from contextlib import suppress
2122from numbers import Real
2223from pathlib import Path
23- from time import time
24+ from time import monotonic , time
2425
2526import numpy as np
2627import simplekml
4344# this is the only format it can both resume from and overwrite safely.
4445_SIMULATION_LOG_SUFFIX = ".txt"
4546
47+ # Which simulation a row belongs to. Every check on a finished run reads it.
48+ _SIMULATION_INDEX_KEY = "index"
49+
50+ # How a manager that has gone away answers a proxy call.
51+ _MANAGER_IS_GONE = (OSError , EOFError )
52+
53+ # Bounded, so a lock its dead holder never released cannot pin this worker.
54+ _REPORT_LOCK_SECONDS = 5.0
55+
56+ # Longer than the exit-code path: a worker that only read the event is healthy
57+ # and leaving at the end of the simulation in hand, not blocked on a dead lock.
58+ _REPORTED_FAILURE_GRACE_SECONDS = 60.0
59+
4660
4761def _refuse_logs_this_run_cannot_write (
4862 input_file , output_file , error_file , export_config = None
@@ -300,6 +314,14 @@ def simulate(
300314 -------
301315 None
302316
317+ Raises
318+ ------
319+ RuntimeError
320+ If a parallel run does not finish. A worker that ends badly, one
321+ that reports a failure, and logs that do not hold every simulation
322+ asked for are each refused, since a run that lost work must not be
323+ reported as one that completed.
324+
303325 Notes
304326 -----
305327 If you need to stop the simulations after starting them, you can
@@ -471,22 +493,26 @@ def __run_in_parallel(self, n_workers=None):
471493 processes = []
472494 seeds = np .random .SeedSequence ().spawn (n_workers )
473495
474- for seed in seeds :
475- sim_producer = multiprocess .Process (
476- target = self .__sim_producer ,
477- args = (
478- seed ,
479- sim_monitor ,
480- mutex ,
481- simulation_error_event ,
482- ),
483- )
484- processes .append (sim_producer )
485- sim_producer .start ()
486-
487496 try :
488- for sim_producer in processes :
489- sim_producer .join ()
497+ for seed in seeds :
498+ sim_producer = multiprocess .Process (
499+ target = self .__sim_producer ,
500+ args = (
501+ seed ,
502+ sim_monitor ,
503+ mutex ,
504+ simulation_error_event ,
505+ ),
506+ )
507+ sim_producer .start ()
508+ # Started first: one that never did cannot be joined, and
509+ # a later start failing still has to bring these down.
510+ processes .append (sim_producer )
511+
512+ _join_the_workers (processes , simulation_error_event )
513+
514+ # Before the event: a killed worker never sets it.
515+ _refuse_a_worker_that_did_not_finish (processes )
490516
491517 # Handle error from the child processes
492518 if simulation_error_event .is_set ():
@@ -496,15 +522,21 @@ def __run_in_parallel(self, n_workers=None):
496522 "for more information."
497523 )
498524
525+ # An exit code cannot show a worker that left between
526+ # claiming an index and recording it.
527+ _refuse_logs_missing_a_simulation (
528+ self .input_file , self .output_file , self .number_of_simulations
529+ )
530+
499531 sim_monitor .print_final_status ()
500532
501533 # Handle error from the main process
502534 # pylint: disable=broad-except
503535 except (Exception , KeyboardInterrupt ) as error :
504- simulation_error_event . set ()
505-
506- for sim_producer in processes :
507- sim_producer . join ( )
536+ # Bounded here too. An unbounded join undid the bound above.
537+ _stop_the_workers_still_running (
538+ processes , simulation_error_event , _SHUTDOWN_GRACE_SECONDS
539+ )
508540
509541 if not isinstance (error , KeyboardInterrupt ):
510542 raise error
@@ -531,6 +563,8 @@ def __sim_producer(self, seed, sim_monitor, mutex, error_event): # pylint: disa
531563 error_event : multiprocess.Event
532564 Event signaling an error occurred during the simulation.
533565 """
566+ # The handler reads both, and a failure above the loop precedes them.
567+ sim_idx , inputs_json = None , ""
534568 try :
535569 # Ensure Processes generate different random numbers
536570 self .environment ._set_stochastic (seed )
@@ -567,18 +601,48 @@ def __sim_producer(self, seed, sim_monitor, mutex, error_event): # pylint: disa
567601 finally :
568602 mutex .release ()
569603
570- except Exception : # pylint: disable=broad-except
571- mutex .acquire ()
572- with open (self .error_file , "a" , encoding = "utf-8" ) as f :
573- f .write (inputs_json )
604+ # Nothing is in flight between two simulations, nor are these.
605+ sim_idx , inputs_json = None , ""
574606
575- # See note above: must use print() to remain visible from a
576- # multiprocessing worker process.
577- _SimMonitor .reprint (
578- f"Error on iteration { sim_idx } :\n { traceback .format_exc ()} "
579- )
607+ except Exception : # pylint: disable=broad-except
608+ if not self .__report_a_failed_simulation (
609+ sim_idx , inputs_json , mutex , error_event
610+ ):
611+ # The event could not be set; the exit code is what is left.
612+ raise
613+
614+ def __report_a_failed_simulation (self , sim_idx , inputs_json , mutex , error_event ):
615+ """Write down and announce a simulation this worker could not finish.
616+
617+ The event goes first and from outside the lock, since a worker that
618+ cannot write its diagnostics still has to be able to stop the others.
619+ Each step under the lock is suppressed on its own: a full disk would
620+ otherwise replace the failure being reported, and the lock is a
621+ manager's, so ending while holding it leaves the next worker waiting
622+ on a process that no longer exists.
623+ """
624+ details = traceback .format_exc ()
625+ where = "worker startup" if sim_idx is None else f"iteration { sim_idx } "
626+ announced = False
627+ with suppress (_MANAGER_IS_GONE ):
580628 error_event .set ()
581- mutex .release ()
629+ announced = True
630+
631+ held = False
632+ with suppress (* _MANAGER_IS_GONE ):
633+ held = mutex .acquire (timeout = _REPORT_LOCK_SECONDS )
634+ try :
635+ with suppress (OSError ):
636+ with open (self .error_file , "a" , encoding = "utf-8" ) as f :
637+ f .write (_worker_failure_record (where , details , inputs_json ))
638+ with suppress (OSError , ValueError ):
639+ # Must use print() to remain visible from a worker process.
640+ _SimMonitor .reprint (f"Error on { where } :\n { details } " )
641+ finally :
642+ if held :
643+ with suppress (* _MANAGER_IS_GONE ):
644+ mutex .release ()
645+ return announced
582646
583647 def __run_single_simulation (self ):
584648 """Runs a single simulation and returns the inputs and outputs.
@@ -983,6 +1047,13 @@ def _check_data_collector(self, data_collector):
9831047 "Invalid 'data_collector' key! "
9841048 f"Variable names overwrites 'export_list' key '{ key } '."
9851049 )
1050+ if key == _SIMULATION_INDEX_KEY :
1051+ raise ValueError (
1052+ f"Invalid 'data_collector' key '{ key } '! It is the "
1053+ f"number of the simulation the row belongs to, which "
1054+ f"is written after the collectors run and cannot be "
1055+ f"replaced by one."
1056+ )
9861057 if not callable (callback ):
9871058 raise ValueError (
9881059 f"Invalid value in 'data_collector' for key '{ key } '! "
@@ -1755,6 +1826,192 @@ def export_errors_to_json(self, filename):
17551826 self ._write_log_to_json (self .errors_log , filename )
17561827
17571828
1829+ # Prompt enough to notice a dead worker, cheap enough over a run of hours.
1830+ _JOIN_POLL_SECONDS = 0.2
1831+ _SHUTDOWN_GRACE_SECONDS = 5.0
1832+
1833+
1834+ def _ended_badly (worker ):
1835+ """Whether a worker has stopped, and stopped for the wrong reason."""
1836+ return worker .exitcode not in (None , 0 )
1837+
1838+
1839+ def _a_failure_was_reported (error_event ):
1840+ """Whether a worker has said it failed, false if it cannot be asked."""
1841+ with suppress (* _MANAGER_IS_GONE ):
1842+ return error_event .is_set ()
1843+ return False
1844+
1845+
1846+ def _wait_for_the_workers (processes , seconds ):
1847+ """Join every worker against one shared deadline, not one each.
1848+
1849+ Monotonic, since a clock correction would move a wall-clock deadline.
1850+ """
1851+ deadline = monotonic () + seconds
1852+ for worker in processes :
1853+ worker .join (timeout = max (0.0 , deadline - monotonic ()))
1854+
1855+
1856+ def _stop_the_workers_still_running (processes , error_event , grace_period ):
1857+ """Ask the rest to stop, end what cannot, kill what outlives that.
1858+
1859+ Asked first because a worker between simulations reads the event and leaves
1860+ with its logs intact. One blocked on a lock its dead sibling was holding
1861+ never reaches that check. Terminate runs no handlers, so it comes second,
1862+ and a worker can still ignore it.
1863+ """
1864+ with suppress (_MANAGER_IS_GONE ):
1865+ error_event .set ()
1866+ _wait_for_the_workers (processes , grace_period )
1867+
1868+ for worker in processes :
1869+ if worker .is_alive ():
1870+ worker .terminate ()
1871+ _wait_for_the_workers (processes , grace_period )
1872+
1873+ for worker in processes :
1874+ if worker .is_alive ():
1875+ worker .kill ()
1876+ _wait_for_the_workers (processes , grace_period )
1877+
1878+
1879+ def _join_the_workers (processes , error_event , grace_period = _SHUTDOWN_GRACE_SECONDS ):
1880+ """Wait for the workers, and stop once one of them has failed.
1881+
1882+ A reported failure ends the wait as well as a bad exit code, since a
1883+ worker that reports one leaves cleanly and says nothing through its exit
1884+ status. Its siblings read the event between simulations, but one blocked
1885+ on a lock nobody owns never reaches that check, and the run is already
1886+ short a simulation either way, so the wait is bounded here rather than
1887+ left to them. The reported path gets the longer grace: those siblings are
1888+ working, not stuck.
1889+
1890+ Slowness alone ends nothing. With no failure reported a healthy worker is
1891+ given as long as it needs.
1892+ """
1893+ while any (worker .is_alive () for worker in processes ):
1894+ for worker in processes :
1895+ worker .join (timeout = _JOIN_POLL_SECONDS )
1896+ if any (_ended_badly (worker ) for worker in processes ):
1897+ _stop_the_workers_still_running (processes , error_event , grace_period )
1898+ return
1899+ if _a_failure_was_reported (error_event ):
1900+ _stop_the_workers_still_running (
1901+ processes , error_event , _REPORTED_FAILURE_GRACE_SECONDS
1902+ )
1903+ return
1904+
1905+
1906+ def _worker_failure_record (where , details , inputs_json = "" ):
1907+ """A row saying what failed, and what the simulation had drawn so far.
1908+
1909+ The inputs alone left the error file with no stage and no traceback, which
1910+ is what the caller is sent there to read.
1911+ """
1912+ record = {"index" : None , "stage" : where , "error" : details }
1913+ with suppress (ValueError ):
1914+ drawn = json .loads (inputs_json )
1915+ if isinstance (drawn , dict ):
1916+ record ["index" ] = drawn .get ("index" )
1917+ record ["inputs" ] = drawn
1918+ return json .dumps (record ) + "\n "
1919+
1920+
1921+ def _indices_a_log_holds (path ):
1922+ """Every index a log records, in order, and ``None`` for a row it cannot."""
1923+ found = []
1924+ with open (path , "r" , encoding = "utf-8" ) as recorded :
1925+ for line in recorded :
1926+ if not line .strip ():
1927+ continue
1928+ try :
1929+ index = json .loads (line )["index" ]
1930+ except (ValueError , KeyError , TypeError ):
1931+ found .append (None )
1932+ continue
1933+ usable = (
1934+ isinstance (index , int ) and not isinstance (index , bool ) and index >= 0
1935+ )
1936+ found .append (index if usable else None )
1937+ return found
1938+
1939+
1940+ def _refuse_logs_missing_a_simulation (input_file , output_file , target ):
1941+ """Raise unless both logs hold every simulation the run was asked for.
1942+
1943+ An exit code says how a worker ended, never whether the index it had
1944+ already claimed reached the logs, and the monitor counts claims rather than
1945+ rows. A worker that leaves between the two is invisible to everything else
1946+ here, so the logs themselves are what the run is judged on.
1947+
1948+ Rows numbered past the target are left alone: an append given a smaller
1949+ target than the checkpoint already holds is an append question, not a lost
1950+ simulation. What each log holds still has to be the consecutive run it
1951+ claims to be, so its indices are required to be exactly as many as its
1952+ rows, which refuses a stray number and a hole without needing to be told
1953+ how long the checkpoint was. The two logs must also agree row for row,
1954+ since a record goes into both under one lock. Streamed rather than read
1955+ through ``_read_log_file``, which would hold every row in memory.
1956+ """
1957+ wanted = set (range (target ))
1958+ recorded = {}
1959+ for label , path in (("input" , input_file ), ("output" , output_file )):
1960+ found = _indices_a_log_holds (path )
1961+ recorded [label ] = found
1962+ held = set (found )
1963+ if None in held :
1964+ raise RuntimeError (
1965+ f"The run is incomplete: the { label } log has rows that cannot "
1966+ f"be read, so what it holds cannot be established."
1967+ )
1968+ if len (found ) != len (held ):
1969+ raise RuntimeError (
1970+ f"The run is incomplete: the { label } log records "
1971+ f"{ len (found ) - len (held )} simulation(s) more than once."
1972+ )
1973+ missing = sorted (wanted - held )
1974+ if missing :
1975+ raise RuntimeError (
1976+ f"The run is incomplete: the { label } log is missing "
1977+ f"{ len (missing )} of { target } simulations, the first being "
1978+ f"{ missing [0 ]} ."
1979+ )
1980+ strays = sorted (held - set (range (len (found ))))
1981+ if strays :
1982+ raise RuntimeError (
1983+ f"The run is incomplete: the { label } log numbers a simulation "
1984+ f"{ strays [0 ]} , past the { len (found )} it holds, so what it "
1985+ f"records is not one run of consecutive simulations."
1986+ )
1987+
1988+ if recorded ["input" ] != recorded ["output" ]:
1989+ raise RuntimeError (
1990+ "The run is incomplete: the input and output logs do not record "
1991+ "the same simulations in the same order. A record is written to "
1992+ "both under one lock, so they hold two different runs."
1993+ )
1994+
1995+
1996+ def _refuse_a_worker_that_did_not_finish (processes ):
1997+ """Raise if any worker left without exiting cleanly.
1998+
1999+ A negative code is the signal that ended it, ``None`` one still running.
2000+ """
2001+ unfinished = [
2002+ f"worker { position } with exit code { process .exitcode } "
2003+ for position , process in enumerate (processes )
2004+ if process .exitcode != 0
2005+ ]
2006+ if not unfinished :
2007+ return
2008+ raise RuntimeError (
2009+ f"The run is incomplete: { ', ' .join (unfinished )} . A worker that ends "
2010+ "this way records nothing and cannot say why, so the simulations it "
2011+ "held are missing from the results."
2012+ )
2013+
2014+
17582015def _import_multiprocess ():
17592016 """Import the necessary modules and submodules for the
17602017 multiprocess library.
0 commit comments