[Git][ghc/ghc][wip/dcoutts/issue-27105-stopTicker] 19 commits: Promote HAVE_PREEMPTION from Timer.c to OSThreads.h
Duncan Coutts pushed to branch wip/dcoutts/issue-27105-stopTicker at Glasgow Haskell Compiler / GHC Commits: 84cee6bb by Duncan Coutts at 2026-06-08T13:07:29+02:00 Promote HAVE_PREEMPTION from Timer.c to OSThreads.h We will want to know about HAVE_PREEMPTION in more places. HAVE_PREEMPTION tells us that we do have OS threads available, irrespective of whether THREADED is defined. In particular, HAVE_PREEMPTION is defined on all proper OSs, but not on WASM (and hyopthetically may not be true on some other platforms like micro-controllers, RTOSs, VM hypervisors etc). - - - - - 1e270fa6 by Duncan Coutts at 2026-06-08T13:07:31+02:00 Define ACQUIRE_LOCK_ALWAYS and friends Fix issue #27335 Like the atomic _ALWAYS variants, these lock actions are always defined, rather than being dependent on whether we are in the THREADED case. All the "normal" LOCK macros are defined to be no-ops when !THREADED. The use case for the _ALWAYS variants is where we are using OS threads even in the non-threaded RTS. This includes everything to do with the timer/ticker thread, which is used in the non-threaded RTS too. In particular, we will want to use this for eventlog things, because the timer thread performs eventlogging concurrently with the main capability, even in the non-threaded RTS. - - - - - 10806a8c by Duncan Coutts at 2026-06-08T13:07:31+02:00 Use ACQUIRE/RELEASE_LOCK_ALWAYS with eventBufMutex Even in the non-threaded RTS the eventBufMutex is needed by both the main capability and the timer/ticker thread, so always use the mutex. This should fix #25165 which is about the main capability and the timer thread posting events to the eventlog buffer concurrently and thereby corrupting the buffer data. - - - - - a5e2d9c6 by Duncan Coutts at 2026-06-08T13:07:31+02:00 Expose eventBufMutex in the EventLog interface/header We will need it in forkProcess to ensure we don't write to the global eventlog buffer concurrently with trying to flush eventlog buffers and do the fork(). - - - - - 7f881fdd by Duncan Coutts at 2026-06-08T13:07:31+02:00 Split flushAllCapsEventsBufs into safe and unlocked version Following the convention that unlocked versions have a trailing _ underscore in their name. This one requires the caller to hold the eventlog global buffer mutex. We will need this in forkProcess. - - - - - 45c43c44 by Duncan Coutts at 2026-06-08T13:07:31+02:00 Remove redundant use of stopTimer in setNumCapabilities Historically, the comment here was: We must stop the interval timer while we are changing the capabilities array lest handle_tick may try to context switch an old capability. See #17289. and We must disable the timer while we do this since the tick handler may call contextSwitchAllCapabilities, which may see the capabilities array as we free it. What this refers to is that historically, when changing the number of capabilities, the array of capabilities was reallocated to a new size, allocating new ones and freeing the old ones, thus invalidating all existing capbility pointers. Strangely, for good measure the code used to call stopTimer twice (hence the two similar comments above). However, since commit a3eccf06292dd666b24606251a52da2b466a9612, the capabilities array is no longer reallocated. Instead the array is allcoated once on RTS startup to the maximum size it could ever be allowed to be, and then capabilities get enabled/disabled at runtime. So the capability pointers never become invalid anymore. At worst, they may point to capabilities that are disabled. Thus we no longer need to stop the timer (twice) while we change the number of enabled capabilities. This also partially solves issue #27105, which notes that stopTimer is being used as if it were synchronous, when it is not. At least for this case, the solution is that stopTimer is not needed at all! - - - - - 55c052e2 by Duncan Coutts at 2026-06-08T13:07:31+02:00 Remove redundant use of stopTimer in forkProcess but replace it with taking the eventlog buffer lock during the fork. Fixes issue #27105 The original reason to block the timer during a fork was that historically the timer was implemented using a periodic timer signal, and the signal itself would interrupt the fork system call (returning EINTR). For large processes (where fork() takes a while) this could permanently livelock: the timer always would go off before the fork could complete, which got retried in a loop forever. The timer is no longer implemented as a unix signal, but uses threads. Thus the original problem no longer exists. The only remaining reason to block the timer tick is to prevent actions taken by the tick from interfering with the delicate process involved in fork (taking a load of locks and pausing everything). The only thing we need to do is to prevent the eventlog from being written to or flushed while the fork is taking place. To achieve this all we need to do is hold the mutex for the global eventlog buffer. This removes the last use of stopTimer that expects stopTimer to work synchronously (which it was not) and thus solves issue #27105. To be clear, we solve issue #27105 not by making stopTimer synchronous, but by eliminating the use sites that expected it to be synchronous. - - - - - 3d2a4e6e by Duncan Coutts at 2026-06-10T09:59:48+01:00 Add a test for thread scheduler fairness It also tests that the interval timer and context switching works. We also test that fairness is lost when the context switching interval is too coarse for the duration of the test. We add this test before doing surgery on the interval timer, so we have decent coverage. - - - - - 5b29323b by Duncan Coutts at 2026-06-10T10:59:05+01:00 Make exported stop/startTimer no-ops, and rename internal functions Specifically, internally rename: stop/startTimer to pause/unpauseTimer stop/startTicker to pause/unpauseTicker and keep stop/startTimer as exported functions, but now as no-ops. In the past the stop/startTicker actions were used incorrectly as if they were synchronous, which they are not. See issue #27105. We now document pause/unpackTicker as being async and not to be used for the purpose of concurrency safety. The existing stop/startTimer (note Timer not Ticker, the Timer calls the Ticker!) are also exported from the RTS as a public API. This was historically because the ticker used signals and it was important to suspend the timer signel over a process fork. So these functions were exported to be used by the process and unix libraries. We cannot just remove the RTS exports, but we now make them no-ops, and they can be removed from the process and unix library later. This was already documented in a changelog.d entry no-more-timer-signal but due to changes during the MR process the change to make stop/startTicker into no-ops didn't make it into the earlier MR. - - - - - a847e252 by Duncan Coutts at 2026-06-10T11:02:43+01:00 Make exitTicker/exitTimer unconditionally synchronous We never use them asynchronously, and we should never need to do so. And update some related comments. - - - - - e3705818 by Duncan Coutts at 2026-06-11T00:22:04+01:00 posix ticker: update and improve comments on (un)pause and exit Clarify what is async vs sync. - - - - - 2de52c61 by Duncan Coutts at 2026-06-11T00:22:34+01:00 posix ticker: split out ppoll/select helper functions Move the #ifdefs out of the main code body by introducing local helper functions and types, which themselves have two implementations (with a common API) based on ppoll or select. This helps improve clarity/readability. - - - - - 499de471 by Duncan Coutts at 2026-06-11T00:22:52+01:00 posix ticker: improve the implementation The existing implementation supported pausing and exiting, with the implementation of pausing reling on a mutex and condition variable. It needed to check the pause and stop shared variables on every iteration. It relies on ppoll or select, to wait on the timeout and also wait on an interrupt fd. The interrupt fd was only used for prompt exit/shutdown, and not for pausing or other notification. The pause only needed a lock and a memory operation, but the pause was not prompt. The resume used a lock, and signaling a cond var. The new implementation uses a somewhat more regular design: every notification is done by setting a shared variable and interrupting/notifying the ticker via the fd. The ticker thread does not need to check any shared variables on normal timer expiry, only when it recevies notification. This may be a micro-optimisation, but the tick occurs 100 times a second by default so any improvements in the hot path are amplified. When the ticker thread does receive notification it can check the various shared variables and update its local state. The blocking relies on using ppoll/select but without a timeout. This avoids the condition var and also allows further notifications when paused (also used for unpausing). This design can be extended with further notification types if needed by using and checking further shared vars (or making existing shared vars an enum or counter). This may be used in future for additional notifications to the ticker thread. This will likely be used to proxy wakeUpRts from a single handler context for example. And this approach, avoiding mutexes, is compatible with use from signal handlers. So overall, it's: * slightly simpler / more regular; * easier to extend with additional notifications; * probably slightly more efficient (but a micro-optimisation); * and supports calling notification from signal handlers - - - - - 993b9399 by Duncan Coutts at 2026-06-11T00:23:46+01:00 posix ticker: further minor local renaming for code clarity Improve the clarity with better choice of names for several local vars and function. - - - - - df53d0e3 by Duncan Coutts at 2026-06-11T10:17:21+01:00 win32 ticker: split out local helper functions - - - - - 0f01836b by Duncan Coutts at 2026-06-11T10:19:40+01:00 win32 ticker: provide guarantee about concurrency and idempotency Use a lock to ensure pause/unpause can be used concurrently. Use a paused variable, protected by the lock, to ensure that pause and unpause are both idempotent. This is what the portable API expects. - - - - - e0c747e9 by Duncan Coutts at 2026-06-11T10:22:52+01:00 win32 ticker: make the initial tick be after one wait interval There is no need to tick immediately. This is consistent with the posix implementation. - - - - - c6a0ceb6 by Duncan Coutts at 2026-06-11T10:22:52+01:00 ticker: remove now-unnecessary layer of enable/disable There was an atomic variable used to block *part* of the actions of the tick handler. This still did not make stopTimer synchronous, even for the part of the the handle_tick actions it covered. It also added a more expensive (sequentuially consistent) atomic operation in the hot path for the handle_tick action, whereas our new design requires no atomic ops at all. Now that we have eliminate the need for synchronous stop/startTicker, we don't need this not-quite-working-anyway atomic protocol. The new pause/unpauseTicker is explicitly asynchronous and idempotent. - - - - - 50dd42da by Duncan Coutts at 2026-06-11T10:22:52+01:00 ticker: add TODOs about issue #27250: too much being done from handle_tick The handle_tick should not perform I/O, block, perform long-running operations or call arbitrary user code. Unfortunately, everything to do with the eventlog (at the moment) falls into all those categories. - - - - - 14 changed files: - rts/Capability.c - rts/RtsStartup.c - rts/Schedule.c - rts/Ticker.h - rts/Timer.c - rts/Timer.h - rts/eventlog/EventLog.c - rts/eventlog/EventLog.h - rts/include/rts/OSThreads.h - rts/include/rts/Timer.h - rts/posix/Ticker.c - rts/win32/Ticker.c - + testsuite/tests/concurrent/should_run/T27105.hs - testsuite/tests/concurrent/should_run/all.T Changes: ===================================== rts/Capability.c ===================================== @@ -443,13 +443,6 @@ void moreCapabilities (uint32_t from USED_IF_THREADS, uint32_t to USED_IF_THREADS) { #if defined(THREADED_RTS) - // We must disable the timer while we do this since the tick handler may - // call contextSwitchAllCapabilities, which may see the capabilities array - // as we free it. The alternative would be to protect the capabilities - // array with a lock but this seems more expensive than necessary. - // See #17289. - stopTimer(); - if (to == 1) { // THREADED_RTS must work on builds that don't have a mutable // BaseReg (eg. unregisterised), so in this case @@ -470,8 +463,6 @@ moreCapabilities (uint32_t from USED_IF_THREADS, uint32_t to USED_IF_THREADS) } debugTrace(DEBUG_sched, "allocated %d more capabilities", to - from); - - startTimer(); #endif } ===================================== rts/RtsStartup.c ===================================== @@ -415,8 +415,8 @@ hs_init_ghc(int *argc, char **argv[], RtsConfig rts_config) traceInitEvent(dumpIPEToEventLog); initHeapProfiling(); - /* start the virtual timer 'subsystem'. */ - startTimer(); + /* start the timer (after initTimer above) */ + unpauseTimer(); #if defined(RTS_USER_SIGNALS) if (RtsFlags.MiscFlags.install_signal_handlers) { @@ -512,14 +512,12 @@ hs_exit_(bool wait_foreign) } #endif - /* stop the ticker */ - stopTimer(); - /* - * it is quite important that we wait here as some timer implementations - * (e.g. pthread) may fire even after we exit, which may segfault as we've - * already freed the capabilities. + /* We rely on the guarantee that exitTimer stops the timer synchronously, + * which ensures the timer handler does not get run again after this point. + * We are about to start freeing resources used by the timer handler (like + * the capabilities, eventlog and profiling data structures). */ - exitTimer(true); + exitTimer(); /* * Dump the ticky counter definitions ===================================== rts/Schedule.c ===================================== @@ -37,6 +37,7 @@ #include "win32/AsyncWinIO.h" #endif #include "Trace.h" +#include "eventlog/EventLog.h" #include "RaiseAsync.h" #include "Threads.h" #include "Timer.h" @@ -454,7 +455,7 @@ run_thread: prev = setRecentActivity(ACTIVITY_YES); if (prev == ACTIVITY_DONE_GC) { #if !defined(PROFILING) - startTimer(); + unpauseTimer(); #endif } break; @@ -1935,7 +1936,7 @@ delete_threads_and_gc: // it will get re-enabled if we run any threads after the GC. setRecentActivity(ACTIVITY_DONE_GC); #if !defined(PROFILING) - stopTimer(); + pauseTimer(); #endif break; } @@ -2100,24 +2101,31 @@ forkProcess(HsStablePtr *entry ACQUIRE_LOCK(&all_tasks_mutex); #endif - stopTimer(); // See #4074 - #if defined(TRACING) - flushAllCapsEventsBufs(); // so that child won't inherit dirty file buffers +#if defined(HAVE_PREEMPTION) + // We must hold the eventlog global mutex over the fork to prevent the + // timer thread from trying to post events. While holding the mutex we need + // to flush the eventlogs (global and per-cap) so that child won't inherit + // dirty eventlog buffers or file buffers. + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); +#endif + flushAllCapsEventsBufs_(); #endif pid = fork(); if (pid) { // parent - startTimer(); // #4074 - RELEASE_LOCK(&sched_mutex); RELEASE_LOCK(&sm_mutex); RELEASE_LOCK(&stable_ptr_mutex); RELEASE_LOCK(&stable_name_mutex); RELEASE_LOCK(&task->lock); +#if defined(TRACING) && defined(HAVE_PREEMPTION) + RELEASE_LOCK_ALWAYS(&eventBufMutex); +#endif + #if defined(THREADED_RTS) /* N.B. releaseCapability_ below may need to take all_tasks_mutex */ RELEASE_LOCK(&all_tasks_mutex); @@ -2224,8 +2232,8 @@ forkProcess(HsStablePtr *entry generations[g].threads = END_TSO_QUEUE; } - // On Unix, all timers are reset in the child, so we need to start - // the timer again. + // The timer thread is not present in the child process, so we need + // to initialise the timer again. initTimer(); // TODO: need to trace various other things in the child @@ -2236,7 +2244,7 @@ forkProcess(HsStablePtr *entry // start timer after the IOManager is initialized // (the idle GC may wake up the IOManager) - startTimer(); + unpauseTimer(); // Install toplevel exception handlers, so interruption // signal will be sent to the main thread. @@ -2303,12 +2311,6 @@ setNumCapabilities (uint32_t new_n_capabilities USED_IF_THREADS) cap = rts_lock(); task = cap->running_task; - - // N.B. We must stop the interval timer while we are changing the - // capabilities array lest handle_tick may try to context switch - // an old capability. See #17289. - stopTimer(); - stopAllCapabilities(&cap, task); if (new_n_capabilities < enabled_capabilities) @@ -2364,9 +2366,7 @@ setNumCapabilities (uint32_t new_n_capabilities USED_IF_THREADS) tracingAddCapabilities(n_capabilities, new_n_capabilities); #endif - // Resize the capabilities array - // NB. after this, capabilities points somewhere new. Any pointers - // of type (Capability *) are now invalid. + // Allocate and initialise the extra capabilities moreCapabilities(n_capabilities, new_n_capabilities); // Resize and update storage manager data structures @@ -2394,8 +2394,6 @@ setNumCapabilities (uint32_t new_n_capabilities USED_IF_THREADS) // Notify IO manager that the number of capabilities has changed. notifyIOManagerCapabilitiesChanged(&cap); - startTimer(); - rts_unlock(cap); #endif // THREADED_RTS ===================================== rts/Ticker.h ===================================== @@ -12,9 +12,44 @@ typedef void (*TickProc)(int); -void initTicker (Time interval, TickProc handle_tick); -void startTicker (void); -void stopTicker (void); -void exitTicker (bool wait); +/* The ticker is initialised in a paused state. Use unpauseTicker to start. */ +void initTicker(Time interval, TickProc handle_tick); + +/* Stop and terminate the ticker. It does not need to be stopped first. + * The exitTicker action is *synchronous*. When it returns the caller is + * guaranteed that the tick action is blocked. + */ +void exitTicker(void); + +/* Pause and unpause (resume) the ticker. + * + * The pauseTicker and unpauseTicker actions are *asynchronous*. After calling + * pauseTicker, the ticker will pause eventually, but there may be another tick + * action before it does pause (and theoretically there could be several but + * in practice this is unlikely). Similarly, after calling unpauseTicker the + * ticker will start up again eventually, but there is an unspecified delay + * between the unpause and the next tick action (but in practice it is short). + * + * This should be used for the purpose of *efficiency*: to avoid unnecessary + * OS thread wakeups caused by the ticker. + * + * These should *not* be used for the purpose of *concurrency safety*: to + * prevent the tick action from running concurrently with some other critical + * section. The synchronous case is not provided because it is not currently + * needed (and proper locking is often a better solution anyway). + * + * The pairing of unpauseTicker and the handle_tick action form a + * synchonises-with relation: values written before unpauseTicker can be + * read from the resulting handle_tick action. + * + * It *is* safe to call these functions from within the tick handler itself. + * + * It is safe to use these functions concurrently from multiple threads, but + * note that they *are* idempotent. This means it is not appropriate to use + * paired pause/unpause calls concurrently. They can be used by threads based + * on consistent use of some shared state or observation. + */ +void pauseTicker(void); +void unpauseTicker(void); #include "EndPrivate.h" ===================================== rts/Timer.c ===================================== @@ -28,20 +28,6 @@ #include "RtsSignals.h" #include "rts/EventLogWriter.h" -// See Note [No timer on wasm32] -#if !defined(wasm32_HOST_ARCH) -#define HAVE_PREEMPTION -#endif - -// This global counter is used to allow multiple threads to stop the -// timer temporarily with a stopTimer()/startTimer() pair. If -// timer_enabled == 0 timer is enabled -// timer_disabled == N, N > 0 timer is disabled by N threads -// When timer_enabled makes a transition to 0, we enable the timer, -// and when it makes a transition to non-0 we disable it. - -static StgWord timer_disabled; - /* ticks left before next pre-emptive context switch */ static int ticks_to_ctxt_switch = 0; @@ -112,9 +98,9 @@ static void handle_tick(int unused STG_UNUSED) { - handleProfTick(); - if (RtsFlags.ConcFlags.ctxtSwitchTicks > 0 - && SEQ_CST_LOAD_ALWAYS(&timer_disabled) == 0) + handleProfTick(); // Bad or worse: see issue #27250. + + if (RtsFlags.ConcFlags.ctxtSwitchTicks > 0) { ticks_to_ctxt_switch--; if (ticks_to_ctxt_switch <= 0) { @@ -128,7 +114,7 @@ handle_tick(int unused STG_UNUSED) ticks_to_eventlog_flush--; if (ticks_to_eventlog_flush <= 0) { ticks_to_eventlog_flush = RtsFlags.TraceFlags.eventlogFlushTicks; - flushEventLog(NULL); + flushEventLog(NULL); // Bad or worse: see issue #27250. } } #endif @@ -153,7 +139,7 @@ handle_tick(int unused STG_UNUSED) RtsFlags.MiscFlags.tickInterval; #if defined(THREADED_RTS) wakeUpRts(); - // The scheduler will call stopTimer() when it has done + // The scheduler will call pauseTimer() when it has done // the GC. #endif } else { @@ -165,10 +151,10 @@ handle_tick(int unused STG_UNUSED) #if defined(PROFILING) if (!(RtsFlags.ProfFlags.doHeapProfile || RtsFlags.CcFlags.doCostCentres)) { - stopTimer(); + pauseTimer(); } #else - stopTimer(); + pauseTimer(); #endif } } else { @@ -181,48 +167,49 @@ handle_tick(int unused STG_UNUSED) } } -void -initTimer(void) +void initTimer(void) { #if defined(HAVE_PREEMPTION) initProfTimer(); if (RtsFlags.MiscFlags.tickInterval != 0) { initTicker(RtsFlags.MiscFlags.tickInterval, handle_tick); } - SEQ_CST_STORE_ALWAYS(&timer_disabled, 1); #endif } -void -startTimer(void) +/* Deprecated exported functions. Now no-ops. + * Historically they were used by the process and unix libraries to disable + * the signal-based interval timer, since otherwise the timer signal would + * keep going off in the child process and confusing everything. The interval + * timer no longer uses signals, so there is no need any more for libraries to + * disable the timer. Also, the timer internal API has changed. + */ +void stopTimer(void) { /* no-op */ } +void startTimer(void) { /* no-op */ } + +void pauseTimer(void) { #if defined(HAVE_PREEMPTION) - if (SEQ_CST_SUB_ALWAYS(&timer_disabled, 1) == 0) { - if (RtsFlags.MiscFlags.tickInterval != 0) { - startTicker(); - } + if (RtsFlags.MiscFlags.tickInterval != 0) { + pauseTicker(); } #endif } -void -stopTimer(void) +void unpauseTimer(void) { #if defined(HAVE_PREEMPTION) - if (SEQ_CST_ADD_ALWAYS(&timer_disabled, 1) == 1) { - if (RtsFlags.MiscFlags.tickInterval != 0) { - stopTicker(); - } + if (RtsFlags.MiscFlags.tickInterval != 0) { + unpauseTicker(); } #endif } -void -exitTimer (bool wait) +void exitTimer (void) { #if defined(HAVE_PREEMPTION) if (RtsFlags.MiscFlags.tickInterval != 0) { - exitTicker(wait); + exitTicker(); } #endif } ===================================== rts/Timer.h ===================================== @@ -8,5 +8,12 @@ #pragma once -RTS_PRIVATE void initTimer (void); -RTS_PRIVATE void exitTimer (bool wait); +#include "BeginPrivate.h" + +void initTimer(void); +void exitTimer(void); + +void pauseTimer(void); +void unpauseTimer(void); + +#include "EndPrivate.h" ===================================== rts/eventlog/EventLog.c ===================================== @@ -129,8 +129,11 @@ typedef struct _EventsBuf { static EventsBuf *capEventBuf; // one EventsBuf for each Capability static EventsBuf eventBuf; // an EventsBuf not associated with any Capability -#if defined(THREADED_RTS) -static Mutex eventBufMutex; // protected by this mutex +#if defined(HAVE_PREEMPTION) +// Note that this mutex is used even in the non-threaded RTS, since the timer +// thread posts events and flushes. So _all_ uses of this mutex must use +// ACQUIRE_LOCK_ALWAYS/RELEASE_LOCK_ALWAYS. +Mutex eventBufMutex; // protects eventBuf above #endif // Event type @@ -393,8 +396,10 @@ initEventLogging(void) moreCapEventBufs(0, get_n_capabilities()); initEventsBuf(&eventBuf, EVENT_LOG_SIZE, (EventCapNo)(-1)); -#if defined(THREADED_RTS) +#if defined(HAVE_PREEMPTION) initMutex(&eventBufMutex); +#endif +#if defined(THREADED_RTS) initMutex(&state_change_mutex); #endif } @@ -416,7 +421,7 @@ startEventLogging_(void) { initEventLogWriter(); - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); postHeaderEvents(); /* @@ -425,7 +430,7 @@ startEventLogging_(void) */ printAndClearEventBuf(&eventBuf); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); return true; } @@ -495,7 +500,7 @@ endEventLogging(void) flushEventLog_(NULL); - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); // Mark end of events (data). postEventTypeNum(&eventBuf, EVENT_DATA_END); @@ -503,7 +508,7 @@ endEventLogging(void) // Flush the end of data marker. printAndClearEventBuf(&eventBuf); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); stopEventLogWriter(); event_log_writer = NULL; @@ -666,7 +671,7 @@ void postCapEvent (EventTypeNum tag, EventCapNo capno) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, tag); postEventHeader(&eventBuf, tag); @@ -685,14 +690,14 @@ postCapEvent (EventTypeNum tag, barf("postCapEvent: unknown event tag %d", tag); } - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postCapsetEvent (EventTypeNum tag, EventCapsetID capset, StgWord info) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, tag); postEventHeader(&eventBuf, tag); @@ -726,7 +731,7 @@ void postCapsetEvent (EventTypeNum tag, barf("postCapsetEvent: unknown event tag %d", tag); } - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postCapsetStrEvent (EventTypeNum tag, @@ -740,14 +745,14 @@ void postCapsetStrEvent (EventTypeNum tag, return; } - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); if (!hasRoomForVariableEvent(&eventBuf, size)){ printAndClearEventBuf(&eventBuf); if (!hasRoomForVariableEvent(&eventBuf, size)){ errorBelch("Event size exceeds buffer size, bail out"); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); return; } } @@ -758,7 +763,7 @@ void postCapsetStrEvent (EventTypeNum tag, postBuf(&eventBuf, (StgWord8*) msg, strsize); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postCapsetVecEvent (EventTypeNum tag, @@ -783,14 +788,14 @@ void postCapsetVecEvent (EventTypeNum tag, } } - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); if (!hasRoomForVariableEvent(&eventBuf, size)){ printAndClearEventBuf(&eventBuf); if(!hasRoomForVariableEvent(&eventBuf, size)){ errorBelch("Event size exceeds buffer size, bail out"); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); return; } } @@ -804,7 +809,7 @@ void postCapsetVecEvent (EventTypeNum tag, postBuf(&eventBuf, (StgWord8*) argv[i], 1 + strlen(argv[i])); } - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postWallClockTime (EventCapsetID capset) @@ -813,7 +818,7 @@ void postWallClockTime (EventCapsetID capset) StgWord64 sec; StgWord32 nsec; - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); /* The EVENT_WALL_CLOCK_TIME event is intended to allow programs reading the eventlog to match up the event timestamps with wall @@ -846,7 +851,7 @@ void postWallClockTime (EventCapsetID capset) postWord64(&eventBuf, sec); postWord32(&eventBuf, nsec); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } /* @@ -885,7 +890,7 @@ void postEventHeapInfo (EventCapsetID heap_capset, W_ mblockSize, W_ blockSize) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_HEAP_INFO_GHC); postEventHeader(&eventBuf, EVENT_HEAP_INFO_GHC); @@ -899,7 +904,7 @@ void postEventHeapInfo (EventCapsetID heap_capset, postWord64(&eventBuf, mblockSize); postWord64(&eventBuf, blockSize); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postEventGcStats (Capability *cap, @@ -952,7 +957,7 @@ void postTaskCreateEvent (EventTaskId taskId, EventCapNo capno, EventKernelThreadId tid) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_TASK_CREATE); postEventHeader(&eventBuf, EVENT_TASK_CREATE); @@ -961,14 +966,14 @@ void postTaskCreateEvent (EventTaskId taskId, postCapNo(&eventBuf, capno); postKernelThreadId(&eventBuf, tid); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postTaskMigrateEvent (EventTaskId taskId, EventCapNo capno, EventCapNo new_capno) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_TASK_MIGRATE); postEventHeader(&eventBuf, EVENT_TASK_MIGRATE); @@ -977,28 +982,28 @@ void postTaskMigrateEvent (EventTaskId taskId, postCapNo(&eventBuf, capno); postCapNo(&eventBuf, new_capno); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postTaskDeleteEvent (EventTaskId taskId) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_TASK_DELETE); postEventHeader(&eventBuf, EVENT_TASK_DELETE); /* EVENT_TASK_DELETE (taskID) */ postTaskId(&eventBuf, taskId); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postEventNoCap (EventTypeNum tag) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, tag); postEventHeader(&eventBuf, tag); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void @@ -1042,9 +1047,9 @@ void postLogMsg(EventsBuf *eb, EventTypeNum type, char *msg, va_list ap) void postMsg(char *msg, va_list ap) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); postLogMsg(&eventBuf, EVENT_LOG_MSG, msg, ap); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postCapMsg(Capability *cap, char *msg, va_list ap) @@ -1138,32 +1143,32 @@ void postConcUpdRemSetFlush(Capability *cap) void postConcMarkEnd(StgWord32 marked_obj_count) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_CONC_MARK_END); postEventHeader(&eventBuf, EVENT_CONC_MARK_END); postWord32(&eventBuf, marked_obj_count); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postNonmovingHeapCensus(uint16_t blk_size, const struct NonmovingAllocCensus *census) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); postEventHeader(&eventBuf, EVENT_NONMOVING_HEAP_CENSUS); postWord16(&eventBuf, blk_size); postWord32(&eventBuf, census->n_active_segs); postWord32(&eventBuf, census->n_filled_segs); postWord32(&eventBuf, census->n_live_blocks); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postNonmovingPrunedSegments(uint32_t pruned_segments, uint32_t free_segments) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); postEventHeader(&eventBuf, EVENT_NONMOVING_PRUNED_SEGMENTS); postWord32(&eventBuf, pruned_segments); postWord32(&eventBuf, free_segments); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void closeBlockMarker (EventsBuf *ebuf) @@ -1224,7 +1229,7 @@ static HeapProfBreakdown getHeapProfBreakdown(void) void postHeapProfBegin(void) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); PROFILING_FLAGS *flags = &RtsFlags.ProfFlags; StgWord modSelector_len = flags->modSelector ? strlen(flags->modSelector) : 0; @@ -1258,42 +1263,42 @@ void postHeapProfBegin(void) postStringLen(&eventBuf, flags->ccsSelector, ccsSelector_len); postStringLen(&eventBuf, flags->retainerSelector, retainerSelector_len); postStringLen(&eventBuf, flags->bioSelector, bioSelector_len); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postHeapProfSampleBegin(StgInt era) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_HEAP_PROF_SAMPLE_BEGIN); postEventHeader(&eventBuf, EVENT_HEAP_PROF_SAMPLE_BEGIN); postWord64(&eventBuf, era); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postHeapBioProfSampleBegin(StgInt era, StgWord64 time) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_HEAP_BIO_PROF_SAMPLE_BEGIN); postEventHeader(&eventBuf, EVENT_HEAP_BIO_PROF_SAMPLE_BEGIN); postWord64(&eventBuf, era); postWord64(&eventBuf, time); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postHeapProfSampleEnd(StgInt era) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_HEAP_PROF_SAMPLE_END); postEventHeader(&eventBuf, EVENT_HEAP_PROF_SAMPLE_END); postWord64(&eventBuf, era); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postHeapProfSampleString(const char *label, StgWord64 residency) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); StgWord label_len = strlen(label); StgWord len = 1+8+label_len+1; CHECK(!ensureRoomForVariableEvent(&eventBuf, len)); @@ -1303,7 +1308,7 @@ void postHeapProfSampleString(const char *label, postWord8(&eventBuf, 0); postWord64(&eventBuf, residency); postStringLen(&eventBuf, label, label_len); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } #if defined(PROFILING) @@ -1313,7 +1318,7 @@ void postHeapProfCostCentre(StgWord32 ccID, const char *srcloc, StgBool is_caf) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); StgWord label_len = strlen(label); StgWord module_len = strlen(module); StgWord srcloc_len = strlen(srcloc); @@ -1326,13 +1331,13 @@ void postHeapProfCostCentre(StgWord32 ccID, postStringLen(&eventBuf, module, module_len); postStringLen(&eventBuf, srcloc, srcloc_len); postWord8(&eventBuf, is_caf); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void postHeapProfSampleCostCentre(CostCentreStack *stack, StgWord64 residency) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); StgWord depth = 0; CostCentreStack *ccs; for (ccs = stack; ccs != NULL && ccs != CCS_MAIN; ccs = ccs->prevStack) @@ -1351,7 +1356,7 @@ void postHeapProfSampleCostCentre(CostCentreStack *stack, depth>0 && ccs != NULL && ccs != CCS_MAIN; ccs = ccs->prevStack, depth--) postWord32(&eventBuf, ccs->cc->ccID); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } @@ -1359,7 +1364,7 @@ void postProfSampleCostCentre(Capability *cap, CostCentreStack *stack, StgWord64 tick) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); StgWord depth = 0; CostCentreStack *ccs; for (ccs = stack; ccs != NULL && ccs != CCS_MAIN; ccs = ccs->prevStack) @@ -1377,7 +1382,7 @@ void postProfSampleCostCentre(Capability *cap, depth>0 && ccs != NULL && ccs != CCS_MAIN; ccs = ccs->prevStack, depth--) postWord32(&eventBuf, ccs->cc->ccID); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } // This event is output at the start of profiling so the tick interval can @@ -1385,11 +1390,11 @@ void postProfSampleCostCentre(Capability *cap, // can be calculated from how many samples there are. void postProfBegin(void) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); postEventHeader(&eventBuf, EVENT_PROF_BEGIN); // The interval that each tick was sampled, in nanoseconds postWord64(&eventBuf, TimeToNS(RtsFlags.MiscFlags.tickInterval)); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } #endif /* PROFILING */ @@ -1415,11 +1420,11 @@ static void postTickyCounterDef(EventsBuf *eb, StgEntCounter *p) void postTickyCounterDefs(StgEntCounter *counters) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); for (StgEntCounter *p = counters; p != NULL; p = p->link) { postTickyCounterDef(&eventBuf, p); } - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } static void postTickyCounterSample(EventsBuf *eb, StgEntCounter *p) @@ -1443,13 +1448,13 @@ static void postTickyCounterSample(EventsBuf *eb, StgEntCounter *p) void postTickyCounterSamples(StgEntCounter *counters) { - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); ensureRoomForEvent(&eventBuf, EVENT_TICKY_COUNTER_SAMPLE); postEventHeader(&eventBuf, EVENT_TICKY_COUNTER_BEGIN_SAMPLE); for (StgEntCounter *p = counters; p != NULL; p = p->link) { postTickyCounterSample(&eventBuf, p); } - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } #endif /* TICKY_TICKY */ void postIPE(const InfoProvEnt *ipe) @@ -1459,7 +1464,7 @@ void postIPE(const InfoProvEnt *ipe) // See Note [Maximum event length]. const StgWord MAX_IPE_STRING_LEN = 65535; - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); StgWord table_name_len = MIN(strlen(ipe->prov.table_name), MAX_IPE_STRING_LEN); StgWord closure_desc_len = MIN(strlen(closure_desc_buf), MAX_IPE_STRING_LEN); StgWord ty_desc_len = MIN(strlen(ipe->prov.ty_desc), MAX_IPE_STRING_LEN); @@ -1489,7 +1494,7 @@ void postIPE(const InfoProvEnt *ipe) postBuf(&eventBuf, &colon, 1); postStringLen(&eventBuf, ipe->prov.src_span, src_span_len); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); } void printAndClearEventBuf (EventsBuf *ebuf) @@ -1601,14 +1606,21 @@ void flushLocalEventsBuf(Capability *cap) // Flush all capabilities' event buffers when we already hold all capabilities. // Used during forkProcess. void flushAllCapsEventsBufs(void) +{ + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); + flushAllCapsEventsBufs_(); + RELEASE_LOCK_ALWAYS(&eventBufMutex); +} + +// Unsafe version that does not acquire/release eventBufMutex. You must +// hold the eventBufMutex, which you must acquire with ACQUIRE_LOCK_ALWAYS! +void flushAllCapsEventsBufs_(void) { if (!event_log_writer) { return; } - ACQUIRE_LOCK(&eventBufMutex); printAndClearEventBuf(&eventBuf); - RELEASE_LOCK(&eventBufMutex); for (unsigned int i=0; i < getNumCapabilities(); i++) { flushLocalEventsBuf(getCapability(i)); @@ -1641,9 +1653,9 @@ static void flushEventLog_(Capability **cap USED_IF_THREADS) return; } - ACQUIRE_LOCK(&eventBufMutex); + ACQUIRE_LOCK_ALWAYS(&eventBufMutex); printAndClearEventBuf(&eventBuf); - RELEASE_LOCK(&eventBufMutex); + RELEASE_LOCK_ALWAYS(&eventBufMutex); #if defined(THREADED_RTS) Task *task = newBoundTask(); ===================================== rts/eventlog/EventLog.h ===================================== @@ -18,6 +18,13 @@ #if defined(TRACING) extern bool eventlog_enabled; +#if defined(HAVE_PREEMPTION) +// Avoid using this mutex directly if at all possible. It is needed in the +// implementation of forkProcess. +// +// All uses of this mutex must use ACQUIRE_LOCK_ALWAYS/RELEASE_LOCK_ALWAYS. +extern Mutex eventBufMutex; +#endif void initEventLogging(void); void restartEventLogging(void); @@ -27,6 +34,7 @@ void abortEventLogging(void); // #4512 - after fork child needs to abort void moreCapEventBufs (uint32_t from, uint32_t to); void flushLocalEventsBuf(Capability *cap); void flushAllCapsEventsBufs(void); +void flushAllCapsEventsBufs_(void); void flushAllEventsBufs(Capability *cap); typedef void (*EventlogInitPost)(void); ===================================== rts/include/rts/OSThreads.h ===================================== @@ -14,6 +14,46 @@ #pragma once +/* Note [Threads and preemption] + ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + All full-fat OSs that GHC works on have OS threads, and we use them even in + the non-threaded RTS for a few features: + * Haskell thread preemption; + * sample-based profiling; + * idle GC; + * periodic eventlog flushing. + + We use defined(HAVE_PREEMPTION) to decide if these features are implemented + via OS threads. + + On platforms like WASM/js we do not have OS threads in any conventional + sense, and the features above are either not available or are implemented + differently. See Note [No timer on wasm32]. + + In future if GHC is ported to platforms like bare-metal micro-controllers, + RTOSs or to run directly under hypervisors then such platforms may also not + have threads available and they should not define HAVE_PREEMPTION here. Or + for some micro-controller RTOSs like Zeypher one may have a choice about + whether to use threads or not (at a size cost). Here would be the right + place to control whether the feature list above is supported. + */ +#if defined(wasm32_HOST_ARCH) + // See Note [No timer on wasm32] + // To confuse matters, WASM _does_ have pthread.h but it doesnt work. +#elif defined(HAVE_PTHREAD_H) || defined(HAVE_WINDOWS_H) +#define HAVE_PREEMPTION +#else +#error Decide if this platform has threads and pre-emption or not. +#endif +// And JS does all of this differently, without using this bit of the RTS. + +// Configuration sanity check +#if defined(THREADED_RTS) && !defined(HAVE_PREEMPTION) +//TODO we would like to be able to assert this: +// #error Configuration error: THREADED_RTS should imply HAVE_PREEMPTION +// however at the moment we cannot due to issue #27346. +#endif + #if defined(HAVE_PTHREAD_H) && !defined(mingw32_HOST_OS) #if defined(CMINUSMINUS) @@ -210,9 +250,29 @@ extern bool timedWaitCondition ( Condition* pCond, Mutex* pMut, Time timeout) // // Mutexes // +// Even in the non-threaded RTS we use threads and mutexes! In particular the +// timer/ticker is implemented using a thread. And using threads needs locks. +// In particular we need locks for the data shared between the timer/ticker +// thread and the thread running the main capability. +#if defined(HAVE_PREEMPTION) extern void initMutex ( Mutex* pMut ); extern void closeMutex ( Mutex* pMut ); +// The "always" variants do locking in the threaded and non-threaded RTS. +// The normal variants below are no-ops in the non-threaded RTS. +#define ACQUIRE_LOCK_ALWAYS(l) OS_ACQUIRE_LOCK(l) +#define TRY_ACQUIRE_LOCK_ALWAYS(l) OS_TRY_ACQUIRE_LOCK(l) +#define RELEASE_LOCK_ALWAYS(l) OS_RELEASE_LOCK(l) +#define ASSERT_LOCK_HELD_ALWAYS(l) OS_ASSERT_LOCK_HELD(l) +#else +// And just to be a bit confusing, the always variants are still no-ops when we +// do not HAVE_PREEMPTION, since then we don't have threads or mutexes at all. +#define ACQUIRE_LOCK_ALWAYS(l) +#define TRY_ACQUIRE_LOCK_ALWAYS(l) 0 +#define RELEASE_LOCK_ALWAYS(l) +#define ASSERT_LOCK_HELD_ALWAYS(l) +#endif + // Processors and affinity void setThreadAffinity (uint32_t n, uint32_t m); void setThreadNode (uint32_t node); @@ -228,6 +288,7 @@ void releaseThreadNode (void); #else +// No-ops in the non-threaded RTS. See also the _ALWAYS variants above. #define ACQUIRE_LOCK(l) #define TRY_ACQUIRE_LOCK(l) 0 #define RELEASE_LOCK(l) ===================================== rts/include/rts/Timer.h ===================================== @@ -13,6 +13,6 @@ #pragma once -void startTimer (void); -void stopTimer (void); +void startTimer (void); // Deprecated: see issue #27073 +void stopTimer (void); // Deprecated: see issue #27073 int rtsTimerSignal (void); // Deprecated: see issue #27073 ===================================== rts/posix/Ticker.c ===================================== @@ -103,120 +103,112 @@ #include <unistd.h> #include <fcntl.h> -static Time itimer_interval = DEFAULT_TICK_INTERVAL; -// Should we be firing ticks? -// Writers to this must hold the mutex below. -static bool stopped = false; +// Forward declarations of local types and helper functions to hide the +// difference between ppoll() and select() +#if defined(HAVE_DECL_PPOLL) && HAVE_DECL_PPOLL == 1 +typedef struct timespec timeout; // for ppoll() +typedef struct { struct pollfd pollfds[1]; } fdset; +#else +typedef struct timeval timeout; // for select() +typedef struct { int fd; fd_set selectfds; } fdset; // need to stash fd +#endif +static void poll_init_timeout(timeout *tv, Time t); +static void poll_init_fdset(fdset *fds, int fd); // single fd only +// poll_*_timeout returns >0 if fd ready, ==0 if timeout, <0 if error +static int poll_no_timeout(fdset *fdset); +static int poll_with_timeout(fdset *fdset, timeout *t); + -// should the ticker thread exit? -// This can be set without holding the mutex. -static bool exited = true; +static Time ticker_interval = DEFAULT_TICK_INTERVAL; -// Signaled when we want to (re)start the timer -static Condition start_cond; -static Mutex mutex; -static OSThreadId thread; +// Atomic variable used by client threads to communicate that they want the +// ticker thread to pause. This communication is one-way, with no +// acknowledgement. +static bool pause_request; -// fds for interrupting the ticker -static int interruptfd_r = -1, interruptfd_w = -1; +// Atomic variable used by other threads to communicate that they want the +// ticker thread to exit. +static bool exit_request; -static void *itimer_thread_func(void *_handle_tick) +// Used to wait for the ticker thread to terminate after asking it to exit. +static OSThreadId ticker_thread_id; + +// Fds used with sendFdWakeup to notify the ticker thread that any of the +// *_request variables above have been set. +static int notifyfd_r = -1, notifyfd_w = -1; + +static void *ticker_thread_func(void *_handle_tick) { TickProc handle_tick = _handle_tick; -#if defined(HAVE_DECL_PPOLL) && HAVE_DECL_PPOLL == 1 - struct pollfd pollfds[1]; - - pollfds[0].fd = interruptfd_r; - pollfds[0].events = POLLIN; + // Thread-local view of our state. We compare these with the corresponding + // atomic shared variables used to request state changes. + bool paused = true; // updated from atomic shared var pause_request + bool exit = false; // updated from atomic shared var exit_request + // Note that we start paused. - struct timespec ts = { .tv_sec = TimeToSeconds(itimer_interval) - , .tv_nsec = TimeToNS(itimer_interval) % 1000000000 - }; -#else - fd_set selectfds; - FD_ZERO(&selectfds); - FD_SET(interruptfd_r, &selectfds); - - struct timeval tv = { .tv_sec = TimeToSeconds(itimer_interval) - /* convert remainder time in nanoseconds - to microseconds, rounding up: */ - , .tv_usec = ((TimeToNS(itimer_interval) % 1000000000) - + 999) / 1000 - }; -#endif + timeout timeout; + fdset fdset; + poll_init_timeout(&timeout, ticker_interval); + poll_init_fdset(&fdset, notifyfd_r); - // Relaxed is sufficient: If we don't see that exited was set in one iteration we will - // see it next time. - while (!RELAXED_LOAD_ALWAYS(&exited)) { + while (!exit) { -#if defined(HAVE_DECL_PPOLL) && HAVE_DECL_PPOLL == 1 - int nfds = 1; - int nready = ppoll(pollfds, nfds, &ts, NULL); -#else - struct timeval tv_tmp = tv; // copy since select may change this value. - int nfds = interruptfd_r+1; - int nready = select(nfds, &selectfds, NULL, NULL, &tv_tmp); -#endif - // In either case (ppoll or select), the result nready is the number - // of fds that are ready. - if (RTS_LIKELY(nready == 0)) { - // Timer expired, not interrupted, continue. - } else if (nready > 0) { - // We only monitor one fd (the interruptfd_r), so we know - // it is that fd that is ready without any further checks. - collectFdWakeup(interruptfd_r); - // No further action needed, continue on to handling the final tick - // and then stop. - - // Note that we rely on sendFdWakeup and select/poll to provide the - // happens-before relation. So if 'exited' was set before calling - // sendFdWakeup, then we should be able to reliably read it after. - // And thus reading 'exited' in the while loop guard is ok. + int notify; + if (paused) { + notify = poll_no_timeout(&fdset); } else { - // While the RTS attempts to mask signals, some foreign libraries - // that rely on signal delivery may unmask them. Consequently we - // may see EINTR. See #24610. - if (errno != EINTR) { - sysErrorBelch("Ticker: poll failed: %s", strerror(errno)); - } + notify = poll_with_timeout(&fdset, &timeout); } - // first try a cheap test - if (RELAXED_LOAD_ALWAYS(&stopped)) { - OS_ACQUIRE_LOCK(&mutex); - // should we really stop? - if (stopped) { - waitCondition(&start_cond, &mutex); - } - OS_RELEASE_LOCK(&mutex); - } else { + if (RTS_LIKELY(notify == 0)) { + // The time expired, no state change notification. handle_tick(0); + + } else if (notify > 0) { + // State change notification, check the request variables. + + // We rely on sendFdWakeup and select/poll to provide the + // happens-before relation. So if the request variables are set + // before calling sendFdWakeup, then we should be able to reliably + // read them here afterwards. + collectFdWakeup(notifyfd_r); + + paused = ACQUIRE_LOAD_ALWAYS(&pause_request); + exit = RELAXED_LOAD_ALWAYS(&exit_request); + } else if (errno != EINTR) { + // While the RTS attempts to mask signals, some foreign libraries + // that rely on signal delivery may unmask them. Consequently we + // may see EINTR. See #24610. + sysErrorBelch("Ticker: poll failed: %s", strerror(errno)); } } return NULL; } +/* Initialise the ticker on startup or re-initialise the ticker after a fork(). + * In the fork case, the thread will not be present, but fds are inherited. + * + * The ticker is started in the paused state. Use unpauseTicker to continue. + */ void initTicker (Time interval, TickProc handle_tick) { - itimer_interval = interval; - stopped = true; - exited = false; + ticker_interval = interval; + pause_request = true; + exit_request = false; + #if defined(HAVE_SIGNAL_H) sigset_t mask, omask; int sigret; #endif int ret; - initCondition(&start_cond); - initMutex(&mutex); - /* Open the interrupt fd synchronously. * - * We used to do it in itimer_thread_func (i.e. in the timer thread) but it + * We used to do it in ticker_thread_func (i.e. in the timer thread) but it * meant that some user code could run before it and get confused by the * allocation of the timerfd. * @@ -226,11 +218,11 @@ initTicker (Time interval, TickProc handle_tick) * descriptor closed by the first call! (see #20618) */ - if (interruptfd_r != -1) { + if (notifyfd_r != -1) { // don't leak the old file descriptors after a fork (#25280) - closeFdWakeup(interruptfd_r, interruptfd_w); + closeFdWakeup(notifyfd_r, notifyfd_w); } - newFdWakeup(&interruptfd_r, &interruptfd_w); + newFdWakeup(¬ifyfd_r, ¬ifyfd_w); /* * Create the thread with all blockable signals blocked, leaving signal @@ -242,7 +234,7 @@ initTicker (Time interval, TickProc handle_tick) sigfillset(&mask); sigret = pthread_sigmask(SIG_SETMASK, &mask, &omask); #endif - ret = createAttachedOSThread(&thread, "ghc_ticker", itimer_thread_func, (void*)handle_tick); + ret = createAttachedOSThread(&ticker_thread_id, "ghc_ticker", ticker_thread_func, (void*)handle_tick); #if defined(HAVE_SIGNAL_H) if (sigret == 0) pthread_sigmask(SIG_SETMASK, &omask, NULL); @@ -253,47 +245,99 @@ initTicker (Time interval, TickProc handle_tick) } } -void -startTicker(void) +/* Asynchronous. Idempotent. */ +void unpauseTicker(void) { - OS_ACQUIRE_LOCK(&mutex); - RELAXED_STORE(&stopped, false); - signalCondition(&start_cond); - OS_RELEASE_LOCK(&mutex); + RELEASE_STORE_ALWAYS(&pause_request, false); + sendFdWakeup(notifyfd_w); } -/* There may be at most one additional tick fired after a call to this */ -void -stopTicker(void) +/* Asynchronous. Idempotent. + * There may be at additional ticks fired after a call to this, but it will + * usually stop quickly. + */ +void pauseTicker(void) { - OS_ACQUIRE_LOCK(&mutex); - RELAXED_STORE(&stopped, true); - OS_RELEASE_LOCK(&mutex); + RELEASE_STORE_ALWAYS(&pause_request, true); + sendFdWakeup(notifyfd_w); } -/* There may be at most one additional tick fired after a call to this */ -void -exitTicker (bool wait) +/* Synchronous. Not idempotent. + * The ticker is guaranteed stopped after this. + */ +void exitTicker(void) { - ASSERT(!SEQ_CST_LOAD(&exited)); - SEQ_CST_STORE(&exited, true); - // ensure that ticker wakes up if stopped - startTicker(); - sendFdWakeup(interruptfd_w); - - // wait for ticker to terminate if necessary - if (wait) { - if (pthread_join(thread, NULL)) { - sysErrorBelch("Ticker: Failed to join: %s", strerror(errno)); - } - closeFdWakeup(interruptfd_r, interruptfd_w); - closeMutex(&mutex); - closeCondition(&start_cond); - } else { - pthread_detach(thread); + ASSERT(!RELAXED_LOAD_ALWAYS(&exit_request)); + RELEASE_STORE_ALWAYS(&exit_request, true); + sendFdWakeup(notifyfd_w); + + // wait for ticker thread to terminate + if (pthread_join(ticker_thread_id, NULL)) { + sysErrorBelch("Ticker: Failed to join: %s", strerror(errno)); } + closeFdWakeup(notifyfd_r, notifyfd_w); +} + +/* Implementation of the local helpers, to hide the difference between ppoll() + * and select(). + */ +#if defined(HAVE_DECL_PPOLL) && HAVE_DECL_PPOLL == 1 +static void poll_init_timeout(timeout *tv, Time t) +{ + tv->tv_sec = TimeToSeconds(t); + tv->tv_nsec = TimeToNS(t) % 1000000000; } +static void poll_init_fdset(fdset *fds, int fd) +{ + fds->pollfds[0].fd = fd; + fds->pollfds[0].events = POLLIN; +} + +static int poll_no_timeout(fdset *fds) +{ + int nfds = 1; + return ppoll(fds->pollfds, nfds, NULL, NULL); +} + +static int poll_with_timeout(fdset *fds, timeout *ts) +{ + int nfds = 1; + return ppoll(fds->pollfds, nfds, ts, NULL); +} + +#else // select() + +static void poll_init_timeout(timeout *tv, Time t) +{ + tv->tv_sec = TimeToSeconds(t); + /* convert remainder time in nanoseconds to microseconds, rounding up: */ + tv->tv_usec = ((TimeToNS(t) % 1000000000) + 999) / 1000; +} + +static void poll_init_fdset(fdset *fds, int fd) +{ + /* select() overwrites the fdset so we must rebuild it every time. */ + FD_ZERO(&fds->selectfds); + FD_SET(fd, &fds->selectfds); + fds->fd = fd; +} + +static int poll_no_timeout(fdset *fds) +{ + int nfds = fds->fd+1; + return select(nfds, &fds->selectfds, NULL, NULL, NULL); +} + +static int poll_with_timeout(fdset *fds, timeout *tv) +{ + struct timeval tv_tmp = *tv; // copy since select may change this value. + int nfds = fds->fd+1; + return select(nfds, &fds->selectfds, NULL, NULL, &tv_tmp); +} +#endif + +/* This is obsolete, but is used in the unix package for now */ int rtsTimerSignal(void) { ===================================== rts/win32/Ticker.c ===================================== @@ -11,7 +11,11 @@ static TickProc tick_proc = NULL; static HANDLE timer_queue = NULL; + +static Mutex lock; // To protect the timer and paused var below static HANDLE timer = NULL; +static bool paused; + static Time tick_interval = 0; static VOID CALLBACK tick_callback( @@ -36,12 +40,19 @@ static VOID CALLBACK tick_callback( // This seems to be the case starting at some point during the // Windows 7 lifetime and any newer versions of windows. +// Forward decls +static void startTicker(void); +static void stopTicker(bool synchronous); + void initTicker (Time interval, TickProc handle_tick) { + ASSERT(timer_queue == NULL); tick_interval = interval; tick_proc = handle_tick; + OS_INIT_LOCK(&lock); + paused = true; // starts paused timer_queue = CreateTimerQueue(); if (timer_queue == NULL) { sysErrorBelch("CreateTimerQueue"); @@ -49,39 +60,81 @@ initTicker (Time interval, TickProc handle_tick) } } +// Asynchronous. Idempotent. void -startTicker(void) +unpauseTicker(void) { - BOOL r; - - r = CreateTimerQueueTimer(&timer, - timer_queue, - tick_callback, - 0, - 0, - TimeToMS(tick_interval), // ms - WT_EXECUTEINTIMERTHREAD); - if (r == 0) { - sysErrorBelch("CreateTimerQueueTimer"); - stg_exit(EXIT_FAILURE); + OS_ACQUIRE_LOCK(&lock); + if (paused) { + startTicker(); } + paused = false; + OS_RELEASE_LOCK(&lock); } +// Asynchronous. Idempotent. void -stopTicker(void) +pauseTicker(void) { - if (timer_queue != NULL && timer != NULL) { - DeleteTimerQueueTimer(timer_queue, timer, NULL); - timer = NULL; + OS_ACQUIRE_LOCK(&lock); + if (!paused) { + /* pauseTicker is called from within the handle_tick, so stopping + * the ticker here /must/ be asynchronous or we will deadlock! */ + stopTicker(false /* asynchronous */); } + paused = true; + OS_RELEASE_LOCK(&lock); } void -exitTicker (bool wait) +exitTicker(void) { - stopTicker(); - if (timer_queue != NULL) { - DeleteTimerQueueEx(timer_queue, wait ? INVALID_HANDLE_VALUE : NULL); - timer_queue = NULL; + ASSERT(timer_queue != NULL); + + OS_ACQUIRE_LOCK(&lock); + if (!paused) { + stopTicker(true /* synchronous */); + } + OS_RELEASE_LOCK(&lock); + + // From the docs for DeleteTimerQueueEx: + // If this parameter is INVALID_HANDLE_VALUE, the function waits + // for all callback functions to complete before returning. + // This is a belt-and-braces approach to ensuring exitTicker is synchronous, + // since stopTicker(true) is already synchronous and there's only one timer. + HANDLE completion = INVALID_HANDLE_VALUE; + DeleteTimerQueueEx(timer_queue, completion); + timer_queue = NULL; +} + +static void startTicker(void) { + ASSERT(timer_queue != NULL && timer == NULL); + DWORD interval = TimeToMS(tick_interval); // ms + BOOL r = CreateTimerQueueTimer(&timer, + timer_queue, + tick_callback, + NULL, // callback param + interval, // inital interval + interval, // recurrant interval + WT_EXECUTEINTIMERTHREAD); + //TODO: using WT_EXECUTEINTIMERTHREAD is fine for context switching, and + // plausibly also ok for profile sampling but is way out for eventlog + // flushing. The eventlog flush does a global synchronisation of all + // capabilities and I/O! And with eventlog providers, it calls arbitrary + // user code. This is not ok! See issue #27250. + if (r == 0) { + sysErrorBelch("CreateTimerQueueTimer"); + stg_exit(EXIT_FAILURE); } + ASSERT(timer != NULL); +} + +static void stopTicker(bool synchronous) { + ASSERT(timer_queue != NULL && timer != NULL); + // From the docs for DeleteTimerQueueTimer: + // If this parameter is INVALID_HANDLE_VALUE, the function waits for any + // running timer callback functions to complete before returning. + HANDLE completion = synchronous ? INVALID_HANDLE_VALUE : NULL; + DeleteTimerQueueTimer(timer_queue, timer, completion); + timer = NULL; } ===================================== testsuite/tests/concurrent/should_run/T27105.hs ===================================== @@ -0,0 +1,114 @@ +{-# OPTIONS_GHC -fno-omit-yields #-} + +import Control.Monad +import Control.Monad.ST +import Control.Concurrent +import Control.Exception +import System.Exit +import System.Mem +import GHC.Arr +import Prelude hiding (init) + +-- Test thread fairness: +-- run two cpu-bound threads concurrently for a second, +-- each counts how many operations it can perform until signaled to stop +-- expect a balance between the two with no more than a 75% imperfection. +-- Yes, 75%! On the CI machines we occasionally observe extraordinary levels +-- of unfairness: nearly 60% in some cases. We don't want this to become a +-- fragile test that is ignored, so we use an extreme bound. This should still +-- catch gross breakage. +-- +-- This _should_ detect if the interval timer is not working, or if thread +-- context switching is messed up. We can expect failure if we force a +-- contex switch interval of more than half the test time, i.e. more than 0.5s +-- +-- We run the test twice, with allocating and non-allocating worker threads. +-- The -fno-omit-yields above is crucial for worker_nonalloc below, or it never +-- gets interrupted and thus no context switches. + +main :: IO () +main = do + test worker_alloc + performMajorGC + test worker_nonalloc + +test :: Worker -> IO () +test worker = do + stop <- newEmptyMVar + res1 <- newEmptyMVar + res2 <- newEmptyMVar + _ <- forkIO (worker stop >>= putMVar res1) + _ <- forkIO (worker stop >>= putMVar res2) + threadDelay 300_000 + -- Let them run for 300ms. The default context switch interval is 20ms. + -- This gives time for 15 context switches, so this _should_ be enough + -- to get less than 10% unfairness. And on most platforms it is enough. + -- But OSX! Oh OSX! How do I loath thee? Let me count++ the ways. + -- To avoid a fragile test, we use a 75% unfairness threshold. + putMVar stop () + count1 <- takeMVar res1 + count2 <- takeMVar res2 + let balance :: Double + balance = abs ((fromIntegral count1 - fromIntegral count2) + / fromIntegral count2) + when (balance > 0.75) $ do + putStrLn "Schedule fairness more than 75% tolerance:" + putStrLn $ "imperfection: " ++ show (balance * 100) ++ "%" + putStrLn $ "work counts: " ++ show (count1, count2) + exitFailure + +type Worker = MVar () -> IO Int + +-- count how many iterations we can calculate until we're signaled to stop +worker_template :: IO a -> (a -> IO ()) -> MVar () -> IO Int +worker_template init iter stop = do + a <- init + go a 0 + where + go a !count = do + ok <- tryReadMVar stop + case ok of + Just () -> return count + Nothing -> do + iter a + go a (count + 1) + + +-- the allocating worker +{-# NOINLINE worker_alloc #-} +worker_alloc :: Worker +worker_alloc = + worker_template + (return 18) + (\n -> evaluate (fib n) >> return ()) + +-- by forcing this to be Integer we cause lots of allocation! +fib :: Integer -> Integer +fib 0 = 0 +fib 1 = 1 +fib n = fib (n-1) + fib (n-2) + + +-- the non-allocating worker +{-# NOINLINE worker_nonalloc #-} +worker_nonalloc :: Worker +worker_nonalloc = + worker_template + (stToIO $ newSTArray (0,50_000) 42) + (\arr -> stToIO $ arrrev arr) + +arrrev :: STArray s Int Int -> ST s () +arrrev arr = + let (i,j) = boundsSTArray arr + in arrrev_go arr i j + +{-# NOINLINE arrrev_go #-} +arrrev_go :: STArray s Int Int -> Int -> Int -> ST s () +arrrev_go !_ !i !j | i >= j = return () +arrrev_go !arr !i !j = do + x <- readSTArray arr i + y <- readSTArray arr j + writeSTArray arr i y + writeSTArray arr j x + arrrev_go arr (i+1) (j-1) + ===================================== testsuite/tests/concurrent/should_run/all.T ===================================== @@ -325,3 +325,15 @@ test('T26341b' # test uses pipe operations which are not supported by the JS/wasm backends , when(arch('wasm32') or arch('javascript'), skip) , compile_and_run, ['-package process']) + +# Scheduler (very rough) fairness +test('T27105', + [when(arch('wasm32'), skip), # same reason as T367_letnoescape + run_timeout_multiplier(0.05)], # we expect this to run for ~2s + compile_and_run, ['']) +test('T27105_fail', + [when(arch('wasm32'), skip), + # And we can expect it to fail if we context switch too coarsely + extra_run_opts('+RTS -C0.2 -RTS'), expect_fail, + run_timeout_multiplier(0.05)], + multimod_compile_and_run, ['T27105.hs', '']) View it on GitLab: https://gitlab.haskell.org/ghc/ghc/-/compare/6a667450ea06ed9c978caf09c169f3e... -- View it on GitLab: https://gitlab.haskell.org/ghc/ghc/-/compare/6a667450ea06ed9c978caf09c169f3e... You're receiving this email because of your account on gitlab.haskell.org.
participants (1)
-
Duncan Coutts (@dcoutts)