Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
3 changes: 2 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -6,4 +6,5 @@ mergin
htmlcov
.vscode
.env
config.yaml
config.yaml
.python-version
5 changes: 5 additions & 0 deletions config.py
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,11 @@ def validate_config(config):
):
raise ConfigError("Config error: `include_tables` parameter should be a list")

if "daemon" in config and "max_retries" in config.daemon:
max_retries = config.daemon.max_retries
if isinstance(max_retries, bool) or not isinstance(max_retries, int) or max_retries < 0:
raise ConfigError("Config error: `max_retries` in `daemon` must be a non-negative integer.")

if "notification" in config:
settings = [
"smtp_server",
Expand Down
2 changes: 2 additions & 0 deletions config.yaml.default
Original file line number Diff line number Diff line change
Expand Up @@ -29,3 +29,5 @@ connections:
daemon:
# How often to synchronize (in seconds)
sleep_time: 10
# Number of consecutive failed retries of the daemon start (login, init) or unexpected errors before the daemon exits (0 = never exit)
max_retries: 10
108 changes: 104 additions & 4 deletions dbsync.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,15 @@
License: MIT
"""

import datetime
import getpass
import json
import os
import shutil
import string
import subprocess
import tempfile
import typing
import random
import uuid
import re
Expand Down Expand Up @@ -52,6 +54,9 @@

FORCE_INIT_MESSAGE = "Running `dbsync_deamon.py` with `--force-init` should fix the issue."

# auth token with less validity left (in seconds) is not reused and new login is done instead
TOKEN_MIN_VALIDITY = 3600


class DbSyncError(Exception):
default_print_password = "password='*****'"
Expand Down Expand Up @@ -593,16 +598,108 @@ def _validate_local_project_id(
)


def create_mergin_client():
"""Create instance of MerginClient"""
_check_has_password()
class AuthTokenStore:
"""Stores Mergin Maps auth token in the working directory, so it can be reused
when DB sync is restarted instead of new login"""

def __init__(self, working_dir: str, url: str, username: str):
self.path = pathlib.Path(working_dir) / ".mergin_auth.json"
self.url = url
self.username = username

@classmethod
def from_config(cls) -> "AuthTokenStore":
return cls(config.working_dir, config.mergin.url, config.mergin.username)

def load(self) -> typing.Optional[str]:
"""Returns stored token if it was issued for this server and user"""
try:
with open(self.path) as f:
data = json.load(f)
if data["url"] == self.url and data["username"] == self.username:
return data["token"]
except FileNotFoundError:
pass
except (OSError, ValueError, KeyError, TypeError) as e:
logging.warning(f"Unable to read stored Mergin Maps auth token: {e}")
return None

def save(self, token: str) -> None:
try:
self.path.parent.mkdir(parents=True, exist_ok=True)
# readable only by the owner (on POSIX systems)
with open(os.open(self.path, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600), "w") as f:
json.dump({"url": self.url, "username": self.username, "token": token}, f)
except OSError as e:
logging.warning(f"Unable to store Mergin Maps auth token: {e}")

def remove(self) -> None:
self.path.unlink(missing_ok=True)


def auth_token_expires_soon(mc: MerginClient) -> bool:
delta = mc._auth_session["expire"] - datetime.datetime.now(datetime.timezone.utc)
return delta.total_seconds() < TOKEN_MIN_VALIDITY


def auth_token_rejected(mc: MerginClient) -> bool:
"""Checks whether the server rejects auth token of the client although it has not expired yet.
Other errors (e.g. server unavailable) are ignored."""
try:
return MerginClient(
mc.user_info()
except ClientError as e:
return e.http_error == 401
except Exception:
pass
return False


def _create_mergin_client_from_stored_token() -> typing.Optional[MerginClient]:
"""Creates MerginClient using stored auth token, returns None if there is no valid stored token"""
token_store = AuthTokenStore.from_config()
token = token_store.load()
if not token:
return None

try:
mc = MerginClient(
config.mergin.url,
auth_token=token,
login=config.mergin.username,
password=config.mergin.password,
plugin_version=f"DB-sync/{__version__}",
)
except ClientError as e:
logging.warning(f"Stored Mergin Maps auth token is invalid: {e}")
token_store.remove()
return None

if auth_token_expires_soon(mc):
return None

if auth_token_rejected(mc):
logging.debug("Stored Mergin Maps auth token was rejected by the server")
token_store.remove()
return None

logging.debug("Using stored Mergin Maps auth token")
return mc


def create_mergin_client():
"""Create instance of MerginClient, reusing stored auth token if possible"""
_check_has_password()
try:
mc = _create_mergin_client_from_stored_token()
if mc is None:
mc = MerginClient(
config.mergin.url,
login=config.mergin.username,
password=config.mergin.password,
plugin_version=f"DB-sync/{__version__}",
)
AuthTokenStore.from_config().save(mc._auth_session["token"])
return mc
except LoginError as e:
# this could be auth failure, but could be also server problem (e.g. worker crash)
raise DbSyncError(
Expand Down Expand Up @@ -1468,6 +1565,9 @@ def clean(conn_cfg, mc):
except FileNotFoundError as e:
raise DbSyncError("Unable to remove working directory: " + str(e))

# keep the auth token, so it can be reused after restart (e.g. when --force-init is kept in container command)
AuthTokenStore.from_config().save(mc._auth_session["token"])

if from_db:
temp_folder = pathlib.Path(config.working_dir).parent / "project_to_delete_sync_file"
try:
Expand Down
153 changes: 111 additions & 42 deletions dbsync_daemon.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,14 @@
from smtp_functions import send_email
from version import __version__

# upper limit (in seconds) of the wait time between retries after repeated failures
MAX_RETRY_WAIT = 600
# upper limit (in seconds) of the wait time for failed logins - credentials rejected
# by the server can not be fixed by retrying and too frequent failed logins may lock the account
MAX_LOGIN_RETRY_WAIT = 3600
# default number of consecutive failed retries of startup (login / init) or unexpected errors before the daemon exits
DEFAULT_MAX_RETRIES = 10


def is_pyinstaller() -> bool:
if (
Expand Down Expand Up @@ -127,6 +135,7 @@ def main():
validate_config(config)
except ConfigError as e:
handle_error_and_exit(e)
max_retries = config.get("daemon.max_retries", DEFAULT_MAX_RETRIES)

send_notifications = "notification" in config

Expand All @@ -145,21 +154,17 @@ def main():
if args.force_init and args.skip_init:
handle_error_and_exit("Cannot use `--force-init` with `--skip-init` Initialization is required. ")

logging.debug("Logging in to Mergin...")

mc = dbsync.create_mergin_client()
if args.single_run:
try:
logging.debug("Logging in to Mergin...")
mc = dbsync.create_mergin_client()

if args.force_init:
dbsync.dbsync_clean(mc)
if args.force_init:
dbsync.dbsync_clean(mc)

if args.single_run:
if not args.skip_init:
try:
if not args.skip_init:
dbsync.dbsync_init(mc)
except dbsync.DbSyncError as e:
handle_error_and_exit(e)

try:
logging.debug("Trying to pull")
dbsync.dbsync_pull(mc)

Expand All @@ -170,46 +175,110 @@ def main():
handle_error_and_exit(e)

else:
if not args.skip_init:
try:
dbsync.dbsync_init(mc)
except dbsync.DbSyncError as e:
handle_error_and_exit(e)
run_daemon(args, sleep_time, max_retries, send_notifications)

last_email_sent = None

while True:
print(datetime.datetime.now())
def retry_wait_time(sleep_time: int, failures: int, max_wait: int = MAX_RETRY_WAIT) -> int:
"""Seconds to wait before the next attempt after `failures` consecutive failed attempts.

try:
logging.debug("Trying to pull")
dbsync.dbsync_pull(mc)
Doubles with each failure up to max_wait, but never less than sleep_time."""
wait_time = sleep_time * 2 ** (failures - 1)
return min(wait_time, max(sleep_time, max_wait))

logging.debug("Trying to push")
dbsync.dbsync_push(mc)

# check mergin client token expiration
delta = mc._auth_session["expire"] - datetime.datetime.now(datetime.timezone.utc)
if delta.total_seconds() < 3600:
mc = dbsync.create_mergin_client()
def run_daemon(args, sleep_time: int, max_retries: int, send_notifications: bool) -> None:
"""Keep syncing until killed. Failures (including login and init) are retried within
this process with exponential backoff instead of exiting, so the daemon does not
log in again on every restart by the container / service manager.

except dbsync.DbSyncError as e:
logging.error(str(e))
if send_notifications:
if "minimal_email_interval" in config.notification:
min_time_delta_hr = config.notification.minimal_email_interval
else:
min_time_delta_hr = 4
Sync errors after a successful start are retried indefinitely. Startup failures (login, clean, init)
and unexpected errors make the daemon exit after `max_retries` consecutive failed retries (0 = never exit)."""
mc = None
cleaned = not args.force_init
initialized = args.skip_init
started = False
failures = 0
fatal_failures = 0
last_email_sent = None

if (
last_email_sent is None
or (datetime.datetime.now() - last_email_sent).total_seconds() > min_time_delta_hr * 3600
):
send_email(str(e), config)
last_email_sent = datetime.datetime.now()
while True:
print(datetime.datetime.now())
login_failed = False

try:
if mc is None:
logging.debug("Logging in to Mergin...")
mc = dbsync.create_mergin_client()

if not cleaned:
dbsync.dbsync_clean(mc)
cleaned = True

if not initialized:
dbsync.dbsync_init(mc)
initialized = True

logging.debug("Trying to pull")
dbsync.dbsync_pull(mc)

logging.debug("Trying to push")
dbsync.dbsync_push(mc)
started = True

# check mergin client token expiration
if dbsync.auth_token_expires_soon(mc):
mc = dbsync.create_mergin_client()

failures = 0
fatal_failures = 0

except Exception as e:
failures += 1
# client is created by login as the first step, so no client means the login has failed
login_failed = mc is None
if not started or not isinstance(e, dbsync.DbSyncError):
fatal_failures += 1
if isinstance(e, dbsync.DbSyncError):
error_msg = str(e)
logging.error(error_msg)
else:
error_msg = f"Unexpected error: {e!r}"
logging.exception(error_msg)

# server may reject the token before it expires, log in again on the next attempt in such case
if mc is not None and dbsync.auth_token_rejected(mc):
logging.warning("Mergin Maps auth token was rejected by the server, going to log in again")
dbsync.AuthTokenStore.from_config().remove()
mc = None

giving_up = max_retries and fatal_failures > max_retries
if giving_up:
error_msg = f"Giving up after {max_retries} retries, the daemon will exit: {error_msg}"

if send_notifications:
if "minimal_email_interval" in config.notification:
min_time_delta_hr = config.notification.minimal_email_interval
else:
min_time_delta_hr = 4

if (
giving_up
or last_email_sent is None
or (datetime.datetime.now() - last_email_sent).total_seconds() > min_time_delta_hr * 3600
):
send_email(error_msg, config)
last_email_sent = datetime.datetime.now()

if giving_up:
handle_error_and_exit(error_msg)

if failures:
wait_time = retry_wait_time(sleep_time, failures, MAX_LOGIN_RETRY_WAIT if login_failed else MAX_RETRY_WAIT)
logging.debug(f"Failed attempt #{failures}, going to sleep for {wait_time} seconds before retrying")
else:
wait_time = sleep_time
logging.debug("Going to sleep")
time.sleep(sleep_time)
time.sleep(wait_time)


if __name__ == "__main__":
Expand Down
Loading
Loading