diff --git a/.github/ISSUE_TEMPLATE/bug_report.md b/.github/ISSUE_TEMPLATE/bug_report.md new file mode 100644 index 0000000..266fd1f --- /dev/null +++ b/.github/ISSUE_TEMPLATE/bug_report.md @@ -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. \ No newline at end of file diff --git a/.github/ISSUE_TEMPLATE/feature_request.md b/.github/ISSUE_TEMPLATE/feature_request.md new file mode 100644 index 0000000..972fe3d --- /dev/null +++ b/.github/ISSUE_TEMPLATE/feature_request.md @@ -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. \ No newline at end of file diff --git a/.github/ISSUE_TEMPLATE/question.md b/.github/ISSUE_TEMPLATE/question.md new file mode 100644 index 0000000..1b086a5 --- /dev/null +++ b/.github/ISSUE_TEMPLATE/question.md @@ -0,0 +1,8 @@ +--- +name: Question +about: Ask what you want about the project +title: "[QUESTION] - TITLE" +labels: question +assignees: cartepro + +--- \ No newline at end of file diff --git a/.github/pull_request_template.md b/.github/pull_request_template.md new file mode 100644 index 0000000..4624b67 --- /dev/null +++ b/.github/pull_request_template.md @@ -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 \ No newline at end of file diff --git a/CMakeLists.txt b/CMakeLists.txt index 243779c..179e468 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -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) @@ -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" ) @@ -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 $ @@ -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() diff --git a/include/kronkflow/macros/types.h b/include/kronkflow/macros/types.h index 644f45a..b9f9f30 100644 --- a/include/kronkflow/macros/types.h +++ b/include/kronkflow/macros/types.h @@ -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 */ diff --git a/include/kronkflow/task.h b/include/kronkflow/task.h index 2016361..10661fa 100644 --- a/include/kronkflow/task.h +++ b/include/kronkflow/task.h @@ -8,6 +8,7 @@ #define PROPHECY_TASKS_H #include "kronkflow/macros/optimization.h" #include "kronkflow/macros/types.h" + #include /////////////////////////////////////////////////////////////////////////////// /** @@ -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 @@ -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; /////////////////////////////////////////////////////////////////////////////// @@ -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); /////////////////////////////////////////////////////////////////////////////// diff --git a/src/scheduler/clear/scheduler_clear.c b/src/scheduler/clear/scheduler_clear.c index c0597ed..449fe06 100644 --- a/src/scheduler/clear/scheduler_clear.c +++ b/src/scheduler/clear/scheduler_clear.c @@ -5,6 +5,7 @@ ** Clear a scheduler */ #include "../scheduler.h" +#include "dynarray.h" #include #include @@ -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; diff --git a/src/scheduler/init/scheduler_init.c b/src/scheduler/init/scheduler_init.c index 241dbcf..8235050 100644 --- a/src/scheduler/init/scheduler_init.c +++ b/src/scheduler/init/scheduler_init.c @@ -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 int kfScheduler_init( @@ -24,5 +26,6 @@ int kfScheduler_init( sch->size = size; sch->count = 0; sch->tick = 0; + kuDynarray_init(&sch->staged, 2, kfTask *); return 0; } diff --git a/src/scheduler/scheduler.h b/src/scheduler/scheduler.h index 2214590..454ad46 100644 --- a/src/scheduler/scheduler.h +++ b/src/scheduler/scheduler.h @@ -7,7 +7,16 @@ #ifndef PROPHECY_SCHEDULER_IMPL_H #define PROPHECY_SCHEDULER_IMPL_H #include "../task/task.h" + #include "kronkflow/task.h" #include + #include + +typedef struct prophecy_stage_data_s { + + kfTask* tasks; + bool passed; + +} kfStageData; /////////////////////////////////////////////////////////////////////////////// /** @@ -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; /////////////////////////////////////////////////////////////////////////////// diff --git a/src/scheduler/scheduler_update.c b/src/scheduler/scheduler_update.c index a84df3c..d15542f 100644 --- a/src/scheduler/scheduler_update.c +++ b/src/scheduler/scheduler_update.c @@ -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 +#include #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, @@ -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; } diff --git a/src/task/task.h b/src/task/task.h index 8000760..c1930f1 100644 --- a/src/task/task.h +++ b/src/task/task.h @@ -8,10 +8,9 @@ #define PROPHECY_TASK_IMPL_H #include "kronkflow/macros/types.h" #include "kronkflow/task.h" + #include #include - - /////////////////////////////////////////////////////////////////////////////// /** * @struct prophecy_task_s @@ -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; /////////////////////////////////////////////////////////////////////////////// diff --git a/src/task/task_create.c b/src/task/task_create.c index 1333df5..01c47e7 100644 --- a/src/task/task_create.c +++ b/src/task/task_create.c @@ -25,5 +25,7 @@ kfTask kfTask_create( .clearer = opt->clearer, .target = delay, .interval = interval, + .stage = opt->stage, + .masks = opt->masks, }; } diff --git a/src/task/task_opt.c b/src/task/task_opt.c index 7e97e8f..41c7594 100644 --- a/src/task/task_opt.c +++ b/src/task/task_opt.c @@ -13,12 +13,16 @@ inline kfTaskOpt kfTask_opt( kfHandler handler, void *data, - kfClearer clearer + kfClearer clearer, + kfStageId stageId, + kfRWMasks rwmask ) { return (kfTaskOpt){ .handler = handler, .data = data, - .clearer = clearer + .clearer = clearer, + .stage = stageId, + .masks = rwmask }; } diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt new file mode 100644 index 0000000..2474632 --- /dev/null +++ b/tests/CMakeLists.txt @@ -0,0 +1,19 @@ +file(GLOB_RECURSE kronkflow_TEST_SOURCES + "${CMAKE_CURRENT_SOURCE_DIR}/*.c" +) + +add_executable(kronkflow_tests ${kronkflow_TEST_SOURCES}) + +target_link_libraries(kronkflow_tests PRIVATE + kronkflow-static + kronklab-static +) + +target_compile_options(kronkflow_tests PRIVATE + -Wall + -Wextra +) + +target_compile_features(kronkflow_tests PRIVATE c_std_11) + +add_test(NAME kronkflow COMMAND kronkflow_tests) diff --git a/tests/scheduler_regression_test.c b/tests/scheduler_regression_test.c new file mode 100644 index 0000000..006cda8 --- /dev/null +++ b/tests/scheduler_regression_test.c @@ -0,0 +1,183 @@ +/* +** FREE PROJECT, 2026 +** PROPHECY +** File description: +** Regression tests for the scheduler public API +*/ +#include +#include +#include +#include + +typedef struct { + int calls; + int cleared; +} Counters; + +static kfBool count_call( + void *context, + void *data +) +{ + (void)data; + ((Counters *)context)->calls++; + return kfTrue; +} + +static void count_clear( + void *data +) +{ + ((Counters *)data)->cleared++; +} + +Test(scheduler, create_and_destroy_empty) +{ + kfScheduler *sch = kfScheduler_create(4); + + AssertNotNull(sch, "create should succeed"); + kfScheduler_destroy(sch); +} + +Test(scheduler, destroy_after_buckets_used) +{ + kfScheduler *sch = kfScheduler_create(4); + Counters counters = { 0, 0 }; + kfTaskOpt opt = kfTask_opt(&count_call, NULL, NULL, 0, (kfRWMasks){ 0, 0 }); + + AssertNotNull(sch, "create should succeed"); + opt.stage = 0; + kfScheduler_addTask(sch, opt, 1, 0); + opt.stage = 3; + kfScheduler_addTask(sch, opt, 1, 0); + opt.stage = 7; + kfScheduler_addTask(sch, opt, 1, 0); + kfScheduler_tick(sch, &counters); + AssertEq(counters.calls, 3, "all three staged tasks should have run"); + /* Regression: freeing kuDynarray_at(sch->staged, i) instead of + ** sch->staged[i] used to double-free the outer bucket array. */ + kfScheduler_destroy(sch); +} + +Test(scheduler, tick_return_counts_executed) +{ + kfScheduler *sch = kfScheduler_create(4); + Counters counters = { 0, 0 }; + kfTaskOpt opt = kfTask_opt(&count_call, NULL, NULL, 0, (kfRWMasks){ 0, 0 }); + size_t done; + + kfScheduler_addTask(sch, opt, 1, 0); + kfScheduler_addTask(sch, opt, 1, 0); + /* Regression: sch->staged's load was never advanced past 0, so the + ** stage-processing loop used to iterate zero times and every tick + ** silently ran nothing. */ + done = kfScheduler_tick(sch, &counters); + AssertEq(done, (size_t)2, "tick should report the number of tasks it ran"); + AssertEq(counters.calls, 2, "handlers should have actually been invoked"); + kfScheduler_destroy(sch); +} + +Test(scheduler, task_not_due_not_executed) +{ + kfScheduler *sch = kfScheduler_create(4); + Counters counters = { 0, 0 }; + kfTaskOpt opt = kfTask_opt(&count_call, NULL, NULL, 0, (kfRWMasks){ 0, 0 }); + size_t i; + + kfScheduler_addTask(sch, opt, 5, 0); + for (i = 0; i < 4; i++) { + kfScheduler_tick(sch, &counters); + } + AssertEq(counters.calls, 0, "should not run before its target tick"); + kfScheduler_tick(sch, &counters); + AssertEq(counters.calls, 1, "should run exactly on its target tick"); + kfScheduler_destroy(sch); +} + +Test(scheduler, interval_zero_runs_once) +{ + kfScheduler *sch = kfScheduler_create(4); + Counters counters = { 0, 0 }; + kfTaskOpt opt = kfTask_opt(&count_call, NULL, NULL, 0, (kfRWMasks){ 0, 0 }); + size_t i; + + kfScheduler_addTask(sch, opt, 1, 0); + for (i = 0; i < 20; i++) { + kfScheduler_tick(sch, &counters); + } + AssertEq(counters.calls, 1, "a task with interval 0 must not be rescheduled"); + kfScheduler_destroy(sch); +} + +Test(scheduler, interval_task_reschedules) +{ + kfScheduler *sch = kfScheduler_create(4); + Counters counters = { 0, 0 }; + kfTaskOpt opt = kfTask_opt(&count_call, NULL, NULL, 0, (kfRWMasks){ 0, 0 }); + size_t i; + + kfScheduler_addTask(sch, opt, 1, 3); + for (i = 0; i < 30; i++) { + kfScheduler_tick(sch, &counters); + } + AssertEq(counters.calls, 10, + "a task rescheduled every 3 ticks starting at tick 1 should run " + "at ticks 1,4,7,...,28 within 30 ticks"); + kfScheduler_destroy(sch); +} + +Test(scheduler, non_reschedule_calls_clearer) +{ + kfScheduler *sch = kfScheduler_create(4); + Counters counters = { 0, 0 }; + kfTaskOpt opt = kfTask_opt(&count_call, &counters, &count_clear, 0, (kfRWMasks){ 0, 0 }); + + kfScheduler_addTask(sch, opt, 1, 0); + kfScheduler_tick(sch, &counters); + AssertEq(counters.cleared, 1, "clearer should run once the task is not rescheduled"); + kfScheduler_destroy(sch); +} + +Test(scheduler, reschedule_skips_clearer) +{ + kfScheduler *sch = kfScheduler_create(4); + Counters counters = { 0, 0 }; + kfTaskOpt opt = kfTask_opt(&count_call, &counters, &count_clear, 0, (kfRWMasks){ 0, 0 }); + + kfScheduler_addTask(sch, opt, 1, 2); + kfScheduler_tick(sch, &counters); + AssertEq(counters.cleared, 0, "a rescheduled task's data is still owned, clearer must not fire"); + kfScheduler_destroy(sch); +} + +Test(scheduler, addTask_null_sch_safe) +{ + kfTaskOpt opt = kfTask_opt(&count_call, NULL, NULL, 0, (kfRWMasks){ 0, 0 }); + + AssertEq(kfScheduler_addTask(NULL, opt, 1, 0), (kfTaskID)0, + "adding to a null scheduler should fail cleanly"); +} + +Test(scheduler, tick_null_sch_safe) +{ + AssertEq(kfScheduler_tick(NULL, NULL), (size_t)0, + "ticking a null scheduler should return zero"); +} + +Test(scheduler, destroy_null_sch_safe) +{ + kfScheduler_destroy(NULL); +} + +Test(scheduler, tick_count_advances) +{ + kfScheduler *sch = kfScheduler_create(4); + size_t i; + + AssertEq(kfScheduler_currentTick(sch), (kfTick)0, "a fresh scheduler starts at tick 0"); + for (i = 0; i < 5; i++) { + kfScheduler_tick(sch, NULL); + } + AssertEq(kfScheduler_currentTick(sch), (kfTick)5, "five ticks should advance the clock by five"); + kfScheduler_destroy(sch); +} diff --git a/tests/staging_test.c b/tests/staging_test.c new file mode 100644 index 0000000..f5fa2c9 --- /dev/null +++ b/tests/staging_test.c @@ -0,0 +1,174 @@ +/* +** FREE PROJECT, 2026 +** PROPHECY +** File description: +** Tests for staged (per-stage bucketed) execution +*/ +#include +#include +#include +#include +#include + +typedef struct { + int order[64]; + int len; +} StageOrder; + +static kfBool record_stage( + void *context, + void *data +) +{ + StageOrder *rec = context; + long marker = (long)(intptr_t)data; + + if (rec->len < 64) { + rec->order[rec->len++] = (int)marker; + } + return kfTrue; +} + +static kfTaskOpt marked_task( + kfStageId stage, + long marker +) +{ + return kfTask_opt(&record_stage, (void *)(intptr_t)marker, NULL, + stage, (kfRWMasks){ 0, 0 }); +} + +Test(staging, sorted_by_stage_not_insertion) +{ + kfScheduler *sch = kfScheduler_create(4); + StageOrder rec = { { 0 }, 0 }; + + kfScheduler_addTask(sch, marked_task(3, 3), 1, 0); + kfScheduler_addTask(sch, marked_task(0, 0), 1, 0); + kfScheduler_addTask(sch, marked_task(5, 5), 1, 0); + kfScheduler_addTask(sch, marked_task(1, 1), 1, 0); + + kfScheduler_tick(sch, &rec); + + AssertEq(rec.len, 4, "all four tasks should have run"); + AssertEq(rec.order[0], 0, "stage 0 should run first"); + AssertEq(rec.order[1], 1, "stage 1 should run second"); + AssertEq(rec.order[2], 3, "stage 3 should run third"); + AssertEq(rec.order[3], 5, "stage 5 should run last"); + kfScheduler_destroy(sch); +} + +Test(staging, same_stage_runs_together) +{ + kfScheduler *sch = kfScheduler_create(4); + StageOrder rec = { { 0 }, 0 }; + int seen_11 = 0; + int seen_22 = 0; + int seen_33 = 0; + int i; + + kfScheduler_addTask(sch, marked_task(2, 11), 1, 0); + kfScheduler_addTask(sch, marked_task(2, 22), 1, 0); + kfScheduler_addTask(sch, marked_task(2, 33), 1, 0); + kfScheduler_addTask(sch, marked_task(5, 99), 1, 0); + + kfScheduler_tick(sch, &rec); + + AssertEq(rec.len, 4, "all four tasks should have run"); + AssertEq(rec.order[3], 99, "the later stage must run after every stage-2 task"); + /* The min-heap gives no ordering guarantee among tasks sharing the same + ** target tick, so same-stage siblings may land in any relative order - + ** only check that all three actually ran, as a set. */ + for (i = 0; i < 3; i++) { + if (rec.order[i] == 11) { + seen_11 = 1; + } + if (rec.order[i] == 22) { + seen_22 = 1; + } + if (rec.order[i] == 33) { + seen_33 = 1; + } + } + AssertEq(seen_11, 1, "task 11 should have run"); + AssertEq(seen_22, 1, "task 22 should have run"); + AssertEq(seen_33, 1, "task 33 should have run"); + kfScheduler_destroy(sch); +} + +Test(staging, sparse_stage_no_phantom_tasks) +{ + kfScheduler *sch = kfScheduler_create(4); + StageOrder rec = { { 0 }, 0 }; + size_t done; + + /* Regression: the outer bucket array's freshly-grown capacity was + ** never zeroed, so lower/unused stage slots held garbage pointers + ** instead of NULL. */ + kfScheduler_addTask(sch, marked_task(9, 9), 1, 0); + done = kfScheduler_tick(sch, &rec); + + AssertEq(done, (size_t)1, "only the single staged task should run"); + AssertEq(rec.len, 1, "no phantom task from unused lower stages should run"); + AssertEq(rec.order[0], 9, "the one task should still run at its own stage"); + kfScheduler_destroy(sch); +} + +Test(staging, buckets_cleared_and_reused) +{ + kfScheduler *sch = kfScheduler_create(4); + StageOrder rec = { { 0 }, 0 }; + + kfScheduler_addTask(sch, marked_task(2, 1), 1, 0); + kfScheduler_tick(sch, &rec); + AssertEq(rec.len, 1, "the first tick should run the stage-2 task once"); + + kfScheduler_tick(sch, &rec); + AssertEq(rec.len, 1, "an empty tick must not replay the previous tick's bucket"); + + kfScheduler_addTask(sch, marked_task(2, 2), 1, 0); + kfScheduler_tick(sch, &rec); + AssertEq(rec.len, 2, "a new task pushed into a reused bucket should run exactly once"); + kfScheduler_destroy(sch); +} + +Test(staging, growing_stage_range_stays_ok) +{ + kfScheduler *sch = kfScheduler_create(4); + StageOrder rec = { { 0 }, 0 }; + size_t stage; + size_t total_done = 0; + + /* Forces sch->staged to grow its capacity on almost every iteration, + ** exercising the resize -> zero-new-slots -> sync-load path. */ + for (stage = 0; stage < 40; stage++) { + kfScheduler_addTask(sch, marked_task((kfStageId)stage, (long)stage), 1, 0); + total_done += kfScheduler_tick(sch, &rec); + } + AssertEq(total_done, (size_t)40, + "every one of the 40 growth ticks should execute exactly one task"); + AssertEq(rec.len, 40, "all 40 tasks across growing stage ids should have run"); + kfScheduler_destroy(sch); +} + +Test(staging, stage_order_after_reschedule) +{ + kfScheduler *sch = kfScheduler_create(4); + StageOrder rec = { { 0 }, 0 }; + + kfScheduler_addTask(sch, marked_task(5, 5), 1, 2); + kfScheduler_addTask(sch, marked_task(1, 1), 3, 0); + + kfScheduler_tick(sch, &rec); + AssertEq(rec.len, 1, "only the stage-5 task is due on tick 1"); + AssertEq(rec.order[0], 5, "stage 5 ran alone on tick 1"); + + kfScheduler_tick(sch, &rec); + AssertEq(rec.len, 1, "nothing new is due on tick 2"); + + kfScheduler_tick(sch, &rec); + AssertEq(rec.len, 3, "tick 3 reruns stage 5 (rescheduled) and stage 1, in stage order"); + AssertEq(rec.order[1], 1, "stage 1 must run before the rescheduled stage 5"); + AssertEq(rec.order[2], 5, "the rescheduled stage-5 task runs after stage 1"); + kfScheduler_destroy(sch); +}