mirror of
https://github.com/DJ2LS/FreeDATA
synced 2024-05-14 08:04:33 +00:00
adjusted some db related things
This commit is contained in:
parent
beec229360
commit
2685c7d5d5
|
@ -3,12 +3,12 @@
|
||||||
import structlog
|
import structlog
|
||||||
import lzma
|
import lzma
|
||||||
import gzip
|
import gzip
|
||||||
from message_p2p import MessageP2P
|
from message_p2p import message_received
|
||||||
from message_system_db_manager import DatabaseManager
|
|
||||||
|
|
||||||
class ARQDataTypeHandler:
|
class ARQDataTypeHandler:
|
||||||
def __init__(self):
|
def __init__(self, event_manager):
|
||||||
self.logger = structlog.get_logger(type(self).__name__)
|
self.logger = structlog.get_logger(type(self).__name__)
|
||||||
|
self.event_manager = event_manager
|
||||||
self.handlers = {
|
self.handlers = {
|
||||||
"raw": {
|
"raw": {
|
||||||
'prepare': self.prepare_raw,
|
'prepare': self.prepare_raw,
|
||||||
|
@ -82,9 +82,5 @@ class ARQDataTypeHandler:
|
||||||
def handle_p2pmsg_lzma(self, data):
|
def handle_p2pmsg_lzma(self, data):
|
||||||
decompressed_data = lzma.decompress(data)
|
decompressed_data = lzma.decompress(data)
|
||||||
self.log(f"Handling LZMA compressed P2PMSG data: {len(decompressed_data)} Bytes from {len(data)} Bytes")
|
self.log(f"Handling LZMA compressed P2PMSG data: {len(decompressed_data)} Bytes from {len(data)} Bytes")
|
||||||
decompressed_json_string = decompressed_data.decode('utf-8')
|
message_received(self.event_manager, decompressed_data)
|
||||||
received_message_obj = MessageP2P.from_payload(decompressed_json_string)
|
|
||||||
received_message_dict = MessageP2P.to_dict(received_message_obj, received=True)
|
|
||||||
result = DatabaseManager().add_message(received_message_dict)
|
|
||||||
|
|
||||||
return decompressed_data
|
return decompressed_data
|
||||||
|
|
|
@ -46,7 +46,7 @@ class ARQSession():
|
||||||
self.frame_factory = data_frame_factory.DataFrameFactory(self.config)
|
self.frame_factory = data_frame_factory.DataFrameFactory(self.config)
|
||||||
self.event_frame_received = threading.Event()
|
self.event_frame_received = threading.Event()
|
||||||
|
|
||||||
self.arq_data_type_handler = ARQDataTypeHandler()
|
self.arq_data_type_handler = ARQDataTypeHandler(self.event_manager)
|
||||||
self.id = None
|
self.id = None
|
||||||
self.session_started = time.time()
|
self.session_started = time.time()
|
||||||
self.session_ended = 0
|
self.session_ended = 0
|
||||||
|
|
|
@ -15,7 +15,7 @@ class TxCommand():
|
||||||
self.event_manager = event_manager
|
self.event_manager = event_manager
|
||||||
self.set_params_from_api(apiParams)
|
self.set_params_from_api(apiParams)
|
||||||
self.frame_factory = DataFrameFactory(config)
|
self.frame_factory = DataFrameFactory(config)
|
||||||
self.arq_data_type_handler = ARQDataTypeHandler()
|
self.arq_data_type_handler = ARQDataTypeHandler(event_manager)
|
||||||
|
|
||||||
def set_params_from_api(self, apiParams):
|
def set_params_from_api(self, apiParams):
|
||||||
pass
|
pass
|
||||||
|
|
|
@ -5,6 +5,7 @@ from queue import Queue
|
||||||
from arq_session_iss import ARQSessionISS
|
from arq_session_iss import ARQSessionISS
|
||||||
from message_p2p import MessageP2P
|
from message_p2p import MessageP2P
|
||||||
from arq_data_type_handler import ARQDataTypeHandler
|
from arq_data_type_handler import ARQDataTypeHandler
|
||||||
|
from message_system_db_manager import DatabaseManager
|
||||||
|
|
||||||
class SendMessageCommand(TxCommand):
|
class SendMessageCommand(TxCommand):
|
||||||
"""Command to send a P2P message using an ARQ transfer session
|
"""Command to send a P2P message using an ARQ transfer session
|
||||||
|
@ -16,9 +17,15 @@ class SendMessageCommand(TxCommand):
|
||||||
|
|
||||||
def transmit(self, modem):
|
def transmit(self, modem):
|
||||||
# Convert JSON string to bytes (using UTF-8 encoding)
|
# Convert JSON string to bytes (using UTF-8 encoding)
|
||||||
|
|
||||||
|
DatabaseManager().add_message(self.message.to_dict())
|
||||||
|
|
||||||
payload = self.message.to_payload().encode('utf-8')
|
payload = self.message.to_payload().encode('utf-8')
|
||||||
json_bytearray = bytearray(payload)
|
json_bytearray = bytearray(payload)
|
||||||
data, data_type = self.arq_data_type_handler.prepare(json_bytearray, 'p2pmsg_lzma')
|
data, data_type = self.arq_data_type_handler.prepare(json_bytearray, 'p2pmsg_lzma')
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
iss = ARQSessionISS(self.config,
|
iss = ARQSessionISS(self.config,
|
||||||
modem,
|
modem,
|
||||||
self.message.destination,
|
self.message.destination,
|
||||||
|
|
|
@ -89,4 +89,7 @@ class EventManager:
|
||||||
|
|
||||||
def modem_failed(self):
|
def modem_failed(self):
|
||||||
event = {"modem": "failed"}
|
event = {"modem": "failed"}
|
||||||
self.broadcast(event)
|
self.broadcast(event)
|
||||||
|
|
||||||
|
def freedata_message_db_change(self):
|
||||||
|
self.broadcast({"message-db": "changed"})
|
|
@ -2,6 +2,14 @@ import datetime
|
||||||
import api_validations
|
import api_validations
|
||||||
import base64
|
import base64
|
||||||
import json
|
import json
|
||||||
|
from message_system_db_manager import DatabaseManager
|
||||||
|
|
||||||
|
|
||||||
|
def message_received(event_manager, data):
|
||||||
|
decompressed_json_string = data.decode('utf-8')
|
||||||
|
received_message_obj = MessageP2P.from_payload(decompressed_json_string)
|
||||||
|
received_message_dict = MessageP2P.to_dict(received_message_obj, received=True)
|
||||||
|
DatabaseManager(event_manager).add_message(received_message_dict)
|
||||||
|
|
||||||
|
|
||||||
class MessageP2P:
|
class MessageP2P:
|
||||||
|
@ -58,17 +66,12 @@ class MessageP2P:
|
||||||
"""Make a dictionary out of the message data
|
"""Make a dictionary out of the message data
|
||||||
"""
|
"""
|
||||||
|
|
||||||
if received:
|
|
||||||
direction = 'receive'
|
|
||||||
else:
|
|
||||||
direction = 'transmit'
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
'id': self.get_id(),
|
'id': self.get_id(),
|
||||||
'origin': self.origin,
|
'origin': self.origin,
|
||||||
'destination': self.destination,
|
'destination': self.destination,
|
||||||
'body': self.body,
|
'body': self.body,
|
||||||
'direction': direction,
|
'direction': 'receive' if received else 'transmit',
|
||||||
'attachments': list(map(self.__encode_attachment__, self.attachments)),
|
'attachments': list(map(self.__encode_attachment__, self.attachments)),
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -76,3 +79,4 @@ class MessageP2P:
|
||||||
"""Make a byte array ready to be sent out of the message data"""
|
"""Make a byte array ready to be sent out of the message data"""
|
||||||
json_string = json.dumps(self.to_dict())
|
json_string = json.dumps(self.to_dict())
|
||||||
return json_string
|
return json_string
|
||||||
|
|
||||||
|
|
|
@ -1,4 +1,5 @@
|
||||||
# database_manager.py
|
# database_manager.py
|
||||||
|
import sqlite3
|
||||||
|
|
||||||
from sqlalchemy import create_engine
|
from sqlalchemy import create_engine
|
||||||
from sqlalchemy.orm import scoped_session, sessionmaker
|
from sqlalchemy.orm import scoped_session, sessionmaker
|
||||||
|
@ -8,9 +9,10 @@ from datetime import datetime
|
||||||
import json
|
import json
|
||||||
import structlog
|
import structlog
|
||||||
|
|
||||||
|
|
||||||
class DatabaseManager:
|
class DatabaseManager:
|
||||||
def __init__(self, uri='sqlite:///freedata-messages.db'):
|
def __init__(self, event_manger, uri='sqlite:///freedata-messages.db'):
|
||||||
|
self.event_manager = event_manger
|
||||||
|
|
||||||
self.engine = create_engine(uri, echo=False)
|
self.engine = create_engine(uri, echo=False)
|
||||||
self.thread_local = local()
|
self.thread_local = local()
|
||||||
self.session_factory = sessionmaker(bind=self.engine)
|
self.session_factory = sessionmaker(bind=self.engine)
|
||||||
|
@ -18,6 +20,32 @@ class DatabaseManager:
|
||||||
|
|
||||||
self.logger = structlog.get_logger(type(self).__name__)
|
self.logger = structlog.get_logger(type(self).__name__)
|
||||||
|
|
||||||
|
def initialize_default_values(self):
|
||||||
|
session = self.get_thread_scoped_session()
|
||||||
|
try:
|
||||||
|
statuses = [
|
||||||
|
"transmitting",
|
||||||
|
"transmitted",
|
||||||
|
"received",
|
||||||
|
"failed",
|
||||||
|
"failed_checksum",
|
||||||
|
"aborted"
|
||||||
|
]
|
||||||
|
|
||||||
|
# Add default statuses if they don't exist
|
||||||
|
for status_name in statuses:
|
||||||
|
existing_status = session.query(Status).filter_by(name=status_name).first()
|
||||||
|
if not existing_status:
|
||||||
|
new_status = Status(name=status_name)
|
||||||
|
session.add(new_status)
|
||||||
|
|
||||||
|
session.commit()
|
||||||
|
self.log("Initialized database")
|
||||||
|
except Exception as e:
|
||||||
|
session.rollback()
|
||||||
|
self.log(f"An error occurred while initializing default values: {e}", isWarning=True)
|
||||||
|
finally:
|
||||||
|
session.remove()
|
||||||
|
|
||||||
def log(self, message, isWarning=False):
|
def log(self, message, isWarning=False):
|
||||||
msg = f"[{type(self).__name__}]: {message}"
|
msg = f"[{type(self).__name__}]: {message}"
|
||||||
|
@ -59,7 +87,6 @@ class DatabaseManager:
|
||||||
|
|
||||||
# Parse the timestamp from the message ID
|
# Parse the timestamp from the message ID
|
||||||
timestamp = datetime.fromisoformat(message_data['id'].split('_')[2])
|
timestamp = datetime.fromisoformat(message_data['id'].split('_')[2])
|
||||||
|
|
||||||
# Create the P2PMessage instance
|
# Create the P2PMessage instance
|
||||||
new_message = P2PMessage(
|
new_message = P2PMessage(
|
||||||
id=message_data['id'],
|
id=message_data['id'],
|
||||||
|
@ -67,6 +94,7 @@ class DatabaseManager:
|
||||||
destination_callsign=destination.callsign,
|
destination_callsign=destination.callsign,
|
||||||
body=message_data['body'],
|
body=message_data['body'],
|
||||||
timestamp=timestamp,
|
timestamp=timestamp,
|
||||||
|
direction=message_data['direction'],
|
||||||
status_id=status.id if status else None
|
status_id=status.id if status else None
|
||||||
)
|
)
|
||||||
|
|
||||||
|
@ -83,6 +111,7 @@ class DatabaseManager:
|
||||||
session.commit()
|
session.commit()
|
||||||
|
|
||||||
self.log(f"Added data to database: {new_message.id}")
|
self.log(f"Added data to database: {new_message.id}")
|
||||||
|
self.event_manager.freedata_message_db_change()
|
||||||
return new_message.id
|
return new_message.id
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
session.rollback()
|
session.rollback()
|
||||||
|
@ -95,12 +124,19 @@ class DatabaseManager:
|
||||||
try:
|
try:
|
||||||
messages = session.query(P2PMessage).all()
|
messages = session.query(P2PMessage).all()
|
||||||
return [message.to_dict() for message in messages]
|
return [message.to_dict() for message in messages]
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
raise e
|
self.log(f"error fetching database messages with error: {e}", isWarning=True)
|
||||||
|
self.log(f"---> please delete or update existing database", isWarning=True)
|
||||||
|
|
||||||
|
return False
|
||||||
|
|
||||||
finally:
|
finally:
|
||||||
session.remove()
|
session.remove()
|
||||||
|
|
||||||
def get_all_messages_json(self):
|
def get_all_messages_json(self):
|
||||||
messages_dict = self.get_all_messages()
|
messages_dict = self.get_all_messages()
|
||||||
messages_with_header = {'total_messages' : len(messages_dict), 'messages' : messages_dict}
|
if messages_dict:
|
||||||
return json.dumps(messages_with_header) # Convert to JSON string
|
messages_with_header = {'total_messages' : len(messages_dict), 'messages' : messages_dict}
|
||||||
|
return json.dumps(messages_with_header) # Convert to JSON string
|
||||||
|
return json.dumps({'error': 'fetching messages from database'})
|
|
@ -8,8 +8,8 @@ Base = declarative_base()
|
||||||
class Station(Base):
|
class Station(Base):
|
||||||
__tablename__ = 'station'
|
__tablename__ = 'station'
|
||||||
callsign = Column(String, primary_key=True)
|
callsign = Column(String, primary_key=True)
|
||||||
location = Column(String, nullable=True)
|
location = Column(JSON, nullable=True)
|
||||||
info = Column(String, nullable=True)
|
info = Column(JSON, nullable=True)
|
||||||
|
|
||||||
class Status(Base):
|
class Status(Base):
|
||||||
__tablename__ = 'status'
|
__tablename__ = 'status'
|
||||||
|
@ -20,15 +20,14 @@ class P2PMessage(Base):
|
||||||
__tablename__ = 'p2p_message'
|
__tablename__ = 'p2p_message'
|
||||||
id = Column(String, primary_key=True)
|
id = Column(String, primary_key=True)
|
||||||
origin_callsign = Column(String, ForeignKey('station.callsign'))
|
origin_callsign = Column(String, ForeignKey('station.callsign'))
|
||||||
via = Column(String, nullable=True)
|
via_callsign = Column(String, ForeignKey('station.callsign'), nullable=True)
|
||||||
destination_callsign = Column(String, ForeignKey('station.callsign'))
|
destination_callsign = Column(String, ForeignKey('station.callsign'))
|
||||||
body = Column(String)
|
body = Column(String, nullable=True)
|
||||||
attachments = relationship('Attachment', backref='p2p_message')
|
attachments = relationship('Attachment', backref='p2p_message')
|
||||||
timestamp = Column(DateTime)
|
timestamp = Column(DateTime)
|
||||||
timestamp_sent = Column(DateTime, nullable=True)
|
|
||||||
status_id = Column(Integer, ForeignKey('status.id'), nullable=True)
|
status_id = Column(Integer, ForeignKey('status.id'), nullable=True)
|
||||||
status = relationship('Status', backref='p2p_messages')
|
status = relationship('Status', backref='p2p_messages')
|
||||||
direction = Column(String, nullable=True)
|
direction = Column(String)
|
||||||
statistics = Column(JSON, nullable=True)
|
statistics = Column(JSON, nullable=True)
|
||||||
|
|
||||||
def to_dict(self):
|
def to_dict(self):
|
||||||
|
@ -36,12 +35,11 @@ class P2PMessage(Base):
|
||||||
'id': self.id,
|
'id': self.id,
|
||||||
'timestamp': self.timestamp.isoformat() if self.timestamp else None,
|
'timestamp': self.timestamp.isoformat() if self.timestamp else None,
|
||||||
'origin': self.origin_callsign,
|
'origin': self.origin_callsign,
|
||||||
'via': self.via,
|
'via': self.via_callsign,
|
||||||
'destination': self.destination_callsign,
|
'destination': self.destination_callsign,
|
||||||
'direction': self.direction,
|
'direction': self.direction,
|
||||||
'body': self.body,
|
'body': self.body,
|
||||||
'attachments': [attachment.to_dict() for attachment in self.attachments],
|
'attachments': [attachment.to_dict() for attachment in self.attachments],
|
||||||
'timestamp_sent': self.timestamp_sent.isoformat() if self.timestamp_sent else None,
|
|
||||||
'status': self.status.name if self.status else None,
|
'status': self.status.name if self.status else None,
|
||||||
'statistics': self.statistics
|
'statistics': self.statistics
|
||||||
}
|
}
|
||||||
|
|
|
@ -240,7 +240,7 @@ def get_post_radio():
|
||||||
@app.route('/freedata/messages', methods=['POST', 'GET'])
|
@app.route('/freedata/messages', methods=['POST', 'GET'])
|
||||||
def get_post_freedata_message():
|
def get_post_freedata_message():
|
||||||
if request.method in ['GET']:
|
if request.method in ['GET']:
|
||||||
result = DatabaseManager().get_all_messages_json()
|
result = DatabaseManager(app.event_manager).get_all_messages_json()
|
||||||
return api_response(result)
|
return api_response(result)
|
||||||
if enqueue_tx_command(command_message_send.SendMessageCommand, request.json):
|
if enqueue_tx_command(command_message_send.SendMessageCommand, request.json):
|
||||||
return api_response(request.json)
|
return api_response(request.json)
|
||||||
|
@ -296,6 +296,8 @@ if __name__ == "__main__":
|
||||||
app.service_manager = service_manager.SM(app)
|
app.service_manager = service_manager.SM(app)
|
||||||
# start modem service
|
# start modem service
|
||||||
app.modem_service.put("start")
|
app.modem_service.put("start")
|
||||||
|
# initialize databse default values
|
||||||
|
DatabaseManager(app.event_manager).initialize_default_values()
|
||||||
|
|
||||||
wsm.startThreads(app)
|
wsm.startThreads(app)
|
||||||
app.run()
|
app.run()
|
||||||
|
|
|
@ -2,16 +2,21 @@ import sys
|
||||||
sys.path.append('modem')
|
sys.path.append('modem')
|
||||||
|
|
||||||
import unittest
|
import unittest
|
||||||
|
import queue
|
||||||
from arq_data_type_handler import ARQDataTypeHandler
|
from arq_data_type_handler import ARQDataTypeHandler
|
||||||
|
from event_manager import EventManager
|
||||||
|
|
||||||
|
|
||||||
class TestDispatcher(unittest.TestCase):
|
class TestDispatcher(unittest.TestCase):
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def setUpClass(cls):
|
def setUpClass(cls):
|
||||||
cls.arq_data_type_handler = ARQDataTypeHandler()
|
cls.event_queue = queue.Queue()
|
||||||
|
cls.event_manager = EventManager([cls.event_queue])
|
||||||
|
cls.arq_data_type_handler = ARQDataTypeHandler(cls.event_manager)
|
||||||
|
|
||||||
|
|
||||||
def testDataTypeHandlerRaw(self):
|
def testDataTypeHevent_managerandlerRaw(self):
|
||||||
# Example usage
|
# Example usage
|
||||||
example_data = b"Hello FreeDATA!"
|
example_data = b"Hello FreeDATA!"
|
||||||
formatted_data, type_byte = self.arq_data_type_handler.prepare(example_data, "raw")
|
formatted_data, type_byte = self.arq_data_type_handler.prepare(example_data, "raw")
|
||||||
|
|
|
@ -6,7 +6,8 @@ import unittest
|
||||||
from config import CONFIG
|
from config import CONFIG
|
||||||
from message_p2p import MessageP2P
|
from message_p2p import MessageP2P
|
||||||
from message_system_db_manager import DatabaseManager
|
from message_system_db_manager import DatabaseManager
|
||||||
|
from event_manager import EventManager
|
||||||
|
import queue
|
||||||
|
|
||||||
class TestDataFrameFactory(unittest.TestCase):
|
class TestDataFrameFactory(unittest.TestCase):
|
||||||
|
|
||||||
|
@ -14,8 +15,11 @@ class TestDataFrameFactory(unittest.TestCase):
|
||||||
def setUpClass(cls):
|
def setUpClass(cls):
|
||||||
config_manager = CONFIG('modem/config.ini.example')
|
config_manager = CONFIG('modem/config.ini.example')
|
||||||
cls.config = config_manager.read()
|
cls.config = config_manager.read()
|
||||||
|
|
||||||
|
cls.event_queue = queue.Queue()
|
||||||
|
cls.event_manager = EventManager([cls.event_queue])
|
||||||
cls.mycall = f"{cls.config['STATION']['mycall']}-{cls.config['STATION']['myssid']}"
|
cls.mycall = f"{cls.config['STATION']['mycall']}-{cls.config['STATION']['myssid']}"
|
||||||
cls.database_manager = DatabaseManager(uri='sqlite:///:memory:')
|
cls.database_manager = DatabaseManager(cls.event_manager, uri='sqlite:///:memory:')
|
||||||
|
|
||||||
def testFromApiParams(self):
|
def testFromApiParams(self):
|
||||||
api_params = {
|
api_params = {
|
||||||
|
@ -34,7 +38,20 @@ class TestDataFrameFactory(unittest.TestCase):
|
||||||
}
|
}
|
||||||
message = MessageP2P(self.mycall, 'DJ2LS-3', 'Hello World!', [attachment])
|
message = MessageP2P(self.mycall, 'DJ2LS-3', 'Hello World!', [attachment])
|
||||||
payload = message.to_payload()
|
payload = message.to_payload()
|
||||||
print(payload)
|
received_message = MessageP2P.from_payload(payload)
|
||||||
|
self.assertEqual(message.origin, received_message.origin)
|
||||||
|
self.assertEqual(message.destination, received_message.destination)
|
||||||
|
self.assertCountEqual(message.attachments, received_message.attachments)
|
||||||
|
self.assertEqual(attachment['data'], received_message.attachments[0]['data'])
|
||||||
|
|
||||||
|
def testToPayloadWithAttachmentAndDatabase(self):
|
||||||
|
attachment = {
|
||||||
|
'name': 'test.gif',
|
||||||
|
'type': 'image/gif',
|
||||||
|
'data': np.random.bytes(1024)
|
||||||
|
}
|
||||||
|
message = MessageP2P(self.mycall, 'DJ2LS-3', 'Hello World!', [attachment])
|
||||||
|
payload = message.to_payload()
|
||||||
received_message = MessageP2P.from_payload(payload)
|
received_message = MessageP2P.from_payload(payload)
|
||||||
received_message_dict = MessageP2P.to_dict(received_message, received=True)
|
received_message_dict = MessageP2P.to_dict(received_message, received=True)
|
||||||
self.database_manager.add_message(received_message_dict)
|
self.database_manager.add_message(received_message_dict)
|
||||||
|
@ -47,5 +64,6 @@ class TestDataFrameFactory(unittest.TestCase):
|
||||||
result = self.database_manager.get_all_messages()
|
result = self.database_manager.get_all_messages()
|
||||||
self.assertEqual(result[0]["destination"], message.destination)
|
self.assertEqual(result[0]["destination"], message.destination)
|
||||||
|
|
||||||
|
|
||||||
if __name__ == '__main__':
|
if __name__ == '__main__':
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|
Loading…
Reference in a new issue