diff --git a/src/os/posix/src/os-impl-queues.c b/src/os/posix/src/os-impl-queues.c index f72ced728..63fe33add 100644 --- a/src/os/posix/src/os-impl-queues.c +++ b/src/os/posix/src/os-impl-queues.c @@ -30,6 +30,12 @@ #include "os-posix.h" #include "bsp-impl.h" +#ifdef __linux__ +#include +#include +#include +#endif + #include "os-impl-queues.h" #include "os-shared-queue.h" #include "os-shared-idmap.h" @@ -41,6 +47,161 @@ OS_impl_queue_internal_record_t OS_impl_queue_table[OS_MAX_QUEUES]; MESSAGE QUEUE API ***************************************************************************************/ +#ifdef __linux__ +/*---------------------------------------------------------------- + * + * Purpose: Local helper function + * + * This function accept time interval, msecs, as an input and + * computes the absolute time at which this time interval will expire. + * The absolute time is programmed into a struct. + * + *-----------------------------------------------------------------*/ +static int OS_Posix_LinuxMqMonotonicDeadline_Impl(uint32 msecs, struct timespec *tm) +{ + if (clock_gettime(CLOCK_MONOTONIC, tm) != 0) + { + return -1; + } + + /* add the delay to the current time */ + tm->tv_sec += (time_t)(msecs / 1000); + /* convert residue ( msecs ) to nanoseconds */ + tm->tv_nsec += (msecs % 1000) * 1000000L; + + if (tm->tv_nsec >= 1000000000L) + { + tm->tv_nsec -= 1000000000L; + tm->tv_sec++; + } + + return 0; +} +#endif + +#ifdef __linux__ +/*--------------------------------------------------------------------------------------- + Name: OS_Posix_LinuxMqReceiveUntilMonotonicDeadline_Impl + + Purpose: Linux-only helper that waits for queue data using a monotonic deadline. + Linux documents pollable mqueue descriptors; this is not portable POSIX behavior. + + ----------------------------------------------------------------------------------------*/ +static ssize_t OS_Posix_LinuxMqReceiveUntilMonotonicDeadline_Impl(mqd_t mqd, char *buf, size_t len, unsigned *prio, + const struct timespec *deadline) +{ + struct timespec now; + struct timespec immediate_ts = {0, 0}; + struct pollfd pfd; + + pfd.fd = (int)mqd; + pfd.events = POLLIN; + pfd.revents = 0; + + /* Loop until success, timeout, or non-retryable error. */ + while (1) + { + int64_t rem_sec; + int64_t rem_nsec; + int64_t rem_ms; + int timeout_ms; + int rc; + ssize_t size; + + if (clock_gettime(CLOCK_MONOTONIC, &now) != 0) + { + return -1; + } + + if ((now.tv_sec > deadline->tv_sec) || (now.tv_sec == deadline->tv_sec && now.tv_nsec >= deadline->tv_nsec)) + { + errno = ETIMEDOUT; + return -1; + } + + rem_sec = deadline->tv_sec - now.tv_sec; + rem_nsec = deadline->tv_nsec - now.tv_nsec; + if (rem_nsec < 0) + { + --rem_sec; + rem_nsec += 1000000000L; + } + + /* Round up so poll() does not undersleep. */ + rem_ms = (rem_sec * 1000) + (rem_nsec + 999999) / 1000000; + if (rem_ms > INT_MAX) + { + timeout_ms = INT_MAX; + } + else if (rem_ms < 0) + { + timeout_ms = 0; + } + else + { + timeout_ms = (int)rem_ms; + } + + rc = poll(&pfd, 1, timeout_ms); + if (rc < 0) + { + if (errno == EINTR) + { + continue; + } + + return -1; + } + + /* Recheck the deadline because poll() uses millisecond granularity. */ + if (rc == 0) + { + if (clock_gettime(CLOCK_MONOTONIC, &now) != 0) + { + return -1; + } + + if ((now.tv_sec > deadline->tv_sec) || + (now.tv_sec == deadline->tv_sec && now.tv_nsec >= deadline->tv_nsec)) + { + errno = ETIMEDOUT; + return -1; + } + + continue; + } + + if ((pfd.revents & POLLIN) == 0) + { + if ((pfd.revents & POLLNVAL) != 0) + { + errno = EBADF; + } + else + { + errno = EIO; + } + + return -1; + } + + /* Use an expired timeout so the receive step cannot block. */ + size = mq_timedreceive(mqd, buf, len, prio, &immediate_ts); + if (size >= 0) + { + return size; + } + + if (errno == EINTR || errno == EAGAIN || errno == ETIMEDOUT) + { + continue; + } + + return -1; + } +} +#endif + /*--------------------------------------------------------------------------------------- Name: OS_Posix_QueueAPI_Impl_Init @@ -190,6 +351,9 @@ int32 OS_QueueGet_Impl(const OS_object_token_t *token, void *data, size_t size, int32 return_code; ssize_t sizeCopied; struct timespec ts; +#ifdef __linux__ + struct timespec monotonic_deadline; +#endif OS_impl_queue_internal_record_t *impl; impl = OS_OBJECT_TABLE_GET(OS_impl_queue_table, *token); @@ -211,6 +375,27 @@ int32 OS_QueueGet_Impl(const OS_object_token_t *token, void *data, size_t size, } else { +#ifdef __linux__ + if (timeout == OS_CHECK) + { + /* + * NOTE - a prior implementation of OS_CHECK would check the mq_attr for a nonzero depth + * and then call mq_receive(). This is insufficient since another thread might do the same + * thing at the same time in which case one thread will read and the other will block. + * + * Calling mq_timedreceive with a zero timeout effectively does the same thing in the typical + * case, but for the case where two threads do a simultaneous read, one will get the message + * while the other will NOT block (as expected). + */ + memset(&ts, 0, sizeof(ts)); + sizeCopied = mq_timedreceive(impl->id, data, size, NULL, &ts); + } + else if (OS_Posix_LinuxMqMonotonicDeadline_Impl(timeout, &monotonic_deadline) == 0) + { + sizeCopied = OS_Posix_LinuxMqReceiveUntilMonotonicDeadline_Impl(impl->id, data, size, NULL, + &monotonic_deadline); + } +#else /* * NOTE - a prior implementation of OS_CHECK would check the mq_attr for a nonzero depth * and then call mq_receive(). This is insufficient since another thread might do the same @@ -237,6 +422,7 @@ int32 OS_QueueGet_Impl(const OS_object_token_t *token, void *data, size_t size, { sizeCopied = mq_timedreceive(impl->id, data, size, NULL, &ts); } while (timeout != OS_CHECK && sizeCopied < 0 && errno == EINTR); +#endif } /* END timeout */ diff --git a/src/tests/queue-test/queue-test.c b/src/tests/queue-test/queue-test.c index 3e49130be..42f6ab858 100644 --- a/src/tests/queue-test/queue-test.c +++ b/src/tests/queue-test/queue-test.c @@ -19,7 +19,11 @@ /* ** Queue read timeout test */ +#include #include +#include +#include +#include #include "common_types.h" #include "osapi.h" #include "utassert.h" @@ -29,12 +33,16 @@ /* Define setup and check functions for UT assert */ void QueueTimeoutSetup(void); void QueueTimeoutCheck(void); +void QueueTimeoutTimeJumpSetup(void); #define MSGQ_DEPTH 50 #define MSGQ_SIZE sizeof(uint32) #define MSGQ_TOTAL 10 #define MSGQ_BURST 3 +#define TIMEJUMP_SECONDS 10 +#define TIMEJUMP_INTERVAL_SECONDS 5 + /* Task 1 */ #define TASK_1_STACK_SIZE 4096 #define TASK_1_PRIORITY 101 @@ -50,6 +58,104 @@ uint32 task_2_stack[TASK_2_STACK_SIZE]; osal_id_t task_2_id; osal_id_t msgq_id; +static bool queue_test_timejump_created; +static bool queue_test_timejump_skip; + +#ifdef __linux__ +static bool queue_test_timejump_applied; +static int queue_test_timejump_errno; +static bool queue_test_saved_time_valid; +static struct timespec queue_test_saved_realtime; +static struct timespec queue_test_saved_monotonic; + +static void QueueTest_AddTimespec(const struct timespec *left, const struct timespec *right, struct timespec *result) +{ + result->tv_sec = left->tv_sec + right->tv_sec; + result->tv_nsec = left->tv_nsec + right->tv_nsec; + + if (result->tv_nsec >= 1000000000L) + { + result->tv_nsec -= 1000000000L; + ++result->tv_sec; + } +} + +static void QueueTest_SubtractTimespec(const struct timespec *left, const struct timespec *right, + struct timespec *result) +{ + result->tv_sec = left->tv_sec - right->tv_sec; + result->tv_nsec = left->tv_nsec - right->tv_nsec; + + if (result->tv_nsec < 0) + { + result->tv_nsec += 1000000000L; + --result->tv_sec; + } +} + +static void QueueTimeoutRestoreRealtime(void) +{ + struct timespec current_monotonic; + struct timespec elapsed_time; + struct timespec restore_realtime; + int status; + + /* Remove the artificial jump while preserving real elapsed time. */ + if (!queue_test_saved_time_valid || !queue_test_timejump_applied) + { + return; + } + + status = clock_gettime(CLOCK_MONOTONIC, ¤t_monotonic); + UtAssert_True(status == 0, "clock_gettime CLOCK_MONOTONIC Rc=0"); + if (status != 0) + { + return; + } + + QueueTest_SubtractTimespec(¤t_monotonic, &queue_test_saved_monotonic, &elapsed_time); + QueueTest_AddTimespec(&queue_test_saved_realtime, &elapsed_time, &restore_realtime); + + status = clock_settime(CLOCK_REALTIME, &restore_realtime); + UtAssert_True(status == 0, "clock_settime CLOCK_REALTIME restore Rc=0"); + if (status == 0) + { + queue_test_timejump_applied = false; + } +} + +static void QueueTimeJumpTask(void) +{ + struct timespec realtime_clock_timeval; + + /* Apply one jump, then idle until teardown deletes this task. */ + OS_TaskDelay(TIMEJUMP_INTERVAL_SECONDS * 1000); + + if (clock_gettime(CLOCK_REALTIME, &realtime_clock_timeval) != 0) + { + queue_test_timejump_errno = errno; + } + else + { + realtime_clock_timeval.tv_sec += TIMEJUMP_SECONDS; + + if (clock_settime(CLOCK_REALTIME, &realtime_clock_timeval) != 0) + { + queue_test_timejump_errno = errno; + } + else + { + queue_test_timejump_applied = true; + } + } + + while (1) + { + OS_TaskDelay(1000); + } +} +#endif + uint32 timer_counter; osal_id_t timer_id; uint32 timer_start = 10000; @@ -89,7 +195,6 @@ void task_1(void) else if (status == OS_QUEUE_TIMEOUT) { ++task_1_timeouts; - OS_printf("TASK 1: Timeout on Queue! Timer counter = %u\n", (unsigned int)timer_counter); } else { @@ -100,15 +205,35 @@ void task_1(void) } } -void QueueTimeoutCheck(void) +static void QueueTimeoutCheckInternal(void) { int32 status; uint32 limit; + if (queue_test_timejump_skip) + { + return; + } + status = OS_TimerDelete(timer_id); UtAssert_True(status == OS_SUCCESS, "Timer delete Rc=%d", (int)status); status = OS_TaskDelete(task_1_id); UtAssert_True(status == OS_SUCCESS, "Task 1 delete Rc=%d", (int)status); + if (queue_test_timejump_created) + { + status = OS_TaskDelete(task_2_id); + UtAssert_True(status == OS_SUCCESS, "Task 2 delete Rc=%d", (int)status); + } + +#ifdef __linux__ + if (queue_test_timejump_created) + { + UtAssert_True(queue_test_timejump_errno == 0, "Timejump task errno=%d", queue_test_timejump_errno); + UtAssert_True(queue_test_timejump_applied, "Timejump task applied"); + QueueTimeoutRestoreRealtime(); + } +#endif + status = OS_QueueDelete(msgq_id); UtAssert_True(status == OS_SUCCESS, "Queue 1 delete Rc=%d", (int)status); @@ -130,7 +255,7 @@ void QueueTimeoutCheck(void) (unsigned int)limit); } -void QueueTimeoutSetup(void) +static void QueueTimeoutSetupInternal(bool enable_timejump) { int32 status; uint32 accuracy = 0; @@ -138,6 +263,47 @@ void QueueTimeoutSetup(void) task_1_failures = 0; task_1_messages = 0; task_1_timeouts = 0; + queue_test_timejump_created = false; + queue_test_timejump_skip = false; + timer_counter = 0; + +#ifdef __linux__ + queue_test_timejump_applied = false; + queue_test_timejump_errno = 0; + queue_test_saved_time_valid = false; +#endif + + if (enable_timejump) + { +#ifndef __linux__ + queue_test_timejump_skip = true; + UtAssert_WARN("Timejump test skipped: Linux only"); + return; +#else + if (geteuid() != 0) + { + queue_test_timejump_skip = true; + UtAssert_WARN("Timejump test skipped: requires root"); + return; + } + + status = clock_gettime(CLOCK_REALTIME, &queue_test_saved_realtime); + UtAssert_True(status == 0, "clock_gettime CLOCK_REALTIME Rc=0"); + if (status != 0) + { + return; + } + + status = clock_gettime(CLOCK_MONOTONIC, &queue_test_saved_monotonic); + UtAssert_True(status == 0, "clock_gettime CLOCK_MONOTONIC Rc=0"); + if (status != 0) + { + return; + } + + queue_test_saved_time_valid = true; +#endif + } status = OS_QueueCreate(&msgq_id, "MsgQ", OSAL_BLOCKCOUNT_C(MSGQ_DEPTH), OSAL_SIZE_C(MSGQ_SIZE), 0); UtAssert_True(status == OS_SUCCESS, "MsgQ create Id=%lx Rc=%d", OS_ObjectIdToInteger(msgq_id), (int)status); @@ -149,6 +315,23 @@ void QueueTimeoutSetup(void) OSAL_PRIORITY_C(TASK_1_PRIORITY), 0); UtAssert_True(status == OS_SUCCESS, "Task 1 create Id=%lx Rc=%d", OS_ObjectIdToInteger(task_1_id), (int)status); + if (enable_timejump) + { +#ifdef __linux__ + /* + ** Create the time jumper task. + */ + status = OS_TaskCreate(&task_2_id, "Task 2", QueueTimeJumpTask, OSAL_STACKPTR_C(task_2_stack), + sizeof(task_2_stack), OSAL_PRIORITY_C(TASK_2_PRIORITY), 0); + UtAssert_True(status == OS_SUCCESS, "Task 2 create Id=%lx Rc=%d", OS_ObjectIdToInteger(task_2_id), + (int)status); + if (status == OS_SUCCESS) + { + queue_test_timejump_created = true; + } +#endif + } + /* ** Create a timer */ @@ -167,6 +350,29 @@ void QueueTimeoutSetup(void) { OS_TaskDelay(100); } + +#ifdef __linux__ + if (enable_timejump) + { + UtAssert_True(queue_test_timejump_errno == 0, "Timejump task errno=%d", queue_test_timejump_errno); + UtAssert_True(queue_test_timejump_applied, "Timejump task applied"); + } +#endif +} + +void QueueTimeoutSetup(void) +{ + QueueTimeoutSetupInternal(false); +} + +void QueueTimeoutTimeJumpSetup(void) +{ + QueueTimeoutSetupInternal(true); +} + +void QueueTimeoutCheck(void) +{ + QueueTimeoutCheckInternal(); } void QueueMessageCheck(void) @@ -251,5 +457,6 @@ void UtTest_Setup(void) * Register the test setup and check routines in UT assert */ UtTest_Add(QueueTimeoutCheck, QueueTimeoutSetup, NULL, "QueueTimeoutTest"); + UtTest_Add(QueueTimeoutCheck, QueueTimeoutTimeJumpSetup, NULL, "QueueTimeoutTimeJumpTest"); UtTest_Add(QueueMessageCheck, QueueMessageSetup, NULL, "QueueMessageCheck"); }