From c2007e39168df103eb66fc1a6fe114fe60926e64 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 3 Sep 2026 23:08:18 -0400 Subject: [PATCH 01/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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/33] 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; From 67cce53aab48f4906cb8cac28aac1c8b83ff4198 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Wed, 9 Sep 2026 19:58:43 -0400 Subject: [PATCH 20/33] null yaml bugfix --- .../java/org/jlab/clas/reco/ReconMutil.java | 50 ++++++++++++------- validation/advanced-tests/run-eb-tests.sh | 2 +- 2 files changed, 32 insertions(+), 20 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 b0776004bc..1e824771e5 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 @@ -98,7 +98,10 @@ final class ReconMutil { * @param input names of input files to read */ void launch(int[] threads, String output, String... input) { + reset(); + + // spawn all the threads: readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); writerThread = CompletableFuture.runAsync(() -> { write(output); }); for (int i=0; i { decode(j); })); procThreads.offer(CompletableFuture.runAsync(() -> { process(j); })); } + + // wait for the writer to be done: while (!writerThread.isDone()) { sleep(100); + + // cleanup completed parallel threads: for (CompletableFuture f : decoThreads) if (f.isDone()) decoThreads.remove(f); for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); + + // perform scaling test: if (threads.length > 1 && rethreadThread == null && writeEvents > 100) { rethreadThread = CompletableFuture.runAsync(() -> { rethread(BENCH_SECONDS,threads); }); rethreadThread.join(); reset(); } } - if (!parser.getOption("-P").isDefault()) { - PostProcessor pp = new PostProcessor(parser.getInputList(), false, false); - pp.processFile(output, output); - } + + //if (!parser.getOption("-P").isDefault()) { + // PostProcessor pp = new PostProcessor(parser.getInputList(), false, false); + // pp.processFile(output, output); + //} } /** @@ -339,23 +349,25 @@ 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 (yaml != null && 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(); + if (json.has("wildcard")) { + if (json.optBoolean("schema_filter",true)) { + int schemaSize = w.getSchemaFactory().getSchemaList().size(); + for (int i=0; i Date: Thu, 10 Sep 2026 14:59:32 -0400 Subject: [PATCH 21/33] remove throwers and catchers --- .../org/jlab/detector/decode/CLASDecoderPool.java | 4 ++-- .../src/main/java/org/jlab/clas/reco/ReconMutil.java | 11 ++++------- 2 files changed, 6 insertions(+), 9 deletions(-) diff --git a/common-tools/clas-detector/src/main/java/org/jlab/detector/decode/CLASDecoderPool.java b/common-tools/clas-detector/src/main/java/org/jlab/detector/decode/CLASDecoderPool.java index 80430583ee..a28135d906 100644 --- a/common-tools/clas-detector/src/main/java/org/jlab/detector/decode/CLASDecoderPool.java +++ b/common-tools/clas-detector/src/main/java/org/jlab/detector/decode/CLASDecoderPool.java @@ -34,11 +34,11 @@ public CLASDecoderPool(int size, String variation, String timestamp) { } } - public CLASDecoder take() throws InterruptedException { + public CLASDecoder take() { return pool.poll(); } - public void put(CLASDecoder decoder) throws InterruptedException { + public void put(CLASDecoder decoder) { pool.offer(decoder); } 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 1e824771e5..1e7be24bb6 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 @@ -233,6 +233,7 @@ void process(int thread) { t.setEventTag(1); Benchmark.getInstance().pause("serial"); output.add(e); + if (!t.isEmpty()) output.add(t); } writeQueue.offer(output); @@ -311,13 +312,9 @@ HipoDataEvent decode(ByteBuffer bytes) { 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; } + CLASDecoder d = decoders.take(); + HipoDataEvent hipo = d.getDecodedDataEvenet(evio); + decoders.put(d); Benchmark.getInstance().pause("deco"); return hipo; } From 6574f292264de42715e3bbc8a619aa1842bc6427 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 17:23:34 -0400 Subject: [PATCH 22/33] add postproc to serial --- .../jlab/detector/serial/SerialHoncho.java | 105 +++++++++++++++--- 1 file changed, 87 insertions(+), 18 deletions(-) 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 index 86a4868432..2b057a5ba2 100644 --- 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 @@ -4,8 +4,11 @@ import java.util.TreeSet; import org.jlab.detector.calib.utils.ConstantsManager; import org.jlab.detector.decode.CLASDecoder; +import org.jlab.detector.helicity.HelicityBit; import org.jlab.detector.helicity.HelicitySequence; +import org.jlab.detector.helicity.HelicitySequenceDelayed; import org.jlab.detector.helicity.HelicityState; +import org.jlab.detector.scalers.DaqScalers; import org.jlab.detector.scalers.DaqScalersSequence; import org.jlab.jnp.hipo4.data.Bank; import org.jlab.jnp.hipo4.data.Event; @@ -20,15 +23,6 @@ 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(); @@ -55,6 +49,22 @@ public synchronized Event read(Event event) { helicities.add(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman)); return CLASDecoder.createTaggedEvent(event, runConfig, tag1banks); } + + public void process(Event event) { + Bank cfg = new Bank(schema.getSchema("RUN::config")); + Bank evt = new Bank(schema.getSchema("REC::Event")); + event.read(cfg); + event.read(evt); + if (cfg.getRows() > 0) { + processEventUnix(event, cfg); + if (evt.getRows() > 0) { + event.remove(evt.getSchema()); + processHelicity(event, cfg, evt); + processScalers(cfg, evt); + event.write(evt); + } + } + } public void finish(HipoWriterSorted writer) { writer.addEvent(getUnixEvent(runConfig),1); @@ -66,6 +76,40 @@ public void clear() { scalers.clear(100); } + public DaqScalersSequence getScalers() { + return scalers; + } + + public TreeSet getHelicities() { + return helicities; + } + + public HelicitySequenceDelayed getHelicitySequence() { + // FIXME: autogenerate if necessary + HelicitySequenceDelayed h = new HelicitySequenceDelayed( + conman.getConstants(4013, "/runcontrol/helicity").getIntValue("delay",0,0,0)); + h.addStream(helicities); + return h; + } + + public ConstantsManager getConstantsManager() { + return conman; + } + + public SchemaFactory getSchemaFactory() { + return schema; + } + + SchemaFactory schema; + Bank[] tag1banks; + // FIXME: store Schema for banks; + Bank runConfig; + Bank helicityAdc; + ConstantsManager conman; + TreeMap eventUnix; + TreeSet helicities; + DaqScalersSequence scalers; + Event getUnixEvent(Bank config) { Bank unix = new Bank(schema.getSchema("RUN::unix")); unix.setRows(eventUnix.size()); @@ -81,19 +125,44 @@ Event getUnixEvent(Bank config) { return e; } - public DaqScalersSequence getScalers() { - return scalers; + int getUnixTime(Bank runConfig) { + if (runConfig.getRows() < 1) { + Integer key = eventUnix.floorKey(runConfig.getInt("event",0)); + if (key != null) { + Integer unix = eventUnix.get(key); + if (unix != null) return unix; + } + } + return 0; } - - public TreeSet getHelicities() { - return helicities; + + void processEventUnix(Event event, Bank runConfig) { + int ut = getUnixTime(runConfig); + event.remove(runConfig.getSchema()); + runConfig.putInt("unixtime", 0, ut); + event.write(runConfig); } - public ConstantsManager getConstantsManager() { - return conman; + void processScalers(Bank runConfig, Bank recEvent) { + DaqScalers ds = scalers.get(runConfig.getLong("timestamp", 0)); + if (ds != null) { + recEvent.putFloat("beamCharge",0, (float) ds.dsc2.getBeamChargeGated()); + recEvent.putDouble("liveTime",0,ds.dsc2.getLivetime()); + } } - public SchemaFactory getSchemaFactory() { - return schema; + void processHelicity(Event event, Bank runcfg, Bank recevt) { + HelicityBit hb = getHelicitySequence().search(runcfg.getLong("timestamp", 0)); + HelicityBit hbraw = getHelicitySequence().getHalfWavePlate() ? HelicityBit.getFlipped(hb) : hb; + recevt.putByte("helicity",0,hb.value()); + recevt.putByte("helicityRaw",0,hbraw.value()); + Bank helScaler = new Bank(schema.getSchema("HEL::scaler")); + event.read(helScaler); + if (helScaler.getRows()>0) { + event.remove(schema.getSchema("HEL::scaler")); + SerialUtil.assignScalerHelicity(runcfg.getLong("timestamp",0), helScaler, getHelicitySequence()); + event.write(helScaler); + } } + } From 252abb05d36e445d8b0472fb14aa3e5f0f6e445f Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 17:24:32 -0400 Subject: [PATCH 23/33] replace duplicate serial.read w/ read+process --- .../java/org/jlab/clas/reco/ReconMutil.java | 33 ++++++++++--------- 1 file changed, 18 insertions(+), 15 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 1e7be24bb6..d7ecbaf7a7 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 @@ -14,6 +14,7 @@ import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.atomic.AtomicInteger; import java.util.logging.Level; import java.util.logging.Logger; import org.jlab.clara.engine.EngineData; @@ -21,7 +22,6 @@ 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; @@ -39,7 +39,8 @@ import org.json.JSONObject; /** - * + * FIXME: add tagged bank counter for completino decision + * * @author baltzell */ final class ReconMutil { @@ -82,6 +83,7 @@ final class ReconMutil { // Progress counters: int readEvents; int writeEvents; + AtomicInteger taggedEvents; int failEvents; int fileEvents; int maxFileEvents; @@ -193,9 +195,14 @@ void decode(int thread) { HipoDataEvent event = input.get(i) instanceof ByteBuffer ? decode((ByteBuffer)input.get(i)) : new HipoDataEvent(((Event)input.get(i)), schema); - Event taggedEvent = serial.read(event.getHipoEvent()); - if (!taggedEvent.isEmpty()) output.add(new HipoDataEvent(taggedEvent, schema)); output.add(event); + Benchmark.getInstance().resume("serial"); + Event taggedEvent = serial.read(event.getHipoEvent()); + if (!taggedEvent.isEmpty()) { + output.add(new HipoDataEvent(taggedEvent, schema)); + taggedEvents.incrementAndGet(); + } + Benchmark.getInstance().pause("serial"); } procQueue.offer(output); } @@ -215,26 +222,20 @@ void process(int thread) { sleep(100); } else { - // put the event back on the queue if we're rethreading: //if (rethreadThread != null && !rethreadThread.isDone()) readQueue.offer(o); List output = new ArrayList<>(input.size()); for (int i=0; i engine : engines.entrySet()) { Benchmark.getInstance().resume(engine.getValue().getName()); try { engine.getValue().processDataEvent(input.get(i)); } catch (Exception ex) { ex.printStackTrace(); } Benchmark.getInstance().pause(engine.getValue().getName()); } - - Benchmark.getInstance().resume("serial"); Event e = input.get(i).getHipoEvent(); - Event t = serial.read(e); - t.setEventTag(1); - Benchmark.getInstance().pause("serial"); + Benchmark.getInstance().resume(thread,"post"); + serial.process(e); + Benchmark.getInstance().pause(thread,"post"); output.add(e); - if (!t.isEmpty()) - output.add(t); } writeQueue.offer(output); } @@ -410,8 +411,9 @@ List read(List chunk) { void close() { serial.finish(writer); writer.close(); - System.out.println(String.format("recon-mutil :: read/write/diff = %d/%d/%d", - readEvents, writeEvents, readEvents-writeEvents)); + System.out.println(Benchmark.getInstance()); + System.out.println(String.format("recon-mutil :: read/write/tagged/diff = %d/%d/%d/%d", + readEvents, writeEvents, taggedEvents.get(), writeEvents-readEvents)); } /** @@ -430,6 +432,7 @@ void reset() { readEvents = 0; writeEvents = 0; failEvents = 0; + taggedEvents.set(0); } /** From d0c451fd0f6815b95882cbe45e8fbeac18a8a8d7 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 19:00:00 -0400 Subject: [PATCH 24/33] remove post-processing from Clas12Writer --- common-tools/clara-io/pom.xml | 8 +-- .../java/org/jlab/io/clara/Clas12Writer.java | 56 ------------------- 2 files changed, 1 insertion(+), 63 deletions(-) diff --git a/common-tools/clara-io/pom.xml b/common-tools/clara-io/pom.xml index f1c3f6f0e1..3c9ddb5948 100644 --- a/common-tools/clara-io/pom.xml +++ b/common-tools/clara-io/pom.xml @@ -36,12 +36,6 @@ jnp-hipo4 - - org.jlab.clas - clas-io - 14.2.0-SNAPSHOT - - org.jlab.clas clas-detector @@ -50,7 +44,7 @@ org.jlab.clas - clas-utils + clas-io 14.2.0-SNAPSHOT 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 edf5596040..5d5a917add 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 @@ -1,21 +1,11 @@ package org.jlab.io.clara; -import java.io.File; import java.nio.file.Path; -import java.util.List; -import java.util.TreeMap; -import java.util.TreeSet; import org.jlab.clara.std.services.EventWriterException; import org.jlab.detector.calib.utils.ConstantsManager; -import org.jlab.detector.helicity.HelicitySequenceDelayed; -import org.jlab.detector.helicity.HelicityState; -import org.jlab.detector.scalers.DaqScalersSequence; -import org.jlab.detector.serial.PostProcessor; 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; -import org.jlab.jnp.hipo4.io.HipoReader; import org.jlab.jnp.hipo4.io.HipoWriterSorted; import org.jlab.jnp.utils.file.FileUtils; import org.json.JSONObject; @@ -25,7 +15,6 @@ * 1. Copies certain banks on-the-fly to new tag-1 events * 2. Caches helicity states, scaler readouts, and unix time * 3. Writes HEL::flip, RUN/HEL::scaler, and RUN::unix to new tag-1 events - * 4. Runs post-processing, writing tag-1 information to all events * 5. Adds .hipo to the output filename, if necessary * * @author baltzell @@ -33,19 +22,15 @@ public class Clas12Writer extends HipoToHipoWriter { SerialHoncho serial; - Bank runConfig; ConstantsManager conman; 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")); conman = new ConstantsManager(); 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")); } @@ -74,47 +59,6 @@ protected void writeEvent(Object event) throws EventWriterException { protected void closeWriter() { serial.finish(writer); super.closeWriter(); - if (postprocess) postprocess(); serial.clear(); } - - /** - * Get the first valid run number from a RUN::config bank. - * @return run - */ - private int getRunNumber() { - Event e = new Event(); - HipoReader r = new HipoReader(); - r.open(filename); - while (r.hasNext()) { - r.nextEvent(e); - e.read(runConfig); - if (runConfig.getRows()>0 && runConfig.getInt("run",0)>0) - return runConfig.getInt("run",0); - } - return 0; - } - - /** - * 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(serial.getHelicities()); - PostProcessor p = new PostProcessor(List.of(filename), fullSchema, helicity, serial.getScalers()); - HipoReader r = new HipoReader(); - r.open(filename); - Event e = new Event(); - writer.open("pp_"+filename); - while (r.hasNext()) { - r.nextEvent(e); - p.processEvent(e); - HipoToHipoWriter.writeEvent(writer, e, schemaBankList); - } - writer.close(); - new File(filename).delete(); - new File("pp_"+filename).renameTo(new File(filename)); - } - } From e17fcd6af2ccc7fa1cdb29dd08ab0ef7045d0c92 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 19:00:30 -0400 Subject: [PATCH 25/33] convert to HelicitySequence, oopts --- .../jlab/detector/serial/SerialHoncho.java | 43 +++++++++---------- 1 file changed, 20 insertions(+), 23 deletions(-) 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 index 2b057a5ba2..b86fc5040b 100644 --- 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 @@ -1,7 +1,6 @@ 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.HelicityBit; @@ -29,7 +28,6 @@ public SerialHoncho(SchemaFactory schema) { 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]; @@ -42,11 +40,17 @@ public synchronized Event read(Event event) { event.read(runConfig); event.read(helicityAdc); if (runConfig.getRows() > 0) { + if (run <= 0 && runConfig.getInt("run", 0) > 0) { + run = runConfig.getInt("run",0); + helicitySequence = new HelicitySequenceDelayed( + conman.getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,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)); + if (helicitySequence != null) + helicitySequence.addState(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman)); return CLASDecoder.createTaggedEvent(event, runConfig, tag1banks); } @@ -68,11 +72,11 @@ public void process(Event event) { public void finish(HipoWriterSorted writer) { writer.addEvent(getUnixEvent(runConfig),1); - HelicitySequence.writeFlips(schema, writer, helicities); + helicitySequence.writeFlips(writer, 1); } public void clear() { - while (helicities.size() > 100) helicities.pollFirst(); + //helicitySequence.clear(100); scalers.clear(100); } @@ -80,16 +84,8 @@ public DaqScalersSequence getScalers() { return scalers; } - public TreeSet getHelicities() { - return helicities; - } - - public HelicitySequenceDelayed getHelicitySequence() { - // FIXME: autogenerate if necessary - HelicitySequenceDelayed h = new HelicitySequenceDelayed( - conman.getConstants(4013, "/runcontrol/helicity").getIntValue("delay",0,0,0)); - h.addStream(helicities); - return h; + public HelicitySequence getHelicities() { + return helicitySequence; } public ConstantsManager getConstantsManager() { @@ -99,7 +95,7 @@ public ConstantsManager getConstantsManager() { public SchemaFactory getSchemaFactory() { return schema; } - + SchemaFactory schema; Bank[] tag1banks; // FIXME: store Schema for banks; @@ -107,8 +103,9 @@ public SchemaFactory getSchemaFactory() { Bank helicityAdc; ConstantsManager conman; TreeMap eventUnix; - TreeSet helicities; + HelicitySequence helicitySequence; DaqScalersSequence scalers; + int run; Event getUnixEvent(Bank config) { Bank unix = new Bank(schema.getSchema("RUN::unix")); @@ -151,16 +148,16 @@ void processScalers(Bank runConfig, Bank recEvent) { } } - void processHelicity(Event event, Bank runcfg, Bank recevt) { - HelicityBit hb = getHelicitySequence().search(runcfg.getLong("timestamp", 0)); - HelicityBit hbraw = getHelicitySequence().getHalfWavePlate() ? HelicityBit.getFlipped(hb) : hb; - recevt.putByte("helicity",0,hb.value()); - recevt.putByte("helicityRaw",0,hbraw.value()); + void processHelicity(Event event, Bank runConfig, Bank recEvent) { + HelicityBit hb = helicitySequence.search(runConfig.getLong("timestamp", 0)); + HelicityBit hbraw = helicitySequence.getHalfWavePlate() ? HelicityBit.getFlipped(hb) : hb; + recEvent.putByte("helicity",0,hb.value()); + recEvent.putByte("helicityRaw",0,hbraw.value()); Bank helScaler = new Bank(schema.getSchema("HEL::scaler")); event.read(helScaler); if (helScaler.getRows()>0) { event.remove(schema.getSchema("HEL::scaler")); - SerialUtil.assignScalerHelicity(runcfg.getLong("timestamp",0), helScaler, getHelicitySequence()); + SerialUtil.assignScalerHelicity(runConfig.getLong("timestamp",0), helScaler, helicitySequence); event.write(helScaler); } } From 7d62e4d06c7f5595cb96713975fa8d6040da452d Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 19:00:54 -0400 Subject: [PATCH 26/33] prime the serial buffer --- .../java/org/jlab/clas/reco/ReconMutil.java | 28 +++++++++++-------- 1 file changed, 17 insertions(+), 11 deletions(-) diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java index d7ecbaf7a7..5c9942dc85 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 @@ -83,10 +83,10 @@ final class ReconMutil { // Progress counters: int readEvents; int writeEvents; - AtomicInteger taggedEvents; int failEvents; int fileEvents; int maxFileEvents; + AtomicInteger taggedEvents = new AtomicInteger(); ProgressPrintout progress = new ProgressPrintout(); ReconMutil(OptionParser parser) { @@ -102,6 +102,8 @@ final class ReconMutil { void launch(int[] threads, String output, String... input) { reset(); + + System.out.println(String.format("recon-mutil:: Spawning %d+++ Threads...",threads[0])); // spawn all the threads: readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); @@ -114,7 +116,12 @@ void launch(int[] threads, String output, String... input) { // wait for the writer to be done: while (!writerThread.isDone()) { - sleep(100); + sleep(1000); + + //System.out.println(String.format("recon-mutil:: read(%b)/[deco(%d)]/proc(%d)/tag/write(%b)", + // readerThread.isDone(), decoThreads.size(), procThreads.size(), writerThread.isDone())); + //System.out.println(String.format("recon-util:: %d-%d/%d/%d/%d", readEvents, + // readQueue.size(), procQueue.size(), taggedEvents.get(), writeQueue.size())); // cleanup completed parallel threads: for (CompletableFuture f : decoThreads) @@ -129,11 +136,6 @@ 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); - //} } /** @@ -155,7 +157,7 @@ void read(int threads, String... input) { if (reader != null) { // sleep instead of overfilling the read queue: - if (readQueue.size() > CHUNKS_PER_QUEUE*threads) sleep(1000); + if (false) sleep(100);//readQueue.size() > CHUNKS_PER_QUEUE*threads) sleep(1000); // read next event into chunk, and fill queue if chunk full: else output = read(output); @@ -215,6 +217,10 @@ void decode(int thread) { */ void process(int thread) { while (true) { + if (serial.getScalers().size() < 10 || taggedEvents.get() < 100) { + sleep(100); + continue; + } List input = procQueue.poll(); if (input == null) { if (decoThreads.isEmpty() && procQueue.isEmpty() && @@ -232,9 +238,6 @@ void process(int thread) { Benchmark.getInstance().pause(engine.getValue().getName()); } Event e = input.get(i).getHipoEvent(); - Benchmark.getInstance().resume(thread,"post"); - serial.process(e); - Benchmark.getInstance().pause(thread,"post"); output.add(e); } writeQueue.offer(output); @@ -259,6 +262,9 @@ void write(String output) { } else { for (int i=0; i 0 || schemaBankList.isEmpty()) From a9e1a2a9dbb614b205f88befb57fa56b6dbb73e1 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 19:03:18 -0400 Subject: [PATCH 27/33] set read queue limit at 100k events, 2 GB --- .../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 5c9942dc85..1c69b509d8 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 @@ -156,8 +156,8 @@ void read(int threads, String... input) { if (reader != null) { - // sleep instead of overfilling the read queue: - if (false) sleep(100);//readQueue.size() > CHUNKS_PER_QUEUE*threads) sleep(1000); + // sleep instead of overfilling the read queue (100K events, ~2GB): + if (readEvents > 1e5) sleep(1000); // read next event into chunk, and fill queue if chunk full: else output = read(output); From 5d4b5877b4ae320930bbb76fc50435b633580881 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 22:18:31 -0400 Subject: [PATCH 28/33] cleanup --- .../jlab/detector/serial/SerialHoncho.java | 53 ++++++++++++------- 1 file changed, 35 insertions(+), 18 deletions(-) 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 index b86fc5040b..6fe28171ef 100644 --- 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 @@ -1,6 +1,7 @@ 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.HelicityBit; @@ -21,7 +22,17 @@ 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; // FIXME: store Schema for banks; + Bank helicityAdc; + ConstantsManager conman; + TreeMap eventUnix; + HelicitySequence helicitySequence; + TreeSet helicities; + DaqScalersSequence scalers; + int run; + public SerialHoncho(SchemaFactory schema) { this.schema = schema; conman = new ConstantsManager(); @@ -29,6 +40,7 @@ public SerialHoncho(SchemaFactory schema) { runConfig = new Bank(schema.getSchema("RUN::config")); helicityAdc = new Bank(schema.getSchema("HEL::adc")); scalers = new DaqScalersSequence(schema); + helicities = new TreeSet<>(); eventUnix = new TreeMap<>(); tag1banks = new Bank[TAG1BANKS.length]; for (int i=0; i 0 && evno > 0) eventUnix.put(evno, unix); } - if (helicitySequence != null) - helicitySequence.addState(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman)); + if (helicitySequence != null) { + HelicityState state = HelicityState.createFromFadcBank(helicityAdc, runConfig, conman); + helicities.add(state); + helicitySequence.addState(state); + } return CLASDecoder.createTaggedEvent(event, runConfig, tag1banks); } @@ -72,19 +87,21 @@ public void process(Event event) { public void finish(HipoWriterSorted writer) { writer.addEvent(getUnixEvent(runConfig),1); + // FIXME: mark written flips and don't write them again helicitySequence.writeFlips(writer, 1); } public void clear() { - //helicitySequence.clear(100); - scalers.clear(100); + eventUnix.clear(); + helicities.clear(); + scalers.clear(); } - + public DaqScalersSequence getScalers() { return scalers; } - public HelicitySequence getHelicities() { + public HelicitySequence getHelicitySequence() { return helicitySequence; } @@ -95,18 +112,18 @@ public ConstantsManager getConstantsManager() { public SchemaFactory getSchemaFactory() { return schema; } + + public TreeSet getHelicities() { + return helicities; + } + + HelicitySequence createHelicitySequence() { + HelicitySequence seq = new HelicitySequenceDelayed( + conman.getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,0)); + seq.addStream(helicities); + return seq; + } - SchemaFactory schema; - Bank[] tag1banks; - // FIXME: store Schema for banks; - Bank runConfig; - Bank helicityAdc; - ConstantsManager conman; - TreeMap eventUnix; - HelicitySequence helicitySequence; - DaqScalersSequence scalers; - int run; - Event getUnixEvent(Bank config) { Bank unix = new Bank(schema.getSchema("RUN::unix")); unix.setRows(eventUnix.size()); From e7100b4ad5e15ef52902d41b4c4410d081680a0f Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Thu, 10 Sep 2026 22:19:03 -0400 Subject: [PATCH 29/33] add pausing --- .../clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java | 3 ++- 1 file changed, 2 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 1c69b509d8..1ae7d07684 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 @@ -73,6 +73,7 @@ final class ReconMutil { ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> procQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); + boolean paused = false; // Static parameters: int maxEvents; @@ -217,7 +218,7 @@ void decode(int thread) { */ void process(int thread) { while (true) { - if (serial.getScalers().size() < 10 || taggedEvents.get() < 100) { + if (paused || serial.getScalers().size() < 10 || taggedEvents.get() < 100) { sleep(100); continue; } From 37f0766f62407781e0f113abc32ef0a84f89d979 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 11 Sep 2026 11:29:48 -0400 Subject: [PATCH 30/33] cleanup --- .../clas-reco/src/main/java/org/jlab/clas/reco/ReconMutil.java | 2 -- 1 file changed, 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 1ae7d07684..d9a58eeeee 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 @@ -39,7 +39,6 @@ import org.json.JSONObject; /** - * FIXME: add tagged bank counter for completino decision * * @author baltzell */ @@ -47,7 +46,6 @@ final class ReconMutil { // Performance parameters: final int BENCH_SECONDS = 30; - final int CHUNKS_PER_QUEUE = 100; final int EVENTS_PER_CHUNK = 100; // File I/O: From b5cbe3bdd51bfa52574e7d983977f8ad57f9683f Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 11 Sep 2026 19:53:56 -0400 Subject: [PATCH 31/33] working state (with yaml) --- .../jlab/detector/serial/SerialHoncho.java | 119 +++++++++++------- .../java/org/jlab/clas/reco/ReconMutil.java | 36 +++--- 2 files changed, 94 insertions(+), 61 deletions(-) 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 index 6fe28171ef..59754e8172 100644 --- 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 @@ -1,7 +1,12 @@ package org.jlab.detector.serial; +import java.util.Arrays; +import java.util.Iterator; +import java.util.List; +import java.util.ListIterator; import java.util.TreeMap; import java.util.TreeSet; +import java.util.stream.Collectors; import org.jlab.detector.calib.utils.ConstantsManager; import org.jlab.detector.decode.CLASDecoder; import org.jlab.detector.helicity.HelicityBit; @@ -12,6 +17,7 @@ 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.Schema; import org.jlab.jnp.hipo4.data.SchemaFactory; import org.jlab.jnp.hipo4.io.HipoWriterSorted; @@ -23,71 +29,85 @@ 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; // FIXME: store Schema for banks; - Bank helicityAdc; + Schema[] tag1banks; + Schema runConfig; + Schema recEvent; + Schema helScaler; + Schema helicityAdc; ConstantsManager conman; TreeMap eventUnix; HelicitySequence helicitySequence; TreeSet helicities; DaqScalersSequence scalers; - int run; + int run = 0; 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")); + runConfig = schema.getSchema("RUN::config"); + recEvent = schema.getSchema("REC::Event"); + helicityAdc = schema.getSchema("HEL::adc"); + helScaler = schema.getSchema("HEL::scaler"); scalers = new DaqScalersSequence(schema); helicities = new TreeSet<>(); eventUnix = new TreeMap<>(); - tag1banks = new Bank[TAG1BANKS.length]; + tag1banks = new Schema[TAG1BANKS.length]; for (int i=0; i 0) { - if (run <= 0 && runConfig.getInt("run", 0) > 0) { - run = runConfig.getInt("run",0); - helicitySequence = new HelicitySequenceDelayed( - conman.getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,0)); - } - int unix = runConfig.getInt("unixtime",0); - int evno = runConfig.getInt("event",0); + event.read(cfg); + event.read(hel); + if (cfg.getRows() > 0) { + if (run <= 0 && cfg.getInt("run", 0) > 0) + run = cfg.getInt("run",0); + int unix = cfg.getInt("unixtime",0); + int evno = cfg.getInt("event",0); if (unix > 0 && evno > 0) eventUnix.put(evno, unix); } - if (helicitySequence != null) { - HelicityState state = HelicityState.createFromFadcBank(helicityAdc, runConfig, conman); - helicities.add(state); - helicitySequence.addState(state); - } - return CLASDecoder.createTaggedEvent(event, runConfig, tag1banks); + helicities.add(HelicityState.createFromFadcBank(hel, cfg, conman)); + return CLASDecoder.createTaggedEvent(event, cfg, createTaggedBanks(tag1banks)); } - + public void process(Event event) { - Bank cfg = new Bank(schema.getSchema("RUN::config")); - Bank evt = new Bank(schema.getSchema("REC::Event")); + Bank cfg = new Bank(runConfig); + Bank evt = new Bank(recEvent); event.read(cfg); event.read(evt); if (cfg.getRows() > 0) { processEventUnix(event, cfg); if (evt.getRows() > 0) { event.remove(evt.getSchema()); - processHelicity(event, cfg, evt); + //processHelicity(event, cfg, evt); processScalers(cfg, evt); event.write(evt); } } } + public void prune() { + // remove identical helicities in the stream: + HelicityState prev = null; + Iterator iter = (ListIterator)helicities.iterator(); + while (iter.hasNext()) { + HelicityState next = iter.next(); + if (prev != null && prev == next) + helicities.remove(next); + } + scalers.clear((int)1e5); + // trim helicities to 100 million events, ~1 run, ~1 GB: + //while (helicities.size() < 1e8) helicities.pollFirst(); + } + public void finish(HipoWriterSorted writer) { - writer.addEvent(getUnixEvent(runConfig),1); - // FIXME: mark written flips and don't write them again + Bank cfg = new Bank(runConfig, 1); + cfg.putInt("run",0,run); + writer.addEvent(getUnixEvent(cfg),1); helicitySequence.writeFlips(writer, 1); } @@ -95,15 +115,12 @@ public void clear() { eventUnix.clear(); helicities.clear(); scalers.clear(); + helicitySequence = null; } public DaqScalersSequence getScalers() { return scalers; } - - public HelicitySequence getHelicitySequence() { - return helicitySequence; - } public ConstantsManager getConstantsManager() { return conman; @@ -117,14 +134,13 @@ public TreeSet getHelicities() { return helicities; } - HelicitySequence createHelicitySequence() { - HelicitySequence seq = new HelicitySequenceDelayed( - conman.getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,0)); - seq.addStream(helicities); - return seq; + public void updateHelicitySequence() { + helicitySequence = new HelicitySequenceDelayed( + conman.getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,0)); + helicitySequence.addStream(helicities); } - Event getUnixEvent(Bank config) { + Event getUnixEvent(Bank runConfig) { Bank unix = new Bank(schema.getSchema("RUN::unix")); unix.setRows(eventUnix.size()); int row = 0; @@ -134,7 +150,7 @@ Event getUnixEvent(Bank config) { row++; } Event e = new Event(); - e.write(config); + e.write(runConfig); e.write(unix); return e; } @@ -170,13 +186,22 @@ void processHelicity(Event event, Bank runConfig, Bank recEvent) { HelicityBit hbraw = helicitySequence.getHalfWavePlate() ? HelicityBit.getFlipped(hb) : hb; recEvent.putByte("helicity",0,hb.value()); recEvent.putByte("helicityRaw",0,hbraw.value()); - Bank helScaler = new Bank(schema.getSchema("HEL::scaler")); - event.read(helScaler); - if (helScaler.getRows()>0) { - event.remove(schema.getSchema("HEL::scaler")); - SerialUtil.assignScalerHelicity(runConfig.getLong("timestamp",0), helScaler, helicitySequence); - event.write(helScaler); + Bank scaler = new Bank(helScaler); + event.read(scaler); + if (scaler.getRows()>0) { + event.remove(helScaler); + SerialUtil.assignScalerHelicity(runConfig.getLong("timestamp",0), scaler, helicitySequence); + event.write(scaler); } } + static Bank[] createTaggedBanks(Schema[] tag1banks) { + List lbank = Arrays.asList(tag1banks).stream().map(s -> new Bank(s)).collect(Collectors.toList()); + ListIterator ibank = lbank.listIterator(); + Bank[] banks = new Bank[lbank.size()]; + while (ibank.hasNext()) + banks[ibank.nextIndex()] = ibank.next(); + return banks; + } + } 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 d9a58eeeee..6082c0db05 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 @@ -44,6 +44,8 @@ */ final class ReconMutil { + boolean DEBUG = true; + // Performance parameters: final int BENCH_SECONDS = 30; final int EVENTS_PER_CHUNK = 100; @@ -68,7 +70,7 @@ final class ReconMutil { ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>(); // Queues: - ConcurrentLinkedQueue> readQueue = new ConcurrentLinkedQueue<>(); + ConcurrentLinkedQueue> decoQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> procQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); boolean paused = false; @@ -117,10 +119,12 @@ void launch(int[] threads, String output, String... input) { while (!writerThread.isDone()) { sleep(1000); - //System.out.println(String.format("recon-mutil:: read(%b)/[deco(%d)]/proc(%d)/tag/write(%b)", - // readerThread.isDone(), decoThreads.size(), procThreads.size(), writerThread.isDone())); - //System.out.println(String.format("recon-util:: %d-%d/%d/%d/%d", readEvents, - // readQueue.size(), procQueue.size(), taggedEvents.get(), writeQueue.size())); + if (DEBUG){ + System.out.println(String.format("recon-mutil:: read(%b)/[deco(%d)]/proc(%d)/tag/write(%b)", + readerThread.isDone(), decoThreads.size(), procThreads.size(), writerThread.isDone())); + System.out.println(String.format("recon-util:: %d-%d/%d/%d/%d", readEvents, + decoQueue.size(), procQueue.size(), taggedEvents.get(), writeQueue.size())); + } // cleanup completed parallel threads: for (CompletableFuture f : decoThreads) @@ -172,7 +176,7 @@ void read(int threads, String... input) { // write leftover, partial chunk: if (!output.isEmpty()) { readEvents += output.size(); - readQueue.offer(output); + decoQueue.offer(output); } if (reader instanceof EvioSource evio) evio.close(); @@ -184,10 +188,9 @@ void read(int threads, String... input) { */ void decode(int thread) { while (true) { - List input = readQueue.poll(); + List input = decoQueue.poll(); if (input == null) { - if (readerThread.isDone() && readQueue.isEmpty() && - writeEvents+skipEvents+failEvents >= readEvents) break; + if (decoQueue.isEmpty() && readerThread.isDone() && decoQueue.isEmpty()) break; sleep(100); } else { @@ -220,6 +223,10 @@ void process(int thread) { sleep(100); continue; } + if (procQueue.isEmpty() && decoThreads.isEmpty() && procQueue.isEmpty()) { + if (writeEvents+skipEvents+failEvents >= readEvents) break; + sleep(100); + } List input = procQueue.poll(); if (input == null) { if (decoThreads.isEmpty() && procQueue.isEmpty() && @@ -237,6 +244,10 @@ void process(int thread) { Benchmark.getInstance().pause(engine.getValue().getName()); } Event e = input.get(i).getHipoEvent(); + Benchmark.getInstance().resume("post"); + serial.process(e); + serial.process(e); + Benchmark.getInstance().pause("post"); output.add(e); } writeQueue.offer(output); @@ -261,9 +272,6 @@ void write(String output) { } else { for (int i=0; i 0 || schemaBankList.isEmpty()) @@ -401,7 +409,7 @@ List read(List chunk) { if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { chunk.add(o); if (chunk.size() >= EVENTS_PER_CHUNK) { - readQueue.offer(chunk); + decoQueue.offer(chunk); readEvents += chunk.size(); chunk = new ArrayList<>(EVENTS_PER_CHUNK); } @@ -431,7 +439,7 @@ void reset() { writerThread.cancel(true); close(); } - readQueue = new ConcurrentLinkedQueue<>(); + decoQueue = new ConcurrentLinkedQueue<>(); writeQueue = new ConcurrentLinkedQueue<>(); procThreads = new ConcurrentLinkedQueue(); readEvents = 0; From 5fc6715ff96057fdbd8782207a351e0393030b98 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Fri, 11 Sep 2026 20:22:28 -0400 Subject: [PATCH 32/33] fix oops --- .../java/org/jlab/detector/serial/SerialHoncho.java | 10 +++++----- .../src/main/java/org/jlab/clas/reco/ReconMutil.java | 11 +++++++---- 2 files changed, 12 insertions(+), 9 deletions(-) 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 index 59754e8172..579af51788 100644 --- 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 @@ -83,7 +83,7 @@ public void process(Event event) { processEventUnix(event, cfg); if (evt.getRows() > 0) { event.remove(evt.getSchema()); - //processHelicity(event, cfg, evt); + processHelicity(event, cfg, evt); processScalers(cfg, evt); event.write(evt); } @@ -118,6 +118,10 @@ public void clear() { helicitySequence = null; } + public TreeSet getHelicities() { + return helicities; + } + public DaqScalersSequence getScalers() { return scalers; } @@ -130,10 +134,6 @@ public SchemaFactory getSchemaFactory() { return schema; } - public TreeSet getHelicities() { - return helicities; - } - public void updateHelicitySequence() { helicitySequence = new HelicitySequenceDelayed( conman.getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,0)); 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 6082c0db05..ae5be5f440 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 @@ -44,7 +44,7 @@ */ final class ReconMutil { - boolean DEBUG = true; + boolean DEBUG = false; // Performance parameters: final int BENCH_SECONDS = 30; @@ -190,7 +190,11 @@ void decode(int thread) { while (true) { List input = decoQueue.poll(); if (input == null) { - if (decoQueue.isEmpty() && readerThread.isDone() && decoQueue.isEmpty()) break; + if (decoQueue.isEmpty() && readerThread.isDone() && decoQueue.isEmpty()) { + if (thread == 0) serial.updateHelicitySequence(); + System.err.println("recon-mutil:: Helicity sequence updated"); + break; + } sleep(100); } else { @@ -224,7 +228,7 @@ void process(int thread) { continue; } if (procQueue.isEmpty() && decoThreads.isEmpty() && procQueue.isEmpty()) { - if (writeEvents+skipEvents+failEvents >= readEvents) break; + if (writeEvents+skipEvents+failEvents >= readEvents+taggedEvents.get()) break; sleep(100); } List input = procQueue.poll(); @@ -246,7 +250,6 @@ void process(int thread) { Event e = input.get(i).getHipoEvent(); Benchmark.getInstance().resume("post"); serial.process(e); - serial.process(e); Benchmark.getInstance().pause("post"); output.add(e); } From dceb40e20c5a4b56e91d797bf93cfde176613ea3 Mon Sep 17 00:00:00 2001 From: Nathan Baltzell Date: Sun, 13 Sep 2026 19:34:53 -0400 Subject: [PATCH 33/33] cleanup, add volatiles --- .../jlab/detector/serial/SerialHoncho.java | 13 ++- .../java/org/jlab/clas/reco/ReconMutil.java | 90 ++++++++++--------- 2 files changed, 55 insertions(+), 48 deletions(-) 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 index 579af51788..eef2aba68f 100644 --- 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 @@ -35,10 +35,10 @@ public class SerialHoncho { Schema helScaler; Schema helicityAdc; ConstantsManager conman; - TreeMap eventUnix; - HelicitySequence helicitySequence; - TreeSet helicities; - DaqScalersSequence scalers; + volatile TreeMap eventUnix; + volatile HelicitySequence helicitySequence; + volatile TreeSet helicities; + volatile DaqScalersSequence scalers; int run = 0; public SerialHoncho(SchemaFactory schema) { @@ -63,9 +63,8 @@ public synchronized Event read(Event event) { scalers.add(event); event.read(cfg); event.read(hel); - if (cfg.getRows() > 0) { - if (run <= 0 && cfg.getInt("run", 0) > 0) - run = cfg.getInt("run",0); + if (cfg.getRows() > 0 && cfg.getInt("run", 0) > 0) { + run = cfg.getInt("run",0); int unix = cfg.getInt("unixtime",0); int evno = cfg.getInt("event",0); if (unix > 0 && evno > 0) eventUnix.put(evno, unix); 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 ae5be5f440..5ec2ddb503 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 @@ -44,7 +44,7 @@ */ final class ReconMutil { - boolean DEBUG = false; + boolean DEBUG = true; // Performance parameters: final int BENCH_SECONDS = 30; @@ -58,7 +58,7 @@ final class ReconMutil { static { schema.initFromDirectory(ClasUtilsFile.getResourceDir("CLAS12DIR","etc/bankdefs/hipo4")); } // Processors: - SerialHoncho serial; + volatile SerialHoncho serial; CLASDecoderPool decoders = new CLASDecoderPool(64,"default",null); Map engines = new LinkedHashMap<>(); @@ -73,8 +73,7 @@ final class ReconMutil { ConcurrentLinkedQueue> decoQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> procQueue = new ConcurrentLinkedQueue<>(); ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>(); - boolean paused = false; - + // Static parameters: int maxEvents; int skipEvents; @@ -82,14 +81,17 @@ final class ReconMutil { OptionParser parser; // Progress counters: - int readEvents; - int writeEvents; - int failEvents; - int fileEvents; - int maxFileEvents; - AtomicInteger taggedEvents = new AtomicInteger(); - ProgressPrintout progress = new ProgressPrintout(); - + volatile int readEvents; + volatile int writeEvents; + volatile int failEvents; + volatile int fileEvents; + volatile int maxFileEvents; + volatile AtomicInteger taggedEvents = new AtomicInteger(); + volatile ProgressPrintout progress = new ProgressPrintout(); + + // Control flags: + volatile boolean paused = true; + ReconMutil(OptionParser parser) { init(parser); } @@ -104,7 +106,7 @@ void launch(int[] threads, String output, String... input) { reset(); - System.out.println(String.format("recon-mutil:: Spawning %d+++ Threads...",threads[0])); + System.out.println(String.format("recon-mutil:: spawning 2*%d+2 threads",threads[0])); // spawn all the threads: readerThread = CompletableFuture.runAsync(() -> { read(threads[0], input); }); @@ -114,30 +116,20 @@ void launch(int[] threads, String output, String... input) { decoThreads.offer(CompletableFuture.runAsync(() -> { decode(j); })); procThreads.offer(CompletableFuture.runAsync(() -> { process(j); })); } - - // wait for the writer to be done: + + // perform scaling test: + if (threads.length > 1 && rethreadThread == null && writeEvents > 100) { + rethreadThread = CompletableFuture.runAsync(() -> { rethread(BENCH_SECONDS,threads); }); + rethreadThread.join(); + reset(); + } + + // wait for finish: while (!writerThread.isDone()) { sleep(1000); - - if (DEBUG){ - System.out.println(String.format("recon-mutil:: read(%b)/[deco(%d)]/proc(%d)/tag/write(%b)", - readerThread.isDone(), decoThreads.size(), procThreads.size(), writerThread.isDone())); - System.out.println(String.format("recon-util:: %d-%d/%d/%d/%d", readEvents, - decoQueue.size(), procQueue.size(), taggedEvents.get(), writeQueue.size())); - } - - // cleanup completed parallel threads: - for (CompletableFuture f : decoThreads) - if (f.isDone()) decoThreads.remove(f); - for (CompletableFuture f : procThreads) - if (f.isDone()) procThreads.remove(f); - - // perform scaling test: - if (threads.length > 1 && rethreadThread == null && writeEvents > 100) { - rethreadThread = CompletableFuture.runAsync(() -> { rethread(BENCH_SECONDS,threads); }); - rethreadThread.join(); - reset(); - } + if (DEBUG) show(); + for (CompletableFuture f : decoThreads) if (f.isDone()) decoThreads.remove(f); + for (CompletableFuture f : procThreads) if (f.isDone()) procThreads.remove(f); } } @@ -187,14 +179,12 @@ void read(int threads, String... input) { * @param thread thread number */ void decode(int thread) { + int serials = 0; while (true) { List input = decoQueue.poll(); if (input == null) { - if (decoQueue.isEmpty() && readerThread.isDone() && decoQueue.isEmpty()) { - if (thread == 0) serial.updateHelicitySequence(); - System.err.println("recon-mutil:: Helicity sequence updated"); + if (decoQueue.isEmpty() && readerThread.isDone() && decoQueue.isEmpty()) break; - } sleep(100); } else { @@ -210,6 +200,12 @@ void decode(int thread) { output.add(new HipoDataEvent(taggedEvent, schema)); taggedEvents.incrementAndGet(); } + if (thread == 0 && ++serials % 10000 == 0) { + paused = true; + sleep(1000); + serial.updateHelicitySequence(); + paused = false; + } Benchmark.getInstance().pause("serial"); } procQueue.offer(output); @@ -223,7 +219,7 @@ void decode(int thread) { */ void process(int thread) { while (true) { - if (paused || serial.getScalers().size() < 10 || taggedEvents.get() < 100) { + if (paused) { sleep(100); continue; } @@ -407,7 +403,7 @@ List read(List chunk) { } else { Event event = new Event(); - o = ((HipoReader)reader).getEvent(event, fileEvents); + o = ((HipoReader)reader).getEvent(event, ++fileEvents); } if (o != null && (skipEvents < 1 || readEvents > skipEvents)) { chunk.add(o); @@ -436,7 +432,9 @@ void close() { * Forcefully shutdown all threads, close files, and reset queues and counters. */ void reset() { + paused = true; for (CompletableFuture f : procThreads) f.cancel(true); + for (CompletableFuture f : decoThreads) f.cancel(true); if (readerThread != null) readerThread.cancel(true); if (writerThread != null) { writerThread.cancel(true); @@ -527,6 +525,16 @@ else if (!parser.getOption("-c").isDefault()) { } } + void show() { + String s1 = String.format("threads(r/d/p/w)=(%b/%d/%d/%b)", + readerThread.isDone(), decoThreads.size(), procThreads.size(), writerThread.isDone()); + String s2 = String.format(" queues(d/p/w)=(%d/%d/%d)", + decoQueue.size(), procQueue.size(), writeQueue.size()); + String s3 = String.format(" events(r/w/t)=(%d/%d/%d)", + readEvents, writeEvents, taggedEvents.get()); + System.out.println("recon-mutil:: "+s1+" "+s2+" "+s3); + } + /** * The command-line entry-point known as "recon-mutil". * @param args command-line arguments