#!/usr/bin/python3 import json import sys from urllib import request, parse import yaml def _get_unicode(data, force=False): ''' This function is taken as-is from https://github.com/AxiomExergy/influx-client/blob/master/influx/line_protocol.py Licensed by Axiom Exergy under the terms of the Apache License 2.0 ''' """Try to return a text aka unicode object from the given data.""" if isinstance(data, bytes): return data.decode('utf-8') elif data is None: return '' elif force: return str(data) else: return data def _escape_tag(tag): ''' This function is taken as-is from https://github.com/AxiomExergy/influx-client/blob/master/influx/line_protocol.py Licensed by Axiom Exergy under the terms of the Apache License 2.0 ''' tag = _get_unicode(tag, force=True) return tag.replace( "\\", "\\\\" ).replace( " ", "\\ " ).replace( ",", "\\," ).replace( "=", "\\=" ) def _escape_tag_value(value): ''' This function is taken as-is from https://github.com/AxiomExergy/influx-client/blob/master/influx/line_protocol.py Licensed by Axiom Exergy under the terms of the Apache License 2.0 ''' ret = _escape_tag(value) if ret.endswith('\\'): ret += ' ' return ret def run_query(query, config): encoded = parse.quote(query['prom_query']) url = config['prometheus']['url'] + '/api/v1/query?query=' + encoded res = request.urlopen(url) response = json.loads(res.read().decode()) if 'status' not in response or 'data' not in response: raise ValueError() if response['status'] != 'success': raise ValueError() if response['data']['resultType'] != 'vector': raise ValueError() lines = [] for result in response['data']['result']: time = int(result['value'][0]) value = float(result['value'][1]) line = _escape_tag(query['influx_series']) for lkey, lvalue in query['influx_labels'].items(): lkey = _escape_tag(lkey) lvalue = _escape_tag_value(lvalue) line += f',{lkey}={lvalue}' line += f' value={_escape_tag_value(value)} {time}' lines.append(line) ib = config['influxdb']['url'] + '/write' user = parse.quote(config['influxdb']['user']) pw = parse.quote(config['influxdb']['pass']) db = parse.quote(query['influx_db']) rp = parse.quote(query['influx_retention']) iq = f'u={user}&p={pw}&db={db}&rp={rp}&precision=s' url = f'{ib}?{iq}' body = '\n'.join(lines).encode() request.urlopen(url, body) if __name__ == '__main__': with open(sys.argv[1], 'r') as c: config = yaml.safe_load(c) for query in config['queries']: run_query(query, config)