cln-plugins/backup/server.py

147 lines
5 KiB
Python
Raw Permalink Normal View History

2022-12-27 13:52:45 +01:00
import logging
import socket
import struct
2020-12-30 21:58:39 +01:00
import json
import sys
2020-12-30 21:58:39 +01:00
from typing import Tuple
from backend import Backend
2024-07-01 16:34:29 +02:00
from protocol import (
PacketType,
PKT_CHANGE_TYPES,
change_from_packet,
packet_from_change,
send_packet,
recv_packet,
)
2022-12-27 13:52:45 +01:00
2020-12-30 21:58:39 +01:00
class SystemdHandler(logging.Handler):
PREFIX = {
# EMERG <0>
# ALERT <1>
logging.CRITICAL: "<2>",
logging.ERROR: "<3>",
logging.WARNING: "<4>",
# NOTICE <5>
logging.INFO: "<6>",
logging.DEBUG: "<7>",
2024-07-01 16:34:29 +02:00
logging.NOTSET: "<7>",
}
def __init__(self, stream=sys.stdout):
self.stream = stream
logging.Handler.__init__(self)
def emit(self, record):
try:
msg = self.PREFIX[record.levelno] + self.format(record) + "\n"
self.stream.write(msg)
self.stream.flush()
except Exception:
self.handleError(record)
2022-12-27 13:52:45 +01:00
def setup_server_logging(mode, level):
root_logger = logging.getLogger()
root_logger.setLevel(level.upper())
mode = mode.lower()
2024-07-01 16:34:29 +02:00
if mode == "systemd":
# replace handler with systemd one
root_logger.handlers = []
root_logger.addHandler(SystemdHandler())
else:
2024-07-01 16:34:29 +02:00
assert mode == "plain"
2022-12-27 13:52:45 +01:00
2020-12-30 21:58:39 +01:00
class SocketServer:
def __init__(self, addr: Tuple[str, int], backend: Backend) -> None:
self.backend = backend
self.addr = addr
self.bind = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.bind.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
self.bind.bind(addr)
def _send_packet(self, typ: int, payload: bytes) -> None:
send_packet(self.sock, typ, payload)
def _recv_packet(self) -> Tuple[int, bytes]:
return recv_packet(self.sock)
def _handle_conn(self, conn) -> None:
# Can only handle one connection at a time
2024-07-01 16:34:29 +02:00
logging.info("Servicing incoming connection")
2020-12-30 21:58:39 +01:00
self.sock = conn
while True:
try:
(typ, payload) = self._recv_packet()
2022-12-27 13:52:45 +01:00
except IOError:
2024-07-01 16:34:29 +02:00
logging.info("Connection closed")
2020-12-30 21:58:39 +01:00
break
if typ in PKT_CHANGE_TYPES:
change = change_from_packet(typ, payload)
if typ == PacketType.CHANGE:
2024-07-01 16:34:29 +02:00
logging.debug("Received CHANGE {}".format(change.version))
2020-12-30 21:58:39 +01:00
else:
2024-07-01 16:34:29 +02:00
logging.info("Received SNAPSHOT {}".format(change.version))
2020-12-30 21:58:39 +01:00
self.backend.add_change(change)
2024-07-01 16:34:29 +02:00
self._send_packet(
PacketType.ACK, struct.pack("!I", self.backend.version)
)
2020-12-30 21:58:39 +01:00
elif typ == PacketType.REWIND:
2024-07-01 16:34:29 +02:00
logging.info("Received REWIND")
(to_version,) = struct.unpack("!I", payload)
2020-12-30 21:58:39 +01:00
if to_version != self.backend.prev_version:
2024-07-01 16:34:29 +02:00
logging.info("Cannot rewind to version {}".format(to_version))
self._send_packet(
PacketType.NACK, struct.pack("!I", self.backend.version)
)
2020-12-30 21:58:39 +01:00
else:
self.backend.rewind()
2024-07-01 16:34:29 +02:00
self._send_packet(
PacketType.ACK, struct.pack("!I", self.backend.version)
)
2020-12-30 21:58:39 +01:00
elif typ == PacketType.REQ_METADATA:
2024-07-01 16:34:29 +02:00
logging.debug("Received REQ_METADATA")
blob = struct.pack(
"!IIIQ",
0x01,
self.backend.version,
self.backend.prev_version,
self.backend.version_count,
)
2020-12-30 21:58:39 +01:00
self._send_packet(PacketType.METADATA, blob)
elif typ == PacketType.RESTORE:
2024-07-01 16:34:29 +02:00
logging.info("Received RESTORE")
2020-12-30 21:58:39 +01:00
for change in self.backend.stream_changes():
(typ, payload) = packet_from_change(change)
self._send_packet(typ, payload)
2024-07-01 16:34:29 +02:00
self._send_packet(PacketType.DONE, b"")
2020-12-30 21:58:39 +01:00
elif typ == PacketType.COMPACT:
2024-07-01 16:34:29 +02:00
logging.info("Received COMPACT")
2020-12-30 21:58:39 +01:00
stats = self.backend.compact()
self._send_packet(PacketType.COMPACT_RES, json.dumps(stats).encode())
elif typ == PacketType.ACK:
2024-07-01 16:34:29 +02:00
logging.debug("Received ACK")
2020-12-30 21:58:39 +01:00
elif typ == PacketType.NACK:
2024-07-01 16:34:29 +02:00
logging.debug("Received NACK")
2020-12-30 21:58:39 +01:00
elif typ == PacketType.METADATA:
2024-07-01 16:34:29 +02:00
logging.debug("Received METADATA")
2020-12-30 21:58:39 +01:00
elif typ == PacketType.COMPACT_RES:
2024-07-01 16:34:29 +02:00
logging.debug("Received COMPACT_RES")
2020-12-30 21:58:39 +01:00
else:
2024-07-01 16:34:29 +02:00
raise Exception("Unknown or unexpected packet type {}".format(typ))
2020-12-30 21:58:39 +01:00
self.conn = None
def run(self) -> None:
self.bind.listen(1)
2024-07-01 16:34:29 +02:00
logging.info("Waiting for connection on {}".format(self.addr))
2020-12-30 21:58:39 +01:00
while True:
conn, _ = self.bind.accept()
try:
self._handle_conn(conn)
2022-12-27 13:52:45 +01:00
except Exception:
2024-07-01 16:34:29 +02:00
logging.exception("Got exception")
2020-12-30 21:58:39 +01:00
finally:
conn.close()