Skip to content
Open
Show file tree
Hide file tree
Changes from 7 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .idea/.gitignore

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

9 changes: 9 additions & 0 deletions lib_python3/origin/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
#

from origin.origin_current_time import current_time
from origin.origin_data_types import data_types
from origin.origin_registration_validation import registration_validation

TIMESTAMP = "measurement_time"

from origin.client import *
20 changes: 20 additions & 0 deletions lib_python3/origin/client/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
class float_field:
pass


class integer_field:
pass


class string_field:
pass


class file_field:
pass


from origin.client import origin_server_connection
ServerConnection = origin_server_connection.ServerConnection
from origin.client import origin_server
Server = origin_server.Server
146 changes: 146 additions & 0 deletions lib_python3/origin/client/origin_server.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
from origin.client import ServerConnection

from origin.client import float_field
from origin.client import integer_field
from origin.client import string_field
from origin import data_types, registration_validation


import zmq
import struct
import json


def decode(measurement_type: str) -> str:
"""
Decodes a measurement type of it happens to be a specific class
TODO is this necessary? Those classes are empty.
Args:
measurement_type : string describing the data type a measurement should be. If not in
data_types, checks and translates for a 'float_field', 'integer_field' or
'integer_field' class

Returns:
The data type description string. If measurement_type is a valid key for data_types,
the argument is returned, other wise the argument is translated (if valid) or a KeyError
is raised

Raises:
KeyError if measurement_type is not a valid data type
"""
try:
tmp = data_types[measurement_type]
except KeyError:
if measurement_type == float_field:
return "float"
if measurement_type == integer_field:
return "int"
if measurement_type == string_field:
return "string"
raise
else:
return measurement_type


def declaration_formatter(stream, records, key_order)-> bytes:
dec_str = [stream]
for key in key_order:
dec_str.append(':'.join([key, records[key]]))
return bytes(','.join(dec_str), encoding="utf8")


def format_stream_declaration(stream, records, key_order, encoding_format):
measurements = records.keys()
sent_dict = {}
for m in measurements:
try:
decoded_type = decode(records[m])
except KeyError as e:
print(f"{records[m]} is not a valid data type. Programming error. Error should be "
f"caught before this")
return None
else:
sent_dict[m] = decoded_type
if (encoding_format is not None) and (encoding_format.lower() == "json"):
if key_order is not None:
msg = (
"Warning: JSON formatting selected and a key order has been "
"defined. JSON object order is not gaurenteed therefore it is "
"not recommended to use binary data packets."
)
print(msg)
# make deterministic
return json.dumps((stream, sent_dict), sort_keys=True)
else:
return declaration_formatter(stream, sent_dict, key_order)


class Server:
def __init__(self, config):
self.config = config

def ping(self):
return True

def register_stream(self, stream, records, key_order=None, data_format=None, timeout=1000):
valid = registration_validation(stream, records, key_order)

if not valid:
print("invalid stream declaration")
return None

port = self.config.get('Server', "register_port")
msgport = self.config.get('Server', "measure_port")
if data_format == 'json':
port = self.config.get('Server', "json_register_port")
msgport = self.config.get('Server', "json_measure_port")

context = zmq.Context()
socket = context.socket(zmq.REQ)
socket.setsockopt(zmq.RCVTIMEO, timeout)
socket.setsockopt(zmq.LINGER, 0)
host = self.config.get('Server', "ip")
socket.connect("tcp://%s:%s" % (host, port))

if (key_order is None) and (data_format is None):
key_order = records.keys()
registerComm = format_stream_declaration(stream, records, key_order, data_format)

if registerComm is None:
if data_format is None:
data_format = "comma-separated values"
print("can't format stream into {}".format(data_format))
return None

socket.send(registerComm, zmq.NOBLOCK)
try:
confirmation = socket.recv()
except Exception as e:
print(f"Problem registering stream: {stream}\nError = {e}")
print("Server did not respond in time") # TODO how do we know ??
return None
return_code, msg = bytes(confirmation).split(b',', 1)
print(return_code, msg)

if int(return_code) != 0:
print(f"Problem registering stream {stream}")
print(msg)
return None

streamID, version = struct.unpack("!II", msg)
print(f"successfully registered with streamID: {streamID}, version: {version}")
socket.close() # I 'm pretty sure we need this here

# error checking
socket_data = context.socket(zmq.PUSH)
socket_data.connect("tcp://%s:%s" % (host, msgport))
return ServerConnection(
self.config,
stream,
streamID,
key_order,
data_format,
records,
context,
socket_data
)
101 changes: 101 additions & 0 deletions lib_python3/origin/client/origin_server_connection.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
import json
import sys
import zmq
import struct
import ctypes
import traceback

from origin import data_types, TIMESTAMP


# returns string and size tuple
def make_format_string(config, key_order, records):
try:
ts_type = config.get('Server',"timestamp_type")
except KeyError:
ts_type = "uint"
ts_size = data_types[ts_type]["size"]
fstr = "!I" + data_types[ts_type]["format_char"]# network byte order

data_length = ts_size
for entry in key_order:
data_length += data_types[records[entry]]["size"]
fstr += data_types[records[entry]]["format_char"]
return fstr, data_length


class ServerConnection:
def __init__(self, config, stream, streamID, key_order, data_format, records, context, socket):
self.config = config
self.stream = stream
self.streamID = streamID
self.key_order = key_order
try:
self.data_format = data_format.lower()
except AttributeError:
self.data_format = None
self.records = records
self.context = context
self.socket = socket

self.socket.setsockopt(zmq.SNDTIMEO,2000)

if key_order is None:
self.format_string, self.data_size = (None, None)
else:
self.format_string, self.data_size = make_format_string(self.config, key_order, records)

def send(self,**kwargs):
msg_data = [ self.streamID ]
try:
msg_data.append(kwargs[TIMESTAMP])
except KeyError:
#print "No timestamp specified, server will timestamp on arrival"
msg_data.append(0)

if self.data_format is None:
for k in self.key_order:
msg_data.append(kwargs[k])
try:
self.socket.send(self.format_record(msg_data), zmq.NOBLOCK)
except zmq.Again:
print("Connection to Server Failed")
#self.socket.close()
exit(1)
except Exception as e:
print(f"Uncaught exception: {e}")
print('-'*60)
traceback.print_exc(file=sys.stdout)
print('-'*60)
print("Exiting")
#self.socket.close()
exit(1)

elif self.data_format == "json":
msg_data[0] = self.stream
msg_map = {}
for k in kwargs.keys():
if k != TIMESTAMP:
msg_map[k] = kwargs[k]
msg_data.append(msg_map)
try:
self.socket.send(json.dumps(msg_data), zmq.NOBLOCK)
except zmq.Again:
print("Connection to Server Failed")
#self.socket.close()
exit(1)
except Exception as e:
print(f"Uncaught exception: {e}")
print('-' * 60)
traceback.print_exc(file=sys.stdout)
print('-' * 60)
print("Exiting")
# self.socket.close()
exit(1)

def close(self):
print("closing socket")
self.socket.close()

def format_record(self, data):
return struct.pack( self.format_string, *data )
18 changes: 18 additions & 0 deletions lib_python3/origin/origin_current_time.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
import calendar
import time
from configparser import ConfigParser


def current_time(config: ConfigParser):
"""
Figures out the current time in the format that origin wants
Args:
config : config object
Returns:
time in format desired by origin
"""

#Unix time (in UTC)
if config.get('Server', "timestamp_type") == "uint64":
return int(time.time()*2**32)
return calendar.timegm(time.gmtime())
Loading