From 79e56c5dc1aaf6dc1488750b298689dbb656b0fd Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 11:57:21 -0400 Subject: [PATCH 01/27] multi it --- bin/recon-mutil | 13 ++ .../jlab/clas/reco/EngineMultiProcessor.java | 147 ++++++++++++++++++ .../org/jlab/clas/reco/EngineProcessor.java | 11 +- validation/advanced-tests/run-eb-tests.sh | 2 +- 4 files changed, 167 insertions(+), 6 deletions(-) create mode 100755 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 100755 index 0000000000..f9be37b69a --- /dev/null +++ b/bin/recon-mutil @@ -0,0 +1,13 @@ +#!/bin/bash + +. `dirname $0`/../libexec/env.sh + +split_cli $@ + +export MALLOC_ARENA_MAX=1 + +java ${JAVA_OPTS-} -Xmx1536m -Xms1024m -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..c70cda54c6 --- /dev/null +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -0,0 +1,147 @@ +package org.jlab.clas.reco; + +import java.util.LinkedList; +import java.util.logging.Logger; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import org.json.JSONObject; +import org.jlab.io.base.DataEvent; +import org.jlab.io.base.DataSource; +import org.jlab.io.evio.EvioSource; +import org.jlab.io.hipo.HipoDataSource; +import org.jlab.io.hipo.HipoDataSync; +import org.jlab.utils.ClaraYaml; +import org.jlab.utils.options.OptionParser; + +/** + * + * @author baltzell + */ +public class EngineMultiProcessor { + + int threads; + DataSource reader; + HipoDataSync writer; + LinkedList queue; + EngineProcessor engines; + ExecutorService executor; + + public EngineMultiProcessor(int threads) { + this.threads = threads; + engines = new EngineProcessor(); + queue = new LinkedList<>(); + executor = Executors.newFixedThreadPool(threads); + } + + public void open(String input, String output, int events, int skip) { + writer = new HipoDataSync(); + writer.setCompressionType(2); + if (input.endsWith(".hipo")) reader = new HipoDataSource(); + else reader = new EvioSource(); + reader.open(input); + writer.open(output); + engines.updateDictionary((HipoDataSource)reader, writer); + while (queue.size() < threads && reader.hasEvent()) + queue.offer(reader.getNextEvent()); + for (int i=0; i { + engines.processEvent(event); + return event; + }, executor).thenRun(() -> callback(event)); + } + + public synchronized void callback(DataEvent event) { + if (event != null) writer.writeEvent(event); + if (reader.hasEvent()) queue.offer(reader.getNextEvent()); + else if (queue.isEmpty()) { + executor.shutdown(); + writer.close(); + reader.close(); + } + if (!queue.isEmpty()) spawn(queue.poll()); + } + + public static void main(String[] args) { + + 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"); + parser.addOption("-y","0","yaml file"); + parser.addOption("-u","true","update dictionary from writer ? "); + parser.addOption("-S",null,"schema directory"); + parser.addOption("-B",null,"background file"); + parser.addOption("-P",null,"preload file for post-processing"); + parser.addOption("-R","0","rebuild scalers"); + parser.addOption("-H","0","restream helicity"); + parser.addOption("-t","1","threads"); + parser.parse(args); + parser.syncLogLevel(Logger.getLogger(EngineMultiProcessor.class.getPackage().getName())); + + EngineMultiProcessor multi = new EngineMultiProcessor(parser.getOption("-t").intValue()); + + if (parser.getOption("-u").stringValue().contains("false")) + multi.engines.updateDictionary = false; + + // service list from YAML: + if(!parser.getOption("-y").stringValue().equals("0")) { + ClaraYaml yaml = new ClaraYaml(parser.getOption("-y").stringValue()); + if (yaml.schemaDirectory() != null) { + multi.engines.setBanksToKeep(yaml.schemaDirectory()); + } + for (JSONObject service : yaml.services()) { + JSONObject cfg = yaml.filter(service.getString("name")); + if (cfg.length() > 0) { + multi.engines.addEngine(service.getString("name"),service.getString("class"),cfg.toString()); + } else { + multi.engines.addEngine(service.getString("name"),service.getString("class")); + } + } + } + // built-in service list: + else if (parser.getOption("-c").intValue() > 0) { + if (parser.getOption("-c").intValue() > 2) { + multi.engines.initCaloDebug(); + } else if(parser.getOption("-c").intValue() == 2) { + multi.engines.initAll(); + } else { + multi.engines.initDefault(); + } + } + // user-defined service list: + else { + for(String engine : parser.getInputList()){ + System.out.println("Adding reconstruction engine " + engine); + multi.engines.addEngine(engine); + } + } + + // command-line schema overrides YAML: + if (parser.getOption("-S").stringValue() != null) + multi.engines.setBanksToKeep(parser.getOption("-S").stringValue()); + + // command-line filename for background merging overrides YAML: + if (parser.getOption("-B").stringValue() != null) + multi.engines.setBackgroundFiles(parser.getOption("-B").stringValue()); + + // command-line filename for post-processing overrides YAML: + if (parser.getOption("-P").stringValue() != null) { + multi.engines.setPreloadFiles(parser.getOption("-P").stringValue(), + parser.getOption("-H").intValue()!=0, + parser.getOption("-R").intValue()!=0); + } + + multi.open(parser.getOption("-i").stringValue(), + parser.getOption("-o").stringValue(), + parser.getOption("-s").intValue(), + parser.getOption("-n").intValue()); + } + +} 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 5f17e9b6d4..ca1e3580a5 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 @@ -38,10 +38,11 @@ public class EngineProcessor { private 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"); + protected boolean updateDictionary = true; + private CLASDecoder4 decoder = new CLASDecoder4(); public EngineProcessor(){} @@ -55,7 +56,7 @@ private ReconstructionEngine findEngine(String clazz) { return null; } - private void setBackgroundFiles(String filenames) { + protected void setBackgroundFiles(String filenames) { if (findEngine(ENGINE_CLASS_BG) == null) { LOGGER.info("Adding BackgroundEngine for -B option."); addEngine("BG",ENGINE_CLASS_BG); @@ -64,7 +65,7 @@ private void setBackgroundFiles(String filenames) { findEngine(ENGINE_CLASS_BG).init(); } - private void setPreloadFiles(String filenames, boolean restream, boolean rebuild) { + protected void setPreloadFiles(String filenames, boolean restream, boolean rebuild) { if (findEngine(ENGINE_CLASS_PP) == null) { LOGGER.info("Adding PostprocEngine for -P option."); addEngine("BG",ENGINE_CLASS_PP); @@ -75,7 +76,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(); @@ -90,7 +91,7 @@ private void updateDictionary(HipoDataSource source, HipoDataSync sync){ } } - private void setBanksToKeep(String schemaDirectory) { + protected void setBanksToKeep(String schemaDirectory) { if (!Files.isDirectory((new File(schemaDirectory)).toPath())) { LOGGER.log(Level.SEVERE, "Invalid schema directory, aborting: "+schemaDirectory); System.exit(1); diff --git a/validation/advanced-tests/run-eb-tests.sh b/validation/advanced-tests/run-eb-tests.sh index 05c7470db5..dcaa9166da 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 4 -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 ed45dc0f202b002f8e40217d025ab8fa3d9f5abf Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 13:02:31 -0400 Subject: [PATCH 02/27] change it up --- bin/recon-mutil | 2 +- .../jlab/clas/reco/EngineMultiProcessor.java | 61 ++++++++++++------- validation/advanced-tests/run-eb-tests.sh | 2 +- 3 files changed, 40 insertions(+), 25 deletions(-) diff --git a/bin/recon-mutil b/bin/recon-mutil index f9be37b69a..3138056c13 100755 --- a/bin/recon-mutil +++ b/bin/recon-mutil @@ -6,7 +6,7 @@ split_cli $@ export MALLOC_ARENA_MAX=1 -java ${JAVA_OPTS-} -Xmx1536m -Xms1024m -XX:+UseSerialGC ${jvm_options[@]} \ +java ${JAVA_OPTS-} -Xmx10536m -Xms6024m -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 index c70cda54c6..5603ebf44d 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 @@ -1,8 +1,8 @@ package org.jlab.clas.reco; -import java.util.LinkedList; import java.util.logging.Logger; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import org.json.JSONObject; @@ -23,16 +23,37 @@ public class EngineMultiProcessor { int threads; DataSource reader; HipoDataSync writer; - LinkedList queue; + ConcurrentLinkedQueue readQueue; + ConcurrentLinkedQueue writeQueue; EngineProcessor engines; ExecutorService executor; public EngineMultiProcessor(int threads) { this.threads = threads; engines = new EngineProcessor(); - queue = new LinkedList<>(); + readQueue = new ConcurrentLinkedQueue<>(); + writeQueue = new ConcurrentLinkedQueue<>(); executor = Executors.newFixedThreadPool(threads); } + void read() throws InterruptedException { + while (reader.hasEvent()) { + if (readQueue.size() > 10*threads) Thread.sleep(1000); + else readQueue.offer(reader.getNextEvent()); + } + } + void process() { + while (!readQueue.isEmpty()) { + DataEvent event = readQueue.poll(); + engines.processEvent(event); + writeQueue.offer(event); + } + } + void write() throws InterruptedException { + while (true) { + if (!writeQueue.isEmpty()) writer.writeEvent(writeQueue.poll()); + else Thread.sleep(1000); + } + } public void open(String input, String output, int events, int skip) { writer = new HipoDataSync(); @@ -42,27 +63,21 @@ public void open(String input, String output, int events, int skip) { reader.open(input); writer.open(output); engines.updateDictionary((HipoDataSource)reader, writer); - while (queue.size() < threads && reader.hasEvent()) - queue.offer(reader.getNextEvent()); - for (int i=0; i { - engines.processEvent(event); - return event; - }, executor).thenRun(() -> callback(event)); - } - - public synchronized void callback(DataEvent event) { - if (event != null) writer.writeEvent(event); - if (reader.hasEvent()) queue.offer(reader.getNextEvent()); - else if (queue.isEmpty()) { - executor.shutdown(); - writer.close(); - reader.close(); - } - if (!queue.isEmpty()) spawn(queue.poll()); + try { read(); } catch (InterruptedException ex) {} + return true; + }, executor); + CompletableFuture.supplyAsync(() -> { + try { write(); } catch (InterruptedException ex) {} + return true; + }, executor); + try { Thread.sleep(1000); } catch (InterruptedException ex) {} + System.err.println("DOGGIES"); + process(); + System.err.println("KITTIES"); + executor.shutdownNow(); + writer.close(); + reader.close(); } public static void main(String[] args) { diff --git a/validation/advanced-tests/run-eb-tests.sh b/validation/advanced-tests/run-eb-tests.sh index dcaa9166da..656d703b4a 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-mutil -t 4 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 +../../coatjava/bin/recon-mutil -t 24 -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 052fd3fe142e36e11b50ed6b0985addeed6be836 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 13:35:15 -0400 Subject: [PATCH 03/27] hmm --- .../jlab/clas/reco/EngineMultiProcessor.java | 19 +++++++++++++------ 1 file changed, 13 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 5603ebf44d..68c94f3eea 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 @@ -5,6 +5,7 @@ import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.ThreadPoolExecutor; import org.json.JSONObject; import org.jlab.io.base.DataEvent; import org.jlab.io.base.DataSource; @@ -50,11 +51,11 @@ void process() { } void write() throws InterruptedException { while (true) { - if (!writeQueue.isEmpty()) writer.writeEvent(writeQueue.poll()); + if (!writeQueue.isEmpty()) + writer.writeEvent(writeQueue.poll()); else Thread.sleep(1000); } } - public void open(String input, String output, int events, int skip) { writer = new HipoDataSync(); writer.setCompressionType(2); @@ -71,10 +72,16 @@ public void open(String input, String output, int events, int skip) { try { write(); } catch (InterruptedException ex) {} return true; }, executor); - try { Thread.sleep(1000); } catch (InterruptedException ex) {} - System.err.println("DOGGIES"); - process(); - System.err.println("KITTIES"); + try { Thread.sleep(5000); } catch (InterruptedException ex) {} + for (int i=0; i { + process(); + return true; + }, executor); + } + while (!readQueue.isEmpty() || !writeQueue.isEmpty()) { + try { Thread.sleep(1000); } catch (InterruptedException ex) {} + } executor.shutdownNow(); writer.close(); reader.close(); From 3583bce62292d5c12059bee79b16082b5bfd4d8a Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 13:54:02 -0400 Subject: [PATCH 04/27] protect with another queue --- .../java/org/jlab/clas/reco/EngineMultiProcessor.java | 8 ++++++-- 1 file changed, 6 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 68c94f3eea..db09d10f36 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,6 +26,7 @@ public class EngineMultiProcessor { HipoDataSync writer; ConcurrentLinkedQueue readQueue; ConcurrentLinkedQueue writeQueue; + ConcurrentLinkedQueue procQueue; EngineProcessor engines; ExecutorService executor; @@ -34,18 +35,21 @@ public EngineMultiProcessor(int threads) { engines = new EngineProcessor(); readQueue = new ConcurrentLinkedQueue<>(); writeQueue = new ConcurrentLinkedQueue<>(); + procQueue = new ConcurrentLinkedQueue<>(); executor = Executors.newFixedThreadPool(threads); } void read() throws InterruptedException { while (reader.hasEvent()) { - if (readQueue.size() > 10*threads) Thread.sleep(1000); + if (readQueue.size() > 10*threads) Thread.sleep(100); else readQueue.offer(reader.getNextEvent()); } } void process() { while (!readQueue.isEmpty()) { DataEvent event = readQueue.poll(); + procQueue.offer(event); engines.processEvent(event); + procQueue.remove(event); writeQueue.offer(event); } } @@ -79,7 +83,7 @@ public void open(String input, String output, int events, int skip) { return true; }, executor); } - while (!readQueue.isEmpty() || !writeQueue.isEmpty()) { + while (!readQueue.isEmpty() || !writeQueue.isEmpty() || !procQueue.isEmpty()) { try { Thread.sleep(1000); } catch (InterruptedException ex) {} } executor.shutdownNow(); From ddf76f70f133b14e54e35a5fafd7f12faff994b6 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 13:58:19 -0400 Subject: [PATCH 05/27] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 63 +++++++++---------- 1 file changed, 31 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 db09d10f36..360a9e82d9 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 @@ -5,7 +5,6 @@ import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.ThreadPoolExecutor; import org.json.JSONObject; import org.jlab.io.base.DataEvent; import org.jlab.io.base.DataSource; @@ -24,42 +23,17 @@ public class EngineMultiProcessor { int threads; DataSource reader; HipoDataSync writer; - ConcurrentLinkedQueue readQueue; - ConcurrentLinkedQueue writeQueue; - ConcurrentLinkedQueue procQueue; - EngineProcessor engines; + ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); + EngineProcessor engines = new EngineProcessor(); ExecutorService executor; public EngineMultiProcessor(int threads) { this.threads = threads; - engines = new EngineProcessor(); - readQueue = new ConcurrentLinkedQueue<>(); - writeQueue = new ConcurrentLinkedQueue<>(); - procQueue = new ConcurrentLinkedQueue<>(); executor = Executors.newFixedThreadPool(threads); } - void read() throws InterruptedException { - while (reader.hasEvent()) { - if (readQueue.size() > 10*threads) Thread.sleep(100); - else readQueue.offer(reader.getNextEvent()); - } - } - void process() { - while (!readQueue.isEmpty()) { - DataEvent event = readQueue.poll(); - procQueue.offer(event); - engines.processEvent(event); - procQueue.remove(event); - writeQueue.offer(event); - } - } - void write() throws InterruptedException { - while (true) { - if (!writeQueue.isEmpty()) - writer.writeEvent(writeQueue.poll()); - else Thread.sleep(1000); - } - } + public void open(String input, String output, int events, int skip) { writer = new HipoDataSync(); writer.setCompressionType(2); @@ -77,7 +51,7 @@ public void open(String input, String output, int events, int skip) { return true; }, executor); try { Thread.sleep(5000); } catch (InterruptedException ex) {} - for (int i=0; i { process(); return true; @@ -91,6 +65,31 @@ public void open(String input, String output, int events, int skip) { reader.close(); } + void read() throws InterruptedException { + while (reader.hasEvent()) { + if (readQueue.size() > 10*threads) Thread.sleep(100); + else readQueue.offer(reader.getNextEvent()); + } + } + + void process() { + while (!readQueue.isEmpty()) { + DataEvent event = readQueue.poll(); + procQueue.offer(event); + engines.processEvent(event); + procQueue.remove(event); + writeQueue.offer(event); + } + } + + void write() throws InterruptedException { + while (true) { + if (!writeQueue.isEmpty()) + writer.writeEvent(writeQueue.poll()); + else Thread.sleep(1000); + } + } + public static void main(String[] args) { OptionParser parser = new OptionParser("recon-util"); From 622029223be647f811a7818840220ebf42e1f045 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 14:10:54 -0400 Subject: [PATCH 06/27] cleanup --- bin/recon-mutil | 13 -- .../jlab/clas/reco/EngineMultiProcessor.java | 172 ------------------ .../org/jlab/clas/reco/EngineProcessor.java | 83 ++++++++- 3 files changed, 81 insertions(+), 187 deletions(-) delete mode 100755 bin/recon-mutil delete 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 deleted file mode 100755 index 3138056c13..0000000000 --- a/bin/recon-mutil +++ /dev/null @@ -1,13 +0,0 @@ -#!/bin/bash - -. `dirname $0`/../libexec/env.sh - -split_cli $@ - -export MALLOC_ARENA_MAX=1 - -java ${JAVA_OPTS-} -Xmx10536m -Xms6024m -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 deleted file mode 100644 index 360a9e82d9..0000000000 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java +++ /dev/null @@ -1,172 +0,0 @@ -package org.jlab.clas.reco; - -import java.util.logging.Logger; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import org.json.JSONObject; -import org.jlab.io.base.DataEvent; -import org.jlab.io.base.DataSource; -import org.jlab.io.evio.EvioSource; -import org.jlab.io.hipo.HipoDataSource; -import org.jlab.io.hipo.HipoDataSync; -import org.jlab.utils.ClaraYaml; -import org.jlab.utils.options.OptionParser; - -/** - * - * @author baltzell - */ -public class EngineMultiProcessor { - - int threads; - DataSource reader; - HipoDataSync writer; - ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); - EngineProcessor engines = new EngineProcessor(); - ExecutorService executor; - - public EngineMultiProcessor(int threads) { - this.threads = threads; - executor = Executors.newFixedThreadPool(threads); - } - - public void open(String input, String output, int events, int skip) { - writer = new HipoDataSync(); - writer.setCompressionType(2); - if (input.endsWith(".hipo")) reader = new HipoDataSource(); - else reader = new EvioSource(); - reader.open(input); - writer.open(output); - engines.updateDictionary((HipoDataSource)reader, writer); - CompletableFuture.supplyAsync(() -> { - try { read(); } catch (InterruptedException ex) {} - return true; - }, executor); - CompletableFuture.supplyAsync(() -> { - try { write(); } catch (InterruptedException ex) {} - return true; - }, executor); - try { Thread.sleep(5000); } catch (InterruptedException ex) {} - for (int i=0; i { - process(); - return true; - }, executor); - } - while (!readQueue.isEmpty() || !writeQueue.isEmpty() || !procQueue.isEmpty()) { - try { Thread.sleep(1000); } catch (InterruptedException ex) {} - } - executor.shutdownNow(); - writer.close(); - reader.close(); - } - - void read() throws InterruptedException { - while (reader.hasEvent()) { - if (readQueue.size() > 10*threads) Thread.sleep(100); - else readQueue.offer(reader.getNextEvent()); - } - } - - void process() { - while (!readQueue.isEmpty()) { - DataEvent event = readQueue.poll(); - procQueue.offer(event); - engines.processEvent(event); - procQueue.remove(event); - writeQueue.offer(event); - } - } - - void write() throws InterruptedException { - while (true) { - if (!writeQueue.isEmpty()) - writer.writeEvent(writeQueue.poll()); - else Thread.sleep(1000); - } - } - - public static void main(String[] args) { - - 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"); - parser.addOption("-y","0","yaml file"); - parser.addOption("-u","true","update dictionary from writer ? "); - parser.addOption("-S",null,"schema directory"); - parser.addOption("-B",null,"background file"); - parser.addOption("-P",null,"preload file for post-processing"); - parser.addOption("-R","0","rebuild scalers"); - parser.addOption("-H","0","restream helicity"); - parser.addOption("-t","1","threads"); - parser.parse(args); - parser.syncLogLevel(Logger.getLogger(EngineMultiProcessor.class.getPackage().getName())); - - EngineMultiProcessor multi = new EngineMultiProcessor(parser.getOption("-t").intValue()); - - if (parser.getOption("-u").stringValue().contains("false")) - multi.engines.updateDictionary = false; - - // service list from YAML: - if(!parser.getOption("-y").stringValue().equals("0")) { - ClaraYaml yaml = new ClaraYaml(parser.getOption("-y").stringValue()); - if (yaml.schemaDirectory() != null) { - multi.engines.setBanksToKeep(yaml.schemaDirectory()); - } - for (JSONObject service : yaml.services()) { - JSONObject cfg = yaml.filter(service.getString("name")); - if (cfg.length() > 0) { - multi.engines.addEngine(service.getString("name"),service.getString("class"),cfg.toString()); - } else { - multi.engines.addEngine(service.getString("name"),service.getString("class")); - } - } - } - // built-in service list: - else if (parser.getOption("-c").intValue() > 0) { - if (parser.getOption("-c").intValue() > 2) { - multi.engines.initCaloDebug(); - } else if(parser.getOption("-c").intValue() == 2) { - multi.engines.initAll(); - } else { - multi.engines.initDefault(); - } - } - // user-defined service list: - else { - for(String engine : parser.getInputList()){ - System.out.println("Adding reconstruction engine " + engine); - multi.engines.addEngine(engine); - } - } - - // command-line schema overrides YAML: - if (parser.getOption("-S").stringValue() != null) - multi.engines.setBanksToKeep(parser.getOption("-S").stringValue()); - - // command-line filename for background merging overrides YAML: - if (parser.getOption("-B").stringValue() != null) - multi.engines.setBackgroundFiles(parser.getOption("-B").stringValue()); - - // command-line filename for post-processing overrides YAML: - if (parser.getOption("-P").stringValue() != null) { - multi.engines.setPreloadFiles(parser.getOption("-P").stringValue(), - parser.getOption("-H").intValue()!=0, - parser.getOption("-R").intValue()!=0); - } - - multi.open(parser.getOption("-i").stringValue(), - parser.getOption("-o").stringValue(), - parser.getOption("-s").intValue(), - parser.getOption("-n").intValue()); - } - -} 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 ca1e3580a5..4603fde2ed 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 @@ -17,8 +17,13 @@ import org.jlab.clara.engine.EngineData; import org.jlab.clara.engine.EngineDataType; import java.util.Arrays; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import org.jlab.coda.jevio.EvioException; import org.jlab.detector.decode.CLASDecoder4; +import org.jlab.io.base.DataSource; import org.jlab.io.evio.EvioDataEvent; import org.jlab.io.evio.EvioSource; import org.jlab.io.hipo.HipoDataEvent; @@ -380,6 +385,7 @@ public static void main(String[] args){ 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("-t","1","number of threads"); parser.addOption("-y","0","yaml file"); parser.addOption("-u","true","update dictionary from writer ? "); parser.addOption("-S",null,"schema directory"); @@ -396,7 +402,7 @@ public static void main(String[] args){ String inputFile = parser.getOption("-i").stringValue(); String outputFile = parser.getOption("-o").stringValue(); - EngineProcessor proc = new EngineProcessor(); + EngineMultiProcessor proc = new EngineMultiProcessor(parser.getOption("-t").intValue()); int config = parser.getOption("-c").intValue(); int nskip = parser.getOption("-s").intValue(); @@ -451,7 +457,80 @@ else if (config>0){ parser.getOption("-R").intValue()!=0); } - proc.processFile(inputFile,outputFile,nskip,nevents); + proc.open(inputFile,outputFile,nskip,nevents); + //proc.processFile(inputFile,outputFile,nskip,nevents); } + public static class EngineMultiProcessor extends EngineProcessor { + + int threads; + DataSource reader; + HipoDataSync writer; + ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); + ExecutorService executor; + + public EngineMultiProcessor(int threads) { + super(); + this.threads = threads; + executor = Executors.newFixedThreadPool(threads); + } + + public void open(String input, String output, int events, int skip) { + writer = new HipoDataSync(); + writer.setCompressionType(2); + if (input.endsWith(".hipo")) reader = new HipoDataSource(); + else reader = new EvioSource(); + reader.open(input); + writer.open(output); + updateDictionary((HipoDataSource)reader, writer); + CompletableFuture.supplyAsync(() -> { + try { read(); } catch (InterruptedException ex) {} + return true; + }, executor); + CompletableFuture.supplyAsync(() -> { + try { write(); } catch (InterruptedException ex) {} + return true; + }, executor); + try { Thread.sleep(5000); } catch (InterruptedException ex) {} + for (int i=0; i { + process(); + return true; + }, executor); + } + while (!readQueue.isEmpty() || !writeQueue.isEmpty() || !procQueue.isEmpty()) { + try { Thread.sleep(1000); } catch (InterruptedException ex) {} + } + executor.shutdownNow(); + writer.close(); + reader.close(); + } + + void read() throws InterruptedException { + while (reader.hasEvent()) { + if (readQueue.size() > 10*threads) Thread.sleep(100); + else readQueue.offer(reader.getNextEvent()); + } + } + + void process() { + while (!readQueue.isEmpty()) { + DataEvent event = readQueue.poll(); + procQueue.offer(event); + processEvent(event); + procQueue.remove(event); + writeQueue.offer(event); + } + } + + void write() throws InterruptedException { + while (true) { + if (!writeQueue.isEmpty()) + writer.writeEvent(writeQueue.poll()); + else Thread.sleep(1000); + } + } + } } From c4994e90f39ab412bba14ca8a681a62dd783d4ea Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 14:11:27 -0400 Subject: [PATCH 07/27] cleanup --- validation/advanced-tests/run-eb-tests.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/validation/advanced-tests/run-eb-tests.sh b/validation/advanced-tests/run-eb-tests.sh index 656d703b4a..6bcd2522be 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-mutil -t 24 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 +../../coatjava/bin/recon-util -t 2 -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 00b135cbb79fec6cccc51a169a9682860008cf76 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 14:12:55 -0400 Subject: [PATCH 08/27] cleanup --- .../src/main/java/org/jlab/clas/reco/EngineProcessor.java | 6 +++--- 1 file changed, 3 insertions(+), 3 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 4603fde2ed..ceaaab9dd1 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 @@ -457,8 +457,7 @@ else if (config>0){ parser.getOption("-R").intValue()!=0); } - proc.open(inputFile,outputFile,nskip,nevents); - //proc.processFile(inputFile,outputFile,nskip,nevents); + proc.processFile(inputFile,outputFile,nskip,nevents); } public static class EngineMultiProcessor extends EngineProcessor { @@ -477,7 +476,8 @@ public EngineMultiProcessor(int threads) { executor = Executors.newFixedThreadPool(threads); } - public void open(String input, String output, int events, int skip) { + @Override + public void processFile(String input, String output, int events, int skip) { writer = new HipoDataSync(); writer.setCompressionType(2); if (input.endsWith(".hipo")) reader = new HipoDataSource(); From ebc0b85a438eb01445cbeeaac0a409f13ed05bd0 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 15:54:24 -0400 Subject: [PATCH 09/27] cleanup --- bin/recon-mutil | 13 ++ .../jlab/clas/reco/EngineMultiProcessor.java | 189 ++++++++++++++++++ .../org/jlab/clas/reco/EngineProcessor.java | 89 +-------- validation/advanced-tests/run-eb-tests.sh | 3 +- 4 files changed, 210 insertions(+), 84 deletions(-) create mode 100755 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 100755 index 0000000000..f9be37b69a --- /dev/null +++ b/bin/recon-mutil @@ -0,0 +1,13 @@ +#!/bin/bash + +. `dirname $0`/../libexec/env.sh + +split_cli $@ + +export MALLOC_ARENA_MAX=1 + +java ${JAVA_OPTS-} -Xmx1536m -Xms1024m -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..795342fbf8 --- /dev/null +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineMultiProcessor.java @@ -0,0 +1,189 @@ +package org.jlab.clas.reco; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.Executors; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.logging.Logger; +import org.jlab.coda.jevio.EvioException; +import org.jlab.io.base.DataEvent; +import org.jlab.io.base.DataSource; +import org.jlab.io.evio.EvioSource; +import org.jlab.io.hipo.HipoDataSource; +import org.jlab.io.hipo.HipoDataSync; +import org.jlab.utils.ClaraYaml; +import org.jlab.utils.options.OptionParser; +import org.json.JSONObject; + +/** + * + * @author baltzell + */ +public class EngineMultiProcessor extends EngineProcessor { + + DataSource reader; + HipoDataSync writer; + ThreadPoolExecutor executor; + ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); + + public EngineMultiProcessor(int threads) { + super(); + executor = (ThreadPoolExecutor)Executors.newFixedThreadPool(threads + 2); + } + + void open(String input, String output) { + writer = new HipoDataSync(); + writer.setCompressionType(2); + if (input.endsWith(".hipo")) reader = new HipoDataSource(); + else reader = new EvioSource(); + reader.open(input); + writer.open(output); + updateDictionary((HipoDataSource)reader, writer); + } + + void read() throws InterruptedException, EvioException { + while (reader.hasEvent()) { + if (readQueue.size() > 10*executor.getMaximumPoolSize()) Thread.sleep(100); + else readQueue.offer(reader.getNextEvent()); + } + } + + void process() { + while (!readQueue.isEmpty()) { + DataEvent event = readQueue.poll(); + procQueue.offer(event); + processEvent(event); + procQueue.remove(event); + writeQueue.offer(event); + } + } + + void write() throws InterruptedException { + while (true) { + if (writeQueue.isEmpty()) Thread.sleep(1000); + else writer.writeEvent(writeQueue.poll()); + } + } + + @Override + public void processFile(String input, String output, int events, int skip) { + + // create new reader and writer: + open(input, output); + + // start reader thread: + CompletableFuture.supplyAsync(() -> { + try { read(); } catch (InterruptedException | EvioException ex) {} + return true; + }, executor); + + // start writer thread: + CompletableFuture.supplyAsync(() -> { + try { write(); } catch (InterruptedException ex) {} + return true; + }, executor); + + // prime the queue: + while (readQueue.size() < 2*executor.getMaximumPoolSize()) {} + + // start processor threads: + for (int i=0; i { + process(); + return true; + }, executor); + } + + // wait for queues to be empty: + while (!readQueue.isEmpty() || !writeQueue.isEmpty() || !procQueue.isEmpty()) { + try { Thread.sleep(1000); } catch (InterruptedException ex) {} + } + + // close shop: + executor.shutdownNow(); + writer.close(); + reader.close(); + } + + public static void main(String[] args) { + + OptionParser parser = new OptionParser("recon-util"); + 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("-t","1","number of threads"); + parser.addOption("-y","0","yaml file"); + parser.addOption("-u","true","update dictionary from writer ? "); + parser.addOption("-S",null,"schema directory"); + parser.addOption("-B",null,"background file"); + parser.addOption("-P",null,"preload file for post-processing"); + parser.addOption("-R","0","rebuild scalers"); + parser.addOption("-H","0","restream helicity"); + parser.setRequiresInputList(false); + parser.parse(args); + parser.syncLogLevel(Logger.getLogger(EngineProcessor.class.getPackage().getName())); + + EngineMultiProcessor proc = new EngineMultiProcessor(parser.getOption("-t").intValue()); + + if (parser.getOption("-u").stringValue().contains("false")) { + proc.updateDictionary = false; + } + + // services from YAML: + if (!parser.getOption("-y").stringValue().equals("0")) { + ClaraYaml yaml = new ClaraYaml(parser.getOption("-y").stringValue()); + if (yaml.schemaDirectory() != null) { + proc.setBanksToKeep(yaml.schemaDirectory()); + } + for (JSONObject service : yaml.services()) { + JSONObject cfg = yaml.filter(service.getString("name")); + if (cfg.length() > 0) { + proc.addEngine(service.getString("name"),service.getString("class"),cfg.toString()); + } else { + proc.addEngine(service.getString("name"),service.getString("class")); + } + } + } + // services from builtin configurations: + else if (parser.getOption("-c").intValue() > 0){ + if(parser.getOption("-c").intValue() > 2){ + proc.initCaloDebug(); + } else if(parser.getOption("-c").intValue() == 2){ + proc.initAll(); + } else { + proc.initDefault(); + } + } + // user-defined services: + else { + for(String engine : parser.getInputList()){ + System.out.println("Adding reconstruction engine " + engine); + proc.addEngine(engine); + } + } + + // command-line schema overrides YAML: + if (parser.getOption("-S").stringValue() != null) + proc.setBanksToKeep(parser.getOption("-S").stringValue()); + + // command-line filename for background merging overrides YAML: + if (parser.getOption("-B").stringValue() != null) + proc.setBackgroundFiles(parser.getOption("-B").stringValue()); + + // command-line filename for post-processing overrides YAML: + if (parser.getOption("-P").stringValue() != null) { + proc.setPreloadFiles(parser.getOption("-P").stringValue(), + parser.getOption("-H").intValue()!=0, + parser.getOption("-R").intValue()!=0); + } + + proc.processFile(parser.getOption("-i").stringValue(), + parser.getOption("-o").stringValue(), + parser.getOption("-n").intValue(), + parser.getOption("-s").intValue()); + } +} 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 ceaaab9dd1..9d140e61ee 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 @@ -17,13 +17,8 @@ import org.jlab.clara.engine.EngineData; import org.jlab.clara.engine.EngineDataType; import java.util.Arrays; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import org.jlab.coda.jevio.EvioException; import org.jlab.detector.decode.CLASDecoder4; -import org.jlab.io.base.DataSource; import org.jlab.io.evio.EvioDataEvent; import org.jlab.io.evio.EvioSource; import org.jlab.io.hipo.HipoDataEvent; @@ -48,8 +43,10 @@ public class EngineProcessor { protected boolean updateDictionary = true; - private CLASDecoder4 decoder = new CLASDecoder4(); + protected CLASDecoder4 decoder = new CLASDecoder4(); + int eventsRead = 0; + public EngineProcessor(){} private ReconstructionEngine findEngine(String clazz) { @@ -309,7 +306,7 @@ public void processEvent(DataEvent event, HipoDataSync writer) { public void processFile(HipoDataSource reader, HipoDataSync writer, int skipEvents, int maxEvents) { if (updateDictionary==true) updateDictionary(reader, writer); ProgressPrintout progress = new ProgressPrintout(); - int eventsRead = 0; + eventsRead = 0; while (reader.hasEvent()) { DataEvent event = reader.getNextEvent(); eventsRead++; @@ -322,7 +319,7 @@ public void processFile(HipoDataSource reader, HipoDataSync writer, int skipEven public void processFile(EvioSource reader, HipoDataSync writer, int skipEvents, int maxEvents) { ProgressPrintout progress = new ProgressPrintout(); - int eventsRead = 0; + eventsRead = 0; while (reader.hasEvent()) { eventsRead++; try { @@ -385,7 +382,6 @@ public static void main(String[] args){ 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("-t","1","number of threads"); parser.addOption("-y","0","yaml file"); parser.addOption("-u","true","update dictionary from writer ? "); parser.addOption("-S",null,"schema directory"); @@ -402,7 +398,7 @@ public static void main(String[] args){ String inputFile = parser.getOption("-i").stringValue(); String outputFile = parser.getOption("-o").stringValue(); - EngineMultiProcessor proc = new EngineMultiProcessor(parser.getOption("-t").intValue()); + EngineProcessor proc = new EngineProcessor(); int config = parser.getOption("-c").intValue(); int nskip = parser.getOption("-s").intValue(); @@ -460,77 +456,4 @@ else if (config>0){ proc.processFile(inputFile,outputFile,nskip,nevents); } - public static class EngineMultiProcessor extends EngineProcessor { - - int threads; - DataSource reader; - HipoDataSync writer; - ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); - ExecutorService executor; - - public EngineMultiProcessor(int threads) { - super(); - this.threads = threads; - executor = Executors.newFixedThreadPool(threads); - } - - @Override - public void processFile(String input, String output, int events, int skip) { - writer = new HipoDataSync(); - writer.setCompressionType(2); - if (input.endsWith(".hipo")) reader = new HipoDataSource(); - else reader = new EvioSource(); - reader.open(input); - writer.open(output); - updateDictionary((HipoDataSource)reader, writer); - CompletableFuture.supplyAsync(() -> { - try { read(); } catch (InterruptedException ex) {} - return true; - }, executor); - CompletableFuture.supplyAsync(() -> { - try { write(); } catch (InterruptedException ex) {} - return true; - }, executor); - try { Thread.sleep(5000); } catch (InterruptedException ex) {} - for (int i=0; i { - process(); - return true; - }, executor); - } - while (!readQueue.isEmpty() || !writeQueue.isEmpty() || !procQueue.isEmpty()) { - try { Thread.sleep(1000); } catch (InterruptedException ex) {} - } - executor.shutdownNow(); - writer.close(); - reader.close(); - } - - void read() throws InterruptedException { - while (reader.hasEvent()) { - if (readQueue.size() > 10*threads) Thread.sleep(100); - else readQueue.offer(reader.getNextEvent()); - } - } - - void process() { - while (!readQueue.isEmpty()) { - DataEvent event = readQueue.poll(); - procQueue.offer(event); - processEvent(event); - procQueue.remove(event); - writeQueue.offer(event); - } - } - - void write() throws InterruptedException { - while (true) { - if (!writeQueue.isEmpty()) - writer.writeEvent(writeQueue.poll()); - else Thread.sleep(1000); - } - } - } } diff --git a/validation/advanced-tests/run-eb-tests.sh b/validation/advanced-tests/run-eb-tests.sh index 6bcd2522be..fc758ef487 100755 --- a/validation/advanced-tests/run-eb-tests.sh +++ b/validation/advanced-tests/run-eb-tests.sh @@ -49,7 +49,8 @@ if [ $? != 0 ] ; then echo "EBTwoTrackTest compilation failure" ; exit 1 ; fi # run reconstruction: rm -f out_${stub}.hipo -../../coatjava/bin/recon-util -t 2 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 +#../../coatjava/bin/recon-util -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 -- -Xmx25000m -Xms3000m +../../coatjava/bin/recon-mutil -t 12 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 -- -Xmx25000m -Xms3000m # run EB tests: java -Xmx1536m -Xms1024m -cp $classPath -DINPUTFILE=out_${stub}.hipo eb.EBTwoTrackTest From 834557e825468f591289000f57ac4a854d5b68a5 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 15:55:58 -0400 Subject: [PATCH 10/27] relax memory --- bin/recon-mutil | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/bin/recon-mutil b/bin/recon-mutil index f9be37b69a..622579382b 100755 --- a/bin/recon-mutil +++ b/bin/recon-mutil @@ -6,7 +6,7 @@ split_cli $@ export MALLOC_ARENA_MAX=1 -java ${JAVA_OPTS-} -Xmx1536m -Xms1024m -XX:+UseSerialGC ${jvm_options[@]} \ +java ${JAVA_OPTS-} -Xms2048m -XX:+UseSerialGC ${jvm_options[@]} \ -cp ${COATJAVA_CLASSPATH:-''} \ org.jlab.clas.reco.EngineMultiProcessor \ ${class_options[@]} From 35237abf0647784955f2b943df1e48c2b37a2bdf Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 16:29:37 -0400 Subject: [PATCH 11/27] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 51 ++++++++++++++----- 1 file changed, 38 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 795342fbf8..d47e90a5ee 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 @@ -1,5 +1,6 @@ package org.jlab.clas.reco; +import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.Executors; @@ -20,7 +21,11 @@ * @author baltzell */ public class EngineMultiProcessor extends EngineProcessor { - + + int maxEvents = 0; + int skipEvents = 0; + int readEvents = 0; + List inputs; DataSource reader; HipoDataSync writer; ThreadPoolExecutor executor; @@ -32,21 +37,42 @@ public EngineMultiProcessor(int threads) { super(); executor = (ThreadPoolExecutor)Executors.newFixedThreadPool(threads + 2); } + + public EngineMultiProcessor(int threads, int events, int skip) { + super(); + this.maxEvents = events; + this.skipEvents = skip; + executor = (ThreadPoolExecutor)Executors.newFixedThreadPool(threads + 2); + } - void open(String input, String output) { + void open(String output, String... input) { + for (int i=0; i 10*executor.getMaximumPoolSize()) Thread.sleep(100); - else readQueue.offer(reader.getNextEvent()); + while (true) { + if (readEvents > 0 && maxEvents > readEvents) break; + if (reader.hasEvent()) { + if (readQueue.size() > 10*executor.getMaximumPoolSize()) Thread.sleep(100); + else { + readEvents++; + DataEvent event = reader.getNextEvent(); + if (skipEvents < 0 || readEvents > skipEvents) readQueue.offer(event); + } + } + else if (inputs.isEmpty()) break; + else open(inputs.remove(0)); } } @@ -62,16 +88,15 @@ void process() { void write() throws InterruptedException { while (true) { - if (writeQueue.isEmpty()) Thread.sleep(1000); + if (writeQueue.isEmpty()) Thread.sleep(100); else writer.writeEvent(writeQueue.poll()); } } - @Override - public void processFile(String input, String output, int events, int skip) { + public void processFiles(String output, String... input) { // create new reader and writer: - open(input, output); + open(output, input); // start reader thread: CompletableFuture.supplyAsync(() -> { @@ -86,7 +111,7 @@ public void processFile(String input, String output, int events, int skip) { }, executor); // prime the queue: - while (readQueue.size() < 2*executor.getMaximumPoolSize()) {} + while (readQueue.size() < executor.getMaximumPoolSize()) {} // start processor threads: for (int i=0; i Date: Tue, 18 Aug 2026 16:57:46 -0400 Subject: [PATCH 12/27] add decoder engine to the defaults --- .../src/main/java/org/jlab/clas/reco/EngineProcessor.java | 6 ++++-- 1 file changed, 4 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 9d140e61ee..a295453296 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 @@ -116,12 +116,13 @@ private void removeBanks(DataEvent event) { public void initDefault(){ String[] names = new String[]{ - "MAGFIELDS", + "DECO","MAGFIELDS", "DCCR","DCHB","FTOFHB","EC","HTCC","EBHB", "DCTB","FTOFTB","EBTB","VTX" }; String[] services = new String[]{ + "org.jlab.clas.reco.DecoderEngine", "org.jlab.clas.swimtools.MagFieldsEngine", "org.jlab.service.dc.DCHBClustering", "org.jlab.service.dc.DCHBPostClusterConv", @@ -142,7 +143,7 @@ public void initDefault(){ public void initAll(){ String[] names = new String[]{ - "MAGFIELDS", + "DECO","MAGFIELDS", "FTCAL", "FTHODO", "FTTRK", "FTEB", "URWT", "DCCR", "DCHB","FTOFHB","EC","RASTER", "CVTFP","CTOF","CND","BAND", @@ -152,6 +153,7 @@ public void initAll(){ }; String[] services = new String[]{ + "org.jlab.clas.reco.DecoderEngine", "org.jlab.clas.swimtools.MagFieldsEngine", "org.jlab.rec.ft.cal.FTCALEngine", "org.jlab.rec.ft.hodo.FTHODOEngine", From 0cf6bd52c91e6fc99c1544974a1a475555aca1bf Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 17:06:03 -0400 Subject: [PATCH 13/27] fix --- .../org/jlab/clas/reco/EngineMultiProcessor.java | 14 ++++++-------- 1 file changed, 6 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 d47e90a5ee..97b8085a19 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 @@ -1,6 +1,6 @@ package org.jlab.clas.reco; -import java.util.List; +import java.util.ArrayList; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.Executors; @@ -25,7 +25,7 @@ public class EngineMultiProcessor extends EngineProcessor { int maxEvents = 0; int skipEvents = 0; int readEvents = 0; - List inputs; + ArrayList inputs = new ArrayList<>(); DataSource reader; HipoDataSync writer; ThreadPoolExecutor executor; @@ -56,10 +56,10 @@ void open(String output, String... input) { void open(String input) { if (input.endsWith(".hipo")) reader = new HipoDataSource(); else reader = new EvioSource(); - reader.open(inputs.remove(0)); + reader.open(input); updateDictionary((HipoDataSource)reader, writer); } - + void read() throws InterruptedException, EvioException { while (true) { if (readEvents > 0 && maxEvents > readEvents) break; @@ -206,9 +206,7 @@ else if (parser.getOption("-c").intValue() > 0){ parser.getOption("-R").intValue()!=0); } - proc.processFile(parser.getOption("-i").stringValue(), - parser.getOption("-o").stringValue(), - parser.getOption("-n").intValue(), - parser.getOption("-s").intValue()); + proc.processFiles(parser.getOption("-o").stringValue(), + parser.getOption("-i").stringValue()); } } From 7a18991b443fe2fc7a584407607e0a2067ae2b7b Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 17:53:23 -0400 Subject: [PATCH 14/27] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 47 +++++++++---------- 1 file changed, 22 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 97b8085a19..54af0cf906 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 @@ -22,13 +22,14 @@ */ public class EngineMultiProcessor extends EngineProcessor { + DataSource reader; + HipoDataSync writer; + ThreadPoolExecutor executor; + int maxEvents = 0; int skipEvents = 0; int readEvents = 0; ArrayList inputs = new ArrayList<>(); - DataSource reader; - HipoDataSync writer; - ThreadPoolExecutor executor; ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); @@ -45,34 +46,27 @@ public EngineMultiProcessor(int threads, int events, int skip) { executor = (ThreadPoolExecutor)Executors.newFixedThreadPool(threads + 2); } - void open(String output, String... input) { - for (int i=0; i 0 && maxEvents > readEvents) break; - if (reader.hasEvent()) { - if (readQueue.size() > 10*executor.getMaximumPoolSize()) Thread.sleep(100); - else { - readEvents++; - DataEvent event = reader.getNextEvent(); - if (skipEvents < 0 || readEvents > skipEvents) readQueue.offer(event); + while (maxEvents < 1 || readEvents <= maxEvents) { + if (reader != null && reader.hasEvent()) { + if (readQueue.size() > 10*executor.getMaximumPoolSize()) { + Thread.sleep(100); + continue; } + DataEvent event = reader.getNextEvent(); + readEvents++; + if (skipEvents < 1 || readEvents > skipEvents) + readQueue.offer(event); } else if (inputs.isEmpty()) break; - else open(inputs.remove(0)); + else open(); } } @@ -81,8 +75,8 @@ void process() { DataEvent event = readQueue.poll(); procQueue.offer(event); processEvent(event); - procQueue.remove(event); writeQueue.offer(event); + procQueue.remove(event); } } @@ -96,7 +90,10 @@ void write() throws InterruptedException { public void processFiles(String output, String... input) { // create new reader and writer: - open(output, input); + for (int i=0; i { From 39ba65e695fd3fe867e0c46b6a5e820e001a1712 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 18:00:58 -0400 Subject: [PATCH 15/27] 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 54af0cf906..250792935c 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,7 +38,7 @@ public EngineMultiProcessor(int threads) { super(); executor = (ThreadPoolExecutor)Executors.newFixedThreadPool(threads + 2); } - + public EngineMultiProcessor(int threads, int events, int skip) { super(); this.maxEvents = events; @@ -54,7 +54,7 @@ void open() { } void read() throws InterruptedException, EvioException { - while (maxEvents < 1 || readEvents <= maxEvents) { + while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { if (readQueue.size() > 10*executor.getMaximumPoolSize()) { Thread.sleep(100); @@ -128,7 +128,7 @@ public void processFiles(String output, String... input) { writer.close(); reader.close(); } - + public static void main(String[] args) { OptionParser parser = new OptionParser("recon-util"); From c86af455ed175eb356b64e6582a9ef6e5ef5d41c Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 18:26:28 -0400 Subject: [PATCH 16/27] fix --- .../jlab/clas/reco/EngineMultiProcessor.java | 65 +++++++++++-------- 1 file changed, 37 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 250792935c..7b05ae14d3 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 @@ -3,9 +3,8 @@ import java.util.ArrayList; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.Executors; -import java.util.concurrent.ThreadPoolExecutor; import java.util.logging.Logger; +import org.json.JSONObject; import org.jlab.coda.jevio.EvioException; import org.jlab.io.base.DataEvent; import org.jlab.io.base.DataSource; @@ -14,7 +13,6 @@ import org.jlab.io.hipo.HipoDataSync; import org.jlab.utils.ClaraYaml; import org.jlab.utils.options.OptionParser; -import org.json.JSONObject; /** * @@ -24,8 +22,9 @@ public class EngineMultiProcessor extends EngineProcessor { DataSource reader; HipoDataSync writer; - ThreadPoolExecutor executor; + CompletableFuture readerThread; + int threads = 4; int maxEvents = 0; int skipEvents = 0; int readEvents = 0; @@ -36,14 +35,12 @@ public class EngineMultiProcessor extends EngineProcessor { public EngineMultiProcessor(int threads) { super(); - executor = (ThreadPoolExecutor)Executors.newFixedThreadPool(threads + 2); } public EngineMultiProcessor(int threads, int events, int skip) { super(); this.maxEvents = events; this.skipEvents = skip; - executor = (ThreadPoolExecutor)Executors.newFixedThreadPool(threads + 2); } void open() { @@ -56,7 +53,7 @@ void open() { void read() throws InterruptedException, EvioException { while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { - if (readQueue.size() > 10*executor.getMaximumPoolSize()) { + if (readQueue.size() > 10*threads) { Thread.sleep(100); continue; } @@ -70,19 +67,32 @@ void read() throws InterruptedException, EvioException { } } - void process() { - while (!readQueue.isEmpty()) { - DataEvent event = readQueue.poll(); - procQueue.offer(event); - processEvent(event); - writeQueue.offer(event); - procQueue.remove(event); + void process() throws InterruptedException { + while (true) { + if (readQueue.isEmpty()) { + if (readerThread.isDone()) break; + Thread.sleep(100); + } + else { + DataEvent event = readQueue.poll(); + procQueue.offer(event); + processEvent(event); + writeQueue.offer(event); + procQueue.remove(event); + } } } void write() throws InterruptedException { while (true) { - if (writeQueue.isEmpty()) Thread.sleep(100); + if (writeQueue.isEmpty()) { + if (readerThread.isDone()) { + if (readQueue.isEmpty() && procQueue.isEmpty()) { + break; + } + } + Thread.sleep(100); + } else writer.writeEvent(writeQueue.poll()); } } @@ -96,35 +106,34 @@ public void processFiles(String output, String... input) { writer.open(output); // start reader thread: - CompletableFuture.supplyAsync(() -> { + readerThread = CompletableFuture.supplyAsync(() -> { try { read(); } catch (InterruptedException | EvioException ex) {} return true; - }, executor); - + }); + // start writer thread: - CompletableFuture.supplyAsync(() -> { + CompletableFuture writerThread = CompletableFuture.supplyAsync(() -> { try { write(); } catch (InterruptedException ex) {} return true; - }, executor); + }); // prime the queue: - while (readQueue.size() < executor.getMaximumPoolSize()) {} + while (readQueue.size() < threads) {} // start processor threads: - for (int i=0; i { - process(); + try { process(); } catch (InterruptedException ex) {} return true; - }, executor); + }); } - // wait for queues to be empty: - while (!readQueue.isEmpty() || !writeQueue.isEmpty() || !procQueue.isEmpty()) { - try { Thread.sleep(1000); } catch (InterruptedException ex) {} + // wait for finish: + while (!writerThread.isDone()) { + try { Thread.sleep(100); } catch (InterruptedException ex) {} } // close shop: - executor.shutdownNow(); writer.close(); reader.close(); } From e07a73f8291bee9fa302eebe2fa1b3ceb0029b0b Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Tue, 18 Aug 2026 18:34:43 -0400 Subject: [PATCH 17/27] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 17 ++++++++--------- 1 file changed, 8 insertions(+), 9 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 7b05ae14d3..b29322985f 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 @@ -24,7 +24,7 @@ public class EngineMultiProcessor extends EngineProcessor { HipoDataSync writer; CompletableFuture readerThread; - int threads = 4; + int threads; int maxEvents = 0; int skipEvents = 0; int readEvents = 0; @@ -35,10 +35,12 @@ 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.maxEvents = events; this.skipEvents = skip; } @@ -88,6 +90,7 @@ void write() throws InterruptedException { if (writeQueue.isEmpty()) { if (readerThread.isDone()) { if (readQueue.isEmpty() && procQueue.isEmpty()) { + writer.close(); break; } } @@ -97,7 +100,7 @@ void write() throws InterruptedException { } } - public void processFiles(String output, String... input) { + public void process(String output, String... input) { // create new reader and writer: for (int i=0; i 0){ parser.getOption("-R").intValue()!=0); } - proc.processFiles(parser.getOption("-o").stringValue(), - parser.getOption("-i").stringValue()); + proc.process(parser.getOption("-o").stringValue(), parser.getOption("-i").stringValue()); } } From 1bb25086a4dec9656d9d62f1d5eeae16a4913fd6 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Wed, 19 Aug 2026 12:30:08 -0400 Subject: [PATCH 18/27] priming no longer needed --- .../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 b29322985f..04953648a2 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 @@ -102,7 +102,7 @@ void write() throws InterruptedException { public void process(String output, String... input) { - // create new reader and writer: + // add inputs and create writer: for (int i=0; i { From 5ed42a8a2a32878266230e7019fa8aba52afeec4 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Wed, 19 Aug 2026 12:35:41 -0400 Subject: [PATCH 19/27] add progress printout --- .../java/org/jlab/clas/reco/EngineMultiProcessor.java | 8 +++++++- 1 file changed, 7 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 04953648a2..7316cc043b 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 @@ -12,6 +12,7 @@ import org.jlab.io.hipo.HipoDataSource; import org.jlab.io.hipo.HipoDataSync; import org.jlab.utils.ClaraYaml; +import org.jlab.utils.benchmark.ProgressPrintout; import org.jlab.utils.options.OptionParser; /** @@ -29,6 +30,7 @@ public class EngineMultiProcessor extends EngineProcessor { int skipEvents = 0; int readEvents = 0; ArrayList inputs = new ArrayList<>(); + ProgressPrintout progress = new ProgressPrintout(); ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); @@ -90,13 +92,17 @@ void write() throws InterruptedException { if (writeQueue.isEmpty()) { if (readerThread.isDone()) { if (readQueue.isEmpty() && procQueue.isEmpty()) { + progress.showStatus(); writer.close(); break; } } Thread.sleep(100); } - else writer.writeEvent(writeQueue.poll()); + else { + writer.writeEvent(writeQueue.poll()); + progress.updateStatus(); + } } } From 2f3cfd966882bd71478b93811ea286d59b329e3c Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Wed, 19 Aug 2026 13:53:41 -0400 Subject: [PATCH 20/27] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 42 +++++++++---------- 1 file changed, 20 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 7316cc043b..63ed313c9c 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 @@ -47,13 +47,6 @@ public EngineMultiProcessor(int threads, int events, int skip) { this.skipEvents = skip; } - void open() { - if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); - else reader = new EvioSource(); - reader.open(inputs.remove(0)); - updateDictionary((HipoDataSource)reader, writer); - } - void read() throws InterruptedException, EvioException { while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { @@ -67,22 +60,11 @@ void read() throws InterruptedException, EvioException { readQueue.offer(event); } else if (inputs.isEmpty()) break; - else open(); - } - } - - void process() throws InterruptedException { - while (true) { - if (readQueue.isEmpty()) { - if (readerThread.isDone()) break; - Thread.sleep(100); - } else { - DataEvent event = readQueue.poll(); - procQueue.offer(event); - processEvent(event); - writeQueue.offer(event); - procQueue.remove(event); + if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); + else reader = new EvioSource(); + reader.open(inputs.remove(0)); + updateDictionary((HipoDataSource)reader, writer); } } } @@ -106,6 +88,22 @@ void write() throws InterruptedException { } } + void process() throws InterruptedException { + while (true) { + if (readQueue.isEmpty()) { + if (readerThread.isDone()) break; + Thread.sleep(100); + } + else { + DataEvent event = readQueue.poll(); + procQueue.offer(event); + processEvent(event); + writeQueue.offer(event); + procQueue.remove(event); + } + } + } + public void process(String output, String... input) { // add inputs and create writer: From 21166183b10c31eb6a252068bb73e733172eda72 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Wed, 19 Aug 2026 15:12:19 -0400 Subject: [PATCH 21/27] break it up --- .../jlab/clas/reco/EngineMultiProcessor.java | 108 +++--------------- .../org/jlab/clas/reco/EngineProcessor.java | 103 +++++++++-------- validation/advanced-tests/run-eb-tests.sh | 3 +- 3 files changed, 77 insertions(+), 137 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 63ed313c9c..0566396ff2 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 @@ -1,17 +1,15 @@ package org.jlab.clas.reco; import java.util.ArrayList; +import java.util.Arrays; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.logging.Logger; -import org.json.JSONObject; import org.jlab.coda.jevio.EvioException; import org.jlab.io.base.DataEvent; import org.jlab.io.base.DataSource; import org.jlab.io.evio.EvioSource; import org.jlab.io.hipo.HipoDataSource; import org.jlab.io.hipo.HipoDataSync; -import org.jlab.utils.ClaraYaml; import org.jlab.utils.benchmark.ProgressPrintout; import org.jlab.utils.options.OptionParser; @@ -35,6 +33,13 @@ public class EngineMultiProcessor extends EngineProcessor { ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); + public EngineMultiProcessor(OptionParser parser) { + super(parser); + threads = parser.getOption("-t").intValue(); + maxEvents = parser.getOption("-n").intValue(); + skipEvents = parser.getOption("-s").intValue(); + } + public EngineMultiProcessor(int threads) { super(); this.threads = threads; @@ -47,7 +52,8 @@ public EngineMultiProcessor(int threads, int events, int skip) { this.skipEvents = skip; } - void read() throws InterruptedException, EvioException { + void read(String... input) throws InterruptedException, EvioException { + inputs.addAll(Arrays.asList(input)); while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { if (readQueue.size() > 10*threads) { @@ -69,7 +75,10 @@ void read() throws InterruptedException, EvioException { } } - void write() throws InterruptedException { + void write(String output) throws InterruptedException { + writer = new HipoDataSync(); + writer.setCompressionType(2); + writer.open(output); while (true) { if (writeQueue.isEmpty()) { if (readerThread.isDone()) { @@ -105,25 +114,16 @@ void process() throws InterruptedException { } public void process(String output, String... input) { - - // add inputs and create writer: - for (int i=0; i { - try { read(); } catch (InterruptedException | EvioException ex) {} + try { read(input); } catch (InterruptedException | EvioException ex) {} return true; }); - // start writer thread: CompletableFuture writerThread = CompletableFuture.supplyAsync(() -> { - try { write(); } catch (InterruptedException ex) {} + try { write(output); } catch (InterruptedException ex) {} return true; }); - // start processor threads: for (int i=0; i { @@ -131,7 +131,6 @@ public void process(String output, String... input) { return true; }); } - // wait for finish: while (!writerThread.isDone()) { try { Thread.sleep(100); } catch (InterruptedException ex) {} @@ -139,80 +138,11 @@ public void process(String output, String... input) { } public static void main(String[] args) { - - OptionParser parser = new OptionParser("recon-util"); - 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("-t","1","number of threads"); - parser.addOption("-y","0","yaml file"); - parser.addOption("-u","true","update dictionary from writer ? "); - parser.addOption("-S",null,"schema directory"); - parser.addOption("-B",null,"background file"); - parser.addOption("-P",null,"preload file for post-processing"); - parser.addOption("-R","0","rebuild scalers"); - parser.addOption("-H","0","restream helicity"); - parser.setRequiresInputList(false); + OptionParser parser = EngineProcessor.parser(); + parser.addOption("-t","4","number of threads"); parser.parse(args); - parser.syncLogLevel(Logger.getLogger(EngineProcessor.class.getPackage().getName())); - - EngineMultiProcessor proc = new EngineMultiProcessor(parser.getOption("-t").intValue(), - parser.getOption("-n").intValue(), parser.getOption("-s").intValue()); - - if (parser.getOption("-u").stringValue().contains("false")) { - proc.updateDictionary = false; - } - - // services from YAML: - if (!parser.getOption("-y").stringValue().equals("0")) { - ClaraYaml yaml = new ClaraYaml(parser.getOption("-y").stringValue()); - if (yaml.schemaDirectory() != null) { - proc.setBanksToKeep(yaml.schemaDirectory()); - } - for (JSONObject service : yaml.services()) { - JSONObject cfg = yaml.filter(service.getString("name")); - if (cfg.length() > 0) { - proc.addEngine(service.getString("name"),service.getString("class"),cfg.toString()); - } else { - proc.addEngine(service.getString("name"),service.getString("class")); - } - } - } - // services from builtin configurations: - else if (parser.getOption("-c").intValue() > 0){ - if(parser.getOption("-c").intValue() > 2){ - proc.initCaloDebug(); - } else if(parser.getOption("-c").intValue() == 2){ - proc.initAll(); - } else { - proc.initDefault(); - } - } - // user-defined services: - else { - for(String engine : parser.getInputList()){ - System.out.println("Adding reconstruction engine " + engine); - proc.addEngine(engine); - } - } - - // command-line schema overrides YAML: - if (parser.getOption("-S").stringValue() != null) - proc.setBanksToKeep(parser.getOption("-S").stringValue()); - - // command-line filename for background merging overrides YAML: - if (parser.getOption("-B").stringValue() != null) - proc.setBackgroundFiles(parser.getOption("-B").stringValue()); - - // command-line filename for post-processing overrides YAML: - if (parser.getOption("-P").stringValue() != null) { - proc.setPreloadFiles(parser.getOption("-P").stringValue(), - parser.getOption("-H").intValue()!=0, - parser.getOption("-R").intValue()!=0); - } + 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 a295453296..28d3ed55a2 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 @@ -49,6 +49,10 @@ public class EngineProcessor { public EngineProcessor(){} + public EngineProcessor(OptionParser parser) { + init(parser); + } + private ReconstructionEngine findEngine(String clazz) { for (String k : processorEngines.keySet()) { if (processorEngines.get(k).getClass().getName().equals(clazz)) { @@ -58,6 +62,21 @@ private ReconstructionEngine findEngine(String clazz) { return null; } + protected void parseYaml(String filename) { + ClaraYaml yaml = new ClaraYaml(filename); + if (yaml.schemaDirectory() != null) { + setBanksToKeep(yaml.schemaDirectory()); + } + for (JSONObject service : yaml.services()) { + JSONObject cfg = yaml.filter(service.getString("name")); + if (cfg.length() > 0) { + addEngine(service.getString("name"),service.getString("class"),cfg.toString()); + } else { + addEngine(service.getString("name"),service.getString("class")); + } + } + } + protected void setBackgroundFiles(String filenames) { if (findEngine(ENGINE_CLASS_BG) == null) { LOGGER.info("Adding BackgroundEngine for -B option."); @@ -270,7 +289,7 @@ public void addEngine(String clazz) { /** * Initialize all the engines in the chain. */ - public void init(){ + public final void init(){ System.out.println("\n\n\n "); for(Map.Entry entry : this.processorEngines.entrySet()){ System.out.println(String.format(" >>>>>> (*) initializing : %8s : %s",entry.getKey(), @@ -376,11 +395,16 @@ public void show(){ } } - public static void main(String[] args){ + protected final void init(int config) { + if (config > 2) initCaloDebug(); + else if(config == 2) initAll(); + else initDefault(); + } + + protected static OptionParser parser() { 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"); @@ -391,71 +415,58 @@ public static void main(String[] args){ 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; + } - parser.parse(args); + protected final void init(OptionParser parser) { parser.syncLogLevel(LOGGER); + if (parser.getOption("-u").stringValue().contains("false")) + updateDictionary = false; - List services = parser.getInputList(); - - String inputFile = parser.getOption("-i").stringValue(); - String outputFile = parser.getOption("-o").stringValue(); - - EngineProcessor proc = new EngineProcessor(); - - int config = parser.getOption("-c").intValue(); - int nskip = parser.getOption("-s").intValue(); - int nevents = parser.getOption("-n").intValue(); - String yamlFileName = parser.getOption("-y").stringValue(); - - String update = parser.getOption("-u").stringValue(); - if(update.contains("false")==true) proc.updateDictionary = false; - - if(!yamlFileName.equals("0")) { - ClaraYaml yaml = new ClaraYaml(yamlFileName); - if (yaml.schemaDirectory() != null) { - proc.setBanksToKeep(yaml.schemaDirectory()); - } - for (JSONObject service : yaml.services()) { - JSONObject cfg = yaml.filter(service.getString("name")); - if (cfg.length() > 0) { - proc.addEngine(service.getString("name"),service.getString("class"),cfg.toString()); - } else { - proc.addEngine(service.getString("name"),service.getString("class")); - } - } + // read services and schema from YAML: + if (!parser.getOption("-y").stringValue().equals("0")) { + parseYaml(parser.getOption("-y").stringValue()); } - else if (config>0){ - if(config>2){ - proc.initCaloDebug(); - } else if(config==2){ - proc.initAll(); - } else { - proc.initDefault(); - } + // builtin configuration: + else if (parser.getOption("-c").intValue() > 0) { + init(parser.getOption("-c").intValue()); } + // user-defined services: else { - for(String engine : services){ + for(String engine : parser.getInputList()) { System.out.println("Adding reconstruction engine " + engine); - proc.addEngine(engine); + addEngine(engine); } } // command-line schema overrides YAML: if (parser.getOption("-S").stringValue() != null) - proc.setBanksToKeep(parser.getOption("-S").stringValue()); + setBanksToKeep(parser.getOption("-S").stringValue()); // command-line filename for background merging overrides YAML: if (parser.getOption("-B").stringValue() != null) - proc.setBackgroundFiles(parser.getOption("-B").stringValue()); + setBackgroundFiles(parser.getOption("-B").stringValue()); // command-line filename for post-processing overrides YAML: if (parser.getOption("-P").stringValue() != null) { - proc.setPreloadFiles(parser.getOption("-P").stringValue(), + setPreloadFiles(parser.getOption("-P").stringValue(), parser.getOption("-H").intValue()!=0, parser.getOption("-R").intValue()!=0); } + } + + public static void main(String[] args) { + + OptionParser parser = EngineProcessor.parser(); + parser.parse(args); + + EngineProcessor proc = new EngineProcessor(parser); - proc.processFile(inputFile,outputFile,nskip,nevents); + proc.processFile(parser.getOption("-i").stringValue(), + parser.getOption("-o").stringValue(), + parser.getOption("-s").intValue(), + parser.getOption("-n").intValue()); } } diff --git a/validation/advanced-tests/run-eb-tests.sh b/validation/advanced-tests/run-eb-tests.sh index fc758ef487..a996fb638e 100755 --- a/validation/advanced-tests/run-eb-tests.sh +++ b/validation/advanced-tests/run-eb-tests.sh @@ -49,8 +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 -- -Xmx25000m -Xms3000m -../../coatjava/bin/recon-mutil -t 12 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 -- -Xmx25000m -Xms3000m +../../coatjava/bin/recon-mutil -t 4 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 -- -Xmx25000m -Xms3000m # run EB tests: java -Xmx1536m -Xms1024m -cp $classPath -DINPUTFILE=out_${stub}.hipo eb.EBTwoTrackTest From a7bde538f8294ddd78ae16673cc56879473a2807 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Wed, 19 Aug 2026 18:53:22 -0400 Subject: [PATCH 22/27] rigorous queue checking --- .../jlab/clas/reco/EngineMultiProcessor.java | 32 +++++++++++-------- 1 file changed, 18 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 0566396ff2..5137179f04 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 @@ -4,6 +4,7 @@ import java.util.Arrays; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.stream.Collectors; import org.jlab.coda.jevio.EvioException; import org.jlab.io.base.DataEvent; import org.jlab.io.base.DataSource; @@ -22,6 +23,8 @@ public class EngineMultiProcessor extends EngineProcessor { DataSource reader; HipoDataSync writer; CompletableFuture readerThread; + CompletableFuture writerThread; + ArrayList procThreads = new ArrayList<>(); int threads; int maxEvents = 0; @@ -33,13 +36,6 @@ public class EngineMultiProcessor extends EngineProcessor { ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); - public EngineMultiProcessor(OptionParser parser) { - super(parser); - threads = parser.getOption("-t").intValue(); - maxEvents = parser.getOption("-n").intValue(); - skipEvents = parser.getOption("-s").intValue(); - } - public EngineMultiProcessor(int threads) { super(); this.threads = threads; @@ -52,11 +48,18 @@ public EngineMultiProcessor(int threads, int events, int skip) { this.skipEvents = skip; } + public EngineMultiProcessor(OptionParser parser) { + super(parser); + threads = parser.getOption("-t").intValue(); + maxEvents = parser.getOption("-n").intValue(); + skipEvents = parser.getOption("-s").intValue(); + } + void read(String... input) throws InterruptedException, EvioException { inputs.addAll(Arrays.asList(input)); while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { - if (readQueue.size() > 10*threads) { + if (readQueue.size() > 100*threads) { Thread.sleep(100); continue; } @@ -81,8 +84,8 @@ void write(String output) throws InterruptedException { writer.open(output); while (true) { if (writeQueue.isEmpty()) { - if (readerThread.isDone()) { - if (readQueue.isEmpty() && procQueue.isEmpty()) { + if (procThreads.stream().filter(x->!x.isDone()).collect(Collectors.toList()).isEmpty()) { + if (writeQueue.isEmpty()) { progress.showStatus(); writer.close(); break; @@ -100,7 +103,8 @@ void write(String output) throws InterruptedException { void process() throws InterruptedException { while (true) { if (readQueue.isEmpty()) { - if (readerThread.isDone()) break; + if (readerThread.isDone()) + if (readQueue.isEmpty()) break; Thread.sleep(100); } else { @@ -120,16 +124,16 @@ public void process(String output, String... input) { return true; }); // start writer thread: - CompletableFuture writerThread = CompletableFuture.supplyAsync(() -> { + writerThread = CompletableFuture.supplyAsync(() -> { try { write(output); } catch (InterruptedException ex) {} return true; }); // start processor threads: for (int i=0; i { + procThreads.add(CompletableFuture.supplyAsync(() -> { try { process(); } catch (InterruptedException ex) {} return true; - }); + })); } // wait for finish: while (!writerThread.isDone()) { From 1e6acff1619d13029433861e1cad0fa87cc70133 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Wed, 19 Aug 2026 19:16:55 -0400 Subject: [PATCH 23/27] cleanup --- .../jlab/clas/reco/EngineMultiProcessor.java | 110 ++++++++---------- 1 file changed, 49 insertions(+), 61 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 5137179f04..0132cf496c 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 @@ -5,7 +5,6 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.stream.Collectors; -import org.jlab.coda.jevio.EvioException; import org.jlab.io.base.DataEvent; import org.jlab.io.base.DataSource; import org.jlab.io.evio.EvioSource; @@ -20,22 +19,6 @@ */ public class EngineMultiProcessor extends EngineProcessor { - DataSource reader; - HipoDataSync writer; - CompletableFuture readerThread; - CompletableFuture writerThread; - ArrayList procThreads = new ArrayList<>(); - - int threads; - int maxEvents = 0; - int skipEvents = 0; - int readEvents = 0; - ArrayList inputs = new ArrayList<>(); - ProgressPrintout progress = new ProgressPrintout(); - ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); - public EngineMultiProcessor(int threads) { super(); this.threads = threads; @@ -55,18 +38,54 @@ public EngineMultiProcessor(OptionParser parser) { skipEvents = parser.getOption("-s").intValue(); } - void read(String... input) throws InterruptedException, EvioException { + public static void main(String[] args) { + OptionParser parser = EngineProcessor.parser(); + 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()); + } + + public void process(String output, String... input) { + readerThread = CompletableFuture.runAsync(() -> { read(input); }); + writerThread = CompletableFuture.runAsync(() -> { write(output); }); + for (int i=0; i { process(); })); + while (!writerThread.isDone()) + try { Thread.sleep(100); } catch (InterruptedException ex) {} + } + + DataSource reader; + HipoDataSync writer; + CompletableFuture readerThread; + CompletableFuture writerThread; + + int threads; + int maxEvents = 0; + int skipEvents = 0; + int readEvents = 0; + + ArrayList procThreads = new ArrayList<>(); + ArrayList inputs = new ArrayList<>(); + ProgressPrintout progress = new ProgressPrintout(); + ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); + + void read(String... input) { inputs.addAll(Arrays.asList(input)); while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { if (readQueue.size() > 100*threads) { - Thread.sleep(100); - continue; + try { Thread.sleep(100); } + catch (InterruptedException ex) {} + } + else { + DataEvent event = reader.getNextEvent(); + readEvents++; + if (skipEvents < 1 || readEvents > skipEvents) + readQueue.offer(event); } - DataEvent event = reader.getNextEvent(); - readEvents++; - if (skipEvents < 1 || readEvents > skipEvents) - readQueue.offer(event); } else if (inputs.isEmpty()) break; else { @@ -78,7 +97,7 @@ void read(String... input) throws InterruptedException, EvioException { } } - void write(String output) throws InterruptedException { + void write(String output) { writer = new HipoDataSync(); writer.setCompressionType(2); writer.open(output); @@ -91,7 +110,8 @@ void write(String output) throws InterruptedException { break; } } - Thread.sleep(100); + try { Thread.sleep(100); } + catch (InterruptedException ex) {} } else { writer.writeEvent(writeQueue.poll()); @@ -100,12 +120,13 @@ void write(String output) throws InterruptedException { } } - void process() throws InterruptedException { + void process() { while (true) { if (readQueue.isEmpty()) { if (readerThread.isDone()) if (readQueue.isEmpty()) break; - Thread.sleep(100); + try { Thread.sleep(100); } + catch (InterruptedException ex) {} } else { DataEvent event = readQueue.poll(); @@ -116,37 +137,4 @@ void process() throws InterruptedException { } } } - - public void process(String output, String... input) { - // start reader thread: - readerThread = CompletableFuture.supplyAsync(() -> { - try { read(input); } catch (InterruptedException | EvioException ex) {} - return true; - }); - // start writer thread: - writerThread = CompletableFuture.supplyAsync(() -> { - try { write(output); } catch (InterruptedException ex) {} - return true; - }); - // start processor threads: - for (int i=0; i { - try { process(); } catch (InterruptedException ex) {} - return true; - })); - } - // wait for finish: - while (!writerThread.isDone()) { - try { Thread.sleep(100); } catch (InterruptedException ex) {} - } - } - - public static void main(String[] args) { - OptionParser parser = EngineProcessor.parser(); - 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()); - } } From e44c7877c6628c207b1014fc61458b446fec0b41 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 20 Aug 2026 19:19:42 -0400 Subject: [PATCH 24/27] make the pool a reusable class --- .../org/jlab/clas/reco/DecoderEngine.java | 45 +++++++++++-------- 1 file changed, 27 insertions(+), 18 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/DecoderEngine.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/DecoderEngine.java index 0d10debe60..8b4c99b051 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/DecoderEngine.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/DecoderEngine.java @@ -23,13 +23,34 @@ */ public class DecoderEngine implements Engine { + public static class DecoderPool { + BlockingQueue pool; + int constantsShared = 64; + public DecoderPool(int size, String variation, String timestamp) { + pool = new ArrayBlockingQueue<>(size); + CLASDecoder d0 = null; + for (int i=0; i ED_TYPES = ClaraUtil.buildDataTypes( Clas12Types.EVIO,Clas12Types.HIPO,EngineDataType.JSON,EngineDataType.STRING); SchemaFactory schema; - BlockingQueue pool; - int constantsShared = 64; + DecoderPool pool; public DecoderEngine() { schema = new SchemaFactory(); @@ -57,22 +78,10 @@ public void destroy() {} @Override public EngineData configure(EngineData ed) { - JSONObject json = new JSONObject(ed.getData()); - pool = new ArrayBlockingQueue<>(POOL_SIZE); - CLASDecoder d0 = null; - for (int i=0; i Date: Thu, 20 Aug 2026 19:20:09 -0400 Subject: [PATCH 25/27] use threadpool --- .../jlab/clas/reco/EngineMultiProcessor.java | 71 ++++++++++++------- .../org/jlab/clas/reco/EngineProcessor.java | 37 ++++++---- 2 files changed, 70 insertions(+), 38 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 0132cf496c..629aace26e 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 @@ -1,12 +1,16 @@ package org.jlab.clas.reco; +import java.nio.ByteBuffer; +import java.nio.ByteOrder; import java.util.ArrayList; import java.util.Arrays; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.stream.Collectors; +import org.jlab.coda.jevio.EvioException; 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.HipoDataSource; import org.jlab.io.hipo.HipoDataSync; @@ -18,26 +22,26 @@ * @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.maxEvents = events; this.skipEvents = skip; } - + public EngineMultiProcessor(OptionParser parser) { super(parser); threads = parser.getOption("-t").intValue(); maxEvents = parser.getOption("-n").intValue(); skipEvents = parser.getOption("-s").intValue(); } - + public static void main(String[] args) { OptionParser parser = EngineProcessor.parser(); parser.addOption("-t","4","number of threads"); @@ -45,7 +49,7 @@ public static void main(String[] args) { EngineMultiProcessor proc = new EngineMultiProcessor(parser); proc.process(parser.getOption("-o").stringValue(), parser.getOption("-i").stringValue()); } - + public void process(String output, String... input) { readerThread = CompletableFuture.runAsync(() -> { read(input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); @@ -54,49 +58,60 @@ public void process(String output, String... input) { while (!writerThread.isDone()) try { Thread.sleep(100); } catch (InterruptedException ex) {} } - + DataSource reader; HipoDataSync writer; CompletableFuture readerThread; CompletableFuture writerThread; - + int threads; int maxEvents = 0; int skipEvents = 0; int readEvents = 0; - + ArrayList procThreads = new ArrayList<>(); ArrayList inputs = new ArrayList<>(); ProgressPrintout progress = new ProgressPrintout(); - ConcurrentLinkedQueue readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue evioQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue hipoQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); - + ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); + void read(String... input) { inputs.addAll(Arrays.asList(input)); while (maxEvents < 1 || readEvents < maxEvents) { if (reader != null && reader.hasEvent()) { - if (readQueue.size() > 100*threads) { + if (evioQueue.size()+hipoQueue.size() > 100*threads) { try { Thread.sleep(100); } catch (InterruptedException ex) {} } else { - DataEvent event = reader.getNextEvent(); readEvents++; - if (skipEvents < 1 || readEvents > skipEvents) - readQueue.offer(event); + 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); + } } } else if (inputs.isEmpty()) break; else { + readEvents = 0; if (inputs.get(0).endsWith(".hipo")) reader = new HipoDataSource(); else reader = new EvioSource(); reader.open(inputs.remove(0)); - updateDictionary((HipoDataSource)reader, writer); + if (reader instanceof HipoDataSource) + updateDictionary((HipoDataSource)reader, writer); + else + maxEvents = ((EvioSource)reader).getEventCount(); } } } - + void write(String output) { writer = new HipoDataSync(); writer.setCompressionType(2); @@ -115,26 +130,34 @@ void write(String output) { } else { writer.writeEvent(writeQueue.poll()); - progress.updateStatus(); + if (readEvents > 20) progress.updateStatus(); } } } void process() { while (true) { - if (readQueue.isEmpty()) { + if (evioQueue.isEmpty() && hipoQueue.isEmpty()) { if (readerThread.isDone()) - if (readQueue.isEmpty()) break; + if (evioQueue.isEmpty() && hipoQueue.isEmpty()) break; try { Thread.sleep(100); } catch (InterruptedException ex) {} } - else { - DataEvent event = readQueue.poll(); - procQueue.offer(event); + else if (!evioQueue.isEmpty()) { + ByteBuffer bb = evioQueue.poll(); + procQueue.offer(bb); + DataEvent event = new EvioDataEvent(bb.array(), ByteOrder.LITTLE_ENDIAN); + processEvent(event); + writeQueue.offer(event); + procQueue.remove(bb); + } + else if (!hipoQueue.isEmpty()){ + DataEvent event = hipoQueue.poll(); + procQueue.offer(event); processEvent(event); writeQueue.offer(event); procQueue.remove(event); - } + } } } } 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 28d3ed55a2..a5186b5e02 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 @@ -17,8 +17,9 @@ import org.jlab.clara.engine.EngineData; import org.jlab.clara.engine.EngineDataType; import java.util.Arrays; +import org.jlab.clas.reco.DecoderEngine.DecoderPool; import org.jlab.coda.jevio.EvioException; -import org.jlab.detector.decode.CLASDecoder4; +import org.jlab.detector.decode.CLASDecoder; import org.jlab.io.evio.EvioDataEvent; import org.jlab.io.evio.EvioSource; import org.jlab.io.hipo.HipoDataEvent; @@ -42,8 +43,8 @@ public class EngineProcessor { private final List schemaExempt = Arrays.asList("RUN::config","DC::tdc"); protected boolean updateDictionary = true; - - protected CLASDecoder4 decoder = new CLASDecoder4(); + + protected DecoderPool decoders = new DecoderPool(64,"default",null); int eventsRead = 0; @@ -135,13 +136,12 @@ private void removeBanks(DataEvent event) { public void initDefault(){ String[] names = new String[]{ - "DECO","MAGFIELDS", + "MAGFIELDS", "DCCR","DCHB","FTOFHB","EC","HTCC","EBHB", "DCTB","FTOFTB","EBTB","VTX" }; String[] services = new String[]{ - "org.jlab.clas.reco.DecoderEngine", "org.jlab.clas.swimtools.MagFieldsEngine", "org.jlab.service.dc.DCHBClustering", "org.jlab.service.dc.DCHBPostClusterConv", @@ -162,7 +162,7 @@ public void initDefault(){ public void initAll(){ String[] names = new String[]{ - "DECO","MAGFIELDS", + "MAGFIELDS", "FTCAL", "FTHODO", "FTTRK", "FTEB", "URWT", "DCCR", "DCHB","FTOFHB","EC","RASTER", "CVTFP","CTOF","CND","BAND", @@ -172,7 +172,6 @@ public void initAll(){ }; String[] services = new String[]{ - "org.jlab.clas.reco.DecoderEngine", "org.jlab.clas.swimtools.MagFieldsEngine", "org.jlab.rec.ft.cal.FTCALEngine", "org.jlab.rec.ft.hodo.FTHODOEngine", @@ -211,7 +210,7 @@ public void initAll(){ } } - public void initCaloDebug(){ + public void initCaloDebug(){ String[] names = new String[]{ "EC","EB" @@ -303,7 +302,16 @@ public final void init(){ * process a single event through the chain. * @param event */ - public void processEvent(DataEvent event){ + public DataEvent processEvent(DataEvent event) { + if (event instanceof EvioDataEvent evio) { + try { + CLASDecoder d = decoders.take(); + Event hipo = d.getDecodedEvent(evio, -1, ++eventsRead, null, null); + event = new HipoDataEvent(hipo, d.getSchemaFactory()); + decoders.put(d); + } + catch (InterruptedException ex) { ex.printStackTrace(); return null; } + } for(Map.Entry engine : this.processorEngines.entrySet()){ try { engine.getValue().processDataEvent(event); @@ -312,6 +320,7 @@ public void processEvent(DataEvent event){ e.printStackTrace(); } } + return event; } public void processFile(String file, String output){ @@ -320,8 +329,10 @@ public void processFile(String file, String output){ public void processEvent(DataEvent event, HipoDataSync writer) { processEvent(event); - removeBanks(event); - writer.writeEvent(event); + if (event instanceof HipoDataEvent) { + removeBanks(event); + writer.writeEvent(event); + } } public void processFile(HipoDataSource reader, HipoDataSync writer, int skipEvents, int maxEvents) { @@ -347,9 +358,7 @@ 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); + processEvent(evio, writer); } if (maxEvents > 0 && eventsRead > maxEvents+skipEvents) break; } catch (EvioException ex) { From 971a777bcea02c34aa6aa0f967512012f6284a3d Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 20 Aug 2026 19:27:19 -0400 Subject: [PATCH 26/27] give it its own file --- .../org/jlab/detector/decode/DecoderPool.java | 37 +++++++++++++++++++ .../org/jlab/clas/reco/DecoderEngine.java | 25 +------------ .../org/jlab/clas/reco/EngineProcessor.java | 2 +- 3 files changed, 39 insertions(+), 25 deletions(-) create mode 100644 common-tools/clas-detector/src/main/java/org/jlab/detector/decode/DecoderPool.java diff --git a/common-tools/clas-detector/src/main/java/org/jlab/detector/decode/DecoderPool.java b/common-tools/clas-detector/src/main/java/org/jlab/detector/decode/DecoderPool.java new file mode 100644 index 0000000000..7fd7eeeb90 --- /dev/null +++ b/common-tools/clas-detector/src/main/java/org/jlab/detector/decode/DecoderPool.java @@ -0,0 +1,37 @@ +package org.jlab.detector.decode; + +import java.util.concurrent.ArrayBlockingQueue; + +/** + * + * @author baltzell + */ +public class DecoderPool { + + ArrayBlockingQueue pool; + int constantsShared = 64; + + public DecoderPool(int size, String variation, String timestamp) { + pool = new ArrayBlockingQueue<>(size); + CLASDecoder d0 = null; + for (int i=0; i pool; - int constantsShared = 64; - public DecoderPool(int size, String variation, String timestamp) { - pool = new ArrayBlockingQueue<>(size); - CLASDecoder d0 = null; - for (int i=0; i ED_TYPES = ClaraUtil.buildDataTypes( Clas12Types.EVIO,Clas12Types.HIPO,EngineDataType.JSON,EngineDataType.STRING); 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 a5186b5e02..46d38862f7 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 @@ -17,9 +17,9 @@ import org.jlab.clara.engine.EngineData; import org.jlab.clara.engine.EngineDataType; import java.util.Arrays; -import org.jlab.clas.reco.DecoderEngine.DecoderPool; import org.jlab.coda.jevio.EvioException; import org.jlab.detector.decode.CLASDecoder; +import org.jlab.detector.decode.DecoderPool; import org.jlab.io.evio.EvioDataEvent; import org.jlab.io.evio.EvioSource; import org.jlab.io.hipo.HipoDataEvent; From 7958719a56195a182e7de9b7217f7fd3429929e7 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 20 Aug 2026 20:01:31 -0400 Subject: [PATCH 27/27] remove no-longer-used queue --- .../main/java/org/jlab/clas/reco/EngineMultiProcessor.java | 5 ----- 1 file changed, 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 629aace26e..507b8bc466 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 @@ -75,7 +75,6 @@ public void process(String output, String... input) { ConcurrentLinkedQueue evioQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue hipoQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue writeQueue = new ConcurrentLinkedQueue<>(); - ConcurrentLinkedQueue procQueue = new ConcurrentLinkedQueue<>(); void read(String... input) { inputs.addAll(Arrays.asList(input)); @@ -145,18 +144,14 @@ void process() { } else if (!evioQueue.isEmpty()) { ByteBuffer bb = evioQueue.poll(); - procQueue.offer(bb); DataEvent event = new EvioDataEvent(bb.array(), ByteOrder.LITTLE_ENDIAN); processEvent(event); writeQueue.offer(event); - procQueue.remove(bb); } else if (!hipoQueue.isEmpty()){ DataEvent event = hipoQueue.poll(); - procQueue.offer(event); processEvent(event); writeQueue.offer(event); - procQueue.remove(event); } } }