diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml
index 4967588f6e..3cef54e4a3 100644
--- a/.github/workflows/ci.yml
+++ b/.github/workflows/ci.yml
@@ -224,20 +224,15 @@ jobs:
clas_018779.evio.00001
- name: untar build
run: |
- tar xzvf clara.tar.gz
tar xzvf coatjava.tar.gz
- run: ls
- name: run test
- run: ./bin/run-clara -y ./etc/services/rgd-clarode.yml -t 4 -n 500 -c ./clara -o ./tmp ./clas_018779.evio.00001
- - name: ls tmp
- run: ls -lhtr tmp
- - name: rename
- run: mv -v tmp/rec_clas_018779.evio.00001.hipo rec.hipo
+ run: ./coatjava/bin/recon-mutil -t 4 -n 500 -y etc/services/rgd-clarode.yml -o rec_clas_018779.evio.00001.hipo -i clas_018779.evio.00001
- uses: actions/upload-artifact@v7
with:
name: test_clara_result
retention-days: 1
- path: rec.hipo
+ path: rec_clas_018779.evio.00001.hipo
test_coatjava:
needs: [ build ]
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/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[@]}
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 a094b939cd..e1ebfe310e 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.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;
@@ -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);
- if (!t.isEmpty()) writer.addEvent(t, 1);
+ Event tagged = serial.read(((Event)event));
+ if (!tagged.isEmpty()) writer.addEvent(tagged, 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());
+ 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 86%
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..e95ae19575 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,22 +1,22 @@
-package org.jlab.analysis.postprocess;
+package org.jlab.detector.serial;
import java.util.List;
import java.util.TreeMap;
+import org.jlab.io.base.DataBank;
+import org.jlab.io.base.DataEvent;
+import org.jlab.io.hipo.HipoDataBank;
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.io.base.DataBank;
-import org.jlab.io.base.DataEvent;
-import org.jlab.io.hipo.HipoDataBank;
+import org.jlab.jnp.hipo4.io.HipoWriterSorted;
-import org.jlab.detector.calib.utils.ConstantsManager;
import org.jlab.detector.scalers.DaqScalers;
import org.jlab.detector.scalers.DaqScalersSequence;
import org.jlab.detector.helicity.HelicityBit;
import org.jlab.detector.helicity.HelicitySequenceDelayed;
-import org.jlab.jnp.hipo4.io.HipoWriterSorted;
+import org.jlab.detector.calib.utils.ConstantsManager;
import org.jlab.utils.options.OptionParser;
import org.jlab.utils.system.ClasUtilsFile;
@@ -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;
@@ -274,4 +274,20 @@ public static void main(String args[]) {
if (writer != null) writer.close();
}
+ public void postprocess(int run, String input, SerialHoncho serial, HipoWriterSorted writer, List banks) {
+ int d = serial.getConstantsManager().getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,0);
+ HelicitySequenceDelayed helicity = new HelicitySequenceDelayed(d);
+ helicity.addStream(serial.getHelicities());
+ PostProcessor p = new PostProcessor(List.of(input), serial.getSchemaFactory(), helicity, serial.getScalers());
+ HipoReader r = new HipoReader();
+ r.open(input);
+ Event e = new Event();
+ while (r.hasNext()) {
+ r.nextEvent(e);
+ p.processEvent(e);
+ int tag = e.getEventTag();
+ if (tag==1 || banks.isEmpty()) writer.addEvent(e,tag);
+ else writer.addEvent(e.reduceEvent(banks),tag);
+ }
+ }
}
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;
+ }
+}
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
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..caa868e6a8 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
@@ -161,9 +161,13 @@ 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 = "services";
+ if (serviceName.equals("reader") || serviceName.equals("writer")) {
+ section = "io-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));
}
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..58ee6b17e8 100644
--- a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java
+++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/EngineProcessor.java
@@ -18,11 +18,10 @@
import org.jlab.clara.engine.EngineDataType;
import java.util.Arrays;
import org.jlab.coda.jevio.EvioException;
-import org.jlab.detector.decode.CLASDecoder4;
+import org.jlab.detector.decode.CLASDecoder;
+import org.jlab.detector.decode.CLASDecoderPool;
import org.jlab.io.evio.EvioDataEvent;
import org.jlab.io.evio.EvioSource;
-import org.jlab.io.hipo.HipoDataEvent;
-import org.jlab.jnp.hipo4.data.Event;
import org.jlab.jnp.hipo4.data.SchemaFactory;
import org.json.JSONObject;
import org.jlab.utils.ClaraYaml;
@@ -36,13 +35,13 @@ public class EngineProcessor {
public static final String ENGINE_CLASS_BG = "org.jlab.service.bg.BackgroundEngine";
public static final String ENGINE_CLASS_PP = "org.jlab.service.postproc.PostprocEngine";
- private final Map processorEngines = new LinkedHashMap<>();
+ protected final Map processorEngines = new LinkedHashMap<>();
private static final Logger LOGGER = Logger.getLogger(EngineProcessor.class.getPackage().getName());
private boolean updateDictionary = true;
private SchemaFactory banksToKeep = null;
private final List schemaExempt = Arrays.asList("RUN::config","DC::tdc");
- private CLASDecoder4 decoder = new CLASDecoder4();
+ protected final CLASDecoderPool decoders = new CLASDecoderPool(64,"default",null);
public EngineProcessor(){}
@@ -94,7 +93,7 @@ private void setPreloadFiles(String filenames, boolean restream, boolean rebuild
findEngine(ENGINE_CLASS_PP).init();
}
- private void updateDictionary(HipoDataSource source, HipoDataSync sync){
+ protected void updateDictionary(HipoDataSource source, HipoDataSync sync){
SchemaFactory fsync = sync.getWriter().getSchemaFactory();
SchemaFactory fsrc = source.getReader().getSchemaFactory();
List schemaList = fsync.getSchemaKeys();
@@ -255,7 +254,8 @@ public void addEngine(String name, String clazz, String jsonConf) {
}
this.processorEngines.put(name == null ? engine.getName() : name, engine);
} else {
- LOGGER.log(Level.SEVERE, ">>>> ERROR: class is not a reconstruction engine : {0}", clazz);
+ LOGGER.log( clazz.contains("DecoderEngine") ? Level.INFO : Level.SEVERE,
+ "Class is not a reconstruction engine : {0}", clazz);
}
} catch (ClassNotFoundException | InstantiationException | IllegalAccessException ex) {
@@ -342,9 +342,13 @@ public void processFile(EvioSource reader, HipoDataSync writer, int skipEvents,
ByteBuffer bb = reader.getEventBuffer(eventsRead, true);
if (skipEvents <= 0 || eventsRead > skipEvents) {
EvioDataEvent evio = new EvioDataEvent(bb.array(), ByteOrder.LITTLE_ENDIAN);
- Event hipo = decoder.getDecodedEvent(evio, -1, eventsRead, null, null);
- HipoDataEvent hipo2 = new HipoDataEvent(hipo, decoder.getSchemaFactory());
- processEvent(hipo2, writer);
+ try {
+ CLASDecoder d = decoders.take();
+ processEvent(d.getDecodedDataEvenet(evio), writer);
+ decoders.put(d);
+ } catch (InterruptedException ex) {
+ System.getLogger(EngineProcessor.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex);
+ }
}
if (maxEvents > 0 && eventsRead > maxEvents+skipEvents) break;
} catch (EvioException ex) {
@@ -394,10 +398,9 @@ protected static OptionParser getParser() {
OptionParser parser = new OptionParser("recon-util");
parser.addRequired("-o","output.hipo");
parser.addRequired("-i","input.evio/hipo");
- parser.setRequiresInputList(false);
parser.addOption("-c","0","use default configuration [0 - no, 1 - yes/default, 2 - all services] ");
- parser.addOption("-s","-1","number of events to skip");
- parser.addOption("-n","-1","number of events to process");
+ 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");
@@ -405,6 +408,7 @@ protected static OptionParser getParser() {
parser.addOption("-P",null,"preload file for post-processing");
parser.addOption("-R","0","rebuild scalers");
parser.addOption("-H","0","restream helicity");
+ parser.setRequiresInputList(false);
return parser;
}
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