+# -*- coding: utf-8 -*-
+
+# Copyright 2018 Telefonica S.A.
+#
+# Licensed under the Apache License, Version 2.0 (the "License");
+# you may not use this file except in compliance with the License.
+# You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
+# implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+
import logging
import os
import yaml
import asyncio
from osm_common.msgbase import MsgBase, MsgException
from time import sleep
+from http import HTTPStatus
__author__ = "Alfonso Tierno <alfonso.tiernosepulveda@telefonica.com>"
self.logger = logging.getLogger(logger_name)
self.path = None
# create a different file for each topic
- self.files = {}
+ self.files_read = {}
+ self.files_write = {}
self.buffer = {}
def connect(self, config):
except MsgException:
raise
except Exception as e: # TODO refine
- raise MsgException(str(e))
+ raise MsgException(str(e), http_code=HTTPStatus.INTERNAL_SERVER_ERROR)
def disconnect(self):
- for f in self.files.values():
+ for f in self.files_read.values():
+ try:
+ f.close()
+ except Exception: # TODO refine
+ pass
+ for f in self.files_write.values():
try:
f.close()
except Exception: # TODO refine
:return: None or raises and exception
"""
try:
- if topic not in self.files:
- self.files[topic] = open(self.path + topic, "a+")
- yaml.safe_dump({key: msg}, self.files[topic], default_flow_style=True, width=20000)
- self.files[topic].flush()
+ if topic not in self.files_write:
+ self.files_write[topic] = open(self.path + topic, "a+")
+ yaml.safe_dump({key: msg}, self.files_write[topic], default_flow_style=True, width=20000)
+ self.files_write[topic].flush()
except Exception as e: # TODO refine
- raise MsgException(str(e))
+ raise MsgException(str(e), HTTPStatus.INTERNAL_SERVER_ERROR)
def read(self, topic, blocks=True):
"""
topic_list = (topic, )
while True:
for single_topic in topic_list:
- if single_topic not in self.files:
- self.files[single_topic] = open(self.path + single_topic, "a+")
+ if single_topic not in self.files_read:
+ self.files_read[single_topic] = open(self.path + single_topic, "a+")
self.buffer[single_topic] = ""
- self.buffer[single_topic] += self.files[single_topic].readline()
+ self.buffer[single_topic] += self.files_read[single_topic].readline()
if not self.buffer[single_topic].endswith("\n"):
continue
msg_dict = yaml.load(self.buffer[single_topic])
return None
sleep(2)
except Exception as e: # TODO refine
- raise MsgException(str(e))
+ raise MsgException(str(e), HTTPStatus.INTERNAL_SERVER_ERROR)
async def aioread(self, topic, loop):
"""
except MsgException:
raise
except Exception as e: # TODO refine
- raise MsgException(str(e))
+ raise MsgException(str(e), HTTPStatus.INTERNAL_SERVER_ERROR)
+
+ async def aiowrite(self, topic, key, msg, loop=None):
+ """
+ Asyncio write. It blocks
+ :param topic: str
+ :param key: str
+ :param msg: message, can be str or yaml
+ :param loop: asyncio loop
+ :return: nothing if ok or raises an exception
+ """
+ return self.write(topic, key, msg)