#!/usr/bin/python3
# mqtt2logger - forward the MQTT logs of ESPHome and Tasmota devices to rsyslog and systemd-journald
#
# Copyright (C) 2026 Thomas Wagner <wagner-thomas@gmx.at>
# SPDX-License-Identifier: GPL-2.0-or-later
"""Subscribe to the log topics of ESPHome and Tasmota devices on an MQTT
broker and forward each log line to rsyslog (RFC 5424, keeping the device as
host name) and to systemd-journald (native protocol), with the severity and
time stamp found in the message where there is one.
"""

import argparse
import collections
import datetime
import logging
import os
import re
import signal
import socket
import ssl
import struct
import sys
import threading
import time

VERSION = "0.1"

LEVELS = ["emerg", "alert", "crit", "err", "warning", "notice", "info", "debug"]
LEVEL_ALIASES = {"emergency": 0, "panic": 0, "critical": 2, "error": 3, "warn": 4, "informational": 6}
FACILITIES = {"kern": 0, "user": 1, "mail": 2, "daemon": 3, "auth": 4, "syslog": 5, "lpr": 6, "news": 7,
              "uucp": 8, "cron": 9, "authpriv": 10, "ftp": 11, **{"local%d" % n: 16 + n for n in range(8)}}

# ESPHome's log levels, as the letter in "[W][tag:line]: message"
ESPHOME_LEVELS = {"E": 3, "W": 4, "I": 6, "C": 6, "D": 7, "V": 7, "VV": 7}

DEFAULTS = {
    "mqtt": {"host": "localhost", "port": None, "username": None, "password_file": None, "password": None,
             "tls": False, "ca_cert": None, "insecure": False, "client_id": None, "keepalive": 60, "qos": 1,
             "clean_session": False},
    "sources": [
        {"type": "esphome", "topic": "+/debug"},
        {"type": "tasmota", "topic": "stat/+/LOGGING"},
    ],
    "ignore_retained": True,
    "tasmota": {"priority": "info"},
    "priority_rules": [],
    "timestamps": {"max_age": 3600, "max_future": 120},
    "rsyslog": {"enabled": True, "address": "tcp://127.0.0.1:10514", "facility": "local0",
                "hostname": "{device}", "app_name": "{firmware}"},
    "journal": {"enabled": False, "socket": "/run/systemd/journal/socket", "identifier": "{device}",
                "facility": "local0"},
    "queue_size": 10000,
}

log = logging.getLogger("mqtt2logger")


class ConfigError(Exception):
    pass


def parse_level(text):
    value = str(text).strip().lower()
    if value.isdigit() and int(value) <= 7:
        return int(value)
    if value in LEVELS:
        return LEVELS.index(value)
    if value in LEVEL_ALIASES:
        return LEVEL_ALIASES[value]
    raise ConfigError("unknown severity %r" % text)


def parse_facility(text):
    value = str(text).strip().lower()
    if value.isdigit() and int(value) <= 23:
        return int(value)
    if value in FACILITIES:
        return FACILITIES[value]
    raise ConfigError("unknown facility %r" % text)


# ---------------------------------------------------------------------------
# log records and the parsers of the two firmwares
# ---------------------------------------------------------------------------

class Record:
    """One log line of a device."""
    __slots__ = ("device", "firmware", "tag", "message", "priority", "time", "device_time", "topic")

    def __init__(self, device, firmware, tag, message, priority, time, device_time, topic):
        self.device = device
        self.firmware = firmware
        self.tag = tag
        self.message = message
        self.priority = priority
        self.time = time                # the time stamp to log, seconds since the epoch
        self.device_time = device_time  # whether it came from the device
        self.topic = topic


ANSI_RE = re.compile(r"\x1b\[[0-9;]*[A-Za-z]")
ESPHOME_RE = re.compile(r"^\[(VV|[EWICDV])\]\[([^\]]*?)(?::(\d+))?\]:?\s?(.*)$", re.S)
TASMOTA_RE = re.compile(r"^(\d\d):(\d\d):(\d\d)\.(\d{3})(?:-\d+)?(?:/\d+)?\s+(?:([A-Z0-9]{2,5}): )?(.*)$", re.S)


def parse_esphome(text, received):
    """ESPHome publishes "[W][modbus_controller:083]: message" in ANSI colours, without a time stamp."""
    text = ANSI_RE.sub("", text).strip()
    match = ESPHOME_RE.match(text)
    if not match:
        return [(None, "", line, received, False) for line in text.splitlines() if line.strip()]
    level, tag, _line, message = match.groups()
    lines = [line for line in message.splitlines() if line.strip()] or [""]
    return [(ESPHOME_LEVELS[level], tag, line, received, False) for line in lines]


def tasmota_time(hour, minute, second, millis, received, timestamps):
    """The device's local time of day on the day that puts it closest to the reception, if plausible.

    Tasmota counts from midnight of 1970-01-01 until its clock is synchronized,
    which makes such times implausible.
    """
    base = datetime.datetime.fromtimestamp(received)
    best = None
    for days in (-1, 0, 1):
        day = base + datetime.timedelta(days=days)
        try:
            candidate = day.replace(hour=hour, minute=minute, second=second, microsecond=millis * 1000).timestamp()
        except ValueError:
            continue
        if best is None or abs(candidate - received) < abs(best - received):
            best = candidate
    if best is None or not received - timestamps["max_age"] <= best <= received + timestamps["max_future"]:
        return received, False
    return best, True


def parse_tasmota(text, received, timestamps):
    """Tasmota publishes "13:45:21.999-123 MQT: message", the time of day without severity."""
    lines = []
    for line in text.replace("\r", "").splitlines():
        if not line.strip():
            continue
        match = TASMOTA_RE.match(line)
        if not match:
            lines.append((None, "", line, received, False))
            continue
        hour, minute, second, millis, tag, message = match.groups()
        stamp, from_device = tasmota_time(int(hour), int(minute), int(second), int(millis), received, timestamps)
        lines.append((None, tag or "", message, stamp, from_device))
    return lines


# ---------------------------------------------------------------------------
# priority rules
# ---------------------------------------------------------------------------

RULE_FIELDS = {"device": "device", "host": "device", "firmware": "firmware", "tag": "tag", "component": "tag",
               "message": "message", "content": "message", "topic": "topic"}


class PriorityRule:
    def __init__(self, spec):
        if not isinstance(spec, dict) or "priority" not in spec:
            raise ConfigError("a priority rule needs a priority and at least one of %s: %r"
                              % (", ".join(sorted(set(RULE_FIELDS.values()))), spec))
        self.priority = parse_level(spec["priority"])
        self.conditions = []
        for key, value in spec.items():
            if key == "priority":
                continue
            field = RULE_FIELDS.get(str(key).lower())
            if field is None:
                raise ConfigError("unknown field %r in priority rule" % key)
            try:
                self.conditions.append((field, re.compile(str(value))))
            except re.error as err:
                raise ConfigError("invalid regular expression %r: %s" % (value, err))
        if not self.conditions:
            raise ConfigError("a priority rule needs at least one condition: %r" % spec)

    def matches(self, record):
        return all(regex.search(getattr(record, field) or "") for field, regex in self.conditions)


# ---------------------------------------------------------------------------
# sources: topics and which device and firmware they belong to
# ---------------------------------------------------------------------------

class Source:
    def __init__(self, spec):
        if not isinstance(spec, dict) or not spec.get("topic"):
            raise ConfigError("a source needs a topic: %r" % spec)
        self.firmware = str(spec.get("type", "esphome")).lower()
        if self.firmware not in ("esphome", "tasmota"):
            raise ConfigError("unknown source type %r, use esphome or tasmota" % spec.get("type"))
        self.topic = spec["topic"]
        self.levels = self.topic.split("/")
        wildcards = [index for index, level in enumerate(self.levels) if level == "+"]
        self.device_level = spec.get("device_level", wildcards[0] if wildcards else None)
        self.device = spec.get("device")
        self.priority = parse_level(spec["priority"]) if spec.get("priority") is not None else None

    def matches(self, topic):
        parts = topic.split("/")
        for index, level in enumerate(self.levels):
            if level == "#":
                return True
            if index >= len(parts) or (level != "+" and level != parts[index]):
                return False
        return len(parts) == len(self.levels)

    def device_of(self, topic):
        if self.device:
            return self.device
        parts = topic.split("/")
        if self.device_level is not None and self.device_level < len(parts):
            return parts[self.device_level]
        return topic


# ---------------------------------------------------------------------------
# outputs
# ---------------------------------------------------------------------------

def render(template, record):
    return template.format(device=record.device, firmware=record.firmware, tag=record.tag or "-",
                           topic=record.topic)


def syslog_token(text, limit):
    """A header field of RFC 5424: printable US-ASCII without spaces, or "-"."""
    token = "".join(char for char in text if 33 <= ord(char) <= 126)[:limit]
    return token or "-"


def rfc5424(record, facility, hostname, app_name):
    stamp = datetime.datetime.fromtimestamp(record.time).astimezone().isoformat(timespec="microseconds")
    priority = facility * 8 + (6 if record.priority is None else record.priority)
    header = "<%d>1 %s %s %s - %s -" % (priority, stamp, syslog_token(hostname, 255),
                                       syslog_token(app_name, 48), syslog_token(record.tag, 32))
    return header + " " + record.message.replace("\n", " ")


def journal_fields(record, identifier, facility):
    fields = [("MESSAGE", record.message), ("PRIORITY", str(6 if record.priority is None else record.priority)),
              ("SYSLOG_IDENTIFIER", identifier), ("SYSLOG_FACILITY", str(facility)),
              ("DEVICE", record.device), ("DEVICE_FIRMWARE", record.firmware), ("MQTT_TOPIC", record.topic)]
    if record.tag:
        fields.append(("DEVICE_COMPONENT", record.tag))
    if record.device_time:
        # journald stamps an entry with its own time; the device's goes in a field
        fields.append(("DEVICE_TIMESTAMP", "%d" % round(record.time * 1e6)))
        fields.append(("SYSLOG_TIMESTAMP",
                       datetime.datetime.fromtimestamp(record.time).astimezone().isoformat(timespec="milliseconds")))
    return fields


def journal_datagram(fields):
    """The native protocol of systemd-journald; values with a newline are sent with their length."""
    out = bytearray()
    for name, value in fields:
        data = value.encode("utf-8", "replace")
        if b"\n" in data:
            out += name.encode() + b"\n" + struct.pack("<Q", len(data)) + data + b"\n"
        else:
            out += name.encode() + b"=" + data + b"\n"
    return bytes(out)


class Output(threading.Thread):
    """Delivers records from its own queue, retrying while the receiver is unavailable."""

    name_text = "output"

    def __init__(self, queue_size):
        super().__init__(daemon=True)
        self.queue = collections.deque()
        self.queue_size = queue_size
        self.condition = threading.Condition()
        self.stopping = False
        self.dropped = 0
        self.sent = 0

    def put(self, record):
        with self.condition:
            if len(self.queue) >= self.queue_size:
                self.queue.popleft()
                self.dropped += 1
                if self.dropped == 1 or self.dropped % 1000 == 0:
                    log.error("%s: queue full, %d message(s) dropped", self.name_text, self.dropped)
            self.queue.append(record)
            self.condition.notify()

    def stop(self):
        with self.condition:
            self.stopping = True
            self.condition.notify()

    def run(self):
        delay = 1
        while True:
            with self.condition:
                while not self.queue and not self.stopping:
                    self.condition.wait()
                if not self.queue:
                    break
                record = self.queue[0]
            try:
                self.send(record)
            except OSError as err:
                if self.stopping:
                    log.error("%s: %d message(s) not delivered at shutdown: %s", self.name_text, len(self.queue),
                              err)
                    break
                log.warning("%s: %s, retrying in %ds", self.name_text, err, delay)
                self.close()
                time.sleep(delay)
                delay = min(delay * 2, 60)
                continue
            delay = 1
            self.sent += 1
            with self.condition:
                self.queue.popleft()
        self.close()

    def send(self, record):
        raise NotImplementedError

    def close(self):
        pass


class RsyslogOutput(Output):
    name_text = "rsyslog"

    def __init__(self, config, queue_size):
        super().__init__(queue_size)
        self.facility = parse_facility(config.get("facility", "local0"))
        self.hostname = config.get("hostname", "{device}")
        self.app_name = config.get("app_name", "{firmware}")
        address = str(config.get("address", "tcp://127.0.0.1:10514"))
        match = re.match(r"^(tcp|udp)://(\[[^\]]+\]|[^:/]+):(\d+)$", address)
        if match:
            self.transport, host, port = match.group(1), match.group(2).strip("[]"), int(match.group(3))
            self.target = (host, port)
        elif address.startswith("unix:"):
            self.transport, self.target = "unix", address[5:]
        else:
            raise ConfigError("rsyslog address is tcp://HOST:PORT, udp://HOST:PORT or unix:PATH, not %r" % address)
        self.sock = None

    def connect(self):
        if self.transport == "tcp":
            self.sock = socket.create_connection(self.target, timeout=30)
        elif self.transport == "udp":
            family = socket.getaddrinfo(self.target[0], self.target[1], 0, socket.SOCK_DGRAM)[0][0]
            self.sock = socket.socket(family, socket.SOCK_DGRAM)
            self.sock.connect(self.target)
        else:
            self.sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM)
            self.sock.connect(self.target)

    def send(self, record):
        if self.sock is None:
            self.connect()
        line = rfc5424(record, self.facility, render(self.hostname, record), render(self.app_name, record))
        data = line.encode("utf-8", "replace")
        if self.transport == "tcp":
            # octet counting framing of RFC 6587, so messages may contain anything
            self.sock.sendall(b"%d " % len(data) + data)
        else:
            self.sock.send(data)

    def close(self):
        if self.sock is not None:
            try:
                self.sock.close()
            except OSError:
                pass
            self.sock = None


class JournalOutput(Output):
    name_text = "journal"

    def __init__(self, config, queue_size):
        super().__init__(queue_size)
        self.path = config.get("socket", "/run/systemd/journal/socket")
        self.identifier = config.get("identifier", "{device}")
        self.facility = parse_facility(config.get("facility", "local0"))
        self.sock = None

    def send(self, record):
        if self.sock is None:
            self.sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM)
        self.sock.sendto(journal_datagram(journal_fields(record, render(self.identifier, record), self.facility)),
                         self.path)

    def close(self):
        if self.sock is not None:
            self.sock.close()
            self.sock = None


# ---------------------------------------------------------------------------
# the daemon
# ---------------------------------------------------------------------------

class Forwarder:
    def __init__(self, config):
        self.sources = [Source(spec) for spec in config["sources"]]
        if not self.sources:
            raise ConfigError("no sources configured")
        self.rules = [PriorityRule(spec) for spec in config.get("priority_rules") or []]
        self.tasmota_priority = parse_level(config["tasmota"].get("priority", "info"))
        self.timestamps = {"max_age": float(config["timestamps"]["max_age"]),
                           "max_future": float(config["timestamps"]["max_future"])}
        self.ignore_retained = bool(config.get("ignore_retained", True))
        self.outputs = []
        if enabled(config["rsyslog"].get("enabled")):
            self.outputs.append(RsyslogOutput(config["rsyslog"], int(config["queue_size"])))
        if enabled(config["journal"].get("enabled")):
            self.outputs.append(JournalOutput(config["journal"], int(config["queue_size"])))
        if not self.outputs:
            raise ConfigError("neither rsyslog nor journal output is enabled")
        self.received = 0

    def records(self, topic, payload, received):
        source = next((source for source in self.sources if source.matches(topic)), None)
        if source is None:
            return []
        text = payload.decode("utf-8", "replace")
        if source.firmware == "esphome":
            lines = parse_esphome(text, received)
        else:
            lines = parse_tasmota(text, received, self.timestamps)
        device = source.device_of(topic)
        result = []
        for priority, tag, message, stamp, from_device in lines:
            if priority is None:
                priority = source.priority if source.priority is not None else (
                    self.tasmota_priority if source.firmware == "tasmota" else 6)
            record = Record(device, source.firmware, tag, message, priority, stamp, from_device, topic)
            for rule in self.rules:
                if rule.matches(record):
                    record.priority = rule.priority
                    break
            result.append(record)
        return result

    def on_message(self, client, userdata, message):
        if message.retain and self.ignore_retained:
            # a retained message is the last one before we subscribed, possibly long ago
            return
        self.received += 1
        for record in self.records(message.topic, message.payload, time.time()):
            log.debug("%s %s %s [%s]: %s", record.device, record.firmware, LEVELS[record.priority],
                      record.tag, record.message)
            for output in self.outputs:
                output.put(record)


def enabled(value):
    if isinstance(value, str):
        return value.strip().lower() in ("1", "yes", "true", "on")
    return bool(value)


def read_password(mqtt):
    path = mqtt.get("password_file") or os.environ.get("MQTT_PASSWORD_FILE")
    if path:
        try:
            with open(path) as handle:
                return handle.read().rstrip("\r\n")
        except OSError as err:
            raise ConfigError("cannot read password file %s: %s" % (path, err.strerror))
    return mqtt.get("password") or os.environ.get("MQTT_PASSWORD")


def create_client(mqtt, client_id, clean_session):
    import paho.mqtt.client as paho
    if hasattr(paho, "CallbackAPIVersion"):
        return paho.Client(paho.CallbackAPIVersion.VERSION2, client_id=client_id, clean_session=clean_session)
    return paho.Client(client_id=client_id, clean_session=clean_session)


def is_failure(reason_code):
    return reason_code.is_failure if hasattr(reason_code, "is_failure") else reason_code != 0


def run_daemon(config):
    forwarder = Forwarder(config)
    try:
        import paho.mqtt.client  # noqa: F401
    except ImportError:
        raise ConfigError("mqtt2logger needs paho-mqtt (python3-paho-mqtt)")
    mqtt = config["mqtt"]
    client_id = mqtt.get("client_id") or "mqtt2logger-%s" % socket.gethostname()
    clean_session = bool(mqtt.get("clean_session", False))
    qos = int(mqtt.get("qos", 1))
    client = create_client(mqtt, client_id, clean_session)
    state = {"stopping": False}

    def on_connect(client, userdata, flags, reason_code, properties=None):
        if is_failure(reason_code):
            log.error("broker refused the connection: %s", reason_code)
            return
        session = flags.get("session present") if isinstance(flags, dict) else getattr(flags, "session_present", None)
        log.info("connected to %s:%s%s", mqtt["host"], port, ", session resumed" if session else "")
        for source in forwarder.sources:
            client.subscribe(source.topic, qos=qos)
            log.info("subscribed to %s (%s)", source.topic, source.firmware)

    def on_disconnect(client, userdata, *args):
        if not state["stopping"]:
            log.warning("disconnected from the broker, reconnecting")

    client.on_connect = on_connect
    client.on_disconnect = on_disconnect
    client.on_message = forwarder.on_message
    client.reconnect_delay_set(min_delay=1, max_delay=60)
    password = read_password(mqtt)
    if mqtt.get("username") or password:
        client.username_pw_set(mqtt.get("username"), password)
    use_tls = mqtt.get("tls") or mqtt.get("ca_cert")
    if use_tls:
        client.tls_set(ca_certs=mqtt.get("ca_cert"),
                       cert_reqs=ssl.CERT_NONE if mqtt.get("insecure") else ssl.CERT_REQUIRED)
        if mqtt.get("insecure"):
            client.tls_insecure_set(True)
    port = int(mqtt.get("port") or (8883 if use_tls else 1883))

    for output in forwarder.outputs:
        output.start()

    def stop(signum, frame):
        state["stopping"] = True
        client.disconnect()

    signal.signal(signal.SIGTERM, stop)
    signal.signal(signal.SIGINT, stop)

    log.info("mqtt2logger %s forwarding to %s", VERSION, ", ".join(output.name_text for output in forwarder.outputs))
    client.connect_async(mqtt["host"], port, keepalive=int(mqtt.get("keepalive", 60)))
    client.loop_forever(retry_first_connection=True)

    # deliver what is still queued, for at most 10 seconds
    for output in forwarder.outputs:
        output.stop()
    deadline = time.time() + 10
    for output in forwarder.outputs:
        output.join(max(0.1, deadline - time.time()))
        log.info("%s: %d message(s) sent, %d dropped, %d left", output.name_text, output.sent, output.dropped,
                 len(output.queue))
    return 0


# ---------------------------------------------------------------------------
# configuration
# ---------------------------------------------------------------------------

def merge_config(base, override):
    result = dict(base)
    for key, value in override.items():
        if isinstance(value, dict) and isinstance(result.get(key), dict):
            result[key] = merge_config(result[key], value)
        else:
            result[key] = value
    return result


def load_yaml(path):
    try:
        import yaml
    except ImportError:
        raise ConfigError("reading %s needs PyYAML (python3-PyYAML / python3-yaml)" % path)
    try:
        with open(path, encoding="utf-8") as handle:
            data = yaml.safe_load(handle) or {}
    except OSError as err:
        raise ConfigError("cannot read %s: %s" % (path, err.strerror))
    except yaml.YAMLError as err:
        raise ConfigError("invalid YAML in %s: %s" % (path, err))
    if not isinstance(data, dict):
        raise ConfigError("%s must contain a mapping" % path)
    return data


def build_parser():
    parser = argparse.ArgumentParser(
        prog="mqtt2logger",
        description="Forward the MQTT logs of ESPHome and Tasmota devices to rsyslog and systemd-journald, with "
                    "the severity and time stamp found in the messages.",
        epilog="Options given here override the configuration file.")
    parser.add_argument("-V", "--version", action="version", version="mqtt2logger %s" % VERSION)
    parser.add_argument("-C", "--config", help="YAML configuration file")
    parser.add_argument("-H", "--mqtt-host", help="MQTT broker (default: localhost)")
    parser.add_argument("-p", "--mqtt-port", type=int, help="MQTT port (default: 1883, 8883 with TLS)")
    parser.add_argument("-u", "--mqtt-username")
    parser.add_argument("--mqtt-password-file", help="file holding the MQTT password (also: $MQTT_PASSWORD_FILE)")
    parser.add_argument("--tls", action="store_true", help="connect to the broker with TLS")
    parser.add_argument("--ca-cert", help="CA bundle for the broker's certificate; implies --tls")
    parser.add_argument("--client-id", help="MQTT client id; keep it fixed for the persistent session "
                                            "(default: mqtt2logger-HOSTNAME)")
    parser.add_argument("--esphome-topic", action="append", metavar="TOPIC",
                        help="topic of ESPHome logs, + marks the device (default: +/debug); repeatable")
    parser.add_argument("--tasmota-topic", action="append", metavar="TOPIC",
                        help="topic of Tasmota logs, + marks the device (default: stat/+/LOGGING); repeatable")
    parser.add_argument("--rsyslog", metavar="ADDRESS",
                        help="forward to rsyslog at tcp://HOST:PORT, udp://HOST:PORT or unix:PATH "
                             "(default: tcp://127.0.0.1:10514)")
    parser.add_argument("--no-rsyslog", action="store_true", help="do not forward to rsyslog")
    parser.add_argument("--journal", action="store_true", help="also forward to systemd-journald")
    parser.add_argument("--no-journal", action="store_true", help="do not forward to systemd-journald")
    parser.add_argument("-v", "--verbose", action="store_true", help="log each forwarded message")
    return parser


def load_config(args):
    config = merge_config(DEFAULTS, load_yaml(args.config)) if args.config else merge_config(DEFAULTS, {})
    mqtt = dict(config["mqtt"])
    for key, value in (("host", args.mqtt_host), ("port", args.mqtt_port), ("username", args.mqtt_username),
                       ("password_file", args.mqtt_password_file), ("ca_cert", args.ca_cert),
                       ("client_id", args.client_id)):
        if value is not None:
            mqtt[key] = value
    if args.tls:
        mqtt["tls"] = True
    config["mqtt"] = mqtt
    if args.esphome_topic or args.tasmota_topic:
        config["sources"] = ([{"type": "esphome", "topic": topic} for topic in args.esphome_topic or []]
                             + [{"type": "tasmota", "topic": topic} for topic in args.tasmota_topic or []])
    if args.rsyslog:
        config["rsyslog"] = dict(config["rsyslog"], address=args.rsyslog, enabled=True)
    if args.no_rsyslog:
        config["rsyslog"] = dict(config["rsyslog"], enabled=False)
    if args.journal:
        config["journal"] = dict(config["journal"], enabled=True)
    if args.no_journal:
        config["journal"] = dict(config["journal"], enabled=False)
    return config


def main(argv=None):
    args = build_parser().parse_args(argv)
    logging.basicConfig(level=logging.DEBUG if args.verbose else logging.INFO,
                        format="%(levelname)s: %(message)s", stream=sys.stderr)
    try:
        config = load_config(args)
        return run_daemon(config)
    except ConfigError as err:
        log.error("%s", err)
        return 2


if __name__ == "__main__":
    sys.exit(main())
