diff --git a/sf_daq_broker/broker_manager.py b/sf_daq_broker/broker_manager.py index 897eefc..b665493 100644 --- a/sf_daq_broker/broker_manager.py +++ b/sf_daq_broker/broker_manager.py @@ -190,6 +190,10 @@ class BrokerManager(object): request.get("channels_list"), config.OUTPUT_FILE_SUFFIX_DATA_BUFFER) + send_write_request(broker_config.TAG_DATA3BUFFER, + request.get("channels_list"), + config.OUTPUT_FILE_SUFFIX_DATA3_BUFFER) + send_write_request(broker_config.TAG_IMAGEBUFFER, request.get("camera_list"), config.OUTPUT_FILE_SUFFIX_IMAGE_BUFFER) diff --git a/sf_daq_broker/config.py b/sf_daq_broker/config.py index 02c87c3..25d7115 100644 --- a/sf_daq_broker/config.py +++ b/sf_daq_broker/config.py @@ -10,9 +10,10 @@ AUDIT_FILE_TIME_FORMAT = "%Y%m%d-%H%M%S" DEFAULT_AUDIT_FILENAME = "/var/log/sf_databuffer_audit.log" -DATA_API_QUERY_ADDRESS = "http://sf-data-api-02.psi.ch/query" +#DATA_API_QUERY_ADDRESS = "http://sf-data-api-02.psi.ch/query" +DATA_API_QUERY_ADDRESS = "https://data-api.psi.ch/sf/query" IMAGE_API_QUERY_ADDRESS = ["http://172.27.0.14:8080/api/v1/query", "http://172.27.0.15:8080/api/v1/query"] -DATA_API3_QUERY_ADDRESS = "http://sf-daqbuf-21:8371/api/1.0.1/query" +DATA_API3_QUERY_ADDRESS = "http://sf-daqbuf-33:8371/api/1" EPICS_QUERY_ADDRESS = "https://data-api.psi.ch/sf" DATA_BACKEND = "sf-databuffer" @@ -23,5 +24,6 @@ TRANSFORM_PULSE_ID_TO_TIMESTAMP_QUERY = False SEPARATE_CAMERA_CHANNELS = True OUTPUT_FILE_SUFFIX_DATA_BUFFER = "BSREAD" +OUTPUT_FILE_SUFFIX_DATA3_BUFFER = "BSDATA" OUTPUT_FILE_SUFFIX_IMAGE_BUFFER = "CAMERAS" OUTPUT_FILE_SUFFIX_EPICS = "PVCHANNELS" diff --git a/sf_daq_broker/rabbitmq/config.py b/sf_daq_broker/rabbitmq/config.py index 03c7935..ac53520 100644 --- a/sf_daq_broker/rabbitmq/config.py +++ b/sf_daq_broker/rabbitmq/config.py @@ -11,4 +11,5 @@ DEFAULT_QUEUE = "write_request" # Name of the queue for bsread write requests. TAG_DATABUFFER = "databuffer" TAG_IMAGEBUFFER = "imagebuffer" -TAG_EPICS = "epics" \ No newline at end of file +TAG_EPICS = "epics" +TAG_DATA3BUFFER = "data3buffer" diff --git a/sf_daq_broker/utils.py b/sf_daq_broker/utils.py index 5f08acf..0b49db7 100644 --- a/sf_daq_broker/utils.py +++ b/sf_daq_broker/utils.py @@ -6,7 +6,7 @@ import requests from sf_daq_broker import config -_logger = getLogger(__name__) +_logger = getLogger("broker_writer") def get_data_api_request(channels, start_pulse_id, stop_pulse_id): @@ -50,7 +50,7 @@ def transform_range_from_pulse_id_to_timestamp(data_api_request): mapping_request = {'range': {'startPulseId': data_api_request["range"]["startPulseId"], 'endPulseId': data_api_request["range"]["endPulseId"]+1}} - mapping_response = requests.post(url=config.DATA_API_QUERY_ADDRESS + "/mapping", json=mapping_request).json() + mapping_response = requests.post(url=config.DATA_API_QUERY_ADDRESS + "/mapping", json=mapping_request, timeout=10).json() _logger.info("Response to mapping request: %s", mapping_response) @@ -63,6 +63,7 @@ def transform_range_from_pulse_id_to_timestamp(data_api_request): _logger.info("Transformed request to startSeconds and endSeconds. %s" % new_data_api_request) except Exception as e: + _logger.error(e) raise RuntimeError("Cannot retrieve the pulse_id to timestamp mapping.") from e return new_data_api_request diff --git a/sf_daq_broker/writer/bsread_writer.py b/sf_daq_broker/writer/bsread_writer.py index 6c4b291..05e6203 100644 --- a/sf_daq_broker/writer/bsread_writer.py +++ b/sf_daq_broker/writer/bsread_writer.py @@ -7,6 +7,7 @@ import h5py import numpy import requests from random import randrange +from copy import deepcopy from sf_daq_broker import config, utils from sf_daq_broker.writer.utils import channel_type_deserializer_mapping @@ -29,7 +30,11 @@ def write_from_databuffer(data_api_request, output_file, metadata): start_time = time() - response = requests.post(url=config.DATA_API_QUERY_ADDRESS, json=data_api_request) + new_data_api_request = deepcopy(data_api_request) +# new_data_api_request["range"]["startPulseId"] -= 1 +# new_data_api_request["range"]["endPulseId"] += 1 + + response = requests.post(url=config.DATA_API_QUERY_ADDRESS, json=new_data_api_request, timeout=1000) data = json.loads(response.content) if not data: @@ -91,6 +96,45 @@ def write_from_imagebuffer(data_api_request, output_file, parameters): _logger.error(e) +def write_from_databuffer_api3(data_api_request, output_file, parameters): + import data_api3.h5 as h5 + import pytz + + _logger.debug("Data3 API request: %s", data_api_request) + + data_api_request_timestamp = utils.transform_range_from_pulse_id_to_timestamp(data_api_request) + + channels = [channel["name"] for channel in data_api_request_timestamp["channels"]] + + start = datetime.fromtimestamp(float(data_api_request_timestamp["range"]["startSeconds"])).astimezone( + pytz.timezone('UTC')).strftime("%Y-%m-%dT%H:%M:%S.%fZ") # isoformat() # "2019-12-13T09:00:00.00 + end = datetime.fromtimestamp(float(data_api_request_timestamp["range"]["endSeconds"])).astimezone( + pytz.timezone('UTC')).strftime("%Y-%m-%dT%H:%M:%S.%fZ") # isoformat() # "2019-12-13T09:00:00.00 + + query = { + "channels": channels, + "range": { + "type": "date", + "startDate": start, + "endDate": end + } + } + + data_buffer_url = config.DATA_API3_QUERY_ADDRESS + + _logger.debug("Requesting '%s' to output_file %s from %s " % + (query, output_file, data_buffer_url)) + + start_time = time() + + try: + h5.request(query, filename=output_file, baseurl=data_buffer_url, default_backend=config.DATA_BACKEND) + _logger.info("Data download and writing took %s seconds." % (time() - start_time)) + except Exception as e: + _logger.error("Got exception from data_api3") + _logger.error(e) + + class BsreadH5Writer(object): def __init__(self, output_file, metadata): diff --git a/sf_daq_broker/writer/start.py b/sf_daq_broker/writer/start.py index f8fccae..b850fe9 100644 --- a/sf_daq_broker/writer/start.py +++ b/sf_daq_broker/writer/start.py @@ -11,7 +11,7 @@ from pika import BlockingConnection, ConnectionParameters, BasicProperties from sf_daq_broker import config, utils import sf_daq_broker.rabbitmq.config as broker_config from sf_daq_broker.utils import get_data_api_request -from sf_daq_broker.writer.bsread_writer import write_from_imagebuffer, write_from_databuffer +from sf_daq_broker.writer.bsread_writer import write_from_imagebuffer, write_from_databuffer, write_from_databuffer_api3 from sf_daq_broker.writer.epics_writer import write_epics_pvs _logger = logging.getLogger("broker_writer") @@ -70,6 +70,8 @@ def process_request(request): if writer_type == broker_config.TAG_DATABUFFER: logger_data_api = logging.getLogger("data_api") + elif writer_type == broker_config.TAG_DATA3BUFFER: + logger_data_api = logging.getLogger("data_api3") elif writer_type == broker_config.TAG_IMAGEBUFFER: logger_data_api = logging.getLogger("data_api3") elif writer_type == broker_config.TAG_EPICS: @@ -99,6 +101,10 @@ def process_request(request): _logger.info("Using databuffer writer.") write_from_databuffer(get_data_api_request(channels, start_pulse_id, stop_pulse_id), output_file, metadata) + elif writer_type == broker_config.TAG_DATA3BUFFER: + _logger.info("Using data_api3 databuffer writer.") + write_from_databuffer_api3(get_data_api_request(channels, start_pulse_id, stop_pulse_id), output_file, metadata) + elif writer_type == broker_config.TAG_IMAGEBUFFER: _logger.info("Using imagebuffer writer.") write_from_imagebuffer(get_data_api_request(channels, start_pulse_id, stop_pulse_id), output_file, metadata)