From 469455c0b50f3f419964c5938f814fa40cd19a0b Mon Sep 17 00:00:00 2001 From: jessiethw Date: Wed, 27 May 2026 11:26:15 -0500 Subject: [PATCH 01/27] remove outdated script --- .../scripts/combine_results_classic.py | 591 ------------------ 1 file changed, 591 deletions(-) delete mode 100644 fast_response/scripts/combine_results_classic.py 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 From 6e1a6ebe17100ed67949aad3014228f55a7a8bff Mon Sep 17 00:00:00 2001 From: jessiethw Date: Thu, 28 May 2026 04:51:51 -0500 Subject: [PATCH 02/27] initialize consumer in kafka --- fast_response/listeners/gcn_listener.py | 56 ++++++++++++++----------- 1 file changed, 31 insertions(+), 25 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 912df3e3..307e7079 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -4,15 +4,29 @@ Author: Alex Pizzuto Date: July 2020 ''' - -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) - -def process_gcn(payload, root): +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 +from gcn_kafka import Consumer + +with open('/home/jthwaites/private/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, + config={'max.poll.interval.ms':1800000}) + +consumer.subscribe(['gcn.notices.icecube.gold_bronze_track_alerts']) + +def process_gcn(record): #payload,root analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') if analysis_path is None: try: @@ -149,19 +163,6 @@ def process_gcn(payload, root): print(e) 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 @@ -180,14 +181,19 @@ def process_gcn(payload, root): args = parser.parse_args() if args.run_live: - print("Listening for GCNs . . . ") - gcn.listen(handler=process_gcn) + print("Listening for IC Alert GCNs . . . ") + while True: + for message in consumer.consume(timeout=1): + if message.error(): + print(message.error()) + continue + value = message.value() + print(value) 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() From 0001f150f5f4c56dac9ceaebb47bbb8dd1b15fbd Mon Sep 17 00:00:00 2001 From: jessiethw Date: Thu, 28 May 2026 07:39:41 -0500 Subject: [PATCH 03/27] more options for test listener --- fast_response/scripts/test_listen_kafka.py | 49 +++++++++++++++------- 1 file changed, 35 insertions(+), 14 deletions(-) diff --git a/fast_response/scripts/test_listen_kafka.py b/fast_response/scripts/test_listen_kafka.py index fd4d9b79..84652397 100644 --- a/fast_response/scripts/test_listen_kafka.py +++ b/fast_response/scripts/test_listen_kafka.py @@ -6,34 +6,51 @@ import json import 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() +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' consumer = Consumer(client_id=client_id, client_secret=client_secret, domain=domain) -#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'] -consumer.subscribe([topic]) +consumer.subscribe(topics) logger = logging.getLogger() logger.setLevel(logging.INFO) -logger.warning("checking for {}, connecting to GCN".format(topic)) +logger.warning("checking for alerts, connecting to GCN")#.format(topic)) while True: for message in consumer.consume(timeout=1): @@ -44,4 +61,8 @@ alert_dict = json.loads(value.decode('utf-8')) print(json.dumps(alert_dict, indent=2)) + + if args.save_out: + with open('test_alert.json', 'w') as f: + json.dump(alert_dict, f) From 7955299e5a32f566bb15973c775cf617c5e5ba8c Mon Sep 17 00:00:00 2001 From: jessiethw Date: Thu, 28 May 2026 08:12:08 -0500 Subject: [PATCH 04/27] modernize alert listener to classic voevent over kafka --- fast_response/listeners/gcn_listener.py | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 307e7079..c7c0da32 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -24,7 +24,11 @@ client_secret=client_secret, config={'max.poll.interval.ms':1800000}) -consumer.subscribe(['gcn.notices.icecube.gold_bronze_track_alerts']) +#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(record): #payload,root analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') @@ -44,10 +48,10 @@ def process_gcn(record): #payload,root # 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. . . ") @@ -188,7 +192,9 @@ def process_gcn(record): #payload,root print(message.error()) continue value = message.value() - print(value) + print('Found GCN on topic {}'.format(message.topic())) + notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) + parse_notice(notice) else: try: import fast_response From 5e910980697cef79e85d6111ca14d2fa2fe12e7c Mon Sep 17 00:00:00 2001 From: jessiethw Date: Thu, 28 May 2026 08:14:46 -0500 Subject: [PATCH 05/27] fix function for test xmls --- fast_response/listeners/gcn_listener.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index c7c0da32..d12da9cd 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -209,10 +209,10 @@ def process_gcn(record): #payload,root payload = open(sample_skymap_path \ + 'sample_astrotrack_alert_2021.xml', 'rb').read() root = lxml.etree.fromstring(payload) - process_gcn(payload, root) + process_gcn(root) else: print("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) From f7bfaf2806eb46728ca606a6aa2e8e70c8061db4 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Mon, 15 Jun 2026 12:35:32 -0500 Subject: [PATCH 06/27] fix consumer issues, use logger package --- fast_response/listeners/gcn_listener.py | 71 +++++++++++++++---------- 1 file changed, 44 insertions(+), 27 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index d12da9cd..08df4d59 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -1,9 +1,14 @@ +#!/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 + Updated Date: June 2026 ''' + +import logging +from gcn_kafka import Consumer import os, subprocess, time, pwd, argparse import healpy as hp import numpy as np @@ -14,15 +19,22 @@ from glob import glob from fast_response.slack_posters.slack import slackbot import pandas as pd -from gcn_kafka import Consumer + +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' + consumer = Consumer(client_id=client_id, client_secret=client_secret, - config={'max.poll.interval.ms':1800000}) + domain='gcn.nasa.gov', + #config={'max.poll.interval.ms':1800000}, + ) #consumer.subscribe(['gcn.notices.icecube.gold_bronze_track_alerts']) # stick with voevent for now for all 3 @@ -37,13 +49,13 @@ def process_gcn(record): #payload,root import fast_response analysis_path = 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('###########################################################################') - exit() + raise Exception(e) # Read all of the VOEvent parameters from the "What" section. params = {elem.attrib['name']: @@ -53,8 +65,8 @@ def process_gcn(record): #payload,root stream = params['Stream'] 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]) @@ -67,8 +79,8 @@ def process_gcn(record): #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'] @@ -88,7 +100,7 @@ def process_gcn(record): #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.) @@ -100,12 +112,12 @@ def process_gcn(record): #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 @@ -116,13 +128,17 @@ def process_gcn(record): #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)) + # for now, since we are still modernizing... just dump params and return + print(params) + return + subprocess.call([command, '--skymap={}'.format(skymap), '--time={}'.format(str(event_mjd)), '--alert_id={}'.format(run_id+':'+event_id), @@ -141,7 +157,7 @@ def process_gcn(record): #payload,root subprocess.call([analysis_path+'document.py', '--path', dir_2d[0]]) doc=True except: - print('Failed to document to private webpage') + logger.warning('Failed to document to private webpage') try: shifters = pd.read_csv(os.path.join(analysis_path,'../slack_posters/fra_shifters.csv'), @@ -164,11 +180,13 @@ def process_gcn(record): #payload,root bot.post_short_msg(done_message) except Exception as e: - print(e) + logger.warning('Failed to push to private webpage') + logger.warning(e) if __name__ == '__main__': 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 @@ -185,14 +203,14 @@ def process_gcn(record): #payload,root args = parser.parse_args() if args.run_live: - print("Listening for IC Alert GCNs . . . ") + logger.warning("Listening for IC Alert GCNs . . . ") while True: for message in consumer.consume(timeout=1): if message.error(): - print(message.error()) + logger.warning(message.error()) continue value = message.value() - print('Found GCN on topic {}'.format(message.topic())) + logger.warning('Found GCN on topic {}'.format(message.topic())) notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) parse_notice(notice) else: @@ -200,18 +218,17 @@ def process_gcn(record): #payload,root import fast_response sample_skymap_path=os.path.dirname(fast_response.__file__) +'/sample_skymaps/' except Exception as e: - print(e) - print('Cannot find path to sample skymaps') - exit() + logger.error('Cannot find path to sample skymaps') + raise Exception(e) if not args.test_cascade: - print("Running on sample track . . . ") + 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(root) else: - print("Running on sample cascade . . . ") + logger.info("Running on sample cascade . . . ") payload = open(sample_skymap_path \ + 'sample_cascade.txt', 'rb').read() root = lxml.etree.fromstring(payload) From 8cba5653e4540e4ec39d13e6be0313689558236f Mon Sep 17 00:00:00 2001 From: jessiethw Date: Mon, 15 Jun 2026 17:36:59 -0500 Subject: [PATCH 07/27] switch to classic over kafka --- fast_response/listeners/gw_gcn_listener.py | 95 ++++++++++++---------- 1 file changed, 52 insertions(+), 43 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 8692431b..6e8e74c0 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -3,22 +3,46 @@ ''' 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 + Updated Date: June 2026 ''' -import gcn -import sys -import pickle +#import gcn +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 lxml.etree +import argparse, time, wget +from astropy.time import Time +from datetime import datetime +from fast_response.slack_posters.slack import slackbot + +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') + +consumer = Consumer(client_id=client_id, + client_secret=client_secret, + domain='gcn.nasa.gov', + #config={'max.poll.interval.ms':1800000}, + ) + +consumer.subscribe(['gcn.classic.voevent.LVC_EARLY_WARNING', + 'gcn.classic.voevent.LVC_INITIAL', + 'gcn.classic.voevent.LVC_PRELIMINARY', + 'gcn.classic.voevent.LVC_RETRACTION', + #'gcn.classic.voevent.LVC_TEST', + 'gcn.classic.voevent.LVC_UPDATE']) + +def process_gcn(record): #payload, root): AlertTime=datetime.utcnow().isoformat() log_file.flush() @@ -40,15 +64,15 @@ def process_gcn(payload, root): # 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] + for elem in record.iterfind('.//Param')} + name = record.attrib['ivorn'].split('#')[1] # 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' + record.attrib['role']='test' log_file.flush() #return else: @@ -56,7 +80,7 @@ def process_gcn(payload, root): print('No significance parameter found in LVK GCN.') log_file.flush() # 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 record.attrib['role']!='observation': return print('\n' +'INCOMING ALERT FOUND: ',datetime.utcnow()) @@ -76,7 +100,7 @@ def process_gcn(payload, root): print('Could not determine type of event') merger_type = None - if root.attrib['role']=='observation' and not mock: + if record.attrib['role']=='observation' and not mock: ## Call everyone because it's a real event! call_command=['/home/jthwaites/private/make_call.py', f'--name={name}'] @@ -95,13 +119,13 @@ def process_gcn(payload, root): 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': + if mock and record.attrib['role']=='observation': print('Listener in heartbeat mode found real event. Returning...') log_file.flush() return # Read trigger time of event - eventtime = root.find('.//ISOTime').text + eventtime = record.find('.//ISOTime').text event_mjd = Time(eventtime, format='isot').mjd print(f'Alert MJD: {event_mjd}') print('GW merger time: %s \n' % Time(eventtime, format='isot').iso) @@ -164,7 +188,7 @@ def process_gcn(payload, root): log_file.flush() return - if root.attrib['role'] != 'observation': + if record.attrib['role'] != 'observation': name=name+'_test' print('Running on scrambled data') log_file.flush() @@ -184,7 +208,7 @@ def process_gcn(payload, root): 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 record.attrib['role'] == 'observation': try: subprocess.call([webpage_update, '--gw', f'--path={output}']) @@ -214,7 +238,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 record.attrib['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)): @@ -227,11 +251,11 @@ def process_gcn(payload, root): pickle.dump(gw_latency, file, protocol=pickle.HIGHEST_PROTOCOL) #save xml and skymap, for later - et = lxml.etree.ElementTree(root) + et = lxml.etree.ElementTree(record) et.write(os.path.join(output, '{}-{}-{}.xml'.format(params['GraceID'], params['Pkt_Ser_Num'], params['AlertType'])), pretty_print=True) - if root.attrib['role'] != 'observation': + if record.attrib['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 @@ -240,28 +264,13 @@ def process_gcn(payload, root): log_file.flush() 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, + parser.add_argument('--log_path', default='/home/jthwaites/public_html/FastResponse/', type=str, help='Redirect output to a log file with this path') parser.add_argument('--test_path', default='S191216ap_update.xml', type=str, help='Skymap for use in testing listener') @@ -305,12 +314,12 @@ def process_gcn(payload, root): #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) + record = lxml.etree.fromstring(payload) mock=args.heartbeat #test runs on scrambles, observation runs on unblinded data if not args.test_o3: - root.attrib['role']='test' + record.attrib['role']='test' mock=True - process_gcn(payload, root) + process_gcn(record) From 2b33bbee738ca46a1532660e16739ce6268dad1d Mon Sep 17 00:00:00 2001 From: jessiethw Date: Mon, 15 Jun 2026 17:51:37 -0500 Subject: [PATCH 08/27] use logging --- fast_response/listeners/gw_gcn_listener.py | 95 +++++++++------------- 1 file changed, 37 insertions(+), 58 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 6e8e74c0..8100894e 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -21,9 +21,7 @@ from datetime import datetime from fast_response.slack_posters.slack import slackbot -logger = logging.getLogger() -logger.setLevel(logging.INFO) -logger.warning("Connecting to GCN as Consumer") +print("Connecting to GCN as Consumer") with open('/home/jthwaites/private/tokens/kafka_token.txt') as f: client_id = f.readline().rstrip('\n') @@ -43,23 +41,21 @@ 'gcn.classic.voevent.LVC_UPDATE']) def process_gcn(record): #payload, root): - 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() + raise Exception(e) # Read all of the VOEvent parameters from the "What" section. params = {elem.attrib['name']: @@ -71,20 +67,16 @@ def process_gcn(record): #payload, root): if 'Significant' in params.keys(): if int(params['Significant'])==0: #not significant, do not run - print(f'Found a subthreshold event {name}') + logger.warning(f'Found a subthreshold event {name}') record.attrib['role']='test' - log_file.flush() - #return 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 record.attrib['role']!='observation': return - print('\n' +'INCOMING ALERT FOUND: ',datetime.utcnow()) - log_file.flush() + logger.warning('\n' +'INCOMING ALERT FOUND: ',datetime.utcnow()) #get type of event (burst, bbh, nsbh, bns) try: @@ -97,7 +89,7 @@ def process_gcn(record): #payload, root): probs = {j: float(params[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 record.attrib['role']=='observation' and not mock: @@ -112,24 +104,20 @@ def process_gcn(record): #payload, root): try: subprocess.call(call_command) - #print('Call here.') except Exception as e: - print('Call failed.') - print(e) - log_file.flush() + logger.error('Call failed!') + logger.error(e) # want heartbeat listener not to run on real events, otherwise it overwrites the main listener output if mock and record.attrib['role']=='observation': - print('Listener in heartbeat mode found real event. Returning...') - log_file.flush() + logger.info('Listener in heartbeat mode found real event. Skipping...') return # Read trigger time of event eventtime = record.find('.//ISOTime').text 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: %s \n' % Time(eventtime, format='isot').iso) current_mjd = Time(datetime.utcnow(), scale='utc').mjd needed_delay = 1000./84600./2. @@ -140,10 +128,9 @@ def process_gcn(record): #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 @@ -167,8 +154,7 @@ def process_gcn(record): #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] @@ -178,24 +164,22 @@ def process_gcn(record): #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 record.attrib['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)) + #### FOR NOW: testing + return subprocess.call([command, '--skymap={}'.format(skymap), '--time={}'.format(str(event_mjd)), @@ -221,9 +205,8 @@ def process_gcn(record): #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 @@ -260,8 +243,7 @@ def process_gcn(record): #payload, root): 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: ',output) if __name__ == '__main__': @@ -271,7 +253,7 @@ def process_gcn(record): #payload, root): 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='/home/jthwaites/public_html/FastResponse/', type=str, - help='Redirect output to a log file with this path') + help='Redirect 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', @@ -283,25 +265,21 @@ def process_gcn(record): #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 = logging.getLogger() + logging.basicConfig(filename=logfile, level=logging.INFO) + logger.warning("Listening for GCNs . . . ") mock=args.heartbeat - print('Starting heartbeat listener') if mock else print('Running on REAL events only') - log_file.flush() + logger.info('Starting heartbeat listener') if mock else logger.info('Running on REAL events only') gcn.listen(handler=process_gcn) - else: - print("Offline testing . . . ") - log_file.flush() + else: + logger = logging.getLogger() + logger.setLevel(logging.INFO) + logger.warning("Offline testing . . . ") ### FOR OFFLINE TESTING try: @@ -309,7 +287,8 @@ def process_gcn(record): #payload, root): #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) + logger.error('Failed to find sample skymap paths!') + logger.error(e) sample_skymap_path='/data/user/jthwaites/o3-gw-skymaps/' #payload = open(os.path.join(sample_skymap_path,args.test_path), 'rb').read() From 99693fb1e90260dbd63de06fef4c421e48382501 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Mon, 15 Jun 2026 18:08:43 -0500 Subject: [PATCH 09/27] oops add correct parsing for kafka --- fast_response/listeners/gw_gcn_listener.py | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 8100894e..589f3137 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -267,14 +267,24 @@ def process_gcn(record): #payload, root): if args.run_live: print(f'Logging to file: {logfile}') + #logging.basicConfig(filename=logfile, level=logging.INFO, filemode='a+') + logger = logging.getLogger() - logging.basicConfig(filename=logfile, level=logging.INFO) + logger.setLevel(logging.INFO) logger.warning("Listening for GCNs . . . ") mock=args.heartbeat logger.info('Starting heartbeat listener') if mock else logger.info('Running on REAL events only') - gcn.listen(handler=process_gcn) + while True: + for message in consumer.consume(timeout=1): + if message.error(): + logger.warning(message.error()) + continue + value = message.value() + logger.warning('Found GCN on topic {}'.format(message.topic())) + notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) + parse_notice(notice) else: logger = logging.getLogger() From be8c62a55a5d4c4d9288d042ca52de13ea259dc9 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Tue, 16 Jun 2026 17:45:21 -0500 Subject: [PATCH 10/27] add capability to nicely log to file --- fast_response/listeners/gw_gcn_listener.py | 38 +++++++++++++++------- 1 file changed, 26 insertions(+), 12 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 589f3137..7b1e5fc5 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -40,6 +40,15 @@ #'gcn.classic.voevent.LVC_TEST', 'gcn.classic.voevent.LVC_UPDATE']) +# make a new logging.FileHandler that can flush as we go +class LogFileWriter(logging.FileHandler): + '''Make a new logging handler that uses a file + and flushes all messges to the file as they are written + ''' + def emit(self, record): + super().emit(record) + self.flush() + def process_gcn(record): #payload, root): AlertTime=datetime.utcnow().isoformat() analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') @@ -267,30 +276,35 @@ def process_gcn(record): #payload, root): if args.run_live: print(f'Logging to file: {logfile}') - #logging.basicConfig(filename=logfile, level=logging.INFO, filemode='a+') - + logger = logging.getLogger() logger.setLevel(logging.INFO) + # adds the logfile as an additional logger. this will also log to stout + logger.addHandler(LogFileWriter(logfile, mode='a+')) logger.warning("Listening for GCNs . . . ") mock=args.heartbeat logger.info('Starting heartbeat listener') if mock else logger.info('Running on REAL events only') - while True: - for message in consumer.consume(timeout=1): - if message.error(): - logger.warning(message.error()) - continue - value = message.value() - logger.warning('Found GCN on topic {}'.format(message.topic())) - notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) - parse_notice(notice) + try: + while True: + for message in consumer.consume(timeout=1): + if message.error(): + logger.warning(message.error()) + continue + value = message.value() + logger.warning('Found GCN on topic {}'.format(message.topic())) + notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) + parse_notice(notice) + except KeyboardInterrupt: + # make sure the logfile gets shutdown correctly and file closed + logger.shutdown() else: logger = logging.getLogger() logger.setLevel(logging.INFO) logger.warning("Offline testing . . . ") - + ### FOR OFFLINE TESTING try: import fast_response From 7d980732bda0cd4c04f43edcce9843ba7d99e129 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Wed, 17 Jun 2026 11:17:41 -0500 Subject: [PATCH 11/27] fix formatter for logging --- fast_response/listeners/gw_gcn_listener.py | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 7b1e5fc5..58ac7729 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -36,8 +36,7 @@ consumer.subscribe(['gcn.classic.voevent.LVC_EARLY_WARNING', 'gcn.classic.voevent.LVC_INITIAL', 'gcn.classic.voevent.LVC_PRELIMINARY', - 'gcn.classic.voevent.LVC_RETRACTION', - #'gcn.classic.voevent.LVC_TEST', + #'gcn.classic.voevent.LVC_RETRACTION', 'gcn.classic.voevent.LVC_UPDATE']) # make a new logging.FileHandler that can flush as we go @@ -279,8 +278,10 @@ def process_gcn(record): #payload, root): logger = logging.getLogger() logger.setLevel(logging.INFO) + filelogger = LogFileWriter(logfile, mode='a+') + filelogger.setFormatter(logging.Formatter(fmt='[%(asctime)s] %(levelname)s %(message)s', datefmt='%Y/%m/%d %H:%M:%S')) # adds the logfile as an additional logger. this will also log to stout - logger.addHandler(LogFileWriter(logfile, mode='a+')) + logger.addHandler(filelogger) logger.warning("Listening for GCNs . . . ") mock=args.heartbeat @@ -295,10 +296,10 @@ def process_gcn(record): #payload, root): value = message.value() logger.warning('Found GCN on topic {}'.format(message.topic())) notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) - parse_notice(notice) + process_gcn(notice) except KeyboardInterrupt: # make sure the logfile gets shutdown correctly and file closed - logger.shutdown() + logging.shutdown() else: logger = logging.getLogger() From de16d84c07190d4ea1b141ee0e3163d39c846b3c Mon Sep 17 00:00:00 2001 From: jessiethw Date: Thu, 18 Jun 2026 18:10:03 -0500 Subject: [PATCH 12/27] move logger utils to their own class --- fast_response/comm_utils/logger_util.py | 82 +++++++++++++++++++++++++ 1 file changed, 82 insertions(+) create mode 100644 fast_response/comm_utils/logger_util.py diff --git a/fast_response/comm_utils/logger_util.py b/fast_response/comm_utils/logger_util.py new file mode 100644 index 00000000..520733ee --- /dev/null +++ b/fast_response/comm_utils/logger_util.py @@ -0,0 +1,82 @@ + +''' +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) + logger.setFormatter(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) From 0e6335fb7384d41805035dcd5a17d576e6baf524 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Thu, 18 Jun 2026 18:11:51 -0500 Subject: [PATCH 13/27] fix lxml issue with encoding --- fast_response/listeners/gcn_listener.py | 6 +++--- fast_response/listeners/gw_gcn_listener.py | 15 +++++++++------ 2 files changed, 12 insertions(+), 9 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 08df4d59..d0928007 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -209,10 +209,10 @@ def process_gcn(record): #payload,root if message.error(): logger.warning(message.error()) continue - value = message.value() + value = message.value().decode('utf-8') + value = value.replace("","") #lxml doesn't like this line logger.warning('Found GCN on topic {}'.format(message.topic())) - notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) - parse_notice(notice) + notice = lxml.etree.fromstring(value) else: try: import fast_response diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 58ac7729..3160deb3 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -36,8 +36,10 @@ consumer.subscribe(['gcn.classic.voevent.LVC_EARLY_WARNING', 'gcn.classic.voevent.LVC_INITIAL', 'gcn.classic.voevent.LVC_PRELIMINARY', - #'gcn.classic.voevent.LVC_RETRACTION', - 'gcn.classic.voevent.LVC_UPDATE']) + 'gcn.classic.voevent.LVC_RETRACTION', + 'gcn.classic.voevent.LVC_UPDATE', + 'igwn.gwalert' + ]) # make a new logging.FileHandler that can flush as we go class LogFileWriter(logging.FileHandler): @@ -279,7 +281,7 @@ def process_gcn(record): #payload, root): logger = logging.getLogger() logger.setLevel(logging.INFO) filelogger = LogFileWriter(logfile, mode='a+') - filelogger.setFormatter(logging.Formatter(fmt='[%(asctime)s] %(levelname)s %(message)s', datefmt='%Y/%m/%d %H:%M:%S')) + filelogger.setFormatter(logging.Formatter(fmt='[%(asctime)s] %(levelname)s\t %(message)s', datefmt='%Y/%m/%d %H:%M:%S')) # adds the logfile as an additional logger. this will also log to stout logger.addHandler(filelogger) logger.warning("Listening for GCNs . . . ") @@ -293,10 +295,11 @@ def process_gcn(record): #payload, root): if message.error(): logger.warning(message.error()) continue - value = message.value() + value = message.value().decode('utf-8') + value = value.replace("","") #lxml doesn't like this line logger.warning('Found GCN on topic {}'.format(message.topic())) - notice = lxml.etree.fromstring(value.decode('utf-8').encode('ascii')) - process_gcn(notice) + notice = lxml.etree.fromstring(value) + #process_gcn(notice) except KeyboardInterrupt: # make sure the logfile gets shutdown correctly and file closed logging.shutdown() From 56d0f4c8bfdbbab99cf944bb1124501b4f543168 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Fri, 10 Jul 2026 10:44:29 -0500 Subject: [PATCH 14/27] fix lxml issues --- fast_response/listeners/gw_gcn_listener.py | 2 +- fast_response/scripts/test_listen_kafka.py | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index 3160deb3..b247e1ca 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -296,7 +296,7 @@ def process_gcn(record): #payload, root): logger.warning(message.error()) continue value = message.value().decode('utf-8') - value = value.replace("","") #lxml doesn't like this line + value = value.replace("encoding='UTF-8'","") #lxml doesn't like this line logger.warning('Found GCN on topic {}'.format(message.topic())) notice = lxml.etree.fromstring(value) #process_gcn(notice) diff --git a/fast_response/scripts/test_listen_kafka.py b/fast_response/scripts/test_listen_kafka.py index 84652397..7308b433 100644 --- a/fast_response/scripts/test_listen_kafka.py +++ b/fast_response/scripts/test_listen_kafka.py @@ -58,6 +58,7 @@ print(message.error()) continue value = message.value() + logger.warning('Found GCN on topic {}'.format(message.topic())) alert_dict = json.loads(value.decode('utf-8')) print(json.dumps(alert_dict, indent=2)) From 52c0f28cb9c3038a776ffacde8b63776a07cd086 Mon Sep 17 00:00:00 2001 From: Alicia Mand Date: Fri, 10 Jul 2026 14:54:34 -0500 Subject: [PATCH 15/27] bugfix for file parsing --- fast_response/listeners/gcn_listener.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index d0928007..4c76e593 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -209,8 +209,9 @@ def process_gcn(record): #payload,root 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().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) else: From 7eb337d3a07daddecad4b21cbb9490b5bb4e5453 Mon Sep 17 00:00:00 2001 From: Alicia Mand Date: Tue, 11 Aug 2026 12:48:16 -0500 Subject: [PATCH 16/27] Added error handling and posting error --- fast_response/listeners/gcn_listener.py | 37 +++++++++++++++++++++---- 1 file changed, 32 insertions(+), 5 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 4c76e593..025e6ad0 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -3,8 +3,8 @@ ''' Script to automatically receive GCN notices for IceCube alert events and run followup accordingly - Author: Alex Pizzuto, Jessie Thwaites - Updated Date: June 2026 + Author: Alex Pizzuto, Jessie Thwaites, Alicia Mand + Updated Date: August 2026 ''' import logging @@ -183,6 +183,23 @@ def process_gcn(record): #payload,root logger.warning('Failed to push to private webpage') logger.warning(e) +def post_error(): + analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') + bot = slackbot('fra-shifting') + message = "Error processing alert information. Please run FRA manually in 24 hours." + 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] Date: Tue, 25 Aug 2026 10:30:36 -0500 Subject: [PATCH 17/27] prevent ipv6 errors (realtime uses ipv4) --- fast_response/listeners/gcn_listener.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 025e6ad0..4a7c01f6 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -29,10 +29,11 @@ client_secret = f.readline().rstrip('\n') domain = 'gcn.nasa.gov' - +config = {'broker.address.family': 'v4'} consumer = Consumer(client_id=client_id, client_secret=client_secret, domain='gcn.nasa.gov', + config=config, #config={'max.poll.interval.ms':1800000}, ) From cdeccd61cf2d2442ee93dc9e6eb044f19945d7e3 Mon Sep 17 00:00:00 2001 From: Alicia Mand Date: Tue, 25 Aug 2026 10:43:00 -0500 Subject: [PATCH 18/27] enable running followups --- fast_response/listeners/gcn_listener.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 4a7c01f6..1e496f7b 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -136,9 +136,8 @@ def process_gcn(record): #payload,root logger.info('\nRunning {} --skymap={} --time={} --alert_id={} --suffix={}'.format( command, skymap, str(event_mjd), run_id+':'+event_id, suffix)) - # for now, since we are still modernizing... just dump params and return + print(params) - return subprocess.call([command, '--skymap={}'.format(skymap), '--time={}'.format(str(event_mjd)), From cbe5a2cea7f281216d48fead6507ada23d7ac0e0 Mon Sep 17 00:00:00 2001 From: Alicia Mand Date: Tue, 25 Aug 2026 10:55:31 -0500 Subject: [PATCH 19/27] slightly more sophisticated error msgs --- fast_response/listeners/gcn_listener.py | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index 1e496f7b..ffe73082 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -51,6 +51,7 @@ def process_gcn(record): #payload,root analysis_path = os.path.dirname(fast_response.__file__) + '/scripts/' except Exception as 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 ') @@ -157,6 +158,7 @@ def process_gcn(record): #payload,root subprocess.call([analysis_path+'document.py', '--path', dir_2d[0]]) doc=True except: + post_error("Failed to run document command") logger.warning('Failed to document to private webpage') try: @@ -180,18 +182,22 @@ def process_gcn(record): #payload,root bot.post_short_msg(done_message) except Exception as e: + post_error("Failed to push results to private webpage") logger.warning('Failed to push to private webpage') logger.warning(e) -def post_error(): +def post_error(errMsg=None): analysis_path = os.environ.get('FAST_RESPONSE_SCRIPTS') bot = slackbot('fra-shifting') - message = "Error processing alert information. Please run FRA manually in 24 hours." + 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] Date: Fri, 28 Aug 2026 14:08:25 -0500 Subject: [PATCH 20/27] reduce loglevel to prevent flooded output --- fast_response/listeners/gcn_listener.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/fast_response/listeners/gcn_listener.py b/fast_response/listeners/gcn_listener.py index ffe73082..c79daffb 100644 --- a/fast_response/listeners/gcn_listener.py +++ b/fast_response/listeners/gcn_listener.py @@ -29,7 +29,8 @@ client_secret = f.readline().rstrip('\n') domain = 'gcn.nasa.gov' -config = {'broker.address.family': 'v4'} +config = {'broker.address.family': 'v4', + 'log_level': 0} consumer = Consumer(client_id=client_id, client_secret=client_secret, domain='gcn.nasa.gov', From 2b39dd1df0f88a9b37f4c4ee2fbb6db7d0674dcd Mon Sep 17 00:00:00 2001 From: jessiethw Date: Sun, 20 Sep 2026 12:27:18 -0500 Subject: [PATCH 21/27] move logger to slack poster directory --- fast_response/{comm_utils => slack_posters}/logger_util.py | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename fast_response/{comm_utils => slack_posters}/logger_util.py (100%) diff --git a/fast_response/comm_utils/logger_util.py b/fast_response/slack_posters/logger_util.py similarity index 100% rename from fast_response/comm_utils/logger_util.py rename to fast_response/slack_posters/logger_util.py From 0ba3ac3ef1fb50776e73a75e57ce108a3b73800e Mon Sep 17 00:00:00 2001 From: jessiethw Date: Sun, 20 Sep 2026 12:27:59 -0500 Subject: [PATCH 22/27] switch to json params in listener --- fast_response/listeners/gw_gcn_listener.py | 143 +++++++++------------ 1 file changed, 58 insertions(+), 85 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index b247e1ca..e69aa325 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -4,10 +4,9 @@ to run realtime neutrino follow-up Author: Raamis Hussain, Jessie Thwaites, MJ Romfoe - Updated Date: June 2026 + Last Updated: Sept 2026 ''' -#import gcn import logging from gcn_kafka import Consumer import sys, pickle, os, subprocess, pwd @@ -15,11 +14,13 @@ from dateutil.relativedelta import relativedelta import healpy as hp import numpy as np -import lxml.etree 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") @@ -30,33 +31,16 @@ consumer = Consumer(client_id=client_id, client_secret=client_secret, domain='gcn.nasa.gov', - #config={'max.poll.interval.ms':1800000}, ) -consumer.subscribe(['gcn.classic.voevent.LVC_EARLY_WARNING', - 'gcn.classic.voevent.LVC_INITIAL', - 'gcn.classic.voevent.LVC_PRELIMINARY', - 'gcn.classic.voevent.LVC_RETRACTION', - 'gcn.classic.voevent.LVC_UPDATE', - 'igwn.gwalert' - ]) - -# make a new logging.FileHandler that can flush as we go -class LogFileWriter(logging.FileHandler): - '''Make a new logging handler that uses a file - and flushes all messges to the file as they are written - ''' - def emit(self, record): - super().emit(record) - self.flush() - -def process_gcn(record): #payload, root): +consumer.subscribe(['igwn.gwalert']) + +def process_gcn(params, mock=False): AlertTime=datetime.utcnow().isoformat() 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: logger.error('Error finding FRA package!!') @@ -67,42 +51,43 @@ def process_gcn(record): #payload, root): print('###########################################################################') raise Exception(e) - # Read all of the VOEvent parameters from the "What" section. - params = {elem.attrib['name']: - elem.attrib['value'] - for elem in record.iterfind('.//Param')} - name = record.attrib['ivorn'].split('#')[1] + name = params['superevent_id'] + '-' + params['alert_type'].lower() + params['role'] = 'observation' if params['event']['search'] is not 'MDC' else 'test' # only run on significant events - if 'Significant' in params.keys(): - if int(params['Significant'])==0: - #not significant, do not run + 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}') - record.attrib['role']='test' + params['role']='test' else: # O3 does not have this parameter, this should only happen for testing 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 record.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 logger.warning('\n' +'INCOMING ALERT FOUND: ',datetime.utcnow()) #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: logger.warning('Could not determine type of event') merger_type = None - if record.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}'] @@ -117,17 +102,12 @@ def process_gcn(record): #payload, root): except Exception as e: logger.error('Call failed!') logger.error(e) - - # want heartbeat listener not to run on real events, otherwise it overwrites the main listener output - if mock and record.attrib['role']=='observation': - logger.info('Listener in heartbeat mode found real event. Skipping...') - return # Read trigger time of event - eventtime = record.find('.//ISOTime').text + eventtime = params['event']['time'] event_mjd = Time(eventtime, format='isot').mjd logger.info(f'Alert MJD: {event_mjd}') - logger.info('GW merger time: %s \n' % Time(eventtime, format='isot').iso) + 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. @@ -145,7 +125,8 @@ def process_gcn(record): #payload, root): 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 @@ -182,7 +163,7 @@ def process_gcn(record): #payload, root): f'args: --time {event_mjd} --name {name} --skymap PATH_TO_SKYMAP') return - if record.attrib['role'] != 'observation': + if params['role'] != 'observation': name=name+'_test' logger.info('Running on scrambled data') command = os.path.join(analysis_path, 'run_gw_followup.py') @@ -191,18 +172,19 @@ def process_gcn(record): #payload, root): #### FOR NOW: testing return - 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 record.attrib['role'] == 'observation': + if not mock and params['role'] == 'observation': try: subprocess.call([webpage_update, '--gw', f'--path={output}']) @@ -231,7 +213,7 @@ def process_gcn(record): #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 record.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)): @@ -243,12 +225,11 @@ def process_gcn(record): #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(record) - et.write(os.path.join(output, '{}-{}-{}.xml'.format(params['GraceID'], - params['Pkt_Ser_Num'], params['AlertType'])), pretty_print=True) - - if record.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 @@ -263,7 +244,7 @@ def process_gcn(record): #payload, root): 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='/home/jthwaites/public_html/FastResponse/', type=str, - help='Redirect output to a log file with this path. Note: this is only used when running live') + 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', @@ -278,12 +259,7 @@ def process_gcn(record): #payload, root): if args.run_live: print(f'Logging to file: {logfile}') - logger = logging.getLogger() - logger.setLevel(logging.INFO) - filelogger = LogFileWriter(logfile, mode='a+') - filelogger.setFormatter(logging.Formatter(fmt='[%(asctime)s] %(levelname)s\t %(message)s', datefmt='%Y/%m/%d %H:%M:%S')) - # adds the logfile as an additional logger. this will also log to stout - logger.addHandler(filelogger) + logger = FRA_Logger(file=logfile) logger.warning("Listening for GCNs . . . ") mock=args.heartbeat @@ -296,37 +272,34 @@ def process_gcn(record): #payload, root): logger.warning(message.error()) continue value = message.value().decode('utf-8') - value = value.replace("encoding='UTF-8'","") #lxml doesn't like this line logger.warning('Found GCN on topic {}'.format(message.topic())) - notice = lxml.etree.fromstring(value) - #process_gcn(notice) + notice = json.loads(value) + process_gcn(notice,mock=mock) except KeyboardInterrupt: # make sure the logfile gets shutdown correctly and file closed logging.shutdown() else: - logger = logging.getLogger() - logger.setLevel(logging.INFO) + logger = FRA_Logger() logger.warning("Offline testing . . . ") - ### 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: - logger.error('Failed to find sample skymap paths!') - logger.error(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() - record = lxml.etree.fromstring(payload) + # 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: - record.attrib['role']='test' + params['event']['search'] = 'MDC' mock=True - process_gcn(record) + process_gcn(params, mock=mock) From d82dbe4dc7bfdd4441d34807ea3190518718be62 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Sun, 20 Sep 2026 12:34:44 -0500 Subject: [PATCH 23/27] add loggers to docs --- doc/source/slack.rst | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) 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: From f0e22442479ef04d7ed2cfff7591a280859101d0 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Mon, 28 Sep 2026 17:26:12 -0500 Subject: [PATCH 24/27] handle retractions better --- fast_response/listeners/gw_gcn_listener.py | 28 +++++++++++++++------- 1 file changed, 20 insertions(+), 8 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index e69aa325..a59c246f 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -28,9 +28,15 @@ 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']) @@ -52,8 +58,14 @@ def process_gcn(params, mock=False): raise Exception(e) name = params['superevent_id'] + '-' + params['alert_type'].lower() - params['role'] = 'observation' if params['event']['search'] is not 'MDC' else 'test' - + if params['alert_type'].lower() == 'retraction': + print('Error! Listener does not run on Retractions. Returning...') + 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['event']: if not params['event']['significant']: @@ -71,8 +83,7 @@ def process_gcn(params, mock=False): logger.info('Listener in heartbeat mode found real event. Skipping...') return - logger.warning('\n' +'INCOMING ALERT FOUND: ',datetime.utcnow()) - + 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['event']['group'] == 'Burst': @@ -169,8 +180,6 @@ def process_gcn(params, mock=False): command = os.path.join(analysis_path, 'run_gw_followup.py') logger.info('Running {}'.format(command)) - #### FOR NOW: testing - return subprocess.call([ command, @@ -259,7 +268,7 @@ def process_gcn(params, mock=False): if args.run_live: print(f'Logging to file: {logfile}') - logger = FRA_Logger(file=logfile) + logger = FRA_Logger(file=logfile).logger logger.warning("Listening for GCNs . . . ") mock=args.heartbeat @@ -274,6 +283,9 @@ def process_gcn(params, mock=False): 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 @@ -294,7 +306,7 @@ def process_gcn(params, mock=False): sys.exit() logger.info('Running offline on file: {}'.format(test_file)) - params = json.loads(test_file)) + params = json.loads(test_file) mock=args.heartbeat #test runs on scrambles, observation runs on unblinded data From 92613d3c0a678e496d0987230c0ddc13506c4943 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Mon, 28 Sep 2026 17:26:56 -0500 Subject: [PATCH 25/27] fix error messages --- .../scripts/combine_results_kafka.py | 8 +++- fast_response/scripts/test_listen_kafka.py | 45 +++++++++++++------ fast_response/slack_posters/logger_util.py | 3 +- 3 files changed, 39 insertions(+), 17 deletions(-) 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 7308b433..2decc511 100644 --- a/fast_response/scripts/test_listen_kafka.py +++ b/fast_response/scripts/test_listen_kafka.py @@ -3,8 +3,7 @@ import logging from gcn_kafka import Consumer from icecube import realtime_tools -import json -import argparse +import json, argparse parser = argparse.ArgumentParser(description='test listener for icecube kafka fra/llama results') parser.add_argument('--test_domain', action='store_true', default=False, @@ -17,6 +16,9 @@ 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: @@ -31,9 +33,13 @@ 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 + ) # choose topics to listen to if args.classic: @@ -44,12 +50,10 @@ else: topics =['gcn.notices.icecube.gold_bronze_track_alerts', 'gcn.notices.icecube.test.gold_bronze_track_alerts', - 'gcn.notices.icecube.lvk_nu_track_search'] + 'gcn.notices.icecube.lvk_nu_track_search', + 'igwn.gwalert'] consumer.subscribe(topics) - -logger = logging.getLogger() -logger.setLevel(logging.INFO) logger.warning("checking for alerts, connecting to GCN")#.format(topic)) while True: @@ -59,11 +63,24 @@ continue value = message.value() logger.warning('Found GCN on topic {}'.format(message.topic())) - - alert_dict = json.loads(value.decode('utf-8')) - print(json.dumps(alert_dict, indent=2)) - if args.save_out: - with open('test_alert.json', 'w') as f: - json.dump(alert_dict, f) - + 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 + diff --git a/fast_response/slack_posters/logger_util.py b/fast_response/slack_posters/logger_util.py index 520733ee..efd9c1de 100644 --- a/fast_response/slack_posters/logger_util.py +++ b/fast_response/slack_posters/logger_util.py @@ -33,7 +33,8 @@ def __init__(self, file='', **fmt_kwargs): logger.setLevel(logging.INFO) self.logger_format(**fmt_kwargs) - logger.setFormatter(self.logger_fmt) + 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!') From cf02cce4a3f432db88fe6bed6d8121c95216d957 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Wed, 30 Sep 2026 12:52:20 -0500 Subject: [PATCH 26/27] add time of lvk notice sent to name for unique id --- fast_response/listeners/gw_gcn_listener.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/fast_response/listeners/gw_gcn_listener.py b/fast_response/listeners/gw_gcn_listener.py index a59c246f..652e0ce6 100755 --- a/fast_response/listeners/gw_gcn_listener.py +++ b/fast_response/listeners/gw_gcn_listener.py @@ -57,9 +57,9 @@ def process_gcn(params, mock=False): print('###########################################################################') raise Exception(e) - name = params['superevent_id'] + '-' + params['alert_type'].lower() + name = params['superevent_id'] + '-'+ +params['time_created'].replace(':','').replace('-','') + '-' + params['alert_type'].lower() if params['alert_type'].lower() == 'retraction': - print('Error! Listener does not run on Retractions. Returning...') + 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 @@ -243,7 +243,7 @@ def process_gcn(params, mock=False): subprocess.call(['mv',output, '/data/user/jthwaites/o4-mocks/']) output = '/data/user/jthwaites/o4-mocks/' + eventtime[0:10].replace('-','_')+'_'+name - logger.info('Output directory: ',output) + logger.info('Output directory: {}'.format(output)) if __name__ == '__main__': From 157d22788f0870774e7482125635cdfaab617380 Mon Sep 17 00:00:00 2001 From: jessiethw Date: Wed, 30 Sep 2026 12:58:13 -0500 Subject: [PATCH 27/27] correct requirements for kafka --- requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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