-
Notifications
You must be signed in to change notification settings - Fork 2
Update all listeners to GCN-Kafka #74
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
469455c
6e1a6eb
0001f15
7955299
5e91098
f7bfaf2
8cba565
2b33bbe
99693fb
be8c62a
7d98073
de16d84
0e6335f
56d0f4c
52c0f28
7eb337d
e6f904b
cdeccd6
cbe5a2c
ed48d14
2b39dd1
0ba3ac3
d82dbe4
f0e2244
92613d3
cf02cce
157d227
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,9 +1,17 @@ | ||
| Slack Posters | ||
| Auxillary utility (slack posters, loggers) | ||
| ======================================================= | ||
| .. note:: | ||
| class for posting results, links, and files to slack | ||
| classes for posting results, links, and files to slack, as well as making nice loggers for files. | ||
|
|
||
| slackbot class | ||
| -------------------------------- | ||
| .. autoclass:: fast_response.slack_posters.slack.slackbot | ||
| :members: | ||
|
|
||
| logger classes | ||
| ------------------------------- | ||
| .. autoclass:: fast_response.slack_posters.logger_util.FRA_Logger | ||
| :members: | ||
|
|
||
| .. autoclass:: fast_response.slack_posters.logger_util.LogFileWriter | ||
| :members: |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,42 +1,75 @@ | ||
| #!/usr/bin/env python | ||
|
|
||
| ''' Script to automatically receive GCN notices for IceCube | ||
| alert events and run followup accordingly | ||
|
|
||
| Author: Alex Pizzuto | ||
| Date: July 2020 | ||
| Author: Alex Pizzuto, Jessie Thwaites, Alicia Mand | ||
| Updated Date: August 2026 | ||
| ''' | ||
|
|
||
| from itertools import count | ||
| import gcn | ||
| @gcn.handlers.include_notice_types( | ||
| gcn.notice_types.ICECUBE_ASTROTRACK_GOLD, | ||
| gcn.notice_types.ICECUBE_ASTROTRACK_BRONZE, | ||
| gcn.notice_types.ICECUBE_CASCADE) | ||
| import logging | ||
| from gcn_kafka import Consumer | ||
| import os, subprocess, time, pwd, argparse | ||
| import healpy as hp | ||
| import numpy as np | ||
| import lxml.etree | ||
| from astropy.time import Time | ||
| from datetime import datetime | ||
| from dateutil.parser import parse | ||
| from glob import glob | ||
| from fast_response.slack_posters.slack import slackbot | ||
| import pandas as pd | ||
|
|
||
| logger = logging.getLogger() | ||
| logger.setLevel(logging.INFO) | ||
| logger.warning("Connecting to GCN as Consumer") | ||
|
|
||
| with open('/home/jthwaites/private/tokens/kafka_token.txt') as f: | ||
| client_id = f.readline().rstrip('\n') | ||
| client_secret = f.readline().rstrip('\n') | ||
|
|
||
| domain = 'gcn.nasa.gov' | ||
| config = {'broker.address.family': 'v4', | ||
| 'log_level': 0} | ||
| consumer = Consumer(client_id=client_id, | ||
| client_secret=client_secret, | ||
| domain='gcn.nasa.gov', | ||
| config=config, | ||
| #config={'max.poll.interval.ms':1800000}, | ||
| ) | ||
|
|
||
| #consumer.subscribe(['gcn.notices.icecube.gold_bronze_track_alerts']) | ||
| # stick with voevent for now for all 3 | ||
| consumer.subscribe(['gcn.classic.voevent.ICECUBE_ASTROTRACK_BRONZE', | ||
| 'gcn.classic.voevent.ICECUBE_ASTROTRACK_GOLD', | ||
| 'gcn.classic.voevent.ICECUBE_CASCADE']) | ||
|
|
||
| def process_gcn(payload, root): | ||
| def process_gcn(record): #payload,root | ||
| analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') | ||
| if analysis_path is None: | ||
| try: | ||
| import fast_response | ||
| analysis_path = os.path.dirname(fast_response.__file__) + '/scripts/' | ||
| except Exception as e: | ||
| print(e) | ||
| logger.error('Error finding FRA package!!') | ||
| post_error("Error finding FRA package") | ||
| print('###########################################################################') | ||
| print('CANNOT FIND ENVIRONMENT VARIABLE POINTING TO REALTIME FAST RESPONSE PACKAGE\n') | ||
| print('You can either (1) install fast_response via pip or ') | ||
| print('(2) put \'export FAST_RESPONSE_SCRIPTS=/path/to/fra/scripts\' in your bashrc') | ||
| print('###########################################################################') | ||
| exit() | ||
| raise Exception(e) | ||
|
|
||
| # Read all of the VOEvent parameters from the "What" section. | ||
| params = {elem.attrib['name']: | ||
| elem.attrib['value'] | ||
| for elem in root.iterfind('.//Param')} | ||
| for elem in record.iterfind('.//Param')} | ||
|
|
||
| stream = params['Stream'] | ||
| eventtime = root.find('.//ISOTime').text | ||
| eventtime = record.find('.//ISOTime').text | ||
| if stream == '26': | ||
| print("INCOMING ALERT: ",datetime.utcnow()) | ||
| print("Detected cascade type alert, running cascade followup. . . ") | ||
| logger.warning("INCOMING ALERT: ",datetime.utcnow()) | ||
| logger.warning("Detected cascade type alert, running cascade followup. . . ") | ||
| alert_type='cascade' | ||
| event_name='IceCube-Cascade_{}{}{}'.format(eventtime[2:4],eventtime[5:7],eventtime[8:10]) | ||
|
|
||
|
|
@@ -49,8 +82,8 @@ def process_gcn(payload, root): | |
| if int(params['Rev']) !=0: | ||
| return | ||
|
|
||
| print("INCOMING ALERT: ",datetime.utcnow()) | ||
| print("Found track type alert, running track followup. . . ") | ||
| logger.warning("INCOMING ALERT: ",datetime.utcnow()) | ||
| logger.warning("Found track type alert, running track followup. . . ") | ||
|
|
||
| event_id = params['event_id'] | ||
| run_id = params['run_id'] | ||
|
|
@@ -70,7 +103,7 @@ def process_gcn(payload, root): | |
| needed_delay = 1. | ||
| current_delay = current_mjd - event_mjd | ||
| while current_delay < needed_delay: | ||
| print("Need to wait another {:.1f} seconds before running".format( | ||
| logger.info("Need to wait another {:.1f} seconds before running".format( | ||
| (needed_delay - current_delay)*86400.) | ||
| ) | ||
| time.sleep((needed_delay - current_delay)*86400.) | ||
|
|
@@ -82,12 +115,12 @@ def process_gcn(payload, root): | |
| skymap_f = glob(base_skymap_path \ | ||
| + f'run{int(run_id):08d}.evt{int(event_id):012d}.*probability.fits.gz') | ||
| if len(skymap_f) == 0: | ||
| print("COULD NOT FIND THE SKYMAP FILE FOR V2 TRACK ALERT EVENT") | ||
| logger.error("COULD NOT FIND THE SKYMAP FILE FOR V2 TRACK ALERT EVENT") | ||
| return | ||
| elif len(skymap_f) == 1: | ||
| skymap = skymap_f[0] | ||
| else: | ||
| print("TOO MANY OPTIONS FOR THE SKYMAP FILE FOR V2 TRACK ALERT EVENT") | ||
| logger.error("TOO MANY OPTIONS FOR THE SKYMAP FILE FOR V2 TRACK ALERT EVENT") | ||
| return | ||
|
|
||
| #checking for events on the same day: looks for existing output files from previous runs | ||
|
|
@@ -98,13 +131,16 @@ def process_gcn(payload, root): | |
| elif count_dir==2: suffix='B' | ||
| elif count_dir==4: suffix='C' | ||
| else: | ||
| print("COULD NOT DETERMINE EVENT SUFFIX") | ||
| print("check for other events on the same day and re-run with args:") | ||
| print('--skymap={} --time={} --alert_id={}'.format(skymap, str(event_mjd), run_id+':'+event_id)) | ||
| logger.error("COULD NOT DETERMINE EVENT SUFFIX") | ||
| logger.error("check for other events on the same day and re-run with args:") | ||
| logger.error('--skymap={} --time={} --alert_id={}'.format(skymap, str(event_mjd), run_id+':'+event_id)) | ||
| return | ||
|
|
||
| print('\nRunning {} --skymap={} --time={} --alert_id={} --suffix={}'.format( | ||
| logger.info('\nRunning {} --skymap={} --time={} --alert_id={} --suffix={}'.format( | ||
| command, skymap, str(event_mjd), run_id+':'+event_id, suffix)) | ||
|
|
||
| print(params) | ||
|
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. do you need this |
||
|
|
||
| subprocess.call([command, '--skymap={}'.format(skymap), | ||
| '--time={}'.format(str(event_mjd)), | ||
| '--alert_id={}'.format(run_id+':'+event_id), | ||
|
|
@@ -123,7 +159,8 @@ def process_gcn(payload, root): | |
| subprocess.call([analysis_path+'document.py', '--path', dir_2d[0]]) | ||
| doc=True | ||
| except: | ||
| print('Failed to document to private webpage') | ||
| post_error("Failed to run document command") | ||
| logger.warning('Failed to document to private webpage') | ||
|
|
||
| try: | ||
| shifters = pd.read_csv(os.path.join(analysis_path,'../slack_posters/fra_shifters.csv'), | ||
|
|
@@ -146,24 +183,34 @@ def process_gcn(payload, root): | |
|
|
||
| bot.post_short_msg(done_message) | ||
| except Exception as e: | ||
| print(e) | ||
| post_error("Failed to push results to private webpage") | ||
| logger.warning('Failed to push to private webpage') | ||
| logger.warning(e) | ||
|
|
||
| def post_error(errMsg=None): | ||
|
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. could this be added as a general utility to |
||
| analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') | ||
| bot = slackbot('fra-shifting') | ||
| if errMsg != None: | ||
| message = errMsg | ||
| else: | ||
| message = "ERROR in FRA, please check internal listener!" | ||
| try: | ||
| shifters = pd.read_csv(os.path.join(analysis_path, '../slack_posters/fra_shifters.csv'), parse_dates=[0,1]) | ||
| on_shift = '' | ||
| for i in shifters.index: | ||
| if shifters['start'][i] < datetime.utcnow() < shifters['stop'][i]: | ||
| on_shift+='<@{}> '.format(shifters['slack_id'][i]) | ||
| error_message = f"{message} {on_shift} on shift." | ||
| bot.post_short_msg(error_message) | ||
| except Exception as e: | ||
| logger.warning("Failed to post error message") | ||
| logger.warning(e) | ||
| return | ||
|
|
||
| if __name__ == '__main__': | ||
| import os, subprocess | ||
| import healpy as hp | ||
| import numpy as np | ||
| import lxml.etree | ||
| import argparse | ||
| from astropy.time import Time | ||
| from datetime import datetime | ||
| from dateutil.parser import parse | ||
| import time | ||
| from glob import glob | ||
| from fast_response.slack_posters.slack import slackbot | ||
| import pandas as pd | ||
| import pwd | ||
|
|
||
| username = pwd.getpwuid(os.getuid())[0] | ||
|
|
||
| #default for if to document or not: only way to check reports on realtime | ||
| if username == 'realtime': | ||
| document = True | ||
|
|
@@ -175,32 +222,50 @@ def process_gcn(payload, root): | |
| help='Run on live GCNs') | ||
| parser.add_argument('--test_cascade', default=False, action='store_true', | ||
| help='When testing, raise to run a cascade, else track') | ||
| parser.add_argument('--test_error', default=False, action='store_true', | ||
| help='When testing, raise to post an error message') | ||
| parser.add_argument('--document', action='store_true', default=document, | ||
| help='flag to raise to push results to internal webpage') | ||
| args = parser.parse_args() | ||
|
|
||
| if args.run_live: | ||
| print("Listening for GCNs . . . ") | ||
| gcn.listen(handler=process_gcn) | ||
| logger.warning("Listening for IC Alert GCNs . . . ") | ||
| while True: | ||
| for message in consumer.consume(timeout=1): | ||
| if message.error(): | ||
| logger.warning(message.error()) | ||
| continue | ||
| # value = message.value().decode('utf-8') | ||
| # value = value.replace("<?xml version='1.0' encoding='UTF-8'?>","") #lxml doesn't like this line | ||
| value = message.value() | ||
| logger.warning('Found GCN on topic {}'.format(message.topic())) | ||
| notice = lxml.etree.fromstring(value) | ||
| try: | ||
| process_gcn(notice) | ||
| except Exception as e: | ||
| post_error("Could not process GCN") | ||
| logger.warning("Could not process GCN: ", e) | ||
|
|
||
| else: | ||
| try: | ||
| import fast_response | ||
| sample_skymap_path=os.path.dirname(fast_response.__file__) +'/sample_skymaps/' | ||
| except Exception as e: | ||
| #future: possibly point to FRA on /data/ana/ | ||
| print(e) | ||
| print('Cannot find path to sample skymaps') | ||
| exit() | ||
|
|
||
| if not args.test_cascade: | ||
| print("Running on sample track . . . ") | ||
| post_error("Cannot find path to sample skymaps") | ||
|
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. do you want a posted error if testing? I think this won't really be helpful, I recommend just a logger error here, since the user will be running in testing mode to hit the |
||
| logger.error('Cannot find path to sample skymaps') | ||
| raise Exception(e) | ||
| if args.test_error: | ||
| post_error() | ||
|
|
||
| if not args.test_cascade and not args.test_error: | ||
| logger.info("Running on sample track . . . ") | ||
| payload = open(sample_skymap_path \ | ||
| + 'sample_astrotrack_alert_2021.xml', 'rb').read() | ||
| root = lxml.etree.fromstring(payload) | ||
| process_gcn(payload, root) | ||
| else: | ||
| print("Running on sample cascade . . . ") | ||
| process_gcn(root) | ||
| elif args.test_cascade and not args.test_error: | ||
| logger.info("Running on sample cascade . . . ") | ||
| payload = open(sample_skymap_path \ | ||
| + 'sample_cascade.txt', 'rb').read() | ||
| root = lxml.etree.fromstring(payload) | ||
| process_gcn(payload, root) | ||
| process_gcn(root) | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
suggest dropping the commented out stuff here