diff --git a/core/src/CMakeLists.txt b/core/src/CMakeLists.txt index bdd89c56..3055e545 100644 --- a/core/src/CMakeLists.txt +++ b/core/src/CMakeLists.txt @@ -43,13 +43,13 @@ add_executable(plc_main ${CMAKE_SOURCE_DIR}/core/src/plc_app/located_globals.c ${CMAKE_SOURCE_DIR}/core/src/plc_app/journal_buffer.c ${CMAKE_SOURCE_DIR}/core/src/plc_app/debug_write_journal.cpp - ${CMAKE_SOURCE_DIR}/core/src/plc_app/plc_io_cycle.cpp ${CMAKE_SOURCE_DIR}/core/src/plc_app/plc_retain.cpp ${CMAKE_SOURCE_DIR}/core/src/plc_app/plc_retain_file_store.cpp ${CMAKE_SOURCE_DIR}/core/src/plc_app/plc_state_manager.cpp ${CMAKE_SOURCE_DIR}/core/src/plc_app/plc_switch.c ${CMAKE_SOURCE_DIR}/core/src/plc_app/plcapp_manager.c ${CMAKE_SOURCE_DIR}/core/src/plc_app/scan_cycle_manager.c + ${CMAKE_SOURCE_DIR}/core/src/plc_app/task_policy.c ${CMAKE_SOURCE_DIR}/core/src/drivers/plugin_driver.c ${CMAKE_SOURCE_DIR}/core/src/drivers/plugin_config.c ${CMAKE_SOURCE_DIR}/core/src/drivers/vpp_plugin_seal.c diff --git a/core/src/plc_app/debug_write_journal.cpp b/core/src/plc_app/debug_write_journal.cpp index 4b878000..e7b37a00 100644 --- a/core/src/plc_app/debug_write_journal.cpp +++ b/core/src/plc_app/debug_write_journal.cpp @@ -31,6 +31,7 @@ #include "image_tables.h" /* ext_strucpp_debug_set / _write / _locate */ #include "journal_buffer.h" /* journal_write_* / journal_force_set/clear */ +#include "utils/rt_mutex.h" extern "C" { #include "utils/log.h" @@ -61,6 +62,11 @@ std::atomic g_dbgw_count{0}; pthread_mutex_t g_dbgw_lock = PTHREAD_MUTEX_INITIALIZER; bool g_overflow_logged = false; +__attribute__((constructor)) void dbgw_lock_init_pi(void) +{ + rt_mutex_upgrade_static(&g_dbgw_lock, "g_dbgw_lock"); +} + /* LocatedArea (strucpp_abi.hpp): Input=0, Output=1, Memory=2. * LocatedSize: Bit=0, Byte=1, Word=2, DWord=3, LWord=4. */ enum { AREA_INPUT = 0, AREA_OUTPUT = 1, AREA_MEMORY = 2 }; diff --git a/core/src/plc_app/image_tables.cpp b/core/src/plc_app/image_tables.cpp index 79f659fb..d26427e9 100644 --- a/core/src/plc_app/image_tables.cpp +++ b/core/src/plc_app/image_tables.cpp @@ -694,6 +694,21 @@ void image_tables_fill_null_pointers(void) log_info("[image_tables] filled %d NULL slots with backing buffers", filled); } +void image_tables_zero_outputs(void) +{ + for (int i = 0; i < BUFFER_SIZE; ++i) + { + for (int b = 0; b < 8; ++b) + { + if (bool_output[i][b]) *bool_output[i][b] = 0; + } + if (byte_output[i]) *byte_output[i] = 0; + if (int_output[i]) *int_output[i] = 0; + if (dint_output[i]) *dint_output[i] = 0; + if (lint_output[i]) *lint_output[i] = 0; + } +} + void image_tables_clear_null_pointers(void) { // Threaded process-image state: free the dirty-diff snapshot. (The mutexes diff --git a/core/src/plc_app/image_tables.h b/core/src/plc_app/image_tables.h index 9c04e8b2..1df8741f 100644 --- a/core/src/plc_app/image_tables.h +++ b/core/src/plc_app/image_tables.h @@ -153,6 +153,15 @@ extern "C" * --------------------------------------------------------------------- */ void image_tables_clear_null_pointers(void); + /** + * @brief Write 0 to every output image slot (%QX, %QB, %QW, %QD, %QL). + * + * Used on every stop so plugins push de-energised outputs to the hardware + * before they are stopped. Program storage is not touched. Caller must hold + * the image-tables mutex. + */ + void image_tables_zero_outputs(void); + /* ------------------------------------------------------------------------- * Image-tables mutex accessor. Returns a pointer to the runtime-owned * recursive PI mutex that protects the image tables. The runtime locks diff --git a/core/src/plc_app/plc_io_cycle.cpp b/core/src/plc_app/plc_io_cycle.cpp deleted file mode 100644 index 079bbc5c..00000000 --- a/core/src/plc_app/plc_io_cycle.cpp +++ /dev/null @@ -1,48 +0,0 @@ -// SPDX-License-Identifier: MIT -// Copyright (c) 2026 Autonomy® - -// plc_io_cycle.cpp — per-cycle I/O work, split into pre/post halves -// around the fastest IEC task's body. -// -// Decoupled from plc_state_manager.cpp so the housekeeping is in one -// place. Both halves run inside the image-tables critical section. - -#include -#include - -extern "C" { -#include "../drivers/plugin_driver.h" -} - -#include "image_tables.h" -#include "journal_buffer.h" -#include "plc_io_cycle.h" -#include "utils/utils.h" - -extern std::atomic plc_heartbeat; -extern plugin_driver_t *plugin_driver; - -// --- Threaded (process-image) model housekeeping --------------------------- -// The drain runs at every task's copy-in (under the image mutex) so each task -// sees freshly-applied plugin/peer writes. The pre/post halves run only on the -// fastest task: pre opens the plugin cycle window before bodies, post advances -// the scan clock, closes the plugin window, and bumps the global heartbeat / -// scan counter once per scan. - -extern "C" void plc_run_io_cycle_threaded_drain(void) -{ - journal_apply_and_clear(); -} - -extern "C" void plc_run_io_cycle_threaded_pre(void) -{ - if (plugin_driver) plugin_driver_cycle_start(plugin_driver); -} - -extern "C" void plc_run_io_cycle_threaded_post(void) -{ - if (ext_strucpp_advance_time) ext_strucpp_advance_time(base_tick_ns); - if (plugin_driver) plugin_driver_cycle_end(plugin_driver); - plc_heartbeat.store((long)time(nullptr)); - ++scan_counter; -} diff --git a/core/src/plc_app/plc_io_cycle.h b/core/src/plc_app/plc_io_cycle.h deleted file mode 100644 index 3eabd0b3..00000000 --- a/core/src/plc_app/plc_io_cycle.h +++ /dev/null @@ -1,40 +0,0 @@ -// SPDX-License-Identifier: MIT -// Copyright (c) 2026 Autonomy® - -#ifndef OPENPLC_PLC_IO_CYCLE_H -#define OPENPLC_PLC_IO_CYCLE_H - -#ifdef __cplusplus -extern "C" { -#endif - -/* - * I/O cycle helpers — threaded (process-image) model housekeeping. - * - * Encapsulates the work that has to happen once per scan around the - * highest-priority task's body. The fastest task's thread (Phase 6 picks it; - * ctx->is_fastest_task) calls the pre/post halves around its body; every task - * calls the drain at its copy-in. Other task threads just run their bodies. - * - * See docs/strucpp-migration/07-runtime-v4-plugin-and-io.md for the - * topology rationale. - * - * plc_run_io_cycle_threaded_drain() — apply pending journal entries to the - * image. Called at EVERY task's copy-in, under the image mutex, so each - * task sees freshly-applied plugin/peer writes. - * - * plc_run_io_cycle_threaded_pre() — fire plugin cycle_start (fastest task, - * before bodies; no image lock held). - * - * plc_run_io_cycle_threaded_post() — advance time, fire plugin cycle_end, - * update heartbeat, increment scan_counter (fastest task, after bodies). - */ -void plc_run_io_cycle_threaded_drain(void); -void plc_run_io_cycle_threaded_pre(void); -void plc_run_io_cycle_threaded_post(void); - -#ifdef __cplusplus -} -#endif - -#endif /* OPENPLC_PLC_IO_CYCLE_H */ diff --git a/core/src/plc_app/plc_main.c b/core/src/plc_app/plc_main.c index 09087811..97b6e7f7 100644 --- a/core/src/plc_app/plc_main.c +++ b/core/src/plc_app/plc_main.c @@ -5,6 +5,7 @@ #include #include +#include #include #include #include @@ -20,6 +21,7 @@ #include "plc_state_manager.h" #include "plc_switch.h" #include "plcapp_manager.h" +#include "task_policy.h" #include "unix_socket.h" #include "utils/log.h" #include "utils/utils.h" @@ -53,6 +55,7 @@ int main(int argc, char *argv[]) { bool print_debug = false; bool safe_mode = false; + bool after_fault = false; // Check for command line arguments for (int i = 1; i < argc; i++) @@ -69,6 +72,10 @@ int main(int argc, char *argv[]) { safe_mode = true; } + else if (strcmp(argv[i], "--fault") == 0) + { + after_fault = true; + } } // Initialize logging system @@ -124,6 +131,29 @@ int main(int argc, char *argv[]) // and plc_set_state() is now the body of a claimed transition rather than a // setter -- calling it with nothing loaded would just log a failed unload. + bool skip_outputs_off = false; + if (access(PLC_WATCHDOG_FAULT_MARKER, F_OK) == 0) + { + char reason[512] = {0}; + FILE *marker = fopen(PLC_WATCHDOG_FAULT_MARKER, "r"); + if (marker) + { + size_t n = fread(reason, 1, sizeof(reason) - 1, marker); + reason[n] = '\0'; + fclose(marker); + } + skip_outputs_off = strstr(reason, PLC_FAULT_CONTEXT_BOOT_OUTPUTS_OFF) != NULL; + if (unlink(PLC_WATCHDOG_FAULT_MARKER) != 0) + log_warn("Could not remove %s: %s", PLC_WATCHDOG_FAULT_MARKER, strerror(errno)); + safe_mode = true; + after_fault = true; + } + if (after_fault && !safe_mode) + { + log_warn("--fault is only honoured together with --safe-mode; ignoring it"); + after_fault = false; + } + // Initialize watchdog if (watchdog_init() != 0) { @@ -176,6 +206,31 @@ int main(int argc, char *argv[]) } } + // Before the socket exists, so no command can claim a transition underneath. + if (safe_mode) + { + log_info("Runtime started in SAFE MODE - PLC program will not be loaded"); + log_info("Upload a corrected program to recover"); + if (after_fault) + { + log_error("Previous run ended in an unrecoverable watchdog fault"); + plc_force_error_state(); + if (skip_outputs_off) + { + log_error("Outputs not driven off: the previous attempt did not complete"); + } + else if (plc_claim_transition(PLC_STATE_STOPPED)) + { + // Bounded by the watchdog's stop budget; the context breaks a restart loop. + watchdog_set_fault_context(PLC_FAULT_CONTEXT_BOOT_OUTPUTS_OFF); + if (!plc_outputs_off_without_program()) + log_error("Outputs could not be driven off after the watchdog fault"); + watchdog_set_fault_context(NULL); + plc_publish_final_state(PLC_STATE_ERROR); + } + } + } + // Start the command socket only now that the plugin driver is fully built. // Everything the socket can ask for -- START, STOP, PLUGIN_CMD, STATS -- // reaches into the driver, so serving commands before this point was serving @@ -197,11 +252,6 @@ int main(int argc, char *argv[]) // finishes, causing two concurrent load_plc_program() calls — and two // dispatcher threads. plc_begin_transition() also makes the start // asynchronous, which is fine: the main thread just sleeps below. - if (safe_mode) - { - log_info("Runtime started in SAFE MODE - PLC program will not be loaded"); - log_info("Upload a corrected program to recover"); - } // Same gate as any other start, but note what it can and cannot see. A VPP // plugin that owns a physical mode switch is initialised as part of loading // the program — inside the start transition below — so at this point the @@ -212,12 +262,12 @@ int main(int argc, char *argv[]) // reconciliation stops the PLC as soon as the start lands. Safe, but the gate // only bites here for a switch position already known at this point (e.g. one // reported by a plugin the runtime loaded independently of the program). - else if (!plc_switch_allows_run()) + if (!safe_mode && !plc_switch_allows_run()) { log_info("Hardware mode switch is in STOP - PLC left stopped"); log_info("Move the switch to RUN to start the PLC"); } - else if (!plc_begin_transition(PLC_STATE_RUNNING)) + else if (!safe_mode && !plc_begin_transition(PLC_STATE_RUNNING)) { log_error("Failed to initiate PLC start"); } diff --git a/core/src/plc_app/plc_retain_file_store.cpp b/core/src/plc_app/plc_retain_file_store.cpp index 1e9f1446..74d08b15 100644 --- a/core/src/plc_app/plc_retain_file_store.cpp +++ b/core/src/plc_app/plc_retain_file_store.cpp @@ -9,6 +9,7 @@ #include "plc_retain_file_store.h" #include "plc_retain.h" // PLC_RETAIN_PROGRAM_ID_LEN — one definition for both sides +#include "utils/rt_mutex.h" #include #include @@ -43,7 +44,7 @@ static_assert(PROGRAM_ID_LEN == PLC_RETAIN_PROGRAM_ID_LEN, "identity the runtime hands to read() — a shorter or longer header would be " "indistinguishable from a torn write and every load would discard good values."); -std::mutex g_lock; +RtMutex g_lock; std::vector g_pending; bool g_dirty = false; @@ -194,7 +195,7 @@ void commit(const uint8_t *buf, uint16_t len, const std::string &program_id) void discard_stored() { { - std::lock_guard guard(g_lock); + std::lock_guard guard(g_lock); g_pending.clear(); g_dirty = false; } @@ -222,7 +223,7 @@ void flush_loop() * save() every cycle and must never wait on a disk write. The * identity is snapshotted with the bytes so the pair committed * below is the pair that was current at this instant. */ - std::lock_guard guard(g_lock); + std::lock_guard guard(g_lock); if (!g_dirty) continue; snapshot = g_pending; snapshot_id = g_program_md5; @@ -254,7 +255,7 @@ void plc_retain_file_store_stop(void) if (g_flusher.joinable()) g_flusher.join(); /* Final flush: a clean stop should not discard the last interval. */ - std::lock_guard guard(g_lock); + std::lock_guard guard(g_lock); if (g_dirty && !g_pending.empty()) { commit(g_pending.data(), (uint16_t)g_pending.size(), g_program_md5); @@ -278,7 +279,7 @@ int plc_retain_file_store_save(const uint8_t *blob, uint16_t len) { if (!g_enabled.load() || !blob || len == 0 || len > RETAIN_MAX) return -1; - std::lock_guard guard(g_lock); + std::lock_guard guard(g_lock); /* Only mark dirty on an actual change. The runtime deliberately does not * diff — it cannot know what a write costs here — so doing it at this layer * is how a slow medium avoids rewriting an unchanged blob every interval. */ @@ -301,7 +302,7 @@ int plc_retain_file_store_load(const char *program_md5, uint16_t md5_len, uint8_ * a store that just discarded a previous program's values still has to * label the new program's first commit. */ { - std::lock_guard guard(g_lock); + std::lock_guard guard(g_lock); g_program_md5.assign(program_md5, md5_len); } @@ -337,7 +338,7 @@ int plc_retain_file_store_load(const char *program_md5, uint16_t md5_len, uint8_ /* Prime the in-memory copy so the first flush after start does not rewrite * a byte-identical file. */ - std::lock_guard guard(g_lock); + std::lock_guard guard(g_lock); g_pending.assign(out, out + n); g_dirty = false; return 0; @@ -349,7 +350,7 @@ int plc_retain_file_store_flush(void) * plc_retain_file_store_stop(): the PLC can be started again without the * daemon restarting, and joining the thread here would leave the next run * with nothing committing on a timer. */ - std::lock_guard guard(g_lock); + std::lock_guard guard(g_lock); if (!g_dirty || g_pending.empty()) return 0; commit(g_pending.data(), (uint16_t)g_pending.size(), g_program_md5); g_dirty = false; diff --git a/core/src/plc_app/plc_state_manager.cpp b/core/src/plc_app/plc_state_manager.cpp index 86e01e16..148c905b 100644 --- a/core/src/plc_app/plc_state_manager.cpp +++ b/core/src/plc_app/plc_state_manager.cpp @@ -29,6 +29,7 @@ extern "C" { #include "../drivers/plugin_driver.h" +#include "unix_socket.h" } // Runtime-side strucpp ABI mirror — see core/src/lib/strucpp_abi.hpp @@ -41,17 +42,21 @@ extern "C" { #include "plc_state_manager.h" #include "plcapp_manager.h" #include "scan_cycle_manager.h" +#include "task_policy.h" #include "utils/log.h" +#include "utils/rt_mutex.h" #include "utils/utils.h" +#include "utils/watchdog.h" -static PLCState plc_state = PLC_STATE_STOPPED; -static pthread_mutex_t state_mutex = PTHREAD_MUTEX_INITIALIZER; +/* Writers serialise on state_mutex; readers load the atomic without locking. */ +static std::atomic plc_state{PLC_STATE_STOPPED}; +static_assert(std::atomic::is_always_lock_free, "state reads must not lock"); +static pthread_mutex_t state_mutex = PTHREAD_MUTEX_INITIALIZER; struct timespec timer_start; pthread_t plc_thread; PluginManager *plc_program = NULL; -extern std::atomic plc_heartbeat; extern plugin_driver_t *plugin_driver; /* ----------------------------------------------------------------------- @@ -101,6 +106,13 @@ static pthread_mutex_t done_mutex = PTHREAD_MUTEX_INITIALIZER; static pthread_cond_t done_cond; /* initialised once, CLOCK_MONOTONIC */ static pthread_once_t done_cond_once = PTHREAD_ONCE_INIT; +__attribute__((constructor)) static void state_manager_locks_init_pi(void) +{ + rt_mutex_upgrade_static(&state_mutex, "state_mutex"); + rt_mutex_upgrade_static(&plc_tasks_lock, "plc_tasks_lock"); + rt_mutex_upgrade_static(&done_mutex, "done_mutex"); +} + /* One-time init of done_cond on the CLOCK_MONOTONIC clock (the default is * CLOCK_REALTIME, which would mismatch the dispatcher's monotonic deadline and * jump under NTP). Init-once (not per-load) so a crash that skips teardown @@ -138,6 +150,22 @@ static volatile sig_atomic_t bootstrap_crash_sig = 0; static volatile sig_atomic_t bootstrap_holding_mutex = 0; static volatile sig_atomic_t plc_crash_signal = 0; +/* Interval between abort signals while waiting for an aborted task to exit. */ +#define PLC_TASK_ABORT_RESEND_NS 100000000LL + +/* Poll interval while waiting for a task thread to return. */ +#define PLC_TASK_EXIT_POLL_NS 1000000L + +/* Set when a task had to be aborted; the stop then lands ERROR, not STOPPED. */ +static std::atomic g_task_fault{false}; + +/* Longest task interval of the loaded program, for the stop budget. */ +static std::atomic g_longest_interval_ns{0}; + +/* Oldest in-flight first scan (release ns, 0 = none), and the watchdog's trip request for it. */ +static std::atomic g_first_scan_since_ns{0}; +static std::atomic g_first_scan_trip{false}; + /* The SIGUSR1 wake handler is installed once at process init in * plc_main.c (handle_sigusr1). Every task thread relies on EINTR from * pthread_kill(target, SIGUSR1) to break out of clock_nanosleep on @@ -163,6 +191,17 @@ static void plc_crash_handler(int sig) raise(sig); } +/* Jumps out of the scan body to the task's recovery point. No-op outside the scan window. */ +static void plc_abort_handler(int sig) +{ + PlcTaskCtx *ctx = current_task_ctx; + if (ctx && ctx->in_body) + { + ctx->crash_sig = sig; + siglongjmp(ctx->crash_jmp, sig); + } +} + /* Drop whichever runtime lock this task thread currently holds. Mirrors the * signal-handler recovery (the sigsetjmp block below) so a C++ exception * thrown mid-scan can't leave the image mutex locked when the thread unwinds @@ -183,40 +222,34 @@ static void plc_task_release_locks(PlcTaskCtx *ctx) } } -/* ----------------------------------------------------------------------- - * Per-task thread function. - * - * Phase 6 keeps this minimal: SCHED_FIFO priority elevation, optional - * CPU affinity, per-thread crash recovery, then a clock_nanosleep loop - * that runs task->programs[]->run() under the process-image protocol. - * Phase 7 specializes the fastest task by adding housekeeping pre/post. - * --------------------------------------------------------------------- */ -static void *plc_task_thread(void *arg) +/* Per-task thread: FIFO priority from the IEC priority, optional affinity, crash and + * watchdog-abort recovery, then one scan per release posted by the dispatcher. */ +static void plc_task_body(PlcTaskCtx *ctx) { - PlcTaskCtx *ctx = static_cast(arg); current_task_ctx = ctx; pthread_setname_np(pthread_self(), ctx->name); - /* 99 is reserved for the dispatcher, which has to be strictly above every - * worker for its tick never to be delayed by a busy one. A worker allowed to - * reach 99 would only TIE it, and SCHED_FIFO does not time-slice between equal - * priorities: a task that never blocks (an unbounded loop in IEC code) would - * then keep the dispatcher off that CPU entirely, along with anything else - * trying to bring the PLC down. */ - int rt = ctx->priority; - if (rt < 1) rt = 1; - if (rt > 98) rt = 98; + /* IEC 0 (highest) -> FIFO 49; always below the dispatcher and watchdog. */ + bool clamped = false; + int rt = plc_task_fifo_priority(ctx->priority, &clamped); + if (clamped) + { + log_warn("[task %s] IEC priority %d outside %d..%d, clamped", ctx->name, ctx->priority, + PLC_IEC_PRIORITY_MIN, PLC_IEC_PRIORITY_MAX); + } sched_param sp{}; sp.sched_priority = rt; - if (pthread_setschedparam(pthread_self(), SCHED_FIFO, &sp) != 0) + int sp_rc = pthread_setschedparam(pthread_self(), SCHED_FIFO, &sp); + if (sp_rc != 0) { log_warn("[task %s] SCHED_FIFO(%d) failed: %s — running default scheduling", - ctx->name, rt, strerror(errno)); + ctx->name, rt, strerror(sp_rc)); } else { - log_info("[task %s] SCHED_FIFO priority %d", ctx->name, rt); + log_info("[task %s] SCHED_FIFO priority %d (IEC priority %d)", ctx->name, rt, + ctx->priority); } if (ctx->cpu_affinity_mask != 0) @@ -245,24 +278,21 @@ static void *plc_task_thread(void *arg) #endif } - /* Per-thread fault recovery for HARDWARE signals (SIGSEGV/SIGFPE), e.g. a - * raw integer divide-by-zero in IEC code. The signal handler siglongjmp's - * here. Like the C++ exception path below, we isolate per task: release any - * held lock, mark this worker dead so the dispatcher stops releasing it, and - * exit this thread only. The other task threads keep running. (Bodies run on - * private storage, so a fault is contained to this task's data; the shared - * image stays protected by the journal + locks.) */ + /* Recovery point for SIGSEGV/SIGFPE from IEC code and for the watchdog abort + * (PLC_TASK_ABORT_SIGNAL): release held locks, mark the task dead, exit this thread. */ if (sigsetjmp(ctx->crash_jmp, 1) != 0) { + ctx->in_body = 0; plc_task_release_locks(ctx); ctx->alive.store(0, std::memory_order_release); - /* A hardware fault only fires from a scan body, so this thread held an - * in-flight slot — release it (and wake the dispatcher if last) so the - * completion count never leaks. */ + /* Both signals only jump from a scan body, so release the in-flight slot. */ worker_scan_done(); - log_error("[task %s] terminated by signal %d — other tasks keep running", - ctx->name, ctx->crash_sig); - return nullptr; + if (ctx->crash_sig == PLC_TASK_ABORT_SIGNAL) + log_error("[task %s] scan aborted by the watchdog", ctx->name); + else + log_error("[task %s] terminated by signal %d — other tasks keep running", + ctx->name, ctx->crash_sig); + return; } auto *task = static_cast(ctx->task_handle); @@ -323,8 +353,11 @@ static void *plc_task_thread(void *arg) /* 2. Run the bodies. Shared-global access self-serializes on each * global's own mutex inside run() (strucpp GlobalVar); no * runtime-owned global lock and no private copy-in/out. */ + /* The abort may jump only out of IEC code, never out of a runtime lock window. */ + ctx->in_body = 1; for (size_t p = 0; p < task->program_count; ++p) task->programs[p]->run(); + ctx->in_body = 0; /* 3. Copy-out: journal changed located outputs (lock-free; applied * to the image on the next drain — the dispatcher's frame top). */ @@ -337,26 +370,27 @@ static void *plc_task_thread(void *arg) } catch (const std::exception &e) { + ctx->in_body = 0; plc_task_release_locks(ctx); ctx->alive.store(0, std::memory_order_release); worker_scan_done(); /* release the in-flight slot before exiting */ log_error("[task %s] terminated by unhandled exception: %s — " "other tasks keep running", ctx->name, e.what()); - return nullptr; + return; } catch (...) { + ctx->in_body = 0; plc_task_release_locks(ctx); ctx->alive.store(0, std::memory_order_release); worker_scan_done(); /* release the in-flight slot before exiting */ log_error("[task %s] terminated by unknown exception — " "other tasks keep running", ctx->name); - return nullptr; + return; } scan_cycle_tracker_end(&ctx->tracker); - ctx->heartbeat.store((long)time(nullptr), std::memory_order_relaxed); ctx->local_tick.fetch_add(1, std::memory_order_relaxed); /* Signal scan completion LAST (release order): the dispatcher reads * completed vs released to decide whether this worker is idle (safe to @@ -370,37 +404,119 @@ static void *plc_task_thread(void *arg) log_info("[task %s] stopped after %llu scans", ctx->name, (unsigned long long)ctx->local_tick.load()); +} + +static void *plc_task_thread(void *arg) +{ + PlcTaskCtx *ctx = static_cast(arg); + plc_task_body(ctx); + ctx->exited.store(1, std::memory_order_release); return nullptr; } -/* Wake every worker, join them, and destroy the task array. - * - * Two callers: the normal end of the dispatcher loop, and the early-out below - * when bring-up finished but this start is no longer the transition in flight. - * The second one exists so that path tears its workers down instead of leaving - * them parked on a semaphore nobody will ever post again. */ +static bool task_in_scan(const PlcTaskCtx *c) +{ + return c->released.load(std::memory_order_relaxed) != + c->completed.load(std::memory_order_acquire); +} + +static int64_t task_stuck_limit_ns(const PlcTaskCtx *c) +{ + return PLC_TASK_STUCK_PERIODS * c->interval_ns; +} + +static bool task_first_scan_done(const PlcTaskCtx *c) +{ + return c->completed.load(std::memory_order_acquire) > 0; +} + +/* How long the current scan may run after its release before the teardown aborts it. */ +static int64_t task_scan_limit_ns(const PlcTaskCtx *c) +{ + const int64_t first_ns = PLC_FIRST_SCAN_TIMEOUT_MS * NS_PER_MS; + if (!task_first_scan_done(c) && first_ns > task_stuck_limit_ns(c)) + return first_ns; + return task_stuck_limit_ns(c); +} + +/* Waits for the worker to return. An in-flight scan may run until scan_deadline; an idle + * worker (already woken) until idle_deadline. Returns true when it exited. */ +static bool wait_task_exit(const PlcTaskCtx *c, int64_t scan_deadline, int64_t idle_deadline) +{ + const timespec poll = {0, PLC_TASK_EXIT_POLL_NS}; + while (!c->exited.load(std::memory_order_acquire)) + { + const int64_t now = monotonic_ns(); + if (task_in_scan(c) ? now >= scan_deadline : now >= idle_deadline) + return false; + nanosleep(&poll, nullptr); + } + return true; +} + +/* Signals the worker until it leaves its scan; exits the process if it never does. */ +static void abort_task(PlcTaskCtx *c) +{ + if (task_in_scan(c)) + log_error("[task %s] scan still running %lld ms after its release: aborting it", c->name, + (long long)(task_scan_limit_ns(c) / NS_PER_MS)); + else + log_error("[task %s] did not exit within %d ms of being woken", c->name, + PLC_TASK_ABORT_TIMEOUT_MS); + g_task_fault.store(true, std::memory_order_release); + + const timespec poll = {0, PLC_TASK_EXIT_POLL_NS}; + const int64_t give_up = monotonic_ns() + PLC_TASK_ABORT_TIMEOUT_MS * NS_PER_MS; + int64_t next_signal = 0; + while (!c->exited.load(std::memory_order_acquire)) + { + const int64_t now = monotonic_ns(); + if (now >= give_up) + { + char reason[WATCHDOG_REASON_LEN]; + std::snprintf(reason, sizeof reason, "task %s did not exit after the abort", c->name); + watchdog_fatal_exit(reason); + } + if (now >= next_signal) + { + int rc = pthread_kill(c->thread, PLC_TASK_ABORT_SIGNAL); + if (rc != 0) + log_error("[task %s] abort signal failed: %s", c->name, strerror(rc)); + next_signal = now + PLC_TASK_ABORT_RESEND_NS; + } + nanosleep(&poll, nullptr); + } +} + +/* Wakes every worker, lets each in-flight scan run until PLC_TASK_STUCK_PERIODS of its own + * periods after its release, aborts the ones still running, then joins and frees the array. */ static void reap_task_threads(void) { log_info("Stopping %zu PLC task thread(s)", plc_task_count); - /* Wake every worker: post its release semaphore (breaks sem_wait) and - * SIGUSR1 (breaks a syscall). A worker mid-scan finishes, loops to - * sem_wait, consumes the post, observes state != RUNNING, and exits. */ for (size_t i = 0; i < plc_task_count; ++i) { sem_post(&plc_tasks[i].go); pthread_kill(plc_tasks[i].thread, SIGUSR1); } + + const int64_t t_reap = monotonic_ns(); + for (size_t i = 0; i < plc_task_count; ++i) + { + PlcTaskCtx *c = &plc_tasks[i]; + const int64_t scan_deadline = + c->release_ns.load(std::memory_order_acquire) + task_scan_limit_ns(c); + const int64_t idle_deadline = (scan_deadline > t_reap ? scan_deadline : t_reap) + + PLC_TASK_ABORT_TIMEOUT_MS * NS_PER_MS; + if (!wait_task_exit(c, scan_deadline, idle_deadline)) + abort_task(c); + } + for (size_t i = 0; i < plc_task_count; ++i) { pthread_join(plc_tasks[i].thread, nullptr); } - /* Take plc_tasks_lock for the tracker-cleanup + free. A STATS reader - * that started iterating before STOP arrived will block briefly - * waiting for this critical section, then exit because plc_task_count - * is observed as 0. Without the lock, the reader could be midway - * through scan_cycle_tracker_snapshot when we pthread_mutex_destroy - * the tracker's own mutex below — undefined behaviour. */ + /* Under plc_tasks_lock so a STATS reader never sees a destroyed tracker. */ pthread_mutex_lock(&plc_tasks_lock); for (size_t i = 0; i < plc_task_count; ++i) { @@ -413,6 +529,25 @@ static void reap_task_threads(void) pthread_mutex_unlock(&plc_tasks_lock); } +/* Zeroes every output and runs one last I/O frame, then holds it so threaded plugins send it. */ +static void plc_outputs_off(void) +{ + image_lock(); + /* Pending writes were drained by image_lock; reject later ones so none re-energise %Q. */ + journal_cleanup(); + image_tables_zero_outputs(); + image_unlock(); + if (plugin_driver) + { + plugin_driver_cycle_start(plugin_driver); + plugin_driver_cycle_end(plugin_driver); + } + timespec settle = {PLC_OUTPUTS_OFF_SETTLE_MS / 1000, + (long)((PLC_OUTPUTS_OFF_SETTLE_MS % 1000) * NS_PER_MS)}; + nanosleep(&settle, nullptr); + log_info("Outputs forced to 0"); +} + void *plc_cycle_thread(void *arg) { PluginManager *pm = (PluginManager *)arg; @@ -420,6 +555,10 @@ void *plc_cycle_thread(void *arg) plc_crash_signal = 0; bootstrap_crash_sig = 0; bootstrap_holding_mutex = 0; + g_task_fault.store(false, std::memory_order_release); + g_first_scan_since_ns.store(0, std::memory_order_release); + g_first_scan_trip.store(false, std::memory_order_release); + watchdog_dispatcher_stopped(); /* Per-task trackers are initialised below, once we know the task list * and each task's interval. */ @@ -498,6 +637,13 @@ void *plc_cycle_thread(void *arg) sigaction(SIGFPE, &crash_sa, NULL); sigaction(SIGSEGV, &crash_sa, NULL); + struct sigaction abort_sa; + std::memset(&abort_sa, 0, sizeof(abort_sa)); + abort_sa.sa_handler = plc_abort_handler; + sigemptyset(&abort_sa.sa_mask); + if (sigaction(PLC_TASK_ABORT_SIGNAL, &abort_sa, NULL) != 0) + log_error("Failed to install the task abort handler: %s", strerror(errno)); + /* SIGUSR1 wake handler is installed once at process init (plc_main.c). * No per-thread re-installation here — the bootstrap thread inherits * the handler from the process. */ @@ -538,38 +684,17 @@ void *plc_cycle_thread(void *arg) pthread_mutex_unlock(&state_mutex); log_info("PLC State: ERROR"); - /* If the crash happened in the dispatcher loop (after workers were - * spawned), the workers are still alive — wake, join, and free them so - * they aren't orphaned (which would UAF on the next load). The state is - * already ERROR, so each worker exits after its current scan. */ + /* Workers spawned before a dispatcher crash are still alive: reap them (aborting any + * still scanning) so they are not orphaned, then drive outputs off. */ if (plc_tasks && plc_task_count) { - for (size_t i = 0; i < plc_task_count; ++i) - { - sem_post(&plc_tasks[i].go); - pthread_kill(plc_tasks[i].thread, SIGUSR1); - } - for (size_t i = 0; i < plc_task_count; ++i) - pthread_join(plc_tasks[i].thread, nullptr); - pthread_mutex_lock(&plc_tasks_lock); - for (size_t i = 0; i < plc_task_count; ++i) - { - scan_cycle_tracker_cleanup(&plc_tasks[i].tracker); - sem_destroy(&plc_tasks[i].go); - } - std::free(plc_tasks); - plc_tasks = nullptr; - plc_task_count = 0; - pthread_mutex_unlock(&plc_tasks_lock); + reap_task_threads(); + plc_outputs_off(); } return NULL; } - /* Walk the configuration via virtual dispatch and discover the GCD - * base tick + flat task list. Phase 5 keeps a single-thread cycle - * that runs every task in round-robin (each task runs every - * interval/base ticks). Phase 6 will replace this with one thread per - * task on SCHED_FIFO. */ + /* Walk the configuration via virtual dispatch: GCD base tick and flat task list. */ auto *cfg = static_cast(strucpp_config_handle()); if (!cfg) { @@ -646,7 +771,6 @@ void *plc_cycle_thread(void *arg) { size_t flat_idx = 0; - long now_t = (long)time(nullptr); for (size_t r = 0; r < cfg->get_resource_count(); ++r) { for (size_t t = 0; t < resources[r].task_count; ++t) @@ -667,7 +791,6 @@ void *plc_cycle_thread(void *arg) { std::snprintf(ctx->name, sizeof ctx->name, "plc-task-%zu", flat_idx); } - ctx->heartbeat.store(now_t, std::memory_order_relaxed); ctx->local_tick.store(0, std::memory_order_relaxed); /* Dispatcher plumbing. divisor = interval / base_tick (exact; @@ -682,6 +805,10 @@ void *plc_cycle_thread(void *arg) ctx->released.store(0, std::memory_order_relaxed); ctx->completed.store(0, std::memory_order_relaxed); ctx->overrun_count.store(0, std::memory_order_relaxed); + ctx->stuck_ticks.store(0, std::memory_order_relaxed); + ctx->exited.store(0, std::memory_order_relaxed); + ctx->release_ns.store(0, std::memory_order_relaxed); + ctx->in_body = 0; if (scan_cycle_tracker_init(&ctx->tracker, ctx->interval_ns) != 0) { @@ -692,6 +819,13 @@ void *plc_cycle_thread(void *arg) } } + { + int64_t longest = 0; + for (size_t i = 0; i < plc_task_count; ++i) + if (plc_tasks[i].interval_ns > longest) longest = plc_tasks[i].interval_ns; + g_longest_interval_ns.store(longest, std::memory_order_release); + } + /* Pick the fastest task: smallest interval, tie-break by priority, * then by declaration order (which is the iteration order above). */ { @@ -708,7 +842,7 @@ void *plc_cycle_thread(void *arg) } plc_tasks[fastest_idx].is_fastest_task = true; /* Housekeeping no longer rides a real task — the GCD master-tick - * dispatcher owns time/cycle hooks/heartbeat. is_fastest_task is kept + * dispatcher owns time/cycle hooks/watchdog feed. is_fastest_task is kept * only as a STATS hint (the fastest task is the tightest schedule). */ log_info("PLC: fastest task is %s (interval=%lld ns, priority=%d)", plc_tasks[fastest_idx].name, @@ -720,15 +854,14 @@ void *plc_cycle_thread(void *arg) * * Failure mode: if pthread_create succeeds for tasks 0..i-1 and then * fails for task i, the previously-spawned threads are running at - * SCHED_FIFO 99 holding image_tables_mutex and reading from + * SCHED_FIFO (1..PLC_FIFO_TASK_MAX) and reading from * plc_tasks[]. Returning here without cleanup leaves them orphaned * — the next load_plc_program reallocates plc_tasks and the old * threads dereference freed memory. We must: * * 1) flip plc_state to ERROR so the surviving task threads exit * their `while (state == RUNNING)` loop on their next iteration; - * 2) SIGUSR1 each surviving thread to break it out of its - * clock_nanosleep without waiting up to interval_ns; + * 2) post each surviving thread's release semaphore so it wakes; * 3) join all spawned threads before freeing the array. * * After this rollback, plc_tasks is nullptr and plc_task_count is 0, @@ -810,7 +943,7 @@ void *plc_cycle_thread(void *arg) * * This thread is the single time authority. It wakes every base_ns on an * absolute deadline anchored at one t0, and on each tick: - * - bumps the global heartbeat (every tick, so the watchdog sees us); + * - feeds the watchdog (watchdog_feed, every tick); * - computes the due set (task due iff masterTick % divisor == 0); * - on a task-bearing tick: drains the journal (committing the previous * frame's outputs), fires cycle_end (prev frame) then cycle_start (new @@ -820,14 +953,15 @@ void *plc_cycle_thread(void *arg) * worker still running when re-due is an overrun and is simply not * re-released that tick. A faulted worker (alive==0) is skipped forever. * - * Run at SCHED_FIFO 99 — above every worker — so the tick is never delayed - * by a busy worker on a shared CPU. + * Runs at PLC_FIFO_DISPATCHER: above every worker, below the watchdog. * --------------------------------------------------------------------- */ { + pthread_setname_np(pthread_self(), "plc_dispatch"); sched_param dsp{}; - dsp.sched_priority = 99; - if (pthread_setschedparam(pthread_self(), SCHED_FIFO, &dsp) != 0) - log_warn("dispatcher SCHED_FIFO(99) failed: %s", strerror(errno)); + dsp.sched_priority = PLC_FIFO_DISPATCHER; + int rc = pthread_setschedparam(pthread_self(), SCHED_FIFO, &dsp); + if (rc != 0) + log_warn("dispatcher SCHED_FIFO(%d) failed: %s", PLC_FIFO_DISPATCHER, strerror(rc)); } /* Completion-signal condvar shares the CLOCK_MONOTONIC timeline with the @@ -839,9 +973,12 @@ void *plc_cycle_thread(void *arg) log_info("GCD master-tick dispatcher running (base tick %llu ns)", (unsigned long long)base_ns); + watchdog_dispatcher_started((int64_t)base_ns); uint64_t master_tick = 0; bool cycle_end_pending = false; /* a frame's cycle_end not yet fired */ + size_t stuck_idx = SIZE_MAX; + bool fault_stop_claimed = false; timespec next_tick; clock_gettime(CLOCK_MONOTONIC, &next_tick); @@ -850,8 +987,22 @@ void *plc_cycle_thread(void *arg) /* ---- Phase B: the tick (runs at the absolute deadline) ---- */ const int64_t master_time = (int64_t)master_tick * (int64_t)base_ns; - /* Always: feed the global watchdog. */ - plc_heartbeat.store((long)time(nullptr), std::memory_order_relaxed); + watchdog_feed(); + + { + int64_t oldest = 0; + for (size_t i = 0; i < plc_task_count; ++i) + { + PlcTaskCtx *c = &plc_tasks[i]; + if (c->alive.load(std::memory_order_acquire) && task_in_scan(c) && + !task_first_scan_done(c)) + { + const int64_t rel = c->release_ns.load(std::memory_order_acquire); + if (oldest == 0 || rel < oldest) oldest = rel; + } + } + g_first_scan_since_ns.store(oldest, std::memory_order_release); + } /* Which tasks are due this tick? */ bool any_due = false; @@ -911,6 +1062,8 @@ void *plc_cycle_thread(void *arg) * and release. The fetch_add must happen-before sem_post so * the worker's matching worker_scan_done() can never drive * g_tasks_running negative. */ + c->stuck_ticks.store(0, std::memory_order_relaxed); + c->release_ns.store(monotonic_ns(), std::memory_order_release); c->time_at_dispatch = master_time; c->released.store(r + 1, std::memory_order_relaxed); g_tasks_running.fetch_add(1, std::memory_order_acq_rel); @@ -923,6 +1076,13 @@ void *plc_cycle_thread(void *arg) * NOT re-release (binary), so activations never pile up. The * task simply runs at a lower effective rate; the others are * unaffected. Rate-limit the log. */ + /* Ticks, not wall time: replayed ticks after a late dispatcher are missed deadlines too. */ + if (task_first_scan_done(c)) + { + long st = c->stuck_ticks.fetch_add(1, std::memory_order_relaxed) + 1; + if (st >= PLC_TASK_STUCK_PERIODS && stuck_idx == SIZE_MAX) + stuck_idx = i; + } long oc = c->overrun_count.fetch_add(1, std::memory_order_relaxed) + 1; if (oc == 1 || (oc % 50) == 0) log_warn("[task %s] scan overrun #%ld: body exceeds its " @@ -936,6 +1096,39 @@ void *plc_cycle_thread(void *arg) /* This frame owes a cycle_end once its released tasks all finish. */ if (released_any) cycle_end_pending = true; ++scan_counter; + + bool first_scan_trip = false; + if (stuck_idx == SIZE_MAX && g_first_scan_trip.exchange(false, std::memory_order_acq_rel)) + { + for (size_t i = 0; i < plc_task_count; ++i) + { + PlcTaskCtx *c = &plc_tasks[i]; + if (!c->alive.load(std::memory_order_acquire) || !task_in_scan(c) || + task_first_scan_done(c)) + continue; + if (stuck_idx == SIZE_MAX || + c->release_ns.load(std::memory_order_acquire) < + plc_tasks[stuck_idx].release_ns.load(std::memory_order_acquire)) + stuck_idx = i; + } + first_scan_trip = (stuck_idx != SIZE_MAX); + } + + if (stuck_idx != SIZE_MAX) + { + PlcTaskCtx *c = &plc_tasks[stuck_idx]; + if (first_scan_trip) + log_error("[task %s] first scan still running after %d ms: stopping the PLC", + c->name, PLC_FIRST_SCAN_TIMEOUT_MS); + else + log_error("[task %s] stuck in one scan for %d periods (%lld ms): stopping the PLC", + c->name, PLC_TASK_STUCK_PERIODS, + (long long)(task_stuck_limit_ns(c) / NS_PER_MS)); + g_task_fault.store(true, std::memory_order_release); + /* Claim now so no command lands mid-drain; the teardown runs after the reap. */ + fault_stop_claimed = plc_claim_transition(PLC_STATE_STOPPED); + break; + } } ++master_tick; @@ -1000,11 +1193,18 @@ void *plc_cycle_thread(void *arg) pthread_mutex_unlock(&done_mutex); } + watchdog_dispatcher_stopped(); + g_first_scan_since_ns.store(0, std::memory_order_release); reap_task_threads(); + plc_outputs_off(); signal(SIGFPE, SIG_DFL); signal(SIGSEGV, SIG_DFL); + /* The stop's unload joins this thread, so it must run on another one. */ + if (fault_stop_claimed && !plc_complete_claimed_transition_async(PLC_STATE_STOPPED)) + log_error("Could not start the fault stop; the PLC stays in TRANSITIONING_TO_STOP"); + return NULL; } @@ -1119,7 +1319,25 @@ extern "C" int load_plc_program(PluginManager *pm) } } +/* Serialises unloads: a shutdown and a fault stop may both reach here; the second finds nothing. */ +static pthread_mutex_t unload_mutex = PTHREAD_MUTEX_INITIALIZER; + +__attribute__((constructor)) static void unload_mutex_init_pi(void) +{ + rt_mutex_upgrade_static(&unload_mutex, "unload_mutex"); +} + +static int unload_plc_program_locked(PluginManager *pm); + extern "C" int unload_plc_program(PluginManager *pm) +{ + pthread_mutex_lock(&unload_mutex); + int rc = unload_plc_program_locked(pm); + pthread_mutex_unlock(&unload_mutex); + return rc; +} + +static int unload_plc_program_locked(PluginManager *pm) { if (pm && pm == plc_program) { @@ -1173,9 +1391,11 @@ extern "C" int unload_plc_program(PluginManager *pm) log_info("PLC program unloaded successfully"); - /* The teardown is done, so this is the moment STOPPED becomes true. - * plc_publish_final_state keeps ERROR if a task crashed on the way out. */ - plc_publish_final_state(PLC_STATE_STOPPED); + /* An aborted task makes the stop a fault: land ERROR. ERROR also survives STOPPED. */ + if (g_task_fault.exchange(false, std::memory_order_acq_rel)) + plc_publish_final_state(PLC_STATE_ERROR); + else + plc_publish_final_state(PLC_STATE_STOPPED); return 0; } else @@ -1187,10 +1407,61 @@ extern "C" int unload_plc_program(PluginManager *pm) extern "C" PLCState plc_get_state(void) { - pthread_mutex_lock(&state_mutex); - PLCState s = plc_state; - pthread_mutex_unlock(&state_mutex); - return s; + return plc_state.load(std::memory_order_acquire); +} + +extern "C" int64_t plc_stop_budget_ms(void) +{ + int64_t grace_ms = + PLC_TASK_STUCK_PERIODS * g_longest_interval_ns.load(std::memory_order_acquire) / NS_PER_MS; + if (grace_ms < PLC_FIRST_SCAN_TIMEOUT_MS) + grace_ms = PLC_FIRST_SCAN_TIMEOUT_MS; + return grace_ms + PLC_OUTPUTS_OFF_SETTLE_MS + PLC_STOP_TEARDOWN_ALLOWANCE_MS; +} + +extern "C" bool plc_outputs_off_without_program(void) +{ + if (!plugin_driver) + { + log_warn("No plugin driver: outputs cannot be driven off"); + return false; + } + if (plugin_driver_update_config(plugin_driver, "./plugins.conf") != 0 || + plugin_driver_append_config(plugin_driver, "./vpp_plugins.conf") != 0) + { + log_error("[PLUGIN]: Could not load the plugin configuration to drive outputs off"); + return false; + } + if (plugin_driver_init(plugin_driver) != 0) + { + plugin_driver_cleanup_init(plugin_driver); + log_error("[PLUGIN]: Plugin init failed: outputs cannot be driven off"); + return false; + } + + pthread_mutex_t *itm = image_tables_mutex(); + pthread_mutex_lock(itm); + image_tables_fill_null_pointers(); + pthread_mutex_unlock(itm); + + plugin_driver_start(plugin_driver); + plc_outputs_off(); + plugin_driver_stop(plugin_driver); + + pthread_mutex_lock(itm); + image_tables_clear_null_pointers(); + pthread_mutex_unlock(itm); + return true; +} + +extern "C" int64_t plc_first_scan_pending_since_ns(void) +{ + return g_first_scan_since_ns.load(std::memory_order_acquire); +} + +extern "C" void plc_request_first_scan_trip(void) +{ + g_first_scan_trip.store(true, std::memory_order_release); } extern "C" bool plc_state_is_transitioning(void) diff --git a/core/src/plc_app/plc_state_manager.h b/core/src/plc_app/plc_state_manager.h index e1d0c129..eca51f26 100644 --- a/core/src/plc_app/plc_state_manager.h +++ b/core/src/plc_app/plc_state_manager.h @@ -21,11 +21,13 @@ #include typedef std::atomic plc_atomic_long_t; typedef std::atomic plc_atomic_u64_t; +typedef std::atomic plc_atomic_i64_t; extern "C" { #else #include typedef atomic_long plc_atomic_long_t; typedef atomic_uint_least64_t plc_atomic_u64_t; +typedef atomic_int_least64_t plc_atomic_i64_t; #endif /** @@ -109,12 +111,15 @@ typedef struct PlcTaskCtx plc_atomic_long_t released; plc_atomic_long_t completed; plc_atomic_long_t overrun_count; + plc_atomic_long_t stuck_ticks; /* consecutive due ticks found still in one scan */ + plc_atomic_long_t exited; /* 1 once the thread function has returned */ + plc_atomic_i64_t release_ns; /* CLOCK_MONOTONIC time of the last release */ sigjmp_buf crash_jmp; volatile sig_atomic_t crash_sig; volatile sig_atomic_t holding_mutex; /* image-tables mutex held (crash unlock) */ + volatile sig_atomic_t in_body; /* inside IEC program code, where the abort may jump out */ - plc_atomic_long_t heartbeat; plc_atomic_u64_t local_tick; /* Per-task scan/cycle/latency tracker. Each task thread updates its @@ -210,20 +215,46 @@ void plc_publish_final_state(PLCState final_state); */ bool plc_publish_running_if_claimed(void); -/** @brief True while a transition is in flight (either direction). */ -bool plc_state_is_transitioning(void); +/** + * @brief Longest a stop may take before the watchdog exits the process. + * + * PLC_TASK_STUCK_PERIODS times the longest task interval of the loaded program, + * plus the outputs-off settle time and PLC_STOP_TEARDOWN_ALLOWANCE_MS. + */ +int64_t plc_stop_budget_ms(void); -/* How long a state change may plausibly take before something is wrong. +/** + * @brief Drive every output to 0 with no program loaded. + * + * Used on the safe-mode boot after a watchdog exit: loads and starts the configured + * plugins (including VPP board plugins), runs the same outputs-off sequence as a + * stop, then stops the plugins again. The caller must hold a claimed stop. * - * ONE bound, two consumers, deliberately ordered: transition_worker stops waiting - * to observe the landing at PLC_TRANSITION_LANDING_TIMEOUT_MS, and the watchdog - * forces ERROR strictly later. Two independent numbers is how the watchdog came to - * fire 30 s before the runtime itself had given up -- ending a transition while its - * worker was still executing it. + * @return true when the plugins were brought up and the zeroed outputs were pushed + */ +bool plc_outputs_off_without_program(void); + +/** + * @brief CLOCK_MONOTONIC release time (ns) of the oldest first scan still running, or 0. * - * Generous on purpose: a start brings plugins up (SPI base scans, fieldbus probes, - * certificate generation) and a stop joins task threads. The bound is here to catch - * a transition that will never finish, not to police a slow one. */ + * Read by the watchdog to bound first scans, which the dispatcher does not count. + */ +int64_t plc_first_scan_pending_since_ns(void); + +/** + * @brief Ask the dispatcher to trip on the oldest first scan still running. + * + * Called by the watchdog when that scan exceeds PLC_FIRST_SCAN_TIMEOUT_MS. The + * dispatcher then stops the PLC exactly as for a stuck task. + */ +void plc_request_first_scan_trip(void); + +/** @brief True while a transition is in flight (either direction). */ +bool plc_state_is_transitioning(void); + +/* transition_worker stops waiting for a landing at PLC_TRANSITION_LANDING_TIMEOUT_MS. + * The watchdog forces ERROR on a start stuck past PLC_TRANSITION_STUCK_TIMEOUT_MS; a + * stuck stop is bounded by plc_stop_budget_ms() and ends in watchdog_fatal_exit(). */ #define PLC_TRANSITION_LANDING_TIMEOUT_MS 90000 #define PLC_TRANSITION_STUCK_TIMEOUT_MS (PLC_TRANSITION_LANDING_TIMEOUT_MS + 30000) diff --git a/core/src/plc_app/python_loader.c b/core/src/plc_app/python_loader.c index cf7451bc..732c9830 100644 --- a/core/src/plc_app/python_loader.c +++ b/core/src/plc_app/python_loader.c @@ -22,6 +22,7 @@ #include #include "include/iec_python.h" +#include "task_policy.h" // Function pointers for logging - set by python_loader_set_loggers() // These are always initialized by symbols_init() before any Python FB code runs @@ -137,9 +138,33 @@ int create_shm_name(char *buf, size_t size) return 0; } +static int python_block_loader_unmasked(const char *script_name, const char *script_content, + char *shm_name, size_t shm_in_size, size_t shm_out_size, + void **shm_in_ptr, void **shm_out_ptr, pid_t pid); + +/* The watchdog abort must not jump out of fork/shm/stdio; it is resent until delivered. */ int python_block_loader(const char *script_name, const char *script_content, char *shm_name, size_t shm_in_size, size_t shm_out_size, void **shm_in_ptr, void **shm_out_ptr, pid_t pid) +{ + sigset_t abort_set, old_set; + sigemptyset(&abort_set); + sigaddset(&abort_set, PLC_TASK_ABORT_SIGNAL); + int mask_rc = pthread_sigmask(SIG_BLOCK, &abort_set, &old_set); + if (mask_rc != 0) + LOG_ERROR("[Python loader] pthread_sigmask failed: %s", strerror(mask_rc)); + + int rc = python_block_loader_unmasked(script_name, script_content, shm_name, shm_in_size, + shm_out_size, shm_in_ptr, shm_out_ptr, pid); + + if (mask_rc == 0) + pthread_sigmask(SIG_SETMASK, &old_set, NULL); + return rc; +} + +static int python_block_loader_unmasked(const char *script_name, const char *script_content, + char *shm_name, size_t shm_in_size, size_t shm_out_size, + void **shm_in_ptr, void **shm_out_ptr, pid_t pid) { char shm_in_name[256]; char shm_out_name[256]; @@ -273,6 +298,12 @@ int python_block_loader(const char *script_name, const char *script_content, cha dup2(pipefd[1], STDERR_FILENO); close(pipefd[1]); + // The abort mask set by python_block_loader must not leak into the Python process. + sigset_t abort_set; + sigemptyset(&abort_set); + sigaddset(&abort_set, PLC_TASK_ABORT_SIGNAL); + sigprocmask(SIG_UNBLOCK, &abort_set, NULL); + // Execute Python with unbuffered output execlp("python3", "python3", "-u", script_name, (char *)NULL); diff --git a/core/src/plc_app/task_policy.c b/core/src/plc_app/task_policy.c new file mode 100644 index 00000000..2beca914 --- /dev/null +++ b/core/src/plc_app/task_policy.c @@ -0,0 +1,23 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Autonomy® + +#include "task_policy.h" + +int plc_task_fifo_priority(int iec_priority, bool *clamped) +{ + int p = iec_priority; + bool out = false; + if (p < PLC_IEC_PRIORITY_MIN) + { + p = PLC_IEC_PRIORITY_MIN; + out = true; + } + else if (p > PLC_IEC_PRIORITY_MAX) + { + p = PLC_IEC_PRIORITY_MAX; + out = true; + } + if (clamped) + *clamped = out; + return PLC_FIFO_TASK_MAX - p; +} diff --git a/core/src/plc_app/task_policy.h b/core/src/plc_app/task_policy.h new file mode 100644 index 00000000..0fcb0388 --- /dev/null +++ b/core/src/plc_app/task_policy.h @@ -0,0 +1,85 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Autonomy® + +#ifndef TASK_POLICY_H +#define TASK_POLICY_H + +#include +#include +#include + +#ifdef __cplusplus +extern "C" +{ +#endif + +/* SCHED_FIFO priority of the main watchdog thread. Above everything else. */ +#define PLC_FIFO_WATCHDOG 99 + +/* SCHED_FIFO priority of the GCD master-tick dispatcher. */ +#define PLC_FIFO_DISPATCHER 98 + +/* Highest SCHED_FIFO priority an IEC task may get. Kept below the PREEMPT_RT + * default of 50 for threaded IRQ handlers. */ +#define PLC_FIFO_TASK_MAX 49 + +/* A task still in one scan after this many of its own periods is stuck. Also + * the grace each task gets to finish its scan when the PLC stops. */ +#define PLC_TASK_STUCK_PERIODS 10 + +/* A task's first scan is exempt from the stuck count; the main watchdog trips it only + * past this bound, so a hung initialisation is not permanent. */ +#define PLC_FIRST_SCAN_TIMEOUT_MS 10000 + +/* How long zeroed outputs are held on stop before plugins are stopped, so + * plugins polling the image from their own threads can send them. */ +#define PLC_OUTPUTS_OFF_SETTLE_MS 500 + +/* How long an aborted task, or an idle one that was woken, may take to exit + * before the runtime gives up and exits with PLC_EXIT_WATCHDOG_FAULT. */ +#define PLC_TASK_ABORT_TIMEOUT_MS 2000 + +/* Fixed allowance for the stop teardown (plugin stop, unload) on top of the + * task grace. Past the whole budget the watchdog exits the process. */ +#define PLC_STOP_TEARDOWN_ALLOWANCE_MS 30000 + +/* A dispatcher with no tick for max(10 base ticks, this) is stalled. */ +#define PLC_DISPATCHER_STALL_MIN_MS 1000 + +/* Process exit code for an unrecoverable watchdog fault. The webserver + * (webserver/runtimemanager.py) restarts the runtime in safe mode on it. */ +#define PLC_EXIT_WATCHDOG_FAULT 42 + +/* Written before the watchdog exit and consumed at the next boot, which then + * starts in safe mode reporting ERROR even when the exit code was not seen. */ +#define PLC_WATCHDOG_FAULT_MARKER "/run/runtime/watchdog_fault" + +/* Fault context written into the marker when the safe-mode boot's outputs-off + * itself hangs; the next boot then skips it instead of looping. */ +#define PLC_FAULT_CONTEXT_BOOT_OUTPUTS_OFF "safe-mode boot outputs-off" + +/* Signal the teardown sends to a task still in its scan after its grace. Runtime + * code reached from IEC bodies blocks it around fork/allocation-heavy sections. */ +#define PLC_TASK_ABORT_SIGNAL SIGUSR2 + +/* IEC TASK priority range accepted by the runtime. 0 is the highest. */ +#define PLC_IEC_PRIORITY_MIN 0 +#define PLC_IEC_PRIORITY_MAX 48 + + /** + * @brief Map an IEC 61131-3 TASK priority to a SCHED_FIFO priority. + * + * IEC priority 0 is the highest and maps to PLC_FIFO_TASK_MAX (49); IEC + * priority 48 maps to 1. Values outside 0..48 are clamped to the nearest end. + * + * @param iec_priority priority declared on the IEC TASK + * @param clamped set to true when iec_priority was out of range; may be NULL + * @return SCHED_FIFO priority in 1..PLC_FIFO_TASK_MAX + */ + int plc_task_fifo_priority(int iec_priority, bool *clamped); + +#ifdef __cplusplus +} +#endif + +#endif // TASK_POLICY_H diff --git a/core/src/plc_app/unix_socket.c b/core/src/plc_app/unix_socket.c index c44b9daf..e9506b41 100644 --- a/core/src/plc_app/unix_socket.c +++ b/core/src/plc_app/unix_socket.c @@ -3,6 +3,7 @@ #include #include +#include #include #include #include @@ -130,6 +131,46 @@ static void *transition_worker(void *arg) return NULL; } +static bool spawn_transition_worker(PLCState target) +{ + PLCState *arg = malloc(sizeof(PLCState)); + if (!arg) + { + log_error("Failed to allocate transition argument"); + return false; + } + *arg = target; + + /* Explicit SCHED_OTHER: the dispatcher (FIFO 98) also spawns this worker for a fault stop. */ + pthread_attr_t attr; + int rc = pthread_attr_init(&attr); + if (rc != 0) + { + log_error("Failed to init transition thread attributes (%s)", strerror(rc)); + free(arg); + return false; + } + struct sched_param sp = {.sched_priority = 0}; + if (pthread_attr_setinheritsched(&attr, PTHREAD_EXPLICIT_SCHED) != 0 || + pthread_attr_setschedpolicy(&attr, SCHED_OTHER) != 0 || + pthread_attr_setschedparam(&attr, &sp) != 0) + { + log_warn("Transition thread inherits the caller's scheduling"); + } + + pthread_t tid; + rc = pthread_create(&tid, &attr, transition_worker, arg); + pthread_attr_destroy(&attr); + if (rc != 0) + { + log_error("Failed to create transition thread (%s)", strerror(rc)); + free(arg); + return false; + } + pthread_detach(tid); + return true; +} + // Start a background thread that performs the (potentially slow) state // transition. Returns false when the request was refused; otherwise the // transition is under way (or, if the worker could not be spawned, has already @@ -173,27 +214,19 @@ bool plc_begin_transition(PLCState target) // Completing it here blocks this caller for the duration -- the socket is // single-client, so the editor waits -- which on a thread-or-memory exhaustion // path is the cheaper of the two costs by a wide margin. - PLCState *arg = malloc(sizeof(PLCState)); - if (!arg) - { - log_error("Failed to allocate transition argument — completing the " - "transition on the calling thread"); - return run_transition(target); - } - *arg = target; - - pthread_t tid; - if (pthread_create(&tid, NULL, transition_worker, arg) != 0) + if (!spawn_transition_worker(target)) { - log_error("Failed to create transition thread (%s) — completing the " - "transition on the calling thread", strerror(errno)); - free(arg); + log_error("Completing the transition on the calling thread"); return run_transition(target); } - pthread_detach(tid); return true; } +bool plc_complete_claimed_transition_async(PLCState target) +{ + return spawn_transition_worker(target); +} + // helper: read one line terminated by '\n' from a socket static ssize_t read_line(int fd, char *buffer, size_t max_length) { diff --git a/core/src/plc_app/unix_socket.h b/core/src/plc_app/unix_socket.h index f592c954..823c4f4a 100644 --- a/core/src/plc_app/unix_socket.h +++ b/core/src/plc_app/unix_socket.h @@ -26,4 +26,15 @@ void unix_socket_set_plugin_driver(void *driver); // (same overlap protection: plc_claim_transition refuses while TRANSITIONING). bool plc_begin_transition(PLCState target); +/** + * @brief Run an already claimed transition on a new detached worker thread. + * + * For callers that claimed with plc_claim_transition() and must not run the + * transition themselves, e.g. the dispatcher, which the stop's teardown joins. + * + * @param target PLC_STATE_RUNNING or PLC_STATE_STOPPED + * @return false when the worker could not be spawned; the claim is still held + */ +bool plc_complete_claimed_transition_async(PLCState target); + #endif // UNIX_SOCKET_H diff --git a/core/src/plc_app/utils/log.c b/core/src/plc_app/utils/log.c index 0a235e0a..c18f6175 100644 --- a/core/src/plc_app/utils/log.c +++ b/core/src/plc_app/utils/log.c @@ -2,6 +2,7 @@ // Copyright (c) 2026 Autonomy® #include "log.h" +#include "rt_mutex.h" #include #include #include @@ -20,6 +21,11 @@ static pthread_mutex_t log_mutex = PTHREAD_MUTEX_INITIALIZER; int socket_fd = -1; bool print_logs = false; +__attribute__((constructor)) static void log_mutex_init_pi(void) +{ + rt_mutex_upgrade_static(&log_mutex, "log_mutex"); +} + extern volatile sig_atomic_t keep_running; void log_set_level(LogLevel level) @@ -258,3 +264,52 @@ void log_error(const char *fmt, ...) log_write(LOG_LEVEL_ERROR, fmt, args); va_end(args); } + +void log_emergency(const char *msg) +{ + char line[LOG_MESSAGE_SIZE]; + int n = snprintf(line, sizeof(line), "[FATAL] %s\n", msg); + if (n > 0) + { + size_t len = (size_t)n < sizeof(line) ? (size_t)n : sizeof(line) - 1; + if (write(STDERR_FILENO, line, len) < 0) + { + /* Nothing else to report to. */ + } + } + + if (pthread_mutex_trylock(&log_mutex) != 0) + return; + if (socket_fd >= 0) + { + char escaped[LOG_MESSAGE_SIZE / 2]; + size_t e = 0; + for (const char *p = msg; *p && e + 7 < sizeof(escaped); ++p) + { + unsigned char c = (unsigned char)*p; + if (c < 0x20) + { + e += (size_t)snprintf(escaped + e, sizeof(escaped) - e, "\\u%04x", c); + continue; + } + if (c == '"' || c == '\\') + escaped[e++] = '\\'; + escaped[e++] = (char)c; + } + escaped[e] = '\0'; + + char json[LOG_MESSAGE_SIZE]; + int m = snprintf(json, sizeof(json), + "{\"timestamp\":\"%ld\",\"level\":\"ERROR\",\"message\":\"%s\"}\n", + (long)time(NULL), escaped); + if (m > 0) + { + size_t len = (size_t)m < sizeof(json) ? (size_t)m : sizeof(json) - 1; + if (send(socket_fd, json, len, MSG_DONTWAIT | MSG_NOSIGNAL) < 0) + { + /* Socket full or gone: stderr already has it. */ + } + } + } + pthread_mutex_unlock(&log_mutex); +} diff --git a/core/src/plc_app/utils/log.h b/core/src/plc_app/utils/log.h index 921972ce..f11fcf74 100644 --- a/core/src/plc_app/utils/log.h +++ b/core/src/plc_app/utils/log.h @@ -65,6 +65,17 @@ void log_warn(const char *fmt, ...); */ void log_error(const char *fmt, ...); +/** + * @brief Log an error without ever blocking on a lock. + * + * For fatal paths where the thread holding the log mutex may be dead. Always + * written to stderr; sent to the log socket only if the mutex is free and the + * socket accepts it without blocking. + * + * @param[in] msg The message, no format expansion + */ +void log_emergency(const char *msg); + #ifdef __cplusplus } #endif diff --git a/core/src/plc_app/utils/rt_mutex.h b/core/src/plc_app/utils/rt_mutex.h new file mode 100644 index 00000000..73e14c1e --- /dev/null +++ b/core/src/plc_app/utils/rt_mutex.h @@ -0,0 +1,116 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Autonomy® + +#ifndef RT_MUTEX_H +#define RT_MUTEX_H + +#include +#include +#include +#include + +/* PTHREAD_PRIO_INHERIT is optional in POSIX; MSYS2/Cygwin lack it. */ +#if !defined(__CYGWIN__) && !defined(__MSYS__) && defined(_POSIX_THREAD_PRIO_INHERIT) && \ + _POSIX_THREAD_PRIO_INHERIT > 0 +#define RT_MUTEX_HAS_PI 1 +#else +#define RT_MUTEX_HAS_PI 0 +#endif + +/** + * @brief Initialise a mutex with the priority-inheritance protocol. + * + * A mutex shared between a SCHED_FIFO thread and a SCHED_OTHER thread must use + * priority inheritance, or a preempted low-priority holder can block the + * real-time thread forever on a single CPU. On platforms without PI the mutex + * is initialised with default attributes. + * + * The mutex is left untouched when this fails, so a statically initialised + * mutex (PTHREAD_MUTEX_INITIALIZER) stays usable as a plain mutex. + * + * @param m mutex to initialise; must not be locked or in use + * @return 0 on success, an errno value otherwise + */ +static inline int rt_mutex_init(pthread_mutex_t *m) +{ + pthread_mutexattr_t attr; + int rc = pthread_mutexattr_init(&attr); + if (rc != 0) + return rc; +#if RT_MUTEX_HAS_PI + rc = pthread_mutexattr_setprotocol(&attr, PTHREAD_PRIO_INHERIT); + if (rc != 0) + { + pthread_mutexattr_destroy(&attr); + return rc; + } +#endif + rc = pthread_mutex_init(m, &attr); + pthread_mutexattr_destroy(&attr); + return rc; +} + +/** + * @brief Re-initialise a statically initialised mutex with priority inheritance. + * + * For use from constructors that run before the logger exists: a failure is + * reported on stderr and the mutex keeps its plain static initialisation. + * + * @param m mutex initialised with PTHREAD_MUTEX_INITIALIZER, not yet in use + * @param name label for the stderr report + */ +static inline void rt_mutex_upgrade_static(pthread_mutex_t *m, const char *name) +{ + int rc = rt_mutex_init(m); + if (rc != 0) + fprintf(stderr, "[rt_mutex] %s: PI init failed (%s), using a plain mutex\n", name, + strerror(rc)); +} + +#ifdef __cplusplus +#include +#include +#include + +/** + * @brief Priority-inheritance replacement for std::mutex. + * + * Satisfies BasicLockable, so it works with std::lock_guard and + * std::unique_lock. std::mutex cannot take mutex attributes. + */ +class RtMutex +{ + public: + RtMutex() + { + int rc = rt_mutex_init(&m_); + if (rc != 0) + { + std::fprintf(stderr, "[rt_mutex] PI init failed (%s), using a plain mutex\n", + std::strerror(rc)); + if (pthread_mutex_init(&m_, nullptr) != 0) + std::abort(); + } + } + ~RtMutex() + { + pthread_mutex_destroy(&m_); + } + RtMutex(const RtMutex &) = delete; + RtMutex &operator=(const RtMutex &) = delete; + + void lock() + { + pthread_mutex_lock(&m_); + } + void unlock() + { + pthread_mutex_unlock(&m_); + } + + private: + pthread_mutex_t m_; +}; +#endif + +#endif // RT_MUTEX_H diff --git a/core/src/plc_app/utils/utils.c b/core/src/plc_app/utils/utils.c index 965d2b80..b70345f6 100644 --- a/core/src/plc_app/utils/utils.c +++ b/core/src/plc_app/utils/utils.c @@ -7,6 +7,7 @@ #endif #include "utils.h" +#include "rt_mutex.h" #include #include #include @@ -129,27 +130,16 @@ void lock_memory(void) #endif } +int64_t monotonic_ns(void) +{ + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + return (int64_t)ts.tv_sec * 1000000000LL + ts.tv_nsec; +} + int init_rt_mutex(pthread_mutex_t *mutex) { -#if HAS_REALTIME_FEATURES - pthread_mutexattr_t attr; - if (pthread_mutexattr_init(&attr) != 0) - return -1; - if (pthread_mutexattr_setprotocol(&attr, PTHREAD_PRIO_INHERIT) != 0) - { - pthread_mutexattr_destroy(&attr); - return -1; - } - if (pthread_mutex_init(mutex, &attr) != 0) - { - pthread_mutexattr_destroy(&attr); - return -1; - } - pthread_mutexattr_destroy(&attr); - return 0; -#else - return pthread_mutex_init(mutex, NULL); -#endif + return rt_mutex_init(mutex) == 0 ? 0 : -1; } size_t parse_hex_string(const char *hex_string, uint8_t *data) diff --git a/core/src/plc_app/utils/utils.h b/core/src/plc_app/utils/utils.h index 0559f25e..62c88869 100644 --- a/core/src/plc_app/utils/utils.h +++ b/core/src/plc_app/utils/utils.h @@ -22,7 +22,7 @@ extern "C" { * compute_base_tick_from_config). Default 20 ms before computation. */ extern uint64_t base_tick_ns; -/* Scan counter — incremented once per scan cycle by plc_run_io_cycle_threaded_post. +/* Scan counter — incremented by the dispatcher on every task-bearing tick. * Reported in DEBUG_GET / DEBUG_GET_LIST responses so the editor can * detect cycle boundaries. */ extern unsigned long scan_counter; @@ -31,6 +31,14 @@ extern unsigned long scan_counter; extern char *ext_strucpp_program_md5; +#define NS_PER_MS 1000000LL + +/** + * @brief Current CLOCK_MONOTONIC time in nanoseconds. + * @return nanoseconds since an arbitrary fixed point + */ +int64_t monotonic_ns(void); + /** * @brief Normalize a timespec structure * diff --git a/core/src/plc_app/utils/watchdog.c b/core/src/plc_app/utils/watchdog.c index 4a2c189c..ba0e5668 100644 --- a/core/src/plc_app/utils/watchdog.c +++ b/core/src/plc_app/utils/watchdog.c @@ -1,97 +1,151 @@ // SPDX-License-Identifier: MIT // Copyright (c) 2026 Autonomy® +#ifndef _GNU_SOURCE +#define _GNU_SOURCE +#endif + +#include #include +#include #include #include #include +#include #include #include #include "../plc_state_manager.h" +#include "../task_policy.h" #include "log.h" #include "utils.h" #include "watchdog.h" -atomic_long plc_heartbeat; +/* CLOCK_MONOTONIC ms of the last dispatcher tick; 0 while no dispatcher runs. */ +static atomic_llong g_dispatch_beat_ms; +static atomic_llong g_dispatch_stall_ms; +static const char *_Atomic g_fault_context; + +#define WATCHDOG_TICK_MS 100 + +static int64_t mono_ms(void) +{ + return monotonic_ns() / NS_PER_MS; +} + +void watchdog_feed(void) +{ + atomic_store_explicit(&g_dispatch_beat_ms, mono_ms(), memory_order_relaxed); +} + +void watchdog_dispatcher_started(int64_t period_ns) +{ + int64_t stall = PLC_TASK_STUCK_PERIODS * period_ns / NS_PER_MS; + if (stall < PLC_DISPATCHER_STALL_MIN_MS) + stall = PLC_DISPATCHER_STALL_MIN_MS; + atomic_store_explicit(&g_dispatch_stall_ms, stall, memory_order_relaxed); + watchdog_feed(); +} + +void watchdog_dispatcher_stopped(void) +{ + atomic_store_explicit(&g_dispatch_beat_ms, 0, memory_order_relaxed); +} -/* Watchdog loop period. The stuck-transition bound is NOT defined here: it comes - * from plc_state_manager.h, where it is derived from the same constant the - * transition worker waits on, so this can never fire while the runtime still - * considers the transition to be progressing normally. */ -#define WATCHDOG_TICK_S 2 -#define TRANSITION_STUCK_S (PLC_TRANSITION_STUCK_TIMEOUT_MS / 1000) +void watchdog_set_fault_context(const char *context) +{ + atomic_store(&g_fault_context, context); +} + +void watchdog_fatal_exit(const char *reason) +{ + const char *context = atomic_load(&g_fault_context); + char msg[512]; + snprintf(msg, sizeof(msg), "Watchdog: %s%s%s. Exiting with code %d for a safe-mode restart", + reason, context ? " during " : "", context ? context : "", PLC_EXIT_WATCHDOG_FAULT); + log_emergency(msg); + int fd = open(PLC_WATCHDOG_FAULT_MARKER, O_CREAT | O_WRONLY | O_TRUNC, 0644); + if (fd >= 0) + { + if (write(fd, msg, strnlen(msg, sizeof(msg))) < 0) + { + /* The marker exists even if empty; the boot check only tests for it. */ + } + close(fd); + } + _exit(PLC_EXIT_WATCHDOG_FAULT); +} void *watchdog_thread(void *arg) { (void)arg; - long last = atomic_load(&plc_heartbeat); - int transitioning_ticks = 0; + pthread_setname_np(pthread_self(), "plc_watchdog"); + + struct sched_param sp = {.sched_priority = PLC_FIFO_WATCHDOG}; + int rc = pthread_setschedparam(pthread_self(), SCHED_FIFO, &sp); + if (rc != 0) + log_warn("Watchdog: SCHED_FIFO(%d) failed: %s", PLC_FIFO_WATCHDOG, strerror(rc)); + else + log_info("Watchdog: SCHED_FIFO priority %d", PLC_FIFO_WATCHDOG); + + PLCState last_state = PLC_STATE_STOPPED; + int64_t state_since_ms = 0; + const struct timespec tick = {0, (long)(WATCHDOG_TICK_MS * NS_PER_MS)}; while (1) { - sleep(WATCHDOG_TICK_S); - - PLCState current_state = plc_get_state(); - - // A transition that never publishes a final state would leave the - // runtime in TRANSITIONING forever, refusing every command but PING and - // STATUS — the state is the interlock now, so there is no flag anyone - // could clear to recover. Every path is meant to land a final state; - // this is the backstop for the one that doesn't, turning a silent - // permanent wedge into a reported fault the webserver can act on. - // - // KNOWN LIMITATION: forcing ERROR releases the interlock but does not - // abort the transition, so the worker that failed to land is still - // running -- and a START accepted from ERROR would begin a second one - // over the top of it. Reaching this point at all now takes longer than - // the runtime's own landing bound (see PLC_TRANSITION_STUCK_TIMEOUT_MS), - // which removes the realistic trigger; closing it properly needs the - // transition owner to be able to abort its own work, which is the - // lifecycle-executor refactor and not this function's job. - if (current_state == PLC_STATE_TRANSITIONING_TO_RUN || - current_state == PLC_STATE_TRANSITIONING_TO_STOP) + nanosleep(&tick, NULL); + const PLCState state = plc_get_state(); + const int64_t now = mono_ms(); + + if (state != last_state) + { + last_state = state; + state_since_ms = now; + } + + if (state == PLC_STATE_TRANSITIONING_TO_STOP) { - transitioning_ticks++; - if (transitioning_ticks * WATCHDOG_TICK_S > TRANSITION_STUCK_S) + const int64_t budget = plc_stop_budget_ms(); + if (now - state_since_ms > budget) { - log_error("Watchdog: state change stuck in progress for over %d s — " - "forcing ERROR so the runtime accepts commands again", - TRANSITION_STUCK_S); - plc_force_error_state(); - transitioning_ticks = 0; + char reason[WATCHDOG_REASON_LEN]; + snprintf(reason, sizeof(reason), "stop did not complete within %lld ms", + (long long)budget); + watchdog_fatal_exit(reason); } continue; } - transitioning_ticks = 0; - if (current_state != PLC_STATE_RUNNING) + /* A start that never lands keeps the runtime refusing commands; release it. */ + if (state == PLC_STATE_TRANSITIONING_TO_RUN) { - // Reset tracking when not running so we get a fresh - // baseline when the PLC starts again - if (current_state == PLC_STATE_ERROR) + if (now - state_since_ms > PLC_TRANSITION_STUCK_TIMEOUT_MS) { - last = 0; - atomic_store(&plc_heartbeat, 0); + log_error("Watchdog: start stuck in progress for over %d s — forcing ERROR", + PLC_TRANSITION_STUCK_TIMEOUT_MS / 1000); + plc_force_error_state(); } continue; } - long now = atomic_load(&plc_heartbeat); - if (now == last) + if (state == PLC_STATE_RUNNING) { - log_error("Watchdog: No heartbeat detected - PLC program is unresponsive"); - log_error("The loaded PLC program may contain an infinite loop. " - "Upload a corrected program to recover."); - - // Transition to ERROR state instead of killing the process. - // This keeps the runtime alive so the webserver can still - // communicate with it and upload a new program. - plc_force_error_state(); - continue; - } + const int64_t first_since = plc_first_scan_pending_since_ns(); + if (first_since != 0 && + monotonic_ns() - first_since > (int64_t)PLC_FIRST_SCAN_TIMEOUT_MS * NS_PER_MS) + plc_request_first_scan_trip(); - last = now; + const int64_t beat = atomic_load_explicit(&g_dispatch_beat_ms, memory_order_relaxed); + const int64_t stall = atomic_load_explicit(&g_dispatch_stall_ms, memory_order_relaxed); + if (beat != 0 && now - beat > stall) + { + char reason[WATCHDOG_REASON_LEN]; + snprintf(reason, sizeof(reason), "dispatcher stalled, no tick for %lld ms", + (long long)(now - beat)); + watchdog_fatal_exit(reason); + } + } } return NULL; diff --git a/core/src/plc_app/utils/watchdog.h b/core/src/plc_app/utils/watchdog.h index 6a4a81ce..0d33bf4e 100644 --- a/core/src/plc_app/utils/watchdog.h +++ b/core/src/plc_app/utils/watchdog.h @@ -4,10 +4,57 @@ #ifndef WATCHDOG_H #define WATCHDOG_H -/** - * @brief Initialize the watchdog - * @return int 0 on success, -1 on failure - */ -int watchdog_init(void); +#include + +/* Size of a watchdog_fatal_exit() reason buffer. */ +#define WATCHDOG_REASON_LEN 160 + +#ifdef __cplusplus +extern "C" +{ +#endif + + /** + * @brief Initialize the watchdog + * @return int 0 on success, -1 on failure + */ + int watchdog_init(void); + + /** + * @brief Record a dispatcher tick. Called once per master tick while RUNNING. + */ + void watchdog_feed(void); + + /** + * @brief Arm the dispatcher-stall check for a dispatcher ticking every period_ns. + */ + void watchdog_dispatcher_started(int64_t period_ns); + + /** + * @brief Disarm the dispatcher-stall check once the dispatcher has left its loop. + */ + void watchdog_dispatcher_stopped(void); + + /** + * @brief Name the operation in progress, reported by watchdog_fatal_exit() and + * written into the fault marker. NULL clears it. + * @param context string with static storage duration, or NULL + */ + void watchdog_set_fault_context(const char *context); + + /** + * @brief Last resort when the runtime cannot recover in-process. + * + * Logs reason without blocking on any lock and terminates the process with + * PLC_EXIT_WATCHDOG_FAULT, so the supervisor restarts it in safe mode. + * Outputs are not driven; the hardware safe-state watchdog covers them. + * + * @param reason message for the log; must not be NULL + */ + void watchdog_fatal_exit(const char *reason) __attribute__((noreturn)); + +#ifdef __cplusplus +} +#endif #endif // WATCHDOG_H diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 99c697bf..c9ce0582 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -84,10 +84,11 @@ EMPTY → INIT → RUNNING ⟷ STOPPED → ERROR 1. **Main Thread**: Initialization and signal handling 2. **Unix Socket Thread**: Accepts and processes commands -3. **PLC Cycle Thread**: Executes scan cycles with real-time priority -4. **Stats Thread**: Logs performance metrics -5. **Watchdog Thread**: Monitors heartbeat and terminates on hang -6. **Log Thread**: Manages log socket connection +3. **PLC Cycle Thread / Dispatcher**: Loads the program, then runs the GCD master-tick dispatcher at SCHED_FIFO 98 +4. **Task Threads**: One per IEC TASK, SCHED_FIFO 49 - IEC priority (IEC 0..48, clamped) +5. **Stats Thread**: Logs performance metrics +6. **Watchdog Thread**: SCHED_FIFO 99; bounds stops and detects a stalled dispatcher +7. **Log Thread**: Manages log socket connection ## Real-Time Execution @@ -156,14 +157,40 @@ docker run -v openplc-runtime-data:/var/run/runtime ... ## Watchdog System -The watchdog monitors PLC health by tracking the `plc_heartbeat` atomic variable: - -- **Update Frequency**: Every scan cycle -- **Timeout**: 2 seconds without update -- **Action**: Terminates process if PLC becomes unresponsive -- **State Awareness**: Only monitors during RUNNING state - -**Implementation:** `core/src/plc_app/utils/watchdog.c` +Three layers, from the gentlest to the last resort. The numbers below are the +defaults of the constants in `core/src/plc_app/task_policy.h`. + +1. **Stuck task (dispatcher).** A task found still in the same scan on 10 + consecutive due ticks (10 of its own periods) trips the dispatcher. It claims a + stop, drains every task, and the stop lands in ERROR. The count resets whenever + the task is found idle, so a task whose scans finish within 10 periods never + trips. Ticks replayed after a late dispatcher count as missed periods too. + A task's first scan (initialisation, e.g. Python function block start-up) is + not counted; the main watchdog trips it only if it runs past 10 s. +2. **Drain and abort (dispatcher).** On every stop each in-flight scan may run until + 10 periods after its release (10 s for a first scan). A task still inside IEC + program code then gets `SIGUSR2`, whose handler jumps to the task's recovery + point. A stop that aborted a task lands in ERROR. On every stop, plugin writes + are fenced off (the journal is closed), all `%Q` outputs are written to 0 and + pushed by a final I/O cycle, and held for 500 ms before plugins stop. Runtime + code reached from IEC bodies that must not be interrupted (the Python block + loader) blocks `SIGUSR2`; user C function blocks are not protected, and an abort + landing inside one that holds a C-library lock ends in layer 3. +3. **Process exit (watchdog thread, FIFO 99).** If an aborted task does not exit + within 2 s, a stop exceeds max(10 x the longest task interval, 10 s) + 30.5 s, + or the dispatcher has not ticked for max(10 base ticks, 1 s), the runtime calls + `_exit(42)` after writing `/run/runtime/watchdog_fault`. The webserver restarts + the runtime with `--safe-mode --fault` (the marker forces the same when the exit + code is not visible). That boot reports ERROR, does not load the program, and + first starts the configured plugins just long enough to drive all outputs to 0. + If that step itself hangs, the next exit records it in the marker and the + following boot skips it, relying on the hardware safe-state watchdog. + +All runtime mutexes that can take it use `PTHREAD_PRIO_INHERIT`, and state reads +are lock free, so a starved normal-priority thread cannot block the dispatcher or +the watchdog on a single CPU. + +**Implementation:** `core/src/plc_app/utils/watchdog.c`, `core/src/plc_app/plc_state_manager.cpp` ## Performance Monitoring @@ -188,7 +215,7 @@ Stats are logged every 5 seconds via the stats thread. - Signal handling for SIGINT (graceful shutdown) - State transitions validated before execution - Plugin failures isolated from core runtime -- Watchdog ensures process termination on hang +- Watchdog stops a stuck program, and exits for a safe-mode restart when it cannot ## Directory Structure diff --git a/docs/DEVELOPMENT.md b/docs/DEVELOPMENT.md index 0f3b42b7..5c332e02 100644 --- a/docs/DEVELOPMENT.md +++ b/docs/DEVELOPMENT.md @@ -187,6 +187,7 @@ sudo ./build/plc_main - `--print-logs` - Print logs to stdout in addition to socket - `--print-debug` - Enable debug-level logging - `--safe-mode` - Start in safe mode +- `--fault` - With `--safe-mode`: report ERROR at boot (used after a watchdog exit, code 42) ### Development Mode diff --git a/tests/host/run.sh b/tests/host/run.sh index d1626350..28c9ded4 100755 --- a/tests/host/run.sh +++ b/tests/host/run.sh @@ -29,6 +29,9 @@ trap 'rm -rf "$OUT"' EXIT # test source : extra sources it links TESTS=( "tests/host/test_plc_retain_file_store.cpp:core/src/plc_app/plc_retain_file_store.cpp" + "tests/host/test_rt_mutex.cpp:" + "tests/host/test_task_policy.cpp:core/src/plc_app/task_policy.c" + "tests/host/test_image_outputs.cpp:core/src/plc_app/image_tables.cpp:core/src/plc_app/located_globals.c" ) failures=0 @@ -38,8 +41,20 @@ for entry in "${TESTS[@]}"; do name=$(basename "$test_src" .cpp) printf '\n=== %s ===\n' "$name" + # C sources are compiled as C, not C++. + objs=() + for dep in ${deps//:/ }; do + if [[ "$dep" == *.c ]]; then + obj="$OUT/$(basename "$dep" .c).o" + # shellcheck disable=SC2086 + ${CC:-cc} -std=gnu11 -Wall -Wextra -g $INCLUDES -c "$dep" -o "$obj" || { objs=(); break; } + objs+=("$obj") + else + objs+=("$dep") + fi + done # shellcheck disable=SC2086 - if ! $CXX $CXXFLAGS $INCLUDES "$test_src" ${deps//:/ } -o "$OUT/$name" -lpthread; then + if ! $CXX $CXXFLAGS $INCLUDES "$test_src" ${objs[@]+"${objs[@]}"} -o "$OUT/$name" -lpthread; then echo " FAIL $name did not compile" failures=$((failures + 1)) continue diff --git a/tests/host/test_image_outputs.cpp b/tests/host/test_image_outputs.cpp new file mode 100644 index 00000000..b3ada8cb --- /dev/null +++ b/tests/host/test_image_outputs.cpp @@ -0,0 +1,86 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Autonomy® + +// Host test: image_tables_zero_outputs clears every output slot and nothing else. + +#include "image_tables.h" +#include "journal_buffer.h" + +#include +#include + +extern "C" { +void log_info(const char *, ...) {} +void log_warn(const char *, ...) {} +void log_error(const char *, ...) {} +void log_debug(const char *, ...) {} +void *plugin_manager_get_symbol(PluginManager *, const char *) { return nullptr; } +void *plugin_manager_try_get_symbol(PluginManager *, const char *) { return nullptr; } +void journal_apply_and_clear(void) {} +int journal_write_bool(journal_buffer_type_t, uint16_t, uint8_t, bool) { return 0; } +int journal_write_byte(journal_buffer_type_t, uint16_t, uint8_t) { return 0; } +int journal_write_int(journal_buffer_type_t, uint16_t, uint16_t) { return 0; } +int journal_write_dint(journal_buffer_type_t, uint16_t, uint32_t) { return 0; } +int journal_write_lint(journal_buffer_type_t, uint16_t, uint64_t) { return 0; } +uint64_t base_tick_ns = 0; +char *ext_strucpp_program_md5 = nullptr; +} + +static int g_failures = 0; + +#define CHECK(cond, what) \ + do \ + { \ + if (!(cond)) \ + { \ + std::printf(" FAIL %s\n", what); \ + g_failures++; \ + } \ + } while (0) + +int main() +{ + std::printf("image_tables: outputs forced to 0 on stop\n"); + + image_tables_fill_null_pointers(); + for (int i = 0; i < BUFFER_SIZE; ++i) + { + for (int b = 0; b < 8; ++b) + { + *bool_output[i][b] = 1; + *bool_input[i][b] = 1; + *bool_memory[i][b] = 1; + } + *byte_output[i] = 0xAB; + *int_output[i] = 0xABCD; + *dint_output[i] = 0xABCDEF01u; + *lint_output[i] = 0xABCDEF0123456789ull; + *int_input[i] = 7; + *int_memory[i] = 9; + } + + image_tables_zero_outputs(); + + bool outputs_zero = true, others_kept = true; + for (int i = 0; i < BUFFER_SIZE; ++i) + { + for (int b = 0; b < 8; ++b) + { + outputs_zero &= *bool_output[i][b] == 0; + others_kept &= *bool_input[i][b] == 1 && *bool_memory[i][b] == 1; + } + outputs_zero &= *byte_output[i] == 0 && *int_output[i] == 0 && *dint_output[i] == 0 && + *lint_output[i] == 0; + others_kept &= *int_input[i] == 7 && *int_memory[i] == 9; + } + CHECK(outputs_zero, "every %Q slot is 0"); + CHECK(others_kept, "inputs and memory are untouched"); + + image_tables_clear_null_pointers(); + image_tables_zero_outputs(); + CHECK(true, "zeroing with unbound (NULL) slots does not crash"); + + if (g_failures == 0) + std::printf("all cases passed\n"); + return g_failures == 0 ? 0 : 1; +} diff --git a/tests/host/test_rt_mutex.cpp b/tests/host/test_rt_mutex.cpp new file mode 100644 index 00000000..61d2b0b7 --- /dev/null +++ b/tests/host/test_rt_mutex.cpp @@ -0,0 +1,66 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Autonomy® + +// Host test: rt_mutex_init and RtMutex produce priority-inheritance mutexes. + +#include "utils/rt_mutex.h" + +#include +#include +#include + +static int g_failures = 0; + +#define CHECK(cond, what) \ + do \ + { \ + if (!(cond)) \ + { \ + std::printf(" FAIL %s\n", what); \ + g_failures++; \ + } \ + } while (0) + +#if defined(__GLIBC__) && RT_MUTEX_HAS_PI +/* glibc internal kind bit for PTHREAD_PRIO_INHERIT mutexes. */ +#define GLIBC_MUTEX_PRIO_INHERIT_NP 0x20 +static bool is_pi(const pthread_mutex_t *m) +{ + return (m->__data.__kind & GLIBC_MUTEX_PRIO_INHERIT_NP) != 0; +} +#endif + +int main() +{ + std::printf("rt_mutex: priority-inheritance initialisation\n"); + + pthread_mutex_t m = PTHREAD_MUTEX_INITIALIZER; + CHECK(rt_mutex_init(&m) == 0, "rt_mutex_init returns 0"); +#if defined(__GLIBC__) && RT_MUTEX_HAS_PI + CHECK(is_pi(&m), "rt_mutex_init sets PTHREAD_PRIO_INHERIT"); + pthread_mutex_t plain = PTHREAD_MUTEX_INITIALIZER; + CHECK(!is_pi(&plain), "negative control: static initializer is not PI"); +#endif + CHECK(pthread_mutex_lock(&m) == 0, "PI mutex locks"); + CHECK(pthread_mutex_trylock(&m) != 0, "PI mutex is exclusive"); + CHECK(pthread_mutex_unlock(&m) == 0, "PI mutex unlocks"); + pthread_mutex_destroy(&m); + + RtMutex rm; + int counter = 0; + auto work = [&]() { + for (int i = 0; i < 100000; ++i) + { + std::lock_guard g(rm); + counter++; + } + }; + std::thread a(work), b(work); + a.join(); + b.join(); + CHECK(counter == 200000, "RtMutex serialises std::lock_guard users"); + + if (g_failures == 0) + std::printf("all cases passed\n"); + return g_failures == 0 ? 0 : 1; +} diff --git a/tests/host/test_task_policy.cpp b/tests/host/test_task_policy.cpp new file mode 100644 index 00000000..2e653fb9 --- /dev/null +++ b/tests/host/test_task_policy.cpp @@ -0,0 +1,51 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Autonomy® + +// Host test: IEC TASK priority to SCHED_FIFO mapping. + +#include "task_policy.h" + +#include + +static int g_failures = 0; + +#define CHECK(cond, what) \ + do \ + { \ + if (!(cond)) \ + { \ + std::printf(" FAIL %s\n", what); \ + g_failures++; \ + } \ + } while (0) + +int main() +{ + std::printf("task_policy: IEC priority mapping\n"); + bool clamped = true; + + CHECK(plc_task_fifo_priority(0, &clamped) == 49 && !clamped, "IEC 0 -> FIFO 49"); + CHECK(plc_task_fifo_priority(1, &clamped) == 48 && !clamped, "IEC 1 -> FIFO 48"); + CHECK(plc_task_fifo_priority(48, &clamped) == 1 && !clamped, "IEC 48 -> FIFO 1"); + CHECK(plc_task_fifo_priority(-1, &clamped) == 49 && clamped, "IEC -1 clamps to FIFO 49"); + CHECK(plc_task_fifo_priority(-1000, &clamped) == 49 && clamped, "IEC -1000 clamps to 49"); + CHECK(plc_task_fifo_priority(49, &clamped) == 1 && clamped, "IEC 49 clamps to FIFO 1"); + CHECK(plc_task_fifo_priority(1000, &clamped) == 1 && clamped, "IEC 1000 clamps to FIFO 1"); + CHECK(plc_task_fifo_priority(10, nullptr) == 39, "NULL clamped flag accepted"); + + for (int p = -100; p <= 100; ++p) + { + int f = plc_task_fifo_priority(p, nullptr); + if (f < 1 || f > PLC_FIFO_TASK_MAX || f >= PLC_FIFO_DISPATCHER) + { + std::printf(" FAIL IEC %d -> FIFO %d out of task range\n", p, f); + g_failures++; + } + } + CHECK(PLC_FIFO_WATCHDOG > PLC_FIFO_DISPATCHER, "watchdog above dispatcher"); + CHECK(PLC_FIFO_DISPATCHER > PLC_FIFO_TASK_MAX, "dispatcher above every task"); + + if (g_failures == 0) + std::printf("all cases passed\n"); + return g_failures == 0 ? 0 : 1; +} diff --git a/tests/lifecycle/README.md b/tests/lifecycle/README.md index 5bff45e9..852dc533 100644 --- a/tests/lifecycle/README.md +++ b/tests/lifecycle/README.md @@ -75,8 +75,9 @@ Not covered: transition is honoured, but the branch that hands the movement record back only runs when the corrective transition is *refused* by a request that slipped in first — a race that cannot be forced without a hook in the runtime. -- **The watchdog forcing ERROR on a stuck transition.** Needs a transition that - outlives `PLC_TRANSITION_STUCK_TIMEOUT_MS` (2 minutes), so it belongs in a slow - opt-in run rather than here. -- **A runaway IEC task.** Terminating one needs the forced-abort ladder that does - not exist yet; today a program with an unbounded loop wedges the stop. +- **The watchdog forcing ERROR on a stuck start.** Needs a start that outlives + `PLC_TRANSITION_STUCK_TIMEOUT_MS` (2 minutes), so it belongs in a slow opt-in run + rather than here. +- **A runaway IEC task.** The stuck-task watchdog (see "Watchdog System" in + `docs/ARCHITECTURE.md`) is not exercised here; its tests are the host tests under + `tests/host/` and end-to-end runs against real uploaded programs. diff --git a/tests/pytest/runtimemanager/__init__.py b/tests/pytest/runtimemanager/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/tests/pytest/runtimemanager/test_watchdog_exit.py b/tests/pytest/runtimemanager/test_watchdog_exit.py new file mode 100644 index 00000000..0366528b --- /dev/null +++ b/tests/pytest/runtimemanager/test_watchdog_exit.py @@ -0,0 +1,97 @@ +# SPDX-License-Identifier: MIT +# Copyright (c) 2026 Autonomy® + +"""RuntimeManager restarts plc_main in safe mode after a watchdog fault exit.""" + +# pylint: disable=protected-access,redefined-outer-name + +import subprocess +from pathlib import Path +from unittest.mock import MagicMock + +import pytest + +from webserver import runtimemanager as rm + + +class _ExitedProcess(subprocess.Popen): + """Popen stand-in for a process that already exited with a given code.""" + + def __init__(self, code: int | None) -> None: # pylint: disable=super-init-not-called + self._code = code + self._child_created = False + self.returncode = code + + def poll(self) -> int | None: + return self._code + + +@pytest.fixture +def manager(monkeypatch: pytest.MonkeyPatch) -> rm.RuntimeManager: + mgr = rm.RuntimeManager("/bin/false", "/tmp/x.sock", "/tmp/y.sock") + monkeypatch.setattr(mgr, "_safe_stop_log_server", MagicMock()) + monkeypatch.setattr(mgr, "_safe_close_runtime_socket", MagicMock()) + monkeypatch.setattr(mgr, "_start_runtime_process", MagicMock()) + return mgr + + +def test_exit_code_constant_matches_runtime() -> None: + header_path = Path(__file__).resolve().parents[3] / "core/src/plc_app/task_policy.h" + with open(header_path, encoding="utf-8") as header_file: + header = header_file.read() + assert f"#define PLC_EXIT_WATCHDOG_FAULT {rm.RUNTIME_EXIT_WATCHDOG_FAULT}" in header + + +def test_watchdog_fault_restarts_in_safe_mode_on_first_exit(manager: rm.RuntimeManager) -> None: + manager.process = _ExitedProcess(rm.RUNTIME_EXIT_WATCHDOG_FAULT) + manager._handle_runtime_exit() + manager._start_runtime_process.assert_called_once_with(safe_mode=True, after_fault=True) + assert manager._safe_mode is True + assert manager._crash_times == [] + + +def test_other_exit_restarts_normally(manager: rm.RuntimeManager) -> None: + manager.process = _ExitedProcess(1) + manager._handle_runtime_exit() + manager._start_runtime_process.assert_called_once_with(safe_mode=False) + assert manager._safe_mode is False + + +def test_safe_mode_sticks_after_a_later_crash(manager: rm.RuntimeManager) -> None: + manager.process = _ExitedProcess(rm.RUNTIME_EXIT_WATCHDOG_FAULT) + manager._handle_runtime_exit() + manager.process = _ExitedProcess(-9) + manager._handle_runtime_exit() + manager._start_runtime_process.assert_called_with(safe_mode=True) + + +def test_upload_clears_safe_mode(manager: rm.RuntimeManager) -> None: + manager.process = _ExitedProcess(rm.RUNTIME_EXIT_WATCHDOG_FAULT) + manager._handle_runtime_exit() + manager.reset_crash_tracking() + manager.process = _ExitedProcess(1) + manager._handle_runtime_exit() + manager._start_runtime_process.assert_called_with(safe_mode=False) + + +def test_rapid_crashes_still_enter_safe_mode_without_fault_flag(manager: rm.RuntimeManager) -> None: + for _ in range(rm.MAX_RAPID_CRASHES): + manager.process = _ExitedProcess(-11) + manager._handle_runtime_exit() + manager._start_runtime_process.assert_called_with(safe_mode=True) + assert manager._safe_mode is True + + +def test_start_runtime_process_passes_fault_flag(monkeypatch: pytest.MonkeyPatch) -> None: + mgr = rm.RuntimeManager("/opt/plc_main", "/tmp/x.sock", "/tmp/y.sock") + monkeypatch.setattr(mgr, "_safe_start_log_server", MagicMock()) + monkeypatch.setattr(mgr, "_safe_connect_runtime_socket", MagicMock()) + monkeypatch.setattr(rm.time, "sleep", lambda _s: None) + popen = MagicMock() + monkeypatch.setattr(rm.subprocess, "Popen", popen) + + mgr._start_runtime_process(safe_mode=True, after_fault=True) + assert popen.call_args.args[0] == ["/opt/plc_main", "--safe-mode", "--fault"] + + mgr._start_runtime_process(safe_mode=True) + assert popen.call_args.args[0] == ["/opt/plc_main", "--safe-mode"] diff --git a/webserver/runtimemanager.py b/webserver/runtimemanager.py index 919cc14b..3a1ecbab 100644 --- a/webserver/runtimemanager.py +++ b/webserver/runtimemanager.py @@ -31,6 +31,10 @@ MAX_RAPID_CRASHES = 3 RAPID_CRASH_WINDOW = 30 # seconds +# plc_main exit code for an unrecoverable watchdog fault (PLC_EXIT_WATCHDOG_FAULT in +# core/src/plc_app/task_policy.h). Restart straight into safe mode, reporting ERROR. +RUNTIME_EXIT_WATCHDOG_FAULT = 42 + # How long to let the runtime shut down gracefully after SIGTERM before killing # it. Has to exceed the worst-case graceful stop: the runtime waits for a state # change already in flight to land (a boot start with plugin bring-up is ~4 s on @@ -178,8 +182,8 @@ def is_runtime_alive(self): return True return False - def _start_runtime_process(self, safe_mode=False): - """Start the runtime process, optionally in safe mode.""" + def _start_runtime_process(self, safe_mode: bool = False, after_fault: bool = False) -> None: + """Start the runtime process, optionally in safe mode after a watchdog fault.""" self._safe_start_log_server() try: cmd = [self.runtime_path] @@ -187,6 +191,8 @@ def _start_runtime_process(self, safe_mode=False): cmd.append("--print-debug") if safe_mode: cmd.append("--safe-mode") + if after_fault: + cmd.append("--fault") self.process = subprocess.Popen(cmd) except (OSError, subprocess.SubprocessError) as e: logger.error("Failed to start PLC runtime process: %s", e) @@ -203,6 +209,50 @@ def _record_crash_and_check_safe_mode(self): self._crash_times.append(now) return len(self._crash_times) >= MAX_RAPID_CRASHES + def _runtime_exit_code(self) -> int | None: + """Exit code of the runtime process, when it was started by us and has exited.""" + if isinstance(self.process, subprocess.Popen): + return self.process.poll() + return None + + def _handle_runtime_exit(self) -> None: + """Restart a runtime that exited: safe mode on a watchdog fault or repeated crashes.""" + exit_code = self._runtime_exit_code() + self._safe_stop_log_server() + self._safe_close_runtime_socket() + + if exit_code == RUNTIME_EXIT_WATCHDOG_FAULT: + logger.error( + "PLC runtime exited after an unrecoverable watchdog fault (code %d). " + "Restarting in SAFE MODE - PLC program will NOT be loaded. " + "Upload a corrected program to recover.", + exit_code, + ) + with self._crash_lock: + self._safe_mode = True + self._start_runtime_process(safe_mode=True, after_fault=True) + return + + logger.warning("PLC runtime process died unexpectedly (exit code %s)", exit_code) + if self._record_crash_and_check_safe_mode(): + with self._crash_lock: + if not self._safe_mode: + logger.error( + "PLC program caused %d crashes within %d seconds. " + "Restarting runtime in SAFE MODE - " + "PLC program will NOT be loaded. " + "Upload a corrected program to recover.", + MAX_RAPID_CRASHES, + RAPID_CRASH_WINDOW, + ) + self._safe_mode = True + self._start_runtime_process(safe_mode=True) + else: + with self._crash_lock: + stay_safe = self._safe_mode + logger.warning("Restarting PLC runtime%s...", " in SAFE MODE" if stay_safe else "") + self._start_runtime_process(safe_mode=stay_safe) + def _monitor(self): """ Monitor the PLC runtime process and restart if it dies. @@ -210,26 +260,7 @@ def _monitor(self): """ while self.running: if not self.is_runtime_alive(): - logger.warning("PLC runtime process died unexpectedly") - self._safe_stop_log_server() - self._safe_close_runtime_socket() - - if self._record_crash_and_check_safe_mode(): - with self._crash_lock: - if not self._safe_mode: - logger.error( - "PLC program caused %d crashes within %d seconds. " - "Restarting runtime in SAFE MODE - " - "PLC program will NOT be loaded. " - "Upload a corrected program to recover.", - MAX_RAPID_CRASHES, - RAPID_CRASH_WINDOW, - ) - self._safe_mode = True - self._start_runtime_process(safe_mode=True) - else: - logger.warning("Restarting PLC runtime...") - self._start_runtime_process(safe_mode=False) + self._handle_runtime_exit() else: # Make sure log server and socket are connected if not self.log_server.running: @@ -392,10 +423,10 @@ def send_plugin_command(self, plugin_name: str, command_json: str, timeout: floa # Parse: "PLUGIN_CMD:OK:{json}" or "PLUGIN_CMD:ERROR:{json}" if response.startswith("PLUGIN_CMD:OK:"): - json_str = response[len("PLUGIN_CMD:OK:"):] + json_str = response[len("PLUGIN_CMD:OK:") :] return json.loads(json_str) elif response.startswith("PLUGIN_CMD:ERROR:"): - json_str = response[len("PLUGIN_CMD:ERROR:"):] + json_str = response[len("PLUGIN_CMD:ERROR:") :] return json.loads(json_str) else: return {"error": f"Unexpected response: {response[:200]}"}