diff --git a/.gitignore b/.gitignore index b3c4f4b..5859bcd 100644 --- a/.gitignore +++ b/.gitignore @@ -6,4 +6,5 @@ mergin htmlcov .vscode .env -config.yaml \ No newline at end of file +config.yaml +.python-version \ No newline at end of file diff --git a/config.py b/config.py index 6b39550..d443523 100644 --- a/config.py +++ b/config.py @@ -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", diff --git a/config.yaml.default b/config.yaml.default index 1d5f03b..8df1c24 100644 --- a/config.yaml.default +++ b/config.yaml.default @@ -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 diff --git a/dbsync.py b/dbsync.py index 1214e21..5553c46 100644 --- a/dbsync.py +++ b/dbsync.py @@ -6,6 +6,7 @@ License: MIT """ +import datetime import getpass import json import os @@ -13,6 +14,7 @@ import string import subprocess import tempfile +import typing import random import uuid import re @@ -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='*****'" @@ -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( @@ -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: diff --git a/dbsync_daemon.py b/dbsync_daemon.py index 64cc932..01de4b7 100644 --- a/dbsync_daemon.py +++ b/dbsync_daemon.py @@ -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 ( @@ -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 @@ -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) @@ -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__": diff --git a/docs/using.md b/docs/using.md index a058734..6b9c5cd 100644 --- a/docs/using.md +++ b/docs/using.md @@ -62,13 +62,31 @@ daemon: sleep_time: 10 ``` +When running as a daemon, the tool does not exit immediately when login, initialization or synchronization fails. The failed step is retried +with exponential backoff, starting at `sleep_time` and doubling after each consecutive failure up to 10 minutes (or `sleep_time` +if it is longer). Failed logins are retried less often, up to 1 hour, as rejected credentials can not be fixed by retrying. +Once a sync succeeds, the regular `sleep_time` interval is used again. + +Synchronization errors after the daemon has successfully started (e.g. temporary network or server issues) are retried indefinitely. +If the start of the daemon (login or initialization) keeps failing, or unexpected errors keep occurring, the daemon exits +after `max_retries` consecutive failed retries (10 by default, which takes about 50 minutes with `sleep_time: 10`, +or longer if the login is failing). +Set `max_retries: 0` to never exit: + +```yaml +daemon: + sleep_time: 10 + # Number of consecutive failed retries of the daemon start or unexpected errors before the daemon exits (0 = never exit) + max_retries: 10 +``` + ## Useful command line options - `config_file_name.yaml` The file name with path of yaml config can be provided. By default the tool uses `config.yaml` file from the current directory. - `--force-init` forces reinitialization of the sync. Drops dbsync schemas from database and the sync file and inits them all from scratch. This should be used to fix issues with dbsync init. -- `--single-run` instead of running the daemon indefinitely, performs just one single run. Such run consists of initialization, pull and push steps. +- `--single-run` instead of running the daemon indefinitely, performs just one single run. Such run consists of initialization, pull and push steps. Unlike the daemon, it exits with a non-zero code on failure. Avoid running it in a tight loop (e.g. from cron every minute) - use the daemon instead. - `--skip-init` allows skipping the initialization of sync step. Should be only used if you know, what you are doing, otherwise issues are likely to occur. diff --git a/requirements-dev.txt b/requirements-dev.txt index 98837e0..ceecf02 100644 --- a/requirements-dev.txt +++ b/requirements-dev.txt @@ -1,2 +1,3 @@ pytest>=6.2 -pytest-cov>=3.0 \ No newline at end of file +pytest-cov>=3.0 +pytest-mock>=3.10 diff --git a/test/test_auth_token.py b/test/test_auth_token.py new file mode 100644 index 0000000..449982f --- /dev/null +++ b/test/test_auth_token.py @@ -0,0 +1,193 @@ +import datetime +import json +import os +import stat + +import pytest +from mergin import ClientError, MerginClient + +import dbsync +from config import config + +from .conftest import _reset_config + +STORED_TOKEN = "Bearer stored" +NEW_TOKEN = "Bearer new" +MALFORMED_TOKEN = "Bearer malformed" + + +@pytest.fixture +def token_config(tmp_path): + """Configures Mergin credentials and working directory where the auth token is stored""" + config.update( + { + "MERGIN__URL": "https://mergin.example.com", + "MERGIN__USERNAME": "user", + "MERGIN__PASSWORD": "pwd", + "WORKING_DIR": str(tmp_path), + } + ) + return tmp_path / ".mergin_auth.json" + + +@pytest.fixture +def mergin_client(mocker): + """Mocks MerginClient: stored token is valid for `stored_token_validity` hours, login issues NEW_TOKEN + valid for 12 hours, MALFORMED_TOKEN can not be decoded. Server response to token validation can be changed + via `user_info` mock.""" + state = mocker.Mock(stored_token_validity=12, user_info=mocker.Mock()) + + def create_client(url, auth_token=None, login=None, password=None, plugin_version=None): + if auth_token == MALFORMED_TOKEN: + raise ClientError("Auth token error") + mc = mocker.MagicMock() + now = datetime.datetime.now(datetime.timezone.utc) + if auth_token: + validity = state.stored_token_validity + mc._auth_session = {"token": auth_token, "expire": now + datetime.timedelta(hours=validity)} + else: + mc._auth_session = {"token": NEW_TOKEN, "expire": now + datetime.timedelta(hours=12)} + mc.user_info = state.user_info + return mc + + state.cls = mocker.patch("dbsync.MerginClient", side_effect=create_client) + return state + + +def _store_token(path, url="https://mergin.example.com", username="user", token=STORED_TOKEN): + path.write_text(json.dumps({"url": url, "username": username, "token": token})) + + +def _logins(mergin_client): + """Number of MerginClient instances created with login (i.e. without stored token)""" + return sum(1 for c in mergin_client.cls.call_args_list if not c.kwargs.get("auth_token")) + + +def test_login_stores_token(token_config, mergin_client): + """Without stored token, DB sync logs in and stores the new token readable only by the owner""" + mc = dbsync.create_mergin_client() + + assert mc._auth_session["token"] == NEW_TOKEN + assert _logins(mergin_client) == 1 + assert json.loads(token_config.read_text()) == { + "url": "https://mergin.example.com", + "username": "user", + "token": NEW_TOKEN, + } + if os.name == "posix": + assert stat.S_IMODE(token_config.stat().st_mode) == 0o600 + + +def test_stored_token_is_reused(token_config, mergin_client): + """Valid stored token is used without new login, after checking the server still accepts it""" + _store_token(token_config) + + mc = dbsync.create_mergin_client() + + assert mc._auth_session["token"] == STORED_TOKEN + assert _logins(mergin_client) == 0 + mergin_client.user_info.assert_called_once() + + +@pytest.mark.parametrize( + "stored, stored_token_validity, user_info_error", + [ + (None, 12, None), + ({"url": "https://other.example.com"}, 12, None), + ({"username": "other"}, 12, None), + ("not a json", 12, None), + ({"token": MALFORMED_TOKEN}, 12, None), + ({}, 0.5, None), + ({}, 12, ClientError("Unauthorized", http_error=401)), + ], + ids=[ + "no-token", + "other-server", + "other-user", + "corrupted-file", + "malformed-token", + "expires-soon", + "rejected-by-server", + ], +) +def test_stored_token_is_not_used(token_config, mergin_client, stored, stored_token_validity, user_info_error): + """Stored token is not used when missing, issued for other server / user, unreadable, malformed, + about to expire or rejected by the server - new login is done and its token is stored instead""" + if isinstance(stored, dict): + _store_token(token_config, **stored) + elif stored: + token_config.write_text(stored) + mergin_client.stored_token_validity = stored_token_validity + mergin_client.user_info.side_effect = user_info_error + + mc = dbsync.create_mergin_client() + + assert mc._auth_session["token"] == NEW_TOKEN + assert _logins(mergin_client) == 1 + assert json.loads(token_config.read_text())["token"] == NEW_TOKEN + + +def test_login_works_when_token_can_not_be_stored(token_config, mergin_client, mocker, caplog): + """Failure to store the token (e.g. read-only working directory) is only logged, login still succeeds""" + mocker.patch("dbsync.os.open", side_effect=PermissionError("Read-only file system")) + + mc = dbsync.create_mergin_client() + + assert mc._auth_session["token"] == NEW_TOKEN + assert not token_config.exists() + assert "Unable to store Mergin Maps auth token: Read-only file system" in caplog.text + + +def test_stored_token_server_error(token_config, mergin_client): + """Server error other than 401 when validating stored token does not cause new login, + the stored token is used (and the sync itself fails and is retried later if the server is really unavailable)""" + _store_token(token_config) + mergin_client.user_info.side_effect = ClientError("Service unavailable", http_error=503) + + mc = dbsync.create_mergin_client() + + assert mc._auth_session["token"] == STORED_TOKEN + assert _logins(mergin_client) == 0 + assert json.loads(token_config.read_text())["token"] == STORED_TOKEN + + +def test_clean_keeps_stored_token(token_config, mergin_client, mocker): + """Cleaning (--force-init) removes the working directory but keeps the auth token, + so restart with --force-init does not need new login""" + mocker.patch("dbsync.psycopg2.connect") + mocker.patch("dbsync._drop_schema") + config.update({"init_from": "gpkg", "CONNECTIONS": [{"conn_info": "", "modified": "main", "base": "base"}]}) + project_dir = token_config.parent / "project" + project_dir.mkdir() + mc = dbsync.create_mergin_client() + + dbsync.dbsync_clean(mc) + + assert not project_dir.exists() + assert json.loads(token_config.read_text())["token"] == NEW_TOKEN + + +def test_stored_token_with_server(tmp_path, mocker): + """Integration test with real server: stored token is reused without new login, + tampered token is rejected by the server and new login is done""" + _reset_config() + config.update({"WORKING_DIR": str(tmp_path)}) + token_path = tmp_path / ".mergin_auth.json" + login = mocker.spy(MerginClient, "login") + + mc = dbsync.create_mergin_client() + assert login.call_count == 1 + assert token_path.exists() + + mc = dbsync.create_mergin_client() + assert login.call_count == 1 + # stored token is accepted by the server (configured login can be either username or email) + user_info = mc.user_info() + assert config.mergin.username in (user_info["username"], user_info["email"]) + + stored = json.loads(token_path.read_text()) + token_path.write_text(json.dumps({**stored, "token": stored["token"][:-4] + "abcd"})) + + mc = dbsync.create_mergin_client() + assert login.call_count == 2 + assert json.loads(token_path.read_text())["token"] != stored["token"][:-4] + "abcd" diff --git a/test/test_basic.py b/test/test_basic.py index b5d31f3..56a1707 100644 --- a/test/test_basic.py +++ b/test/test_basic.py @@ -974,8 +974,8 @@ def test_dbsync_clean_from_gpkg( dbsync_clean(mc) - # after the dbsync_clean nothing exists - assert pathlib.Path(config.working_dir).exists() is False + # after the dbsync_clean nothing exists, except the stored auth token + assert [p.name for p in pathlib.Path(config.working_dir).iterdir()] == [".mergin_auth.json"] assert ( _check_schema_exists( conn, diff --git a/test/test_config.py b/test/test_config.py index 0275bf7..cceaa26 100644 --- a/test/test_config.py +++ b/test/test_config.py @@ -389,3 +389,20 @@ def test_skip_and_include_tables_mutually_exclusive(): config.update({"CONNECTIONS": [{**base, "skip_tables": ["a"], "include_tables": ["b"]}]}) with pytest.raises(ConfigError, match="cannot both be set"): validate_config(config) + + +def test_config_daemon_max_retries(): + """`max_retries` in `daemon` section must be a non-negative integer""" + _reset_config() + config.update({"DAEMON__MAX_RETRIES": 5}) + validate_config(config) + + config.update({"DAEMON__MAX_RETRIES": 0}) + validate_config(config) + + for value in [-1, "abc", 1.5, True]: + config.update({"DAEMON__MAX_RETRIES": value}) + with pytest.raises(ConfigError, match="`max_retries` in `daemon` must be a non-negative integer"): + validate_config(config) + + config.unset("DAEMON__MAX_RETRIES", force=True) diff --git a/test/test_daemon.py b/test/test_daemon.py new file mode 100644 index 0000000..c7360f9 --- /dev/null +++ b/test/test_daemon.py @@ -0,0 +1,224 @@ +import argparse +import datetime +import os + +import pytest +from mergin import ClientError, MerginClient + +import dbsync +import dbsync_daemon +from config import config + +from .conftest import DB_CONNINFO, TEST_DATA_DIR, init_sync_from_geopackage + + +class StopDaemon(Exception): + """Raised from mocked sleep to break out of the infinite daemon loop""" + + +@pytest.fixture +def run_daemon(mocker): + """Returns function running the daemon loop until it exits or sleeps `iterations` times. + + Sleeping is mocked, the function returns list of the requested sleep times and whether the daemon exited.""" + + def run( + iterations=100, + force_init=False, + skip_init=False, + sleep_time=10, + max_retries=10, + send_notifications=False, + on_sleep=None, + ): + sleeps = [] + + def fake_sleep(seconds): + sleeps.append(seconds) + if on_sleep: + on_sleep(len(sleeps)) + if len(sleeps) >= iterations: + raise StopDaemon() + + mocker.patch("dbsync_daemon.time.sleep", side_effect=fake_sleep) + args = argparse.Namespace(force_init=force_init, skip_init=skip_init) + try: + dbsync_daemon.run_daemon(args, sleep_time, max_retries, send_notifications) + except StopDaemon: + return sleeps, False + except SystemExit: + return sleeps, True + + return run + + +@pytest.fixture +def dbsync_mocks(mocker): + """Mocks Mergin client creation, dbsync steps and notification emails""" + mc = mocker.MagicMock() + mc._auth_session = {"expire": datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(hours=12)} + mocker.patch("dbsync_daemon.config").notification = {} + return mocker.Mock( + mc=mc, + create_client=mocker.patch("dbsync.create_mergin_client", return_value=mc), + clean=mocker.patch("dbsync.dbsync_clean"), + init=mocker.patch("dbsync.dbsync_init"), + pull=mocker.patch("dbsync.dbsync_pull"), + push=mocker.patch("dbsync.dbsync_push"), + send_email=mocker.patch("dbsync_daemon.send_email"), + token_store=mocker.patch("dbsync.AuthTokenStore"), + ) + + +def test_retry_wait_time(): + """Wait time doubles with each consecutive failure, starting at sleep time and capped at max wait + (MAX_RETRY_WAIT by default) or sleep time if longer""" + assert [dbsync_daemon.retry_wait_time(10, f) for f in range(1, 9)] == [10, 20, 40, 80, 160, 320, 600, 600] + # sleep time longer than the max retry wait is respected + assert dbsync_daemon.retry_wait_time(3600, 5) == 3600 + # no overflow with many failures + assert dbsync_daemon.retry_wait_time(10, 10**6) == dbsync_daemon.MAX_RETRY_WAIT + # longer limit for failed logins + assert [dbsync_daemon.retry_wait_time(10, f, 3600) for f in range(7, 11)] == [640, 1280, 2560, 3600] + + +def test_daemon_startup_is_retried(run_daemon, dbsync_mocks): + """Failed login and init are retried in the same process with backoff: login is repeated only until it succeeds, + --force-init cleaning is done only once and regular sleep time is used again once the sync succeeds""" + dbsync_mocks.create_client.side_effect = [dbsync.DbSyncError("login failed")] * 2 + [dbsync_mocks.mc] + dbsync_mocks.init.side_effect = [dbsync.DbSyncError("init failed")] * 2 + [None] + + sleeps, exited = run_daemon(iterations=6, force_init=True) + + assert not exited + assert sleeps == [10, 20, 40, 80, 10, 10] + assert dbsync_mocks.create_client.call_count == 3 + dbsync_mocks.clean.assert_called_once() + assert dbsync_mocks.init.call_count == 3 + assert dbsync_mocks.pull.call_count == 2 + + +def test_daemon_keeps_running_after_start(run_daemon, dbsync_mocks): + """Sync errors after a successful start (e.g. server unavailable) never make the daemon exit, only back off. + A single unexpected error does not stop it either. No new login is done in any case.""" + dbsync_mocks.pull.side_effect = ( + [None] + [dbsync.DbSyncError("server unavailable")] * 8 + [RuntimeError("boom")] + [None] + ) + + sleeps, exited = run_daemon(iterations=11, max_retries=2) + + assert not exited + assert sleeps == [10, 10, 20, 40, 80, 160, 320, 600, 600, 600, 10] + dbsync_mocks.create_client.assert_called_once() + dbsync_mocks.init.assert_called_once() + assert dbsync_mocks.push.call_count == 2 + + +@pytest.mark.parametrize( + "failing_step, errors, max_retries, expected_sleeps, expected_exit", + [ + ("init", [dbsync.DbSyncError("init failed")] * 4, 3, [10, 20, 40], True), + ("pull", [None] + [RuntimeError("boom")] * 3, 2, [10, 10, 20], True), + ( + "create_client", + dbsync.DbSyncError("login failed"), + 0, + [10, 20, 40, 80, 160, 320, 640, 1280, 2560, 3600], + False, + ), + ], + ids=["startup-failures", "unexpected-errors", "never-exit"], +) +def test_daemon_gives_up_after_max_retries( + run_daemon, dbsync_mocks, failing_step, errors, max_retries, expected_sleeps, expected_exit +): + """Daemon exits after `max_retries` consecutive failed retries of the start or of unexpected errors + (never with max_retries 0). Notification email is always sent when giving up, regardless of the minimal + email interval.""" + getattr(dbsync_mocks, failing_step).side_effect = errors + + # daemon that should not exit is stopped after the expected number of sleeps + iterations = 100 if expected_exit else len(expected_sleeps) + sleeps, exited = run_daemon(iterations=iterations, max_retries=max_retries, send_notifications=True) + + assert exited == expected_exit + assert sleeps == expected_sleeps + if expected_exit: + # first failure and giving up, the ones in between are suppressed by the minimal email interval + assert dbsync_mocks.send_email.call_count == 2 + assert dbsync_mocks.send_email.call_args.args[0].startswith(f"Giving up after {max_retries} retries") + else: + dbsync_mocks.send_email.assert_called_once() + + +def test_daemon_logs_in_again_when_token_rejected(run_daemon, dbsync_mocks): + """When the server rejects the auth token during sync, stored token is removed and new login is done + on the next attempt. Other failures do not cause new login.""" + dbsync_mocks.pull.side_effect = [None, dbsync.DbSyncError("server unavailable"), dbsync.DbSyncError("401"), None] + dbsync_mocks.mc.user_info.side_effect = [None, ClientError("Unauthorized", http_error=401)] + + sleeps, exited = run_daemon(iterations=4) + + assert not exited + assert sleeps == [10, 10, 20, 10] + assert dbsync_mocks.create_client.call_count == 2 + dbsync_mocks.token_store.from_config.return_value.remove.assert_called_once() + + +def test_daemon_credentials_rejected_after_start(run_daemon, dbsync_mocks): + """When the token is rejected after a successful start and the new login fails too (e.g. user was deactivated), + the daemon keeps running but retries the login with longer wait times (up to MAX_LOGIN_RETRY_WAIT)""" + dbsync_mocks.create_client.side_effect = [dbsync_mocks.mc] + [dbsync.DbSyncError("login failed")] * 20 + dbsync_mocks.pull.side_effect = [None, dbsync.DbSyncError("401")] + dbsync_mocks.mc.user_info.side_effect = ClientError("Unauthorized", http_error=401) + + sleeps, exited = run_daemon(iterations=13, max_retries=2) + + assert not exited + assert sleeps == [10, 10, 20, 40, 80, 160, 320, 640, 1280, 2560, 3600, 3600, 3600] + # first login + login retried after each sleep since the token was rejected + assert dbsync_mocks.create_client.call_count == 12 + + +def test_daemon_recovers_from_init_failure(mc: MerginClient, run_daemon, mocker): + """Integration test with real server and database: init fails on database connection for two attempts, + the daemon retries without logging in again and continues syncing once the database connection is fixed.""" + init_sync_from_geopackage(mc, "test_daemon_recovery", os.path.join(TEST_DATA_DIR, "base.gpkg")) + connection = dict(config.connections[0]) + config.update({"CONNECTIONS": [{**connection, "conn_info": DB_CONNINFO + " password=wrong"}]}) + + def fix_db_connection(sleep_count): + if sleep_count == 2: + config.update({"CONNECTIONS": [connection]}) + + login = mocker.spy(MerginClient, "login") + pull = mocker.spy(dbsync, "dbsync_pull") + + sleeps, exited = run_daemon(iterations=4, max_retries=5, on_sleep=fix_db_connection) + + assert not exited + assert sleeps == [10, 20, 10, 10] + assert login.call_count == 1 + assert pull.call_count == 2 + + +def test_daemon_recovers_from_rejected_token(mc: MerginClient, run_daemon, mocker): + """Integration test with real server and database: the server starts rejecting the auth token + of the running daemon, the daemon logs in again and continues syncing.""" + init_sync_from_geopackage(mc, "test_daemon_token", os.path.join(TEST_DATA_DIR, "base.gpkg")) + login = mocker.spy(MerginClient, "login") + create_client = mocker.spy(dbsync, "create_mergin_client") + pull = mocker.spy(dbsync, "dbsync_pull") + + def invalidate_token(sleep_count): + if sleep_count == 1: + daemon_mc = create_client.spy_return + daemon_mc._auth_session["token"] = daemon_mc._auth_session["token"][:-4] + "abcd" + + sleeps, exited = run_daemon(iterations=3, on_sleep=invalidate_token) + + assert not exited + assert sleeps == [10, 10, 10] + assert login.call_count == 2 + assert pull.call_count == 3 + assert pull.spy_exception is None