diff --git a/Makefile.am b/Makefile.am index 4fbb13a2..980ab551 100644 --- a/Makefile.am +++ b/Makefile.am @@ -54,6 +54,8 @@ swupd_SOURCES = \ src/lib/strings.h \ src/lib/sys.c \ src/lib/sys.h \ + src/lib/thread_pool.c \ + src/lib/thread_pool.h \ src/lock.c \ src/lib/log.c \ src/lib/log.h \ @@ -91,7 +93,8 @@ swupd_LDADD = \ $(openssl_LIBS) \ $(curl_LIBS) \ $(bsdiff_LIBS) \ - $(libarchive_LIBS) + $(libarchive_LIBS) \ + $(pthread_LIBS) verifytime_SOURCES = src/verifytime.h \ src/verifytime.c \ diff --git a/configure.ac b/configure.ac index 17355462..a574f22e 100644 --- a/configure.ac +++ b/configure.ac @@ -22,6 +22,7 @@ PKG_CHECK_MODULES([zlib], [zlib]) PKG_CHECK_MODULES([curl], [libcurl]) PKG_CHECK_MODULES([openssl], [libcrypto >= 1.0.1]) PKG_CHECK_MODULES([libarchive], [libarchive]) +AC_CHECK_LIB([pthread], [pthread_create]) # Program checks diff --git a/src/lib/thread_pool.c b/src/lib/thread_pool.c new file mode 100644 index 00000000..0055c610 --- /dev/null +++ b/src/lib/thread_pool.c @@ -0,0 +1,152 @@ +/* + * Software Updater - client side + * + * Copyright © 2018 Intel Corporation. + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation, version 2 or later of the License. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with this program. If not, see . + * + * Authors: + * Otavio Pontes + * + */ + +#define _GNU_SOURCE + +#include +#include +#include +#include +#include +#include + +#include "macros.h" +#include "thread_pool.h" + +/* + * To keep this thread pool implementation simple, pipes were used instead of + * mutex locks using the fact that read() and write() operations are thread safe. + * So in order to schedule a task to be executed we just need to write() a task + * in to the pipe and the first thread to read() that data will execute the + * processing. + */ + +struct tp { + int num_threads; + int pipe_fds[2]; + pthread_t threads[]; +}; + +struct task { + tp_task_run_t run; + void *data; +}; + +static void *thread_run(void *data) +{ + int *fd; + ssize_t r; + struct task task; + + fd = data; + + while (1) { + r = read(*fd, &task, sizeof(task)); + if (r == 0) { + // EOF: Close thread + return NULL; + } else if (r < 0 || r != sizeof(task)) { + if (errno == EINTR) { + // Not an error, we can continue + continue; + } + fprintf(stderr, "Error - Thread communication failed: %d - %s\n", + errno, strerror(errno)); + return NULL; + } + + // Run task + task.run(task.data); + } + return NULL; +} + +struct tp *tp_start(int num_threads) +{ + int i; + struct tp *tp; + + tp = calloc(1, sizeof(struct tp) + num_threads * sizeof(pthread_t *)); + ON_NULL_ABORT(tp); + + tp->num_threads = num_threads; + if (pipe2(tp->pipe_fds, O_DIRECT | O_CLOEXEC) < 0) { + goto error; + } + + // Create threads + for (i = 0; i < num_threads; i++) { + if (pthread_create(&tp->threads[i], NULL, thread_run, &tp->pipe_fds[0]) != 0) { + tp->num_threads = i; + goto error_threads; + } + } + + return tp; + +error_threads: + tp_complete(tp); + return NULL; + +error: + free(tp); + return NULL; +} + +int tp_task_schedule(struct tp *tp, tp_task_run_t run, void *data) +{ + int r = -1; + struct task task; + + task.run = run; + task.data = data; + + while (r < 0) { + r = write(tp->pipe_fds[1], &task, sizeof(task)); + if (r < 0) { + if (errno == EINTR) { + // Not an error, we can continue + continue; + } + fprintf(stderr, "Error: thread pool task scheduling failed: %d - %s", + errno, strerror(errno)); + return r; + } + } + + return 0; +} + +void tp_complete(struct tp *tp) +{ + int i; + + // Close pipe so threads will get an EOF when all tasks are completed + close(tp->pipe_fds[1]); + + for (i = 0; i < tp->num_threads; i++) { + pthread_join(tp->threads[i], NULL); + } + + close(tp->pipe_fds[0]); + free(tp); +} diff --git a/src/lib/thread_pool.h b/src/lib/thread_pool.h new file mode 100644 index 00000000..d06e7d5b --- /dev/null +++ b/src/lib/thread_pool.h @@ -0,0 +1,42 @@ +#ifndef __THREAD_POOL__ +#define __THREAD_POOL__ + +#include + +/* + * Really simple and small thread pool implementation for swupd. + */ + +struct tp; + +/* + * Thread task function type definition. + */ +typedef void (*tp_task_run_t)(void *data); + +/* + * Create a new thread pool. + * + * num_threads: The number of threads to be created. + * + * Note: Free thread pool data with tp_free() + */ +struct tp *tp_start(int num_threads); + +/* + * Schedule a task to be run in this thread pool. + * + * run: Callback to be executed in a new thread. + * data: Data informed to callback. + */ +int tp_task_schedule(struct tp *tp, tp_task_run_t run, void *data); + +/* + * Wait for all scheduled tasks to be completed, finishes all threads and + * release all memory used by the thread pool. + */ +void tp_complete(struct tp *tp); + +//TODO: Implement a tp_wait() function + +#endif