mirror of
https://github.com/clearlinux/swupd-client.git
synced 2026-10-04 15:58:22 +00:00
tp: Add an initial implementation of a Thread Pool
Simple implementation of a thread pool to be used initially for parallel downloads. Only major features needed were implemented. A tp_wait() function is to be contributed later. Signed-off-by: Otavio Pontes <otavio.pontes@intel.com> Signed-off-by: Brian J Lovin <brian.j.lovin@intel.com>
This commit is contained in:
+4
-1
@@ -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 \
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 <http://www.gnu.org/licenses/>.
|
||||
*
|
||||
* Authors:
|
||||
* Otavio Pontes <otavio.pontes@intel.com>
|
||||
*
|
||||
*/
|
||||
|
||||
#define _GNU_SOURCE
|
||||
|
||||
#include <errno.h>
|
||||
#include <fcntl.h>
|
||||
#include <pthread.h>
|
||||
#include <stdio.h>
|
||||
#include <string.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#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);
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
#ifndef __THREAD_POOL__
|
||||
#define __THREAD_POOL__
|
||||
|
||||
#include <stdlib.h>
|
||||
|
||||
/*
|
||||
* 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
|
||||
Reference in New Issue
Block a user