diff --git a/sf_daq_broker/writer/epics_writer.py b/sf_daq_broker/writer/epics_writer.py index 7714c3a..1bf2196 100644 --- a/sf_daq_broker/writer/epics_writer.py +++ b/sf_daq_broker/writer/epics_writer.py @@ -5,15 +5,16 @@ import logging import dateutil import pytz import requests -import data_api -logger = logging.getLogger("broker_writer") +_logger = logging.getLogger("broker_writer") DATA_API_QUERY_URL = "https://data-api.psi.ch/sf/query" def write_epics_pvs(output_file, start_pulse_id, stop_pulse_id, metadata, epics_pvs): + import data_api + start_date = get_pulse_id_date_mapping(start_pulse_id) stop_date = get_pulse_id_date_mapping(stop_pulse_id) @@ -21,15 +22,18 @@ def write_epics_pvs(output_file, start_pulse_id, stop_pulse_id, metadata, epics_ # TODO: Merge metadata to data. if len(data) > 0: - logger.info("Persist data to hdf5 file") + _logger.info("Persist data to hdf5 file") data_api.to_hdf5(data, output_file, overwrite=True, compression=None, shuffle=False) else: - logger.error("No data retrieved") + _logger.error("No data retrieved") open(output_file + "_NO_DATA", 'a').close() def get_data(channel_list, start=None, stop=None, base_url=None): - logger.info("Requesting range %s to %s for channels: %s" % (start, stop, channel_list)) + + import data_api + + _logger.info("Requesting range %s to %s for channels: %s" % (start, stop, channel_list)) query = {"range": {"startDate": datetime.datetime.isoformat(start), "endDate": datetime.datetime.isoformat(stop), @@ -37,13 +41,13 @@ def get_data(channel_list, start=None, stop=None, base_url=None): "channels": channel_list, "fields": ["pulseId", "globalSeconds", "globalDate", "value", "eventCount"]} - logger.debug(query) + _logger.debug(query) response = requests.post(DATA_API_QUERY_URL, json=query) # Check for successful return of data if response.status_code != 200: - logger.info("Data retrievali failed, sleep for another time and try") + _logger.info("Data retrievali failed, sleep for another time and try") itry = 0 while itry < 5: @@ -53,12 +57,12 @@ def get_data(channel_list, start=None, stop=None, base_url=None): if response.status_code == 200: break - logger.info("Data retrieval failed, post attempt %d" % itry) + _logger.info("Data retrieval failed, post attempt %d" % itry) if response.status_code != 200: raise RuntimeError("Unable to retrieve data from server: ", response) - logger.info("Data retieval is successful") + _logger.info("Data retieval is successful") data = response.json() @@ -67,7 +71,7 @@ def get_data(channel_list, start=None, stop=None, base_url=None): def get_pulse_id_date_mapping(pulse_id): # See https://jira.psi.ch/browse/ATEST-897 for more details ... - logger.info("Retrieve pulse-id/date mapping for pulse_id %s" % pulse_id) + _logger.info("Retrieve pulse-id/date mapping for pulse_id %s" % pulse_id) try: @@ -93,7 +97,7 @@ def get_pulse_id_date_mapping(pulse_id): "Didn't get good responce from data_api : %s " % data) if not pulse_id == data[0]["data"][0]["pulseId"]: - logger.info("retrieval failed") + _logger.info("retrieval failed") if c == 0: ref_date = data[0]["data"][0]["globalDate"] ref_date = dateutil.parser.parse(ref_date) @@ -107,7 +111,7 @@ def get_pulse_id_date_mapping(pulse_id): delta_date = check_date - now_date s = delta_date.seconds - logger.info("retry in " + str(s) + " seconds ") + _logger.info("retry in " + str(s) + " seconds ") if not s <= 0: time.sleep(s) continue diff --git a/sf_daq_broker/writer/start.py b/sf_daq_broker/writer/start.py index 2d6fe9b..89d77ee 100644 --- a/sf_daq_broker/writer/start.py +++ b/sf_daq_broker/writer/start.py @@ -68,6 +68,9 @@ def process_request(request): file_handler.setLevel(logging.INFO) _logger.addHandler(file_handler) + logger_data_api = logging.getLogger("DataApiClient") + logger_data_api.addHandler(file_handler) + try: _logger.info("Request for %s to write %s from pulse_id %s to %s" % (writer_type, output_file, start_pulse_id, stop_pulse_id)) @@ -109,6 +112,7 @@ def process_request(request): finally: if file_handler: _logger.removeHandler(file_handler) + logger_data_api.removeHandler(file_handler) def update_status(channel, body, action, file, message=None): @@ -215,7 +219,15 @@ def run(): writer_id_format = '{broker_writer_%s}' % args.writer_id logs_format = '[%(levelname)s] %(message)s' - logging.basicConfig(level=args.log_level, format=writer_id_format + logs_format) + #logging.basicConfig(level=args.log_level, format=f'{writer_id_format} {logs_format}') + _logger.setLevel(args.log_level) + stream_handler = logging.StreamHandler() + stream_handler.setLevel(args.log_level) + + formatter = logging.Formatter(writer_id_format + logs_format) + stream_handler.setFormatter(formatter) + + _logger.addHandler(stream_handler) logging.getLogger("pika").setLevel(logging.WARNING)