Skip to content
Merged

Dev #10

Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 33 additions & 0 deletions .github/ISSUE_TEMPLATE/bug_report.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
---
name: Bug report
about: Create a report of the bug you've encountered
title: "[BUG] - ISSUE_TITLE"
labels: bug
assignees: cartepro

---

## Description

A clear and concise description of what the bug is.

## Steps to reproduce

Steps to reproduce the bug you encountered

## Expected behaviour

A clear and concise description of what you expected to happen.

## Environment

- OS: [e.g. GNU/Linux Debian 10]
- Architecture [Arm, x86_64, ...]
- Rust version
- library version
- Protocol used
- Remote server version and name

## Additional information

Add any other context about the problem here.
23 changes: 23 additions & 0 deletions .github/ISSUE_TEMPLATE/feature_request.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
---
name: Feature request
about: Suggest an idea to improve zappy
title: "[Feature Request] - FEATURE_TITLE"
labels: "new feature"
assignees: cartepro

---

## Description

Put here a brief introduction to your suggestion.

### Changes

The following changes to the application are expected

- ...

## Implementation

Provide any kind of suggestion you propose on how to implement the feature.
If you have none, delete this section.
8 changes: 8 additions & 0 deletions .github/ISSUE_TEMPLATE/question.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
---
name: Question
about: Ask what you want about the project
title: "[QUESTION] - TITLE"
labels: question
assignees: cartepro

---
17 changes: 17 additions & 0 deletions .github/pull_request_template.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
## Description of the change.

## Related Issue

Closes # ?

## Type of Change

- [ ] Bug fix
- [ ] Refactor
- [ ] Performance improvement
- [ ] New feature
- [ ] Tests
- [ ] CI/CD
- [ ] Documentation

## Tests
28 changes: 28 additions & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
cmake_minimum_required(VERSION 3.15)
project(kronkflow LANGUAGES C)

option(BUILD_TESTS "Build the prophecy test suite (requires kronklab)" OFF)

include(cmake/ProjectOptions.cmake)
include(cmake/CompilerOptions.cmake)
include(cmake/Utils.cmake)
Expand All @@ -9,6 +11,16 @@ if(CMAKE_CURRENT_SOURCE_DIR STREQUAL CMAKE_CURRENT_BINARY_DIR)
message(FATAL_ERROR "Les builds dans le dossier source sont interdits. Créez un dossier de build (ex: 'build/').")
endif()

include(FetchContent)

FetchContent_Declare(
kronkutils_c
GIT_REPOSITORY https://github.com/kronkorp/kronkutils-C
GIT_TAG main
)

FetchContent_MakeAvailable(kronkutils_c)

file(GLOB_RECURSE kronkflow_SOURCES
"${CMAKE_CURRENT_SOURCE_DIR}/src/*.c"
)
Expand All @@ -18,6 +30,8 @@ set(CMAKE_VISIBILITY_INLINES_HIDDEN ON)

add_library(kronkflow-obj OBJECT ${kronkflow_SOURCES})

target_link_libraries(kronkflow-obj PUBLIC kronkutils_static)

target_include_directories(kronkflow-obj
PUBLIC
$<BUILD_INTERFACE:${PROJECT_SOURCE_DIR}/include>
Expand Down Expand Up @@ -64,3 +78,17 @@ install(TARGETS kronkflow-static kronkflow-shared
install(DIRECTORY ${CMAKE_CURRENT_SOURCE_DIR}/include/kronkflow
DESTINATION ${CMAKE_INSTALL_INCLUDEDIR}
)

# NOTE: Tests (opt-in, via -DBUILD_TESTS=ON)

if(BUILD_TESTS)
FetchContent_Declare(
kronklab
GIT_REPOSITORY https://github.com/kronkorp/kronklab.git
GIT_TAG 5584149a3c9ccbb12baaf71cfd591db3b37541ad
)
FetchContent_MakeAvailable(kronklab)

enable_testing()
add_subdirectory(tests)
endif()
1 change: 1 addition & 0 deletions include/kronkflow/macros/types.h
Original file line number Diff line number Diff line change
Expand Up @@ -18,5 +18,6 @@
typedef void (*kfClearer)(void *);
typedef uint64_t kfTick;
typedef size_t kfTaskID;
typedef size_t kfStageId;

#endif /* PROPHECY_MACROS_TYPES_H */
28 changes: 27 additions & 1 deletion include/kronkflow/task.h
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
#define PROPHECY_TASKS_H
#include "kronkflow/macros/optimization.h"
#include "kronkflow/macros/types.h"
#include <stdint.h>

///////////////////////////////////////////////////////////////////////////////
/**
Expand All @@ -18,6 +19,27 @@ typedef struct prophecy_task_s kfTask;
///////////////////////////////////////////////////////////////////////////////


///////////////////////////////////////////////////////////////////////////////
/**
* @struct prophecy_read_write_masks_s
*
* @brief Read mask and write mask for multithreading
*
* @note We're gonna compare tasks in the same stage to find out if we can
* run them in the same threadpool session.
* If read mask & read mask: yes. If read mask & write mask or write mask
* & write mask: no.
*/
///////////////////////////////////////////////////////////////////////////////
typedef struct kronkflow_read_write_masks_s {

uint64_t rmask; //!< The read mask
uint64_t wmask; //!< The write mask

} kfRWMasks;
///////////////////////////////////////////////////////////////////////////////


///////////////////////////////////////////////////////////////////////////////
/**
* @struct prophecy_task_opt_s
Expand All @@ -30,6 +52,8 @@ typedef struct prophecy_task_opt_s {
kfHandler handler; //!< The task handler (function ptr)
void* data; //!< The data to give to the handler
kfClearer clearer; //!< The clearer of the data if allocated
kfStageId stage; //!< The stage ID
kfRWMasks masks; //!< The RW Masks

} kfTaskOpt;
///////////////////////////////////////////////////////////////////////////////
Expand All @@ -42,9 +66,11 @@ typedef struct prophecy_task_opt_s {
* @param handler The function ptr handler kfBool (*)(void *, void *)
* @param data The data to give to the handler
* @param clearer The clearer to clear data if allocated
* @param stageId The stage id (priority) (< is better)
* @param rwmask The read write mask for multithreading
*/
///////////////////////////////////////////////////////////////////////////////
KF_API kfTaskOpt kfTask_opt(kfHandler handler, void *data, kfClearer clearer);
KF_API kfTaskOpt kfTask_opt(kfHandler handler, void *data, kfClearer clearer, kfStageId stageId, kfRWMasks rwmask);
///////////////////////////////////////////////////////////////////////////////


Expand Down
7 changes: 7 additions & 0 deletions src/scheduler/clear/scheduler_clear.c
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
** Clear a scheduler
*/
#include "../scheduler.h"
#include "dynarray.h"
#include <stddef.h>
#include <stdlib.h>

Expand All @@ -23,6 +24,12 @@ void kfScheduler_clear(
if (sch->tasks) {
free(sch->tasks);
}
if (sch->staged) {
for (size_t i = 0; i < kuDynarray_getLoad(sch->staged); ++i) {
kuDynarray_free(sch->staged[i]);
}
kuDynarray_free(sch->staged);
}
sch->count = 0;
sch->size = 0;
sch->tick = 0;
Expand Down
3 changes: 3 additions & 0 deletions src/scheduler/init/scheduler_init.c
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,11 @@
** File description:
** Init the scheduler
*/
#include "dynarray.h"
#include "kronkflow/scheduler.h"
#include "../../task/task.h"
#include "../scheduler.h"
#include "kronkflow/task.h"
#include <stdlib.h>

int kfScheduler_init(
Expand All @@ -24,5 +26,6 @@ int kfScheduler_init(
sch->size = size;
sch->count = 0;
sch->tick = 0;
kuDynarray_init(&sch->staged, 2, kfTask *);
return 0;
}
12 changes: 12 additions & 0 deletions src/scheduler/scheduler.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,16 @@
#ifndef PROPHECY_SCHEDULER_IMPL_H
#define PROPHECY_SCHEDULER_IMPL_H
#include "../task/task.h"
#include "kronkflow/task.h"
#include <stddef.h>
#include <dynarray.h>

typedef struct prophecy_stage_data_s {

kfTask* tasks;
bool passed;

} kfStageData;

///////////////////////////////////////////////////////////////////////////////
/**
Expand All @@ -23,6 +32,9 @@ typedef struct prophecy_scheduler_s {
size_t count; //!< The number of tasks pushed
kfTick tick; //!< The current tick (please tick scheduler at each loop)

// TODO: Can move some tasks to arena. Maybe do it with a dynamic array ?
kfTask **staged; //!< Bucket (vector)

} kfScheduler;
///////////////////////////////////////////////////////////////////////////////

Expand Down
59 changes: 51 additions & 8 deletions src/scheduler/scheduler_update.c
Original file line number Diff line number Diff line change
Expand Up @@ -4,14 +4,42 @@
** File description:
** Update the scheduler
*/
#include "dynarray.h"
#include "kronkflow/macros/optimization.h"
#include "kronkflow/macros/types.h"
#include "kronkflow/task.h"
#include "scheduler.h"
#include <stddef.h>
#include <string.h>
#include "../task/task.h"
#include "../minheap/minheap.h"
#include "kronkflow/scheduler.h"

static void push_to_buckets(
kfScheduler *sch,
kfTask *task
)
{
kfStageId stage = task->stage;
kuDynarrayHeader *header;
size_t oldLoad;

if (stage >= kuDynarray_getSize(sch->staged)) {
kuDynarray_resize(sch->staged, stage + 1);
}
header = kuDynarray_getHeader(sch->staged);
if (stage >= header->load) {
oldLoad = header->load;
memset(sch->staged + oldLoad, 0,
(stage + 1 - oldLoad) * sizeof(*sch->staged));
header->load = stage + 1;
}
if (sch->staged[stage] == NULL) {
kuDynarray_init(&sch->staged[stage], 2, kfTask);
}
kuDynarray_pushBack(sch->staged[stage], *task);
}

KF_API
size_t kfScheduler_tick(
kfScheduler *sch,
Expand All @@ -25,18 +53,33 @@ size_t kfScheduler_tick(
if (!sch) {
return 0;
}
sch->tick++;
++sch->tick;
while (sch->count > 0 && sch->tasks[0].target <= sch->tick) {
ctask = sch->tasks[0];
prMinHeap_remove(sch, 0);
r = ctask.handler(context, ctask.data);
done++;
if (r && ctask.interval > 0) {
ctask.target = ctask.interval;
kfScheduler_addTask(sch, (kfTaskOpt){ctask.handler, ctask.data, ctask.clearer}, ctask.interval, ctask.interval);
} else if (ctask.clearer) {
ctask.clearer(ctask.data);

push_to_buckets(sch, &ctask);
}

for (size_t i = 0; i < kuDynarray_getLoad(sch->staged); ++i) {
for (size_t j = 0; j < kuDynarray_getLoad(sch->staged[i]); ++j) {
ctask = sch->staged[i][j];
r = ctask.handler(context, ctask.data);
++done;
if (r && ctask.interval > 0) {
ctask.target = ctask.interval;
kfScheduler_addTask(sch, (kfTaskOpt){
ctask.handler,
ctask.data,
ctask.clearer,
ctask.stage,
ctask.masks
}, ctask.interval, ctask.interval);
} else if (ctask.clearer) {
ctask.clearer(ctask.data);
}
}
kuDynarray_clear(sch->staged[i]);
}
return done;
}
6 changes: 4 additions & 2 deletions src/task/task.h
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,9 @@
#define PROPHECY_TASK_IMPL_H
#include "kronkflow/macros/types.h"
#include "kronkflow/task.h"
#include <stdint.h>
#include <stddef.h>



///////////////////////////////////////////////////////////////////////////////
/**
* @struct prophecy_task_s
Expand All @@ -28,6 +27,9 @@ typedef struct prophecy_task_s {
kfTick interval; //!< The interval (0 if ponctual, > 0 else)
kfTick target; //!< The tick remainings.

kfStageId stage; //!< Stage id
kfRWMasks masks; //!< Masks

} kfTask;
///////////////////////////////////////////////////////////////////////////////

Expand Down
2 changes: 2 additions & 0 deletions src/task/task_create.c
Original file line number Diff line number Diff line change
Expand Up @@ -25,5 +25,7 @@ kfTask kfTask_create(
.clearer = opt->clearer,
.target = delay,
.interval = interval,
.stage = opt->stage,
.masks = opt->masks,
};
}
Loading
Loading