diff --git a/doc/source/slack.rst b/doc/source/slack.rst index df81c676..daa30230 100644 --- a/doc/source/slack.rst +++ b/doc/source/slack.rst @@ -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: diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 912df3e3..c79daffb 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -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) + 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): + 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("","") #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") + 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) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 8692431b..652e0ce6 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -3,80 +3,102 @@ ''' Script to automatically receive GCN alerts and get LIGO skymaps to run realtime neutrino follow-up - Author: Raamis Hussain, updated by Jessie Thwaites, MJ Romfoe - Date: March 2023 + Author: Raamis Hussain, Jessie Thwaites, MJ Romfoe + Last Updated: Sept 2026 ''' -import gcn -import sys -import pickle +import logging +from gcn_kafka import Consumer +import sys, pickle, os, subprocess, pwd from dateutil.parser import parse from dateutil.relativedelta import relativedelta - -@gcn.handlers.include_notice_types( - gcn.notice_types.LVC_PRELIMINARY, - gcn.notice_types.LVC_INITIAL, - gcn.notice_types.LVC_UPDATE) - -def process_gcn(payload, root): - +import healpy as hp +import numpy as np +import argparse, time, wget +from astropy.time import Time +from datetime import datetime +import fast_response +from fast_response.slack_posters.slack import slackbot +from fast_response.slack_posters.logger_util import FRA_Logger +import json + +print("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') + +config = {'broker.address.family': 'v4', + 'log_level': 0, + 'max.poll.interval.ms': 1800000, + } + +consumer = Consumer(client_id=client_id, + client_secret=client_secret, + domain='gcn.nasa.gov', + config=config, + ) + +consumer.subscribe(['igwn.gwalert']) + +def process_gcn(params, mock=False): AlertTime=datetime.utcnow().isoformat() - log_file.flush() analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') + if analysis_path is None: try: - import fast_response analysis_path = os.path.join(os.path.dirname(fast_response.__file__),'scripts/') except Exception as e: - print(e) + logger.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('###########################################################################') - log_file.flush() - exit() - - # Read all of the VOEvent parameters from the "What" section. - params = {elem.attrib['name']: - elem.attrib['value'] - for elem in root.iterfind('.//Param')} - name = root.attrib['ivorn'].split('#')[1] - + raise Exception(e) + + name = params['superevent_id'] + '-'+ +params['time_created'].replace(':','').replace('-','') + '-' + params['alert_type'].lower() + if params['alert_type'].lower() == 'retraction': + logger.warning('Listener does not run on Retractions. Skipping...') + return + params['role'] = 'observation' if 'MS' in params['superevent_id'] else 'test' + if 'search' in params['event']: # one more check to identify mocks or testing + if params['event']['search'] == 'MDC': + params['role'] = 'test' + # only run on significant events - if 'Significant' in params.keys(): - if int(params['Significant'])==0: - #not significant, do not run - print(f'Found a subthreshold event {name}') - root.attrib['role']='test' - log_file.flush() - #return + if 'significant' in params['event']: + if not params['event']['significant']: + #not significant, do not run on real data + logger.warning(f'Found a subthreshold event {name}') + params['role']='test' else: # O3 does not have this parameter, this should only happen for testing - print('No significance parameter found in LVK GCN.') - log_file.flush() + logger.warning('No significance parameter found in LVK GCN.') # if this is the listener for real events and it gets a mock (or low signficance), skip it - if not mock and root.attrib['role']!='observation': + if not mock and params['role']!='observation': + return + # want heartbeat listener not to run on real events, otherwise it overwrites the main listener output + if mock and params['role']=='observation': + logger.info('Listener in heartbeat mode found real event. Skipping...') return - print('\n' +'INCOMING ALERT FOUND: ',datetime.utcnow()) - log_file.flush() - + logger.warning('\nINCOMING ALERT FOUND: {}'.format(datetime.utcnow().strftime('%Y-%m-%d %H:%M:%S'))) #get type of event (burst, bbh, nsbh, bns) try: - if params['Group'] == 'Burst': + if params['event']['group'] == 'Burst': merger_type = 'Burst' - elif params['Search'] == 'SSM': + elif params['event']['search'] == 'SSM': merger_type='SSM' else: k = ['BNS','NSBH','BBH'] - probs = {j: float(params[j]) for j in k} + probs = {j: float(params['event']['classification'][j]) for j in k} merger_type = max(zip(probs.values(), probs.keys()))[1] except: - print('Could not determine type of event') + logger.warning('Could not determine type of event') merger_type = None - if root.attrib['role']=='observation' and not mock: + if params['role']=='observation' and not mock: ## Call everyone because it's a real event! call_command=['/home/jthwaites/private/make_call.py', f'--name={name}'] @@ -88,24 +110,15 @@ def process_gcn(payload, root): try: subprocess.call(call_command) - #print('Call here.') except Exception as e: - print('Call failed.') - print(e) - log_file.flush() - - # want heartbeat listener not to run on real events, otherwise it overwrites the main listener output - if mock and root.attrib['role']=='observation': - print('Listener in heartbeat mode found real event. Returning...') - log_file.flush() - return + logger.error('Call failed!') + logger.error(e) # Read trigger time of event - eventtime = root.find('.//ISOTime').text + eventtime = params['event']['time'] event_mjd = Time(eventtime, format='isot').mjd - print(f'Alert MJD: {event_mjd}') - print('GW merger time: %s \n' % Time(eventtime, format='isot').iso) - log_file.flush() + logger.info(f'Alert MJD: {event_mjd}') + logger.info('GW merger time: {} \n'.format(Time(eventtime, format='isot').iso)) current_mjd = Time(datetime.utcnow(), scale='utc').mjd needed_delay = 1000./84600./2. @@ -116,15 +129,15 @@ def process_gcn(payload, root): FiveHundred_delay = (needed_delay - current_delay)*86400. 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.) ) - log_file.flush() time.sleep((needed_delay - current_delay)*86400.) current_mjd = Time(datetime.utcnow(), scale='utc').mjd current_delay = current_mjd - event_mjd - skymap = params['skymap_fits'] + skymap_base = 'https://gracedb.ligo.org/api/superevents/{}/files/'.format(params['superevent_id']) + skymap = skymap_base + params['event']['skymap_filename'] # Multiorder Coverage (MOC) map links are distributed over the GCNs. # Download flattened (normal healpy) map from GraceDB @@ -143,8 +156,7 @@ def process_gcn(payload, root): wget.download(new_map, out=os.path.join(os.environ.get('FAST_RESPONSE_OUTPUT'),f'skymaps/{name}_{map_type}{suffix}')) skymap=os.path.join(os.environ.get('FAST_RESPONSE_OUTPUT'),f'skymaps/{name}_{map_type}{suffix}') except: - print('Failed to download flat-resolution skymap. Trying to convert MOC map') - log_file.flush() + logger.warning('Failed to download flat-resolution skymap. Trying to convert MOC map') try: filename=skymap.split('/')[-1] @@ -154,37 +166,34 @@ def process_gcn(payload, root): '--skymap', new_output]) if os.path.exists(new_output.replace('multiorder','converted')): skymap = new_output.replace('multiorder','converted') - print('Successfully converted map: {}'.format(skymap)) - log_file.flush() + logger.info('Successfully converted map: {}'.format(skymap)) else: raise Exception('Failed to convert map.') except: - print('Failed to get skymap in correct format! \nDownload skymap and then re-run script with') - print(f'args: --time {event_mjd} --name {name} --skymap PATH_TO_SKYMAP') - log_file.flush() + logger.error('Failed to get skymap in correct format! \nDownload skymap and then re-run script with' +\ + f'args: --time {event_mjd} --name {name} --skymap PATH_TO_SKYMAP') return - if root.attrib['role'] != 'observation': + if params['role'] != 'observation': name=name+'_test' - print('Running on scrambled data') - log_file.flush() + logger.info('Running on scrambled data') command = os.path.join(analysis_path, 'run_gw_followup.py') - print('Running {}'.format(command)) - log_file.flush() + logger.info('Running {}'.format(command)) - subprocess.call([command, '--skymap={}'.format(skymap), + subprocess.call([ + command, + '--skymap={}'.format(skymap), '--time={}'.format(str(event_mjd)), - '--name={}'.format(name)] - #'--allow_neg_ts=True'] - ) + '--name={}'.format(name) + ]) analysis_start = Time(event_mjd - 500./86400., format='mjd').iso output = os.path.join(os.environ.get('FAST_RESPONSE_OUTPUT'), analysis_start[0:10].replace('-','_')+'_'+name) #update webpages webpage_update = os.path.join(analysis_path,'document.py') - if not mock and root.attrib['role'] == 'observation': + if not mock and params['role'] == 'observation': try: subprocess.call([webpage_update, '--gw', f'--path={output}']) @@ -197,9 +206,8 @@ def process_gcn(payload, root): bot.post_short_msg(slack_message) except Exception as e: - print('Failed to push to (private) webpage.') - print(e) - log_file.flush() + logger.error('Failed to push to (private) webpage.') + logger.error(e) endtime=datetime.utcnow().isoformat() alert_mjd = Time(AlertTime, format='isot').mjd @@ -214,7 +222,7 @@ def process_gcn(payload, root): 'Ligo_Latency': Ligo_late_sec, 'IceCube_Latency': Ice_late_sec, 'Total_Latency': Total_late_sec, 'We_had_to_wait:': FiveHundred_delay} - save_dir = 'latency_o4' if root.attrib['role']=='observation' else 'PickledMocks' + save_dir = 'latency_o4' if params['role']=='observation' else 'PickledMocks' #check for directory to save pickle files and create if needed if not os.path.exists(os.path.join(os.environ.get('FAST_RESPONSE_OUTPUT'),save_dir)): @@ -226,43 +234,26 @@ def process_gcn(payload, root): with open(os.path.join(os.environ.get('FAST_RESPONSE_OUTPUT'), f'{save_dir}/gw_latency_dict_{name}.pickle'), 'wb') as file: pickle.dump(gw_latency, file, protocol=pickle.HIGHEST_PROTOCOL) - #save xml and skymap, for later - et = lxml.etree.ElementTree(root) - et.write(os.path.join(output, '{}-{}-{}.xml'.format(params['GraceID'], - params['Pkt_Ser_Num'], params['AlertType'])), pretty_print=True) - - if root.attrib['role'] != 'observation': + #save notice for later + with open(os.path.join(output, '{}.json'.format(name)), "w") as f: + f.write(json.dumps(params, indent=2)) + + if params['role'] != 'observation': # Move mocks to a seperate folder to avoid swamping FRA output folder subprocess.call(['mv',output, '/data/user/jthwaites/o4-mocks/']) output = '/data/user/jthwaites/o4-mocks/' + eventtime[0:10].replace('-','_')+'_'+name - print('Output directory: ',output) - log_file.flush() + logger.info('Output directory: {}'.format(output)) if __name__ == '__main__': - import os, subprocess, pwd - import healpy as hp - import numpy as np - import lxml.etree - import argparse - import time - from astropy.time import Time - from datetime import datetime - from fast_response.slack_posters.slack import slackbot - import wget - - output_path = '/home/jthwaites/public_html/FastResponse/' - #output_path=os.environ.get('FAST_RESPONSE_OUTPUT') - #if output_path==None: - # output_path=os.getcwd() parser = argparse.ArgumentParser(description='FRA GW followup') parser.add_argument('--run_live', action='store_true', default=False, help='Run on live GCNs') parser.add_argument('--heartbeat', action = 'store_true', default=False, help='Run the listener as a heartbeat, running on mock LVK events only (default=False)') - parser.add_argument('--log_path', default=output_path, type=str, - help='Redirect output to a log file with this path') + parser.add_argument('--log_path', default='/home/jthwaites/public_html/FastResponse/', type=str, + help='Include output to a log file with this path. Note: this is only used when running live') parser.add_argument('--test_path', default='S191216ap_update.xml', type=str, help='Skymap for use in testing listener') parser.add_argument('--test_o3', default=False, action='store_true', @@ -274,43 +265,53 @@ def process_gcn(payload, root): else: logfile=os.path.join(args.log_path,'log.log') - print(f'Logging to file: {logfile}') - original_stdout=sys.stdout - log_file = open(logfile, "a+") - sys.stdout=log_file - sys.stderr=log_file - if args.run_live: - print("Listening for GCNs . . . ") - log_file.flush() + print(f'Logging to file: {logfile}') + + logger = FRA_Logger(file=logfile).logger + logger.warning("Listening for GCNs . . . ") mock=args.heartbeat - print('Starting heartbeat listener') if mock else print('Running on REAL events only') - log_file.flush() - - gcn.listen(handler=process_gcn) - - else: - print("Offline testing . . . ") - log_file.flush() + logger.info('Starting heartbeat listener') if mock else logger.info('Running on REAL events only') - ### FOR OFFLINE TESTING try: - import fast_response - #sample_skymap_path='/data/user/jthwaites/o3-gw-skymaps/' - sample_skymap_path=os.path.join(os.path.dirname(fast_response.__file__),'sample_skymaps/') - except Exception as e: - print(e) - sample_skymap_path='/data/user/jthwaites/o3-gw-skymaps/' - - #payload = open(os.path.join(sample_skymap_path,args.test_path), 'rb').read() - payload = open(args.test_path,'rb').read() - root = lxml.etree.fromstring(payload) + while True: + for message in consumer.consume(timeout=1): + if message.error(): + logger.warning(message.error()) + continue + value = message.value().decode('utf-8') + logger.warning('Found GCN on topic {}'.format(message.topic())) + notice = json.loads(value) + if notice['alert_type'].lower() == 'retraction': + # retractions do not have some required quantities, and should not be run + continue + process_gcn(notice,mock=mock) + except KeyboardInterrupt: + # make sure the logfile gets shutdown correctly and file closed + logging.shutdown() + + else: + logger = FRA_Logger() + logger.warning("Offline testing . . . ") + + # see if we've been passed an absolute path. + if os.path.exists(args.test_path): + test_file = args.test_path + else: + test_file = os.path.join(os.environ.get('I3_SRC'),'realtime_scripts/resources/test', args.test_path) + # if it still isn't found, exit + if not os.path.exists(test_file): + logger.error('Failed to find test file at {}. Check path and try again'.format(test_file)) + sys.exit() + + logger.info('Running offline on file: {}'.format(test_file)) + params = json.loads(test_file) mock=args.heartbeat #test runs on scrambles, observation runs on unblinded data if not args.test_o3: - root.attrib['role']='test' + params['event']['search'] = 'MDC' mock=True - process_gcn(payload, root) + process_gcn(params, mock=mock) diff --git a/fast_response/scripts/combine_results_classic.py b/fast_response/scripts/combine_results_classic.py deleted file mode 100644 index a7df1659..00000000 --- a/fast_response/scripts/combine_results_classic.py +++ /dev/null @@ -1,591 +0,0 @@ -#!/usr/bin/env python - -'''Note: this version listens to GCN CLASSIC. -Main sender is combine_results_kafka.py''' - -import logging -from datetime import datetime -import socket -import requests -import healpy as hp -import matplotlib as mpl -mpl.use('agg') -import matplotlib.pyplot as plt -import io, time, os, glob, subprocess -import urllib.request, urllib.error, urllib.parse -import argparse -import json, pickle -from gcn_kafka import Consumer -import gcn -#from icecube import realtime_tools -import numpy as np -import lxml.etree -from astropy.time import Time -import dateutil.parser -from datetime import datetime -from fast_response.web_utils import updateGW_public - -parser = argparse.ArgumentParser(description='Combine GW-Nu results') -parser.add_argument('--run_live', action='store_true', default=False, - help='Run on live GCNs') -parser.add_argument('--test_path', type=str, default=None, - help='path to test xml file') -parser.add_argument('--max_wait', type=float, default=60., - help='Maximum minutes to wait for LLAMA/UML results before timeout (default=60)') -parser.add_argument('--wait_for_llama', action='store_true', default=False, - help='bool to decide to send llama results with uml, default false while not unblinded') -parser.add_argument('--heartbeat', action='store_true', default=False, - help='bool to save jsons for mocks (saves to save_dir+mocks/)') -parser.add_argument('--save_dir', type=str, default='/home/followup/lvk_followup_output/', - help='Directory to save output json (default=/home/followup/lvk_followup_output/)') -args = parser.parse_args() - -with open('/cvmfs/icecube.opensciencegrid.org/users/jthwaites/tokens/kafka_token.txt') as f: - client_id = f.readline().rstrip('\n') - client_secret = f.readline().rstrip('\n') - -consumer = Consumer(client_id=client_id, - client_secret=client_secret) - -# Subscribe to topics to receive alerts -consumer.subscribe(['gcn.classic.voevent.LVC_PRELIMINARY', - 'gcn.classic.voevent.LVC_INITIAL', - 'gcn.classic.voevent.LVC_UPDATE']) -#consumer.subscribe(['igwn.gwalert']) - -def SendAlert(results=None): - from gcn_kafka import Producer - - if results is None: - logger.fatal('Found no alert to send') - - with open('/cvmfs/icecube.opensciencegrid.org/users/jthwaites/tokens/real_icecube_kafka_prod.txt') as f: - prod_id = f.readline().rstrip('\n') - prod_secret = f.readline().rstrip('\n') - - producer = Producer(#config={'bootstrap.servers': 'kafka3.gcn.nasa.gov'}, - client_id=prod_id, - client_secret=prod_secret, - domain='gcn.nasa.gov') - - if 'MS' in results['ref_id']: - return - topic = 'gcn.notices.icecube.test.lvk_nu_track_search' - else: - topic = 'gcn.notices.icecube.lvk_nu_track_search' - - logger.info('sending to {} on gcn.nasa.gov'.format(topic)) - sendme = json.dumps(results) - producer.produce(topic, sendme.encode()) - ret = producer.flush() - return ret - -def SendTestAlert(results=None): - from gcn_kafka import Producer - - if results is None: - logger.fatal('Found no alert to send') - - with open('/cvmfs/icecube.opensciencegrid.org/users/jthwaites/tokens/test_icecube_kafka_prod.txt') as f: - prod_id = f.readline().rstrip('\n') - prod_secret = f.readline().rstrip('\n') - - producer = Producer(#config={'bootstrap.servers': 'kafka3.gcn.nasa.gov'}, - client_id=prod_id, - client_secret=prod_secret, - domain='test.gcn.nasa.gov') - - topic = 'gcn.notices.icecube.test.lvk_nu_track_search' - logger.info('sending to {} on test.gcn.nasa.gov'.format(topic)) - - sendme = json.dumps(results) - producer.produce(topic, sendme.encode()) - ret = producer.flush() - return ret - -def format_ontime_events_uml(events, event_mjd): - ontime_events={} - for event in events: - ontime_events[event['event']]={ - 'event_dt' : round((event['time']-event_mjd)*86400., 2), - 'localization':{ - 'ra' : round(np.rad2deg(event['ra']), 2), - 'dec' : round(np.rad2deg(event['dec']), 2), - "uncertainty_shape": "circle", - 'ra_uncertainty': round(np.rad2deg(event['sigma']*2.145966),2), - "containment_probability": 0.9, - "systematic_included": False - }, - 'event_pval_generic' : round(event['pvalue'],4) - } - return ontime_events - -def format_ontime_events_llama(events): - ontime_events={} - for event in events: - ontime_events[event['i3event']] = { - 'event_dt' : round(event['dt'],2), - 'localization':{ - 'ra' : round(event['ra'], 2), - 'dec' : round(event['dec'], 2), - "uncertainty_shape": "circle", - 'ra_uncertainty': round(np.rad2deg(np.deg2rad(event['sigma'])*2.145966),3), - "containment_probability": 0.9, - "systematic_included": False - }, - 'event_pval_bayesian': round(event['p_value'],4) - } - return ontime_events - -def format_ontime_events_llama_old(events,event_mjd): - ontime_events={} - for event in events: - ontime_events[event['event']] = { - 'event_dt' : round((event['mjd']-event_mjd)*86400.,2), - 'localization':{ - 'ra' : round(np.rad2deg(event['ra']), 2), - 'dec' : round(np.rad2deg(event['dec']), 2), - "uncertainty_shape": "circle", - 'ra_uncertainty': round(np.rad2deg(event['sigma']*2.145966),3), - "containment_probability": 0.9, - "systematic_included": False - }, - 'event_pval_bayesian': 'null' - } - return ontime_events - -def combine_events(uml_ontime, llama_ontime): - #in the case both return coincident events, want to combine their results - event_ids_all = np.unique(list(uml_ontime.keys())+list(llama_ontime.keys())) - coinc_events = [] - - for id in event_ids_all: - #case: event in both results - if (id in uml_ontime.keys()) and (id in llama_ontime.keys()): - if (uml_ontime[id]['event_pval_generic'] < 0.1) and (llama_ontime[id]['event_pval_bayesian']<0.1): - uml_ontime[id]['event_pval_bayesian'] = llama_ontime[id]['event_pval_bayesian'] - coinc_events.append(uml_ontime[id]) - elif uml_ontime[id]['event_pval_generic'] < 0.1: - uml_ontime[id]['event_pval_bayesian'] = 'null' - coinc_events.append(uml_ontime[id]) - elif llama_ontime[id]['event_pval_bayesian']<0.1: - uml_ontime[id]['event_pval_bayesian'] = llama_ontime[id]['event_pval_bayesian'] - uml_ontime[id]['event_pval_generic'] ='null' - coinc_events.append(uml_ontime[id]) - #case: only in uml - elif id in uml_ontime.keys(): - if uml_ontime[id]['event_pval_generic'] < 0.1: - uml_ontime[id]['event_pval_bayesian'] = 'null' - coinc_events.append(uml_ontime[id]) - #case: only in llama - elif id in llama_ontime.keys(): - if llama_ontime[id]['event_pval_bayesian']<0.1: - llama_ontime[id]['event_pval_generic']='null' - coinc_events.append(llama_ontime[id]) - - return coinc_events - -@gcn.handlers.include_notice_types( - gcn.notice_types.LVC_PRELIMINARY, - gcn.notice_types.LVC_INITIAL, - gcn.notice_types.LVC_UPDATE) - -def parse_notice(payload, record): - logger = logging.getLogger() - wait_for_llama=args.wait_for_llama - heartbeat = args.heartbeat - - if record.attrib['role']!='observation': - fra_results_location = '/data/user/jthwaites/o4-mocks/' - if not heartbeat: - logger.warning('found test event - not in mock mode. returning') - return - else: - logger.warning('found test event') - else: - logger.warning('ALERT FOUND') - fra_results_location = os.environ.get('FAST_RESPONSE_OUTPUT') - - ### GENERAL PARAMETERS ### - #read event information - params = {elem.attrib['name']: - elem.attrib['value'] - for elem in record.iterfind('.//Param')} - - # ignore subthreshold, only run on significant events - subthreshold=False - if 'Significant' in params.keys(): - if int(params['Significant'])==0: - subthreshold=True - logger.warning('low-significance alert found. ') - if params['Group'] == 'Burst': - wait_for_llama = False - m = 'Significant' if not subthreshold else 'Subthreshold' - logger.warning('{} burst alert found. '.format(m)) - if len(params['Instruments'].split(','))==1: - #wait_for_llama = False - logger.warning('One detector event found. ') - - if not wait_for_llama and subthreshold: - #if llama isn't running, don't run on subthreshold - logger.warning('Not waiting for LLAMA. Returning...') - return - - collected_results = {} - collected_results["$schema"]= "https://gcn.nasa.gov/schema/gcn/notices/icecube/LvkNuTrackSearch.schema.json" - collected_results["type"]= "IceCube LVK Alert Nu Track Search" - - eventtime = record.find('.//ISOTime').text - event_mjd = Time(eventtime, format='isot').mjd - name = record.attrib['ivorn'].split('#')[1] - - if record.attrib['role'] != 'observation': - name=name+'_test' - - logger.info("{} alert found, processing GCN".format(name)) - collected_results['ref_id'] = name - - #collected_results['ref_id'] = name.split('-')[0] - #collected_results['id'] = [name.split('-')[2],'Sequence:{}'.format(name.split('-')[1])] - - ### WAIT ### - # wait until the 500s has elapsed for data - current_mjd = Time(datetime.utcnow(), scale='utc').mjd - needed_delay = 500./84600. - #needed_delay = 720./86400. #12 mins while lvk_dropbox is offline - current_delay = current_mjd - event_mjd - - while current_delay < needed_delay: - logger.info("Wait {:.1f} seconds for data".format( - (needed_delay - current_delay)*86400.) - ) - time.sleep((needed_delay - current_delay)*86400.) - current_mjd = Time(datetime.utcnow(), scale='utc').mjd - current_delay = current_mjd - event_mjd - - # start checking for results - results_done = False - start_check_mjd = Time(datetime.utcnow(), scale='utc').mjd - max_delay = max_wait/60./24. #mins to days - - logger.info("Looking for results (max {:.0f} mins)".format( - max_wait) - ) - - while results_done == False: - start_date = Time(dateutil.parser.parse(eventtime)).datetime - start_str = f'{start_date.year:02d}_{start_date.month:02d}_{start_date.day:02d}' - - uml_results_path = os.path.join(fra_results_location, start_str + '_' + name.replace(' ', '_') \ - + '/' + start_str + '_' + name.replace(' ', '_')+'_results.pickle') - uml_results_finished = os.path.exists(uml_results_path) - - if len(params['Instruments'].split(','))==1: - llama_name = '{}.significance_opa_lvc-i3.json'.format(record.attrib['ivorn'].split('#')[1]) - else: - llama_name = '{}.significance_subthreshold_lvc-i3.json'.format(record.attrib['ivorn'].split('#')[1]) - llama_results_path = os.path.join(llama_results_location,llama_name) - if wait_for_llama: - llama_results_finished = os.path.exists(llama_results_path) - else: - llama_results_finished = False - #llama_name = '{}*-significance_subthreshold_lvc-i3.json'.format(record.attrib['ivorn'].split('#')[1]) - #llama_results_path = os.path.join(llama_results_location,llama_name) - #if wait_for_llama: - # llama_results_glob = sorted(glob.glob(llama_results_path)) - # if len(llama_results_glob)>0: - # llama_results_path = llama_results_glob[-1] - # llama_results_finished=True - # else: - # llama_results_finished=False - #else: - # llama_results_finished = False - - if subthreshold and llama_results_finished: - #no UML results in subthreshold case, only LLAMA - results_done = True - uml_results_finished = False - logger.info('found results for LLAMA for subthreshold event, writing notice') - - if uml_results_finished and llama_results_finished: - results_done=True - logger.info('found results for both, writing notice') - elif uml_results_finished and not wait_for_llama: - results_done = True - logger.info('found results for UML, writing notice') - else: - current_mjd = Time(datetime.utcnow(), scale='utc').mjd - if current_mjd - start_check_mjd < max_delay: - time.sleep(10.) - else: - if (uml_results_finished or llama_results_finished): - if uml_results_finished: - logger.warning('LLAMA results not finished in {:.0f} mins. Sending UML only'.format(max_wait)) - results_done=True - if llama_results_finished: - logger.warning('UML results not finished in {:.0f} mins. Sending LLAMA only'.format(max_wait)) - results_done=True - else: - logger.warning('Both analyses not finished after {:.0f} min wait.'.format(max_wait)) - logger.warning('Not sending GCN.') - return - - collected_results['alert_datetime'] = '{}Z'.format(Time(datetime.utcnow(), scale='utc').isot) - collected_results['trigger_time'] = eventtime if 'Z' in eventtime else '{}Z'.format(eventtime) - collected_results['observation_start'] = '{}Z'.format(Time(event_mjd - 500./86400., format = 'mjd').isot) - collected_results['observation_stop'] = '{}Z'.format(Time(event_mjd + 500./86400., format = 'mjd').isot) - collected_results['observation_livetime'] = 1000 - - ### COLLECT RESULTS ### - send_notif = False - - if uml_results_finished: - with open(uml_results_path, 'rb') as f: - uml_results = pickle.load(f) - - if llama_results_finished: - with open(llama_results_path, 'r') as f: - llama_results = json.load(f) - if (record.attrib['role'] == 'observation') and (llama_results['inputs']['neutrino_info'][0]['type']=='blinded'): - logger.warning('LLAMA results blinded for real event! Skipping LLAMA') - llama_results_finished = False - if subthreshold: - logger.warning('LLAMA skipped on subthreshold event, returning') - return - - if uml_results_finished and llama_results_finished: - collected_results['pval_generic'] = round(uml_results['p'],4) - collected_results['pval_bayesian'] = round(llama_results['p_value'],4) - if (collected_results['pval_generic']<0.01) or (collected_results['pval_bayesian']<0.01): - send_notif=True - - if (collected_results['pval_generic']<0.1) or (collected_results['pval_bayesian']<0.1): - uml_ontime = format_ontime_events_uml(uml_results['coincident_events'], event_mjd) - llama_ontime = format_ontime_events_llama(llama_results['single_neutrino']) - coinc_events = combine_events(uml_ontime, llama_ontime) - collected_results['n_events_coincident'] = len(coinc_events) - collected_results['coincident_events'] = coinc_events - - collected_results['most_likely_direction'] = { - 'ra': round(np.rad2deg(uml_results['fit_ra']), 2), - 'dec' : round(np.rad2deg(uml_results['fit_dec']), 2), - } - else: - collected_results['n_events_coincident'] = 0 - - collected_results["neutrino_flux_sensitivity_range"] = { - 'flux_sensitivity' : [ - round(uml_results['sens_range'][0],4), - round(uml_results['sens_range'][1],4) - ], - 'sensitive_energy_range' : [ - round(float('{:.2e}'.format(uml_results['energy_range'][0]))), - round(float('{:.2e}'.format(uml_results['energy_range'][1]))) - ], - } - - elif uml_results_finished: - collected_results['pval_generic'] = round(uml_results['p'],4) - collected_results['pval_bayesian'] = 'null' - if collected_results['pval_generic']<0.01: - send_notif=True - - if collected_results['pval_generic'] <0.1: - uml_ontime = format_ontime_events_uml(uml_results['coincident_events'], event_mjd) - coinc_events=[] - for eventid in uml_ontime.keys(): - if (uml_ontime[eventid]['event_pval_generic'] < 0.1): - uml_ontime[eventid]['event_pval_bayesian'] = 'null' - coinc_events.append(uml_ontime[eventid]) - collected_results['n_events_coincident'] = len(coinc_events) - collected_results['coincident_events'] = coinc_events - - collected_results['most_likely_direction'] = { - 'ra': round(np.rad2deg(uml_results['fit_ra']), 2), - 'dec' : round(np.rad2deg(uml_results['fit_dec']), 2) - } - - else: - collected_results['n_events_coincident'] = 0 - - collected_results["neutrino_flux_sensitivity_range"] = { - 'flux_sensitivity' : [ - round(uml_results['sens_range'][0],4), - round(uml_results['sens_range'][1],4) - ], - 'sensitive_energy_range' :[ - round(float('{:.2e}'.format(uml_results['energy_range'][0]))), - round(float('{:.2e}'.format(uml_results['energy_range'][1]))) - ], - } - - elif llama_results_finished: - collected_results['pval_generic'] = 'null' - collected_results['pval_bayesian'] = round(llama_results['p_value'],4) - if collected_results['pval_bayesian']<0.01: - send_notif=True - - if collected_results['pval_bayesian']<0.1: - llama_ontime = format_ontime_events_llama(llama_results['single_neutrino']) - coinc_events=[] - for eventid in llama_ontime.keys(): - if (llama_ontime[eventid]['event_pval_bayesian'] < 0.1): - llama_ontime[eventid]['event_pval_generic'] = 'null' - coinc_events.append(llama_ontime[eventid]) - collected_results['n_events_coincident'] = len(coinc_events) - collected_results['coincident_events'] = coinc_events - else: - collected_results['n_events_coincident'] = 0 - - collected_results["neutrino_flux_sensitivity_range"] = { - 'flux_sensitivity' : [ - round(llama_results['neutrino_flux_sensitivity_range']['flux_sensitivity'][0],4), - round(llama_results['neutrino_flux_sensitivity_range']['flux_sensitivity'][1],4) - ], - 'sensitive_energy_range' :[ - round(float('{:.2e}'.format(llama_results['neutrino_flux_sensitivity_range'] - ['sensitive_energy_range'][0]))), - round(float('{:.2e}'.format(llama_results['neutrino_flux_sensitivity_range'] - ['sensitive_energy_range'][1]))) - ], - } - - if (collected_results['n_events_coincident'] == 0) and ('coincident_events' in collected_results.keys()): - c = collected_results.pop('coincident_events') - if ('most_likely_direction' in collected_results.keys()): - if (collected_results['n_events_coincident'] == 0): - c = collected_results.pop('most_likely_direction') - if ('most_likely_direction' in collected_results.keys()): - try: - if (collected_results['pval_generic']>0.1): - c = collected_results.pop('most_likely_direction') - except Exception as e: - print(e) - - ### SAVE RESULTS ### - if record.attrib['role']=='observation' and not heartbeat: - with open(os.path.join(save_location, f'{name}_collected_results.json'),'w') as f: - json.dump(collected_results, f, indent = 6) - - logger.info('sending notice') - status = SendAlert(results = collected_results) - st = 'sent' if status == 0 else 'error!' - logger.info('status: {}'.format(status)) - logger.info('{}'.format(st)) - - if status ==0: - with open('/cvmfs/icecube.opensciencegrid.org/users/jthwaites/tokens/gw_token.txt') as f: - my_key = f.readline() - - if not subthreshold: - channels = ['#gwnu-heartbeat', '#gwnu','#alerts'] - else: - channels = ['#gwnu-heartbeat'] - for channel in channels: - with open(os.path.join(save_location, f'{name}_collected_results.json'),'r') as fi: - response = requests.post('https://slack.com/api/files.upload', - timeout=60, - params={'token': my_key}, - data={'filename':'gcn.json', - 'title': f'GCN Notice for {name}', - 'channels': channel}, - files={'file': fi} - ) - if response.ok is True: - logger.info("GCN posted OK to {}".format(channel)) - else: - logger.info("Error posting skymap to {}!".format(channel)) - - collected_results['subthreshold'] = subthreshold - try: - updateGW_public(collected_results) - logger.info('Updated webpage.') - except: - logger.warning('Failed to push to public webpage.') - - if send_notif: - sender_script = os.path.join(os.environ.get('FAST_RESPONSE_SCRIPTS'), - '../slack_posters/lvk_email_sms_notif.py') - try: - subprocess.call([sender_script, '--path_to_gcn', - os.path.join(save_location, f'{name}_collected_results.json')]) - logger.info('Sent alert to ROC for p<0.01') - except: - logger.warning('Failed to send email/SMS notification.') - else: - logger.info('p>0.01: no email/sms sent') - else: - with open(os.path.join(save_location, f'mocks/{name}_collected_results.json'),'w') as f: - json.dump(collected_results, f, indent = 6) - #logger.info('sending test notice') - #status = SendTestAlert(results = collected_results) - #logger.info('status: {}'.format(status)) - - #send the notice to slack (#gw-mock-heartbeat) - with open('/cvmfs/icecube.opensciencegrid.org/users/jthwaites/tokens/gw_token.txt') as f: - my_key = f.readline() - - channel = '#gw-mock-heartbeat' - with open(os.path.join(save_location, f'mocks/{name}_collected_results.json'),'r') as fi: - response = requests.post('https://slack.com/api/files.upload', - timeout=60, - params={'token': my_key}, - data={'filename':'gcn.json', - 'title': f'GCN Notice for {name}', - 'channels': channel}, - files={'file': fi} - ) - if response.ok is True: - logger.info("GCN posted OK to {}".format(channel)) - else: - logger.info("Error posting skymap to {}!".format(channel)) - -logger = logging.getLogger() -logger.setLevel(logging.INFO) -logger.warning("combine_results starting, connecting to GCN") - -#fra_results_location = os.environ.get('FAST_RESPONSE_OUTPUT')#'/data/user/jthwaites/o4-mocks/' -llama_results_location = '/home/followup/lvk_dropbox/' -#llama_results_location = '/home/azhang/public_html/llama/json/' -#save_location = '/home/followup/lvk_followup_output/' #where to save final json -save_location = args.save_dir - -max_wait = args.max_wait -#wait_for_llama = args.wait_for_llama - -if args.run_live: - logger.info('running on live GCNs') - ''' - while True: - for message in consumer.consume(timeout=1): - value = message.value() - if '<' not in value.decode('utf-8')[0]: - #sometimes, we'll get error messages - these make the code fail. skip them - logger.warning(value.decode('utf-8')) - continue - notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) - parse_notice(notice) - logger.info('Done.') - ''' - gcn.listen(handler=parse_notice) - -else: - if args.test_path is None: - paths=glob.glob('/home/jthwaites/FastResponse/*/*xml') - path = paths[1] - #path = '/home/jthwaites/FastResponse/S230522n-preliminary.json,1' - else: - path = args.test_path - - logger.info('running on {}'.format(path)) - - with open(path, 'r') as f: - payload = f.read() - try: - record = lxml.etree.fromstring(payload) - except Exception as e: - print(e) - exit() - - parse_notice(payload, record) -logger.info("done") \ No newline at end of file diff --git a/fast_response/scripts/combine_results_kafka.py b/fast_response/scripts/combine_results_kafka.py index 38ca7f50..83144747 100644 --- a/fast_response/scripts/combine_results_kafka.py +++ b/fast_response/scripts/combine_results_kafka.py @@ -26,9 +26,13 @@ client_id = f.readline().rstrip('\n') client_secret = f.readline().rstrip('\n') +config = {'broker.address.family': 'v4', + 'log_level': 0, + 'max.poll.interval.ms': 1800000, + } consumer = Consumer(client_id=client_id, client_secret=client_secret, - config={'max.poll.interval.ms':1800000}) + config=config) # Subscribe to topics to receive alerts consumer.subscribe(['gcn.classic.voevent.LVC_PRELIMINARY', @@ -599,4 +603,4 @@ def parse_notice(record, wait_for_llama=False, heartbeat=False): exit() parse_notice(record, wait_for_llama=args.wait_for_llama, heartbeat=args.heartbeat) -logger.info("done") \ No newline at end of file +logger.info("done") diff --git a/fast_response/scripts/test_listen_kafka.py b/fast_response/scripts/test_listen_kafka.py index fd4d9b79..2decc511 100644 --- a/fast_response/scripts/test_listen_kafka.py +++ b/fast_response/scripts/test_listen_kafka.py @@ -3,37 +3,58 @@ import logging from gcn_kafka import Consumer from icecube import realtime_tools -import json -import argparse +import json, argparse -## none of these bools work yet. just listening to the real one -#parser = argparse.ArgumentParser(description='test listener for icecube kafka fra/llama results') -#parser.add_argument('--test_domain', type=bool, default=False, -# help='bool to use test.gcn.nasa.gov (default False)') -#parser.add_argument('--test_topic', type=bool, default=True, -# help='listen to gcn.notices.icecube.TEST.lvk_nu_track_search') -#args = parser.parse_args() +parser = argparse.ArgumentParser(description='test listener for icecube kafka fra/llama results') +parser.add_argument('--test_domain', action='store_true', default=False, + help='bool to use test IceCube stream (default False)') +parser.add_argument('--use_prod', action='store_true', default=False, + help='use production token rather than read-only') +parser.add_argument('--save_out', action='store_true', default=False, + help='save the packet as a test file') +parser.add_argument('--classic', action='store_true', default=False, + help='listen to the voevent streams from gcn rather than the kafka') +args = parser.parse_args() +logger = logging.getLogger() +logger.setLevel(logging.INFO) + +if args.use_prod: + token = '/home/jthwaites/private/tokens/real_icecube_kafka_prod.txt' +else: + token = '/home/jthwaites/private/tokens/kafka_token.txt' -with open('/home/jthwaites/private/tokens/kafka_token.txt') as f: +with open(token) as f: client_id = f.readline().rstrip('\n') client_secret = f.readline().rstrip('\n') -#domain = 'test.gcn.nasa.gov' -domain = 'gcn.nasa.gov' +if args.test_domain: + domain = 'test.gcn.nasa.gov' +else: + domain = 'gcn.nasa.gov' +config = {'broker.address.family': 'v4', + 'log_level': 0} consumer = Consumer(client_id=client_id, client_secret=client_secret, - domain=domain) + domain=domain, + config = config + ) -#topic = 'gcn.notices.icecube.test.lvk_nu_track_search' -topic = 'gcn.notices.icecube.lvk_nu_track_search' +# choose topics to listen to +if args.classic: + #no lvk_nu_track_search classic version + topics = ['gcn.classic.voevent.ICECUBE_ASTROTRACK_BRONZE', + 'gcn.classic.voevent.ICECUBE_ASTROTRACK_GOLD', + 'gcn.classic.voevent.ICECUBE_CASCADE'] +else: + topics =['gcn.notices.icecube.gold_bronze_track_alerts', + 'gcn.notices.icecube.test.gold_bronze_track_alerts', + 'gcn.notices.icecube.lvk_nu_track_search', + 'igwn.gwalert'] -consumer.subscribe([topic]) - -logger = logging.getLogger() -logger.setLevel(logging.INFO) -logger.warning("checking for {}, connecting to GCN".format(topic)) +consumer.subscribe(topics) +logger.warning("checking for alerts, connecting to GCN")#.format(topic)) while True: for message in consumer.consume(timeout=1): @@ -41,7 +62,25 @@ print(message.error()) continue value = message.value() + logger.warning('Found GCN on topic {}'.format(message.topic())) + + try: + if args.classic: + alert_xml = value.decode('utf-8') #.encode('ascii') + print(alert_xml) + if args.save_out: + with open('test_alert.xml', 'w') as f: + f.write(alert_xml) + else: + alert_dict = json.loads(value.decode('utf-8')) + if args.save_out: + with open('test_alert.json', 'w') as f: + json.dump(alert_dict, f) + if 'skymap' in alert_dict: + alert_dict.pop('skymap', None) + print(json.dumps(alert_dict, indent=2)) + + except Exception as e: + print('Caught exception when decoding: ', e, '\nSkipping...') + continue - alert_dict = json.loads(value.decode('utf-8')) - print(json.dumps(alert_dict, indent=2)) - diff --git a/fast_response/slack_posters/logger_util.py b/fast_response/slack_posters/logger_util.py new file mode 100644 index 00000000..efd9c1de --- /dev/null +++ b/fast_response/slack_posters/logger_util.py @@ -0,0 +1,83 @@ + +''' +Logging utilities for FRA +''' +import logging, os + +class LogFileWriter(logging.FileHandler): + ''' + Make a new logging handler that writes to a file + and flushes all messges to the file as they are written. + ''' + def emit(self, record): + super().emit(record) + self.flush() + +class FRA_Logger(object): + ''' + Define a logging object for use in FRA listeners. + ''' + + def __init__(self, file='', **fmt_kwargs): + ''' + Initialize a logger object for FRA. + + Parameters: + ----------- + file (str): + Path or file name to log to. If not passed, logs only to stout + fmt_kwargs (optional): + Args passed to self.logger_format + ''' + logger = logging.getLogger() + logger.setLevel(logging.INFO) + + self.logger_format(**fmt_kwargs) + logging.basicConfig(format=self.logger_fmt) + + if file and os.path.dirname(file): + if not os.path.exists(os.path.dirname(file)): + raise Exception('Unable to find parent directory for log file!') + + filelogger = LogFileWriter(file, mode='a+') + filelogger.setFormatter(self.logger_fmt) + + # adds the logfile as an additional logger. this will also log to stout + logger.addHandler(filelogger) + + self.logger = logger + + def logger_format(self, **fmt_kwargs): + ''' + Format the loggers used in FRA + + Optional arguments: + ------------------- + fmt (str): + Full format string to pass to logger.setFormatter. + If this is passed, other format options are ignored and only this format is used + with datefmt if asctime is included in fmt. + datefmt (str): + Format for dates to use. Default %Y/%m/%d %H:%M:%S. + use_time (bool): + include the time in the format (default True) + incl_level (bool): + include the logger level (info, warning, error) (default True) + ''' + + fmt = fmt_kwargs.pop('fmt', None) + datefmt = fmt_kwargs.pop('datefmt', '%Y/%m/%d %H:%M:%S') + if fmt is not None: + self.logger_fmt=logging.Formatter(fmt=fmt, datefmt=datefmt) + return + + fmt = '' + if fmt_kwargs.pop('use_time', True): + fmt = fmt + '[%(asctime)s] ' + if fmt_kwargs.pop('incl_level', True): + fmt = fmt + '%(levelname)s ' + if fmt: # if not empty, add tab break + fmt = fmt + '\t ' + fmt = fmt + '%(message)s' + + self.logger_fmt = logging.Formatter(fmt=fmt, datefmt=datefmt) diff --git a/requirements.txt b/requirements.txt index 5c4c9d3d..6ea76235 100644 --- a/requirements.txt +++ b/requirements.txt @@ -4,6 +4,7 @@ attrs==19.3.0 backports.functools-lru-cache==1.6.1 cycler==0.10.0 funcsigs==1.0.2 +gcn-kafka>=0.3.2 h5py==3.5.0 ipython==7.32.0 healpy==1.13.0 @@ -31,5 +32,4 @@ slack_sdk==3.35.0 subprocess32==3.5.4 zmq==0.0.0 py27hash==1.0.2 -pygcn==1.1.2 wheel==0.37.1