diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml
index c72c5b07a1..c38f661b6e 100644
--- a/.github/workflows/ci.yml
+++ b/.github/workflows/ci.yml
@@ -247,12 +247,10 @@ jobs:
tar xzf clara.tar.gz
tar xzf 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
+ - 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 -o rec.hipo ./clas_018779.evio.00001
- uses: actions/upload-artifact@v7
with:
name: test_clara_result
diff --git a/.gitignore b/.gitignore
index 9f0aad0fba..51862428ed 100644
--- a/.gitignore
+++ b/.gitignore
@@ -58,3 +58,6 @@ hs_err_pid*
bin/evio2hipotest
export.sh
+
+*.txt
+
diff --git a/.gitlab-ci.yml b/.gitlab-ci.yml
index d42ca40d5f..8e4a12eb63 100644
--- a/.gitlab-ci.yml
+++ b/.gitlab-ci.yml
@@ -175,13 +175,24 @@ 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
script:
- - ./clara/plugins/clas12/libexec/profile -e cpu run-clara -c ./clara -y ./etc/services/rgd-clarode.yml -o tmp1 -t 16 -l $EVIOFILE
- - ./clara/plugins/clas12/libexec/profile -e alloc run-clara -c ./clara -y ./etc/services/rgd-clarode.yml -o tmp2 -t 16 -l $EVIOFILE
- - ./clara/plugins/clas12/libexec/profile -e lock run-clara -c ./clara -y ./etc/services/rgd-clarode.yml -o tmp3 -t 16 -l $EVIOFILE
+ - ./clara/plugins/clas12/libexec/profile -e cpu -o _clara run-clara -c ./clara -y ./etc/services/rgd-clarode.yml -o tmp1 -t 16 -l $EVIOFILE
+ - ./clara/plugins/clas12/libexec/profile -e alloc -o _clara run-clara -c ./clara -y ./etc/services/rgd-clarode.yml -o tmp2 -t 16 -l $EVIOFILE
+ - ./clara/plugins/clas12/libexec/profile -e lock -o _clara run-clara -c ./clara -y ./etc/services/rgd-clarode.yml -o tmp3 -t 16 -l $EVIOFILE
+ - ./clara/plugins/clas12/libexec/profile -e cpu -o _mutil recon-mutil -y ./etc/services/rgd-clarode.yml -o mutil1.hipo -t 16 -l SEVERE $EVIOFILE
+ - ./clara/plugins/clas12/libexec/profile -e alloc -o _mutil recon-mutil -y ./etc/services/rgd-clarode.yml -o mutil2.hipo -t 16 -l SEVERE $EVIOFILE
+ - ./clara/plugins/clas12/libexec/profile -e lock -o _mutil recon-mutil -y ./etc/services/rgd-clarode.yml -o mutil3.hipo -t 16 -l SEVERE $EVIOFILE
artifacts:
when: always
expire_in: 1 week
@@ -197,7 +208,7 @@ scaling:
when: always
expire_in: 1 week
paths:
- - scaling.svg
+ - scaling-clara.svg
pages:
stage: deploy
@@ -207,7 +218,7 @@ pages:
- cp -r pages public
- mkdir public/perf
- cp asprof*.html public/perf
- - cp scaling.svg public/perf
+ - cp scaling*.svg public/perf
pages: true
artifacts:
expire_in: 1 week
diff --git a/bin/reco-mutil b/bin/reco-mutil
new file mode 100755
index 0000000000..8d061f4c45
--- /dev/null
+++ b/bin/reco-mutil
@@ -0,0 +1,18 @@
+#!/bin/bash
+
+. `dirname $0`/../libexec/env.sh
+
+split_cli $@
+
+export MALLOC_ARENA_MAX=1
+
+#threads=$(get_threads)
+#node=$(((RANDOM % 2)))
+#cpus=$(get_all_numa_cpus $node)
+#cpus=$(get_numa_cpus $node)
+#taskset -c $cpus \
+
+ java ${JAVA_OPTS-} -XX:MaxRAMPercentage=40.0 -Xms10240m ${jvm_options[@]} \
+ -cp ${COATJAVA_CLASSPATH:-''} \
+ org.jlab.clas.reco.RecoMutil \
+ ${class_options[@]}
diff --git a/bin/recon-mutil b/bin/recon-mutil
new file mode 100755
index 0000000000..f604441243
--- /dev/null
+++ b/bin/recon-mutil
@@ -0,0 +1,16 @@
+#!/bin/bash
+
+. `dirname $0`/../libexec/env.sh
+
+split_cli $@
+
+export MALLOC_ARENA_MAX=12
+
+threads=$(get_threads)
+
+ #java ${JAVA_OPTS-} -Xmx256g -Xms10240m ${jvm_options[@]} \
+ #java ${JAVA_OPTS-} -Xms10240m -XX:+UseParallelGC ${jvm_options[@]} \
+ java ${JAVA_OPTS-} -XX:MaxRAMPercentage=40.0 -Xms10240m ${jvm_options[@]} \
+ -cp ${COATJAVA_CLASSPATH:-''} \
+ org.jlab.clas.reco.ReconMutil \
+ ${class_options[@]}
diff --git a/bin/recon-util b/bin/recon-util
index 90f5a6324f..32843f637f 100755
--- a/bin/recon-util
+++ b/bin/recon-util
@@ -4,7 +4,7 @@
split_cli $@
-export MALLOC_ARENA_MAX=1
+export MALLOC_ARENA_MAX=12
java ${JAVA_OPTS-} -Xmx1536m -Xms1024m -XX:+UseSerialGC ${jvm_options[@]} \
-cp ${COATJAVA_CLASSPATH:-''} \
diff --git a/common-tools/clara-io/pom.xml b/common-tools/clara-io/pom.xml
index 1c3a50f9b6..e192efa6b7 100644
--- a/common-tools/clara-io/pom.xml
+++ b/common-tools/clara-io/pom.xml
@@ -48,12 +48,6 @@
14.2.1-SNAPSHOT
-
- org.jlab.clas
- clas-utils
- 14.2.1-SNAPSHOT
-
-
org.apache.commons
commons-text
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 c262aa3d2b..8b2385f7f3 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,17 +1,13 @@
package org.jlab.io.clara;
-import java.io.File;
import java.nio.file.Path;
-import java.util.List;
import org.jlab.clara.std.services.EventWriterException;
import org.jlab.detector.calib.utils.ConstantsManager;
-import org.jlab.detector.helicity.HelicitySequenceDelayed;
import org.jlab.detector.serial.SerialHoncho;
-import org.jlab.detector.serial.PostProcessor;
+import org.jlab.detector.serial.Occupancer;
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;
@@ -21,18 +17,17 @@
* 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
* 6. Outputs a parallel file of select trigger bits with a raw schema
* @author baltzell
*/
public class Clas12Writer extends HipoToHipoWriter {
+ Occupancer occupancer;
SerialHoncho serial;
- Bank runConfig;
ConstantsManager conman;
SchemaFactory fullSchema;
- boolean postprocess;
+ Bank runConfig;
HipoWriterSorted paraWriter;
long paraCount;
@@ -40,13 +35,13 @@ public class Clas12Writer extends HipoToHipoWriter {
int paraTriggerPrescale;
private void init(JSONObject opts) {
+ occupancer = new Occupancer();
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();
+ runConfig = new Bank(fullSchema.getSchema("RUN::config"));
conman.init("/runcontrol/hwp","/runcontrol/helicity");
- postprocess = opts.optBoolean("postprocess", false);
paraTriggerMask = opts.optLong("paraTriggerMask", 0);
paraTriggerPrescale = opts.optInt("paraTriggerPrescale", 0);
if (opts.has("variation")) conman.setVariation(opts.getString("variation"));
@@ -78,6 +73,7 @@ protected HipoWriterSorted createWriter(Path file, JSONObject opts) throws Event
protected void writeEvent(Object event) throws EventWriterException {
Event t = serial.read((Event)event);
if (!t.isEmpty()) writer.addEvent(t, 1);
+ occupancer.process(((Event)event));
super.writeEvent(event);
writeRaw((Event)event, t);
}
@@ -96,48 +92,8 @@ void writeRaw(Event physics, Event tagged) {
protected void closeWriter() {
serial.closure(writer);
super.closeWriter();
- if (postprocess) postprocess();
serial.clear();
if (paraTriggerMask > 0) paraWriter.close();
}
-
- /**
- * 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));
- }
}
diff --git a/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/Occupancer.java b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/Occupancer.java
new file mode 100644
index 0000000000..92e50b85b1
--- /dev/null
+++ b/common-tools/clas-detector/src/main/java/org/jlab/detector/serial/Occupancer.java
@@ -0,0 +1,183 @@
+package org.jlab.detector.serial;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.nio.file.Paths;
+import java.util.ArrayList;
+import java.util.Map;
+import org.jlab.detector.banks.RawBank;
+import org.jlab.utils.groups.IndexedTable;
+import org.jlab.utils.groups.IndexedTable.IndexedEntry;
+import org.jlab.detector.banks.RawBank.OrderGroups;
+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.utils.system.ClasUtilsFile;
+
+/**
+ * Occupancy bookkeeper based on IndexedTable, with I/O helpers for indexed banks.
+ *
+ * @author baltzell
+ */
+public class Occupancer extends ArrayList {
+
+ static final String BANKDIR = ClasUtilsFile.getResourceDir("CLAS12DIR","etc/bankdefs/hipo4/singles/occupancy");
+ static final SchemaFactory FULL_SCHEMA = new SchemaFactory();
+ static { FULL_SCHEMA.initFromDirectory(ClasUtilsFile.getResourceDir("CLAS12DIR","etc/bankdefs/hipo4")); }
+
+ int prescale = 1000;
+ int nevents = 0;
+
+ public Occupancer() {
+ super();
+ init();
+ }
+
+ public Occupancer(int prescale) {
+ super();
+ this.prescale = prescale;
+ init();
+ }
+
+ private void init() {
+ try {
+ // Initialize OccupanceTable list, by reading occupancy bank filenames
+ // and stripping ".json" suffix and "OCC::" prefix to get hit bank name:
+ addAll(Files.list(Paths.get(BANKDIR))
+ .filter(Files::isRegularFile)
+ .map(p -> p.getFileName().toString())
+ .map(s -> s.substring(0, s.length()-5))
+ .map(s -> s.substring(5, s.length()))
+ .map(OccupanceTable::new).toList());
+ } catch (IOException ex) {
+ System.getLogger(Occupancer.class.getName()).log(System.Logger.Level.ERROR, (String) null, ex);
+ }
+ }
+
+ public boolean process(Event event) {
+ forEach(t -> { t.fill(event, false); });
+ if (++nevents % prescale == 0) {
+ forEach(t -> {
+ if (t.getTable().getRowCount() > 0) {
+ event.write(t.create(nevents, event));
+ }
+ t.reset();
+ });
+ nevents = 0;
+ }
+ return true;
+ }
+
+ public void reset() {
+ forEach(t -> { t.reset(); nevents = 0; });
+ }
+
+ public static class OccupanceTable {
+
+ Schema occSchema;
+ Schema hitSchema;
+ IndexedTable table;
+
+ /**
+ * A 3-index table, e.g., sector/layer/component.
+ * @param hitBank name of the hit bank
+ */
+ public OccupanceTable(String hitBank) {
+ hitSchema = FULL_SCHEMA.getSchema(hitBank);
+ occSchema = FULL_SCHEMA.getSchema("OCC::" + hitBank);
+ table = new IndexedTable(occSchema.hasEntry("order")?4:3, new String[]{"occ/F"});
+ }
+
+ /**
+ * An N-index table.
+ * @param hitBank name of the hit bank
+ * @param indexCount number of inidices in the hit bank
+ */
+ public OccupanceTable(String hitBank, int indexCount) {
+ hitSchema = FULL_SCHEMA.getSchema(hitBank);
+ occSchema = FULL_SCHEMA.getSchema("OCC::" + hitBank);
+ table = new IndexedTable(indexCount, new String[]{"occ/F"});
+ }
+
+ public final IndexedTable getTable() { return table; }
+
+ /**
+ * Zero the occupancy table.
+ */
+ public final void reset() {
+ table = new IndexedTable(table.getList().getIndexSize(), new String[]{"occ/F"});
+ }
+
+ /**
+ * Get the occupancy table, normalized by number of events.
+ * @param events
+ * @return
+ */
+ public final IndexedTable getOccupancy(long events) {
+ IndexedTable t = new IndexedTable(table.getList().getIndexSize(), new String[]{"occ/F"});
+ for (long hash : ((Map)table.getList().getMap()).keySet()) {
+ t.addEntry(IndexedTable.DEFAULT_GENERATOR.getIndices(hash, table.getList().getIndexSize()));
+ t.setDoubleValueByHash((table.getDoubleValueByHash(0, hash))/events, 0, hash);
+ }
+ return t;
+ }
+
+ /**
+ * Fill the occupancy table.
+ * @param weight
+ * @param index
+ */
+ public final void fill(float weight, int... index) {
+ for (int i=0; i m = table.getList().getMap();
+ int i = 0;
+ for (long hash : m.keySet()) {
+ int[] idx = IndexedTable.DEFAULT_GENERATOR.getIndices(hash, table.getList().getIndexSize());
+ for (int j=0; j eventUnix;
volatile HelicitySequence helicitySequence;
@@ -48,6 +51,9 @@ public SerialHoncho(SchemaFactory schema) {
recEvent = schema.getSchema("REC::Event");
helicityAdc = schema.getSchema("HEL::adc");
helScaler = schema.getSchema("HEL::scaler");
+ runScaler = schema.getSchema("RUN::scaler");
+ rawScaler = schema.getSchema("RAW::scaler");
+ epics = schema.getSchema("RAW::epics");
scalers = new DaqScalersSequence(schema);
helicities = new TreeSet<>();
eventUnix = new TreeMap<>();
@@ -56,6 +62,10 @@ public SerialHoncho(SchemaFactory schema) {
tag1banks[i] = schema.getSchema(TAG1BANKS[i]);
}
+ public boolean containsSerial(Event e) {
+ return e.isEmpty() ? false : e.hasBank(rawScaler) || e.hasBank(epics);
+ }
+
/**
* Register an event's serial data and return a (possibly empty) tag-1 event.
* @param event
@@ -182,16 +192,16 @@ void prune() {
* @param runConfig
* @return
*/
- Event getUnixEvent(Bank config) {
+ Event getUnixEvent(Bank runConfig) {
Bank unix = new Bank(schema.getSchema("RUN::unix"), eventUnix.size());
int row = 0;
for (int evno : eventUnix.keySet()) {
unix.putInt("event", row, evno);
- unix.putInt("unixtime",row, eventUnix.get(evno));
+ unix.putInt("unixtime", row, eventUnix.get(evno));
row++;
}
Event e = new Event();
- e.write(config);
+ e.write(runConfig);
e.write(unix);
return e;
}
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 dfa428f7ec..5c920e8d9b 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
@@ -19,6 +19,7 @@
import java.util.Arrays;
import org.jlab.coda.jevio.EvioException;
import org.jlab.detector.decode.CLASDecoder4;
+import org.jlab.detector.serial.Occupancer;
import org.jlab.io.evio.EvioDataEvent;
import org.jlab.io.evio.EvioSource;
import org.jlab.io.hipo.HipoDataEvent;
@@ -42,7 +43,8 @@ public class EngineProcessor {
private SchemaFactory banksToKeep = null;
private final List schemaExempt = Arrays.asList("RUN::config","DC::tdc");
- private CLASDecoder4 decoder = new CLASDecoder4();
+ private final CLASDecoder4 decoder = new CLASDecoder4();
+ private final Occupancer occupancer = new Occupancer();
public EngineProcessor(){}
@@ -134,7 +136,7 @@ public void initDefault(){
String[] names = new String[]{
"MAGFIELDS",
"DCCR","DCHB","FTOFHB","EC","HTCC","EBHB",
- "DCTB","FTOFTB","EBTB","VTX"
+ "DCTB","FTOFTB","EBTB","OCC"
};
String[] services = new String[]{
@@ -148,7 +150,7 @@ public void initDefault(){
"org.jlab.service.dc.DCTBEngine",
"org.jlab.service.ftof.FTOFTBEngine",
"org.jlab.service.eb.EBTBEngine",
- "org.jlab.rec.service.vtx.VTXEngine"
+ "org.jlab.calibration.service.OccupanceEngine",
};
for(int i = 0; i < names.length; i++){
@@ -164,7 +166,7 @@ public void initAll(){
"CVTFP","CTOF","CND","BAND",
"HTCC","LTCC","EBHB",
"DCTB","FMT","FTOFTB","CVT","EBTB",
- "RICHEB","RTPC","AHDC","ATOF","ALERT", "MC","VTX"
+ "RICHEB","RTPC","AHDC","ATOF","ALERT", "MC","VTX","OCC"
};
String[] services = new String[]{
@@ -197,7 +199,8 @@ public void initAll(){
"org.jlab.service.atof.ATOFEngine",
"org.jlab.service.alert.ALERTEngine",
"org.jlab.service.mc.TruthMatch",
- "org.jlab.rec.service.vtx.VTXEngine"
+ "org.jlab.rec.service.vtx.VTXEngine",
+ "org.jlab.calibration.service.OccupanceEngine",
};
if(names.length!=services.length)
LOGGER.log(Level.SEVERE, "initAll : the list of services does not match the list of service names... ");
@@ -316,6 +319,7 @@ public void processFile(String file, String output){
public void processEvent(DataEvent event, HipoDataSync writer) {
processEvent(event);
removeBanks(event);
+ occupancer.process(((HipoDataEvent)event).getHipoEvent());
writer.writeEvent(event);
}
@@ -390,24 +394,6 @@ public void show(){
}
}
- 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","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");
- parser.addOption("-B",null,"background file");
- parser.addOption("-P",null,"preload file for post-processing");
- parser.addOption("-R","0","rebuild scalers");
- parser.addOption("-H","0","restream helicity");
- return parser;
- }
-
protected final void init(OptionParser p) {
p.syncLogLevel(LOGGER);
@@ -449,7 +435,7 @@ else if (p.getOption("-c").intValue() > 0) {
}
public static void main(String[] args) {
- OptionParser parser = EngineProcessor.getParser();
+ OptionParser parser = ReconUtil.getParser();
parser.parse(args);
EngineProcessor proc = new EngineProcessor(parser);
proc.processFile(parser.getOption("-i").stringValue(),
diff --git a/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/Porch.java b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/Porch.java
new file mode 100644
index 0000000000..4b2e31fae6
--- /dev/null
+++ b/common-tools/clas-reco/src/main/java/org/jlab/clas/reco/Porch.java
@@ -0,0 +1,303 @@
+package org.jlab.clas.reco;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.TimerTask;
+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.io.hipo.HipoDataEvent;
+import org.jlab.utils.benchmark.Benchmark;
+import org.jlab.utils.benchmark.BenchmarkTimer;
+import org.jlab.utils.benchmark.ProgressPrintout;
+
+/**
+ * Parallel orchestrator.
+ *
+ */
+public abstract class Porch {
+
+ abstract Object open(String filename);
+ abstract Object read();
+ abstract HipoDataEvent[] decode(int thread, Object event);
+ abstract void process(int thread, HipoDataEvent event);
+ abstract void write(HipoDataEvent event);
+ abstract void readerExit();
+ abstract void decoderExit(int thread);
+ abstract void writerExit();
+ abstract void serial();
+
+ static final String[] BENCHMARK_NAMES = new String[]{"evio","deco","serial","post","write"};
+
+ // Performance parameters:
+ final int BENCH_SECONDS = 30;
+ final int EVENTS_PER_CHUNK = 50;
+
+ // Static parameters:
+ int maxEvents;
+ int skipEvents;
+ String taskset;
+
+ // Threads:
+ CompletableFuture readerThread;
+ CompletableFuture writerThread;
+ ConcurrentLinkedQueue decoThreads = new ConcurrentLinkedQueue<>();
+ ConcurrentLinkedQueue procThreads = new ConcurrentLinkedQueue<>();
+
+ // Queues:
+ ConcurrentLinkedQueue> decoQueue = new ConcurrentLinkedQueue<>();
+ ConcurrentLinkedQueue> procQueue = new ConcurrentLinkedQueue<>();
+ ConcurrentLinkedQueue> writeQueue = new ConcurrentLinkedQueue<>();
+
+ // Progress counters:
+ volatile int readEvents;
+ volatile int failEvents;
+ volatile int fileEvents;
+ volatile int maxFileEvents;
+ volatile int writeEvents;
+ volatile AtomicInteger taggedEvents = new AtomicInteger();
+ volatile ProgressPrintout progress = new ProgressPrintout();
+ volatile Benchmark benchmark = null;
+ TimerTask statsShow = new TimerTask() { @Override public void run() { show(); } };
+
+ // Event reader/source:
+ Object reader;
+
+ /**
+ * The thread launcher and collector.
+ * @param threads number of threads
+ * @param input names of input files to read
+ */
+ public void launch(int[] threads, String... input) {
+
+ reset();
+
+ System.out.println(String.format("recon-mutil:: spawning 2*%d+2 threads",threads[0]));
+
+ // start a period status printout:
+ ReconUtil.runPeriodic(10, statsShow);
+
+ // spawn all the threads:
+ readerThread = CompletableFuture.runAsync(() -> { reader(threads[0], input); });
+ for (int i=0; i { decoder(j); }));
+ }
+ ReconUtil.launch(CompletableFuture.runAsync(() -> { serial(); }));
+ ReconUtil.sleep(1000);
+ for (int i=0; i { processor(j); }));
+ }
+ writerThread = ReconUtil.launch(CompletableFuture.runAsync(() -> { writer(); }));
+
+ // perform scaling test:
+ if (threads.length > 1)
+ ReconUtil.launch(CompletableFuture.runAsync(() -> { rethreader(threads); })).join();
+
+ // wait for finish:
+ writerThread.join();
+ }
+
+ /**
+ * The reader thread.
+ * @param threads
+ * @param input input filenames
+ */
+ final void reader(int threads, String... input) {
+ List inputs = new ArrayList<>(Arrays.asList(input));
+ List