Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
152 changes: 108 additions & 44 deletions MadSpin/interface_madspin.py
Original file line number Diff line number Diff line change
Expand Up @@ -2741,8 +2741,9 @@ def generate_events_mg7(self, decay_dir, nb_event, run_name='run_01'):

The pool is left in the layout the parallel unweighting expects of any
backend -- one LHE file per worker, under ``Events/<run_name>/`` -- so
that mg7 goes down exactly the same paths madevent does. See below for
why LHE rather than the (faster to read) numpy event file.
that mg7 goes down exactly the same paths madevent does. mg7 writes the
per-worker files itself (see below); see further below for why LHE
rather than the (faster to read) numpy event file.
"""
run_card_path = pjoin(decay_dir, 'Cards', 'run_card.toml')
run_card = banner.RunCardMG7(run_card_path)
Expand All @@ -2762,6 +2763,24 @@ def generate_events_mg7(self, decay_dir, nb_event, run_name='run_01'):
# the owner->waiter publish contract, i.e. an extension of the one
# protocol surface where a publish/open race can live. Not worth 5%.
run_card['run']['output_format'] = 'lhe'

# One file per unweighting worker. mg7 writes them itself -- madspace's
# combine_to_lhe deals the events out round-robin as it formats them --
# so the pool is never materialised as one file and then split. That is
# what ``nb_output_files`` asks for; it is mg7's ``nb_unweight_output``.
#
# mg7 still names its own Events/run_NN directory, so the files land
# there and have to be moved to the paths the rest of MadSpin expects of
# any backend (below). The move is N renames within one filesystem, and
# it buys back the atomicity the refill protocol wants: the canonical
# directory never contains a half-written pool, which is the same reason
# _materialise_refill_pool builds in a .part directory.
#
# Each file gets its own <init> block with the channel's partial width,
# so it is still self-describing (see below). Left uncompressed: they
# are read back immediately and gzip is pure overhead here.
nb_split = self._decay_pool_split()
run_card['run']['nb_output_files'] = max(1, int(nb_split))
run_card.write(run_card_path)
with open(pjoin(decay_dir, 'Cards', 'param_card.dat'), 'w') as fsock:
fsock.write(self.banner['slha'])
Expand All @@ -2784,19 +2803,68 @@ def generate_events_mg7(self, decay_dir, nb_event, run_name='run_01'):
'the mg7 decay generator failed in %s (exit %s); see %s'
% (decay_dir, returncode, log_path))
after = set(misc.glob('*', events_dir)) if os.path.isdir(events_dir) else set()
new_runs = sorted(after - before)
# mg7 names its own run directory (Events/run_NN) and drops info.json in
# it, which is what identifies it among anything else that may have
# appeared under Events/ while the launcher ran.
new_runs = sorted(path for path in after - before
if os.path.exists(pjoin(path, 'info.json')))
if not new_runs:
raise self.InvalidCmd(
'the mg7 decay generator produced no run directory in %s' % events_dir)
run_dir = new_runs[-1]

for name in ('events.lhe', 'events.lhe.gz'):
events_path = pjoin(run_dir, name)
if os.path.exists(events_path):
break
targets = []
if nb_split > 1:
# Where the rest of MadSpin expects any backend's pool to be: the
# same paths madevent reaches through ``nb_unweight_output``. Two
# things fall out of the pool being THERE rather than anywhere else:
# a worker opens its own file instead of striding the whole pool
# (which costs it a full parse of every other worker's events), and
# on a refill these paths *are* _refill_pool_paths, so
# _generate_refill_pool finds sources == targets and does no further
# work.
targets = lhe_parser.EventFile.unweight_output_paths(
pjoin(decay_dir, 'Events', run_name, 'unweighted_events.lhe'),
nb_split)
sources = [pjoin(run_dir, 'events_%d.lhe' % i)
for i in range(nb_split)]
single = pjoin(run_dir, 'events.lhe')
target_dir = os.path.dirname(targets[0])
if not os.path.isdir(target_dir):
os.makedirs(target_dir)
if all(os.path.exists(path) for path in sources):
# The whole point of nb_output_files: mg7 already wrote one
# file per worker, so there is nothing to split -- only N
# renames within the same filesystem, which also keep the
# canonical directory from ever holding a half-written pool.
for source, target in zip(sources, targets):
os.replace(source, target)
elif os.path.exists(single):
# A madspace too old to write several files wrote one instead
# (it says so in the log). Deal it out here, which is what
# MadSpin did for every mg7 run before nb_output_files existed.
logger.info('mg7 wrote a single pool; splitting it in python. '
'Rebuild madspace to have mg7 write the per-worker '
'files directly.')
with self._phase('decay_mg7_split'):
self._split_pool_round_robin([single], targets)
try:
os.remove(single)
except OSError:
pass
else:
raise self.InvalidCmd(
'the mg7 decay generator wrote neither %s nor %s (see %s)'
% (sources[0], single, log_path))
banner_path = targets[0]
else:
raise self.InvalidCmd(
'the mg7 decay generator produced no LHE file in %s' % run_dir)
for name in ('events.lhe', 'events.lhe.gz'):
banner_path = pjoin(run_dir, name)
if os.path.exists(banner_path):
break
else:
raise self.InvalidCmd(
'the mg7 decay generator produced no LHE file in %s' % run_dir)

# Timing breakdown, best effort: the launcher reports how long the
# integration itself took, and the rest of the subprocess wall time is
Expand All @@ -2815,44 +2883,15 @@ def generate_events_mg7(self, decay_dir, nb_event, run_name='run_01'):
for stage in run_times.values()
if isinstance(stage, dict)))

pool = lhe_parser.EventFile(events_path)
# <init> of a 1 -> n process carries a width in GeV, not a cross section
pool = lhe_parser.EventFile(banner_path)
# <init> of a 1 -> n process carries a width in GeV, not a cross
# section. Every slice carries the same banner, so slice 0 answers for
# the pool -- which is also what makes a slice openable by path alone.
width = pool.cross

# Deal the pool out into one file per unweighting worker, at the paths
# the rest of MadSpin already expects any backend to use -- the same
# ones madevent reaches through the run_card's ``nb_unweight_output``.
# Two things fall out of writing them THERE rather than anywhere else:
# a worker opens its own file instead of striding the whole pool (which
# costs it a full parse of every other worker's events), and on a refill
# these paths *are* _refill_pool_paths, so _generate_refill_pool finds
# sources == targets and does no further work. Left uncompressed: they
# are read back immediately and gzip is pure overhead here.
nb_split = self._decay_pool_split()
if nb_split <= 1:
if not targets:
pool.seek(0)
return pool, width
pool.close()
targets = lhe_parser.EventFile.unweight_output_paths(
pjoin(decay_dir, 'Events', run_name, 'unweighted_events.lhe'),
nb_split)
target_dir = os.path.dirname(targets[0])
if not os.path.isdir(target_dir):
os.makedirs(target_dir)
with self._phase('decay_mg7_split'):
self._split_pool_round_robin([events_path], targets)
# Between them the slices hold every event of the source -- the split
# reads it to the end or raises -- so keeping the source doubles what a
# pool costs on disk: another 241 MB per channel on the 100k benchmark,
# and a refill adds a generation rather than replacing one. Nothing
# reads it after this point (each slice carries its own banner, and the
# run directory keeps info.json and the generation log), so drop it once
# every slice is actually there.
if all(os.path.exists(path) for path in targets):
try:
os.remove(events_path)
except OSError:
pass
return _ChainedEvents(targets), width

# File in which a decay directory records the partial width that was
Expand Down Expand Up @@ -4518,6 +4557,13 @@ def _generate_decay_entry(self, pdg, job, nb_core, seed_offset, res_path):
# multiplier of its own so a generation child can never land on the
# stream of an unweighting/max-weight worker (those use 7919).
random.seed((int(self.seed) if self.seed else 0) + 104729 * seed_offset)
# Start the child's accounting from zero. It inherited the parent's
# dicts through the fork, and everything already in them has been
# charged in the parent too -- reporting it back would double it.
times = self.__dict__.get('_phase_times')
if times is not None:
times.clear()
self._phase_counts.clear()
self._gen_nb_core = nb_core
self.seed = (int(self.seed) + 1000003 * seed_offset) % (30081 * 30081)
self.options['seed'] = self.seed
Expand All @@ -4533,7 +4579,16 @@ def _generate_decay_entry(self, pdg, job, nb_core, seed_offset, res_path):
for k, v in out.items()),
'width': width,
'channel_widths': dict((str(k), v) for k, v
in channel_widths.items())},
in channel_widths.items()),
# The phase timings too, or every phase of the decay
# generation -- decay_event_generation, the largest
# of them all -- dies here with the child and is
# simply missing from the run's report. See
# _generate_decays for how they are added back.
'phase_times': dict(
self.__dict__.get('_phase_times') or {}),
'phase_counts': dict(
self.__dict__.get('_phase_counts') or {})},
fp)
except Exception as exc:
import traceback
Expand Down Expand Up @@ -4593,6 +4648,15 @@ def _generate_decays(self, gen_jobs, mg5):
if 'error' in data:
raise Exception("MadSpin: the decay generation of pdg %s failed:\n%s"
% (pdg, data.get('tb', data['error'])))
# Put the child's phase timings back into this process's accounting.
# They are a SUM OVER PARTICLES of work that ran concurrently, so
# they are work done, not wall time, and can exceed the wall time of
# the fork that contained them -- which is why they are charged to
# their own phases and never to a parent phase measured out here.
for name, dt in (data.get('phase_times') or {}).items():
self._add_phase(name, dt, count=0)
for name, count in (data.get('phase_counts') or {}).items():
self._add_count(name, count)
channel_widths = dict((int(k), v) for k, v
in data.get('channel_widths', {}).items())
# The per-channel width goes back in with the paths: an mg7 pool has
Expand Down
56 changes: 52 additions & 4 deletions madgraph/iolibs/template_files/mg7/madevent.py
Original file line number Diff line number Diff line change
Expand Up @@ -968,14 +968,62 @@ def generate_events(self) -> None:
)
elif output_format == "lhe":
self.lhe_completer = self.build_lhe_completer()
self.event_generator.combine_to_lhe(
os.path.join(self.run_path, "events.lhe"), self.lhe_completer,
self.build_lhe_meta(),
)
self.combine_to_lhe_files(self.lhe_output_files())
else:
raise ValueError("Unknown output format")
self.save_gridpack()

def lhe_output_files(self) -> list:
"""Where combine_to_lhe should write, as a list of paths.

One ``<run_path>/events.lhe`` normally; with [run] nb_output_files = N,
``<run_path>/events_0.lhe`` ... ``events_{N-1}.lhe``, into which
madspace deals the events round-robin as it formats them. Every one of
them is a complete LHE file, with its own <init> block, so a reader can
open any of them by path alone.

This exists for a consumer that reads the events back in parallel: with
N files each of its workers parses only its own share, instead of
reading the whole stream and skipping the (N-1)/N of it that belongs to
somebody else. It is the mg7 counterpart of madevent's
``nb_unweight_output``, and MadSpin's decay pools are what asked for it.
"""
count = int(self.run_card["run"]["nb_output_files"])
if count <= 1:
return [os.path.join(self.run_path, "events.lhe")]
return [os.path.join(self.run_path, "events_%d.lhe" % i)
for i in range(count)]

def combine_to_lhe_files(self, files) -> None:
"""Write the combined events into ``files``, using madspace's
multi-file combine_to_lhe, which fans them out as it formats them and
so never materialises the whole stream in one file.

madspace is not upgraded in lockstep with this template -- it can be an
older binary wheel -- and one built before that overload exists only
accepts a single path. Rather than reimplement the fan-out here, fall
back to the single ``events.lhe`` such a build can write and say so: a
caller that asked for several files is a caller that already knows how
to split one, and it can see from what is on disk which it got.
"""
meta = self.build_lhe_meta()
single = os.path.join(self.run_path, "events.lhe")
if len(files) == 1:
# the str overload, so an older madspace is fine here too
self.event_generator.combine_to_lhe(files[0], self.lhe_completer, meta)
return
# LHEMultiFileWriter is the class the list overload is built on, so its
# presence is exactly the question "can this madspace write several
# files itself?".
if hasattr(ms, "LHEMultiFileWriter"):
self.event_generator.combine_to_lhe(files, self.lhe_completer, meta)
return
logging.getLogger("madevent").warning(
"nb_output_files = %d, but this madspace build can only write one "
"LHE file; writing %s instead", len(files), single
)
self.event_generator.combine_to_lhe(single, self.lhe_completer, meta)

@staticmethod
def _histogram_mean(hist):
"""Cross-section-weighted mean of a histogrammed observable."""
Expand Down
1 change: 1 addition & 0 deletions madgraph/iolibs/template_files/mg7/run_card.toml
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ cpu_thread_pool_size = %(run.cpu_thread_pool_size)s
gpu_thread_pool_size = %(run.gpu_thread_pool_size)s
combine_thread_pool_size = %(run.combine_thread_pool_size)s
output_format = %(run.output_format)s # options: compact_npy, lhe_npy, lhe
nb_output_files = %(run.nb_output_files)s # split the events over this many LHE files (round-robin, each complete on its own)
verbosity = %(run.verbosity)s # options: silent, pretty, log, auto (pretty if attached to a terminal, log otherwise)
dummy_matrix_element = %(run.dummy_matrix_element)s

Expand Down
2 changes: 2 additions & 0 deletions madgraph/various/banner.py
Original file line number Diff line number Diff line change
Expand Up @@ -6490,6 +6490,8 @@ def default_setup(self):
self.add_toml_param('run', 'combine_thread_pool_size', -1, gridpack=True)
self.add_toml_param('run', 'output_format', "lhe", gridpack=True,
allowed=['compact_npy', 'lhe_npy', 'lhe'])
self.add_toml_param('run', 'nb_output_files', 1, gridpack=True,
comment="split the events over this many LHE files (round-robin, each a complete file with its own <init>); mirrors madevent's nb_unweight_output. Only for output_format = 'lhe'")
self.add_toml_param('run', 'verbosity', "auto", gridpack=True,
allowed=['silent', 'pretty', 'log', 'auto'])
self.add_toml_param('run', 'dummy_matrix_element', False)
Expand Down
10 changes: 10 additions & 0 deletions madspace/include/madspace/driver/event_generator.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,16 @@ class EventGenerator {
LHECompleter& lhe_completer,
const LHEMeta& meta = {}
);
// Same, but dealing the combined events round-robin into several LHE files
// instead of one. Each file is complete on its own -- same header, same
// <init> block with the process cross sections -- so a consumer that wants
// to read the events in parallel can give one file to each of its workers
// and have every worker parse only its own share. See LHEMultiFileWriter.
void combine_to_lhe(
const std::vector<std::string>& file_names,
LHECompleter& lhe_completer,
const LHEMeta& meta = {}
);
GeneratorStatus status() const { return _status; }
std::vector<GeneratorStatus> channel_status() const;
std::vector<Histogram> histograms() const;
Expand Down
40 changes: 40 additions & 0 deletions madspace/include/madspace/driver/lhe_output.hpp
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
#pragma once

#include <fstream>
#include <memory>
#include <random>
#include <string>
#include <unordered_map>
Expand Down Expand Up @@ -193,4 +194,43 @@ class LHEFileWriter {
std::string _buffer;
};

// Writes one event stream into several LHE files at once, dealing the events
// out round-robin. Every file is a complete, self-describing LHE file: it gets
// its own copy of `meta`, so its own <header> and its own <init> block with the
// process cross sections. A consumer can therefore open any one of them by path
// alone, with nothing published beside it.
//
// The point of the fan-out is a reader that wants to consume the events in
// parallel: with N files it opens the one that is its own, instead of reading
// the whole stream and skipping the (N-1)/N of it that belongs to somebody
// else. Round-robin rather than contiguous blocks because the files are then
// balanced to within one event whatever the chunking, and because the split
// cannot correlate with the order the stream happens to arrive in.
//
// Chunked, off-thread formatting is supported through reserve()/file_index():
// reserve() hands out a run of consecutive positions in the stream, file_index()
// says where each of them belongs, and the formatted text can then be handed
// over with write_string() whenever it is ready -- in any order.
class LHEMultiFileWriter {
public:
LHEMultiFileWriter(const std::vector<std::string>& file_names, const LHEMeta& meta);

std::size_t file_count() const { return _writers.size(); }
// Which file the `event_index`-th event of the stream belongs to.
std::size_t file_index(std::size_t event_index) const {
return event_index % _writers.size();
}
// Claim `count` consecutive positions in the stream; returns the first.
std::size_t reserve(std::size_t count);
// Total number of positions claimed so far.
std::size_t event_count() const { return _event_count; }
void write_string(std::size_t file, const std::string& str);
// Convenience for a caller that formats one event at a time.
void write(const LHEEvent& event);

private:
std::vector<std::unique_ptr<LHEFileWriter>> _writers;
std::size_t _event_count = 0;
};

} // namespace madspace
Loading
Loading