diff --git a/cf-reactor/Makefile.am b/cf-reactor/Makefile.am index 9961083ee36..9d9a2e96aae 100644 --- a/cf-reactor/Makefile.am +++ b/cf-reactor/Makefile.am @@ -41,6 +41,7 @@ libcf_reactor_la_SOURCES = \ cf-reactor.c \ file_watcher.c file_watcher.h \ reactor_context.c reactor_context.h \ + reactor_transform.c reactor_transform.h \ stoppable_thread.c stoppable_thread.h \ wakeup_channel.c wakeup_channel.h \ watcher.c watcher.h diff --git a/cf-reactor/cf-reactor.c b/cf-reactor/cf-reactor.c index 55ba9c3d873..e6160e5d611 100644 --- a/cf-reactor/cf-reactor.c +++ b/cf-reactor/cf-reactor.c @@ -40,6 +40,7 @@ #include /* GetSignalPipe, MakeSignalPipe, IsPendingTermination, HandleSignalsForDaemon, ReloadConfigRequested, ClearRequestReloadConfig */ #include #include +#include /*****************************************************************************/ /* Globals */ @@ -244,6 +245,11 @@ static void CheckPolicyUpdates(EvalContext *ctx, Policy **policy, GenericAgentCo GenericAgentDiscoverContext(ctx, config, NULL); *policy = LoadPolicy(ctx, config); + + if (*policy != NULL) + { + KeepReactorPromises(ctx, *policy); + } } /*****************************************************************************/ @@ -326,6 +332,10 @@ int main(int argc, char *argv[]) * below. */ time_t next_tick = time(NULL) + DEFAULT_POLL_INTERVAL_SECS; time_t next_policy_check = time(NULL) + DEFAULT_POLICY_CHECK_INTERVAL_SECS; + + // Setup event watchers from policy + KeepReactorPromises(ctx, policy); + while (!IsPendingTermination()) { int max_fd = ReactorContextSetupFileDescriptors(&reactor_ctx); diff --git a/cf-reactor/reactor_transform.c b/cf-reactor/reactor_transform.c new file mode 100644 index 00000000000..6c8e6c96022 --- /dev/null +++ b/cf-reactor/reactor_transform.c @@ -0,0 +1,202 @@ +/* + Copyright 2026 Northern.tech AS + + This file is part of CFEngine 3 - written and maintained by Northern.tech AS. + + This program is free software; you can redistribute it and/or modify it + under the terms of the GNU General Public License as published by the + Free Software Foundation; version 3. + + This program is distributed in the hope that it will be useful, + but WITHOUT ANY WARRANTY; without even the implied warranty of + MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + GNU General Public License for more details. + + You should have received a copy of the GNU General Public License + along with this program; if not, write to the Free Software + Foundation, Inc., 59 Temple Place - Suite 330, Boston, MA 02111-1307, USA + + To the extent this program is licensed as part of the Enterprise + versions of CFEngine, the applicable Commercial Open Source License + (COSL) may apply to this file if you as a licensee so wish it. See + included file COSL.txt. +*/ + +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +/* Promise types evaluated within `bundle reactor NAME { ... }`. */ +static const char *const REACTOR_TYPESEQUENCE[] = +{ + "meta", + "vars", + "classes", + "events", + NULL +}; + +/* Checks that the bundle referred to by an events promise's 'then' attribute, + * either `then => "name"` or `then => name(args)`, is defined as an agent or + * common bundle. Logs an error and returns false if it is not. */ +static bool ThenBundleExists(const EvalContext *ctx, const Promise *pp, Rval then_rval, const char *key) +{ + const char *bundle_name = NULL; + switch (then_rval.type) + { + case RVAL_TYPE_SCALAR: + bundle_name = RvalScalarValue(then_rval); + break; + case RVAL_TYPE_FNCALL: + bundle_name = RvalFnCallValue(then_rval)->name; + break; + default: + break; + } + + const Bundle *bundle = NULL; + if (bundle_name != NULL) + { + bundle = EvalContextResolveBundleExpression(ctx, PromiseGetPolicy(pp), bundle_name, "agent"); + if (bundle == NULL) + { + bundle = EvalContextResolveBundleExpression(ctx, PromiseGetPolicy(pp), bundle_name, "common"); + } + } + if (bundle == NULL) + { + Log(LOG_LEVEL_ERR, "Reactor events promise '%s' refers to unknown bundle '%s', ignoring", + key, (bundle_name != NULL) ? bundle_name : "(invalid)"); + return false; + } + + return true; +} + +// Temporary limitations: +// - 'when' body: only a single constraint, 'file_deleted', is supported (no OR-ing of constraints yet) +// - 'then': only a single bundle is supported, not a list of bundles yet +// TODO: support multiple 'when' constraints and a list of 'then' bundles +static PromiseResult KeepEventsPromise(EvalContext *ctx, const Promise *pp) +{ + assert(pp != NULL); + + const Bundle *bp = PromiseGetBundle(pp); + char *key = StringFormat("%s:%s:%s", bp->ns, bp->name, pp->promiser); + + const char *path = PromiseGetConstraintAsRval(pp, "file_deleted", RVAL_TYPE_SCALAR); + if (path == NULL) + { + Log(LOG_LEVEL_WARNING, + "Reactor events promise '%s' must specify exactly one file to watch in its 'when' body, ignoring", + key); + free(key); + return PROMISE_RESULT_FAIL; + } + + const Constraint *then_constraint = PromiseGetConstraint(pp, "then"); + if (then_constraint == NULL) + { + Log(LOG_LEVEL_ERR, "Reactor events promise '%s' does not specify a 'then' bundle, ignoring", key); + free(key); + return PROMISE_RESULT_FAIL; + } + + if (!ThenBundleExists(ctx, pp, then_constraint->rval, key)) + { + free(key); + return PROMISE_RESULT_FAIL; + } + + // register watcher + Log(LOG_LEVEL_INFO, "Registering a file_deleted watcher with key '%s', on file '%s'", key, path); + bool kept = WatcherRegister(key, EVENT_FILE_DELETED, FileWatcherStateNew(path), then_constraint->rval, 1); + free(key); + + return (kept) ? PROMISE_RESULT_NOOP : PROMISE_RESULT_FAIL; +} + +static PromiseResult KeepReactorPromise(EvalContext *ctx, const Promise *pp, ARG_UNUSED void *param) +{ + assert(param == NULL); + PromiseBanner(ctx, pp); + + if (StringEqual(PromiseGetPromiseType(pp), "vars") || + StringEqual(PromiseGetPromiseType(pp), "meta")) + { + return VerifyVarPromise(ctx, pp, NULL); + } + + if (StringEqual(PromiseGetPromiseType(pp), "classes")) + { + return VerifyClassPromise(ctx, pp, NULL); + } + + if (StringEqual(PromiseGetPromiseType(pp), "events")) + { + return KeepEventsPromise(ctx, pp); + } + + return PROMISE_RESULT_NOOP; +} + +static void EvaluateReactorBundle(EvalContext *ctx, const Bundle *bp) +{ + assert(bp != NULL); + EvalContextStackPushBundleFrame(ctx, bp, NULL, false, NULL); + + for (int i = 0; REACTOR_TYPESEQUENCE[i] != NULL; i++) + { + const BundleSection *sp = BundleGetSection(bp, REACTOR_TYPESEQUENCE[i]); + if (sp == NULL || SeqLength(sp->promises) == 0) + { + Log(LOG_LEVEL_DEBUG, "No promise type %s in bundle %s", + REACTOR_TYPESEQUENCE[i], bp->name); + continue; + } + + EvalContextStackPushBundleSectionFrame(ctx, sp); + for (size_t j = 0; j < SeqLength(sp->promises); j++) + { + const Promise *pp = SeqAt(sp->promises, j); + ExpandPromise(ctx, pp, KeepReactorPromise, NULL); + } + EvalContextStackPopFrame(ctx); + } + + EvalContextStackPopFrame(ctx); +} + +void KeepReactorPromises(EvalContext *ctx, const Policy *policy) +{ + assert(policy != NULL); + WatcherRegistryClear(); + + for (size_t i = 0; i < SeqLength(policy->bundles); i++) + { + const Bundle *bp = SeqAt(policy->bundles, i); + if (!StringEqual(bp->type, CF_AGENTTYPES[AGENT_TYPE_REACTOR])) + { + continue; + } + + if (RlistLen(bp->args) > 0) + { + Log(LOG_LEVEL_WARNING, + "Cannot implicitly evaluate bundle '%s %s', as this bundle takes arguments.", + bp->type, bp->name); + continue; + } + + EvaluateReactorBundle(ctx, bp); + } +} diff --git a/cf-reactor/reactor_transform.h b/cf-reactor/reactor_transform.h new file mode 100644 index 00000000000..b089a9b53cb --- /dev/null +++ b/cf-reactor/reactor_transform.h @@ -0,0 +1,42 @@ +/* + Copyright 2026 Northern.tech AS + + This file is part of CFEngine 3 - written and maintained by Northern.tech AS. + + This program is free software; you can redistribute it and/or modify it + under the terms of the GNU General Public License as published by the + Free Software Foundation; version 3. + + This program is distributed in the hope that it will be useful, + but WITHOUT ANY WARRANTY; without even the implied warranty of + MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + GNU General Public License for more details. + + You should have received a copy of the GNU General Public License + along with this program; if not, write to the Free Software + Foundation, Inc., 59 Temple Place - Suite 330, Boston, MA 02111-1307, USA + + To the extent this program is licensed as part of the Enterprise + versions of CFEngine, the applicable Commercial Open Source License + (COSL) may apply to this file if you as a licensee so wish it. See + included file COSL.txt. +*/ + +#ifndef CFENGINE_REACTOR_TRANSFORM_H +#define CFENGINE_REACTOR_TRANSFORM_H + +#include +#include + +/** + * @brief Evaluate every `bundle reactor NAME { ... }` in the given policy: + * keep its "meta", "vars" and "classes" promises, and register a watcher + * (see watcher.h) for each of its "events" promises. + * + * Safe to call again for every re-read of the policy: it first discards all + * previously registered watchers (see WatcherRegistryClear()), so the + * result always reflects only the given policy, not a stale prior one. + */ +void KeepReactorPromises(EvalContext *ctx, const Policy *policy); + +#endif diff --git a/cf-reactor/watcher.c b/cf-reactor/watcher.c index c2b693dc653..e0be04da879 100644 --- a/cf-reactor/watcher.c +++ b/cf-reactor/watcher.c @@ -34,6 +34,8 @@ #include #include #include +#include +#include /* Upper bound on how long the watcher thread ever sleeps in one go, so that * IsPendingTermination() is re-checked at least this often during shutdown, @@ -52,11 +54,15 @@ typedef struct time_t next_due; } Watcher; -static void WatcherDestroy(void *item); /* defined below, next to WatcherRegister() */ +static void WatcherDestroy(void *item); +static void DestroyRval(void *item); -// TODO: potential race condition. If the policy is reparsed and watchers are re-registered -// while the watcher thread is iterating over this Seq, -// it may read a Watcher that is being freed or reallocated concurrently +/* Guards `watchers` (and, for the same atomic-rebuild reason, its paired + * `event_to_bundle`) against the watcher thread (WatcherThreadMain()) + * iterating over `watchers` concurrently with the main thread re-registering + * watchers from a freshly (re-)read policy (WatcherRegister(), + * WatcherRegistryClear()). */ +static pthread_mutex_t watchers_mutex = PTHREAD_MUTEX_INITIALIZER; static Seq *watchers = NULL; static Map *event_to_bundle = NULL; @@ -72,7 +78,7 @@ void WatcherRegistryInitialize(void) assert(event_to_bundle == NULL); watchers = SeqNew(4, WatcherDestroy); - event_to_bundle = MapNew(StringHash_untyped, StringEqual_untyped, NULL, NULL); + event_to_bundle = MapNew(StringHash_untyped, StringEqual_untyped, NULL, DestroyRval); } void WatcherRegistryFinalize(void) @@ -83,11 +89,42 @@ void WatcherRegistryFinalize(void) event_to_bundle = NULL; } +void WatcherRegistryClear(void) +{ + assert(watchers != NULL && event_to_bundle != NULL); + + ThreadLock(&watchers_mutex); + + SeqDestroy(watchers); + watchers = SeqNew(4, WatcherDestroy); + + MapDestroy(event_to_bundle); + event_to_bundle = MapNew(StringHash_untyped, StringEqual_untyped, NULL, DestroyRval); + + ThreadUnlock(&watchers_mutex); +} + +static Rval *AllocateRval(Rval val) +{ + Rval *new = xmalloc(sizeof(Rval)); + *new = RvalCopy(val); + return new; +} + +static void DestroyRval(void *item) +{ + Rval *rval = item; + if (rval != NULL) + { + RvalDestroy(*rval); + free(rval); + } +} + // Expects interval to be strictly greater than 0, otherwise the watcher thread will busy spin -void WatcherRegister(const char *key, EventType type, void *state, Bundle *bundle, time_t interval) +bool WatcherRegister(const char *key, EventType type, void *state, Rval val, time_t interval) { assert(key != NULL); - assert(bundle != NULL); assert(watchers != NULL && event_to_bundle != NULL); assert(interval > 0); @@ -106,14 +143,17 @@ void WatcherRegister(const char *key, EventType type, void *state, Bundle *bundl ProgrammingError("Unknown reactor event type %d for watcher '%s'", (int) type, key); } + ThreadLock(&watchers_mutex); + if (MapHasKey(event_to_bundle, key)) { Log(LOG_LEVEL_ERR, "Reactor watcher key '%s' is already registered, ignoring the duplicate", key); + ThreadUnlock(&watchers_mutex); if (destroy_state != NULL) { destroy_state(state); } - return; + return false; } Watcher *w = xmalloc(sizeof(Watcher)); @@ -125,7 +165,10 @@ void WatcherRegister(const char *key, EventType type, void *state, Bundle *bundl w->destroy_state = destroy_state; SeqAppend(watchers, w); - MapInsert(event_to_bundle, w->key, bundle); + MapInsert(event_to_bundle, w->key, AllocateRval(val)); + + ThreadUnlock(&watchers_mutex); + return true; } static void WatcherDestroy(void *item) @@ -154,6 +197,8 @@ static void WatcherThreadMain(StoppableThread *thread, ARG_UNUSED void *unused) time_t sleep_for = MAX_WATCHER_THREAD_SLEEP_SECS; bool any_event = false; + ThreadLock(&watchers_mutex); + for (size_t i = 0; i < SeqLength(watchers); i++) { Watcher *w = SeqAt(watchers, i); @@ -167,9 +212,7 @@ static void WatcherThreadMain(StoppableThread *thread, ARG_UNUSED void *unused) w->next_due = now + w->poll_interval_secs; if (fired) { - /** TODO: potential use after free. If a policy reparse destroys this Watcher (and frees w->key) - * before the queued key is consumed, the consumer will read freed memory */ - ThreadedQueuePush(event_queue, w->key); + ThreadedQueuePush(event_queue, SafeStringDuplicate(w->key)); any_event = true; } } @@ -178,6 +221,8 @@ static void WatcherThreadMain(StoppableThread *thread, ARG_UNUSED void *unused) sleep_for = MIN(sleep_for, until_due); } + ThreadUnlock(&watchers_mutex); + if (any_event) { WakeupChannelNotify(&wakeup_channel); @@ -198,7 +243,7 @@ bool EventWatcherInitialize(int *fd) return false; } - event_queue = ThreadedQueueNew(16, NULL); + event_queue = ThreadedQueueNew(16, free); watcher_thread = StoppableThreadStart(WatcherThreadMain, NULL); if (watcher_thread == NULL) @@ -231,9 +276,22 @@ void EventWatcherHandleEvents(int fd, fd_set *readfds) while (ThreadedQueuePop(event_queue, &item, 0)) { const char *key = item; - ARG_UNUSED const Bundle *bundle = MapGet(event_to_bundle, key); - Log(LOG_LEVEL_NOTICE, "Reactor watcher '%s' fired", key); - // TODO: run bundle + + ThreadLock(&watchers_mutex); + Rval *bundle = MapGet(event_to_bundle, key); + + if (bundle == NULL) + { + Log(LOG_LEVEL_VERBOSE, "Reactor watcher '%s' fired but is no longer registered, ignoring", key); + } + else + { + Log(LOG_LEVEL_NOTICE, "Reactor watcher '%s' fired", key); + // TODO: run bundle + } + + ThreadUnlock(&watchers_mutex); + free(item); } } diff --git a/cf-reactor/watcher.h b/cf-reactor/watcher.h index 4f458386e84..29a043b81d4 100644 --- a/cf-reactor/watcher.h +++ b/cf-reactor/watcher.h @@ -38,16 +38,22 @@ typedef void (*WatcherStateDestroyFn)(void *state); void WatcherRegistryInitialize(void); void WatcherRegistryFinalize(void); +/** + * @brief Discard every currently registered watcher and start over. Call + * this before re-registering watchers from a freshly (re-)read policy. + */ +void WatcherRegistryClear(void); + /** * @brief Register a specific watcher instance. * * @param key the events promise identifier * @param type the type of watcher, defined in when bodies * @param state the data used for by the watcher, depending on the type - * @param bundle the bundle to run on event + * @param val the rval holding the bundle to run on event * @param interval interval between runs */ -void WatcherRegister(const char *key, EventType type, void *state, Bundle *bundle, time_t interval); +bool WatcherRegister(const char *key, EventType type, void *state, Rval val, time_t interval); bool EventWatcherInitialize(int *fd); void EventWatcherHandleEvents(int fd, fd_set *readfds); void EventWatcherFinalize(void); diff --git a/libpromises/mod_reactor.c b/libpromises/mod_reactor.c index c5a1d738ce7..0bc16cd1fc8 100644 --- a/libpromises/mod_reactor.c +++ b/libpromises/mod_reactor.c @@ -31,7 +31,7 @@ static const ConstraintSyntax when_constraints[] = CONSTRAINT_SYNTAX_GLOBAL, /* Row models */ - ConstraintSyntaxNewStringList("files_deleted", CF_ANYSTRING, "List of files to react for on deletion", SYNTAX_STATUS_NORMAL), + ConstraintSyntaxNewString("file_deleted", CF_ANYSTRING, "File to react for on deletion", SYNTAX_STATUS_NORMAL), ConstraintSyntaxNewNull() }; diff --git a/libpromises/syntax.c b/libpromises/syntax.c index 484fb37826a..7bbc28a1b30 100644 --- a/libpromises/syntax.c +++ b/libpromises/syntax.c @@ -397,7 +397,7 @@ SyntaxTypeMatch CheckConstraintTypeMatch(const char *lval, Rval rval, DataType d /* Fn-like objects are assumed to be parameterized bundles in these... */ - checklist = SplitString("bundlesequence,edit_line,edit_xml,usebundle,service_bundle,home_bundle", ','); + checklist = SplitString("bundlesequence,edit_line,edit_xml,usebundle,service_bundle,home_bundle,then", ','); if (!IsItemIn(checklist, lval)) { diff --git a/tests/unit/watcher_test.c b/tests/unit/watcher_test.c index 6b783e3a34d..507bfce37d9 100644 --- a/tests/unit/watcher_test.c +++ b/tests/unit/watcher_test.c @@ -2,7 +2,7 @@ #include #include /* FileWatcherStateNew() */ -#include /* Bundle */ +#include /* Rval */ #include /* LoggingPrivContext, LoggingPrivSetContext() */ #include /* strstr() */ @@ -30,12 +30,11 @@ static void test_watcher_register_single(void) { WatcherRegistryInitialize(); - /* WatcherRegister()/WatcherDestroy() never dereference the bundle - * pointer, they only store it in a map, so a fake non-NULL pointer is - * fine here. */ - Bundle *fake_bundle = (Bundle *) 0x1; + /* WatcherRegister() only deep-copies the rval, it never resolves the + * bundle it names, so the bundle doesn't need to exist. */ + const Rval bundle = { .item = (char *) "test_bundle", .type = RVAL_TYPE_SCALAR }; - WatcherRegister("test-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/test-event"), fake_bundle, 5); + assert_true(WatcherRegister("test-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/test-event"), bundle, 5)); WatcherRegistryFinalize(); } @@ -44,10 +43,10 @@ static void test_watcher_register_duplicate_key_ignored(void) { WatcherRegistryInitialize(); - Bundle *fake_bundle_a = (Bundle *) 0x1; - Bundle *fake_bundle_b = (Bundle *) 0x2; + const Rval bundle_a = { .item = (char *) "bundle_a", .type = RVAL_TYPE_SCALAR }; + const Rval bundle_b = { .item = (char *) "bundle_b", .type = RVAL_TYPE_SCALAR }; - WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-a"), fake_bundle_a, 5); + assert_true(WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-a"), bundle_a, 5)); captured_err_count = 0; captured_err_message[0] = '\0'; @@ -64,11 +63,12 @@ static void test_watcher_register_duplicate_key_ignored(void) /* Registering the same key again must be rejected: an error is logged * (not silently swallowed) and the first registration is kept, not * replaced. */ - WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-b"), fake_bundle_b, 5); + const bool registered = WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-b"), bundle_b, 5); LogSetGlobalLevel(old_level); LoggingPrivSetContext(NULL); + assert_false(registered); assert_int_equal(captured_err_count, 1); assert_true(strstr(captured_err_message, "dup-event") != NULL); assert_true(strstr(captured_err_message, "already registered") != NULL);