diff --git a/src/client/client.c b/src/client/client.c index f2e3875..8d9a8d1 100644 --- a/src/client/client.c +++ b/src/client/client.c @@ -99,12 +99,21 @@ int load_files_multithreaded(void *pipeline_context) { int send_chunks_multithreaded(void *pipeline_context) { PipelineContextSender *context = (PipelineContextSender *)pipeline_context; - Client *client = client_create(); - const char *env_ip = getenv("FASTSYNC_SERVER_IP"); - const char *ip = env_ip ? env_ip : "127.0.0.1"; - const char *env_port = getenv("FASTSYNC_SERVER_PORT"); - int port = env_port ? atoi(env_port) : 8080; - client_connect(client, (char *)ip, port); + Client *client; + if (context->config->transport == TRANSPORT_SSH) { + if (context->config->use_sendfile) { + fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n"); + return 1; + } + client = client_connect_ssh(context->config->ssh_destination); + } else { + client = client_create(); + const char *env_ip = getenv("FASTSYNC_SERVER_IP"); + const char *ip = env_ip ? env_ip : "127.0.0.1"; + const char *env_port = getenv("FASTSYNC_SERVER_PORT"); + int port = env_port ? atoi(env_port) : 8080; + client_connect(client, (char *)ip, port); + } config_send(client->file_descriptor, context->config); while (true) { @@ -114,6 +123,11 @@ int send_chunks_multithreaded(void *pipeline_context) { &context->condition_not_full_loader, &context->loader_done); if (current_chunk == NULL) { send_status(client->file_descriptor, STATUS_FINISHED); + if (receive_status(client->file_descriptor) != STATUS_OK) { + client_disconnect(client); + client_delete(client); + return 1; + } client_disconnect(client); client_delete(client); return thrd_success; @@ -127,12 +141,21 @@ int send_chunks_multithreaded(void *pipeline_context) { } int send_files(Config *config) { - Client *client = client_create(); - const char *env_ip = getenv("FASTSYNC_SERVER_IP"); - const char *ip = env_ip ? env_ip : "127.0.0.1"; - const char *env_port = getenv("FASTSYNC_SERVER_PORT"); - int port = env_port ? atoi(env_port) : 8080; - client_connect(client, (char *)ip, port); + Client *client; + if (config->transport == TRANSPORT_SSH) { + if (config->use_sendfile) { + fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n"); + return 1; + } + client = client_connect_ssh(config->ssh_destination); + } else { + client = client_create(); + const char *env_ip = getenv("FASTSYNC_SERVER_IP"); + const char *ip = env_ip ? env_ip : "127.0.0.1"; + const char *env_port = getenv("FASTSYNC_SERVER_PORT"); + int port = env_port ? atoi(env_port) : 8080; + client_connect(client, (char *)ip, port); + } config_send(client->file_descriptor, config); DirectoryScanner *scanner = directory_scanner_create(config->send_directory, config->use_metadata); Chunk *current_chunk; @@ -176,37 +199,62 @@ int send_files_multithreaded(Config *config) { return 0; } -void handle_arg(char *argument_given, char *argument_to_set, bool *result, - char *message) { - if (strcmp(argument_given, argument_to_set) == 0) { - *result = true; - log_message(LOG_LEVEL_INFO, message); +static int is_remote_dest(const char *s) { + const char *colon = strchr(s, ':'); + if (!colon) return 0; + if (colon == s) return 0; + for (const char *p = s; p < colon; p++) { + if (*p == '/') return 0; } + return 1; } + +static void print_usage(void) { + printf("Usage:\n"); + printf(" fastsync [options] \n"); + printf(" fastsync [options] --source-dir --dest-dir \n"); + printf("\n"); + printf("Destination formats:\n"); + printf(" user@host:/path SSH transport (rsync-style)\n"); + printf(" host:/path SSH transport (current user)\n"); + printf(" /local/path TCP transport (requires server on localhost:8080)\n"); + printf("\n"); + printf("Options:\n"); + printf(" -c [level] Enable compression (level 1-22, default 5)\n"); + printf(" -m Enable multithreading\n"); + printf(" -s Enable chunk serialization\n"); + printf(" -f Enable sendfile (TCP only, not with -c or -s)\n"); + printf(" -M, --preserve Preserve file metadata\n"); + printf(" --source-dir Source directory\n"); + printf(" --dest-dir Destination directory\n"); + printf(" --save-to-disk Write received files to disk\n"); + printf(" --help Show this help\n"); +} + int main(int argc, char *argv[]) { const char *env_source = getenv("FASTSYNC_SOURCE_DIR"); const char *env_dest = getenv("FASTSYNC_DEST_DIR"); const char *env_save = getenv("FASTSYNC_SAVE_TO_DISK"); - char *source_dir = - env_source ? str_dup((char *)env_source) - : str_dup("/home/taptap/Nextcloud/Uni/moodle/B. Schnor: " - "Konzepte Paralleler Programmierung, SoSe 2026"); - char *dest_dir = - env_dest ? str_dup((char *)env_dest) : str_dup("./data_copied"); bool save_to_disk = false; if (env_save && (strcmp(env_save, "true") == 0 || strcmp(env_save, "1") == 0)) { save_to_disk = true; } - Config *config = config_create(str_dup("1.0.0"), source_dir, dest_dir, + Config *config = config_create(str_dup("1.0.0"), NULL, NULL, save_to_disk, false, false, false, false, 5, 20, false); + + int positional_args[2]; + int positional_count = 0; + for (int i = 1; i < argc; i++) { - if (strcmp(argv[i], "-c") == 0) { + if (strcmp(argv[i], "--help") == 0) { + print_usage(); + return 0; + } else if (strcmp(argv[i], "-c") == 0) { config->use_compression = true; log_message(LOG_LEVEL_INFO, "Enabled Compression"); - if (i + 1 < argc) { char *end_ptr; int level = strtol(argv[i + 1], &end_ptr, 10); @@ -214,6 +262,7 @@ int main(int argc, char *argv[]) { config->compression_level = level; log_message(LOG_LEVEL_INFO, "Set Compression level to %d", config->compression_level); + i++; } } } else if (strcmp(argv[i], "--source-dir") == 0 && i + 1 < argc) { @@ -230,11 +279,52 @@ int main(int argc, char *argv[]) { } else if (strcmp(argv[i], "-f") == 0 || strcmp(argv[i], "--sendfile") == 0) { config->use_sendfile = true; log_message(LOG_LEVEL_INFO, "Enabled sendfile"); + } else if (strcmp(argv[i], "-m") == 0) { + config->use_multithreading = true; + log_message(LOG_LEVEL_INFO, "Enabled Multithreading"); + } else if (strcmp(argv[i], "-s") == 0) { + config->use_chunk_serialization = true; + log_message(LOG_LEVEL_INFO, "Enabled Chunk Serialization"); + } else if (argv[i][0] == '-') { + fprintf(stderr, "Unknown option: %s\n", argv[i]); + print_usage(); + return 1; } else { - handle_arg(argv[i], "-m", &config->use_multithreading, - "Enabled Multithreading"); - handle_arg(argv[i], "-s", &config->use_chunk_serialization, - "Enabled Chunk Serialization"); + if (positional_count < 2) + positional_args[positional_count++] = i; + else { + fprintf(stderr, "Unexpected argument: %s\n", argv[i]); + print_usage(); + return 1; + } + } + } + + if (positional_count == 2) { + free(config->send_directory); + config->send_directory = str_dup(argv[positional_args[0]]); + free(config->receive_root_directory); + config->receive_root_directory = str_dup(argv[positional_args[1]]); + config->save_to_disk = true; + + if (is_remote_dest(config->receive_root_directory)) { + config->transport = TRANSPORT_SSH; + config->ssh_destination = str_dup(config->receive_root_directory); + } + } else if (positional_count == 1) { + fprintf(stderr, "Error: missing destination argument\n"); + print_usage(); + return 1; + } else { + if (!config->send_directory) { + config->send_directory = + env_source ? str_dup((char *)env_source) + : str_dup("/home/taptap/Nextcloud/Uni/moodle/B. Schnor: " + "Konzepte Paralleler Programmierung, SoSe 2026"); + } + if (!config->receive_root_directory) { + config->receive_root_directory = + env_dest ? str_dup((char *)env_dest) : str_dup("./data_copied"); } } @@ -243,6 +333,11 @@ int main(int argc, char *argv[]) { return 1; } + if (config->transport == TRANSPORT_SSH && config->use_sendfile) { + fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n"); + return 1; + } + if (config->use_multithreading) return send_files_multithreaded(config); return send_files(config); diff --git a/src/server/server.c b/src/server/server.c index 0a5fb3b..83dba1d 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -10,6 +10,7 @@ #include "utils.h" #include #include +#include #include File *file_receive(Config *config, int file_descriptor) { @@ -161,12 +162,18 @@ void handler(int file_descriptor) { thrd_join(receiver, NULL); thrd_join(writer, NULL); pipeline_context_receiver_destroy(context); + send_status(file_descriptor, STATUS_OK); } else receive_files(config, file_descriptor); close(file_descriptor); } -int main() { +int main(int argc, char *argv[]) { + if (argc > 1 && strcmp(argv[1], "--stdio") == 0) { + io_set_fds(STDIN_FILENO, STDOUT_FILENO); + handler(STDIN_FILENO); + return 0; + } Server *server = server_create(8080); server_listen(server, handler); server_delete(server); diff --git a/src/shared/config.c b/src/shared/config.c index 149b448..4db3409 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -22,6 +22,8 @@ Config *config_create(char *version, char *send_directory, config->compression_level = compression_level; config->num_connections = num_connections; config->use_sendfile = use_sendfile; + config->transport = TRANSPORT_TCP; + config->ssh_destination = NULL; return config; } @@ -29,6 +31,7 @@ void config_delete(Config *config) { free(config->version); free(config->send_directory); free(config->receive_root_directory); + free(config->ssh_destination); free(config); } @@ -63,6 +66,8 @@ Config *config_receive(int file_descriptor) { config->compression_level = receive_int(file_descriptor); config->num_connections = receive_int(file_descriptor); config->use_sendfile = receive_int(file_descriptor); + config->transport = TRANSPORT_TCP; + config->ssh_destination = NULL; send_status(file_descriptor, STATUS_OK); return config; } diff --git a/src/shared/config.h b/src/shared/config.h index 2a645d7..72e0524 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -3,6 +3,11 @@ #include +typedef enum { + TRANSPORT_TCP, + TRANSPORT_SSH +} TransportType; + typedef struct Config { char *version; char *send_directory; @@ -16,6 +21,8 @@ typedef struct Config { bool use_metadata; int compression_level; int num_connections; + TransportType transport; + char *ssh_destination; } Config; Config *config_create(char *version, char *send_directory, diff --git a/src/shared/socket.c b/src/shared/socket.c index 0d6f31e..6e3b8a7 100644 --- a/src/shared/socket.c +++ b/src/shared/socket.c @@ -1,12 +1,134 @@ #include "socket.h" #include "log.h" #include +#include #include #include #include #include +#include +#include #include +static __thread int io_read_fd = -1; +static __thread int io_write_fd = -1; + +void io_set_fds(int read_fd, int write_fd) { + io_read_fd = read_fd; + io_write_fd = write_fd; +} + +static int io_fd(int dir_fd, int file_descriptor) { + return (dir_fd != -1) ? dir_fd : file_descriptor; +} + +typedef struct { + char user[256]; + char host[256]; + char remote_path[4096]; +} RemoteDest; + +static int parse_remote_dest(const char *dest, RemoteDest *r) { + const char *colon = strchr(dest, ':'); + if (!colon) return -1; + + size_t remote_path_len = strlen(colon + 1); + if (remote_path_len >= sizeof(r->remote_path)) return -1; + memcpy(r->remote_path, colon + 1, remote_path_len + 1); + + const char *at = memchr(dest, '@', colon - dest); + if (at) { + size_t user_len = at - dest; + if (user_len >= sizeof(r->user)) return -1; + memcpy(r->user, dest, user_len); + r->user[user_len] = '\0'; + + size_t host_len = colon - at - 1; + if (host_len >= sizeof(r->host)) return -1; + memcpy(r->host, at + 1, host_len); + r->host[host_len] = '\0'; + } else { + r->user[0] = '\0'; + size_t host_len = colon - dest; + if (host_len >= sizeof(r->host)) return -1; + memcpy(r->host, dest, host_len); + r->host[host_len] = '\0'; + } + return 0; +} + +Client *client_connect_ssh(char *destination) { + RemoteDest r; + if (parse_remote_dest(destination, &r) != 0) { + fprintf(stderr, "Invalid remote destination: %s\n", destination); + exit(EXIT_FAILURE); + } + + int sv[2]; + if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) { + perror("socketpair failed"); + exit(EXIT_FAILURE); + } + + int exec_pipe[2]; + if (pipe(exec_pipe) < 0) { + perror("pipe failed"); + exit(EXIT_FAILURE); + } + + pid_t pid = fork(); + if (pid < 0) { + perror("fork failed"); + exit(EXIT_FAILURE); + } + + if (pid == 0) { + close(sv[0]); + close(exec_pipe[0]); + fcntl(exec_pipe[1], F_SETFD, FD_CLOEXEC); + + if (sv[1] != STDIN_FILENO) + dup2(sv[1], STDIN_FILENO); + if (sv[1] != STDOUT_FILENO) + dup2(sv[1], STDOUT_FILENO); + if (sv[1] > 1) close(sv[1]); + + char ssh_user[512]; + if (r.user[0] != '\0') + snprintf(ssh_user, sizeof(ssh_user), "%s@%s", r.user, r.host); + else + snprintf(ssh_user, sizeof(ssh_user), "%s", r.host); + + execlp("ssh", "ssh", "-o", "Compression=no", "-o", + "ControlMaster=no", ssh_user, "fastsync-server", "--stdio", + (char *)NULL); + perror("exec of ssh failed"); + (void)write(exec_pipe[1], "x", 1); + _exit(1); + } + + close(sv[1]); + close(exec_pipe[1]); + + char exec_status; + ssize_t n = read(exec_pipe[0], &exec_status, 1); + close(exec_pipe[0]); + + if (n > 0) { + close(sv[0]); + waitpid(pid, NULL, 0); + fprintf(stderr, "Error: could not launch 'fastsync-server --stdio' on remote\n"); + exit(EXIT_FAILURE); + } + + Client *client = malloc(sizeof(Client)); + client->file_descriptor = sv[0]; + client->address.sin_family = AF_UNIX; + client->address_length = 0; + client->ssh_child_pid = pid; + return client; +} + Server *server_create(int port) { Server *server = (Server *)malloc(sizeof(Server)); if (server == NULL) { @@ -82,6 +204,7 @@ Client *client_create() { client->file_descriptor = file_descriptor; client->address.sin_family = AF_INET; client->address_length = sizeof(client->address); + client->ssh_child_pid = -1; return client; } @@ -100,7 +223,14 @@ void client_connect(Client *client, char *host, int port) { } } -void client_disconnect(Client *client) { close(client->file_descriptor); } +void client_disconnect(Client *client) { + close(client->file_descriptor); + if (client->ssh_child_pid > 0) { + int status; + waitpid(client->ssh_child_pid, &status, 0); + client->ssh_child_pid = -1; + } +} void client_delete(Client *client) { if (client == NULL) @@ -110,12 +240,11 @@ void client_delete(Client *client) { void send_n_data(int file_descriptor, void *data, size_t data_size) { log_message(LOG_LEVEL_DEBUG, " Sending n Data: %zu", data_size); + int fd = io_fd(io_write_fd, file_descriptor); ssize_t total_bytes_send = 0; while (total_bytes_send < data_size) { - printf("Trying: %zu\n", data_size - total_bytes_send); - ssize_t bytes_send = send(file_descriptor, (char *)data + total_bytes_send, - data_size - total_bytes_send, 0); - printf("Bytes send: %zd\n", bytes_send); + ssize_t bytes_send = write(fd, (char *)data + total_bytes_send, + data_size - total_bytes_send); if (bytes_send <= 0) { perror("Could not send data!"); exit(EXIT_FAILURE); @@ -127,11 +256,12 @@ void send_n_data(int file_descriptor, void *data, size_t data_size) { void receive_n_data(int file_descriptor, void *data, size_t data_size) { log_message(LOG_LEVEL_DEBUG, " Receiving n Data: %zu", data_size); + int fd = io_fd(io_read_fd, file_descriptor); size_t total_bytes_received = 0; while (total_bytes_received < data_size) { - long long bytes_received = - recv(file_descriptor, data + total_bytes_received, - data_size - total_bytes_received, 0); + ssize_t bytes_received = + read(fd, (char *)data + total_bytes_received, + data_size - total_bytes_received); if (bytes_received == -1 || bytes_received == 0) { perror("Could not receive bytes!"); exit(EXIT_FAILURE); diff --git a/src/shared/socket.h b/src/shared/socket.h index db694f1..7973edf 100644 --- a/src/shared/socket.h +++ b/src/shared/socket.h @@ -21,13 +21,16 @@ typedef struct Client { struct sockaddr_in address; unsigned int address_length; int file_descriptor; + pid_t ssh_child_pid; } Client; Client *client_create(); void client_disconnect(Client *client); void client_delete(Client *client); void client_connect(Client *client, char *host, int port); +Client *client_connect_ssh(char *destination); +void io_set_fds(int read_fd, int write_fd); void send_n_data(int file_descriptor, void *data, size_t data_size); void receive_n_data(int file_descriptor, void *data, size_t data_size); void send_str(int file_descriptor, char *data); diff --git a/test.py b/test.py index cce5cae..6c48e8a 100755 --- a/test.py +++ b/test.py @@ -51,6 +51,7 @@ BASE_CLIENT_FLAGS = ["--save-to-disk"] TEST_CASES = [ {"name": "Standard", "flags": []}, + {"name": "Posix Args (no flags)", "flags": [], "posix": True}, {"name": "Standard (no metadata)", "flags": [], "use_metadata": False}, {"name": "Multithreading (-m)", "flags": ["-m"]}, {"name": "Compression (-c)", "flags": ["-c"]}, @@ -246,7 +247,10 @@ def run_profile(profile_name, source_dir, dest_dir): results = [] for case in TEST_CASES: flags = BASE_CLIENT_FLAGS + (["-M"] if case.get("use_metadata", True) else []) + case["flags"] - cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags + if case.get("posix"): + cmd = client_prefix + BASE_CLIENT_CMD + [source_dir, dest_dir] + flags + else: + cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags print(f"\n --- {case['name']} ---\n Running: {' '.join(cmd)}") try: r = run_single_test(cmd, case["name"], source_dir, dest_dir) @@ -344,7 +348,48 @@ def format_throughput(bps): return f"{bps:.0f} B/s" +def preflight_checks(): + errors = [] + print("Pre-flight checks:") + print(" [1] --help flag...", end=" ") + r = subprocess.run(BASE_CLIENT_CMD + ["--help"], capture_output=True, text=True) + if r.returncode == 0 and "Usage:" in r.stdout and "SSH transport" in r.stdout: + print("OK") + else: + print("FAIL") + errors.append("--help failed") + + print(" [2] Remote SSH dest detection...", end=" ") + r = subprocess.run(BASE_CLIENT_CMD + ["/x", "somehost:/y"], capture_output=True, text=True, timeout=5) + if r.returncode != 0 and ("ssh" in r.stderr or "Could not receive" in r.stderr or "could not launch" in r.stderr or "Error" in r.stderr): + print("OK (detected as SSH)") + else: + print("FAIL (not detected as SSH dest)") + errors.append("SSH detection failed") + + print(" [3] Server --stdio flag...", end=" ") + r = subprocess.run(["./build/server", "--stdio"], capture_output=True, text=True, timeout=3) + if r.returncode != 0 and ("receiving" in r.stderr or "receiving" in r.stdout or "Receiving" in r.stderr): + print("OK (started in stdio mode)") + else: + print("WARN (stdio exited: rc=%d)" % r.returncode) + + print(" [4] Posix arg syntax (no server, expect failure)...", end=" ") + r = subprocess.run(BASE_CLIENT_CMD + ["/tmp/x", "/tmp/y"], capture_output=True, text=True, timeout=5) + if r.returncode != 0 and "connect" in r.stderr: + print("OK (TCP fallback)") + else: + print("FAIL") + errors.append("Posix arg syntax failed") + + if errors: + print(f"\n {len(errors)} pre-flight check(s) failed: {', '.join(errors)}") + sys.exit(1) + print(" All pre-flight checks passed.\n") + + def main(): + preflight_checks() parser = argparse.ArgumentParser(description="FastSync integration test / benchmark") parser.add_argument("--source-dir", default=DEFAULT_SOURCE_DIR) parser.add_argument("--dest-dir", default=DEFAULT_DEST_DIR) diff --git a/tests/test_config.c b/tests/test_config.c index 664ec70..973a330 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -18,6 +18,20 @@ static void test_config_lifecycle() { EXPECT_FALSE(cfg->use_chunk_serialization); EXPECT_FALSE(cfg->use_compression); EXPECT_EQ_INT(cfg->num_connections, 4); + EXPECT_EQ_INT(cfg->transport, TRANSPORT_TCP); + EXPECT_NULL(cfg->ssh_destination); + config_delete(cfg); +} + +static void test_config_ssh_dest() { + Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("user@host:/dst"), + true, false, false, false, false, 1, 4, false); + cfg->transport = TRANSPORT_SSH; + cfg->ssh_destination = str_dup("user@host:/dst"); + EXPECT_NOT_NULL(cfg); + EXPECT_EQ_INT(cfg->transport, TRANSPORT_SSH); + EXPECT_EQ_STR(cfg->ssh_destination, "user@host:/dst"); + EXPECT_EQ_STR(cfg->receive_root_directory, "user@host:/dst"); config_delete(cfg); } @@ -55,6 +69,7 @@ static void test_pipeline_receiver_lifecycle() { void test_config() { test_config_lifecycle(); + test_config_ssh_dest(); test_pipeline_sender_lifecycle(); test_pipeline_receiver_lifecycle(); }