Skip to content

Commit 8d02fd4

Browse files
authored
prepare for on-the-fly post-processing (#1436)
* prep for on-the-fly post-processing * add comments * make it if the user did not
1 parent f53aab1 commit 8d02fd4

2 files changed

Lines changed: 228 additions & 41 deletions

File tree

‎common-tools/clara-io/src/main/java/org/jlab/io/clara/Clas12Writer.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,7 @@ void writeRaw(Event physics, Event tagged) {
9494

9595
@Override
9696
protected void closeWriter() {
97-
serial.finish(writer);
97+
serial.closure(writer);
9898
super.closeWriter();
9999
if (postprocess) postprocess();
100100
serial.clear();
Lines changed: 227 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,22 @@
11
package org.jlab.detector.serial;
22

3+
import java.util.Arrays;
4+
import java.util.List;
5+
import java.util.ListIterator;
36
import java.util.TreeMap;
47
import java.util.TreeSet;
8+
import java.util.stream.Collectors;
59
import org.jlab.detector.calib.utils.ConstantsManager;
610
import org.jlab.detector.decode.CLASDecoder;
11+
import org.jlab.detector.helicity.HelicityBit;
712
import org.jlab.detector.helicity.HelicitySequence;
13+
import org.jlab.detector.helicity.HelicitySequenceDelayed;
814
import org.jlab.detector.helicity.HelicityState;
15+
import org.jlab.detector.scalers.DaqScalers;
916
import org.jlab.detector.scalers.DaqScalersSequence;
1017
import org.jlab.jnp.hipo4.data.Bank;
1118
import org.jlab.jnp.hipo4.data.Event;
19+
import org.jlab.jnp.hipo4.data.Schema;
1220
import org.jlab.jnp.hipo4.data.SchemaFactory;
1321
import org.jlab.jnp.hipo4.io.HipoWriterSorted;
1422

@@ -19,59 +27,167 @@
1927
public class SerialHoncho {
2028

2129
static final String[] TAG1BANKS = {"RUN::scaler","HEL::scaler","RAW::scaler","RAW::epics","HEL::flip","COAT::config"};
22-
2330
SchemaFactory schema;
24-
Bank[] tag1banks;
25-
Bank runConfig;
26-
Bank helicityAdc;
27-
ConstantsManager conman;
28-
TreeMap<Integer,Integer> eventUnix;
29-
TreeSet<HelicityState> helicities;
30-
DaqScalersSequence scalers;
31-
31+
Schema[] tag1banks;
32+
Schema runConfig;
33+
Schema recEvent;
34+
Schema helScaler;
35+
Schema helicityAdc;
36+
volatile ConstantsManager conman;
37+
volatile TreeMap<Integer,Integer> eventUnix;
38+
volatile HelicitySequence helicitySequence;
39+
volatile TreeSet<HelicityState> helicities;
40+
volatile DaqScalersSequence scalers;
41+
int run = 0;
42+
3243
public SerialHoncho(SchemaFactory schema) {
3344
this.schema = schema;
3445
conman = new ConstantsManager();
3546
conman.init("/runcontrol/hwp","/runcontrol/helicity");
36-
runConfig = new Bank(schema.getSchema("RUN::config"));
37-
helicityAdc = new Bank(schema.getSchema("HEL::adc"));
38-
helicities = new TreeSet<>();
47+
runConfig = schema.getSchema("RUN::config");
48+
recEvent = schema.getSchema("REC::Event");
49+
helicityAdc = schema.getSchema("HEL::adc");
50+
helScaler = schema.getSchema("HEL::scaler");
3951
scalers = new DaqScalersSequence(schema);
52+
helicities = new TreeSet<>();
4053
eventUnix = new TreeMap<>();
41-
tag1banks = new Bank[TAG1BANKS.length];
54+
tag1banks = new Schema[TAG1BANKS.length];
4255
for (int i=0; i<tag1banks.length; ++i)
43-
tag1banks[i] = new Bank(schema.getSchema(TAG1BANKS[i]));
56+
tag1banks[i] = schema.getSchema(TAG1BANKS[i]);
4457
}
4558

46-
public synchronized Event read(Event event) {
47-
scalers.add(event);
48-
event.read(runConfig);
49-
event.read(helicityAdc);
50-
if (runConfig.getRows() > 0) {
51-
int unix = runConfig.getInt("unixtime",0);
52-
int evno = runConfig.getInt("event",0);
53-
if (unix > 0 && evno > 0) eventUnix.put(evno, unix);
59+
/**
60+
* Register an event's serial data and return a (possibly empty) tag-1 event.
61+
* @param event
62+
* @return new tag-1 event
63+
*/
64+
public Event read(Event event) {
65+
Bank cfg = new Bank(runConfig);
66+
Bank hel = new Bank(helicityAdc);
67+
event.read(cfg);
68+
event.read(hel);
69+
read(event, cfg, hel);
70+
return CLASDecoder.createTaggedEvent(event, cfg, createTaggedBanks(tag1banks));
71+
}
72+
73+
/**
74+
* Modify a physics event's helicity and charge information.
75+
* @param event
76+
*/
77+
public void process(Event event) {
78+
Bank cfg = new Bank(runConfig);
79+
Bank evt = new Bank(recEvent);
80+
event.read(cfg);
81+
event.read(evt);
82+
if (cfg.getRows() > 0) {
83+
processEventUnix(event, cfg);
84+
if (evt.getRows() > 0) {
85+
event.remove(evt.getSchema());
86+
processHelicity(event, cfg, evt);
87+
processScalers(scalers.get(cfg.getLong("timestamp",0)), evt);
88+
event.write(evt);
89+
}
5490
}
55-
helicities.add(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman));
56-
return CLASDecoder.createTaggedEvent(event, runConfig, tag1banks);
5791
}
5892

59-
public void finish(HipoWriterSorted writer) {
60-
writer.addEvent(getUnixEvent(runConfig),1);
61-
HelicitySequence.writeFlips(schema, writer, helicities);
93+
/**
94+
* Add helicity and unixtime sequence banks.
95+
* @param writer
96+
*/
97+
public void closure(HipoWriterSorted writer) {
98+
Bank cfg = new Bank(runConfig, 1);
99+
cfg.putInt("run",0,run);
100+
writer.addEvent(getUnixEvent(cfg),1);
101+
if (helicitySequence == null) updateHelicitySequence();
102+
helicitySequence.writeFlips(writer, 1);
62103
}
63104

105+
/**
106+
* Zero all the sequences.
107+
*/
64108
public void clear() {
65-
while (helicities.size() > 100) helicities.pollFirst();
66-
scalers.clear(100);
109+
eventUnix.clear();
110+
helicities.clear();
111+
scalers.clear();
112+
helicitySequence = null;
67113
}
68-
114+
115+
public TreeSet<HelicityState> getHelicities() {
116+
return helicities;
117+
}
118+
119+
public DaqScalersSequence getScalers() {
120+
return scalers;
121+
}
122+
123+
public ConstantsManager getConstantsManager() {
124+
return conman;
125+
}
126+
127+
public SchemaFactory getSchemaFactory() {
128+
return schema;
129+
}
130+
131+
/**
132+
* Recreate the HelicitySequence from the TreeSet of helicity states.
133+
*/
134+
public void updateHelicitySequence() {
135+
helicitySequence = new HelicitySequenceDelayed(
136+
conman.getConstants(run, "/runcontrol/helicity").getIntValue("delay",0,0,0));
137+
helicitySequence.addStream(helicities);
138+
}
139+
140+
/**
141+
*
142+
* @param event
143+
* @param runConfig
144+
* @param helicityAdc
145+
*/
146+
void read(Event event, Bank runConfig, Bank helicityAdc) {
147+
scalers.add(event);
148+
helicities.add(HelicityState.createFromFadcBank(helicityAdc, runConfig, conman));
149+
prune();
150+
if (runConfig.getRows() > 0) {
151+
int r = runConfig.getInt("run", 0);
152+
if (r > 0) {
153+
if (r != run) {
154+
clear();
155+
run = r;
156+
}
157+
}
158+
if (run > 0) {
159+
int unix = runConfig.getInt("unixtime",0);
160+
int evno = runConfig.getInt("event",0);
161+
if (unix > 0 && evno > 0) eventUnix.put(evno, unix);
162+
}
163+
}
164+
}
165+
166+
/**
167+
* Trim down the helicity and scaler FIFOs to reasonable values.
168+
*/
169+
void prune() {
170+
// Estimated size of HelicityState is ~22 bytes.
171+
// 1 million states, ~22 MB, 1 minute at 10 kHz trigger.
172+
if (helicities.size() > 2e6)
173+
pruneHelicities(helicities, (int)1e6);
174+
// Assuming scalers are 50x larger.
175+
// 10,000 events is 2.7 hours at 1 Hz.
176+
if (scalers.size() > 2e4)
177+
scalers.clear((int)1e4);
178+
}
179+
180+
/**
181+
* Get a new event with RUN::unix and RUN::config banks.
182+
* @param runConfig
183+
* @return
184+
*/
69185
Event getUnixEvent(Bank config) {
70186
Bank unix = new Bank(schema.getSchema("RUN::unix"), eventUnix.size());
71187
int row = 0;
72188
for (int evno : eventUnix.keySet()) {
73189
unix.putInt("event", row, evno);
74-
unix.putInt("unixtime", row, eventUnix.get(evno));
190+
unix.putInt("unixtime",row, eventUnix.get(evno));
75191
row++;
76192
}
77193
Event e = new Event();
@@ -80,19 +196,90 @@ Event getUnixEvent(Bank config) {
80196
return e;
81197
}
82198

83-
public DaqScalersSequence getScalers() {
84-
return scalers;
199+
/**
200+
* Get the unix time for a given RUN::config bank.
201+
* @return unix time
202+
*/
203+
int getUnixTime(Bank runConfig) {
204+
if (runConfig.getRows() < 1) {
205+
Integer key = eventUnix.floorKey(runConfig.getInt("event",0));
206+
if (key != null) {
207+
Integer unix = eventUnix.get(key);
208+
if (unix != null) return unix;
209+
}
210+
}
211+
return 0;
85212
}
86-
87-
public TreeSet<HelicityState> getHelicities() {
88-
return helicities;
213+
214+
/**
215+
* Update RUN::config's unixtime, and update the event.
216+
*/
217+
void processEventUnix(Event event, Bank runConfig) {
218+
int ut = getUnixTime(runConfig);
219+
event.remove(runConfig.getSchema());
220+
runConfig.putInt("unixtime", 0, ut);
221+
event.write(runConfig);
89222
}
90223

91-
public ConstantsManager getConstantsManager() {
92-
return conman;
224+
/**
225+
* Update an event's helicity information, in REC::Event and HEL::scaler.
226+
* @param event
227+
* @param runConfig
228+
* @param recEvent
229+
*/
230+
void processHelicity(Event event, Bank runConfig, Bank recEvent) {
231+
HelicityBit hb = helicitySequence.search(runConfig.getLong("timestamp", 0));
232+
HelicityBit hbraw = helicitySequence.getHalfWavePlate() ? HelicityBit.getFlipped(hb) : hb;
233+
recEvent.putByte("helicity",0,hb.value());
234+
recEvent.putByte("helicityRaw",0,hbraw.value());
235+
Bank scaler = new Bank(helScaler);
236+
event.read(scaler);
237+
if (scaler.getRows()>0) {
238+
event.remove(helScaler);
239+
SerialUtil.assignScalerHelicity(runConfig.getLong("timestamp",0), scaler, helicitySequence);
240+
event.write(scaler);
241+
}
93242
}
94243

95-
public SchemaFactory getSchemaFactory() {
96-
return schema;
244+
/**
245+
* Update a REC::EVent's bank beam-charge information.
246+
*/
247+
static void processScalers(DaqScalers ds, Bank recEvent) {
248+
if (ds != null) {
249+
recEvent.putFloat("beamCharge",0, (float) ds.dsc2.getBeamChargeGated());
250+
recEvent.putDouble("liveTime",0,ds.dsc2.getLivetime());
251+
}
252+
}
253+
254+
/**
255+
* For all states deeper than depth, remove consecutive helicity states that
256+
* differ only by timestamp.
257+
* @param helicities
258+
* @param depth
259+
*/
260+
static void pruneHelicities(TreeSet<HelicityState> helicities, int depth) {
261+
HelicityState prev = null;
262+
ListIterator<HelicityState> iter = (ListIterator)helicities.iterator();
263+
final int size = helicities.size();
264+
while (iter.hasNext() && iter.nextIndex() < size-depth) {
265+
HelicityState next = iter.next();
266+
if (prev != null && prev == next)
267+
helicities.remove(next);
268+
}
97269
}
270+
271+
/**
272+
* FIXME: Switch to lists and delete this.
273+
* @param tag1banks
274+
* @return
275+
*/
276+
static Bank[] createTaggedBanks(Schema[] tag1banks) {
277+
List<Bank> lbank = Arrays.asList(tag1banks).stream().map(s -> new Bank(s)).collect(Collectors.toList());
278+
ListIterator<Bank> ibank = lbank.listIterator();
279+
Bank[] banks = new Bank[lbank.size()];
280+
while (ibank.hasNext())
281+
banks[ibank.nextIndex()] = ibank.next();
282+
return banks;
283+
}
284+
98285
}

0 commit comments

Comments
 (0)