diff --git a/src/server/server.c b/src/server/server.c index 494a840..6f3862c 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -2,6 +2,7 @@ #include "charset.h" #include "credentials.h" #include "daemon_conf.h" +#include "daemon_limits.h" #include "delay_updates.h" #include "file.h" #include "identity.h" @@ -56,6 +57,12 @@ static DaemonConf* g_daemon_conf = NULL; * such a module exists. */ static CredentialStore* g_credentials = NULL; +/* Cross-process connection registry (per-module and per-source caps plus the + * shared auth lockout), created once in main BEFORE the accept loop forks and + * shared read-only-by-pointer with every connection child. NULL outside daemon + * mode or when the mapping could not be allocated (global cap + ACLs remain). */ +static DaemonLimitRegistry* g_daemon_limits = NULL; + /* Opaque context threaded through to the config-frame gate: the connection's * SSL object (NULL over plaintext) so the gate can warn when a credential * exchange is not encrypted, plus the super-mode override the gate decides on. @@ -292,6 +299,56 @@ static const DaemonModule* module_gate_lookup_module(const Config* config, const return module; } +/* Index of `module` within the loaded config's module array (the registry's + * per-module counter key). Returns -1 when it cannot be resolved. */ +static int daemon_module_index(const DaemonModule* module) { + if (!g_daemon_conf || !module || module < g_daemon_conf->modules || + module >= g_daemon_conf->modules + g_daemon_conf->module_count) + return -1; + return (int)(module - g_daemon_conf->modules); +} + +/* Shared-registry admission: reserve this connection's slot for the selected + * module and the peer source IP. Enforces the per-module `max connections` and + * the global `max connections per host` across every forked child. Runs before + * auth/ownership so a client that is over a cap is refused before any work. + * The per-source cap is skipped when the peer cannot be classified (host ACLs + * fail closed separately); the module cap still applies. A missing registry + * (allocation failure / non-fork path) fails open -- the global cap and ACLs + * still bound the listener. */ +static const char* module_gate_check_limits(const Config* config, const DaemonModule* module, + ModuleGateContext* gate_ctx) { + if (!g_daemon_limits) + return NULL; + int slot = transport_tcp_current_slot(); + if (slot < 0) + return NULL; /* not on the forked accept-loop path (e.g. --stdio) */ + int module_index = daemon_module_index(module); + if (module_index < 0) + return NULL; + const char* peer = (gate_ctx && gate_ctx->has_peer_ip) ? gate_ctx->peer_ip : ""; + DaemonLimitResult result = + daemon_limits_register(g_daemon_limits, slot, module_index, peer, module->max_connections); + switch (result) { + case DAEMON_LIMIT_OK: + return NULL; + case DAEMON_LIMIT_MODULE_FULL: + log_message(LOG_LEVEL_ERROR, + "daemon module '%s': 'max connections' cap (%d) reached; refusing %s", + config->module, module->max_connections, peer[0] ? peer : "peer"); + return "requested daemon module is at its connection limit"; + case DAEMON_LIMIT_HOST_FULL: + log_message(LOG_LEVEL_ERROR, + "daemon: 'max connections per host' cap (%d) reached for %s; refusing module '%s'", + g_daemon_conf->global.max_connections_per_host, peer[0] ? peer : "peer", + config->module); + return "too many concurrent connections from this host"; + case DAEMON_LIMIT_UNAVAILABLE: + default: + return NULL; + } +} + /* Per-module client-chosen ownership / super-user policy (P7 Wave E hardening): * a daemon module refuses EVERY ownership-affecting request (--numeric-ids, * --chown, --usermap/--groupmap, --fake-super, --copy-as, explicit --super) @@ -396,6 +453,20 @@ static ModuleAuthResult module_gate_authenticate(const Config* config, const Dae ModuleGateContext* gate_ctx, const char** error) { if (module->auth_user_count == 0) return MODULE_AUTH_ACCEPTED; + /* Cross-process lockout: a source that failed too many authentications is + * refused before the challenge is sent (the counter lives in the shared + * registry, so it spans every forked child and survives a child exit). */ + if (g_daemon_limits && gate_ctx && gate_ctx->has_peer_ip) { + int remaining = 0; + if (daemon_limits_auth_locked(g_daemon_limits, gate_ctx->peer_ip, &remaining)) { + log_message(LOG_LEVEL_ERROR, + "daemon module '%s': source %s is locked out after repeated authentication " + "failures (%d s remaining); refusing", + config->module, gate_ctx->peer_ip, remaining); + *error = "too many failed authentication attempts from this host; try again later"; + return MODULE_AUTH_REFUSED; + } + } /* Fail closed: no store -> refuse (server misconfiguration, STATUS_ERROR). */ if (g_credentials == NULL) { log_message(LOG_LEVEL_ERROR, @@ -450,10 +521,16 @@ static ModuleAuthResult module_gate_authenticate(const Config* config, const Dae "daemon module '%s': authentication failed for user '%s' from %s; refusing", config->module, escaped_user ? escaped_user : "(none)", peer); free(escaped_user); - /* Rate-limit online guessing per connection (no delay on success). */ + /* Count the failure in the shared registry (locks the source out once the + * configured threshold is reached) and rate-limit online guessing per + * connection (no delay on success). */ + if (g_daemon_limits && gate_ctx->has_peer_ip) + daemon_limits_auth_record_failure(g_daemon_limits, gate_ctx->peer_ip); daemon_auth_failure_delay(); return MODULE_AUTH_TERMINATED; } + if (g_daemon_limits && gate_ctx->has_peer_ip) + daemon_limits_auth_record_success(g_daemon_limits, gate_ctx->peer_ip); char* escaped_user = output_escape(config->auth_user, config->eight_bit_output); log_message(LOG_LEVEL_INFO, "daemon module '%s': user '%s' from %s authenticated", config->module, escaped_user ? escaped_user : "", @@ -561,6 +638,9 @@ static const char* server_module_gate(const Config* config, void* context) { log_message(LOG_LEVEL_DEBUG, "daemon module '%s': peer address unavailable", config->module); } error = module_gate_check_hosts(config, module, gate_ctx); + if (error) + return error; + error = module_gate_check_limits(config, module, gate_ctx); if (error) return error; error = module_gate_check_ownership(config, module, gate_ctx); @@ -880,7 +960,9 @@ static void print_server_usage(void) { printf(" fastsyncd.conf, else /etc/fastsyncd.conf)\n"); printf(" --dparam=KEY=VALUE Override one global config key on the command line\n"); printf(" (port, motd file, address, max connections,\n"); - printf(" auth failure delay, hosts allow, hosts deny)\n"); + printf(" max connections per host, auth failure delay,\n"); + printf(" auth lockout threshold, auth lockout duration,\n"); + printf(" hosts allow, hosts deny)\n"); printf(" --no-detach Stay in the foreground (default detaches to\n"); printf(" background when running --daemon)\n"); printf(" --password-file=FILE Credential store for modules that declare\n"); @@ -1096,11 +1178,10 @@ int main(int argc, char* argv[]) { "unless the module is intentionally open to the network", g_daemon_conf->modules[i].name); if (g_daemon_conf->modules[i].max_connections > 0) - log_message(LOG_LEVEL_WARNING, - "daemon module '%s': per-module 'max connections' is stored but not enforced " - "per module; the global 'max connections' cap (%d) applies to the whole " - "listener", - g_daemon_conf->modules[i].name, g_daemon_conf->global.max_connections); + log_message(LOG_LEVEL_INFO, + "daemon module '%s': per-module 'max connections' cap = %d (enforced " + "across all connection children)", + g_daemon_conf->modules[i].name, g_daemon_conf->modules[i].max_connections); } /* Daemon credential store (Wave B). --password-file and --early-input * feed the same store, loaded BEFORE the listener forks so every @@ -1144,6 +1225,21 @@ int main(int argc, char* argv[]) { module->name, module->auth_users[j]); } } + /* Shared cross-process registry for the per-module / per-source caps and + * the auth lockout. Created HERE in the parent before any accept-loop + * fork; every connection child inherits the mapping. A failure degrades to + * "registry disabled" (the global cap and host ACLs still apply) rather + * than refusing to start. */ + g_daemon_limits = daemon_limits_create((int)g_daemon_conf->global.max_connections, + g_daemon_conf->module_count, + g_daemon_conf->global.max_connections_per_host, + g_daemon_conf->global.auth_lockout_threshold, + g_daemon_conf->global.auth_lockout_duration_sec); + if (!g_daemon_limits) + log_message(LOG_LEVEL_WARNING, + "daemon: could not allocate the shared connection registry; per-module / " + "per-host caps and the cross-process auth lockout are disabled (the global " + "'max connections' cap and host ACLs still apply)"); } else { if (!configure_authorization(opts.destination_root)) { char* escaped = output_escape(opts.destination_root, false); @@ -1167,6 +1263,8 @@ int main(int argc, char* argv[]) { } if (g_daemon_conf) server_set_max_connections(g_server, (unsigned int)g_daemon_conf->global.max_connections); + if (g_daemon_limits) + server_set_limit_registry(g_server, g_daemon_limits); if (opts.use_tls) { if (!opts.tls_cert || !opts.tls_key || !opts.tls_ca || !opts.client_cn) { fprintf(stderr, "Error: --tls requires --cert, --key, --ca, and --client-cn\n"); @@ -1206,6 +1304,8 @@ int main(int argc, char* argv[]) { release_authorization(); out: + daemon_limits_destroy(g_daemon_limits); + g_daemon_limits = NULL; daemon_conf_free(g_daemon_conf); g_daemon_conf = NULL; credentials_free(g_credentials); diff --git a/src/shared/transport_tcp.c b/src/shared/transport_tcp.c index 2758cbe..e73dbce 100644 --- a/src/shared/transport_tcp.c +++ b/src/shared/transport_tcp.c @@ -1,4 +1,5 @@ #include "transport_tcp.h" +#include "daemon_limits.h" #include "log.h" #include "protocol.h" #include "utils.h" @@ -18,15 +19,25 @@ static volatile sig_atomic_t g_active_connections = 0; +/* Shared registry installed on the active server; the SIGCHLD handler needs a + * file-scope pointer so it can reclaim the dead child's slot. Set once by + * accept_loop before the fork loop (single-threaded parent). */ +static DaemonLimitRegistry* g_limit_registry = NULL; +/* Slot reserved by the parent for the connection child currently being forked. + * Written before fork(), read by the child (which inherits the value). */ +static int g_current_slot = DAEMON_LIMITS_NO_SLOT; + static void tcp_apply_socket_timeout(int fd); static void tcp_enable_nodelay_default(int fd, int family); static void sigchld_handler(int sig) { (void)sig; int saved_errno = errno; - while (waitpid(-1, NULL, WNOHANG) > 0) { + pid_t pid; + while ((pid = waitpid(-1, NULL, WNOHANG)) > 0) { if (g_active_connections > 0) g_active_connections--; + daemon_limits_reclaim_pid(g_limit_registry, (long)pid); } errno = saved_errno; } @@ -108,6 +119,7 @@ Server* server_create_ex(int port, const ServerBindOptions* bind_opts) { server->ssl_ctx = NULL; server->max_connections = 100; server->active_connections = 0; + server->limit_registry = NULL; return server; } @@ -121,6 +133,15 @@ void server_set_max_connections(Server* server, unsigned int max_connections) { server->max_connections = max_connections; } +void server_set_limit_registry(Server* server, struct DaemonLimitRegistry* registry) { + if (server) + server->limit_registry = registry; +} + +int transport_tcp_current_slot(void) { + return g_current_slot; +} + void server_delete(Server** server) { if (server == NULL || *server == NULL) return; @@ -140,6 +161,7 @@ static void accept_loop(Server* server, void (*child_fn)(int, void*), void* chil return; } signal(SIGCHLD, sigchld_handler); + g_limit_registry = server->limit_registry; while (1) { struct sockaddr_storage client_addr; socklen_t client_len = sizeof(client_addr); @@ -159,9 +181,31 @@ static void accept_loop(Server* server, void (*child_fn)(int, void*), void* chil close(fd); continue; } + int slot = DAEMON_LIMITS_NO_SLOT; + if (server->limit_registry) { + slot = daemon_limits_claim_slot(server->limit_registry); + if (slot == DAEMON_LIMITS_NO_SLOT) { + /* The global cap bounds live children, so this only happens when the + * fixed registry is smaller than the configured cap; fail closed. */ + log_message(LOG_LEVEL_WARNING, "Connection registry slots exhausted (max %u), rejecting %s", + server->max_connections, peer); + close(fd); + continue; + } + } log_message(LOG_LEVEL_INFO, "%s from %s", log_fmt, peer); + g_current_slot = slot; + /* Block SIGCHLD across fork() and the parent's pid publication: a child + * that exits immediately must not be reaped before its slot records its + * pid, which would leak the slot and its module/source counts. */ + sigset_t blocked; + sigset_t previous; + sigemptyset(&blocked); + sigaddset(&blocked, SIGCHLD); + sigprocmask(SIG_BLOCK, &blocked, &previous); pid_t pid = fork(); if (pid == 0) { + sigprocmask(SIG_SETMASK, &previous, NULL); /* Connection children must not run the parent's global cleanup(): it * frees state (credentials / daemon conf) that the child's worker * threads may still be reading and closes fd numbers the child could @@ -177,7 +221,13 @@ static void accept_loop(Server* server, void (*child_fn)(int, void*), void* chil _exit(0); } else if (pid > 0) { g_active_connections++; + if (server->limit_registry) + daemon_limits_set_slot_pid(server->limit_registry, slot, (long)pid); + } else if (server->limit_registry) { + /* fork() failed: release the reservation so the slot is not leaked. */ + daemon_limits_reclaim_slot(server->limit_registry, slot); } + sigprocmask(SIG_SETMASK, &previous, NULL); close(fd); } } diff --git a/src/shared/transport_tcp.h b/src/shared/transport_tcp.h index e37b879..c3862a6 100644 --- a/src/shared/transport_tcp.h +++ b/src/shared/transport_tcp.h @@ -7,6 +7,10 @@ #include #include +/* Cross-process daemon registry (daemon_limits.c). Only an opaque pointer is + * stored here so the transport layer does not depend on daemon config. */ +struct DaemonLimitRegistry; + typedef struct Server { struct sockaddr_storage address; unsigned int address_length; @@ -14,6 +18,7 @@ typedef struct Server { void* ssl_ctx; unsigned int max_connections; volatile unsigned int active_connections; + struct DaemonLimitRegistry* limit_registry; } Server; typedef struct Client { @@ -48,6 +53,14 @@ Server* server_create(int port); /* Override the listener's connection cap (the global daemon `max connections` * value). A non-positive value is ignored so the default cap stands. */ void server_set_max_connections(Server* server, unsigned int max_connections); +/* Install the shared per-module / per-source registry used by the accept loop + * to reserve a slot for each forked child. NULL disables the accounting (the + * global cap and ACLs still apply). */ +void server_set_limit_registry(Server* server, struct DaemonLimitRegistry* registry); +/* Slot reserved for the connection child currently running (set by the parent + * before fork, inherited by the child). Returns DAEMON_LIMITS_NO_SLOT (-1) + * outside the accept-loop child path. */ +int transport_tcp_current_slot(void); bool server_listen(Server* server, void (*handler)(int file_descriptor)); void server_accept_loop(Server* server, void (*child_fn)(int, void*), void* child_ctx, const char* log_fmt); diff --git a/tests/integration/test_daemon.py b/tests/integration/test_daemon.py index 98c963d..1ff1d9d 100644 --- a/tests/integration/test_daemon.py +++ b/tests/integration/test_daemon.py @@ -135,7 +135,7 @@ class DaemonManager: self._proc = None self._port = None - def start(self, config_path, port_override=None, extra_args=None): + def start(self, config_path, port_override=None, extra_args=None, log_path=None): self.stop() # When no override is given the daemon binds the config file's `port` # (the plain config-port path); with an override the --dparam path. @@ -146,7 +146,8 @@ class DaemonManager: cmd += ["--dparam", f"port={port_override}"] if extra_args: cmd += extra_args - log_path = os.path.join(TEST_DATA_DIR, "fastsyncd.log") + if log_path is None: + log_path = os.path.join(TEST_DATA_DIR, "fastsyncd.log") log = open(log_path, "w") self._proc = subprocess.Popen( cmd, stdout=log, stderr=log, stdin=subprocess.DEVNULL, start_new_session=True) @@ -1208,3 +1209,77 @@ class TestDaemonTLSAuth: d.stop() os.unlink(client_creds) shutil.rmtree(cert_dir, ignore_errors=True) + + +class TestDaemonConnectionLimits: + """Wave 8: cross-process per-module / per-source connection caps and the + shared auth lockout. Each test boots its own daemon with a unique port so + the shared (per-daemon) registry state is isolated from the module-scoped + `daemon` fixture.""" + + LOCKOUT_CONF = os.path.join(TEST_DATA_DIR, "fastsyncd_lockout.conf") + CAPS_CONF = os.path.join(TEST_DATA_DIR, "fastsyncd_caps.conf") + + @pytest.mark.ci + def test_auth_lockout_is_shared_across_children(self): + """`auth lockout threshold = 1`: the first failed authentication locks the + source out for the cooldown in the SHARED registry, so a subsequent + correct-password attempt (a different forked child) is refused before a + SCRAM challenge is even sent.""" + port = _find_free_port() + with open(self.LOCKOUT_CONF, "w") as f: + f.write("port = %d\n" + "auth lockout threshold = 1\n" + "auth lockout duration = 300\n" + "\n" + "[locked]\n" + "path = %s\n" + "auth users = alice\n" + % (port, AUTH_MODULE)) + d = DaemonManager() + log_path = os.path.join(TEST_DATA_DIR, f"fastsyncd_lockout_{os.getpid()}.log") + try: + d.start(self.LOCKOUT_CONF, port_override=port, extra_args=["--password-file", CRED_FILE], + log_path=log_path) + before = _tree_file_count(AUTH_MODULE) + log_before = os.path.getsize(log_path) if os.path.exists(log_path) else 0 + # First attempt: wrong password -> records failure #1 -> locks. + wrong = _push_with_creds("127.0.0.1::locked", port, "alice", WRONG_PASS) + assert wrong.returncode != 0 + # Second attempt: CORRECT password from the same source must still be + # refused by the shared lockout. + right = _push_with_creds("127.0.0.1::locked", port, "alice", ALICE_PASS) + assert right.returncode != 0, "the shared auth lockout must refuse after threshold" + assert _tree_file_count(AUTH_MODULE) == before, "a locked-out source wrote data" + time.sleep(0.3) + with open(log_path, "rb") as f: + f.seek(log_before) + tail = f.read().decode("utf-8", "replace") + assert "locked out" in tail, tail[-400:] + finally: + d.stop() + + def test_caps_keys_accepted_and_transfer_still_works(self): + """A daemon configured with the new keys (per-host cap, lockout threshold + and duration, per-module cap) starts and serves a normal transfer.""" + port = _find_free_port() + with open(self.CAPS_CONF, "w") as f: + f.write("port = %d\n" + "max connections per host = 5\n" + "auth lockout threshold = 3\n" + "auth lockout duration = 60\n" + "\n" + "[files]\n" + "path = %s\n" + "max connections = 2\n" + % (port, FILES_MODULE)) + d = DaemonManager() + try: + d.start(self.CAPS_CONF, port_override=port) + result = _push("127.0.0.1::files", port) + assert result.returncode == 0, result.stderr or result.stdout + received = get_dest_received_dir(FILES_MODULE, SOURCE_DIR) + _, missing = verify_transfer(SOURCE_DIR, received) + assert not missing, f"missing: {missing[:5]}" + finally: + d.stop()