From c2007e39168df103eb66fc1a6fe114fe60926e64 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 3 Sep 2026 23:08:18 -0400 Subject: [PATCH 01/19] separate functionality --- .../java/org/jlab/utils/benchmark/BenchmarkTimer.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java index 803a754d8b..59918ff499 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java @@ -95,9 +95,11 @@ public double getSeconds(){ @Override public String toString() { - double timePerCall = 0.0; - if (numberOfCalls.get() != 0) timePerCall = getMiliseconds() / numberOfCalls.get(); return String.format("%-15s : #Calls %12d, Total = %12.2f sec, Unit = %12.3f msec", - getName(), numberOfCalls.get(), getSeconds(), timePerCall); + getName(), numberOfCalls.get(), getSeconds(), getTimePerCall()); + } + + public double getTimePerCall() { + return numberOfCalls.get() > 0 ? getMiliseconds() / numberOfCalls.get() : 0; } } From 2c569d3acb7915c8d3a058e651cc23c9b856fe86 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 3 Sep 2026 23:10:26 -0400 Subject: [PATCH 02/19] fix method name --- .../main/java/org/jlab/utils/benchmark/BenchmarkTimer.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java index 59918ff499..8092d63024 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/BenchmarkTimer.java @@ -85,7 +85,7 @@ public void reset(){ isPaused = true; } - public double getMiliseconds(){ + public double getMilliseconds(){ return totalTime.get() / 1.0e6; } @@ -100,6 +100,6 @@ public String toString() { } public double getTimePerCall() { - return numberOfCalls.get() > 0 ? getMiliseconds() / numberOfCalls.get() : 0; + return numberOfCalls.get() > 0 ? getMilliseconds() / numberOfCalls.get() : 0; } } From e2ff5783fcb86ccb1b4c1257a13f9e90d1be1182 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 3 Sep 2026 23:43:43 -0400 Subject: [PATCH 03/19] add accessor --- .../main/java/org/jlab/utils/benchmark/ProgressPrintout.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java index c1ef3626fe..48da5624a0 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java @@ -88,6 +88,10 @@ public String getItemString(String itemname){ } return str.toString(); } + + public int getNumberOfCalls() { + return numberOfCalls; + } public static void main(String[] args){ ProgressPrintout progress = new ProgressPrintout(); From e56685823457e83385f681ec80b80c6217f8a399 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 3 Sep 2026 23:19:26 -0400 Subject: [PATCH 04/19] add serial class --- .../jlab/detector/serial/SerialHoncho.java | 99 +++++++++++++++++++ 1 file changed, 99 insertions(+) create mode 100644 common-tools/clas-detector/src/main/java/org/jlab/detector/serial/SerialHoncho.java diff --git a/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/SerialHoncho.java b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/SerialHoncho.java new file mode 100644 index 0000000000..86a4868432 --- /dev/null +++ b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/SerialHoncho.java @@ -0,0 +1,99 @@ +package org.jlab.detector.serial; + +import java.util.TreeMap; +import java.util.TreeSet; +import org.jlab.detector.calib.utils.ConstantsManager; +import org.jlab.detector.decode.CLASDecoder; +import org.jlab.detector.helicity.HelicitySequence; +import org.jlab.detector.helicity.HelicityState; +import org.jlab.detector.scalers.DaqScalersSequence; +import org.jlab.jnp.hipo4.data.Bank; +import org.jlab.jnp.hipo4.data.Event; +import org.jlab.jnp.hipo4.data.SchemaFactory; +import org.jlab.jnp.hipo4.io.HipoWriterSorted; + +/** + * + * @author baltzell + */ +public class SerialHoncho { + + static final String[] TAG1BANKS = {"RUN::scaler","HEL::scaler","RAW::scaler","RAW::epics","HEL::flip","COAT::config"}; + + SchemaFactory schema; + Bank[] tag1banks; + Bank runConfig; + Bank helicityAdc; + ConstantsManager conman; + TreeMap eventUnix; + TreeSet helicities; + DaqScalersSequence scalers; + + public SerialHoncho(SchemaFactory schema) { + this.schema = schema; + conman = new ConstantsManager(); + conman.init("/runcontrol/hwp","/runcontrol/helicity"); + runConfig = new Bank(schema.getSchema("RUN::config")); + helicityAdc = new Bank(schema.getSchema("HEL::adc")); + helicities = new TreeSet<>(); + scalers = new DaqScalersSequence(schema); + eventUnix = new TreeMap<>(); + tag1banks = new Bank[TAG1BANKS.length]; + for (int i=0; i 0) { + int unix = runConfig.getInt("unixtime",0); + int evno = runConfig.getInt("event",0); + if (unix > 0 && evno > 0) eventUnix.put(evno, unix); + } + helicities.add(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman)); + return CLASDecoder.createTaggedEvent(event, runConfig, tag1banks); + } + + public void finish(HipoWriterSorted writer) { + writer.addEvent(getUnixEvent(runConfig),1); + HelicitySequence.writeFlips(schema, writer, helicities); + } + + public void clear() { + while (helicities.size() > 100) helicities.pollFirst(); + scalers.clear(100); + } + + Event getUnixEvent(Bank config) { + Bank unix = new Bank(schema.getSchema("RUN::unix")); + unix.setRows(eventUnix.size()); + int row = 0; + for (int evno : eventUnix.keySet()) { + unix.putInt("event", row, evno); + unix.putInt("unixtime",row, eventUnix.get(evno)); + row++; + } + Event e = new Event(); + e.write(config); + e.write(unix); + return e; + } + + public DaqScalersSequence getScalers() { + return scalers; + } + + public TreeSet getHelicities() { + return helicities; + } + + public ConstantsManager getConstantsManager() { + return conman; + } + + public SchemaFactory getSchemaFactory() { + return schema; + } +} From 2988f8a336bb3b8b4589c15e0b15fd60a515f4bd Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 09:21:43 -0400 Subject: [PATCH 05/19] separate "serial" functionality for reusability --- .../java/org/jlab/io/clara/Clas12Writer.java | 66 +++---------------- 1 file changed, 8 insertions(+), 58 deletions(-) diff --git a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java index a094b939cd..40bc0e6b3c 100644 --- a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java +++ b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java @@ -3,16 +3,11 @@ import java.io.File; import java.nio.file.Path; import java.util.List; -import java.util.TreeMap; -import java.util.TreeSet; import org.jlab.analysis.postprocess.Processor; import org.jlab.clara.std.services.EventWriterException; import org.jlab.detector.calib.utils.ConstantsManager; -import org.jlab.detector.decode.CLASDecoder4; -import org.jlab.detector.helicity.HelicitySequence; import org.jlab.detector.helicity.HelicitySequenceDelayed; -import org.jlab.detector.helicity.HelicityState; -import org.jlab.detector.scalers.DaqScalersSequence; +import org.jlab.detector.serial.SerialHoncho; import org.jlab.jnp.hipo4.data.Bank; import org.jlab.jnp.hipo4.data.Event; import org.jlab.jnp.hipo4.data.SchemaFactory; @@ -33,34 +28,22 @@ */ public class Clas12Writer extends HipoToHipoWriter { - static final String[] TAG1BANKS = {"RUN::scaler","HEL::scaler","RAW::scaler","RAW::epics","HEL::flip","COAT::config"}; - - Bank[] tag1banks; + SerialHoncho serial; Bank runConfig; - Bank helicityAdc; ConstantsManager conman; - TreeMap eventUnix; - TreeSet helicities; - DaqScalersSequence scalers; SchemaFactory fullSchema; boolean postprocess; private void init(JSONObject opts) { fullSchema = new SchemaFactory(); fullSchema.initFromDirectory(FileUtils.getEnvironmentPath("CLAS12DIR","etc/bankdefs/hipo4")); + serial = new SerialHoncho(fullSchema); runConfig = new Bank(fullSchema.getSchema("RUN::config")); - helicityAdc = new Bank(fullSchema.getSchema("HEL::adc")); - helicities = new TreeSet<>(); - scalers = new DaqScalersSequence(fullSchema); conman = new ConstantsManager(); - eventUnix = new TreeMap<>(); conman.init("/runcontrol/hwp","/runcontrol/helicity"); postprocess = opts.optBoolean("postprocess", false); if (opts.has("variation")) conman.setVariation(opts.getString("variation")); if (opts.has("timestamp")) conman.setTimeStamp(opts.getString("timestamp")); - tag1banks = new Bank[TAG1BANKS.length]; - for (int i=0; i 0) { - int unix = runConfig.getInt("unixtime",0); - int evno = runConfig.getInt("event",0); - if (unix > 0 && evno > 0) eventUnix.put(evno, unix); - } - helicities.add(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman)); - Event t = CLASDecoder4.createTaggedEvent((Event)event, runConfig, tag1banks); + Event t = serial.read((Event)event); if (!t.isEmpty()) writer.addEvent(t, 1); super.writeEvent(event); } @Override protected void closeWriter() { - HelicitySequence.writeFlips(fullSchema, writer, helicities); - writer.addEvent(getUnixEvent(runConfig),1); + serial.finish(writer); super.closeWriter(); if (postprocess) postprocess(); - // keep the latest helicity/scaler reading for the next file: - while (helicities.size() > 60) helicities.pollFirst(); - scalers.clear(10); + serial.clear(); } /** @@ -120,35 +91,14 @@ private int getRunNumber() { return 0; } - /** - * Get a new event with a RUN::unix bank containing event-timestamp mapping, - * and the latest RUN::config bank. - * @param config - * @return - */ - private Event getUnixEvent(Bank config) { - Bank unix = new Bank(fullSchema.getSchema("RUN::unix")); - unix.setRows(eventUnix.size()); - int row = 0; - for (int evno : eventUnix.keySet()) { - unix.putInt("event", row, evno); - unix.putInt("unixtime",row, eventUnix.get(evno)); - row++; - } - Event e = new Event(); - e.write(config); - e.write(unix); - return e; - } - /** * Copy helicity/charge tag-1 information to all events. */ private void postprocess() { int d = conman.getConstants(getRunNumber(), "/runcontrol/helicity").getIntValue("delay",0,0,0); HelicitySequenceDelayed helicity = new HelicitySequenceDelayed(d); - helicity.addStream(helicities); - Processor p = new Processor(List.of(filename), fullSchema, helicity, scalers); + helicity.addStream(serial.getHelicities()); + Processor p = new Processor(List.of(filename), fullSchema, helicity, serial.getScalers()); HipoReader r = new HipoReader(); r.open(filename); Event e = new Event(); From d8b88020c476d0a19fdf401a1902743959941c6e Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 10:18:01 -0400 Subject: [PATCH 06/19] relocate serial/postprocessing utilities to clas-detector --- bin/postprocess2 | 2 +- common-tools/clara-io/pom.xml | 6 ------ .../java/org/jlab/io/clara/Clas12Writer.java | 4 ++-- .../analysis/postprocess/RebuildScalers.java | 3 ++- .../jlab/analysis/postprocess/Tag1ToEvent.java | 10 ++++++---- .../org/jlab/detector/serial/PostProcessor.java} | 16 ++++++++-------- .../org/jlab/detector/serial/SerialUtil.java} | 6 +++--- 7 files changed, 22 insertions(+), 25 deletions(-) rename common-tools/{clas-analysis/src/main/java/org/jlab/analysis/postprocess/Processor.java => clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java} (93%) rename common-tools/{clas-analysis/src/main/java/org/jlab/analysis/postprocess/Util.java => clas-detector/src/main/java/org/jlab/detector/serial/SerialUtil.java} (98%) diff --git a/bin/postprocess2 b/bin/postprocess2 index b10c9ff07a..19219be743 100755 --- a/bin/postprocess2 +++ b/bin/postprocess2 @@ -6,5 +6,5 @@ export MALLOC_ARENA_MAX=1 java ${JAVA_OPTS-} -Xmx768m -Xms768m -XX:+UseSerialGC \ -cp ${COATJAVA_CLASSPATH:-''} \ - org.jlab.analysis.postprocess.Processor \ + org.jlab.detector.serial.PostProcessor \ $* diff --git a/common-tools/clara-io/pom.xml b/common-tools/clara-io/pom.xml index f5f100138e..f1c3f6f0e1 100644 --- a/common-tools/clara-io/pom.xml +++ b/common-tools/clara-io/pom.xml @@ -48,12 +48,6 @@ 14.2.0-SNAPSHOT - - org.jlab.clas - clas-analysis - 14.2.0-SNAPSHOT - - org.jlab.clas clas-utils diff --git a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java index 40bc0e6b3c..5e268de75a 100644 --- a/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java +++ b/common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java @@ -3,7 +3,7 @@ import java.io.File; import java.nio.file.Path; import java.util.List; -import org.jlab.analysis.postprocess.Processor; +import org.jlab.detector.serial.PostProcessor; import org.jlab.clara.std.services.EventWriterException; import org.jlab.detector.calib.utils.ConstantsManager; import org.jlab.detector.helicity.HelicitySequenceDelayed; @@ -98,7 +98,7 @@ private void postprocess() { int d = conman.getConstants(getRunNumber(), "/runcontrol/helicity").getIntValue("delay",0,0,0); HelicitySequenceDelayed helicity = new HelicitySequenceDelayed(d); helicity.addStream(serial.getHelicities()); - Processor p = new Processor(List.of(filename), fullSchema, helicity, serial.getScalers()); + PostProcessor p = new PostProcessor(List.of(filename), fullSchema, helicity, serial.getScalers()); HipoReader r = new HipoReader(); r.open(filename); Event e = new Event(); diff --git a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/RebuildScalers.java b/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/RebuildScalers.java index 525203f641..6fb4fe3d0c 100644 --- a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/RebuildScalers.java +++ b/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/RebuildScalers.java @@ -9,6 +9,7 @@ import org.jlab.detector.scalers.DaqScalers; import org.jlab.detector.helicity.HelicitySequenceManager; import org.jlab.detector.scalers.DaqScalersSequence; +import org.jlab.detector.serial.SerialUtil; import org.jlab.jnp.hipo4.data.Bank; import org.jlab.jnp.hipo4.data.Event; import org.jlab.jnp.hipo4.io.HipoReader; @@ -122,7 +123,7 @@ else if (seq != null) { runScalerBank = ds.createRunBank(writer.getSchemaFactory()); helScalerBank = ds.createHelicityBank(writer.getSchemaFactory()); - Util.assignScalerHelicity(event, helScalerBank, helSeq); + SerialUtil.assignScalerHelicity(event, helScalerBank, helSeq); // put modified HEL/RUN::scaler back in the event: event.write(runScalerBank); diff --git a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Tag1ToEvent.java b/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Tag1ToEvent.java index b85efd3a95..b3293962f2 100644 --- a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Tag1ToEvent.java +++ b/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Tag1ToEvent.java @@ -1,5 +1,6 @@ package org.jlab.analysis.postprocess; +import org.jlab.detector.serial.PostProcessor; import java.util.TreeMap; import java.util.logging.Logger; import org.jlab.clas.reco.ReconstructionEngine; @@ -13,6 +14,7 @@ import org.jlab.detector.scalers.DaqScalersSequence; import org.jlab.detector.helicity.HelicityBit; import org.jlab.detector.helicity.HelicitySequenceDelayed; +import org.jlab.detector.serial.SerialUtil; import org.jlab.jnp.hipo4.data.SchemaFactory; import org.jlab.utils.groups.IndexedTable; import org.jlab.utils.options.OptionParser; @@ -81,7 +83,7 @@ public static void main(String[] args) { LOGGER.info("\n>>> Initializing helicity configuration from CCDB ...\n"); ConstantsManager conman = new ConstantsManager(); conman.init("/runcontrol/hwp","/runcontrol/helicity"); - final int run = Util.getRunNumber(parser.getInputList().get(0)); + final int run = SerialUtil.getRunNumber(parser.getInputList().get(0)); IndexedTable helTable = conman.getConstants(run, "/runcontrol/helicity"); // Initialize the scaler sequence from tag-1 events: @@ -102,7 +104,7 @@ public static void main(String[] args) { } // Initialize the unix-event map: - TreeMap eventUnix = Processor.getEventUnixMap(schema, parser.getInputList()); + TreeMap eventUnix = PostProcessor.getEventUnixMap(schema, parser.getInputList()); // Loop over the input HIPO files: LOGGER.info("\n>>> Starting post-processing ...\n"); @@ -138,7 +140,7 @@ public static void main(String[] args) { if (doHelicityDelay) { recEventBank.putByte("helicity",0,hb.value()); recEventBank.putByte("helicityRaw",0,hbraw.value()); - Util.assignScalerHelicity(runConfigBank.getLong("timestamp",0), helScalerBank, helSeq); + SerialUtil.assignScalerHelicity(runConfigBank.getLong("timestamp",0), helScalerBank, helSeq); } // Write beam charge to REC::Event: @@ -169,7 +171,7 @@ public static void main(String[] args) { writer.addEvent(event, event.getEventTag()); // Copy config banks to new, tag-1 events: - Util.createTag1Events(writer, event, configEvent, configBanks); + SerialUtil.createTag1Events(writer, event, configEvent, configBanks); } reader.close(); diff --git a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Processor.java b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java similarity index 93% rename from common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Processor.java rename to common-tools/clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java index 7652b74d4c..0d0f1b96bc 100644 --- a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Processor.java +++ b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java @@ -1,4 +1,4 @@ -package org.jlab.analysis.postprocess; +package org.jlab.detector.serial; import java.util.List; import java.util.TreeMap; @@ -24,7 +24,7 @@ * * @author baltzell */ -public class Processor { +public class PostProcessor { public static final String CCDB_TABLES[] = {"/runcontrol/fcup","/runcontrol/slm", "/runcontrol/helicity","/daq/config/scalers/dsc1","/runcontrol/hwp"}; @@ -37,7 +37,7 @@ public class Processor { private HelicitySequenceDelayed helicitySequence = null; private TreeMap eventUnix = null; - public Processor(List files, boolean restream, boolean rebuild) { + public PostProcessor(List files, boolean restream, boolean rebuild) { HipoReader r = new HipoReader(); r.open(files.get(0)); schemaFactory = r.getSchemaFactory(); @@ -46,13 +46,13 @@ public Processor(List files, boolean restream, boolean rebuild) { recEvent = new Bank(schemaFactory.getSchema("REC::Event")); conman = new ConstantsManager(); conman.init(CCDB_TABLES); - helicitySequence = Util.getHelicity(files, schemaFactory, restream, conman); + helicitySequence = SerialUtil.getHelicity(files, schemaFactory, restream, conman); if (rebuild) chargeSequence = DaqScalersSequence.rebuildSequence(1, conman, files); else chargeSequence = DaqScalersSequence.readSequence(files); eventUnix = getEventUnixMap(schemaFactory, files); } - public Processor(List files, SchemaFactory schema, HelicitySequenceDelayed h, DaqScalersSequence s) { + public PostProcessor(List files, SchemaFactory schema, HelicitySequenceDelayed h, DaqScalersSequence s) { schemaFactory = schema; helicitySequence = h; chargeSequence = s; @@ -102,7 +102,7 @@ private void processEventHelicity(DataEvent event, DataBank runcfg, DataBank rec DataBank helScaler = event.getBank("HEL::scaler"); if (helScaler.rows()>0) { event.removeBank("HEL::scaler"); - Util.assignScalerHelicity(runcfg.getLong("timestamp",0), ((HipoDataBank)helScaler).getBank(), helicitySequence); + SerialUtil.assignScalerHelicity(runcfg.getLong("timestamp",0), ((HipoDataBank)helScaler).getBank(), helicitySequence); event.appendBank(helScaler); } } @@ -122,7 +122,7 @@ private void processEventHelicity(Event event, Bank runcfg, Bank recevt) { event.read(helScaler); if (helScaler.getRows()>0) { event.remove(schemaFactory.getSchema("HEL::scaler")); - Util.assignScalerHelicity(runcfg.getLong("timestamp",0), helScaler, helicitySequence); + SerialUtil.assignScalerHelicity(runcfg.getLong("timestamp",0), helScaler, helicitySequence); event.write(helScaler); } } @@ -247,7 +247,7 @@ public static void main(String args[]) { boolean restream = !o.getOption("-f").isDefault(); boolean rebuild = !o.getOption("-c").isDefault(); - Processor post = new Processor(o.getInputList(), restream, rebuild); + PostProcessor post = new PostProcessor(o.getInputList(), restream, rebuild); HipoWriterSorted writer = null; diff --git a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Util.java b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/SerialUtil.java similarity index 98% rename from common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Util.java rename to common-tools/clas-detector/src/main/java/org/jlab/detector/serial/SerialUtil.java index 7ac75bcea2..8967760c25 100644 --- a/common-tools/clas-analysis/src/main/java/org/jlab/analysis/postprocess/Util.java +++ b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/SerialUtil.java @@ -1,4 +1,4 @@ -package org.jlab.analysis.postprocess; +package org.jlab.detector.serial; import java.sql.Time; import java.util.Arrays; @@ -26,9 +26,9 @@ * Static utility methods for postprocessing. * @author baltzell */ -class Util { +public class SerialUtil { - static final Logger logger = Logger.getLogger(Util.class.getName()); + static final Logger logger = Logger.getLogger(SerialUtil.class.getName()); /** * Assign the delay-corrected helicity to the HEL::scaler bank's rows From 6522196eecf41154c11437cb96911e78760c48ea Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 11:43:07 -0400 Subject: [PATCH 07/19] default to zero for easier event counting --- .../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 39c93eefd4..dfa428f7ec 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 @@ -310,7 +310,7 @@ public void processEvent(DataEvent event){ } public void processFile(String file, String output){ - this.processFile(file, output, -1, -1); + this.processFile(file, output, 0, 0); } public void processEvent(DataEvent event, HipoDataSync writer) { @@ -396,8 +396,8 @@ protected static OptionParser getParser() { 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("-s","0","number of events to skip"); + parser.addOption("-n","0","number of events to process"); parser.addOption("-y","0","yaml file"); parser.addOption("-u","true","update dictionary from writer ? "); parser.addOption("-S",null,"schema directory"); From b1f3a22a334efb11d423eb9033b3de7340f37085 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 11:39:59 -0400 Subject: [PATCH 08/19] add support for io-services --- .../src/main/java/org/jlab/utils/ClaraYaml.java | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/common-tools/clas-io/src/main/java/org/jlab/utils/ClaraYaml.java b/common-tools/clas-io/src/main/java/org/jlab/utils/ClaraYaml.java index b6525bfb3b..2f4f0afec8 100644 --- a/common-tools/clas-io/src/main/java/org/jlab/utils/ClaraYaml.java +++ b/common-tools/clas-io/src/main/java/org/jlab/utils/ClaraYaml.java @@ -143,9 +143,11 @@ public JSONObject filter(String serviceName) { /** * Emulate the way CLARA parses the full YAML and presents it in EngineData. - * The "global" and "service" subsections in the "configuration" section get - * squashed into one namespace, and service-specific keys override any - * globals of the same name. + * + * The YAML's "global" and "services" configuration sections get squashed + * into one namespace, with service-specific parameters overriding globals + * of the same name. Also, for the special services named "reader" and + * "writer", the "io-services" section is searched instead of "services". * * @param claraJson the full CLARA YAML contents * @param serviceName the name of the service in CLARA YAML (not class name) @@ -161,9 +163,10 @@ public static JSONObject filter(JSONObject claraJson, String serviceName) { ret.accumulate(key, globals.getString(key)); } } - if (config.has("services")) { - if (config.getJSONObject("services").has(serviceName)) { - JSONObject service = config.getJSONObject("services").getJSONObject(serviceName); + String section = serviceName.equals("reader") || serviceName.equals("writer") ? "io-services" : "services"; + if (config.has(section)) { + if (config.getJSONObject(section).has(serviceName)) { + JSONObject service = config.getJSONObject(section).getJSONObject(serviceName); for (String key : service.keySet()) { ret.put(key, service.getString(key)); } From c21f1f6f67b780ca9af871d128540fe6863c3a64 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 18:15:36 -0400 Subject: [PATCH 09/19] add ReconMutil --- .../java/org/jlab/clas/reco/ReconMutil.java | 498 ++++++++++++++++++ 1 file changed, 498 insertions(+) create mode 100644 common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java new file mode 100644 index 0000000000..9a8a21f746 --- /dev/null +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -0,0 +1,498 @@ +package org.jlab.clas.reco; + +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStream; +import java.io.InputStreamReader; +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.logging.Level; +import java.util.logging.Logger; +import org.jlab.clara.engine.EngineData; +import org.jlab.clara.engine.EngineDataType; +import org.jlab.coda.jevio.EvioException; +import org.jlab.detector.decode.CLASDecoder; +import org.jlab.detector.decode.CLASDecoderPool; +import org.jlab.detector.serial.SerialHoncho; +import org.jlab.io.evio.EvioDataEvent; +import org.jlab.io.evio.EvioSource; +import org.jlab.io.hipo.HipoDataEvent; +import org.jlab.jnp.hipo4.data.Bank; +import org.jlab.jnp.hipo4.data.Event; +import org.jlab.jnp.hipo4.data.SchemaFactory; +import org.jlab.jnp.hipo4.io.HipoReader; +import org.jlab.jnp.hipo4.io.HipoWriterSorted; +import org.jlab.utils.ClaraYaml; +import org.jlab.utils.benchmark.Benchmark; +import org.jlab.utils.benchmark.ProgressPrintout; +import org.jlab.utils.options.OptionParser; +import org.jlab.utils.system.ClasUtilsFile; +import org.json.JSONObject; + +/** + * + * @author baltzell + */ +final class ReconMutil { + + // Performance parameters: + final int BENCH_SECONDS = 30; + final int CHUNKS_PER_QUEUE = 100; + final int EVENTS_PER_CHUNK = 100; + + // File reader and writer: + Object reader; + HipoWriterSorted writer; + List schemaBankList; + static final SchemaFactory schema = new SchemaFactory(); + static { schema.initFromDirectory(ClasUtilsFile.getResourceDir("CLAS12DIR","etc/bankdefs/hipo4")); } + + // Processors: + SerialHoncho serial; + Map engines = new LinkedHashMap<>(); + CLASDecoderPool decoders = new CLASDecoderPool(64,"default",null); + + // Threads and queues: + CompletableFuture readerThread; + CompletableFuture writerThread; + CompletableFuture rethreadThread; + ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); + + // Static parameters: + int maxEvents; + int skipEvents; + ClaraYaml yaml; + OptionParser parser; + + // Progress counters: + int readEvents; + int writeEvents; + int failEvents; + int fileEvents; + int maxFileEvents; + ProgressPrintout progress = new ProgressPrintout(); + + ReconMutil(OptionParser parser) { + init(parser); + } + + /** + * The thread launcher and collector. + * @param threads number of threads + * @param output name of output file to write + * @param input names of input files to read + */ + void launch(int[] threads, String output, String... input) { + reset(); + readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); + writerThread = CompletableFuture.runAsync(() -> { write(output); }); + for (int i=0; i { process(j); })); + } + while (!writerThread.isDone()) { + sleep(100); + for (CompletableFuture f : procThreads) + if (f.isDone()) procThreads.remove(f); + if (threads.length > 1 && rethreadThread == null && writeEvents > 100) { + rethreadThread = CompletableFuture.runAsync(() -> { rethread(BENCH_SECONDS,threads); }); + rethreadThread.join(); + reset(); + } + } + } + + /** + * The reader thread. + * @param input input filenames + */ + void read(int threads, String... input) { + + // convert input filenames to a list: + List inputs = new ArrayList<>(Arrays.asList(input)); + + // initialize the event chunk: + List chunk = new ArrayList<>(EVENTS_PER_CHUNK); + + // loop over input events: + while ( (maxEvents < 1 || readEvents < maxEvents) && + (maxFileEvents < 1 || fileEvents < maxFileEvents) ) { + + if (reader != null) { + + // sleep instead of overfilling the read queue: + if (readQueue.size() > CHUNKS_PER_QUEUE*threads) sleep(1000); + + // read next event into chunk, and fill queue if chunk full: + else chunk = read(chunk); + } + + // open the next input file: + else if (!inputs.isEmpty()) open(inputs.removeFirst()); + + // no more events to read: + else break; + } + + // write leftover, partial chunk: + if (!chunk.isEmpty()) { + System.err.println("writing partial chunk: "+chunk.size()); + readEvents += chunk.size(); + readQueue.offer(chunk); + } + + if (reader instanceof EvioSource evio) evio.close(); + } + + /** + * The event processor thread. + * @param thread unique thread number + */ + void process(int thread) { + while (true) { + List o = readQueue.poll(); + if (o == null) { + if (readerThread.isDone() && readQueue.isEmpty() && + writeEvents+skipEvents+failEvents >= readEvents) break; + sleep(100); + } + else { + // put the event back on the queue if we're rethreading: + //if (rethreadThread != null && !rethreadThread.isDone()) readQueue.offer(o); + List chunk = new ArrayList<>(o.size()); + for (int i=0; i engine : engines.entrySet()) { + Benchmark.getInstance().resume(engine.getValue().getName()); + try { engine.getValue().processDataEvent(event); } + catch (Exception ex) { ex.printStackTrace(); } + Benchmark.getInstance().pause(engine.getValue().getName()); + } + + Benchmark.getInstance().resume("serial"); + Event e = event.getHipoEvent(); + Event t = serial.read(e); + t.setEventTag(1); + Benchmark.getInstance().pause("serial"); + chunk.add(e); + chunk.add(t); + } + writeQueue.offer(chunk); + } + } + } + + /** + * The writer thread. + * @param output output filename + */ + void write(String output) { + if (output != null) writer = open(output, yaml); + while (true) { + List e = writeQueue.poll(); + if (e == null) { + if (readerThread.isDone() && procThreads.isEmpty() && writeQueue.isEmpty()) { + close(); + break; + } + sleep(1000); + } + else { + for (int i=0; i 0 || schemaBankList.isEmpty()) + writer.addEvent(e.get(i), e.get(i).getEventTag()); + else + writer.addEvent(e.get(i).reduceEvent(schemaBankList), e.get(i).getEventTag()); + } + Benchmark.getInstance().pause("write"); + progress.updateStatus(); + } + writeEvents += e.size(); + } + } + } + + /** + * The rethreader thread. + * @param seconds delay before switching to next thread count + * @param threads thread counts to use + */ + void rethread(int seconds, int... threads) { + System.out.println("~~~~~~~~~ Rethreading Initiated ~~~~~~~~~"); + for (int i=0; i { process(k); })); + } + while (progress.getNumberOfCalls() < 100) sleep(1000); + sleep(seconds*1000); + System.out.println(String.format("\n~~~~~~~~~ Rethreading Count: %d ~~~~~~~~~\n",threads[i])); + System.out.println(progress.getUpdateString()); + System.out.println(Benchmark.getInstance()); + } + } + + /** + * Decode an event. + * @param bytes the EVIO byte buffer + * @return decoded event + */ + HipoDataEvent decode(ByteBuffer bytes) { + Benchmark.getInstance().resume("evio"); + EvioDataEvent evio = new EvioDataEvent(bytes.array(), ByteOrder.LITTLE_ENDIAN); + Benchmark.getInstance().pause("evio"); + Benchmark.getInstance().resume("deco"); + HipoDataEvent hipo; + try { + CLASDecoder d = decoders.take(); + hipo = d.getDecodedDataEvenet(evio); + decoders.put(d); + } + catch (InterruptedException ex) { hipo = null; } + Benchmark.getInstance().pause("deco"); + return hipo; + } + + /** + * Open a new input HIPO/EVIO event file. + * @param filename + */ + void open(String filename) { + fileEvents = 0; + if (filename.endsWith(".hipo")) { + reader = new HipoReader(); + ((HipoReader)reader).open(filename); + maxFileEvents = ((HipoReader)reader).getEventCount(); + } + else { + reader = new EvioSource(); + ((EvioSource)reader).open(filename); + maxFileEvents = ((EvioSource)reader).getEventCount(); + } + } + + /** + * Open a new writer and initialize its schema. + * @param filename output filename + * @param yaml the configuration + */ + HipoWriterSorted open(String filename, ClaraYaml yaml) { + HipoWriterSorted w = new HipoWriterSorted(); + w.setCompressionType(2); + String d = ClasUtilsFile.getResourceDir("CLAS12DIR", "etc/bankdefs/hipo4"); + if (yaml.getSchemaDirectory() != null) d = yaml.getSchemaDirectory(); + if (!parser.getOption("-S").isDefault()) d = parser.getOption("-S").stringValue(); + SchemaFactory s = new SchemaFactory(); + s.initFromDirectory(d); + JSONObject json = yaml.filter("writer"); + if (json.has("wildcard")) { + SchemaFactory s2 = s.reduce(json.getString("wildcard")); + w.getSchemaFactory().copy(s2); + } + else w.getSchemaFactory().copy(s); + schemaBankList = new ArrayList<>(); + if (json.has("wildcard")) { + if (json.optBoolean("schema_filter",true)) { + int schemaSize = w.getSchemaFactory().getSchemaList().size(); + for (int i=0; i read(List chunk) { + Benchmark.getInstance().resume("read"); + Object o = null; + if (reader instanceof EvioSource evio) { + try { o = evio.getEventBuffer(++fileEvents, true); } + catch (EvioException ex) { + failEvents++; + ex.printStackTrace(); + } + } + else { + Event event = new Event(); + o = ((HipoReader)reader).getEvent(event, fileEvents); + } + if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { + chunk.add(o); + if (chunk.size() >= EVENTS_PER_CHUNK) { + readQueue.offer(chunk); + readEvents += chunk.size(); + chunk = new ArrayList<>(EVENTS_PER_CHUNK); + } + } + Benchmark.getInstance().pause("read"); + return chunk; + } + + /** + * Close the output file. + */ + void close() { + serial.finish(writer); + writer.close(); + System.out.println(String.format("recon-mutil :: read/write/diff = %d/%d/%d", + readEvents, writeEvents, readEvents-writeEvents)); + } + + /** + * Forcefully shutdown all threads, close files, and reset queuess and counters. + */ + void reset() { + for (CompletableFuture f : procThreads) f.cancel(true); + if (readerThread != null) readerThread.cancel(true); + if (writerThread != null) { + writerThread.cancel(true); + close(); + } + readQueue = new ConcurrentLinkedQueue<>(); + writeQueue = new ConcurrentLinkedQueue<>(); + procThreads = new ConcurrentLinkedQueue(); + readEvents = 0; + writeEvents = 0; + failEvents = 0; + } + + /** + * Add a new engine to the list. + * @param label display name + * @param clazz full class name + * @param cfg engine configuration + * @return + */ + ReconstructionEngine addEngine(String label, String clazz, JSONObject cfg) { + ReconstructionEngine engine = null; + try { + Class c = Class.forName(clazz); + if (ReconstructionEngine.class.isAssignableFrom(c)==true){ + engine = (ReconstructionEngine) c.newInstance(); + if (cfg != null && !cfg.toString().equals("null")) { + EngineData input = new EngineData(); + input.setData(EngineDataType.JSON.mimeType(), cfg); + engine.configure(input); + } + else engine.init(); + engines.put(label == null ? engine.getName() : label, engine); + } + else Logger.getLogger(ReconMutil.class.getPackage().getName()) + .log(clazz.contains("DecoderEngine") ? Level.INFO : Level.SEVERE, + "Class is not a reconstruction engine : {0}", clazz); + } catch (ClassNotFoundException | InstantiationException | IllegalAccessException ex) { + Logger.getLogger(ReconMutil.class.getPackage().getName()).log(Level.SEVERE, null, ex); + } + return engine; + } + + /** + * Catch interruptions in sleep. + * @param milliseconds + */ + void sleep(int milliseconds) { + try { Thread.sleep(milliseconds); } + catch (InterruptedException ex) {} + } + + void init(OptionParser parser) { + this.parser = parser; + parser.syncLogLevel(Logger.getLogger(ReconMutil.class.getPackage().getName())); + maxEvents = parser.getOption("-n").intValue(); + skipEvents = parser.getOption("-s").intValue(); + serial = new SerialHoncho(schema); + engines = new LinkedHashMap<>(); + if (!parser.getOption("-y").isDefault()) { + yaml = new ClaraYaml(parser.getOption("-y").stringValue()); + for (JSONObject service : yaml.services()) { + JSONObject cfg = yaml.filter(service.getString("name")); + if (cfg.length() > 0) addEngine(service.getString("name"), service.getString("class"), cfg); + else addEngine(service.getString("name"), service.getString("class"), null); + } + } + else if (!parser.getOption("-c").isDefault()) { + for (String s : parser.getOption("-c").stringValue().split(",")) + addEngine(null, s, null); + } + else { + InputStream is = Thread.currentThread().getContextClassLoader().getResourceAsStream("services.txt"); + BufferedReader br = new BufferedReader(new InputStreamReader(is, StandardCharsets.UTF_8)); + try { + for (String line; (line=br.readLine()) != null;) + addEngine(line.split(" ")[0],line.split(" ")[1],null); + } catch (IOException ex) { + System.getLogger(ReconMutil.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex); + } + } + if (!parser.getOption("-P").isDefault()) { + //ReconstructionEngine pp = addEngine("BG","org.jlab.service.postproc.PostprocEngine",null); + //pp.engineConfigMap.pu("preloadFile"); + //pp.engineConfigMap.pu("restream"); + //pp.engineConfigMap.pu("rebuild"); + } + if (!parser.getOption("-B").isDefault()) { + ReconstructionEngine bg = addEngine("BG","org.jlab.service.bg.BackgroundEngine",null); + bg.engineConfigMap.put("filename",parser.getOption("-B").stringValue()); + } + if (!parser.getOption("-S").isDefault()) { + } + } + + /** + * The command-line entry-point known as "recon-mutil". + * @param args command-line arguments + */ + public static void main(String[] args) { + OptionParser o = EngineProcessor.getParser(); + o.removeOption("-i"); + o.removeOption("-o"); + o.removeOption("-c"); + o.addOption("-t","4","number of threads"); + o.addOption("-o", null, "output file name"); + o.addOption("-c","2","comma-separated engine list"); + o.setRequiresInputList(true); + o.parse(args); + ReconMutil r = new ReconMutil(o); + r.launch(Arrays.stream(o.getOption("-t").stringValue().split(",")).mapToInt(Integer::parseInt).toArray(), + o.getOption("-o").stringValue(), + o.getInputList().stream().toArray(String[]::new)); + } + +} \ No newline at end of file From ea4cc853bc857a5acb754df367f17124c0864fa5 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 18:16:14 -0400 Subject: [PATCH 10/19] add recon-mutil --- bin/recon-mutil | 12 ++++++++++++ 1 file changed, 12 insertions(+) create mode 100755 bin/recon-mutil diff --git a/bin/recon-mutil b/bin/recon-mutil new file mode 100755 index 0000000000..3aaec2480b --- /dev/null +++ b/bin/recon-mutil @@ -0,0 +1,12 @@ +#!/bin/bash + +. `dirname $0`/../libexec/env.sh + +split_cli $@ + +export MALLOC_ARENA_MAX=1 + +java ${JAVA_OPTS-} -Xms10240m -XX:+UseParallelGC ${jvm_options[@]} \ + -cp ${COATJAVA_CLASSPATH:-''} \ + org.jlab.clas.reco.ReconMutil \ + ${class_options[@]} From 152295abb1e7eb4ee1c136ac6bd3babf83dfd0c6 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 18:50:12 -0400 Subject: [PATCH 11/19] fix --- .../java/org/jlab/clas/reco/ReconMutil.java | 3 +- .../resources/org/jlab/clas/reco/services.txt | 31 +++++++++++++++++++ .../org/jlab/utils/options/OptionValue.java | 4 +-- validation/advanced-tests/run-eb-tests.sh | 3 +- 4 files changed, 36 insertions(+), 5 deletions(-) create mode 100644 common-tools/clas-reco/src/main/resources/org/jlab/clas/reco/services.txt diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 9a8a21f746..8b23c04d69 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -146,7 +146,6 @@ void read(int threads, String... input) { // write leftover, partial chunk: if (!chunk.isEmpty()) { - System.err.println("writing partial chunk: "+chunk.size()); readEvents += chunk.size(); readQueue.offer(chunk); } @@ -452,7 +451,7 @@ else if (!parser.getOption("-c").isDefault()) { addEngine(null, s, null); } else { - InputStream is = Thread.currentThread().getContextClassLoader().getResourceAsStream("services.txt"); + InputStream is = ReconMutil.class.getClassLoader().getResourceAsStream("org/jlab/clas/reco/services.txt"); BufferedReader br = new BufferedReader(new InputStreamReader(is, StandardCharsets.UTF_8)); try { for (String line; (line=br.readLine()) != null;) diff --git a/common-tools/clas-reco/src/main/resources/org/jlab/clas/reco/services.txt b/common-tools/clas-reco/src/main/resources/org/jlab/clas/reco/services.txt new file mode 100644 index 0000000000..8a81bedfb6 --- /dev/null +++ b/common-tools/clas-reco/src/main/resources/org/jlab/clas/reco/services.txt @@ -0,0 +1,31 @@ +MAGFIELDS org.jlab.clas.swimtools.MagFieldsEngine +DENOISE org.jlab.service.ai.DCDenoiseEngine +FTCAL org.jlab.rec.ft.cal.FTCALEngine +FTHODO org.jlab.rec.ft.hodo.FTHODOEngine +FTTRK org.jlab.rec.ft.trk.FTTRKEngine +FTEB org.jlab.rec.ft.FTEBEngine +URWT org.jlab.service.urwt.URWTEngine +DCCR org.jlab.service.dc.DCHBClustering +DCHB org.jlab.service.dc.DCHBPostClusterConv +FTOFHB org.jlab.service.ftof.FTOFHBEngine +EC org.jlab.service.ec.ECEngine +RASTER org.jlab.service.raster.RasterEngine +CVT org.jlab.rec.cvt.services.CVTEngine +CTOF org.jlab.service.ctof.CTOFEngine +CND org.jlab.service.cnd.CNDCalibrationEngine +BAND org.jlab.service.band.BANDEngine +HTCC org.jlab.service.htcc.HTCCReconstructionService +LTCC org.jlab.service.ltcc.LTCCEngine +EBHB org.jlab.service.eb.EBHBEngine +DCTB org.jlab.service.dc.DCTBEngine +FMT org.jlab.service.fmt.FMTEngine +FTOFTB org.jlab.service.ftof.FTOFTBEngine +CVT org.jlab.rec.cvt.services.CVTSecondPassEngine +EBTB org.jlab.service.eb.EBTBEngine +RICHEB org.jlab.rec.rich.RICHEBEngine +RTPC org.jlab.service.rtpc.RTPCEngine +AHDC org.jlab.service.ahdc.AHDCEngine +ATOF org.jlab.service.atof.ATOFEngine +ALERT org.jlab.service.alert.ALERTEngine +MC org.jlab.service.mc.TruthMatch +VTX org.jlab.rec.service.vtx.VTXEngine diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionValue.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionValue.java index 684dba4aa7..dee09778c0 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionValue.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/options/OptionValue.java @@ -14,8 +14,8 @@ public class OptionValue { private String optionDefault = "0"; public OptionValue() { } - public OptionValue(String opt) { setOption(opt);} - public OptionValue(String opt,String value) { setOption(opt); setValue(value);} + public OptionValue(String opt) { setOption(opt);} + public OptionValue(String opt,String value) { setOption(opt); setValue(value); setDefault(value); } public final OptionValue setOption(String opt){ this.optionString = opt; return this;} public final OptionValue setValue(String value){ this.optionValue = value; return this;} diff --git a/validation/advanced-tests/run-eb-tests.sh b/validation/advanced-tests/run-eb-tests.sh index 05c7470db5..c95f338d2a 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 -l FINE -i ${input_dir}/${stub}.hipo -o out_${stub}.hipo -c 2 +echo ../../coatjava/bin/recon-mutil -t 6 -l FINE -o out_${stub}.hipo ${input_dir}/${stub}.hipo +exit # run EB tests: java -Xmx1536m -Xms1024m -cp $classPath -DINPUTFILE=out_${stub}.hipo eb.EBTwoTrackTest From dad364d533e21b6a7152722d75ae1592caa8b1ec Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 18:57:17 -0400 Subject: [PATCH 12/19] only printout if interval is positive --- .../java/org/jlab/utils/benchmark/ProgressPrintout.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java index 48da5624a0..767d279760 100644 --- a/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java +++ b/common-tools/clas-utils/src/main/java/org/jlab/utils/benchmark/ProgressPrintout.java @@ -55,14 +55,14 @@ public void showStatus(){ } public void updateStatus(){ - if (++this.numberOfCalls < WARMUP_CALLS) { + if (++this.numberOfCalls < WARMUP_CALLS){ this.previousPrintoutTime = System.currentTimeMillis(); this.startPrintoutTime = this.previousPrintoutTime; } - else { + else if(this.printoutIntervalSeconds > 0){ Long currentTime = System.currentTimeMillis(); Double elapsedTime = (currentTime - this.previousPrintoutTime)*1e-3; - if(elapsedTime >= this.printoutIntervalSeconds){ + if(elapsedTime >= this.printoutIntervalSeconds) { this.previousPrintoutTime = System.currentTimeMillis(); System.out.println(this.getUpdateString()); } From 60ce0df10646801dde6b1a870208479abd73a455 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 19:08:02 -0400 Subject: [PATCH 13/19] update ci for recon-mutil --- .github/workflows/ci.yml | 10 +++++----- .gitlab-ci.yml | 8 ++++++++ 2 files changed, 13 insertions(+), 5 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4967588f6e..2b0bc9ad69 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -227,12 +227,12 @@ jobs: tar xzvf clara.tar.gz tar xzvf coatjava.tar.gz - run: ls - - name: run test - run: ./bin/run-clara -y ./etc/services/rgd-clarode.yml -t 4 -n 500 -c ./clara -o ./tmp ./clas_018779.evio.00001 + - name: run clara + run: ./coatjava/bin/run-clara -y ./etc/services/rgd-clarode.yml -t 4 -n 100 -c ./clara -o ./tmp ./clas_018779.evio.00001 + - name: run mutil + run: ./coatjava/bin/recon-mutil -y ./etc/services/rgd-clarode.yml -t 4 -n 30 -o rec.hipo ./clas_018779.evio.00001 - name: ls tmp - run: ls -lhtr tmp - - name: rename - run: mv -v tmp/rec_clas_018779.evio.00001.hipo rec.hipo + run: ls -lhtr . tmp - uses: actions/upload-artifact@v7 with: name: test_clara_result diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml index d42ca40d5f..a74a9cbd3c 100644 --- a/.gitlab-ci.yml +++ b/.gitlab-ci.yml @@ -175,6 +175,14 @@ clara: - run-clara -v -c $CLARA_HOME -t 4 -y ./etc/services/rgd-clarode.yml -n 30 -o out $EVIOFILE - mv out/rec_$EVIOFILE.hipo claroded.hipo +recon-mutil: + stage: test + needs: [build,download] + dependencies: [build,download] + script: + - tar -xzf coatjava.tar.gz + - recon-mutil -t 4 -y ./etc/services/rgd-clarode.yml -n 30 -o rec_$EVIOFILE.hipo $EVIOFILE + profile: extends: .clon allow_failure: true From 919f999e71d7e4fbc525de53f5aae57616060733 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 19:38:59 -0400 Subject: [PATCH 14/19] bugfix --- .../src/main/java/org/jlab/clas/reco/ReconMutil.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 8b23c04d69..fad7fa8bdf 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -407,7 +407,7 @@ ReconstructionEngine addEngine(String label, String clazz, JSONObject cfg) { engine = (ReconstructionEngine) c.newInstance(); if (cfg != null && !cfg.toString().equals("null")) { EngineData input = new EngineData(); - input.setData(EngineDataType.JSON.mimeType(), cfg); + input.setData(EngineDataType.JSON.mimeType(), cfg.toString()); engine.configure(input); } else engine.init(); @@ -494,4 +494,4 @@ public static void main(String[] args) { o.getInputList().stream().toArray(String[]::new)); } -} \ No newline at end of file +} From d8eb34565f1b69d6ca34c5d1758291e5278a3aa7 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 19:51:39 -0400 Subject: [PATCH 15/19] quiet tarball extraction --- .github/workflows/ci.yml | 20 ++++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 2b0bc9ad69..aaedc49812 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -170,7 +170,7 @@ jobs: with: name: build_ubuntu-latest - name: untar build - run: tar xzvf coatjava.tar.gz + run: tar xzf coatjava.tar.gz - name: spotbugs run: ./build-coatjava.sh --spotbugs --no-maps --no-progress @@ -194,7 +194,7 @@ jobs: path: | clas_018779.evio.00001 - name: untar build - run: tar xzvf coatjava.tar.gz + run: tar xzf coatjava.tar.gz - name: run test run: | ls -lhtr @@ -224,8 +224,8 @@ jobs: clas_018779.evio.00001 - name: untar build run: | - tar xzvf clara.tar.gz - tar xzvf coatjava.tar.gz + tar xzf clara.tar.gz + tar xzf coatjava.tar.gz - run: ls - name: run clara run: ./coatjava/bin/run-clara -y ./etc/services/rgd-clarode.yml -t 4 -n 100 -c ./clara -o ./tmp ./clas_018779.evio.00001 @@ -277,8 +277,8 @@ jobs: name: build_${{ matrix.runner }} - name: untar build run: | - tar xzvf coatjava.tar.gz - tar xzvf clara.tar.gz + tar xzf coatjava.tar.gz + tar xzf clara.tar.gz - name: run test run: | git lfs install @@ -305,8 +305,8 @@ jobs: name: build_macos - name: untar build run: | - tar xzvf coatjava.tar.gz - tar xzvf clara.tar.gz + tar xzf coatjava.tar.gz + tar xzf clara.tar.gz - name: run test run: | git lfs install @@ -333,7 +333,7 @@ jobs: with: name: build_ubuntu-latest - name: untar build - run: tar xzvf coatjava.tar.gz + run: tar xzf coatjava.tar.gz - name: test run-groovy run: coatjava/bin/run-groovy validation/advanced-tests/test-run-groovy.groovy @@ -355,7 +355,7 @@ jobs: with: name: build_ubuntu-latest - name: untar build - run: tar xzvf coatjava.tar.gz + run: tar xzf coatjava.tar.gz - name: hipo2npz run: ./coatjava/bin/hipo2npz rec.hipo rec.npz - name: hipo2npz-dump From 12472a290d158e0ad4080c0b3214ebd4b63d5d74 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 20:16:25 -0400 Subject: [PATCH 16/19] prep pp --- .../org/jlab/detector/serial/PostProcessor.java | 16 ++++++++++++++++ .../main/java/org/jlab/clas/reco/ReconMutil.java | 11 +++++------ 2 files changed, 21 insertions(+), 6 deletions(-) diff --git a/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java index 0d0f1b96bc..a245029e17 100644 --- a/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java +++ b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/PostProcessor.java @@ -231,6 +231,22 @@ public void processEvent(Event event) { } } + public void processFile(String input, String output) { + Event event = new Event(); + HipoReader r = new HipoReader(); + r.open(input); + HipoWriterSorted w = new HipoWriterSorted(); + w.getSchemaFactory().initFromDirectory(ClasUtilsFile.getResourceDir("CLAS12DIR", "etc/bankdefs/hipo4")); + w.setCompressionType(2); + w.open(output); + while (r.hasNext()) { + r.nextEvent(event); + processEvent(event); + if (w != null) w.addEvent(event); + } + r.close(); + } + /** * The "postprocess" program. * @param args diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index fad7fa8bdf..9f7c26851a 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -21,6 +21,7 @@ import org.jlab.coda.jevio.EvioException; import org.jlab.detector.decode.CLASDecoder; import org.jlab.detector.decode.CLASDecoderPool; +import org.jlab.detector.serial.PostProcessor; import org.jlab.detector.serial.SerialHoncho; import org.jlab.io.evio.EvioDataEvent; import org.jlab.io.evio.EvioSource; @@ -110,6 +111,10 @@ void launch(int[] threads, String output, String... input) { reset(); } } + if (!parser.getOption("-P").isDefault()) { + PostProcessor pp = new PostProcessor(parser.getInputList(), false, false); + pp.processFile(output, output); + } } /** @@ -460,12 +465,6 @@ else if (!parser.getOption("-c").isDefault()) { System.getLogger(ReconMutil.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex); } } - if (!parser.getOption("-P").isDefault()) { - //ReconstructionEngine pp = addEngine("BG","org.jlab.service.postproc.PostprocEngine",null); - //pp.engineConfigMap.pu("preloadFile"); - //pp.engineConfigMap.pu("restream"); - //pp.engineConfigMap.pu("rebuild"); - } if (!parser.getOption("-B").isDefault()) { ReconstructionEngine bg = addEngine("BG","org.jlab.service.bg.BackgroundEngine",null); bg.engineConfigMap.put("filename",parser.getOption("-B").stringValue()); From 1a852a7d176c43b7c99b8ae5e3ded2a51e0e4f5f Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 21:49:48 -0400 Subject: [PATCH 17/19] add another queue, to prime for postprocessing --- .../java/org/jlab/clas/reco/ReconMutil.java | 74 ++++++++++++------- 1 file changed, 49 insertions(+), 25 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 9f7c26851a..3e6ac161ad 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -65,8 +65,10 @@ final class ReconMutil { CompletableFuture readerThread; CompletableFuture writerThread; CompletableFuture rethreadThread; + ConcurrentLinkedQueue decoThreads = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue> procQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); // Static parameters: @@ -99,10 +101,13 @@ void launch(int[] threads, String output, String... input) { writerThread = CompletableFuture.runAsync(() -> { write(output); }); for (int i=0; i { decode(j); })); procThreads.offer(CompletableFuture.runAsync(() -> { process(j); })); } while (!writerThread.isDone()) { sleep(100); + for (CompletableFuture f : decoThreads) + if (f.isDone()) decoThreads.remove(f); for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); if (threads.length > 1 && rethreadThread == null && writeEvents > 100) { @@ -127,7 +132,7 @@ void read(int threads, String... input) { List inputs = new ArrayList<>(Arrays.asList(input)); // initialize the event chunk: - List chunk = new ArrayList<>(EVENTS_PER_CHUNK); + List output = new ArrayList<>(EVENTS_PER_CHUNK); // loop over input events: while ( (maxEvents < 1 || readEvents < maxEvents) && @@ -139,7 +144,7 @@ void read(int threads, String... input) { if (readQueue.size() > CHUNKS_PER_QUEUE*threads) sleep(1000); // read next event into chunk, and fill queue if chunk full: - else chunk = read(chunk); + else output = read(output); } // open the next input file: @@ -150,56 +155,75 @@ void read(int threads, String... input) { } // write leftover, partial chunk: - if (!chunk.isEmpty()) { - readEvents += chunk.size(); - readQueue.offer(chunk); + if (!output.isEmpty()) { + readEvents += output.size(); + readQueue.offer(output); } if (reader instanceof EvioSource evio) evio.close(); } /** - * The event processor thread. - * @param thread unique thread number + * The decoder thread. + * @param thread thread number */ - void process(int thread) { + void decode(int thread) { while (true) { - List o = readQueue.poll(); - if (o == null) { + List input = readQueue.poll(); + if (input == null) { if (readerThread.isDone() && readQueue.isEmpty() && writeEvents+skipEvents+failEvents >= readEvents) break; sleep(100); } + else { + List output = new ArrayList<>(input.size()); + for (int i=0; i input = procQueue.poll(); + if (input == null) { + if (decoThreads.isEmpty() && procQueue.isEmpty() && + writeEvents+skipEvents+failEvents >= readEvents) break; + sleep(100); + } else { // put the event back on the queue if we're rethreading: //if (rethreadThread != null && !rethreadThread.isDone()) readQueue.offer(o); - List chunk = new ArrayList<>(o.size()); - for (int i=0; i output = new ArrayList<>(input.size()); + for (int i=0; i engine : engines.entrySet()) { Benchmark.getInstance().resume(engine.getValue().getName()); - try { engine.getValue().processDataEvent(event); } + try { engine.getValue().processDataEvent(input.get(i)); } catch (Exception ex) { ex.printStackTrace(); } Benchmark.getInstance().pause(engine.getValue().getName()); } Benchmark.getInstance().resume("serial"); - Event e = event.getHipoEvent(); + Event e = input.get(i).getHipoEvent(); Event t = serial.read(e); t.setEventTag(1); Benchmark.getInstance().pause("serial"); - chunk.add(e); - chunk.add(t); + output.add(e); + output.add(t); } - writeQueue.offer(chunk); + writeQueue.offer(output); } } } From b6b715d7115b71a8a1716871f024a60d1f859d6b Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 21:56:41 -0400 Subject: [PATCH 18/19] cleanup --- .../src/main/java/org/jlab/clas/reco/ReconMutil.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 3e6ac161ad..48de0c20e7 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -61,12 +61,14 @@ final class ReconMutil { Map engines = new LinkedHashMap<>(); CLASDecoderPool decoders = new CLASDecoderPool(64,"default",null); - // Threads and queues: + // Threads: CompletableFuture readerThread; CompletableFuture writerThread; CompletableFuture rethreadThread; ConcurrentLinkedQueue decoThreads = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); + + // Queues: ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> procQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); From 2afd58d299dcdd22758c51de2d5932da6a2a89b2 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 4 Sep 2026 23:14:52 -0400 Subject: [PATCH 19/19] cleanup --- .../src/main/java/org/jlab/clas/reco/ReconMutil.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index 48de0c20e7..b0776004bc 100644 --- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java +++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java @@ -49,7 +49,7 @@ final class ReconMutil { final int CHUNKS_PER_QUEUE = 100; final int EVENTS_PER_CHUNK = 100; - // File reader and writer: + // File I/O: Object reader; HipoWriterSorted writer; List schemaBankList; @@ -58,8 +58,8 @@ final class ReconMutil { // Processors: SerialHoncho serial; - Map engines = new LinkedHashMap<>(); CLASDecoderPool decoders = new CLASDecoderPool(64,"default",null); + Map engines = new LinkedHashMap<>(); // Threads: CompletableFuture readerThread;