Skip to content

Commit 15ae3a6

Browse files
committed
add support for eainfo streams
1 parent 2140f53 commit 15ae3a6

3 files changed

Lines changed: 43 additions & 1 deletion

File tree

Collectors/DetailedCollector.py

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,12 @@ def __init__(self, *args, **kw):
4242
self._exchange_tpc = self.config.get('AMQP', 'exchange_tpc')
4343
self._wlcg_exchange_tpc = self.config.get('AMQP', 'wlcg_exchange_tpc')
4444

45+
if 'SCITAGS' in self.config:
46+
self._scitags_mapping_file = self.config.get('SCITAGS', 'mapping_file_path')
47+
with open(self._scitags_mapping_file) as scitag_file:
48+
self._scitags_mapping = json.load(scitag_file)
49+
else:
50+
self._scitags_mapping = None
4551

4652
self.last_flush = time.time()
4753
self.seq_data = {}
@@ -131,6 +137,7 @@ def addRecord(self, sid, userID, fileClose, timestamp, addr, openTime, fileToClo
131137
u = self._users[sid][userInfo].get('userinfo', None)
132138
auth = self._users[sid][userInfo].get('authinfo', None)
133139
appinfo = self._users[sid][userInfo].get('appinfo', None)
140+
eainfo = self._users[sid][userInfo].get('eainfo', None)
134141

135142
if u is not None:
136143
hostname = u.host.decode('idna')
@@ -151,6 +158,9 @@ def addRecord(self, sid, userID, fileClose, timestamp, addr, openTime, fileToClo
151158
rec['vo'] = auth.on.decode('utf-8')
152159
if appinfo is not None:
153160
rec['appinfo'] = appinfo
161+
if eainfo is not None:
162+
rec['activity'] = eainfo.activity
163+
rec['experiment'] = eainfo.experiment
154164

155165
except KeyError:
156166
self.logger.exception("File close record from unknown UserID=%i, SID=%s", userID, sid)
@@ -582,6 +592,16 @@ def process(self, data, addr, port):
582592
elif stream_type == "P":
583593
self.process_gstream_tpc(decoded_gstream, addr,sid)
584594

595+
elif header.code == b'U': #eainfo
596+
if self._scitags_mapping:
597+
eainfo = decoding.eaInfo(rest, self._scitags_mapping)
598+
599+
try:
600+
self.logger.debug("Adding new eainfo: %s.", eainfo)
601+
self._users[sid][userInfo]['eainfo'] = eainfo
602+
except KeyError:
603+
self.logger.warning("Received eainfo for a server or user not seen yet")
604+
585605
else:
586606
infolen = len(data) - 4
587607
mm = decoding.mapheader._make(struct.unpack("!I" + str(infolen) + "s", data))

Collectors/decoding.py

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
prginfo = namedtuple("prginfo", ["xfn", "tod", "sz", "at", "ct", "mt", "fn"])
1616
xfrinfo = namedtuple("xfrinfo", ["lfn", "tod", "sz", "tm", "op", "rc", "pd"])
1717
pathinfo = namedtuple("pathinfo", ["userinfo", "path"])
18+
eainfo = namedtuple("eainfo", ["udid", "experiment", "activity"])
1819

1920
fileOpen = namedtuple("fileOpen", ["rectype", "recFlag", "recSize", "fileID", "fileSize", "userID", "fileName"])
2021
fileXfr = namedtuple("fileXfr", ["rectype", "recFlag", "recSize", "fileID", "read", "readv", "write"])
@@ -143,6 +144,21 @@ def xfrInfo(message):
143144
pd = b''
144145
return xfrinfo([lfn, tod, sz, tm, op, rc, pd])
145146

147+
148+
def eaInfo(message, scitags_mapping):
149+
if isinstance(message, str):
150+
message = message.encode('utf-8')
151+
r = message.split(b'&')
152+
153+
udid = r[1].split(b'=')[1]
154+
expc = int(r[2].split(b'=')[1]) - 1
155+
actc = int(r[3].split(b'=')[1]) - 1
156+
157+
experiment = scitags_maping['experiments'][expc]['expName']
158+
activity = scitags_maping['experiments'][expc]['activities'][actc]['activityName']
159+
return eainfo(udid, experiment, activity)
160+
161+
146162
def gStream(message):
147163
# int, int, int64, null terminated string
148164

Collectors/wlcg_converter.py

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -100,11 +100,17 @@ def Convert(source_record):
100100
to_return["user_protocol"] = source_record["protocol"]
101101

102102
# vo
103-
if "vo" in source_record:
103+
if "experiment" in source_record:
104+
to_return["vo"] = source_record["experiment"]
105+
elif "vo" in source_record:
104106
to_return["vo"] = source_record["vo"]
105107

106108
# write_bytes
107109
to_return["write_bytes"] = source_record['write']
110+
111+
# activity
112+
to_return["activity"] = source_record.get('activity', 'unknown')
113+
108114
# remote_access - boolean
109115
# is_transfer (optional)
110116
# user_fqan (optional)

0 commit comments

Comments
 (0)