diff options
| author | dec05eba <dec05eba@protonmail.com> | 2026-08-05 14:56:08 +0200 |
|---|---|---|
| committer | dec05eba <dec05eba@protonmail.com> | 2026-08-05 14:56:08 +0200 |
| commit | fbc31a9561d16418226c0efa65e2a73923516dbd (patch) | |
| tree | 1c515cd13ee8af93a665d0587fb002745a01422d /src/cli | |
| parent | a522748ed2fc04e4ccad26c1da1056472cb28bc2 (diff) | |
ipc: response with video path in save recording/replay
Diffstat (limited to 'src/cli')
| -rw-r--r-- | src/cli/ipc.c | 545 | ||||
| -rw-r--r-- | src/cli/main.c | 72 |
2 files changed, 540 insertions, 77 deletions
diff --git a/src/cli/ipc.c b/src/cli/ipc.c index b9d0c51..ff4b852 100644 --- a/src/cli/ipc.c +++ b/src/cli/ipc.c @@ -17,12 +17,21 @@ #include <sys/stat.h> #include <sys/un.h> +#ifdef __linux__ +#include <sys/epoll.h> +#else +#include <sys/event.h> +#include <sys/time.h> +#endif + #define GSR_IPC_MAX_REQUEST_NAME_SIZE 64 -#define GSR_IPC_MAX_ESCAPED_ERROR_MESSAGE_SIZE (GSR_IPC_MAX_ERROR_MESSAGE_SIZE*6) -#define GSR_IPC_MAX_REPLY_SIZE (GSR_IPC_MAX_ESCAPED_ERROR_MESSAGE_SIZE + 128) -#define GSR_IPC_SEND_TIMEOUT_MILLISECONDS 1000 +#define GSR_IPC_MAX_EVENTS (2 + GSR_IPC_MAX_CLIENTS*2) +#define GSR_IPC_SHUTDOWN_SEND_TIMEOUT_MILLISECONDS 1000 #define GSR_IPC_SOCKET_MODE 0600 +#define GSR_IPC_WAKEUP_QUIT (1 << 0) +#define GSR_IPC_WAKEUP_COMPLETED_REQUEST (1 << 1) + typedef struct { int64_t id; char name[GSR_IPC_MAX_REQUEST_NAME_SIZE]; @@ -30,6 +39,12 @@ typedef struct { bool has_data; } gsr_ipc_request; +typedef struct { + int fd; + bool readable; + bool writable; +} gsr_ipc_event; + static bool string_is_only_whitespace(const char *str, size_t size) { for(size_t i = 0; i < size; ++i) { if(str[i] != ' ' && str[i] != '\t' && str[i] != '\r' && str[i] != '\n') @@ -38,7 +53,103 @@ static bool string_is_only_whitespace(const char *str, size_t size) { return true; } -static bool ipc_send_all(int fd, const char *data, size_t size) { +static bool fd_set_cloexec(int fd) { + const int flags = fcntl(fd, F_GETFD); + return flags != -1 && fcntl(fd, F_SETFD, flags | FD_CLOEXEC) != -1; +} + +static bool fd_set_nonblocking(int fd) { + const int flags = fcntl(fd, F_GETFL); + return flags != -1 && fcntl(fd, F_SETFL, flags | O_NONBLOCK) != -1; +} + +#ifdef __linux__ +static bool ipc_poller_init(gsr_ipc *self) { + self->poll_fd = epoll_create1(EPOLL_CLOEXEC); + if(self->poll_fd == -1) { + gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc: failed to create an epoll instance, error: %s", strerror(errno)); + return false; + } + return true; +} + +static bool ipc_poller_add(gsr_ipc *self, int fd) { + struct epoll_event event; + memset(&event, 0, sizeof(event)); + event.events = EPOLLIN | EPOLLET; + event.data.fd = fd; + return epoll_ctl(self->poll_fd, EPOLL_CTL_ADD, fd, &event) == 0; +} + +static bool ipc_poller_set_write_notify(gsr_ipc *self, int fd, bool enable) { + struct epoll_event event; + memset(&event, 0, sizeof(event)); + event.events = EPOLLIN | EPOLLET | (enable ? EPOLLOUT : 0); + event.data.fd = fd; + return epoll_ctl(self->poll_fd, EPOLL_CTL_MOD, fd, &event) == 0; +} + +/* Returns the number of events, or -1 on failure. Waits until at least one event is available */ +static int ipc_poller_wait(gsr_ipc *self, gsr_ipc_event *events, int events_capacity) { + struct epoll_event platform_events[GSR_IPC_MAX_EVENTS]; + if(events_capacity > GSR_IPC_MAX_EVENTS) + events_capacity = GSR_IPC_MAX_EVENTS; + + const int num_events = epoll_wait(self->poll_fd, platform_events, events_capacity, -1); + if(num_events == -1) + return errno == EINTR ? 0 : -1; + + for(int i = 0; i < num_events; ++i) { + events[i].fd = platform_events[i].data.fd; + events[i].readable = platform_events[i].events & (EPOLLIN | EPOLLHUP | EPOLLERR); + events[i].writable = platform_events[i].events & EPOLLOUT; + } + return num_events; +} +#else +static bool ipc_poller_init(gsr_ipc *self) { + self->poll_fd = kqueue(); + if(self->poll_fd == -1) { + gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc: failed to create a kqueue instance, error: %s", strerror(errno)); + return false; + } + fd_set_cloexec(self->poll_fd); + return true; +} + +static bool ipc_poller_add(gsr_ipc *self, int fd) { + struct kevent change; + EV_SET(&change, fd, EVFILT_READ, EV_ADD | EV_CLEAR, 0, 0, NULL); + return kevent(self->poll_fd, &change, 1, NULL, 0, NULL) != -1; +} + +static bool ipc_poller_set_write_notify(gsr_ipc *self, int fd, bool enable) { + struct kevent change; + EV_SET(&change, fd, EVFILT_WRITE, enable ? (EV_ADD | EV_CLEAR) : EV_DELETE, 0, 0, NULL); + return kevent(self->poll_fd, &change, 1, NULL, 0, NULL) != -1; +} + +/* Returns the number of events, or -1 on failure. Waits until at least one event is available */ +static int ipc_poller_wait(gsr_ipc *self, gsr_ipc_event *events, int events_capacity) { + struct kevent platform_events[GSR_IPC_MAX_EVENTS]; + if(events_capacity > GSR_IPC_MAX_EVENTS) + events_capacity = GSR_IPC_MAX_EVENTS; + + const int num_events = kevent(self->poll_fd, NULL, 0, platform_events, events_capacity, NULL); + if(num_events == -1) + return errno == EINTR ? 0 : -1; + + for(int i = 0; i < num_events; ++i) { + events[i].fd = platform_events[i].ident; + events[i].readable = platform_events[i].filter == EVFILT_READ; + events[i].writable = platform_events[i].filter == EVFILT_WRITE; + } + return num_events; +} +#endif + +/* Only used when the ipc thread exits, to not lose replies that haven't been fully sent yet */ +static bool ipc_send_all_blocking(int fd, const char *data, size_t size) { size_t offset = 0; while(offset < size) { const ssize_t bytes_written = send(fd, data + offset, size - offset, MSG_NOSIGNAL); @@ -56,7 +167,7 @@ static bool ipc_send_all(int fd, const char *data, size_t size) { poll_fd.events = POLLOUT; poll_fd.revents = 0; - const int poll_result = poll(&poll_fd, 1, GSR_IPC_SEND_TIMEOUT_MILLISECONDS); + const int poll_result = poll(&poll_fd, 1, GSR_IPC_SHUTDOWN_SEND_TIMEOUT_MILLISECONDS); if(poll_result == -1 && errno == EINTR) continue; @@ -71,11 +182,76 @@ static bool ipc_send_all(int fd, const char *data, size_t size) { return true; } -static bool ipc_client_send_reply(gsr_ipc_client *client, int64_t id, bool success, const char *error_message) { +static bool ipc_client_send_data(gsr_ipc *self, gsr_ipc_client *client, const char *data, size_t size) { + size_t offset = 0; + if(client->send_buffer_size == 0) { + while(offset < size) { + const ssize_t bytes_written = send(client->fd, data + offset, size - offset, MSG_NOSIGNAL); + if(bytes_written > 0) { + offset += bytes_written; + continue; + } + + if(bytes_written == -1 && errno == EINTR) + continue; + + if(bytes_written == -1 && (errno == EAGAIN || errno == EWOULDBLOCK)) + break; + + return false; + } + } + + const size_t bytes_remaining = size - offset; + if(bytes_remaining == 0) + return true; + + if(client->send_buffer_size + bytes_remaining > sizeof(client->send_buffer)) { + gsr_log(GSR_LOG_LEVEL_WARNING, "gsr_ipc: an ipc client isn't reading replies fast enough, disconnecting it"); + return false; + } + + const bool send_buffer_was_empty = client->send_buffer_size == 0; + memcpy(client->send_buffer + client->send_buffer_size, data + offset, bytes_remaining); + client->send_buffer_size += bytes_remaining; + return !send_buffer_was_empty || ipc_poller_set_write_notify(self, client->fd, true); +} + +static bool ipc_client_flush_send_buffer(gsr_ipc *self, gsr_ipc_client *client) { + if(client->send_buffer_size == 0) + return true; + + size_t offset = 0; + while(offset < client->send_buffer_size) { + const ssize_t bytes_written = send(client->fd, client->send_buffer + offset, client->send_buffer_size - offset, MSG_NOSIGNAL); + if(bytes_written > 0) { + offset += bytes_written; + continue; + } + + if(bytes_written == -1 && errno == EINTR) + continue; + + if(bytes_written == -1 && (errno == EAGAIN || errno == EWOULDBLOCK)) + break; + + return false; + } + + memmove(client->send_buffer, client->send_buffer + offset, client->send_buffer_size - offset); + client->send_buffer_size -= offset; + return client->send_buffer_size != 0 || ipc_poller_set_write_notify(self, client->fd, false); +} + +static bool ipc_client_send_reply(gsr_ipc *self, gsr_ipc_client *client, int64_t id, bool success, const char *error_message, const char *data) { char reply[GSR_IPC_MAX_REPLY_SIZE]; int reply_size = 0; - if(success) { + if(success && data) { + char escaped_data[GSR_IPC_MAX_ESCAPED_DATA_SIZE]; + gsr_json_escape_string(escaped_data, sizeof(escaped_data), data); + reply_size = snprintf(reply, sizeof(reply), "{\"id\":%" PRIi64 ",\"result\":\"ok\",\"data\":\"%s\"}\n", id, escaped_data); + } else if(success) { reply_size = snprintf(reply, sizeof(reply), "{\"id\":%" PRIi64 ",\"result\":\"ok\"}\n", id); } else { char escaped_error_message[GSR_IPC_MAX_ESCAPED_ERROR_MESSAGE_SIZE]; @@ -88,7 +264,7 @@ static bool ipc_client_send_reply(gsr_ipc_client *client, int64_t id, bool succe return false; } - return ipc_send_all(client->fd, reply, reply_size); + return ipc_client_send_data(self, client, reply, reply_size); } static bool ipc_request_parse(char *data, size_t size, gsr_ipc_request *request, char *error_message, size_t error_message_size) { @@ -165,6 +341,76 @@ static bool ipc_request_get_save_replay_seconds(const gsr_ipc_request *request, return true; } +static bool ipc_request_get_set_paused_state(const gsr_ipc_request *request, bool *paused, char *error_message, size_t error_message_size) { + if(request->has_data && request->data.type == SJ_BOOL) { + *paused = gsr_json_string_equals(&request->data, "true"); + return true; + } + + snprintf(error_message, error_message_size, "expected 'data' to be true to pause or false to unpause"); + return false; +} + +static bool ipc_request_name_to_deferred_request_type(const char *name, gsr_ipc_deferred_request_type *type) { + if(strcmp(name, "stop") == 0) { + *type = GSR_IPC_DEFERRED_REQUEST_STOP; + return true; + } + + if(strcmp(name, "save-replay") == 0) { + *type = GSR_IPC_DEFERRED_REQUEST_SAVE_REPLAY; + return true; + } + + if(strcmp(name, "stop-replay-recording") == 0) { + *type = GSR_IPC_DEFERRED_REQUEST_STOP_REPLAY_RECORDING; + return true; + } + + return false; +} + +static const char* deferred_request_already_pending_error(gsr_ipc_deferred_request_type type) { + switch(type) { + case GSR_IPC_DEFERRED_REQUEST_STOP: return "GPU Screen Recorder is already stopping"; + case GSR_IPC_DEFERRED_REQUEST_SAVE_REPLAY: return "a replay is already being saved"; + case GSR_IPC_DEFERRED_REQUEST_STOP_REPLAY_RECORDING: return "the recording is already being stopped"; + case GSR_IPC_DEFERRED_REQUEST_TYPE_COUNT: break; + } + return "the request is already being handled"; +} + +static const char* deferred_request_failed_error(gsr_ipc_deferred_request_type type) { + switch(type) { + case GSR_IPC_DEFERRED_REQUEST_STOP: return "failed to save the recording"; + case GSR_IPC_DEFERRED_REQUEST_SAVE_REPLAY: return "failed to save the replay"; + case GSR_IPC_DEFERRED_REQUEST_STOP_REPLAY_RECORDING: return "failed to save the recording"; + case GSR_IPC_DEFERRED_REQUEST_TYPE_COUNT: break; + } + return "the request failed"; +} + +static bool ipc_set_deferred_request_pending(gsr_ipc *self, gsr_ipc_deferred_request_type type, int client_fd, int64_t request_id) { + pthread_mutex_lock(&self->deferred_requests_mutex); + gsr_ipc_deferred_request *deferred_request = &self->deferred_requests[type]; + const bool was_empty = deferred_request->state == GSR_IPC_DEFERRED_REQUEST_STATE_EMPTY; + if(was_empty) { + deferred_request->state = GSR_IPC_DEFERRED_REQUEST_STATE_PENDING; + deferred_request->client_fd = client_fd; + deferred_request->request_id = request_id; + deferred_request->success = false; + deferred_request->has_filepath = false; + } + pthread_mutex_unlock(&self->deferred_requests_mutex); + return was_empty; +} + +static void ipc_clear_deferred_request(gsr_ipc *self, gsr_ipc_deferred_request_type type) { + pthread_mutex_lock(&self->deferred_requests_mutex); + self->deferred_requests[type].state = GSR_IPC_DEFERRED_REQUEST_STATE_EMPTY; + pthread_mutex_unlock(&self->deferred_requests_mutex); +} + static bool ipc_handle_request(gsr_ipc *self, const gsr_ipc_request *request, char *error_message, size_t error_message_size) { if(strcmp(request->name, "stop") == 0) return self->handlers.stop(error_message, error_message_size, self->handlers.userdata); @@ -172,9 +418,23 @@ static bool ipc_handle_request(gsr_ipc *self, const gsr_ipc_request *request, ch if(strcmp(request->name, "toggle-pause") == 0) return self->handlers.toggle_pause(error_message, error_message_size, self->handlers.userdata); + if(strcmp(request->name, "set-paused") == 0) { + bool paused = false; + if(!ipc_request_get_set_paused_state(request, &paused, error_message, error_message_size)) + return false; + + return self->handlers.set_paused(paused, error_message, error_message_size, self->handlers.userdata); + } + if(strcmp(request->name, "toggle-replay-recording") == 0) return self->handlers.toggle_replay_recording(error_message, error_message_size, self->handlers.userdata); + if(strcmp(request->name, "start-replay-recording") == 0) + return self->handlers.start_replay_recording(error_message, error_message_size, self->handlers.userdata); + + if(strcmp(request->name, "stop-replay-recording") == 0) + return self->handlers.stop_replay_recording(error_message, error_message_size, self->handlers.userdata); + if(strcmp(request->name, "save-replay") == 0) { int seconds = GSR_SAVE_REPLAY_SECONDS_FULL; if(!ipc_request_get_save_replay_seconds(request, &seconds, error_message, error_message_size)) @@ -189,22 +449,35 @@ static bool ipc_handle_request(gsr_ipc *self, const gsr_ipc_request *request, ch static bool ipc_client_on_request(gsr_ipc *self, gsr_ipc_client *client) { if(client->request_too_large) - return ipc_client_send_reply(client, 0, false, "the request is too large"); + return ipc_client_send_reply(self, client, 0, false, "the request is too large", NULL); if(string_is_only_whitespace(client->request, client->request_size)) - return ipc_client_send_reply(client, 0, false, "the request is empty"); + return ipc_client_send_reply(self, client, 0, false, "the request is empty", NULL); char error_message[GSR_IPC_MAX_ERROR_MESSAGE_SIZE]; error_message[0] = '\0'; gsr_ipc_request request; if(!ipc_request_parse(client->request, client->request_size, &request, error_message, sizeof(error_message))) - return ipc_client_send_reply(client, request.id, false, error_message); + return ipc_client_send_reply(self, client, request.id, false, error_message, NULL); - if(!ipc_handle_request(self, &request, error_message, sizeof(error_message))) - return ipc_client_send_reply(client, request.id, false, error_message); + /* The pending deferred request has to be registered before the handler starts the operation, + otherwise the operation could finish before the reply to it gets registered */ + gsr_ipc_deferred_request_type deferred_request_type; + const bool reply_is_deferred = ipc_request_name_to_deferred_request_type(request.name, &deferred_request_type); + if(reply_is_deferred && !ipc_set_deferred_request_pending(self, deferred_request_type, client->fd, request.id)) + return ipc_client_send_reply(self, client, request.id, false, deferred_request_already_pending_error(deferred_request_type), NULL); - return ipc_client_send_reply(client, request.id, true, NULL); + if(!ipc_handle_request(self, &request, error_message, sizeof(error_message))) { + if(reply_is_deferred) + ipc_clear_deferred_request(self, deferred_request_type); + return ipc_client_send_reply(self, client, request.id, false, error_message, NULL); + } + + if(reply_is_deferred) + return true; + + return ipc_client_send_reply(self, client, request.id, true, NULL, NULL); } static bool ipc_client_on_byte(gsr_ipc *self, gsr_ipc_client *client, char c) { @@ -243,16 +516,6 @@ static bool ipc_client_receive(gsr_ipc *self, gsr_ipc_client *client) { } } -static bool fd_set_cloexec(int fd) { - const int flags = fcntl(fd, F_GETFD); - return flags != -1 && fcntl(fd, F_SETFD, flags | FD_CLOEXEC) != -1; -} - -static bool fd_set_nonblocking(int fd) { - const int flags = fcntl(fd, F_GETFL); - return flags != -1 && fcntl(fd, F_SETFL, flags | O_NONBLOCK) != -1; -} - static void ipc_add_client(gsr_ipc *self, int client_fd) { if(self->num_clients == GSR_IPC_MAX_CLIENTS) { gsr_log(GSR_LOG_LEVEL_WARNING, "gsr_ipc: too many ipc clients are connected, rejecting the new connection"); @@ -260,7 +523,7 @@ static void ipc_add_client(gsr_ipc *self, int client_fd) { return; } - if(!fd_set_cloexec(client_fd) || !fd_set_nonblocking(client_fd)) { + if(!fd_set_cloexec(client_fd) || !fd_set_nonblocking(client_fd) || !ipc_poller_add(self, client_fd)) { gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc: failed to setup the ipc client socket, error: %s", strerror(errno)); close(client_fd); return; @@ -270,70 +533,175 @@ static void ipc_add_client(gsr_ipc *self, int client_fd) { client->fd = client_fd; client->request_size = 0; client->request_too_large = false; + client->send_buffer_size = 0; ++self->num_clients; } -static void ipc_accept_client(gsr_ipc *self) { - const int client_fd = accept(self->socket_fd, NULL, NULL); - if(client_fd == -1) - return; +static void ipc_accept_clients(gsr_ipc *self) { + for(;;) { + const int client_fd = accept(self->socket_fd, NULL, NULL); + if(client_fd == -1) { + if(errno == EINTR) + continue; + return; + } - ipc_add_client(self, client_fd); + ipc_add_client(self, client_fd); + } } static void ipc_remove_client(gsr_ipc *self, int index) { - close(self->clients[index].fd); + const int client_fd = self->clients[index].fd; + + pthread_mutex_lock(&self->deferred_requests_mutex); + for(int i = 0; i < GSR_IPC_DEFERRED_REQUEST_TYPE_COUNT; ++i) { + if(self->deferred_requests[i].state != GSR_IPC_DEFERRED_REQUEST_STATE_EMPTY && self->deferred_requests[i].client_fd == client_fd) + self->deferred_requests[i].state = GSR_IPC_DEFERRED_REQUEST_STATE_EMPTY; + } + pthread_mutex_unlock(&self->deferred_requests_mutex); + + close(client_fd); for(int i = index; i < self->num_clients - 1; ++i) { self->clients[i] = self->clients[i + 1]; } --self->num_clients; } -static void* ipc_thread(void *userdata) { - gsr_ipc *self = userdata; - struct pollfd poll_fds[2 + GSR_IPC_MAX_CLIENTS]; +static int ipc_find_client_index_by_fd(const gsr_ipc *self, int fd) { + for(int i = 0; i < self->num_clients; ++i) { + if(self->clients[i].fd == fd) + return i; + } + return -1; +} + +static void ipc_send_completed_request_replies(gsr_ipc *self) { + for(int i = 0; i < GSR_IPC_DEFERRED_REQUEST_TYPE_COUNT; ++i) { + pthread_mutex_lock(&self->deferred_requests_mutex); + const gsr_ipc_deferred_request deferred_request = self->deferred_requests[i]; + if(deferred_request.state == GSR_IPC_DEFERRED_REQUEST_STATE_COMPLETED) + self->deferred_requests[i].state = GSR_IPC_DEFERRED_REQUEST_STATE_EMPTY; + pthread_mutex_unlock(&self->deferred_requests_mutex); + + if(deferred_request.state != GSR_IPC_DEFERRED_REQUEST_STATE_COMPLETED) + continue; + + const int client_index = ipc_find_client_index_by_fd(self, deferred_request.client_fd); + if(client_index == -1) + continue; + + const char *error_message = deferred_request_failed_error(i); + const char *filepath = deferred_request.has_filepath ? deferred_request.filepath : NULL; + if(!ipc_client_send_reply(self, &self->clients[client_index], deferred_request.request_id, deferred_request.success, error_message, filepath)) + ipc_remove_client(self, client_index); + } +} + +static void ipc_fail_pending_requests(gsr_ipc *self) { + for(int i = 0; i < GSR_IPC_DEFERRED_REQUEST_TYPE_COUNT; ++i) { + pthread_mutex_lock(&self->deferred_requests_mutex); + const gsr_ipc_deferred_request deferred_request = self->deferred_requests[i]; + self->deferred_requests[i].state = GSR_IPC_DEFERRED_REQUEST_STATE_EMPTY; + pthread_mutex_unlock(&self->deferred_requests_mutex); + + if(deferred_request.state != GSR_IPC_DEFERRED_REQUEST_STATE_PENDING) + continue; + + const int client_index = ipc_find_client_index_by_fd(self, deferred_request.client_fd); + if(client_index == -1) + continue; + + if(!ipc_client_send_reply(self, &self->clients[client_index], deferred_request.request_id, false, "GPU Screen Recorder exited before the request finished", NULL)) + ipc_remove_client(self, client_index); + } +} + +static void ipc_flush_clients_blocking(gsr_ipc *self) { + for(int i = 0; i < self->num_clients; ++i) { + gsr_ipc_client *client = &self->clients[i]; + if(client->send_buffer_size > 0) + ipc_send_all_blocking(client->fd, client->send_buffer, client->send_buffer_size); + client->send_buffer_size = 0; + } +} +static int ipc_drain_wakeup_pipe(gsr_ipc *self) { + int wakeup_flags = 0; for(;;) { - poll_fds[0].fd = self->wakeup_pipe[0]; - poll_fds[0].events = POLLIN; - poll_fds[0].revents = 0; - poll_fds[1].fd = self->socket_fd; - poll_fds[1].events = POLLIN; - poll_fds[1].revents = 0; - - const int num_polled_clients = self->num_clients; - for(int i = 0; i < num_polled_clients; ++i) { - poll_fds[2 + i].fd = self->clients[i].fd; - poll_fds[2 + i].events = POLLIN; - poll_fds[2 + i].revents = 0; + char buffer[64]; + const ssize_t bytes_read = read(self->wakeup_pipe[0], buffer, sizeof(buffer)); + if(bytes_read == -1 && errno == EINTR) + continue; + + if(bytes_read <= 0) + break; + + for(ssize_t i = 0; i < bytes_read; ++i) { + if(buffer[i] == 'q') + wakeup_flags |= GSR_IPC_WAKEUP_QUIT; + else if(buffer[i] == 'c') + wakeup_flags |= GSR_IPC_WAKEUP_COMPLETED_REQUEST; } + } + return wakeup_flags; +} - if(poll(poll_fds, 2 + num_polled_clients, -1) == -1) { - if(errno == EINTR) - continue; +static void ipc_wakeup_thread(gsr_ipc *self, char wakeup_value) { + ssize_t bytes_written = 0; + do { + bytes_written = write(self->wakeup_pipe[1], &wakeup_value, 1); + } while(bytes_written == -1 && errno == EINTR); - gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc: failed to poll the ipc sockets, error: %s", strerror(errno)); + if(bytes_written == -1) + gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc: failed to wake up the ipc thread, error: %s", strerror(errno)); +} + +static void* ipc_thread(void *userdata) { + gsr_ipc *self = userdata; + gsr_ipc_event events[GSR_IPC_MAX_EVENTS]; + bool running = true; + + while(running) { + const int num_events = ipc_poller_wait(self, events, GSR_IPC_MAX_EVENTS); + if(num_events == -1) { + gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc: failed to wait for ipc events, error: %s", strerror(errno)); break; } - if(poll_fds[0].revents != 0) - break; + for(int i = 0; i < num_events; ++i) { + if(events[i].fd == self->wakeup_pipe[0]) { + const int wakeup_flags = ipc_drain_wakeup_pipe(self); + if(wakeup_flags & GSR_IPC_WAKEUP_COMPLETED_REQUEST) + ipc_send_completed_request_replies(self); + if(wakeup_flags & GSR_IPC_WAKEUP_QUIT) + running = false; + continue; + } - if(poll_fds[1].revents & POLLIN) - ipc_accept_client(self); + if(events[i].fd == self->socket_fd) { + ipc_accept_clients(self); + continue; + } + + const int client_index = ipc_find_client_index_by_fd(self, events[i].fd); + if(client_index == -1) + continue; - for(int i = num_polled_clients - 1; i >= 0; --i) { + gsr_ipc_client *client = &self->clients[client_index]; bool keep_client = true; - if(poll_fds[2 + i].revents & POLLIN) - keep_client = ipc_client_receive(self, &self->clients[i]); - else if(poll_fds[2 + i].revents & (POLLHUP | POLLERR | POLLNVAL)) - keep_client = false; + if(events[i].readable) + keep_client = ipc_client_receive(self, client); + if(keep_client && events[i].writable) + keep_client = ipc_client_flush_send_buffer(self, client); if(!keep_client) - ipc_remove_client(self, i); + ipc_remove_client(self, client_index); } } + ipc_send_completed_request_replies(self); + ipc_fail_pending_requests(self); + ipc_flush_clients_blocking(self); return NULL; } @@ -403,6 +771,11 @@ static void ipc_close(gsr_ipc *self) { } } + if(self->poll_fd != -1) { + close(self->poll_fd); + self->poll_fd = -1; + } + if(self->socket_fd != -1) { close(self->socket_fd); self->socket_fd = -1; @@ -412,11 +785,17 @@ static void ipc_close(gsr_ipc *self) { unlink(self->socket_filepath); self->socket_bound = false; } + + if(self->deferred_requests_mutex_created) { + pthread_mutex_destroy(&self->deferred_requests_mutex); + self->deferred_requests_mutex_created = false; + } } int gsr_ipc_init(gsr_ipc *self, const char *socket_filepath) { memset(self, 0, sizeof(*self)); self->socket_fd = -1; + self->poll_fd = -1; self->wakeup_pipe[0] = -1; self->wakeup_pipe[1] = -1; @@ -430,6 +809,12 @@ int gsr_ipc_init(gsr_ipc *self, const char *socket_filepath) { snprintf(self->socket_filepath, sizeof(self->socket_filepath), "%s", socket_filepath); + if(pthread_mutex_init(&self->deferred_requests_mutex, NULL) != 0) { + gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc_init: failed to create the deferred requests mutex"); + goto err; + } + self->deferred_requests_mutex_created = true; + if(pipe(self->wakeup_pipe) == -1) { gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc_init: failed to create the ipc wakeup pipe, error: %s", strerror(errno)); self->wakeup_pipe[0] = -1; @@ -437,7 +822,7 @@ int gsr_ipc_init(gsr_ipc *self, const char *socket_filepath) { goto err; } - if(!fd_set_cloexec(self->wakeup_pipe[0]) || !fd_set_cloexec(self->wakeup_pipe[1])) { + if(!fd_set_cloexec(self->wakeup_pipe[0]) || !fd_set_cloexec(self->wakeup_pipe[1]) || !fd_set_nonblocking(self->wakeup_pipe[0])) { gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc_init: failed to setup the ipc wakeup pipe, error: %s", strerror(errno)); goto err; } @@ -448,6 +833,14 @@ int gsr_ipc_init(gsr_ipc *self, const char *socket_filepath) { goto err; } + if(!ipc_poller_init(self)) + goto err; + + if(!ipc_poller_add(self, self->wakeup_pipe[0]) || !ipc_poller_add(self, self->socket_fd)) { + gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc_init: failed to register the ipc sockets for events, error: %s", strerror(errno)); + goto err; + } + if(!ipc_bind(self, &addr)) goto err; @@ -500,15 +893,27 @@ void gsr_ipc_stop(gsr_ipc *self) { if(!self->thread_running) return; - const char wakeup_value = 1; - ssize_t bytes_written = 0; - do { - bytes_written = write(self->wakeup_pipe[1], &wakeup_value, 1); - } while(bytes_written == -1 && errno == EINTR); - - if(bytes_written == -1) - gsr_log(GSR_LOG_LEVEL_ERROR, "gsr_ipc_stop: failed to wake up the ipc thread, error: %s", strerror(errno)); - + ipc_wakeup_thread(self, 'q'); pthread_join(self->thread, NULL); self->thread_running = false; } + +void gsr_ipc_complete_request(gsr_ipc *self, gsr_ipc_deferred_request_type type, bool success, const char *filepath) { + if(!self->initialized) + return; + + pthread_mutex_lock(&self->deferred_requests_mutex); + gsr_ipc_deferred_request *deferred_request = &self->deferred_requests[type]; + const bool was_pending = deferred_request->state == GSR_IPC_DEFERRED_REQUEST_STATE_PENDING; + if(was_pending) { + deferred_request->state = GSR_IPC_DEFERRED_REQUEST_STATE_COMPLETED; + deferred_request->success = success; + deferred_request->has_filepath = filepath != NULL; + if(filepath) + snprintf(deferred_request->filepath, sizeof(deferred_request->filepath), "%s", filepath); + } + pthread_mutex_unlock(&self->deferred_requests_mutex); + + if(was_pending) + ipc_wakeup_thread(self, 'c'); +} diff --git a/src/cli/main.c b/src/cli/main.c index d34e8cc..de014a7 100644 --- a/src/cli/main.c +++ b/src/cli/main.c @@ -137,6 +137,17 @@ static bool ipc_toggle_pause_handler(char *error_message, size_t error_message_s return true; } +static bool ipc_set_paused_handler(bool paused, char *error_message, size_t error_message_size, void *userdata) { + const gsr_recorder_settings *settings = userdata; + if(settings->is_replaying) { + snprintf(error_message, error_message_size, "pausing is not supported when recording a replay"); + return false; + } + + gsr_recorder_set_paused(recorder, paused); + return true; +} + static bool ipc_toggle_replay_recording_handler(char *error_message, size_t error_message_size, void *userdata) { const gsr_recorder_settings *settings = userdata; if(!settings->replay_recording_directory) { @@ -148,6 +159,33 @@ static bool ipc_toggle_replay_recording_handler(char *error_message, size_t erro return true; } +static bool ipc_start_replay_recording_handler(char *error_message, size_t error_message_size, void *userdata) { + const gsr_recorder_settings *settings = userdata; + if(!settings->replay_recording_directory) { + snprintf(error_message, error_message_size, "option -ro is required to start a recording"); + return false; + } + + gsr_recorder_start_replay_recording(recorder); + return true; +} + +static bool ipc_stop_replay_recording_handler(char *error_message, size_t error_message_size, void *userdata) { + const gsr_recorder_settings *settings = userdata; + if(!settings->replay_recording_directory) { + snprintf(error_message, error_message_size, "option -ro is required to start a recording"); + return false; + } + + if(!gsr_recorder_is_replay_recording(recorder)) { + snprintf(error_message, error_message_size, "no recording is running"); + return false; + } + + gsr_recorder_stop_replay_recording(recorder); + return true; +} + static bool ipc_save_replay_handler(int seconds, char *error_message, size_t error_message_size, void *userdata) { const gsr_recorder_settings *settings = userdata; if(!settings->is_replaying) { @@ -295,18 +333,26 @@ static void screenshot_saved_callback(const char *filepath, void *userdata) { run_recording_saved_script_async(recording_saved_script, filepath, "screenshot"); } +typedef struct { + const char *recording_saved_script; + gsr_ipc *ipc; +} recorder_callbacks_context; + static void replay_saved_callback(const char *filepath, void *userdata) { - const char *recording_saved_script = userdata; + recorder_callbacks_context *context = userdata; if(!filepath) { printf("gsr error: Failed to save replay\n"); fflush(stdout); + gsr_ipc_complete_request(context->ipc, GSR_IPC_DEFERRED_REQUEST_SAVE_REPLAY, false, NULL); return; } puts(filepath); fflush(stdout); - if(recording_saved_script) - run_recording_saved_script_async(recording_saved_script, filepath, "replay"); + if(context->recording_saved_script) + run_recording_saved_script_async(context->recording_saved_script, filepath, "replay"); + + gsr_ipc_complete_request(context->ipc, GSR_IPC_DEFERRED_REQUEST_SAVE_REPLAY, true, filepath); } static void recording_started_callback(const char *filepath, void *userdata) { @@ -318,17 +364,20 @@ static void recording_started_callback(const char *filepath, void *userdata) { } static void recording_stopped_callback(const char *filepath, void *userdata) { - const char *recording_saved_script = userdata; + recorder_callbacks_context *context = userdata; if(!filepath) { printf("gsr error: Failed to save recording\n"); fflush(stdout); + gsr_ipc_complete_request(context->ipc, GSR_IPC_DEFERRED_REQUEST_STOP_REPLAY_RECORDING, false, NULL); return; } puts(filepath); fflush(stdout); - if(recording_saved_script) - run_recording_saved_script_async(recording_saved_script, filepath, "regular"); + if(context->recording_saved_script) + run_recording_saved_script_async(context->recording_saved_script, filepath, "regular"); + + gsr_ipc_complete_request(context->ipc, GSR_IPC_DEFERRED_REQUEST_STOP_REPLAY_RECORDING, true, filepath); } #ifdef GSR_APP_AUDIO @@ -410,12 +459,16 @@ static int record(args_parser *arg_parser, gsr_windowing *windowing, gsr_capture recorder_params.pipewire_audio = &pipewire_audio; #endif + recorder_callbacks_context callbacks_context; + callbacks_context.recording_saved_script = arg_parser->settings.recording_saved_script; + callbacks_context.ipc = ipc; + gsr_recorder_callbacks callbacks; memset(&callbacks, 0, sizeof(callbacks)); callbacks.replay_saved = replay_saved_callback; callbacks.recording_started = recording_started_callback; callbacks.recording_stopped = recording_stopped_callback; - callbacks.userdata = (void*)arg_parser->settings.recording_saved_script; + callbacks.userdata = &callbacks_context; int error = GSR_ERROR_OK; recorder = gsr_recorder_create(&recorder_params, &callbacks, &error); @@ -427,9 +480,13 @@ static int record(args_parser *arg_parser, gsr_windowing *windowing, gsr_capture gsr_recorder_stop(recorder); gsr_ipc_handlers ipc_handlers; + memset(&ipc_handlers, 0, sizeof(ipc_handlers)); ipc_handlers.stop = ipc_stop_handler; ipc_handlers.toggle_pause = ipc_toggle_pause_handler; + ipc_handlers.set_paused = ipc_set_paused_handler; ipc_handlers.toggle_replay_recording = ipc_toggle_replay_recording_handler; + ipc_handlers.start_replay_recording = ipc_start_replay_recording_handler; + ipc_handlers.stop_replay_recording = ipc_stop_replay_recording_handler; ipc_handlers.save_replay = ipc_save_replay_handler; ipc_handlers.userdata = &arg_parser->settings; @@ -437,6 +494,7 @@ static int record(args_parser *arg_parser, gsr_windowing *windowing, gsr_capture if(run_result == GSR_ERROR_OK) run_result = gsr_recorder_run(recorder); + gsr_ipc_complete_request(ipc, GSR_IPC_DEFERRED_REQUEST_STOP, true, arg_parser->settings.is_replaying ? NULL : arg_parser->settings.filename); gsr_ipc_stop(ipc); gsr_recorder_destroy(recorder); recorder = NULL; |
