From c1c66705db83e6ead04c34f19d5b281697949681 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 14:20:36 -0400 Subject: [PATCH 01/62] relax permissions and use decoder pool and new convenience method --- .../org/jlab/clas/reco/EngineProcessor.java | 21 +++++++++++-------- 1 file changed, 12 insertions(+), 9 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java index 39c93eefd4..1d68ec7e9a 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java @@ -18,11 +18,10 @@ import org.jlab.clara.engine.EngineDataType; import java.util.Arrays; import org.jlab.coda.jevio.EvioException; -import org.jlab.detector.decode.CLASDecoder4; +import org.jlab.detector.decode.CLASDecoder; +import org.jlab.detector.decode.CLASDecoderPool; import org.jlab.io.evio.EvioDataEvent; import org.jlab.io.evio.EvioSource; -import org.jlab.io.hipo.HipoDataEvent; -import org.jlab.jnp.hipo4.data.Event; import org.jlab.jnp.hipo4.data.SchemaFactory; import org.json.JSONObject; import org.jlab.utils.ClaraYaml; @@ -36,13 +35,13 @@ public class EngineProcessor { public static final String ENGINE_CLASS_BG = "org.jlab.service.bg.BackgroundEngine"; public static final String ENGINE_CLASS_PP = "org.jlab.service.postproc.PostprocEngine"; - private final Map processorEngines = new LinkedHashMap<>(); + protected final Map processorEngines = new LinkedHashMap<>(); private static final Logger LOGGER = Logger.getLogger(EngineProcessor.class.getPackage().getName()); private boolean updateDictionary = true; private SchemaFactory banksToKeep = null; private final List schemaExempt = Arrays.asList("RUN::config","DC::tdc"); - private CLASDecoder4 decoder = new CLASDecoder4(); + protected final CLASDecoderPool decoders = new CLASDecoderPool(64,"default",null); public EngineProcessor(){} @@ -94,7 +93,7 @@ private void setPreloadFiles(String filenames, boolean restream, boolean rebuild findEngine(ENGINE_CLASS_PP).init(); } - private void updateDictionary(HipoDataSource source, HipoDataSync sync){ + protected void updateDictionary(HipoDataSource source, HipoDataSync sync){ SchemaFactory fsync = sync.getWriter().getSchemaFactory(); SchemaFactory fsrc = source.getReader().getSchemaFactory(); List schemaList = fsync.getSchemaKeys(); @@ -342,9 +341,13 @@ public void processFile(EvioSource reader, HipoDataSync writer, int skipEvents, ByteBuffer bb = reader.getEventBuffer(eventsRead, true); if (skipEvents <= 0 || eventsRead > skipEvents) { EvioDataEvent evio = new EvioDataEvent(bb.array(), ByteOrder.LITTLE_ENDIAN); - Event hipo = decoder.getDecodedEvent(evio, -1, eventsRead, null, null); - HipoDataEvent hipo2 = new HipoDataEvent(hipo, decoder.getSchemaFactory()); - processEvent(hipo2, writer); + try { + CLASDecoder d = decoders.take(); + processEvent(d.getDecodedDataEvenet(evio), writer); + decoders.put(d); + } catch (InterruptedException ex) { + System.getLogger(EngineProcessor.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex); + } } if (maxEvents > 0 && eventsRead > maxEvents+skipEvents) break; } catch (EvioException ex) { From 76ab3607db2f432cf7da1e110be6ac23dd065c17 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 14:21:10 -0400 Subject: [PATCH 02/62] add multi-threaded recon-util --- bin/recon-mutil | 12 ++ .../jlab/clas/reco/EngineMultiProcessor.java | 193 ++++++++++++++++++ validation/advanced-tests/run-eb-tests.sh | 2 +- 3 files changed, 206 insertions(+), 1 deletion(-) create mode 100644 bin/recon-mutil create mode 100644 common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java diff --git a/bin/recon-mutil b/bin/recon-mutil new file mode 100644 index 0000000000..003b3ca8ca --- /dev/null +++ b/bin/recon-mutil @@ -0,0 +1,12 @@ +#!/bin/bash + +. `dirname $0`/../libexec/env.sh + +split_cli $@ + +export MALLOC_ARENA_MAX=1 + +java ${JAVA_OPTS-} -Xms10240m -XX:+UseSerialGC ${jvm_options[@]} \ + -cp ${COATJAVA_CLASSPATH:-''} \ + org.jlab.clas.reco.EngineMultiProcessor \ + ${class_options[@]} diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java new file mode 100644 index 0000000000..49ccc6956a --- /dev/null +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -0,0 +1,193 @@ +package org.jlab.clas.reco; + +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentLinkedQueue; +import org.jlab.coda.jevio.EvioException; +import org.jlab.detector.decode.CLASDecoder; +import org.jlab.io.base.DataEvent; +import org.jlab.io.base.DataSource; +import org.jlab.io.evio.EvioDataEvent; +import org.jlab.io.evio.EvioSource; +import org.jlab.io.hipo.HipoDataEvent; +import org.jlab.io.hipo.HipoDataSource; +import org.jlab.io.hipo.HipoDataSync; +import org.jlab.utils.benchmark.Benchmark; +import org.jlab.utils.benchmark.ProgressPrintout; + +/** + * + * @author baltzell + */ +public class EngineMultiProcessor extends EngineProcessor { + + public EngineMultiProcessor(int threads) { + super(); + this.threads = threads; + } + + public EngineMultiProcessor(int threads, int events, int skip) { + super(); + this.threads = threads; + this.maxEventsUser = events; + this.skipEvents = skip; + } + + public void process(String output, String... input) { + readerThread = CompletableFuture.runAsync(() -> { read(input); }); + writerThread = CompletableFuture.runAsync(() -> { write(output); }); + for (int i=0; i { process(j); })); + } + while (!writerThread.isDone()) + try { Thread.sleep(100); } catch (InterruptedException ex) {} + } + + DataSource reader; + HipoDataSync writer; + CompletableFuture readerThread; + CompletableFuture writerThread; + + int threads; + int maxEvents = 0; + int maxEventsUser = 0; + int skipEvents = 0; + int readEvents = 0; + int writeEvents = 0; + + ArrayList inputs = new ArrayList<>(); + + ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); + ConcurrentLinkedQueue evioQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue hipoQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); + + ProgressPrintout progress = new ProgressPrintout(); + + void read(String... input) { + inputs.addAll(Arrays.asList(input)); + while (maxEvents < 1 || readEvents < maxEvents) { + if (reader != null && reader.hasEvent()) { + if (evioQueue.size()+hipoQueue.size() > 100*threads) { + try { Thread.sleep(100); } + catch (InterruptedException ex) {} + } + else { + Benchmark.getInstance().resume("read"); + readEvents++; + if (reader instanceof EvioSource evio) { + try { evioQueue.offer(evio.getEventBuffer(readEvents, true)); } + catch (EvioException ex) { ex.printStackTrace(); } + } + else { + DataEvent event = reader.getNextEvent(); + if (skipEvents < 1 || readEvents > skipEvents) + hipoQueue.offer(event); + } + Benchmark.getInstance().pause("read"); + } + } + else if (inputs.isEmpty()) break; + else { + if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); + else reader = new EvioSource(); + reader.open(inputs.remove(0)); + maxEvents = maxEventsUser; + if (reader instanceof HipoDataSource hipo) + updateDictionary(hipo, writer); + else { + int n = ((EvioSource)reader).getEventCount(); + maxEvents = maxEventsUser < n ? maxEventsUser : n; + } + readEvents = 0; + } + } + } + + EvioDataEvent process(ByteBuffer bytes) { + Benchmark.getInstance().resume("EVIO"); + EvioDataEvent e = new EvioDataEvent(bytes.array(), ByteOrder.LITTLE_ENDIAN); + Benchmark.getInstance().pause("EVIO"); + return e; + } + + HipoDataEvent process(EvioDataEvent event) { + Benchmark.getInstance().resume("DECO"); + HipoDataEvent e; + try { + CLASDecoder d = decoders.take(); + e = d.getDecodedDataEvenet(event); + decoders.put(d); + } + catch (InterruptedException ex) { e = null; } + Benchmark.getInstance().pause("DECO"); + return e; + } + + void process(HipoDataEvent event) { + for (Map.Entry engine : processorEngines.entrySet()) { + Benchmark.getInstance().resume(engine.getValue().getName()); + try { engine.getValue().processDataEvent(event); } + catch (Exception ex) { ex.printStackTrace(); } + Benchmark.getInstance().pause(engine.getValue().getName()); + } + } + + void process(int thread) { + while (true) { + if (evioQueue.isEmpty() && hipoQueue.isEmpty()) { + if (readerThread.isDone()) + if (evioQueue.isEmpty() && hipoQueue.isEmpty()) break; + try { Thread.sleep(100); } + catch (InterruptedException ex) {} + } + else { + DataEvent event; + if (!evioQueue.isEmpty()) event = process(process(evioQueue.poll())); + else if (!hipoQueue.isEmpty()) event = hipoQueue.poll(); + else continue; + process((HipoDataEvent)event); + writeQueue.offer(event); + } + } + } + + void write(String output) { + writer = new HipoDataSync(); + writer.setCompressionType(2); + writer.open(output); + while (true) { + if (writeQueue.isEmpty()) { + for (CompletableFuture f : procThreads) + if (f.isDone()) procThreads.remove(f); + if (procThreads.isEmpty()) { + if (writeQueue.isEmpty()) { + writer.close(); + System.out.println(Benchmark.getInstance()); + System.out.println(String.format("recon-mutil::::: Read/Write/Diff = %d/%d/%d", + readEvents, writeEvents, readEvents-writeEvents)); + break; + } + } + try { Thread.sleep(100); } + catch (InterruptedException ex) {} + } + else write(writeQueue.poll()); + } + } + + void write(DataEvent event) { + Benchmark.getInstance().resume("write"); + writer.writeEvent(writeQueue.poll()); + if (writeEvents > 100) progress.updateStatus(); + if (writeEvents == 101) Benchmark.getInstance().printTimer(10); + writeEvents++; + Benchmark.getInstance().pause("write"); + } + +} diff --git a/validation/advanced-tests/run-eb-tests.sh b/validation/advanced-tests/run-eb-tests.sh index 05c7470db5..b21c234a1c 100755 --- a/validation/advanced-tests/run-eb-tests.sh +++ b/validation/advanced-tests/run-eb-tests.sh @@ -49,7 +49,7 @@ if [ $? != 0 ] ; then echo "EBTwoTrackTest compilation failure" ; exit 1 ; fi # run reconstruction: rm -f out_${stub}.hipo -../../coatjava/bin/recon-util -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 +../../coatjava/bin/recon-mutil -t 6 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 # run EB tests: java -Xmx1536m -Xms1024m -cp $classPath -DINPUTFILE=out_${stub}.hipo eb.EBTwoTrackTest From c047bcd9f435ab6b64e88fc8279096d01d143f92 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 14:41:55 -0400 Subject: [PATCH 03/62] break it up --- .../src/main/java/org/jlab/clas/reco/EngineProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java index 1d68ec7e9a..f2bad4fdd9 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java @@ -397,7 +397,6 @@ protected static OptionParser getParser() { OptionParser parser = new OptionParser("recon-util"); parser.addRequired("-o","output.hipo"); parser.addRequired("-i","input.evio/hipo"); - parser.setRequiresInputList(false); parser.addOption("-c","0","use default configuration [0 - no, 1 - yes/default, 2 - all services] "); parser.addOption("-s","-1","number of events to skip"); parser.addOption("-n","-1","number of events to process"); @@ -408,6 +407,7 @@ protected static OptionParser getParser() { parser.addOption("-P",null,"preload file for post-processing"); parser.addOption("-R","0","rebuild scalers"); parser.addOption("-H","0","restream helicity"); + parser.setRequiresInputList(false); return parser; } From fa7f4328f23a3052d74bb30f69cb31d725868c40 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 16:01:08 -0400 Subject: [PATCH 04/62] fininsh --- bin/recon-mutil | 2 +- .../jlab/clas/reco/EngineMultiProcessor.java | 222 ++++++++++-------- .../org/jlab/clas/reco/EngineProcessor.java | 3 +- 3 files changed, 123 insertions(+), 104 deletions(-) mode change 100644 => 100755 bin/recon-mutil diff --git a/bin/recon-mutil b/bin/recon-mutil old mode 100644 new mode 100755 index 003b3ca8ca..533f8130bb --- a/bin/recon-mutil +++ b/bin/recon-mutil @@ -6,7 +6,7 @@ split_cli $@ export MALLOC_ARENA_MAX=1 -java ${JAVA_OPTS-} -Xms10240m -XX:+UseSerialGC ${jvm_options[@]} \ +java ${JAVA_OPTS-} -Xms10240m -XX:+UseParallelGC ${jvm_options[@]} \ -cp ${COATJAVA_CLASSPATH:-''} \ org.jlab.clas.reco.EngineMultiProcessor \ ${class_options[@]} diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 49ccc6956a..acfeab53b0 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -18,6 +18,7 @@ import org.jlab.io.hipo.HipoDataSync; import org.jlab.utils.benchmark.Benchmark; import org.jlab.utils.benchmark.ProgressPrintout; +import org.jlab.utils.options.OptionParser; /** * @@ -25,34 +26,17 @@ */ public class EngineMultiProcessor extends EngineProcessor { - public EngineMultiProcessor(int threads) { - super(); - this.threads = threads; - } - - public EngineMultiProcessor(int threads, int events, int skip) { - super(); - this.threads = threads; - this.maxEventsUser = events; - this.skipEvents = skip; - } - - public void process(String output, String... input) { - readerThread = CompletableFuture.runAsync(() -> { read(input); }); - writerThread = CompletableFuture.runAsync(() -> { write(output); }); - for (int i=0; i { process(j); })); - } - while (!writerThread.isDone()) - try { Thread.sleep(100); } catch (InterruptedException ex) {} - } - DataSource reader; HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; + ArrayList inputs = new ArrayList<>(); + ProgressPrintout progress = new ProgressPrintout(); + ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); + ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); + int threads; int maxEvents = 0; int maxEventsUser = 0; @@ -60,40 +44,60 @@ public void process(String output, String... input) { int readEvents = 0; int writeEvents = 0; - ArrayList inputs = new ArrayList<>(); - - ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); - ConcurrentLinkedQueue evioQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue hipoQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); - - ProgressPrintout progress = new ProgressPrintout(); + public EngineMultiProcessor(OptionParser parser) { + super(parser); + threads = parser.getOption("-t").intValue(); + maxEventsUser = parser.getOption("-n").intValue(); + skipEvents = parser.getOption("-s").intValue(); + } + + /** + * The thread launcher. + * @param output + * @param input + */ + public void process(String output, String... input) { + readerThread = CompletableFuture.runAsync(() -> { read(input); }); + writerThread = CompletableFuture.runAsync(() -> { write(output); }); + for (int i=0; i { process(j); })); + } + while (!writerThread.isDone()) { + for (CompletableFuture f : procThreads) + if (f.isDone()) procThreads.remove(f); + sleep(100); + } + } + /** + * The reader thread. + * @param input input filenames + */ void read(String... input) { inputs.addAll(Arrays.asList(input)); while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { - if (evioQueue.size()+hipoQueue.size() > 100*threads) { - try { Thread.sleep(100); } - catch (InterruptedException ex) {} - } + // sleep instead of overfilling the read queue: + if (readQueue.size() > 100*threads) sleep(100); + // read the next event: else { Benchmark.getInstance().resume("read"); readEvents++; + Object o = null; if (reader instanceof EvioSource evio) { - try { evioQueue.offer(evio.getEventBuffer(readEvents, true)); } + try { o = evio.getEventBuffer(readEvents, true); } catch (EvioException ex) { ex.printStackTrace(); } } - else { - DataEvent event = reader.getNextEvent(); - if (skipEvents < 1 || readEvents > skipEvents) - hipoQueue.offer(event); - } + else o = reader.getNextEvent(); + if (skipEvents < 1 || readEvents > skipEvents) + if (o != null) readQueue.offer(o); Benchmark.getInstance().pause("read"); } } else if (inputs.isEmpty()) break; else { + // open a new input file: if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); else reader = new EvioSource(); reader.open(inputs.remove(0)); @@ -101,6 +105,7 @@ void read(String... input) { if (reader instanceof HipoDataSource hipo) updateDictionary(hipo, writer); else { + // override maxEvents for EVIO: int n = ((EvioSource)reader).getEventCount(); maxEvents = maxEventsUser < n ? maxEventsUser : n; } @@ -108,86 +113,99 @@ void read(String... input) { } } } - - EvioDataEvent process(ByteBuffer bytes) { - Benchmark.getInstance().resume("EVIO"); - EvioDataEvent e = new EvioDataEvent(bytes.array(), ByteOrder.LITTLE_ENDIAN); - Benchmark.getInstance().pause("EVIO"); - return e; - } - - HipoDataEvent process(EvioDataEvent event) { - Benchmark.getInstance().resume("DECO"); - HipoDataEvent e; - try { - CLASDecoder d = decoders.take(); - e = d.getDecodedDataEvenet(event); - decoders.put(d); - } - catch (InterruptedException ex) { e = null; } - Benchmark.getInstance().pause("DECO"); - return e; - } - - void process(HipoDataEvent event) { - for (Map.Entry engine : processorEngines.entrySet()) { - Benchmark.getInstance().resume(engine.getValue().getName()); - try { engine.getValue().processDataEvent(event); } - catch (Exception ex) { ex.printStackTrace(); } - Benchmark.getInstance().pause(engine.getValue().getName()); - } - } - + + /** + * The event processor thread. + * @param thread unique thread number + */ void process(int thread) { while (true) { - if (evioQueue.isEmpty() && hipoQueue.isEmpty()) { - if (readerThread.isDone()) - if (evioQueue.isEmpty() && hipoQueue.isEmpty()) break; - try { Thread.sleep(100); } - catch (InterruptedException ex) {} + if (readQueue.isEmpty()) { + if (readerThread.isDone() && readQueue.isEmpty()) + break; + sleep(100); } else { + Object o = readQueue.poll(); DataEvent event; - if (!evioQueue.isEmpty()) event = process(process(evioQueue.poll())); - else if (!hipoQueue.isEmpty()) event = hipoQueue.poll(); - else continue; - process((HipoDataEvent)event); + // decode if necessary: + if (o instanceof ByteBuffer bb) event = decode(bb); + else event = (HipoDataEvent)o; + // run it through the engine chain: + for (Map.Entry engine : processorEngines.entrySet()) { + Benchmark.getInstance().resume(engine.getValue().getName()); + try { engine.getValue().processDataEvent(event); } + catch (Exception ex) { ex.printStackTrace(); } + Benchmark.getInstance().pause(engine.getValue().getName()); + } writeQueue.offer(event); } } } - + + /** + * The writer thread. + * @param output output filename + */ void write(String output) { writer = new HipoDataSync(); writer.setCompressionType(2); writer.open(output); while (true) { if (writeQueue.isEmpty()) { - for (CompletableFuture f : procThreads) - if (f.isDone()) procThreads.remove(f); - if (procThreads.isEmpty()) { - if (writeQueue.isEmpty()) { - writer.close(); - System.out.println(Benchmark.getInstance()); - System.out.println(String.format("recon-mutil::::: Read/Write/Diff = %d/%d/%d", - readEvents, writeEvents, readEvents-writeEvents)); - break; - } + if (procThreads.isEmpty() && writeQueue.isEmpty()) { + writer.close(); + System.out.println(Benchmark.getInstance()); + System.out.println(String.format("recon-mutil::::: Read/Write/Diff = %d/%d/%d", + readEvents, writeEvents, readEvents-writeEvents)); + break; } - try { Thread.sleep(100); } - catch (InterruptedException ex) {} + sleep(100); + } + else { + DataEvent e = writeQueue.poll(); + Benchmark.getInstance().resume("write"); + writer.writeEvent(e); + if (writeEvents > 100) progress.updateStatus(); + if (writeEvents == 101) Benchmark.getInstance().printTimer(10); + writeEvents++; + Benchmark.getInstance().pause("write"); } - else write(writeQueue.poll()); } } - - void write(DataEvent event) { - Benchmark.getInstance().resume("write"); - writer.writeEvent(writeQueue.poll()); - if (writeEvents > 100) progress.updateStatus(); - if (writeEvents == 101) Benchmark.getInstance().printTimer(10); - writeEvents++; - Benchmark.getInstance().pause("write"); + + /** + * Decoding. + * @param bytes EVIO byte buffer + * @return decoded event + */ + HipoDataEvent decode(ByteBuffer bytes) { + Benchmark.getInstance().resume("EVIO"); + EvioDataEvent evio = new EvioDataEvent(bytes.array(), ByteOrder.LITTLE_ENDIAN); + Benchmark.getInstance().pause("EVIO"); + Benchmark.getInstance().resume("DECO"); + HipoDataEvent hipo; + try { + CLASDecoder d = decoders.take(); + hipo = d.getDecodedDataEvenet(evio); + decoders.put(d); + } + catch (InterruptedException ex) { hipo = null; } + Benchmark.getInstance().pause("DECO"); + return hipo; } + void sleep(int milliseconds) { + try { Thread.sleep(milliseconds); } + catch (InterruptedException ex) {} + } + + public static void main(String[] args) { + OptionParser parser = EngineProcessor.getParser(); + parser.addOption("-t","4","number of threads"); + parser.parse(args); + EngineMultiProcessor proc = new EngineMultiProcessor(parser); + proc.process(parser.getOption("-o").stringValue(), parser.getOption("-i").stringValue()); + } + } diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java index f2bad4fdd9..af1a610fa2 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java @@ -254,7 +254,8 @@ public void addEngine(String name, String clazz, String jsonConf) { } this.processorEngines.put(name == null ? engine.getName() : name, engine); } else { - LOGGER.log(Level.SEVERE, ">>>> ERROR: class is not a reconstruction engine : {0}", clazz); + LOGGER.log( clazz.contains("DecoderEngine") ? Level.INFO : Level.SEVERE, + "Class is not a reconstruction engine : {0}", clazz); } } catch (ClassNotFoundException | InstantiationException | IllegalAccessException ex) { From 11298f18539ac2f77e6fb08bbfe26c68745be25e Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 16:06:27 -0400 Subject: [PATCH 05/62] optimize --- .../java/org/jlab/clas/reco/EngineMultiProcessor.java | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index acfeab53b0..d9e4c7e2c3 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -30,7 +30,6 @@ public class EngineMultiProcessor extends EngineProcessor { HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; - ArrayList inputs = new ArrayList<>(); ProgressPrintout progress = new ProgressPrintout(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); @@ -120,13 +119,13 @@ void read(String... input) { */ void process(int thread) { while (true) { - if (readQueue.isEmpty()) { + Object o = readQueue.poll(); + if (o == null) { if (readerThread.isDone() && readQueue.isEmpty()) break; sleep(100); } else { - Object o = readQueue.poll(); DataEvent event; // decode if necessary: if (o instanceof ByteBuffer bb) event = decode(bb); @@ -152,7 +151,8 @@ void write(String output) { writer.setCompressionType(2); writer.open(output); while (true) { - if (writeQueue.isEmpty()) { + DataEvent e = writeQueue.poll(); + if (e == null) { if (procThreads.isEmpty() && writeQueue.isEmpty()) { writer.close(); System.out.println(Benchmark.getInstance()); @@ -163,7 +163,6 @@ void write(String output) { sleep(100); } else { - DataEvent e = writeQueue.poll(); Benchmark.getInstance().resume("write"); writer.writeEvent(e); if (writeEvents > 100) progress.updateStatus(); @@ -207,5 +206,4 @@ public static void main(String[] args) { EngineMultiProcessor proc = new EngineMultiProcessor(parser); proc.process(parser.getOption("-o").stringValue(), parser.getOption("-i").stringValue()); } - } From efc2394a48165d8ce3ad282116a465e3ff7c8f43 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 17:51:17 -0400 Subject: [PATCH 06/62] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 54 ++++++++++--------- 1 file changed, 30 insertions(+), 24 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index d9e4c7e2c3..2a6621d2ab 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -82,34 +82,22 @@ void read(String... input) { // read the next event: else { Benchmark.getInstance().resume("read"); - readEvents++; Object o = null; if (reader instanceof EvioSource evio) { - try { o = evio.getEventBuffer(readEvents, true); } + try { o = evio.getEventBuffer(readEvents+1, true); } catch (EvioException ex) { ex.printStackTrace(); } } else o = reader.getNextEvent(); if (skipEvents < 1 || readEvents > skipEvents) - if (o != null) readQueue.offer(o); + if (o != null) { + readQueue.offer(o); + readEvents++; + } Benchmark.getInstance().pause("read"); } } else if (inputs.isEmpty()) break; - else { - // open a new input file: - if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); - else reader = new EvioSource(); - reader.open(inputs.remove(0)); - maxEvents = maxEventsUser; - if (reader instanceof HipoDataSource hipo) - updateDictionary(hipo, writer); - else { - // override maxEvents for EVIO: - int n = ((EvioSource)reader).getEventCount(); - maxEvents = maxEventsUser < n ? maxEventsUser : n; - } - readEvents = 0; - } + else open(); } } @@ -122,7 +110,7 @@ void process(int thread) { Object o = readQueue.poll(); if (o == null) { if (readerThread.isDone() && readQueue.isEmpty()) - break; + if (writeEvents >= readEvents) break; sleep(100); } else { @@ -154,10 +142,7 @@ void write(String output) { DataEvent e = writeQueue.poll(); if (e == null) { if (procThreads.isEmpty() && writeQueue.isEmpty()) { - writer.close(); - System.out.println(Benchmark.getInstance()); - System.out.println(String.format("recon-mutil::::: Read/Write/Diff = %d/%d/%d", - readEvents, writeEvents, readEvents-writeEvents)); + close(); break; } sleep(100); @@ -166,7 +151,6 @@ void write(String output) { Benchmark.getInstance().resume("write"); writer.writeEvent(e); if (writeEvents > 100) progress.updateStatus(); - if (writeEvents == 101) Benchmark.getInstance().printTimer(10); writeEvents++; Benchmark.getInstance().pause("write"); } @@ -194,6 +178,28 @@ HipoDataEvent decode(ByteBuffer bytes) { return hipo; } + void open() { + if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); + else reader = new EvioSource(); + reader.open(inputs.remove(0)); + maxEvents = maxEventsUser; + if (reader instanceof HipoDataSource hipo) { + updateDictionary(hipo, writer); + } else { + int n = ((EvioSource)reader).getEventCount(); + maxEvents = maxEventsUser < n ? maxEventsUser : n; + } + readEvents = 0; + writeEvents = 0; + } + + void close() { + writer.close(); + System.out.println(Benchmark.getInstance()); + System.out.println(String.format("recon-mutil::::: Read/Write/Diff = %d/%d/%d", + readEvents, writeEvents, readEvents-writeEvents)); + } + void sleep(int milliseconds) { try { Thread.sleep(milliseconds); } catch (InterruptedException ex) {} From 8c4e160a24d4f7409799e62c1519c7ccb79b50fd Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 18:04:59 -0400 Subject: [PATCH 07/62] cleanup --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 2a6621d2ab..ecfa8534e9 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -55,7 +55,7 @@ public EngineMultiProcessor(OptionParser parser) { * @param output * @param input */ - public void process(String output, String... input) { + public void launch(String output, String... input) { readerThread = CompletableFuture.runAsync(() -> { read(input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); for (int i=0; i skipEvents) + if (skipEvents < 1 || readEvents > skipEvents) { if (o != null) { readQueue.offer(o); readEvents++; } + } Benchmark.getInstance().pause("read"); } } @@ -210,6 +211,6 @@ public static void main(String[] args) { parser.addOption("-t","4","number of threads"); parser.parse(args); EngineMultiProcessor proc = new EngineMultiProcessor(parser); - proc.process(parser.getOption("-o").stringValue(), parser.getOption("-i").stringValue()); + proc.launch(parser.getOption("-o").stringValue(), parser.getOption("-i").stringValue()); } } From 71a0520e0562226e29c0e4d08b15daf9f432600d Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 20:57:28 -0400 Subject: [PATCH 08/62] enlarge read queue --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index ecfa8534e9..829e99ce4b 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -78,7 +78,7 @@ void read(String... input) { while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { // sleep instead of overfilling the read queue: - if (readQueue.size() > 100*threads) sleep(100); + if (readQueue.size() > 1000*threads) sleep(100); // read the next event: else { Benchmark.getInstance().resume("read"); From d8c84873deaf4b86db5b740f22445c587c5d9119 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 27 Aug 2026 21:13:09 -0400 Subject: [PATCH 09/62] cleanup --- .../java/org/jlab/io/clara/Clas12Writer.java | 7 +++++ .../jlab/clas/reco/EngineMultiProcessor.java | 31 ++++++++++++------- 2 files changed, 26 insertions(+), 12 deletions(-) diff --git a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java index a094b939cd..e64dacdfa5 100644 --- a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java +++ b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java @@ -33,6 +33,13 @@ */ public class Clas12Writer extends HipoToHipoWriter { + public static class Poster { + TreeMap eventUnix; + TreeSet helicities; + DaqScalersSequence scalers; + void process(Event e) {} + } + static final String[] TAG1BANKS = {"RUN::scaler","HEL::scaler","RAW::scaler","RAW::epics","HEL::flip","COAT::config"}; Bank[] tag1banks; diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 829e99ce4b..fccd6daa0d 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -2,8 +2,8 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; -import java.util.ArrayList; import java.util.Arrays; +import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; @@ -30,7 +30,6 @@ public class EngineMultiProcessor extends EngineProcessor { HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; - ArrayList inputs = new ArrayList<>(); ProgressPrintout progress = new ProgressPrintout(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); @@ -74,11 +73,17 @@ public void launch(String output, String... input) { * @param input input filenames */ void read(String... input) { - inputs.addAll(Arrays.asList(input)); + + // store the input filenames: + List inputs = Arrays.asList(input); + while (maxEvents < 1 || readEvents < maxEvents) { + if (reader != null && reader.hasEvent()) { + // sleep instead of overfilling the read queue: if (readQueue.size() > 1000*threads) sleep(100); + // read the next event: else { Benchmark.getInstance().resume("read"); @@ -88,17 +93,20 @@ void read(String... input) { catch (EvioException ex) { ex.printStackTrace(); } } else o = reader.getNextEvent(); - if (skipEvents < 1 || readEvents > skipEvents) { - if (o != null) { + if (o != null) { + readEvents++; + if (skipEvents < 1 || readEvents > skipEvents) readQueue.offer(o); - readEvents++; - } } Benchmark.getInstance().pause("read"); } } + + // we're done if there's no more input files: else if (inputs.isEmpty()) break; - else open(); + + // open the next input file: + else open(inputs.removeFirst()); } } @@ -179,10 +187,9 @@ HipoDataEvent decode(ByteBuffer bytes) { return hipo; } - void open() { - if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); - else reader = new EvioSource(); - reader.open(inputs.remove(0)); + void open(String filename) { + reader = filename.endsWith(".hipo") ? new HipoDataSource() : new EvioSource(); + reader.open(filename); maxEvents = maxEventsUser; if (reader instanceof HipoDataSource hipo) { updateDictionary(hipo, writer); From 124a0ee46af515bcc301bdadc6b6eed1a9671d71 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 28 Aug 2026 16:09:51 -0400 Subject: [PATCH 10/62] cleanup --- .../src/main/java/org/jlab/io/clara/Clas12Writer.java | 7 ------- 1 file changed, 7 deletions(-) diff --git a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java index e64dacdfa5..a094b939cd 100644 --- a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java +++ b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java @@ -33,13 +33,6 @@ */ public class Clas12Writer extends HipoToHipoWriter { - public static class Poster { - TreeMap eventUnix; - TreeSet helicities; - DaqScalersSequence scalers; - void process(Event e) {} - } - static final String[] TAG1BANKS = {"RUN::scaler","HEL::scaler","RAW::scaler","RAW::epics","HEL::flip","COAT::config"}; Bank[] tag1banks; From 06c9b0b2ca657eec413afa925e8338516207ba89 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 28 Aug 2026 17:52:50 -0400 Subject: [PATCH 11/62] try this --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4967588f6e..29da402154 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -72,7 +72,7 @@ jobs: - name: build run: | ./build-coatjava.sh --lfs --no-progress -T${{ env.nthreads }} - ./bin/install-clara -b -c ./coatjava ./clara + ./bin/install-clara -c ./coatjava ./clara - name: tar # tarball to preserve permissions run: | tar czvf coatjava.tar.gz coatjava From 90d0c54688db90f25c2be2057bd49bd66bfb0562 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 28 Aug 2026 18:05:15 -0400 Subject: [PATCH 12/62] and the other one ... --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 29da402154..3a39ffcb56 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -102,7 +102,7 @@ jobs: - name: build run: | ./build-coatjava.sh --lfs --no-progress -T${{ env.nthreads }} - ./bin/install-clara -b -c ./coatjava ./clara + ./bin/install-clara -c ./coatjava ./clara - name: tar # tarball to preserve permissions run: | tar czvf coatjava.tar.gz coatjava From 448cdac75a2b99a24d46df5ffae83f2bda6ae20c Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Sun, 30 Aug 2026 18:28:18 -0400 Subject: [PATCH 13/62] switch to recon-mutil in github ci test --- .github/workflows/ci.yml | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 3a39ffcb56..e472c075e9 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -224,20 +224,15 @@ jobs: clas_018779.evio.00001 - name: untar build run: | - tar xzvf clara.tar.gz tar xzvf coatjava.tar.gz - run: ls - name: run test - run: ./bin/run-clara -y ./etc/services/rgd-clarode.yml -t 4 -n 500 -c ./clara -o ./tmp ./clas_018779.evio.00001 - - name: ls tmp - run: ls -lhtr tmp - - name: rename - run: mv -v tmp/rec_clas_018779.evio.00001.hipo rec.hipo + run: ./bin/recon-mutil -t 4 -n 500 -y etc/services/rgd-clarode.yml -o rec_clas_018779.evio.00001.hipo -i clas_018779.evio.00001 - uses: actions/upload-artifact@v7 with: name: test_clara_result retention-days: 1 - path: rec.hipo + path: rec_clas_018779.evio.00001.hipo test_coatjava: needs: [ build ] From 4835468b53992b711ec71c86f41165c8780e55f8 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 11:32:02 -0400 Subject: [PATCH 14/62] fix --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index fccd6daa0d..40acb0f31b 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -2,6 +2,7 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; +import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.Map; @@ -75,7 +76,8 @@ public void launch(String output, String... input) { void read(String... input) { // store the input filenames: - List inputs = Arrays.asList(input); + List inputs = new ArrayList<>(); + inputs.addAll(Arrays.asList(input)); while (maxEvents < 1 || readEvents < maxEvents) { From 04d51ce2e4d37348f529e001c4e8d3e3a79a60a0 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 13:02:01 -0400 Subject: [PATCH 15/62] cleanup --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 40acb0f31b..e08de02173 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -76,8 +76,7 @@ public void launch(String output, String... input) { void read(String... input) { // store the input filenames: - List inputs = new ArrayList<>(); - inputs.addAll(Arrays.asList(input)); + List inputs = new ArrayList<>(Arrays.asList(input)); while (maxEvents < 1 || readEvents < maxEvents) { From 4b54db410225e565cf6b42f4eface572af72d2bc Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 14:31:53 -0400 Subject: [PATCH 16/62] switch to frames --- .../jlab/clas/reco/EngineMultiProcessor.java | 63 ++++++++++++------- 1 file changed, 41 insertions(+), 22 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index e08de02173..eb6524fba1 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -33,8 +33,8 @@ public class EngineMultiProcessor extends EngineProcessor { CompletableFuture writerThread; ProgressPrintout progress = new ProgressPrintout(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); - ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); int threads; int maxEvents = 0; @@ -75,15 +75,18 @@ public void launch(String output, String... input) { */ void read(String... input) { - // store the input filenames: + // convert input filenames to a list: List inputs = new ArrayList<>(Arrays.asList(input)); + // event buffer: + List frame = new ArrayList<>(100); + while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { // sleep instead of overfilling the read queue: - if (readQueue.size() > 1000*threads) sleep(100); + if (readQueue.size() > 100*threads) sleep(100); // read the next event: else { @@ -96,8 +99,13 @@ void read(String... input) { else o = reader.getNextEvent(); if (o != null) { readEvents++; - if (skipEvents < 1 || readEvents > skipEvents) - readQueue.offer(o); + if (skipEvents < 1 || readEvents > skipEvents) { + frame.add(o); + if (frame.size() >= 100) { + readQueue.offer(frame); + frame = new ArrayList<>(100); + } + } } Benchmark.getInstance().pause("read"); } @@ -109,6 +117,7 @@ void read(String... input) { // open the next input file: else open(inputs.removeFirst()); } + if (!frame.isEmpty()) readQueue.offer(frame); } /** @@ -116,28 +125,36 @@ void read(String... input) { * @param thread unique thread number */ void process(int thread) { + List frame = new ArrayList<>(100); while (true) { - Object o = readQueue.poll(); + List o = readQueue.poll(); if (o == null) { if (readerThread.isDone() && readQueue.isEmpty()) if (writeEvents >= readEvents) break; sleep(100); } else { - DataEvent event; - // decode if necessary: - if (o instanceof ByteBuffer bb) event = decode(bb); - else event = (HipoDataEvent)o; - // run it through the engine chain: - for (Map.Entry engine : processorEngines.entrySet()) { - Benchmark.getInstance().resume(engine.getValue().getName()); - try { engine.getValue().processDataEvent(event); } - catch (Exception ex) { ex.printStackTrace(); } - Benchmark.getInstance().pause(engine.getValue().getName()); + for (int i=0; i engine : processorEngines.entrySet()) { + Benchmark.getInstance().resume(engine.getValue().getName()); + try { engine.getValue().processDataEvent(event); } + catch (Exception ex) { ex.printStackTrace(); } + Benchmark.getInstance().pause(engine.getValue().getName()); + } + frame.add(event); + if (frame.size() >= 100) { + writeQueue.offer(frame); + frame = new ArrayList<>(100); + } } - writeQueue.offer(event); } } + if (!frame.isEmpty()) writeQueue.offer(frame); } /** @@ -149,7 +166,7 @@ void write(String output) { writer.setCompressionType(2); writer.open(output); while (true) { - DataEvent e = writeQueue.poll(); + List e = writeQueue.poll(); if (e == null) { if (procThreads.isEmpty() && writeQueue.isEmpty()) { close(); @@ -159,9 +176,11 @@ void write(String output) { } else { Benchmark.getInstance().resume("write"); - writer.writeEvent(e); - if (writeEvents > 100) progress.updateStatus(); - writeEvents++; + for (int i=0; i 100) progress.updateStatus(); + writeEvents++; + } Benchmark.getInstance().pause("write"); } } From 28164eb1b2c96394e83f4d4dcd442dbd9ef27f95 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 14:51:49 -0400 Subject: [PATCH 17/62] try this --- .../jlab/clas/reco/EngineMultiProcessor.java | 20 ++++++++++++++----- 1 file changed, 15 insertions(+), 5 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index eb6524fba1..44dce5ebf4 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -38,6 +38,7 @@ public class EngineMultiProcessor extends EngineProcessor { int threads; int maxEvents = 0; + int failEvents = 0; int maxEventsUser = 0; int skipEvents = 0; int readEvents = 0; @@ -93,12 +94,21 @@ void read(String... input) { Benchmark.getInstance().resume("read"); Object o = null; if (reader instanceof EvioSource evio) { - try { o = evio.getEventBuffer(readEvents+1, true); } - catch (EvioException ex) { ex.printStackTrace(); } + System.err.println("DOGGIES: "+readEvents); + System.err.println("DOGGIES: "+readEvents); + System.err.println("DOGGIES: "+readEvents); + System.err.println("DOGGIES: "+readEvents); + try { o = evio.getEventBuffer(++readEvents, true); } + catch (EvioException ex) { + failEvents++; + ex.printStackTrace(); + } } - else o = reader.getNextEvent(); - if (o != null) { + else { readEvents++; + o = reader.getNextEvent(); + } + if (o != null) { if (skipEvents < 1 || readEvents > skipEvents) { frame.add(o); if (frame.size() >= 100) { @@ -130,7 +140,7 @@ void process(int thread) { List o = readQueue.poll(); if (o == null) { if (readerThread.isDone() && readQueue.isEmpty()) - if (writeEvents >= readEvents) break; + if (writeEvents+skipEvents+failEvents >= readEvents) break; sleep(100); } else { From b5e865a80294d003b7ef0f0da6a1771a05abe3c6 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 14:53:51 -0400 Subject: [PATCH 18/62] try this --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 44dce5ebf4..5fbbd5dbf4 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -94,10 +94,7 @@ void read(String... input) { Benchmark.getInstance().resume("read"); Object o = null; if (reader instanceof EvioSource evio) { - System.err.println("DOGGIES: "+readEvents); - System.err.println("DOGGIES: "+readEvents); - System.err.println("DOGGIES: "+readEvents); - System.err.println("DOGGIES: "+readEvents); + System.err.println("DOGGIES: "+readEvents+"/"+maxEvents); try { o = evio.getEventBuffer(++readEvents, true); } catch (EvioException ex) { failEvents++; From bc8f9e08729ebc1bfda8211bae344492b4f976ff Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 14:58:21 -0400 Subject: [PATCH 19/62] try this --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 5fbbd5dbf4..438c07cbe2 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -223,6 +223,9 @@ void open(String filename) { } else { int n = ((EvioSource)reader).getEventCount(); maxEvents = maxEventsUser < n ? maxEventsUser : n; + System.err.println(maxEventsUser); + System.err.println(maxEvents); + System.exit(1); } readEvents = 0; writeEvents = 0; From 0846c7031439106aeb6c4fd566617a0105f96079 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 15:19:33 -0400 Subject: [PATCH 20/62] protect --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 438c07cbe2..ee2526c098 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -222,10 +222,7 @@ void open(String filename) { updateDictionary(hipo, writer); } else { int n = ((EvioSource)reader).getEventCount(); - maxEvents = maxEventsUser < n ? maxEventsUser : n; - System.err.println(maxEventsUser); - System.err.println(maxEvents); - System.exit(1); + maxEvents = maxEventsUser>0 && maxEventsUser < n ? maxEventsUser : n; } readEvents = 0; writeEvents = 0; From bb64864778c0d14eec5c651b9f07e45ff1df6228 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 15:20:14 -0400 Subject: [PATCH 21/62] cleanup --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 1 - 1 file changed, 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index ee2526c098..0fcf764fc1 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -94,7 +94,6 @@ void read(String... input) { Benchmark.getInstance().resume("read"); Object o = null; if (reader instanceof EvioSource evio) { - System.err.println("DOGGIES: "+readEvents+"/"+maxEvents); try { o = evio.getEventBuffer(++readEvents, true); } catch (EvioException ex) { failEvents++; From 5a52ade048216b6dc3a2ae9d1af3e7620f6e5b92 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 15:27:46 -0400 Subject: [PATCH 22/62] release memory --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 1 + 1 file changed, 1 insertion(+) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 0fcf764fc1..d21cd6c78e 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -124,6 +124,7 @@ void read(String... input) { else open(inputs.removeFirst()); } if (!frame.isEmpty()) readQueue.offer(frame); + reader.close(); } /** From 535ce925ee496a26d8600cc4dbb1a6aa6454e1ff Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 15:43:04 -0400 Subject: [PATCH 23/62] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 20 +++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index d21cd6c78e..d65b4d60dc 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -26,23 +26,26 @@ * @author baltzell */ public class EngineMultiProcessor extends EngineProcessor { + + static final int FRAME_SIZE = 100; DataSource reader; HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; - ProgressPrintout progress = new ProgressPrintout(); + ProgressPrintout progress; ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); int threads; - int maxEvents = 0; - int failEvents = 0; + int failEvents; + int readEvents; + int writeEvents; + int maxEvents; + int maxEventsUser = 0; int skipEvents = 0; - int readEvents = 0; - int writeEvents = 0; public EngineMultiProcessor(OptionParser parser) { super(parser); @@ -80,7 +83,7 @@ void read(String... input) { List inputs = new ArrayList<>(Arrays.asList(input)); // event buffer: - List frame = new ArrayList<>(100); + List frame = new ArrayList<>(FRAME_SIZE); while (maxEvents < 1 || readEvents < maxEvents) { @@ -107,9 +110,9 @@ void read(String... input) { if (o != null) { if (skipEvents < 1 || readEvents > skipEvents) { frame.add(o); - if (frame.size() >= 100) { + if (frame.size() >= FRAME_SIZE) { readQueue.offer(frame); - frame = new ArrayList<>(100); + frame = new ArrayList<>(FRAME_SIZE); } } } @@ -226,6 +229,7 @@ void open(String filename) { } readEvents = 0; writeEvents = 0; + failEvents = 0; } void close() { From ac9d837959d908e885e96aa4196b5cc877849de1 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 15:50:32 -0400 Subject: [PATCH 24/62] fix --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index d65b4d60dc..4fb2eb57e3 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -33,7 +33,7 @@ public class EngineMultiProcessor extends EngineProcessor { HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; - ProgressPrintout progress; + ProgressPrintout progress = new ProgressPrintout(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); From fff88ea06c8ca7c970ff09a90d09e3bd648cf62e Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 15:52:07 -0400 Subject: [PATCH 25/62] fix --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 4fb2eb57e3..0d3fe412a4 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -33,7 +33,7 @@ public class EngineMultiProcessor extends EngineProcessor { HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; - ProgressPrintout progress = new ProgressPrintout(); + ProgressPrintout progress; ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); @@ -188,7 +188,10 @@ void write(String output) { Benchmark.getInstance().resume("write"); for (int i=0; i 100) progress.updateStatus(); + if (writeEvents > 100) { + if (progress == null) progress = new ProgressPrintout(); + progress.updateStatus(); + } writeEvents++; } Benchmark.getInstance().pause("write"); From dd4d5bfe3cae0572edbf710dc0c4e6df224583f7 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 16:46:03 -0400 Subject: [PATCH 26/62] switch scaling test to recon-mutil --- libexec/scaling | 21 ++++++--------------- 1 file changed, 6 insertions(+), 15 deletions(-) diff --git a/libexec/scaling b/libexec/scaling index 6cad8e3450..a85d99c29b 100755 --- a/libexec/scaling +++ b/libexec/scaling @@ -4,7 +4,6 @@ def cli(): import os,shutil,argparse cli = argparse.ArgumentParser(description='CLARA scaling test') cli.add_argument('-y','--yaml', metavar='YAML',help='path to YAML file',required=True) - cli.add_argument('-c','--clara', metavar='DIR',help='CLARA_HOME path (default=$CLARA_HOME)',default=os.getenv('CLARA_HOME',None)) cli.add_argument('-t','--threads',metavar='#',help='threads (default=8,16,24,34,44)',default='8,16,24,34,44') cli.add_argument('-e','--events', metavar='#',help='events per thread (default=555)',default=555,type=int) cli.add_argument('-N','--numa', metavar='#',help='restrict to a NUMA socket (choices=[0,1])',choices=[0,1]) @@ -13,14 +12,8 @@ def cli(): cli.add_argument('datafile', help='input EVIO/HIPO data file') cfg = cli.parse_args() cfg.threads = cfg.threads.split(',') - if not cfg.clara: - cli.error('cannot find CLARA installation via -c or $CLARA_HOME') - elif shutil.which('run-clara'): - cfg.run_clara = shutil.which('run-clara') - elif os.path.exists(clara_home + '/plugins/clas12/bin/run-clara'): - cfg.run_clara = clara_home + '/plugins/clas12/bin/run-clara' - else: - cli.error('cannot find run-clara in $PATH or $CLARA_HOME') + if not shutil.which('recon-mutil'): + cli.error('cannot find recon-mutil in $PATH.') if cfg.cpus: try: cfg.cpus = ','.join([str(int(i)) for i in cfg.cpus.split(',')]) @@ -58,15 +51,13 @@ def benchmark(cfg, threads, log): cmd.extend(['taskset','-c',','.join(cfg.cpus.split(',')[:threads])]) else: cmd.extend(['taskset','-c',cfg.cpus]) - # add the run-clara command: - cmd.extend([cfg.run_clara, - '-c',cfg.clara, + # add the recon-mutil command: + cmd.extend(['recon-mutil', '-n',str(cfg.events*int(threads)), '-t',str(threads), - '-l', '-y',cfg.yaml, - '-o',f'tmp-scaling-{threads}', - cfg.datafile]) + '-o',f'tmp-scaling-{threads}.hipo', + '-i',cfg.datafile]) exiting,benchmarks = False,collections.OrderedDict() for line in run(cmd): cols = line.split() From 056cdc59fe24b6da3451b395cd1637a18b11cff9 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 16:47:43 -0400 Subject: [PATCH 27/62] no need for negatives for #events and #skip --- .../src/main/java/org/jlab/clas/reco/EngineProcessor.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java index af1a610fa2..58ee6b17e8 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java @@ -399,8 +399,8 @@ protected static OptionParser getParser() { parser.addRequired("-o","output.hipo"); parser.addRequired("-i","input.evio/hipo"); parser.addOption("-c","0","use default configuration [0 - no, 1 - yes/default, 2 - all services] "); - parser.addOption("-s","-1","number of events to skip"); - parser.addOption("-n","-1","number of events to process"); + parser.addOption("-s","0","number of events to skip"); + parser.addOption("-n","0","number of events to process"); parser.addOption("-y","0","yaml file"); parser.addOption("-u","true","update dictionary from writer ? "); parser.addOption("-S",null,"schema directory"); From e39db2922a50b36a401b33e697a06812867d4caa Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 17:04:14 -0400 Subject: [PATCH 28/62] cleanup --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 0d3fe412a4..e2ec334f84 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -33,7 +33,7 @@ public class EngineMultiProcessor extends EngineProcessor { HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; - ProgressPrintout progress; + ProgressPrintout progress = new ProgressPrintout(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); @@ -188,10 +188,7 @@ void write(String output) { Benchmark.getInstance().resume("write"); for (int i=0; i 100) { - if (progress == null) progress = new ProgressPrintout(); - progress.updateStatus(); - } + progress.updateStatus(); writeEvents++; } Benchmark.getInstance().pause("write"); From 2405005db89066dc947db0325030763cbc535220 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 17:08:51 -0400 Subject: [PATCH 29/62] cleanup --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index e2ec334f84..dbf4fb4d2f 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -26,7 +26,8 @@ * @author baltzell */ public class EngineMultiProcessor extends EngineProcessor { - + + static final int READ_QUEUE_SIZE = 100; static final int FRAME_SIZE = 100; DataSource reader; @@ -90,7 +91,7 @@ void read(String... input) { if (reader != null && reader.hasEvent()) { // sleep instead of overfilling the read queue: - if (readQueue.size() > 100*threads) sleep(100); + if (readQueue.size() > READ_QUEUE_SIZE*threads) sleep(100); // read the next event: else { From c17674370d3c91148e43c506523fd2e64e8d9dcb Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 17:23:16 -0400 Subject: [PATCH 30/62] debug --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index dbf4fb4d2f..11982a7224 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -140,8 +140,10 @@ void process(int thread) { while (true) { List o = readQueue.poll(); if (o == null) { - if (readerThread.isDone() && readQueue.isEmpty()) + if (readerThread.isDone() && readQueue.isEmpty()) { + System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+" "+readEvents); if (writeEvents+skipEvents+failEvents >= readEvents) break; + } sleep(100); } else { From a08374e4e777501c8250000dd6962db41e5905e3 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 17:49:36 -0400 Subject: [PATCH 31/62] fixup --- .../jlab/clas/reco/EngineMultiProcessor.java | 29 ++++++++++--------- 1 file changed, 15 insertions(+), 14 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 11982a7224..d93d629fcd 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -44,6 +44,7 @@ public class EngineMultiProcessor extends EngineProcessor { int readEvents; int writeEvents; int maxEvents; + int evioEvents; int maxEventsUser = 0; int skipEvents = 0; @@ -98,23 +99,19 @@ void read(String... input) { Benchmark.getInstance().resume("read"); Object o = null; if (reader instanceof EvioSource evio) { - try { o = evio.getEventBuffer(++readEvents, true); } + try { o = evio.getEventBuffer(++evioEvents, true); } catch (EvioException ex) { failEvents++; ex.printStackTrace(); } } - else { - readEvents++; - o = reader.getNextEvent(); - } - if (o != null) { - if (skipEvents < 1 || readEvents > skipEvents) { - frame.add(o); - if (frame.size() >= FRAME_SIZE) { - readQueue.offer(frame); - frame = new ArrayList<>(FRAME_SIZE); - } + else o = reader.getNextEvent(); + if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { + frame.add(o); + if (frame.size() >= FRAME_SIZE) { + readQueue.offer(frame); + readEvents += frame.size(); + frame = new ArrayList<>(FRAME_SIZE); } } Benchmark.getInstance().pause("read"); @@ -127,7 +124,10 @@ void read(String... input) { // open the next input file: else open(inputs.removeFirst()); } - if (!frame.isEmpty()) readQueue.offer(frame); + if (!frame.isEmpty()) { + readQueue.offer(frame); + readEvents += frame.size(); + } reader.close(); } @@ -192,8 +192,8 @@ void write(String output) { for (int i=0; i Date: Mon, 31 Aug 2026 17:58:25 -0400 Subject: [PATCH 32/62] fix --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index d93d629fcd..a2945e4ac4 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -87,7 +87,7 @@ void read(String... input) { // event buffer: List frame = new ArrayList<>(FRAME_SIZE); - while (maxEvents < 1 || readEvents < maxEvents) { + while (maxEvents < 1 || Math.max(readEvents,evioEvents) < maxEvents) { if (reader != null && reader.hasEvent()) { From 65b61d506d7fb36238b15355819842791b56c265 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 18:25:34 -0400 Subject: [PATCH 33/62] try this --- .../jlab/clas/reco/EngineMultiProcessor.java | 86 ++++++++++++------- 1 file changed, 54 insertions(+), 32 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index a2945e4ac4..2488f8cff0 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -26,10 +26,10 @@ * @author baltzell */ public class EngineMultiProcessor extends EngineProcessor { - - static final int READ_QUEUE_SIZE = 100; - static final int FRAME_SIZE = 100; - + + final int MAX_READ_QUEUE = 100; + final int FRAME_SIZE = 100; + DataSource reader; HipoDataSync writer; CompletableFuture readerThread; @@ -62,6 +62,9 @@ public EngineMultiProcessor(OptionParser parser) { * @param input */ public void launch(String output, String... input) { + readEvents = 0; + writeEvents = 0; + failEvents = 0; readerThread = CompletableFuture.runAsync(() -> { read(input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); for (int i=0; i READ_QUEUE_SIZE*threads) sleep(100); + if (readQueue.size() > MAX_READ_QUEUE*threads) sleep(100); // read the next event: - else { - Benchmark.getInstance().resume("read"); - Object o = null; - if (reader instanceof EvioSource evio) { - try { o = evio.getEventBuffer(++evioEvents, true); } - catch (EvioException ex) { - failEvents++; - ex.printStackTrace(); - } - } - else o = reader.getNextEvent(); - if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { - frame.add(o); - if (frame.size() >= FRAME_SIZE) { - readQueue.offer(frame); - readEvents += frame.size(); - frame = new ArrayList<>(FRAME_SIZE); - } - } - Benchmark.getInstance().pause("read"); - } + else frame = read(frame); } // we're done if there's no more input files: @@ -124,7 +111,10 @@ void read(String... input) { // open the next input file: else open(inputs.removeFirst()); } + + // leftover, partial frame: if (!frame.isEmpty()) { + System.err.println("writing partial frame: "+frame.size()); readQueue.offer(frame); readEvents += frame.size(); } @@ -141,7 +131,6 @@ void process(int thread) { List o = readQueue.poll(); if (o == null) { if (readerThread.isDone() && readQueue.isEmpty()) { - System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+" "+readEvents); if (writeEvents+skipEvents+failEvents >= readEvents) break; } sleep(100); @@ -219,8 +208,42 @@ HipoDataEvent decode(ByteBuffer bytes) { Benchmark.getInstance().pause("DECO"); return hipo; } - + + /** + * Read the next event, add it to the frame, and, if the frame is full, + * add the frame to the read queue and return a new, empty frame. + * @frame the frame to fill + * @return the modified frame, or a new one if the frame was full + */ + List read(List frame) { + Benchmark.getInstance().resume("read"); + Object o = null; + if (reader instanceof EvioSource evio) { + try { o = evio.getEventBuffer(++evioEvents, true); } + catch (EvioException ex) { + failEvents++; + ex.printStackTrace(); + } + } + else o = reader.getNextEvent(); + if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { + frame.add(o); + if (frame.size() >= FRAME_SIZE) { + readQueue.offer(frame); + readEvents += frame.size(); + frame = new ArrayList<>(FRAME_SIZE); + } + } + Benchmark.getInstance().pause("read"); + return frame; + } + + /** + * Open the input file and do some initializations. + * @param filename + */ void open(String filename) { + evioEvents = 0; reader = filename.endsWith(".hipo") ? new HipoDataSource() : new EvioSource(); reader.open(filename); maxEvents = maxEventsUser; @@ -230,12 +253,11 @@ void open(String filename) { int n = ((EvioSource)reader).getEventCount(); maxEvents = maxEventsUser>0 && maxEventsUser < n ? maxEventsUser : n; } - readEvents = 0; - writeEvents = 0; - failEvents = 0; - evioEvents = 0; } + /** + * Close the output file and print performance info. + */ void close() { writer.close(); System.out.println(Benchmark.getInstance()); From f93f5327e0bb2c6975f6852ff7ee40e1080fd198 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 18:28:15 -0400 Subject: [PATCH 34/62] try this --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 2488f8cff0..6aeb970873 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -76,8 +76,8 @@ public void launch(String output, String... input) { if (f.isDone()) procThreads.remove(f); sleep(100); if (readerThread.isDone()) { - System.err.println(readQueue.size()+"/"+writeQueue.size()+"/"+procThreads.size()+"/"+writerThread.isDone()); - System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+" "+readEvents); + System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()+","+writerThread.isDone()); + System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+"/"+readEvents); } } } From 48cc2dedf1eea5eabc3f9a60561fbbb234bf79cc Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 18:38:28 -0400 Subject: [PATCH 35/62] fix --- .../jlab/clas/reco/EngineMultiProcessor.java | 26 +++++++++---------- 1 file changed, 13 insertions(+), 13 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 6aeb970873..98a2cb1d28 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -44,7 +44,7 @@ public class EngineMultiProcessor extends EngineProcessor { int readEvents; int writeEvents; int maxEvents; - int evioEvents; + int fileEvents; int maxEventsUser = 0; int skipEvents = 0; @@ -74,7 +74,7 @@ public void launch(String output, String... input) { while (!writerThread.isDone()) { for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); - sleep(100); + sleep(1000); if (readerThread.isDone()) { System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()+","+writerThread.isDone()); System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+"/"+readEvents); @@ -94,7 +94,9 @@ void read(String... input) { // event buffer: List frame = new ArrayList<>(FRAME_SIZE); - while (maxEvents < 1 || Math.max(readEvents,evioEvents) < maxEvents) { + while (maxEventsUser < 1 || readEvents < maxEventsUser) { + + if (maxEvents > 0 && fileEvents < maxEvents) break; if (reader != null && reader.hasEvent()) { @@ -105,18 +107,17 @@ void read(String... input) { else frame = read(frame); } - // we're done if there's no more input files: - else if (inputs.isEmpty()) break; - // open the next input file: - else open(inputs.removeFirst()); + else if (!inputs.isEmpty()) open(inputs.removeFirst()); + + else break; } // leftover, partial frame: if (!frame.isEmpty()) { System.err.println("writing partial frame: "+frame.size()); - readQueue.offer(frame); readEvents += frame.size(); + readQueue.offer(frame); } reader.close(); } @@ -219,7 +220,7 @@ List read(List frame) { Benchmark.getInstance().resume("read"); Object o = null; if (reader instanceof EvioSource evio) { - try { o = evio.getEventBuffer(++evioEvents, true); } + try { o = evio.getEventBuffer(++fileEvents, true); } catch (EvioException ex) { failEvents++; ex.printStackTrace(); @@ -243,15 +244,14 @@ List read(List frame) { * @param filename */ void open(String filename) { - evioEvents = 0; + fileEvents = 0; reader = filename.endsWith(".hipo") ? new HipoDataSource() : new EvioSource(); reader.open(filename); - maxEvents = maxEventsUser; if (reader instanceof HipoDataSource hipo) { + maxEvents = 0; updateDictionary(hipo, writer); } else { - int n = ((EvioSource)reader).getEventCount(); - maxEvents = maxEventsUser>0 && maxEventsUser < n ? maxEventsUser : n; + maxEvents = ((EvioSource)reader).getEventCount(); } } From 4474d7cb0822e270aa3a938c29efd81132970d04 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 18:40:02 -0400 Subject: [PATCH 36/62] fix --- .../java/org/jlab/clas/reco/EngineMultiProcessor.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 98a2cb1d28..a877295656 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -43,8 +43,8 @@ public class EngineMultiProcessor extends EngineProcessor { int failEvents; int readEvents; int writeEvents; - int maxEvents; int fileEvents; + int maxFileEvents; int maxEventsUser = 0; int skipEvents = 0; @@ -96,7 +96,7 @@ void read(String... input) { while (maxEventsUser < 1 || readEvents < maxEventsUser) { - if (maxEvents > 0 && fileEvents < maxEvents) break; + if (maxFileEvents > 0 && fileEvents > maxFileEvents) break; if (reader != null && reader.hasEvent()) { @@ -248,10 +248,10 @@ void open(String filename) { reader = filename.endsWith(".hipo") ? new HipoDataSource() : new EvioSource(); reader.open(filename); if (reader instanceof HipoDataSource hipo) { - maxEvents = 0; + maxFileEvents = 0; updateDictionary(hipo, writer); } else { - maxEvents = ((EvioSource)reader).getEventCount(); + maxFileEvents = ((EvioSource)reader).getEventCount(); } } From 86e98a0903f00509a3b4393a1622b3e283b5c007 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 18:43:19 -0400 Subject: [PATCH 37/62] fix --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index a877295656..851986dc63 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -76,7 +76,8 @@ public void launch(String output, String... input) { if (f.isDone()) procThreads.remove(f); sleep(1000); if (readerThread.isDone()) { - System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()+","+writerThread.isDone()); + System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()); + System.err.println(readerThread.isDone()+"|"+writerThread.isDone()); System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+"/"+readEvents); } } @@ -103,7 +104,7 @@ void read(String... input) { // sleep instead of overfilling the read queue: if (readQueue.size() > MAX_READ_QUEUE*threads) sleep(100); - // read the next event: + // read the next event into the frame: else frame = read(frame); } From 24a3342647d3a1c51ad1c8641c05d78bf9cd5118 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 18:44:50 -0400 Subject: [PATCH 38/62] fix --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 851986dc63..890908a610 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -45,7 +45,6 @@ public class EngineMultiProcessor extends EngineProcessor { int writeEvents; int fileEvents; int maxFileEvents; - int maxEventsUser = 0; int skipEvents = 0; @@ -97,7 +96,7 @@ void read(String... input) { while (maxEventsUser < 1 || readEvents < maxEventsUser) { - if (maxFileEvents > 0 && fileEvents > maxFileEvents) break; + if (maxFileEvents > 0 && fileEvents >= maxFileEvents) break; if (reader != null && reader.hasEvent()) { From 59075e8b62d895ad4578b68c0ca100cfe0971bbf Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 18:51:53 -0400 Subject: [PATCH 39/62] cleanup --- .../org/jlab/clas/reco/EngineMultiProcessor.java | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 890908a610..82d279d75a 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -40,18 +40,19 @@ public class EngineMultiProcessor extends EngineProcessor { ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); int threads; - int failEvents; + int maxEvents; + int skipEvents; + int readEvents; int writeEvents; + int failEvents; int fileEvents; int maxFileEvents; - int maxEventsUser = 0; - int skipEvents = 0; public EngineMultiProcessor(OptionParser parser) { super(parser); threads = parser.getOption("-t").intValue(); - maxEventsUser = parser.getOption("-n").intValue(); + maxEvents = parser.getOption("-n").intValue(); skipEvents = parser.getOption("-s").intValue(); } @@ -75,9 +76,10 @@ public void launch(String output, String... input) { if (f.isDone()) procThreads.remove(f); sleep(1000); if (readerThread.isDone()) { + System.err.println("------------------------------------------"); System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()); - System.err.println(readerThread.isDone()+"|"+writerThread.isDone()); System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+"/"+readEvents); + System.err.println(readerThread.isDone()+"|"+writerThread.isDone()); } } } @@ -94,7 +96,7 @@ void read(String... input) { // event buffer: List frame = new ArrayList<>(FRAME_SIZE); - while (maxEventsUser < 1 || readEvents < maxEventsUser) { + while (maxEvents < 1 || readEvents < maxEvents) { if (maxFileEvents > 0 && fileEvents >= maxFileEvents) break; From b5c130b6384c80162b7f27bf1439f6a7133865c1 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 19:02:19 -0400 Subject: [PATCH 40/62] fix --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 82d279d75a..d42d5505f3 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -77,9 +77,8 @@ public void launch(String output, String... input) { sleep(1000); if (readerThread.isDone()) { System.err.println("------------------------------------------"); - System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()); System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+"/"+readEvents); - System.err.println(readerThread.isDone()+"|"+writerThread.isDone()); + System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()); } } } @@ -133,7 +132,8 @@ void process(int thread) { while (true) { List o = readQueue.poll(); if (o == null) { - if (readerThread.isDone() && readQueue.isEmpty()) { + if (readerThread.isDone()) { + if (readQueue.isEmpty() && !frame.isEmpty()) writeQueue.offer(frame); if (writeEvents+skipEvents+failEvents >= readEvents) break; } sleep(100); From 63704cc6f9d8dcded22c19f5d2afe18477ec592e Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 19:04:05 -0400 Subject: [PATCH 41/62] fix --- .../java/org/jlab/clas/reco/EngineMultiProcessor.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index d42d5505f3..c9135533b6 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -133,7 +133,10 @@ void process(int thread) { List o = readQueue.poll(); if (o == null) { if (readerThread.isDone()) { - if (readQueue.isEmpty() && !frame.isEmpty()) writeQueue.offer(frame); + if (readQueue.isEmpty() && !frame.isEmpty()) { + writeQueue.offer(frame); + frame = new ArrayList<>(FRAME_SIZE); + } if (writeEvents+skipEvents+failEvents >= readEvents) break; } sleep(100); @@ -152,14 +155,13 @@ void process(int thread) { Benchmark.getInstance().pause(engine.getValue().getName()); } frame.add(event); - if (frame.size() >= 100) { + if (frame.size() >= FRAME_SIZE) { writeQueue.offer(frame); - frame = new ArrayList<>(100); + frame = new ArrayList<>(FRAME_SIZE); } } } } - if (!frame.isEmpty()) writeQueue.offer(frame); } /** From 3692d65d639f29cbd09d04378a162d7c49d63878 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 22:46:38 -0400 Subject: [PATCH 42/62] fix --- .../jlab/clas/reco/EngineMultiProcessor.java | 130 ++++++++---------- 1 file changed, 56 insertions(+), 74 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index c9135533b6..511d54b8d8 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -27,18 +27,21 @@ */ public class EngineMultiProcessor extends EngineProcessor { - final int MAX_READ_QUEUE = 100; + final int QUEUE_SIZE = 100; final int FRAME_SIZE = 100; DataSource reader; HipoDataSync writer; + CompletableFuture readerThread; CompletableFuture writerThread; - ProgressPrintout progress = new ProgressPrintout(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); + ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); + ProgressPrintout progress = new ProgressPrintout(); + int threads; int maxEvents; int skipEvents; @@ -55,11 +58,18 @@ public EngineMultiProcessor(OptionParser parser) { maxEvents = parser.getOption("-n").intValue(); skipEvents = parser.getOption("-s").intValue(); } + + public void scaling(String input, int... threads) { + for (int t : threads) { + this.threads = t; + launch(String.format("scaling-%d.hipo",t),input); + } + } /** - * The thread launcher. - * @param output - * @param input + * The thread launcher and collector. + * @param output name of output file to write + * @param input names of input files to read */ public void launch(String output, String... input) { readEvents = 0; @@ -77,8 +87,8 @@ public void launch(String output, String... input) { sleep(1000); if (readerThread.isDone()) { System.err.println("------------------------------------------"); - System.err.println(writeEvents+"/"+skipEvents+"/"+failEvents+"/"+readEvents); - System.err.println(readQueue.size()+","+writeQueue.size()+","+procThreads.size()); + System.err.println(writeEvents+"+"+skipEvents+"+"+failEvents+"=?"+readEvents); + System.err.println(readQueue.size()+" -> "+writeQueue.size()); } } } @@ -92,29 +102,50 @@ void read(String... input) { // convert input filenames to a list: List inputs = new ArrayList<>(Arrays.asList(input)); - // event buffer: + // initialize the event frame: List frame = new ArrayList<>(FRAME_SIZE); - while (maxEvents < 1 || readEvents < maxEvents) { - - if (maxFileEvents > 0 && fileEvents >= maxFileEvents) break; + // loop over input events: + while ( (maxEvents < 1 || readEvents < maxEvents) && + (maxFileEvents < 1 || fileEvents < maxFileEvents) ) { if (reader != null && reader.hasEvent()) { // sleep instead of overfilling the read queue: - if (readQueue.size() > MAX_READ_QUEUE*threads) sleep(100); + if (readQueue.size() > QUEUE_SIZE*threads) sleep(100); // read the next event into the frame: - else frame = read(frame); + else { + Benchmark.getInstance().resume("read"); + Object o = null; + if (reader instanceof EvioSource evio) { + try { o = evio.getEventBuffer(++fileEvents, true); } + catch (EvioException ex) { + failEvents++; + ex.printStackTrace(); + } + } + else o = reader.getNextEvent(); + if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { + frame.add(o); + if (frame.size() >= FRAME_SIZE) { + readQueue.offer(frame); + readEvents += frame.size(); + frame = new ArrayList<>(FRAME_SIZE); + } + } + Benchmark.getInstance().pause("read"); + } } // open the next input file: else if (!inputs.isEmpty()) open(inputs.removeFirst()); - + + // nothing left to do: else break; } - // leftover, partial frame: + // write leftover, partial frame: if (!frame.isEmpty()) { System.err.println("writing partial frame: "+frame.size()); readEvents += frame.size(); @@ -128,20 +159,15 @@ void read(String... input) { * @param thread unique thread number */ void process(int thread) { - List frame = new ArrayList<>(100); while (true) { List o = readQueue.poll(); if (o == null) { - if (readerThread.isDone()) { - if (readQueue.isEmpty() && !frame.isEmpty()) { - writeQueue.offer(frame); - frame = new ArrayList<>(FRAME_SIZE); - } - if (writeEvents+skipEvents+failEvents >= readEvents) break; - } + if (readerThread.isDone() && readQueue.isEmpty() && + writeEvents+skipEvents+failEvents >= readEvents) break; sleep(100); } else { + List frame = new ArrayList<>(o.size()); for (int i=0; i= FRAME_SIZE) { - writeQueue.offer(frame); - frame = new ArrayList<>(FRAME_SIZE); - } } + writeQueue.offer(frame); } } } - + /** * The writer thread. * @param output output filename @@ -179,7 +202,7 @@ void write(String output) { close(); break; } - sleep(100); + sleep(1000); } else { Benchmark.getInstance().resume("write"); @@ -193,16 +216,11 @@ void write(String output) { } } - /** - * Decoding. - * @param bytes EVIO byte buffer - * @return decoded event - */ HipoDataEvent decode(ByteBuffer bytes) { - Benchmark.getInstance().resume("EVIO"); + Benchmark.getInstance().resume("evio"); EvioDataEvent evio = new EvioDataEvent(bytes.array(), ByteOrder.LITTLE_ENDIAN); - Benchmark.getInstance().pause("EVIO"); - Benchmark.getInstance().resume("DECO"); + Benchmark.getInstance().pause("evio"); + Benchmark.getInstance().resume("deco"); HipoDataEvent hipo; try { CLASDecoder d = decoders.take(); @@ -210,43 +228,10 @@ HipoDataEvent decode(ByteBuffer bytes) { decoders.put(d); } catch (InterruptedException ex) { hipo = null; } - Benchmark.getInstance().pause("DECO"); + Benchmark.getInstance().pause("deco"); return hipo; } - /** - * Read the next event, add it to the frame, and, if the frame is full, - * add the frame to the read queue and return a new, empty frame. - * @frame the frame to fill - * @return the modified frame, or a new one if the frame was full - */ - List read(List frame) { - Benchmark.getInstance().resume("read"); - Object o = null; - if (reader instanceof EvioSource evio) { - try { o = evio.getEventBuffer(++fileEvents, true); } - catch (EvioException ex) { - failEvents++; - ex.printStackTrace(); - } - } - else o = reader.getNextEvent(); - if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { - frame.add(o); - if (frame.size() >= FRAME_SIZE) { - readQueue.offer(frame); - readEvents += frame.size(); - frame = new ArrayList<>(FRAME_SIZE); - } - } - Benchmark.getInstance().pause("read"); - return frame; - } - - /** - * Open the input file and do some initializations. - * @param filename - */ void open(String filename) { fileEvents = 0; reader = filename.endsWith(".hipo") ? new HipoDataSource() : new EvioSource(); @@ -259,9 +244,6 @@ void open(String filename) { } } - /** - * Close the output file and print performance info. - */ void close() { writer.close(); System.out.println(Benchmark.getInstance()); From 001ae7aef87527ac4066958ef6e1e9cc0bf75c0d Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Mon, 31 Aug 2026 22:57:45 -0400 Subject: [PATCH 43/62] normalize --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 511d54b8d8..07ea7981f6 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -205,13 +205,13 @@ void write(String output) { sleep(1000); } else { - Benchmark.getInstance().resume("write"); for (int i=0; i Date: Tue, 1 Sep 2026 08:50:21 -0400 Subject: [PATCH 44/62] start splitting --- .../jlab/clas/reco/EngineMultiProcessor.java | 18 +++++++++++------- 1 file changed, 11 insertions(+), 7 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 07ea7981f6..fcabd5dd04 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -37,6 +37,7 @@ public class EngineMultiProcessor extends EngineProcessor { CompletableFuture writerThread; ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); + List> splitQueue; ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); @@ -75,21 +76,18 @@ public void launch(String output, String... input) { readEvents = 0; writeEvents = 0; failEvents = 0; + splitQueue = new ArrayList<>(threads); readerThread = CompletableFuture.runAsync(() -> { read(input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); for (int i=0; i { process(j); })); + splitQueue.add(i, new ConcurrentLinkedQueue<>()); } while (!writerThread.isDone()) { for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); - sleep(1000); - if (readerThread.isDone()) { - System.err.println("------------------------------------------"); - System.err.println(writeEvents+"+"+skipEvents+"+"+failEvents+"=?"+readEvents); - System.err.println(readQueue.size()+" -> "+writeQueue.size()); - } + sleep(100); } } @@ -112,7 +110,7 @@ void read(String... input) { if (reader != null && reader.hasEvent()) { // sleep instead of overfilling the read queue: - if (readQueue.size() > QUEUE_SIZE*threads) sleep(100); + if (readQueue.size() > QUEUE_SIZE*threads) sleep(1000); // read the next event into the frame: else { @@ -208,6 +206,8 @@ void write(String output) { for (int i=0; i Date: Tue, 1 Sep 2026 09:59:36 -0400 Subject: [PATCH 45/62] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 61 +++++++++++-------- 1 file changed, 36 insertions(+), 25 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index fcabd5dd04..05ab9089de 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -2,8 +2,10 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; +import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Arrays; +import java.util.Deque; import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; @@ -35,11 +37,11 @@ public class EngineMultiProcessor extends EngineProcessor { CompletableFuture readerThread; CompletableFuture writerThread; - ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); + Deque procThreads = new ArrayDeque(); - List> splitQueue; ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); + List> splitQueue; ProgressPrintout progress = new ProgressPrintout(); @@ -112,28 +114,8 @@ void read(String... input) { // sleep instead of overfilling the read queue: if (readQueue.size() > QUEUE_SIZE*threads) sleep(1000); - // read the next event into the frame: - else { - Benchmark.getInstance().resume("read"); - Object o = null; - if (reader instanceof EvioSource evio) { - try { o = evio.getEventBuffer(++fileEvents, true); } - catch (EvioException ex) { - failEvents++; - ex.printStackTrace(); - } - } - else o = reader.getNextEvent(); - if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { - frame.add(o); - if (frame.size() >= FRAME_SIZE) { - readQueue.offer(frame); - readEvents += frame.size(); - frame = new ArrayList<>(FRAME_SIZE); - } - } - Benchmark.getInstance().pause("read"); - } + // read next event into frame, and fill queue if frame full: + else frame = read(frame); } // open the next input file: @@ -207,7 +189,7 @@ void write(String output) { Benchmark.getInstance().resume("write"); writer.writeEvent(e.get(i)); for (int j=0; j read(List frame) { + Benchmark.getInstance().resume("read"); + Object o = null; + if (reader instanceof EvioSource evio) { + try { o = evio.getEventBuffer(++fileEvents, true); } + catch (EvioException ex) { + failEvents++; + ex.printStackTrace(); + } + } + else o = reader.getNextEvent(); + if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { + frame.add(o); + if (frame.size() >= FRAME_SIZE) { + readQueue.offer(frame); + readEvents += frame.size(); + frame = new ArrayList<>(FRAME_SIZE); + } + } + Benchmark.getInstance().pause("read"); + return frame; + } + void close() { writer.close(); System.out.println(Benchmark.getInstance()); From cc305d1c1d9011e9295fad7a7df5d101f44feba3 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 10:32:35 -0400 Subject: [PATCH 46/62] reset --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 6 ------ 1 file changed, 6 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 05ab9089de..8d1a30c44c 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -188,8 +188,6 @@ void write(String output) { for (int i=0; i Date: Tue, 1 Sep 2026 10:41:39 -0400 Subject: [PATCH 47/62] build clara from source --- .github/workflows/ci.yml | 4 ++-- build-coatjava.sh | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index e472c075e9..31fcc295ee 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -72,7 +72,7 @@ jobs: - name: build run: | ./build-coatjava.sh --lfs --no-progress -T${{ env.nthreads }} - ./bin/install-clara -c ./coatjava ./clara + ./bin/install-clara -b -c ./coatjava ./clara - name: tar # tarball to preserve permissions run: | tar czvf coatjava.tar.gz coatjava @@ -102,7 +102,7 @@ jobs: - name: build run: | ./build-coatjava.sh --lfs --no-progress -T${{ env.nthreads }} - ./bin/install-clara -c ./coatjava ./clara + ./bin/install-clara -b -c ./coatjava ./clara - name: tar # tarball to preserve permissions run: | tar czvf coatjava.tar.gz coatjava diff --git a/build-coatjava.sh b/build-coatjava.sh index 1586b1fa05..fa3c674506 100755 --- a/build-coatjava.sh +++ b/build-coatjava.sh @@ -403,6 +403,6 @@ done echo "installed coatjava to: $prefix_dir" # install clara -if $installClara; then ./bin/install-clara -c $prefix_dir $clara_home; fi +if $installClara; then ./bin/install-clara -b -c $prefix_dir $clara_home; fi echo "COATJAVA SUCCESSFULLY BUILT !" From c6dca3d805f6524fec023f868b9ad75c92662ae1 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 10:57:49 -0400 Subject: [PATCH 48/62] fix path --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 31fcc295ee..3cef54e4a3 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -227,7 +227,7 @@ jobs: tar xzvf coatjava.tar.gz - run: ls - name: run test - run: ./bin/recon-mutil -t 4 -n 500 -y etc/services/rgd-clarode.yml -o rec_clas_018779.evio.00001.hipo -i clas_018779.evio.00001 + run: ./coatjava/bin/recon-mutil -t 4 -n 500 -y etc/services/rgd-clarode.yml -o rec_clas_018779.evio.00001.hipo -i clas_018779.evio.00001 - uses: actions/upload-artifact@v7 with: name: test_clara_result From 6e2022fb177631fb5eaca1d0581cd11450ed7be5 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 11:15:55 -0400 Subject: [PATCH 49/62] fix --- .../jlab/clas/reco/EngineMultiProcessor.java | 20 +++++++++---------- 1 file changed, 9 insertions(+), 11 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 8d1a30c44c..54dfed1cdc 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -2,10 +2,8 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; -import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Arrays; -import java.util.Deque; import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; @@ -37,7 +35,7 @@ public class EngineMultiProcessor extends EngineProcessor { CompletableFuture readerThread; CompletableFuture writerThread; - Deque procThreads = new ArrayDeque(); + ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); @@ -45,7 +43,6 @@ public class EngineMultiProcessor extends EngineProcessor { ProgressPrintout progress = new ProgressPrintout(); - int threads; int maxEvents; int skipEvents; @@ -57,29 +54,28 @@ public class EngineMultiProcessor extends EngineProcessor { public EngineMultiProcessor(OptionParser parser) { super(parser); - threads = parser.getOption("-t").intValue(); maxEvents = parser.getOption("-n").intValue(); skipEvents = parser.getOption("-s").intValue(); } public void scaling(String input, int... threads) { for (int t : threads) { - this.threads = t; - launch(String.format("scaling-%d.hipo",t),input); + launch(t, String.format("scaling-%d.hipo",t), input); } } /** * The thread launcher and collector. + * @param threads number of threads * @param output name of output file to write * @param input names of input files to read */ - public void launch(String output, String... input) { + public void launch(int threads, String output, String... input) { readEvents = 0; writeEvents = 0; failEvents = 0; splitQueue = new ArrayList<>(threads); - readerThread = CompletableFuture.runAsync(() -> { read(input); }); + readerThread = CompletableFuture.runAsync(() -> { read(threads, input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); for (int i=0; i inputs = new ArrayList<>(Arrays.asList(input)); @@ -270,6 +266,8 @@ public static void main(String[] args) { parser.addOption("-t","4","number of threads"); parser.parse(args); EngineMultiProcessor proc = new EngineMultiProcessor(parser); - proc.launch(parser.getOption("-o").stringValue(), parser.getOption("-i").stringValue()); + proc.launch(parser.getOption("-t").intValue(), + parser.getOption("-o").stringValue(), + parser.getOption("-i").stringValue()); } } From 43390ab6ca5c2dfe5d2254a0bb43097f3e19ae65 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 13:21:39 -0400 Subject: [PATCH 50/62] cleanup --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 54dfed1cdc..a447418518 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -117,7 +117,7 @@ void read(int threads, String... input) { // open the next input file: else if (!inputs.isEmpty()) open(inputs.removeFirst()); - // nothing left to do: + // no more events to read: else break; } From 92adfb7ff509fcbb71377b8cc9c94066a7c6e66c Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 19:10:43 -0400 Subject: [PATCH 51/62] add getter for number of calls --- .../java/org/jlab/utils/benchmark/ProgressPrintout.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java index c1ef3626fe..38e4ce997f 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java @@ -19,7 +19,7 @@ public class ProgressPrintout { private double printoutIntervalSeconds = 10.0; private String printoutLeadingString = ">>>>> progress : "; private Integer numberOfCalls = 0; - + public ProgressPrintout(){ this.previousPrintoutTime = System.currentTimeMillis(); this.startPrintoutTime = System.currentTimeMillis(); @@ -34,6 +34,10 @@ public ProgressPrintout(String name){ public void setInterval(double interval){ this.printoutIntervalSeconds = interval; } + + public int getNumberOfCalls() { + return numberOfCalls; + } public String getUpdateString(){ double totalElapsedTime = (this.previousPrintoutTime-this.startPrintoutTime)*1e-3; From 4b3f0f08b46377a6112a9e3b90d7429641321d98 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 19:10:52 -0400 Subject: [PATCH 52/62] add rethreading --- .../jlab/clas/reco/EngineMultiProcessor.java | 115 +++++++++++++----- 1 file changed, 87 insertions(+), 28 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index a447418518..156c76cbaf 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -27,6 +27,7 @@ */ public class EngineMultiProcessor extends EngineProcessor { + final int BENCH_SECONDS = 30; final int QUEUE_SIZE = 100; final int FRAME_SIZE = 100; @@ -35,8 +36,7 @@ public class EngineMultiProcessor extends EngineProcessor { CompletableFuture readerThread; CompletableFuture writerThread; - ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue(); - + ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); List> splitQueue; @@ -52,42 +52,72 @@ public class EngineMultiProcessor extends EngineProcessor { int fileEvents; int maxFileEvents; + public static void main(String[] args) { + OptionParser cfg = EngineProcessor.getParser(); + cfg.addOption("-t","4","number of threads"); + cfg.parse(args); + EngineMultiProcessor proc = new EngineMultiProcessor(cfg); + proc.launch(Arrays.stream(cfg.getOption("-t").stringValue().split(",")).mapToInt(Integer::parseInt).toArray(), + cfg.getOption("-o").stringValue(), + cfg.getOption("-i").stringValue()); + } + public EngineMultiProcessor(OptionParser parser) { super(parser); maxEvents = parser.getOption("-n").intValue(); skipEvents = parser.getOption("-s").intValue(); } - public void scaling(String input, int... threads) { - for (int t : threads) { - launch(t, String.format("scaling-%d.hipo",t), input); - } - } - /** * The thread launcher and collector. * @param threads number of threads * @param output name of output file to write * @param input names of input files to read */ - public void launch(int threads, String output, String... input) { - readEvents = 0; - writeEvents = 0; - failEvents = 0; - splitQueue = new ArrayList<>(threads); - readerThread = CompletableFuture.runAsync(() -> { read(threads, input); }); + public void launch(int[] threads, String output, String... input) { + shutdown(); + reset(); + splitQueue = new ArrayList<>(threads[0]); + readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); - for (int i=0; i { process(j); })); splitQueue.add(i, new ConcurrentLinkedQueue<>()); } + CompletableFuture rethreadThread = null; while (!writerThread.isDone()) { + sleep(100); for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); - sleep(100); + if (threads.length > 1 && writeEvents > 100 && rethreadThread == null) + rethreadThread = CompletableFuture.runAsync(() -> { rethread(BENCH_SECONDS,threads); }); + } + } + + /** + * Shutdown all data processing, close files. + */ + public void shutdown() { + for (CompletableFuture f : procThreads) f.cancel(true); + if (readerThread != null) readerThread.cancel(true); + if (writerThread != null) { + writerThread.cancel(true); + writer.close(); } } + + /** + * Reset counters and queues. + */ + public void reset() { + readQueue = new ConcurrentLinkedQueue<>(); + writeQueue = new ConcurrentLinkedQueue<>(); + procThreads = new ConcurrentLinkedQueue(); + readEvents = 0; + writeEvents = 0; + failEvents = 0; + } /** * The reader thread. @@ -174,7 +204,7 @@ void write(String output) { while (true) { List e = writeQueue.poll(); if (e == null) { - if (procThreads.isEmpty() && writeQueue.isEmpty()) { + if (readerThread.isDone() && procThreads.isEmpty() && writeQueue.isEmpty()) { close(); break; } @@ -192,6 +222,38 @@ void write(String output) { } } + /** + * The rethread thread. + * @param seconds delay before switching to next thread count + * @param threads thread counts to use + */ + void rethread(int seconds, int... threads) { + System.out.println("<><><> Rethreading spawned ..."); + for (int i=0; i { process(k); })); + } + while (progress.getNumberOfCalls() < 100) sleep(1000); + sleep(seconds*1000); + System.out.println(String.format("\n<><><><><> RETHREAD COUNT: %d\n",threads[i])); + System.out.println(progress.getUpdateString()); + System.out.println(Benchmark.getInstance()); + } + shutdown(); + } + + /** + * Decode an event. + * @param bytes the EVIO byte buffer + * @return decoded event + */ HipoDataEvent decode(ByteBuffer bytes) { Benchmark.getInstance().resume("evio"); EvioDataEvent evio = new EvioDataEvent(bytes.array(), ByteOrder.LITTLE_ENDIAN); @@ -207,7 +269,11 @@ HipoDataEvent decode(ByteBuffer bytes) { Benchmark.getInstance().pause("deco"); return hipo; } - + + /** + * Open a new input HIPO/EVIO event file. + * @param filename + */ void open(String filename) { fileEvents = 0; reader = filename.endsWith(".hipo") ? new HipoDataSource() : new EvioSource(); @@ -249,6 +315,9 @@ List read(List frame) { return frame; } + /** + * Close the output file and print some stuff. + */ void close() { writer.close(); System.out.println(Benchmark.getInstance()); @@ -260,14 +329,4 @@ void sleep(int milliseconds) { try { Thread.sleep(milliseconds); } catch (InterruptedException ex) {} } - - public static void main(String[] args) { - OptionParser parser = EngineProcessor.getParser(); - parser.addOption("-t","4","number of threads"); - parser.parse(args); - EngineMultiProcessor proc = new EngineMultiProcessor(parser); - proc.launch(parser.getOption("-t").intValue(), - parser.getOption("-o").stringValue(), - parser.getOption("-i").stringValue()); - } } From 2fbb0bf18a742dace9a4b256c5a553e19e1534e5 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 19:11:30 -0400 Subject: [PATCH 53/62] cleanup --- .../src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 156c76cbaf..7ef0688d25 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -90,7 +90,7 @@ public void launch(int[] threads, String output, String... input) { sleep(100); for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); - if (threads.length > 1 && writeEvents > 100 && rethreadThread == null) + if (rethreadThread == null && threads.length > 1 && writeEvents > 100) rethreadThread = CompletableFuture.runAsync(() -> { rethread(BENCH_SECONDS,threads); }); } } From 275ff789152c73ccd58a17de3a630d5323164b17 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 19:26:13 -0400 Subject: [PATCH 54/62] reuse events during rethreading --- .../jlab/clas/reco/EngineMultiProcessor.java | 30 +++++++++++-------- 1 file changed, 17 insertions(+), 13 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 7ef0688d25..7dd4783855 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -34,6 +34,7 @@ public class EngineMultiProcessor extends EngineProcessor { DataSource reader; HipoDataSync writer; + CompletableFuture rethreadThread; CompletableFuture readerThread; CompletableFuture writerThread; ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); @@ -85,7 +86,6 @@ public void launch(int[] threads, String output, String... input) { procThreads.offer(CompletableFuture.runAsync(() -> { process(j); })); splitQueue.add(i, new ConcurrentLinkedQueue<>()); } - CompletableFuture rethreadThread = null; while (!writerThread.isDone()) { sleep(100); for (CompletableFuture f : procThreads) @@ -107,18 +107,6 @@ public void shutdown() { } } - /** - * Reset counters and queues. - */ - public void reset() { - readQueue = new ConcurrentLinkedQueue<>(); - writeQueue = new ConcurrentLinkedQueue<>(); - procThreads = new ConcurrentLinkedQueue(); - readEvents = 0; - writeEvents = 0; - failEvents = 0; - } - /** * The reader thread. * @param input input filenames @@ -173,6 +161,8 @@ void process(int thread) { sleep(100); } else { + // put the event back on the queue if we're rethreading: + if (rethreadThread != null && !rethreadThread.isDone()) readQueue.offer(o); List frame = new ArrayList<>(o.size()); for (int i=0; i(); + writeQueue = new ConcurrentLinkedQueue<>(); + procThreads = new ConcurrentLinkedQueue(); + readEvents = 0; + writeEvents = 0; + failEvents = 0; + } + void sleep(int milliseconds) { try { Thread.sleep(milliseconds); } catch (InterruptedException ex) {} From fc96f681b46a2f17957b6048a6c349a853b22300 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 20:06:32 -0400 Subject: [PATCH 55/62] cleanup, more accessibility --- .../java/org/jlab/utils/benchmark/BenchmarkTimer.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java index 803a754d8b..59918ff499 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java @@ -95,9 +95,11 @@ public double getSeconds(){ @Override public String toString() { - double timePerCall = 0.0; - if (numberOfCalls.get() != 0) timePerCall = getMiliseconds() / numberOfCalls.get(); return String.format("%-15s : #Calls %12d, Total = %12.2f sec, Unit = %12.3f msec", - getName(), numberOfCalls.get(), getSeconds(), timePerCall); + getName(), numberOfCalls.get(), getSeconds(), getTimePerCall()); + } + + public double getTimePerCall() { + return numberOfCalls.get() > 0 ? getMiliseconds() / numberOfCalls.get() : 0; } } From 32fc55ea9f725c6971a487be2de6049930c7318b Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 20:06:38 -0400 Subject: [PATCH 56/62] playing --- .../src/main/java/org/jlab/utils/benchmark/Benchmark.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/Benchmark.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/Benchmark.java index 708cf988c5..5b8d5ddf08 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/Benchmark.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/Benchmark.java @@ -81,6 +81,11 @@ public BenchmarkTimer getTotal(String name) { return total; } + public String toCSV() { + return String.join(",",timerStore.entrySet().stream(). + map(b -> String.format("%s:%.4f",b.getKey(),b.getValue().getTimePerCall())).toArray(String[]::new)); + } + @Override public String toString(){ StringBuilder s = new StringBuilder(); From e6cc9c3f5d369bbdfed006177d914fc2c53c6cf0 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 20:26:38 -0400 Subject: [PATCH 57/62] start --- .../java/org/jlab/clas/reco/SerialHoncho.java | 40 +++++++++++++++++++ 1 file changed, 40 insertions(+) create mode 100644 common-tools/clas-reco/src/main/java/org/jlab/clas/reco/SerialHoncho.java diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/SerialHoncho.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/SerialHoncho.java new file mode 100644 index 0000000000..27cdfe737b --- /dev/null +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/SerialHoncho.java @@ -0,0 +1,40 @@ +package org.jlab.clas.reco; + +import java.util.TreeMap; +import java.util.TreeSet; +import org.jlab.detector.calib.utils.ConstantsManager; +import org.jlab.detector.decode.CLASDecoder; +import org.jlab.detector.helicity.HelicityState; +import org.jlab.detector.scalers.DaqScalersSequence; +import org.jlab.jnp.hipo4.data.Bank; +import org.jlab.jnp.hipo4.data.Event; + +/** + * + * @author baltzell + */ +public class SerialHoncho { + + Bank[] tag1banks; + Bank runConfig; + Bank helicityAdc; + ConstantsManager conman; + TreeMap eventUnix; + TreeSet helicities; + DaqScalersSequence scalers; + + public SerialHoncho(){} + + public Event[] doggy(Event event) { + scalers.add(event); + event.read(runConfig); + event.read(helicityAdc); + if (runConfig.getRows() > 0) { + int unix = runConfig.getInt("unixtime",0); + int evno = runConfig.getInt("event",0); + if (unix > 0 && evno > 0) eventUnix.put(evno, unix); + } + helicities.add(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman)); + return new Event[]{CLASDecoder.createTaggedEvent(event, runConfig, tag1banks)}; + } +} From 57c42ac65336c609be0a9fcceed2a0fbb28eee5a Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 22:31:43 -0400 Subject: [PATCH 58/62] more --- .../jlab/clas/reco/EngineMultiProcessor.java | 130 ++++++++++-------- .../utils/benchmark/ProgressPrintout.java | 4 +- .../org/jlab/utils/options/OptionParser.java | 7 +- 3 files changed, 82 insertions(+), 59 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java index 7dd4783855..76021426c6 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -28,12 +28,14 @@ public class EngineMultiProcessor extends EngineProcessor { final int BENCH_SECONDS = 30; - final int QUEUE_SIZE = 100; - final int FRAME_SIZE = 100; + final int CHUNKS_PER_QUEUE = 100; + final int EVENTS_PER_CHUNK = 100; + // File reader and writer: DataSource reader; HipoDataSync writer; + // threads and queeues: CompletableFuture rethreadThread; CompletableFuture readerThread; CompletableFuture writerThread; @@ -41,34 +43,50 @@ public class EngineMultiProcessor extends EngineProcessor { ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); List> splitQueue; - ProgressPrintout progress = new ProgressPrintout(); - + + // static parameters: int maxEvents; int skipEvents; + // progress counters: int readEvents; int writeEvents; int failEvents; int fileEvents; int maxFileEvents; - + + public EngineMultiProcessor(OptionParser parser) { + super(parser); + maxEvents = parser.getOption("-n").intValue(); + skipEvents = parser.getOption("-s").intValue(); + } + + /** + * The "recon-mutil" command-line tool. + * @param args + */ public static void main(String[] args) { - OptionParser cfg = EngineProcessor.getParser(); - cfg.addOption("-t","4","number of threads"); + OptionParser cfg = EngineMultiProcessor.getParser(); cfg.parse(args); EngineMultiProcessor proc = new EngineMultiProcessor(cfg); proc.launch(Arrays.stream(cfg.getOption("-t").stringValue().split(",")).mapToInt(Integer::parseInt).toArray(), cfg.getOption("-o").stringValue(), - cfg.getOption("-i").stringValue()); + cfg.getInputList().stream().toArray(String[]::new)); } - public EngineMultiProcessor(OptionParser parser) { - super(parser); - maxEvents = parser.getOption("-n").intValue(); - skipEvents = parser.getOption("-s").intValue(); + /** + * Add threads option and replace -i option with argument list. + * @return + */ + public static OptionParser getParser() { + OptionParser p = EngineProcessor.getParser(); + p.addOption("-t","4","number of threads"); + p.removeOption("-i"); + p.setRequiresInputList(true); + return p; } - + /** * The thread launcher and collector. * @param threads number of threads @@ -76,7 +94,6 @@ public EngineMultiProcessor(OptionParser parser) { * @param input names of input files to read */ public void launch(int[] threads, String output, String... input) { - shutdown(); reset(); splitQueue = new ArrayList<>(threads[0]); readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); @@ -90,20 +107,11 @@ public void launch(int[] threads, String output, String... input) { sleep(100); for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); - if (rethreadThread == null && threads.length > 1 && writeEvents > 100) + if (threads.length > 1 && rethreadThread == null && writeEvents > 100) { rethreadThread = CompletableFuture.runAsync(() -> { rethread(BENCH_SECONDS,threads); }); - } - } - - /** - * Shutdown all data processing, close files. - */ - public void shutdown() { - for (CompletableFuture f : procThreads) f.cancel(true); - if (readerThread != null) readerThread.cancel(true); - if (writerThread != null) { - writerThread.cancel(true); - writer.close(); + rethreadThread.join(); + reset(); + } } } @@ -116,8 +124,8 @@ void read(int threads, String... input) { // convert input filenames to a list: List inputs = new ArrayList<>(Arrays.asList(input)); - // initialize the event frame: - List frame = new ArrayList<>(FRAME_SIZE); + // initialize the event chunk: + List chunk = new ArrayList<>(EVENTS_PER_CHUNK); // loop over input events: while ( (maxEvents < 1 || readEvents < maxEvents) && @@ -126,10 +134,10 @@ void read(int threads, String... input) { if (reader != null && reader.hasEvent()) { // sleep instead of overfilling the read queue: - if (readQueue.size() > QUEUE_SIZE*threads) sleep(1000); + if (readQueue.size() > CHUNKS_PER_QUEUE*threads) sleep(1000); - // read next event into frame, and fill queue if frame full: - else frame = read(frame); + // read next event into chunk, and fill queue if chunk full: + else chunk = read(chunk); } // open the next input file: @@ -139,11 +147,11 @@ void read(int threads, String... input) { else break; } - // write leftover, partial frame: - if (!frame.isEmpty()) { - System.err.println("writing partial frame: "+frame.size()); - readEvents += frame.size(); - readQueue.offer(frame); + // write leftover, partial chunk: + if (!chunk.isEmpty()) { + System.err.println("writing partial chunk: "+chunk.size()); + readEvents += chunk.size(); + readQueue.offer(chunk); } reader.close(); } @@ -163,7 +171,7 @@ void process(int thread) { else { // put the event back on the queue if we're rethreading: if (rethreadThread != null && !rethreadThread.isDone()) readQueue.offer(o); - List frame = new ArrayList<>(o.size()); + List chunk = new ArrayList<>(o.size()); for (int i=0; i<><> Rethreading spawned ..."); + System.out.println("~~~~~~~~~ Rethreading Initiated ~~~~~~~~~"); for (int i=0; i<><><><> RETHREAD COUNT: %d\n",threads[i])); + System.out.println(String.format("\n~~~~~~~~~ Rethreading Count: %d ~~~~~~~~~\n",threads[i])); System.out.println(progress.getUpdateString()); System.out.println(Benchmark.getInstance()); } - shutdown(); } /** @@ -279,12 +287,12 @@ void open(String filename) { } /** - * Read the next event into the frame, and, if it's full, queue the frame + * Read the next event into the chunk, and, if it's full, queue the chunk * and make a new one. - * @param frame - * @return modified frame + * @param chunk + * @return modified chunk */ - List read(List frame) { + List read(List chunk) { Benchmark.getInstance().resume("read"); Object o = null; if (reader instanceof EvioSource evio) { @@ -296,15 +304,15 @@ List read(List frame) { } else o = reader.getNextEvent(); if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { - frame.add(o); - if (frame.size() >= FRAME_SIZE) { - readQueue.offer(frame); - readEvents += frame.size(); - frame = new ArrayList<>(FRAME_SIZE); + chunk.add(o); + if (chunk.size() >= EVENTS_PER_CHUNK) { + readQueue.offer(chunk); + readEvents += chunk.size(); + chunk = new ArrayList<>(EVENTS_PER_CHUNK); } } Benchmark.getInstance().pause("read"); - return frame; + return chunk; } /** @@ -318,9 +326,15 @@ void close() { } /** - * Reset counters and queues. + * Shutdown all threads, close files, and reset counters. */ void reset() { + for (CompletableFuture f : procThreads) f.cancel(true); + if (readerThread != null) readerThread.cancel(true); + if (writerThread != null) { + writerThread.cancel(true); + writer.close(); + } readQueue = new ConcurrentLinkedQueue<>(); writeQueue = new ConcurrentLinkedQueue<>(); procThreads = new ConcurrentLinkedQueue(); @@ -328,7 +342,11 @@ void reset() { writeEvents = 0; failEvents = 0; } - + + /** + * Catch interruptions in sleep. + * @param milliseconds + */ void sleep(int milliseconds) { try { Thread.sleep(milliseconds); } catch (InterruptedException ex) {} diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java index 38e4ce997f..b13dc5aa1d 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java @@ -63,10 +63,10 @@ public void updateStatus(){ this.previousPrintoutTime = System.currentTimeMillis(); this.startPrintoutTime = this.previousPrintoutTime; } - else { + else if (this.printoutIntervalSeconds > 0) { Long currentTime = System.currentTimeMillis(); Double elapsedTime = (currentTime - this.previousPrintoutTime)*1e-3; - if(elapsedTime >= this.printoutIntervalSeconds){ + if (elapsedTime >= this.printoutIntervalSeconds){ this.previousPrintoutTime = System.currentTimeMillis(); System.out.println(this.getUpdateString()); } diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionParser.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionParser.java index 1bdb2c5222..c362aff883 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionParser.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionParser.java @@ -70,7 +70,12 @@ public void addRequired(String key,String desc){ option.setDescription(desc); requiredOptions.put(key, option); } - + + public void removeOption(String key) { + optionsDescriptors.remove(key); + requiredOptions.remove(key); + } + public void addOption(String key, String defaultValue){ check(key, optionsDescriptors.keySet()); OptionValue option = new OptionValue(key,defaultValue); From 6933b44686bd69bb75da66bf90432838a7cbcffa Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 23:11:37 -0400 Subject: [PATCH 59/62] cleanup --- bin/recon-mutil | 2 +- ...ineMultiProcessor.java => ReconMutil.java} | 56 +++++++++---------- 2 files changed, 27 insertions(+), 31 deletions(-) rename common-tools/clas-reco/src/main/java/org/jlab/clas/reco/{EngineMultiProcessor.java => ReconMutil.java} (93%) diff --git a/bin/recon-mutil b/bin/recon-mutil index 533f8130bb..3aaec2480b 100755 --- a/bin/recon-mutil +++ b/bin/recon-mutil @@ -8,5 +8,5 @@ export MALLOC_ARENA_MAX=1 java ${JAVA_OPTS-} -Xms10240m -XX:+UseParallelGC ${jvm_options[@]} \ -cp ${COATJAVA_CLASSPATH:-''} \ - org.jlab.clas.reco.EngineMultiProcessor \ + org.jlab.clas.reco.ReconMutil \ ${class_options[@]} diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java similarity index 93% rename from common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java rename to common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 76021426c6..e0ec389509 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -25,8 +25,9 @@ * * @author baltzell */ -public class EngineMultiProcessor extends EngineProcessor { +public class ReconMutil extends EngineProcessor { + // Performance parameters: final int BENCH_SECONDS = 30; final int CHUNKS_PER_QUEUE = 100; final int EVENTS_PER_CHUNK = 100; @@ -35,58 +36,53 @@ public class EngineMultiProcessor extends EngineProcessor { DataSource reader; HipoDataSync writer; - // threads and queeues: + // Threads and queues: CompletableFuture rethreadThread; CompletableFuture readerThread; CompletableFuture writerThread; ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); - List> splitQueue; - ProgressPrintout progress = new ProgressPrintout(); - // static parameters: + // Static parameters: int maxEvents; int skipEvents; - // progress counters: + // Progress counters: int readEvents; int writeEvents; int failEvents; int fileEvents; int maxFileEvents; - - public EngineMultiProcessor(OptionParser parser) { + ProgressPrintout progress = new ProgressPrintout(); + + public ReconMutil(OptionParser parser) { super(parser); maxEvents = parser.getOption("-n").intValue(); skipEvents = parser.getOption("-s").intValue(); } + public static OptionParser getParser() { + OptionParser p = EngineProcessor.getParser(); + p.addOption("-t","4","number of threads"); + p.removeOption("-i"); + p.setRequiresInputList(true); + return p; + } + /** - * The "recon-mutil" command-line tool. + * The "recon-mutil" command-line entry-point. * @param args */ public static void main(String[] args) { - OptionParser cfg = EngineMultiProcessor.getParser(); + OptionParser cfg = ReconMutil.getParser(); cfg.parse(args); - EngineMultiProcessor proc = new EngineMultiProcessor(cfg); + ReconMutil proc = new ReconMutil(cfg); proc.launch(Arrays.stream(cfg.getOption("-t").stringValue().split(",")).mapToInt(Integer::parseInt).toArray(), cfg.getOption("-o").stringValue(), cfg.getInputList().stream().toArray(String[]::new)); } - /** - * Add threads option and replace -i option with argument list. - * @return - */ - public static OptionParser getParser() { - OptionParser p = EngineProcessor.getParser(); - p.addOption("-t","4","number of threads"); - p.removeOption("-i"); - p.setRequiresInputList(true); - return p; - } - /** * The thread launcher and collector. * @param threads number of threads @@ -95,13 +91,11 @@ public static OptionParser getParser() { */ public void launch(int[] threads, String output, String... input) { reset(); - splitQueue = new ArrayList<>(threads[0]); readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); for (int i=0; i { process(j); })); - splitQueue.add(i, new ConcurrentLinkedQueue<>()); } while (!writerThread.isDone()) { sleep(100); @@ -196,9 +190,11 @@ void process(int thread) { * @param output output filename */ void write(String output) { - writer = new HipoDataSync(); - writer.setCompressionType(2); - writer.open(output); + if (output != null) { + writer = new HipoDataSync(); + writer.setCompressionType(2); + writer.open(output); + } while (true) { List e = writeQueue.poll(); if (e == null) { @@ -211,7 +207,7 @@ void write(String output) { else { for (int i=0; i Date: Tue, 1 Sep 2026 23:24:13 -0400 Subject: [PATCH 60/62] privatize --- .../java/org/jlab/clas/reco/ReconMutil.java | 19 +++++++------------ 1 file changed, 7 insertions(+), 12 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index e0ec389509..921c65a568 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -37,9 +37,9 @@ public class ReconMutil extends EngineProcessor { HipoDataSync writer; // Threads and queues: - CompletableFuture rethreadThread; CompletableFuture readerThread; CompletableFuture writerThread; + CompletableFuture rethreadThread; ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); @@ -56,26 +56,21 @@ public class ReconMutil extends EngineProcessor { int maxFileEvents; ProgressPrintout progress = new ProgressPrintout(); - public ReconMutil(OptionParser parser) { + ReconMutil(OptionParser parser) { super(parser); maxEvents = parser.getOption("-n").intValue(); skipEvents = parser.getOption("-s").intValue(); } - public static OptionParser getParser() { - OptionParser p = EngineProcessor.getParser(); - p.addOption("-t","4","number of threads"); - p.removeOption("-i"); - p.setRequiresInputList(true); - return p; - } - /** * The "recon-mutil" command-line entry-point. * @param args */ public static void main(String[] args) { - OptionParser cfg = ReconMutil.getParser(); + OptionParser cfg = EngineProcessor.getParser(); + cfg.addOption("-t","4","number of threads"); + cfg.removeOption("-i"); + cfg.setRequiresInputList(true); cfg.parse(args); ReconMutil proc = new ReconMutil(cfg); proc.launch(Arrays.stream(cfg.getOption("-t").stringValue().split(",")).mapToInt(Integer::parseInt).toArray(), @@ -89,7 +84,7 @@ public static void main(String[] args) { * @param output name of output file to write * @param input names of input files to read */ - public void launch(int[] threads, String output, String... input) { + void launch(int[] threads, String output, String... input) { reset(); readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); From 83ed13ebf8c88c1dbc3aed5f8ff50add1d7dfa25 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 23:26:01 -0400 Subject: [PATCH 61/62] cleanup --- .../src/main/java/org/jlab/clas/reco/ReconMutil.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 921c65a568..1fb34f32d9 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -63,8 +63,8 @@ public class ReconMutil extends EngineProcessor { } /** - * The "recon-mutil" command-line entry-point. - * @param args + * The command-line entry-point known as recon-mutil. + * @param args command-line arguments */ public static void main(String[] args) { OptionParser cfg = EngineProcessor.getParser(); From a7bf60d962ce7e511a872a09b454f74b82f404e5 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 1 Sep 2026 23:30:23 -0400 Subject: [PATCH 62/62] cleanup --- .../java/org/jlab/clas/reco/ReconMutil.java | 22 +++++++++---------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 1fb34f32d9..d9fe7aa6ea 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -25,7 +25,7 @@ * * @author baltzell */ -public class ReconMutil extends EngineProcessor { +public final class ReconMutil extends EngineProcessor { // Performance parameters: final int BENCH_SECONDS = 30; @@ -63,19 +63,19 @@ public class ReconMutil extends EngineProcessor { } /** - * The command-line entry-point known as recon-mutil. + * The command-line entry-point known as "recon-mutil". * @param args command-line arguments */ public static void main(String[] args) { - OptionParser cfg = EngineProcessor.getParser(); - cfg.addOption("-t","4","number of threads"); - cfg.removeOption("-i"); - cfg.setRequiresInputList(true); - cfg.parse(args); - ReconMutil proc = new ReconMutil(cfg); - proc.launch(Arrays.stream(cfg.getOption("-t").stringValue().split(",")).mapToInt(Integer::parseInt).toArray(), - cfg.getOption("-o").stringValue(), - cfg.getInputList().stream().toArray(String[]::new)); + OptionParser o = EngineProcessor.getParser(); + o.addOption("-t","4","number of threads"); + o.removeOption("-i"); + o.setRequiresInputList(true); + o.parse(args); + ReconMutil r = new ReconMutil(o); + r.launch(Arrays.stream(o.getOption("-t").stringValue().split(",")).mapToInt(Integer::parseInt).toArray(), + o.getOption("-o").stringValue(), + o.getInputList().stream().toArray(String[]::new)); } /**