From 433fff5ea2e255c114245ef1964e7e548692f112 Mon Sep 17 00:00:00 2001 From: Victor Moene Date: Wed, 30 Sep 2026 11:39:55 +0200 Subject: [PATCH 1/6] Moved ScheduleAgentOperation to agent_operation.c Signed-off-by: Victor Moene --- cf-agent/Makefile.am | 1 + cf-agent/agent_operations.c | 770 ++++++++++++++++++++++++++++++++++++ cf-agent/agent_operations.h | 38 ++ cf-agent/cf-agent.c | 707 +-------------------------------- cf-agent/verify_methods.c | 1 + libpromises/prototypes3.h | 6 - 6 files changed, 811 insertions(+), 712 deletions(-) create mode 100644 cf-agent/agent_operations.c create mode 100644 cf-agent/agent_operations.h diff --git a/cf-agent/Makefile.am b/cf-agent/Makefile.am index 8e1a9815fa1..5770526209d 100644 --- a/cf-agent/Makefile.am +++ b/cf-agent/Makefile.am @@ -73,6 +73,7 @@ libcf_agent_la_LIBADD = ../libpromises/libpromises.la \ libcf_agent_la_SOURCES = \ agent-diagnostics.c agent-diagnostics.h \ + agent_operations.c agent_operations.h \ simulate_mode.c simulate_mode.h \ tokyo_check.c tokyo_check.h \ abstract_dir.c abstract_dir.h \ diff --git a/cf-agent/agent_operations.c b/cf-agent/agent_operations.c new file mode 100644 index 00000000000..02033a6f6dd --- /dev/null +++ b/cf-agent/agent_operations.c @@ -0,0 +1,770 @@ +/* + 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. +*/ + +/* ScheduleAgentOperations() and the promise actuator it uses: this is what + * cf-agent does with a single bundle (evaluate its promises in + * AGENT_TYPESEQUENCE order, keeping each one with KeepAgentPromise()). It is + * kept separate from cf-agent.c so that other components (e.g. cf-reactor, + * running the bundle named in an events promise's "then") can run a bundle + * the same way cf-agent does, without linking cf-agent's main(). */ + +#include +#include + +#include +#include /* ExpandPromise() */ +#include /* SpecialTypeBanner() */ +#include +#include +#include +#include /* IsRegexItemIn() */ +#include /* FullTextMatch() */ +#include +#include +#include +#include /* IsBuiltInPromiseType() */ +#include /* EvaluateCustomPromise() */ +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +extern int PR_KEPT; +extern int PR_REPAIRED; +extern int PR_NOTKEPT; + +int CFA_BACKGROUND = 0; /* GLOBAL_X */ +int CFA_BACKGROUND_LIMIT = 1; /* GLOBAL_P */ + +Item *PROCESSREFRESH = NULL; /* GLOBAL_P */ + +static const char *const AGENT_TYPESEQUENCE[] = +{ + "meta", + "vars", + "defaults", + "classes", /* Maelstrom order 2 */ + "users", + "files", + "packages", + "guest_environments", + "methods", + "processes", + "services", + "commands", + "storage", + "databases", + "reports", + NULL +}; + +static PromiseResult KeepAgentPromise(EvalContext *ctx, const Promise *pp, void *param); +static void NewTypeContext(TypeSequence type); +static void DeleteTypeContext(EvalContext *ctx, TypeSequence type); +static PromiseResult ParallelFindAndVerifyFilesPromises(EvalContext *ctx, const Promise *pp); +static void BannerStatusEnd(PromiseResult status, const char *type, char *name); +static void BannerStatusBegin(const char *type, char *name); +static PromiseResult DefaultVarPromise(EvalContext *ctx, const Promise *pp); +static int NoteBundleCompliance(const Bundle *bundle, int save_pr_kept, int save_pr_repaired, int save_pr_notkept, struct timespec start); + +/** + @brief + Wrapper around DefaultVarPromise to silence cast-function-type compiler warning in ScheduleAgentOperations + */ +static PromiseResult DefaultVarPromiseWrapper(EvalContext *ctx, const Promise *pp, void *param) { + UNUSED(param); + return DefaultVarPromise(ctx, pp); +} + +PromiseResult ScheduleAgentOperations(EvalContext *ctx, const Bundle *bp) +// NB - this function can be called recursively through "methods" +{ + if (EvalContextIsClassicOrder(ctx, bp)) + { + return ScheduleAgentOperationsNormalOrder(ctx, bp); + } + return ScheduleAgentOperationsTopDownOrder(ctx, bp); +} + +PromiseResult ScheduleAgentOperationsNormalOrder(EvalContext *ctx, const Bundle *bp) +{ + assert(bp != NULL); + + int save_pr_kept = PR_KEPT; + int save_pr_repaired = PR_REPAIRED; + int save_pr_notkept = PR_NOTKEPT; + struct timespec start = BeginMeasure(); + + if (PROCESSREFRESH == NULL || (PROCESSREFRESH && IsRegexItemIn(ctx, PROCESSREFRESH, bp->name))) + { + ClearProcessTable(); + } + + PromiseResult result = PROMISE_RESULT_SKIPPED; + + for (int pass = 1; pass < CF_DONEPASSES; pass++) + { + // Evaluate built-in (non-custom) promise types, according to type sequence (normal order): + for (TypeSequence type = 0; AGENT_TYPESEQUENCE[type] != NULL; type++) + { + const BundleSection *sp = BundleGetSection((Bundle *)bp, AGENT_TYPESEQUENCE[type]); + + if (!sp || SeqLength(sp->promises) == 0) + { + continue; + } + + NewTypeContext(type); + + SpecialTypeBanner(type, pass); + EvalContextStackPushBundleSectionFrame(ctx, sp); + + for (size_t ppi = 0; ppi < SeqLength(sp->promises); ppi++) + { + Promise *pp = SeqAt(sp->promises, ppi); + + EvalContextSetPass(ctx, pass); + + PromiseResult promise_result = ExpandPromise(ctx, pp, KeepAgentPromise, NULL); + result = PromiseResultUpdate(result, promise_result); + + if (EvalAborted(ctx) || BundleAbort(ctx)) + { + DeleteTypeContext(ctx, type); + EvalContextStackPopFrame(ctx); + NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); + return result; + } + } + + DeleteTypeContext(ctx, type); + EvalContextStackPopFrame(ctx); + + if (type == TYPE_SEQUENCE_CONTEXTS) + { + BundleResolve(ctx, bp); + BundleResolvePromiseType(ctx, bp, "defaults", DefaultVarPromiseWrapper); + } + } + + // Custom promises are evaluated at the end of an evaluation pass: + const size_t sections = SeqLength(bp->custom_sections); + for (size_t i = 0; i < sections; ++i) + { + BundleSection *section = SeqAt(bp->custom_sections, i); + + EvalContextStackPushBundleSectionFrame(ctx, section); + + const size_t promises = SeqLength(section->promises); + for (size_t ppi = 0; ppi < promises; ppi++) + { + Promise *pp = SeqAt(section->promises, ppi); + + EvalContextSetPass(ctx, pass); + + PromiseResult promise_result = ExpandPromise(ctx, pp, KeepAgentPromise, NULL); + result = PromiseResultUpdate(result, promise_result); + + if (EvalAborted(ctx) || BundleAbort(ctx)) + { + EvalContextStackPopFrame(ctx); + NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); + return result; + } + } + EvalContextStackPopFrame(ctx); + } + } + + NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); + return result; +} + +PromiseResult ScheduleAgentOperationsTopDownOrder(EvalContext *ctx, const Bundle *bp) +{ + assert(bp != NULL); + + int save_pr_kept = PR_KEPT; + int save_pr_repaired = PR_REPAIRED; + int save_pr_notkept = PR_NOTKEPT; + struct timespec start = BeginMeasure(); + + if (PROCESSREFRESH == NULL || (PROCESSREFRESH && IsRegexItemIn(ctx, PROCESSREFRESH, bp->name))) + { + ClearProcessTable(); + } + + PromiseResult result = PROMISE_RESULT_SKIPPED; + for (int pass = 1; pass < CF_DONEPASSES; pass++) + { + const char *last_promise_type = ""; + for (size_t ppi = 0; ppi < SeqLength(bp->all_promises); ppi++) + { + EvalContextSetPass(ctx, pass); + Promise *pp = SeqAt(bp->all_promises, ppi); + BundleSection *parent_section = pp->parent_section; + + if (!StringEqual(last_promise_type, parent_section->promise_type)) + { + SpecialTypeBannerFromString(parent_section->promise_type, pass); + } + last_promise_type = parent_section->promise_type; + + EvalContextStackPushBundleSectionFrame(ctx, parent_section); + + PromiseResult promise_result = ExpandPromise(ctx, pp, KeepAgentPromise, NULL); + result = PromiseResultUpdate(result, promise_result); + if (EvalAborted(ctx) || BundleAbort(ctx)) + { + EvalContextStackPopFrame(ctx); + NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); + return result; + } + EvalContextStackPopFrame(ctx); + } + } + + NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); + return result; +} + +/*********************************************************************/ + +static PromiseResult DefaultVarPromise(EvalContext *ctx, const Promise *pp) +{ + assert(pp != NULL); + + char *regex = PromiseGetConstraintAsRval(pp, "if_match_regex", RVAL_TYPE_SCALAR); + bool okay = true; + + + DataType value_type = CF_DATA_TYPE_NONE; + const void *value = NULL; + { + VarRef *ref = VarRefParseFromScope(pp->promiser, "this"); + value = EvalContextVariableGetPlaintext(ctx, ref, &value_type); + VarRefDestroy(ref); + } + + switch (value_type) + { + case CF_DATA_TYPE_STRING: + case CF_DATA_TYPE_INT: + case CF_DATA_TYPE_REAL: + if (regex && !FullTextMatch(ctx, regex, value)) + { + return PROMISE_RESULT_NOOP; + } + + if (regex == NULL) + { + return PROMISE_RESULT_NOOP; + } + break; + + case CF_DATA_TYPE_STRING_LIST: + case CF_DATA_TYPE_INT_LIST: + case CF_DATA_TYPE_REAL_LIST: + if (regex) + { + for (const Rlist *rp = value; rp != NULL; rp = rp->next) + { + if (FullTextMatch(ctx, regex, RlistScalarValue(rp))) + { + okay = false; + break; + } + } + + if (okay) + { + return PROMISE_RESULT_NOOP; + } + } + break; + + default: + break; + } + + { + VarRef *ref = VarRefParseFromBundle(pp->promiser, PromiseGetBundle(pp)); + EvalContextVariableRemove(ctx, ref); + VarRefDestroy(ref); + } + + return VerifyVarPromise(ctx, pp, NULL); +} + +static void LogVariableValue(const EvalContext *ctx, const Promise *pp) +{ + assert(pp != NULL); + + VarRef *ref = VarRefParseFromBundle(pp->promiser, PromiseGetBundle(pp)); + char *out = NULL; + + DataType type; + const void *var = EvalContextVariableGetPlaintext(ctx, ref, &type); + switch (type) + { + case CF_DATA_TYPE_INT: + case CF_DATA_TYPE_REAL: + case CF_DATA_TYPE_STRING: + out = xstrdup((char *) var); + break; + case CF_DATA_TYPE_INT_LIST: + case CF_DATA_TYPE_REAL_LIST: + case CF_DATA_TYPE_STRING_LIST: + { + size_t siz = CF_BUFSIZE; + size_t len = 0; + out = xcalloc(1, CF_BUFSIZE); + + for (Rlist *rp = (Rlist *) var; rp != NULL; rp = rp->next) + { + const char *s = (char *) rp->val.item; + + if (strlen(s) + len + 3 >= siz) // ", " + NULL + { + out = xrealloc(out, siz + CF_BUFSIZE); + siz += CF_BUFSIZE; + } + + if (len > 0) + { + len += strlcat(out, ", ", siz); + } + + len += strlcat(out, s, siz); + } + break; + } + case CF_DATA_TYPE_CONTAINER: + { + Writer *w = StringWriter(); + JsonWriteCompact(w, (JsonElement *) var); + out = StringWriterClose(w); + break; + } + default: + /* TODO is CF_DATA_TYPE_NONE acceptable? Today all meta variables + * are of this type. */ + /* UnexpectedError("Variable '%s' is of unknown type %d", */ + /* pp->promiser, type); */ + out = xstrdup("NONE"); + break; + } + + Log(LOG_LEVEL_DEBUG, "V: '%s' => '%s'", pp->promiser, out); + free(out); + VarRefDestroy(ref); +} + +static PromiseResult KeepAgentPromise(EvalContext *ctx, const Promise *pp, ARG_UNUSED void *param) +{ + assert(param == NULL); + assert(pp != NULL); + + BannerStatusBegin(PromiseGetPromiseType(pp), pp->promiser); + struct timespec start = BeginMeasure(); + PromiseResult result = PROMISE_RESULT_NOOP; + + if (strcmp("meta", PromiseGetPromiseType(pp)) == 0 || + strcmp("vars", PromiseGetPromiseType(pp)) == 0) + { + Log(LOG_LEVEL_VERBOSE, "V: Computing value of '%s'", pp->promiser); + + result = VerifyVarPromise(ctx, pp, NULL); + if (result != PROMISE_RESULT_FAIL) + { + if (LogGetGlobalLevel() >= LOG_LEVEL_DEBUG) + { + LogVariableValue(ctx, pp); + } + } + } + else if (strcmp("defaults", PromiseGetPromiseType(pp)) == 0) + { + result = DefaultVarPromise(ctx, pp); + } + else if (strcmp("classes", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyClassPromise(ctx, pp, NULL); + } + else if (strcmp("processes", PromiseGetPromiseType(pp)) == 0) + { + if (!LoadProcessTable()) + { + Log(LOG_LEVEL_ERR, "Unable to read the process table - cannot keep processes: type promises"); + return PROMISE_RESULT_FAIL; + } + result = VerifyProcessesPromise(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("storage", PromiseGetPromiseType(pp)) == 0) + { + result = FindAndVerifyStoragePromises(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("packages", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyPackagesPromise(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("users", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyUsersPromise(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + + else if (strcmp("files", PromiseGetPromiseType(pp)) == 0) + { + result = ParallelFindAndVerifyFilesPromises(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("commands", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyExecPromise(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("databases", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyDatabasePromises(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("methods", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyMethodsPromise(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("services", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyServicesPromise(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("guest_environments", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyEnvironmentsPromise(ctx, pp); + if (result != PROMISE_RESULT_SKIPPED) + { + EndMeasurePromise(start, pp); + } + } + else if (strcmp("reports", PromiseGetPromiseType(pp)) == 0) + { + result = VerifyReportPromise(ctx, pp); + } + else if (!IsBuiltInPromiseType(PromiseGetPromiseType(pp))) + { + result = EvaluateCustomPromise(ctx, pp); + } + else + { + result = PROMISE_RESULT_NOOP; + } + + BannerStatusEnd(result, PromiseGetPromiseType(pp), pp->promiser); + EvalContextLogPromiseIterationOutcome(ctx, pp, result); + return result; +} + +static void BannerStatusBegin(const char *type, char *name) +{ + if (StringEqual(type, "vars") || StringEqual(type, "classes")) + { + return; + } + Log(LOG_LEVEL_VERBOSE, "P: BEGIN %s promise (%.30s%s)", + type, name, + (strlen(name) > 30) ? "..." : ""); +} + +static void BannerStatusEnd(PromiseResult status, const char *type, char *name) +{ + if ((strcmp(type, "vars") == 0) || (strcmp(type, "classes") == 0)) + { + return; + } + + switch (status) + { + case PROMISE_RESULT_CHANGE: + Log(LOG_LEVEL_VERBOSE, "A: Promise REPAIRED"); + break; + + case PROMISE_RESULT_TIMEOUT: + Log(LOG_LEVEL_VERBOSE, "A: Promise TIMED-OUT"); + break; + + case PROMISE_RESULT_WARN: + case PROMISE_RESULT_FAIL: + case PROMISE_RESULT_INTERRUPTED: + Log(LOG_LEVEL_VERBOSE, "A: Promise NOT KEPT!"); + break; + + case PROMISE_RESULT_DENIED: + Log(LOG_LEVEL_VERBOSE, "A: Promise NOT KEPT - denied"); + break; + + case PROMISE_RESULT_NOOP: + Log(LOG_LEVEL_VERBOSE, "A: Promise was KEPT"); + break; + default: + return; + break; + } + + Log(LOG_LEVEL_VERBOSE, "P: END %s promise (%.30s%s)", + type, name, + (strlen(name) > 30) ? "..." : ""); +} + +/*********************************************************************/ +/* Type context */ +/*********************************************************************/ + +static void NewTypeContext(TypeSequence type) +{ +// get maxconnections + + switch (type) + { + case TYPE_SEQUENCE_ENVIRONMENTS: + NewEnvironmentsContext(); + break; + + case TYPE_SEQUENCE_FILES: + break; + + case TYPE_SEQUENCE_PROCESSES: + break; + + case TYPE_SEQUENCE_STORAGE: +#ifndef __MINGW32__ // TODO: Run if implemented on Windows + if (SeqLength(GetGlobalMountedFSList())) + { + DeleteMountInfo(GetGlobalMountedFSList()); + SeqClear(GetGlobalMountedFSList()); + } +#endif /* !__MINGW32__ */ + break; + + default: + break; + } + + return; +} + +/*********************************************************************/ + +static void DeleteTypeContext(EvalContext *ctx, TypeSequence type) +{ + switch (type) + { + case TYPE_SEQUENCE_ENVIRONMENTS: + DeleteEnvironmentsContext(); + break; + + case TYPE_SEQUENCE_FILES: + break; + + case TYPE_SEQUENCE_PROCESSES: + break; + + case TYPE_SEQUENCE_STORAGE: + DeleteStorageContext(); + break; + + case TYPE_SEQUENCE_PACKAGES: + ExecuteScheduledPackages(ctx); + CleanScheduledPackages(); + break; + + default: + break; + } +} + +/**************************************************************/ +/* Thread context */ +/**************************************************************/ + +#ifdef __MINGW32__ + +static PromiseResult ParallelFindAndVerifyFilesPromises(EvalContext *ctx, const Promise *pp) +{ + int background = PromiseGetConstraintAsBoolean(ctx, "background", pp); + + if (background) + { + Log(LOG_LEVEL_VERBOSE, "Background processing of files promises is not supported on Windows"); + } + + return FindAndVerifyFilesPromises(ctx, pp); +} + +#else /* !__MINGW32__ */ + +static PromiseResult ParallelFindAndVerifyFilesPromises(EvalContext *ctx, const Promise *pp) +{ + int background = PromiseGetConstraintAsBoolean(ctx, "background", pp); + pid_t child = 1; + PromiseResult result = PROMISE_RESULT_SKIPPED; + + if (background) + { + if (CFA_BACKGROUND < CFA_BACKGROUND_LIMIT) + { + CFA_BACKGROUND++; + Log(LOG_LEVEL_VERBOSE, "Spawning new process..."); + child = fork(); + + if (child == 0) + { + ALARM_PID = -1; + + result = PromiseResultUpdate(result, FindAndVerifyFilesPromises(ctx, pp)); + + Log(LOG_LEVEL_VERBOSE, "Exiting backgrounded promise"); + PromiseRef(LOG_LEVEL_VERBOSE, pp); + _exit(EXIT_SUCCESS); + // TODO: need to solve this + } + } + else + { + Log(LOG_LEVEL_VERBOSE, "Promised parallel execution promised but exceeded the max number of promised background tasks, so serializing"); + background = 0; + } + } + else + { + result = PromiseResultUpdate(result, FindAndVerifyFilesPromises(ctx, pp)); + } + + return result; +} + +#endif /* !__MINGW32__ */ + +/*********************************************************************/ +/* Compliance comp */ +/*********************************************************************/ + +static int NoteBundleCompliance(const Bundle *bundle, int save_pr_kept, int save_pr_repaired, int save_pr_notkept, struct timespec start) +{ + assert(bundle != NULL); + + double delta_pr_kept, delta_pr_repaired, delta_pr_notkept; + double bundle_compliance = 0.0; + + delta_pr_kept = (double) (PR_KEPT - save_pr_kept); + delta_pr_notkept = (double) (PR_NOTKEPT - save_pr_notkept); + delta_pr_repaired = (double) (PR_REPAIRED - save_pr_repaired); + + Log(LOG_LEVEL_VERBOSE, "A: ..................................................."); + Log(LOG_LEVEL_VERBOSE, "A: Bundle Accounting Summary for '%s' in namespace %s", bundle->name, bundle->ns); + + if (delta_pr_kept + delta_pr_notkept + delta_pr_repaired <= 0) + { + Log(LOG_LEVEL_VERBOSE, "A: Zero promises executed for bundle '%s'", bundle->name); + Log(LOG_LEVEL_VERBOSE, "A: ..................................................."); + return PROMISE_RESULT_NOOP; + } + else + { + Log(LOG_LEVEL_VERBOSE, "A: Promises kept in '%s' = %.0lf", bundle->name, delta_pr_kept); + Log(LOG_LEVEL_VERBOSE, "A: Promises not kept in '%s' = %.0lf", bundle->name, delta_pr_notkept); + Log(LOG_LEVEL_VERBOSE, "A: Promises repaired in '%s' = %.0lf", bundle->name, delta_pr_repaired); + + bundle_compliance = (delta_pr_kept + delta_pr_repaired) / (delta_pr_kept + delta_pr_notkept + delta_pr_repaired); + + Log(LOG_LEVEL_VERBOSE, "A: Aggregate compliance (promises kept/repaired) for bundle '%s' = %.1lf%%", + bundle->name, bundle_compliance * 100.0); + + if (LogGetGlobalLevel() >= LOG_LEVEL_INFO) + { + char name[CF_MAXVARSIZE]; + snprintf(name, CF_MAXVARSIZE, "%s:%s", bundle->ns, bundle->name); + EndMeasure(name, start); + } + else + { + EndMeasure(NULL, start); + } + Log(LOG_LEVEL_VERBOSE, "A: ..................................................."); + } + + // return the worst case for the bundle status + + if (delta_pr_notkept > 0) + { + return PROMISE_RESULT_FAIL; + } + + if (delta_pr_repaired > 0) + { + return PROMISE_RESULT_CHANGE; + } + + return PROMISE_RESULT_NOOP; +} diff --git a/cf-agent/agent_operations.h b/cf-agent/agent_operations.h new file mode 100644 index 00000000000..b569e1f89c9 --- /dev/null +++ b/cf-agent/agent_operations.h @@ -0,0 +1,38 @@ +/* + 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_AGENT_OPERATIONS_H +#define CFENGINE_AGENT_OPERATIONS_H + +#include + +PromiseResult ScheduleAgentOperations(EvalContext *ctx, const Bundle *bp); +PromiseResult ScheduleAgentOperationsNormalOrder(EvalContext *ctx, const Bundle *bp); +PromiseResult ScheduleAgentOperationsTopDownOrder(EvalContext *ctx, const Bundle *bp); + +extern int CFA_BACKGROUND; /* GLOBAL_X */ +extern int CFA_BACKGROUND_LIMIT; /* GLOBAL_P, body agent control: max_children */ +extern Item *PROCESSREFRESH; /* GLOBAL_P, body agent control: refresh_processes */ + +#endif diff --git a/cf-agent/cf-agent.c b/cf-agent/cf-agent.c index 4ffe508e72f..4a3a9612f44 100644 --- a/cf-agent/cf-agent.c +++ b/cf-agent/cf-agent.c @@ -29,6 +29,7 @@ #include #include +#include #include #include #include @@ -104,10 +105,6 @@ #include -extern int PR_KEPT; -extern int PR_REPAIRED; -extern int PR_NOTKEPT; - static bool ALLCLASSESREPORT = false; /* GLOBAL_P */ static bool ALWAYS_VALIDATE = false; /* GLOBAL_P */ static bool CFPARANOID = false; /* GLOBAL_P */ @@ -115,31 +112,6 @@ static bool PERFORM_DB_CHECK = false; static const Rlist *ACCESSLIST = NULL; /* GLOBAL_P */ -static int CFA_BACKGROUND = 0; /* GLOBAL_X */ -static int CFA_BACKGROUND_LIMIT = 1; /* GLOBAL_P */ - -static Item *PROCESSREFRESH = NULL; /* GLOBAL_P */ - -static const char *const AGENT_TYPESEQUENCE[] = -{ - "meta", - "vars", - "defaults", - "classes", /* Maelstrom order 2 */ - "users", - "files", - "packages", - "guest_environments", - "methods", - "processes", - "services", - "commands", - "storage", - "databases", - "reports", - NULL -}; - /*******************************************************************/ /* Agent specific variables */ /*******************************************************************/ @@ -151,20 +123,12 @@ static char **TranslateOldBootstrapOptionsConcatenated(int argc, char **argv); static void FreeFixedStringArray(int size, char **array); static void CheckAgentAccess(const Rlist *list, const Policy *policy); static void KeepControlPromises(EvalContext *ctx, const Policy *policy, GenericAgentConfig *config); -static PromiseResult KeepAgentPromise(EvalContext *ctx, const Promise *pp, void *param); -static void NewTypeContext(TypeSequence type); -static void DeleteTypeContext(EvalContext *ctx, TypeSequence type); -static PromiseResult ParallelFindAndVerifyFilesPromises(EvalContext *ctx, const Promise *pp); static bool VerifyBootstrap(bool skip_cf_execd_check); static void KeepPromiseBundles(EvalContext *ctx, const Policy *policy, GenericAgentConfig *config); static void KeepPromises(EvalContext *ctx, const Policy *policy, GenericAgentConfig *config); -static int NoteBundleCompliance(const Bundle *bundle, int save_pr_kept, int save_pr_repaired, int save_pr_notkept, struct timespec start); static void AllClassesReport(const EvalContext *ctx); static bool HasAvahiSupport(void); static int AutomaticBootstrap(GenericAgentConfig *config); -static void BannerStatusEnd(PromiseResult status, const char *type, char *name); -static void BannerStatusBegin(const char *type, char *name); -static PromiseResult DefaultVarPromise(EvalContext *ctx, const Promise *pp); static void WaitForBackgroundProcesses(); /*******************************************************************/ @@ -260,17 +224,6 @@ static const char *const HINTS[] = NULL }; -/** - @brief - Wrapper around DefaultVarPromise to silence cast-function-type compiler warning in ScheduleAgentOperations - */ -static PromiseResult DefaultVarPromiseWrapper(EvalContext *ctx, const Promise *pp, void *param) { - UNUSED(param); - return DefaultVarPromise(ctx, pp); -} - -/*******************************************************************/ - int main(int argc, char *argv[]) { SetupSignalsForAgent(); @@ -1599,160 +1552,6 @@ static void AllClassesReport(const EvalContext *ctx) } } -PromiseResult ScheduleAgentOperations(EvalContext *ctx, const Bundle *bp) -// NB - this function can be called recursively through "methods" -{ - if (EvalContextIsClassicOrder(ctx, bp)) - { - return ScheduleAgentOperationsNormalOrder(ctx, bp); - } - return ScheduleAgentOperationsTopDownOrder(ctx, bp); -} - -PromiseResult ScheduleAgentOperationsNormalOrder(EvalContext *ctx, const Bundle *bp) -{ - assert(bp != NULL); - - int save_pr_kept = PR_KEPT; - int save_pr_repaired = PR_REPAIRED; - int save_pr_notkept = PR_NOTKEPT; - struct timespec start = BeginMeasure(); - - if (PROCESSREFRESH == NULL || (PROCESSREFRESH && IsRegexItemIn(ctx, PROCESSREFRESH, bp->name))) - { - ClearProcessTable(); - } - - PromiseResult result = PROMISE_RESULT_SKIPPED; - - for (int pass = 1; pass < CF_DONEPASSES; pass++) - { - // Evaluate built-in (non-custom) promise types, according to type sequence (normal order): - for (TypeSequence type = 0; AGENT_TYPESEQUENCE[type] != NULL; type++) - { - const BundleSection *sp = BundleGetSection((Bundle *)bp, AGENT_TYPESEQUENCE[type]); - - if (!sp || SeqLength(sp->promises) == 0) - { - continue; - } - - NewTypeContext(type); - - SpecialTypeBanner(type, pass); - EvalContextStackPushBundleSectionFrame(ctx, sp); - - for (size_t ppi = 0; ppi < SeqLength(sp->promises); ppi++) - { - Promise *pp = SeqAt(sp->promises, ppi); - - EvalContextSetPass(ctx, pass); - - PromiseResult promise_result = ExpandPromise(ctx, pp, KeepAgentPromise, NULL); - result = PromiseResultUpdate(result, promise_result); - - if (EvalAborted(ctx) || BundleAbort(ctx)) - { - DeleteTypeContext(ctx, type); - EvalContextStackPopFrame(ctx); - NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); - return result; - } - } - - DeleteTypeContext(ctx, type); - EvalContextStackPopFrame(ctx); - - if (type == TYPE_SEQUENCE_CONTEXTS) - { - BundleResolve(ctx, bp); - BundleResolvePromiseType(ctx, bp, "defaults", DefaultVarPromiseWrapper); - } - } - - // Custom promises are evaluated at the end of an evaluation pass: - const size_t sections = SeqLength(bp->custom_sections); - for (size_t i = 0; i < sections; ++i) - { - BundleSection *section = SeqAt(bp->custom_sections, i); - - EvalContextStackPushBundleSectionFrame(ctx, section); - - const size_t promises = SeqLength(section->promises); - for (size_t ppi = 0; ppi < promises; ppi++) - { - Promise *pp = SeqAt(section->promises, ppi); - - EvalContextSetPass(ctx, pass); - - PromiseResult promise_result = ExpandPromise(ctx, pp, KeepAgentPromise, NULL); - result = PromiseResultUpdate(result, promise_result); - - if (EvalAborted(ctx) || BundleAbort(ctx)) - { - EvalContextStackPopFrame(ctx); - NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); - return result; - } - } - EvalContextStackPopFrame(ctx); - } - } - - NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); - return result; -} - -PromiseResult ScheduleAgentOperationsTopDownOrder(EvalContext *ctx, const Bundle *bp) -{ - assert(bp != NULL); - - int save_pr_kept = PR_KEPT; - int save_pr_repaired = PR_REPAIRED; - int save_pr_notkept = PR_NOTKEPT; - struct timespec start = BeginMeasure(); - - if (PROCESSREFRESH == NULL || (PROCESSREFRESH && IsRegexItemIn(ctx, PROCESSREFRESH, bp->name))) - { - ClearProcessTable(); - } - - PromiseResult result = PROMISE_RESULT_SKIPPED; - for (int pass = 1; pass < CF_DONEPASSES; pass++) - { - const char *last_promise_type = ""; - for (size_t ppi = 0; ppi < SeqLength(bp->all_promises); ppi++) - { - EvalContextSetPass(ctx, pass); - Promise *pp = SeqAt(bp->all_promises, ppi); - BundleSection *parent_section = pp->parent_section; - - if (!StringEqual(last_promise_type, parent_section->promise_type)) - { - SpecialTypeBannerFromString(parent_section->promise_type, pass); - } - last_promise_type = parent_section->promise_type; - - EvalContextStackPushBundleSectionFrame(ctx, parent_section); - - PromiseResult promise_result = ExpandPromise(ctx, pp, KeepAgentPromise, NULL); - result = PromiseResultUpdate(result, promise_result); - if (EvalAborted(ctx) || BundleAbort(ctx)) - { - EvalContextStackPopFrame(ctx); - NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); - return result; - } - EvalContextStackPopFrame(ctx); - } - } - - NoteBundleCompliance(bp, save_pr_kept, save_pr_repaired, save_pr_notkept, start); - return result; -} - -/*********************************************************************/ - #ifdef __MINGW32__ static void CheckAgentAccess(const Rlist *list, const Policy *policy) @@ -1821,449 +1620,6 @@ static void CheckAgentAccess(const Rlist *list, const Policy *policy) /*********************************************************************/ -static PromiseResult DefaultVarPromise(EvalContext *ctx, const Promise *pp) -{ - char *regex = PromiseGetConstraintAsRval(pp, "if_match_regex", RVAL_TYPE_SCALAR); - bool okay = true; - - - DataType value_type = CF_DATA_TYPE_NONE; - const void *value = NULL; - { - VarRef *ref = VarRefParseFromScope(pp->promiser, "this"); - value = EvalContextVariableGetPlaintext(ctx, ref, &value_type); - VarRefDestroy(ref); - } - - switch (value_type) - { - case CF_DATA_TYPE_STRING: - case CF_DATA_TYPE_INT: - case CF_DATA_TYPE_REAL: - if (regex && !FullTextMatch(ctx, regex, value)) - { - return PROMISE_RESULT_NOOP; - } - - if (regex == NULL) - { - return PROMISE_RESULT_NOOP; - } - break; - - case CF_DATA_TYPE_STRING_LIST: - case CF_DATA_TYPE_INT_LIST: - case CF_DATA_TYPE_REAL_LIST: - if (regex) - { - for (const Rlist *rp = value; rp != NULL; rp = rp->next) - { - if (FullTextMatch(ctx, regex, RlistScalarValue(rp))) - { - okay = false; - break; - } - } - - if (okay) - { - return PROMISE_RESULT_NOOP; - } - } - break; - - default: - break; - } - - { - VarRef *ref = VarRefParseFromBundle(pp->promiser, PromiseGetBundle(pp)); - EvalContextVariableRemove(ctx, ref); - VarRefDestroy(ref); - } - - return VerifyVarPromise(ctx, pp, NULL); -} - -static void LogVariableValue(const EvalContext *ctx, const Promise *pp) -{ - VarRef *ref = VarRefParseFromBundle(pp->promiser, PromiseGetBundle(pp)); - char *out = NULL; - - DataType type; - const void *var = EvalContextVariableGetPlaintext(ctx, ref, &type); - switch (type) - { - case CF_DATA_TYPE_INT: - case CF_DATA_TYPE_REAL: - case CF_DATA_TYPE_STRING: - out = xstrdup((char *) var); - break; - case CF_DATA_TYPE_INT_LIST: - case CF_DATA_TYPE_REAL_LIST: - case CF_DATA_TYPE_STRING_LIST: - { - size_t siz = CF_BUFSIZE; - size_t len = 0; - out = xcalloc(1, CF_BUFSIZE); - - for (Rlist *rp = (Rlist *) var; rp != NULL; rp = rp->next) - { - const char *s = (char *) rp->val.item; - - if (strlen(s) + len + 3 >= siz) // ", " + NULL - { - out = xrealloc(out, siz + CF_BUFSIZE); - siz += CF_BUFSIZE; - } - - if (len > 0) - { - len += strlcat(out, ", ", siz); - } - - len += strlcat(out, s, siz); - } - break; - } - case CF_DATA_TYPE_CONTAINER: - { - Writer *w = StringWriter(); - JsonWriteCompact(w, (JsonElement *) var); - out = StringWriterClose(w); - break; - } - default: - /* TODO is CF_DATA_TYPE_NONE acceptable? Today all meta variables - * are of this type. */ - /* UnexpectedError("Variable '%s' is of unknown type %d", */ - /* pp->promiser, type); */ - out = xstrdup("NONE"); - break; - } - - Log(LOG_LEVEL_DEBUG, "V: '%s' => '%s'", pp->promiser, out); - free(out); - VarRefDestroy(ref); -} - -static PromiseResult KeepAgentPromise(EvalContext *ctx, const Promise *pp, ARG_UNUSED void *param) -{ - assert(param == NULL); - assert(pp != NULL); - - BannerStatusBegin(PromiseGetPromiseType(pp), pp->promiser); - struct timespec start = BeginMeasure(); - PromiseResult result = PROMISE_RESULT_NOOP; - - if (strcmp("meta", PromiseGetPromiseType(pp)) == 0 || - strcmp("vars", PromiseGetPromiseType(pp)) == 0) - { - Log(LOG_LEVEL_VERBOSE, "V: Computing value of '%s'", pp->promiser); - - result = VerifyVarPromise(ctx, pp, NULL); - if (result != PROMISE_RESULT_FAIL) - { - if (LogGetGlobalLevel() >= LOG_LEVEL_DEBUG) - { - LogVariableValue(ctx, pp); - } - } - } - else if (strcmp("defaults", PromiseGetPromiseType(pp)) == 0) - { - result = DefaultVarPromise(ctx, pp); - } - else if (strcmp("classes", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyClassPromise(ctx, pp, NULL); - } - else if (strcmp("processes", PromiseGetPromiseType(pp)) == 0) - { - if (!LoadProcessTable()) - { - Log(LOG_LEVEL_ERR, "Unable to read the process table - cannot keep processes: type promises"); - return PROMISE_RESULT_FAIL; - } - result = VerifyProcessesPromise(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("storage", PromiseGetPromiseType(pp)) == 0) - { - result = FindAndVerifyStoragePromises(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("packages", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyPackagesPromise(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("users", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyUsersPromise(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - - else if (strcmp("files", PromiseGetPromiseType(pp)) == 0) - { - result = ParallelFindAndVerifyFilesPromises(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("commands", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyExecPromise(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("databases", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyDatabasePromises(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("methods", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyMethodsPromise(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("services", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyServicesPromise(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("guest_environments", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyEnvironmentsPromise(ctx, pp); - if (result != PROMISE_RESULT_SKIPPED) - { - EndMeasurePromise(start, pp); - } - } - else if (strcmp("reports", PromiseGetPromiseType(pp)) == 0) - { - result = VerifyReportPromise(ctx, pp); - } - else if (!IsBuiltInPromiseType(PromiseGetPromiseType(pp))) - { - result = EvaluateCustomPromise(ctx, pp); - } - else - { - result = PROMISE_RESULT_NOOP; - } - - BannerStatusEnd(result, PromiseGetPromiseType(pp), pp->promiser); - EvalContextLogPromiseIterationOutcome(ctx, pp, result); - return result; -} - -static void BannerStatusBegin(const char *type, char *name) -{ - if (StringEqual(type, "vars") || StringEqual(type, "classes")) - { - return; - } - Log(LOG_LEVEL_VERBOSE, "P: BEGIN %s promise (%.30s%s)", - type, name, - (strlen(name) > 30) ? "..." : ""); -} - -static void BannerStatusEnd(PromiseResult status, const char *type, char *name) -{ - if ((strcmp(type, "vars") == 0) || (strcmp(type, "classes") == 0)) - { - return; - } - - switch (status) - { - case PROMISE_RESULT_CHANGE: - Log(LOG_LEVEL_VERBOSE, "A: Promise REPAIRED"); - break; - - case PROMISE_RESULT_TIMEOUT: - Log(LOG_LEVEL_VERBOSE, "A: Promise TIMED-OUT"); - break; - - case PROMISE_RESULT_WARN: - case PROMISE_RESULT_FAIL: - case PROMISE_RESULT_INTERRUPTED: - Log(LOG_LEVEL_VERBOSE, "A: Promise NOT KEPT!"); - break; - - case PROMISE_RESULT_DENIED: - Log(LOG_LEVEL_VERBOSE, "A: Promise NOT KEPT - denied"); - break; - - case PROMISE_RESULT_NOOP: - Log(LOG_LEVEL_VERBOSE, "A: Promise was KEPT"); - break; - default: - return; - break; - } - - Log(LOG_LEVEL_VERBOSE, "P: END %s promise (%.30s%s)", - type, name, - (strlen(name) > 30) ? "..." : ""); -} - -/*********************************************************************/ -/* Type context */ -/*********************************************************************/ - -static void NewTypeContext(TypeSequence type) -{ -// get maxconnections - - switch (type) - { - case TYPE_SEQUENCE_ENVIRONMENTS: - NewEnvironmentsContext(); - break; - - case TYPE_SEQUENCE_FILES: - break; - - case TYPE_SEQUENCE_PROCESSES: - break; - - case TYPE_SEQUENCE_STORAGE: -#ifndef __MINGW32__ // TODO: Run if implemented on Windows - if (SeqLength(GetGlobalMountedFSList())) - { - DeleteMountInfo(GetGlobalMountedFSList()); - SeqClear(GetGlobalMountedFSList()); - } -#endif /* !__MINGW32__ */ - break; - - default: - break; - } - - return; -} - -/*********************************************************************/ - -static void DeleteTypeContext(EvalContext *ctx, TypeSequence type) -{ - switch (type) - { - case TYPE_SEQUENCE_ENVIRONMENTS: - DeleteEnvironmentsContext(); - break; - - case TYPE_SEQUENCE_FILES: - break; - - case TYPE_SEQUENCE_PROCESSES: - break; - - case TYPE_SEQUENCE_STORAGE: - DeleteStorageContext(); - break; - - case TYPE_SEQUENCE_PACKAGES: - ExecuteScheduledPackages(ctx); - CleanScheduledPackages(); - break; - - default: - break; - } -} - -/**************************************************************/ -/* Thread context */ -/**************************************************************/ - -#ifdef __MINGW32__ - -static PromiseResult ParallelFindAndVerifyFilesPromises(EvalContext *ctx, const Promise *pp) -{ - int background = PromiseGetConstraintAsBoolean(ctx, "background", pp); - - if (background) - { - Log(LOG_LEVEL_VERBOSE, "Background processing of files promises is not supported on Windows"); - } - - return FindAndVerifyFilesPromises(ctx, pp); -} - -#else /* !__MINGW32__ */ - -static PromiseResult ParallelFindAndVerifyFilesPromises(EvalContext *ctx, const Promise *pp) -{ - int background = PromiseGetConstraintAsBoolean(ctx, "background", pp); - pid_t child = 1; - PromiseResult result = PROMISE_RESULT_SKIPPED; - - if (background) - { - if (CFA_BACKGROUND < CFA_BACKGROUND_LIMIT) - { - CFA_BACKGROUND++; - Log(LOG_LEVEL_VERBOSE, "Spawning new process..."); - child = fork(); - - if (child == 0) - { - ALARM_PID = -1; - - result = PromiseResultUpdate(result, FindAndVerifyFilesPromises(ctx, pp)); - - Log(LOG_LEVEL_VERBOSE, "Exiting backgrounded promise"); - PromiseRef(LOG_LEVEL_VERBOSE, pp); - _exit(EXIT_SUCCESS); - // TODO: need to solve this - } - } - else - { - Log(LOG_LEVEL_VERBOSE, "Promised parallel execution promised but exceeded the max number of promised background tasks, so serializing"); - background = 0; - } - } - else - { - result = PromiseResultUpdate(result, FindAndVerifyFilesPromises(ctx, pp)); - } - - return result; -} - -#endif /* !__MINGW32__ */ - -/**************************************************************/ - static bool VerifyBootstrap(bool skip_cf_execd_check) { const char *policy_server = PolicyServerGet(); @@ -2302,67 +1658,6 @@ static bool VerifyBootstrap(bool skip_cf_execd_check) return true; } -/**************************************************************/ -/* Compliance comp */ -/**************************************************************/ - -static int NoteBundleCompliance(const Bundle *bundle, int save_pr_kept, int save_pr_repaired, int save_pr_notkept, struct timespec start) -{ - double delta_pr_kept, delta_pr_repaired, delta_pr_notkept; - double bundle_compliance = 0.0; - - delta_pr_kept = (double) (PR_KEPT - save_pr_kept); - delta_pr_notkept = (double) (PR_NOTKEPT - save_pr_notkept); - delta_pr_repaired = (double) (PR_REPAIRED - save_pr_repaired); - - Log(LOG_LEVEL_VERBOSE, "A: ..................................................."); - Log(LOG_LEVEL_VERBOSE, "A: Bundle Accounting Summary for '%s' in namespace %s", bundle->name, bundle->ns); - - if (delta_pr_kept + delta_pr_notkept + delta_pr_repaired <= 0) - { - Log(LOG_LEVEL_VERBOSE, "A: Zero promises executed for bundle '%s'", bundle->name); - Log(LOG_LEVEL_VERBOSE, "A: ..................................................."); - return PROMISE_RESULT_NOOP; - } - else - { - Log(LOG_LEVEL_VERBOSE, "A: Promises kept in '%s' = %.0lf", bundle->name, delta_pr_kept); - Log(LOG_LEVEL_VERBOSE, "A: Promises not kept in '%s' = %.0lf", bundle->name, delta_pr_notkept); - Log(LOG_LEVEL_VERBOSE, "A: Promises repaired in '%s' = %.0lf", bundle->name, delta_pr_repaired); - - bundle_compliance = (delta_pr_kept + delta_pr_repaired) / (delta_pr_kept + delta_pr_notkept + delta_pr_repaired); - - Log(LOG_LEVEL_VERBOSE, "A: Aggregate compliance (promises kept/repaired) for bundle '%s' = %.1lf%%", - bundle->name, bundle_compliance * 100.0); - - if (LogGetGlobalLevel() >= LOG_LEVEL_INFO) - { - char name[CF_MAXVARSIZE]; - snprintf(name, CF_MAXVARSIZE, "%s:%s", bundle->ns, bundle->name); - EndMeasure(name, start); - } - else - { - EndMeasure(NULL, start); - } - Log(LOG_LEVEL_VERBOSE, "A: ..................................................."); - } - - // return the worst case for the bundle status - - if (delta_pr_notkept > 0) - { - return PROMISE_RESULT_FAIL; - } - - if (delta_pr_repaired > 0) - { - return PROMISE_RESULT_CHANGE; - } - - return PROMISE_RESULT_NOOP; -} - #if defined(HAVE_AVAHI_CLIENT_CLIENT_H) && defined(HAVE_AVAHI_COMMON_ADDRESS_H) static bool HasAvahiSupport(void) diff --git a/cf-agent/verify_methods.c b/cf-agent/verify_methods.c index b19d5952c9c..3669dea3e34 100644 --- a/cf-agent/verify_methods.c +++ b/cf-agent/verify_methods.c @@ -39,6 +39,7 @@ #include #include #include +#include // ScheduleAgentOperations() static void GetReturnValue(EvalContext *ctx, const Bundle *callee, const Promise *caller); diff --git a/libpromises/prototypes3.h b/libpromises/prototypes3.h index 2b963ed2c0a..061e7afc40d 100644 --- a/libpromises/prototypes3.h +++ b/libpromises/prototypes3.h @@ -41,12 +41,6 @@ const char *NameVersion(void); void yyerror(const char *s); -/* agent.c */ - -PromiseResult ScheduleAgentOperations(EvalContext *ctx, const Bundle *bp); -PromiseResult ScheduleAgentOperationsNormalOrder(EvalContext *ctx, const Bundle *bp); -PromiseResult ScheduleAgentOperationsTopDownOrder(EvalContext *ctx, const Bundle *bp); - /* Only for agent.c */ void ConnectionsInit(void); From 87a44070435d5ecea694c549df16abc4fd97d75d Mon Sep 17 00:00:00 2001 From: Victor Moene Date: Tue, 6 Oct 2026 18:50:01 +0200 Subject: [PATCH 2/6] Added TODOs and comments after moving eval from agent to agent_operations.c Linked TODOs to new tickets to fix old bugs/TODOs that were already in cf-agent before moving the code. Signed-off-by: Victor Moene --- cf-agent/agent_operations.c | 71 +++++++++++++++++++++++++++++++------ cf-agent/agent_operations.h | 29 +++++++++++++++ cf-agent/cf-agent.c | 2 +- 3 files changed, 90 insertions(+), 12 deletions(-) diff --git a/cf-agent/agent_operations.c b/cf-agent/agent_operations.c index 02033a6f6dd..58a9a094db9 100644 --- a/cf-agent/agent_operations.c +++ b/cf-agent/agent_operations.c @@ -22,12 +22,18 @@ included file COSL.txt. */ -/* ScheduleAgentOperations() and the promise actuator it uses: this is what - * cf-agent does with a single bundle (evaluate its promises in - * AGENT_TYPESEQUENCE order, keeping each one with KeepAgentPromise()). It is - * kept separate from cf-agent.c so that other components (e.g. cf-reactor, - * running the bundle named in an events promise's "then") can run a bundle - * the same way cf-agent does, without linking cf-agent's main(). */ +/* This file, agent_operations.c, implements ScheduleAgentOperations(), the + * function cf-agent uses to evaluate a single bundle and its promises, either + * in AGENT_TYPESEQUENCE order (normal order) or in the order they are written + * (top-down order). Each promise is handed to KeepAgentPromise(), which calls + * the Verify*() function for its promise type. + * + * It is kept separate from cf-agent.c so that other components can evaluate a + * bundle the same way cf-agent does, without linking cf-agent's main(), e.g. + * cf-reactor running the bundle named in an events promise's 'then'. + * + * TODO: ENT-14705 move this code (and the Verify*() functions it calls) from + * cf-agent to libpromises. */ #include #include @@ -64,9 +70,27 @@ extern int PR_KEPT; extern int PR_REPAIRED; extern int PR_NOTKEPT; +/* See the GLOBALS file for the GLOBAL_* classification. CFA_BACKGROUND_LIMIT + * and PROCESSREFRESH are only set by cf-agent (KeepControlPromises() in + * cf-agent.c), other components using this file keep the defaults below. */ + +/* Number of files promises run in the background (forked) by + * ParallelFindAndVerifyFilesPromises(), also passed to EndAudit() by cf-agent. + * It is never decremented, so CFA_BACKGROUND_LIMIT limits the total number of + * backgrounded promises in the lifetime of the process, not the number of + * concurrent ones. GLOBAL_X: this is a counter changed during evaluation, not + * configuration. It is OK for cf-agent, which evaluates its policy once and + * exits, but in a long-running process (cf-reactor) it is never reset, so + * once the limit is reached all later background promises are serialized. */ int CFA_BACKGROUND = 0; /* GLOBAL_X */ + +/* Max number of backgrounded files promises. + * GLOBAL_P, body agent control: max_children */ int CFA_BACKGROUND_LIMIT = 1; /* GLOBAL_P */ +/* Regexes matching the bundles before which the process table is re-read, or + * NULL to re-read it before every bundle. + * GLOBAL_P, body agent control: refresh_processes */ Item *PROCESSREFRESH = NULL; /* GLOBAL_P */ static const char *const AGENT_TYPESEQUENCE[] = @@ -108,7 +132,7 @@ static PromiseResult DefaultVarPromiseWrapper(EvalContext *ctx, const Promise *p } PromiseResult ScheduleAgentOperations(EvalContext *ctx, const Bundle *bp) -// NB - this function can be called recursively through "methods" +// NOTE: this function can be called recursively through "methods" { if (EvalContextIsClassicOrder(ctx, bp)) { @@ -126,6 +150,11 @@ PromiseResult ScheduleAgentOperationsNormalOrder(EvalContext *ctx, const Bundle int save_pr_notkept = PR_NOTKEPT; struct timespec start = BeginMeasure(); + /* Clear the cached process table so that processes promises re-read it, + * unless refresh_processes is set and does not match this bundle. Note + * that cf-agent only sets PROCESSREFRESH in verbose mode (see the TODO in + * KeepControlPromises() in cf-agent.c), and other components never set + * it, so in practice the table is cleared before every bundle. */ if (PROCESSREFRESH == NULL || (PROCESSREFRESH && IsRegexItemIn(ctx, PROCESSREFRESH, bp->name))) { ClearProcessTable(); @@ -178,7 +207,10 @@ PromiseResult ScheduleAgentOperationsNormalOrder(EvalContext *ctx, const Bundle } } - // Custom promises are evaluated at the end of an evaluation pass: + // Custom promises are evaluated at the end of an evaluation pass. + // NewTypeContext()/DeleteTypeContext() are not called for them, they + // only do something for built-in promise types (guest_environments, + // storage and packages). const size_t sections = SeqLength(bp->custom_sections); for (size_t i = 0; i < sections; ++i) { @@ -220,6 +252,11 @@ PromiseResult ScheduleAgentOperationsTopDownOrder(EvalContext *ctx, const Bundle int save_pr_notkept = PR_NOTKEPT; struct timespec start = BeginMeasure(); + /* Clear the cached process table so that processes promises re-read it, + * unless refresh_processes is set and does not match this bundle. Note + * that cf-agent only sets PROCESSREFRESH in verbose mode (see the TODO in + * KeepControlPromises() in cf-agent.c), and other components never set + * it, so in practice the table is cleared before every bundle. */ if (PROCESSREFRESH == NULL || (PROCESSREFRESH && IsRegexItemIn(ctx, PROCESSREFRESH, bp->name))) { ClearProcessTable(); @@ -228,6 +265,11 @@ PromiseResult ScheduleAgentOperationsTopDownOrder(EvalContext *ctx, const Bundle PromiseResult result = PROMISE_RESULT_SKIPPED; for (int pass = 1; pass < CF_DONEPASSES; pass++) { + // TODO: CFE-ENT-14706 NewTypeContext()/DeleteTypeContext() are not called + // in top-down order. In particular, ExecuteScheduledPackages() is + // never called here, so packages promises using package_method seem + // to be scheduled but never executed. Either call them for each block + // of promises of the same type, or remove the type contexts entirely. const char *last_promise_type = ""; for (size_t ppi = 0; ppi < SeqLength(bp->all_promises); ppi++) { @@ -378,8 +420,8 @@ static void LogVariableValue(const EvalContext *ctx, const Promise *pp) break; } default: - /* TODO is CF_DATA_TYPE_NONE acceptable? Today all meta variables - * are of this type. */ + /* TODO: ENT-14708 is CF_DATA_TYPE_NONE acceptable? Today all meta + * variables are of this type. */ /* UnexpectedError("Variable '%s' is of unknown type %d", */ /* pp->promiser, type); */ out = xstrdup("NONE"); @@ -635,6 +677,10 @@ static void DeleteTypeContext(EvalContext *ctx, TypeSequence type) break; case TYPE_SEQUENCE_PACKAGES: + /* There is no matching case in NewTypeContext(): packages promises + * using package_method are scheduled by VerifyPackagesPromise() and + * all executed here, at the end of the packages section. This does + * not happen in top-down order, see ScheduleAgentOperationsTopDownOrder(). */ ExecuteScheduledPackages(ctx); CleanScheduledPackages(); break; @@ -687,7 +733,10 @@ static PromiseResult ParallelFindAndVerifyFilesPromises(EvalContext *ctx, const Log(LOG_LEVEL_VERBOSE, "Exiting backgrounded promise"); PromiseRef(LOG_LEVEL_VERBOSE, pp); _exit(EXIT_SUCCESS); - // TODO: need to solve this + // TODO: ENT-14707 the result of the promise is lost: the child + // always exits with EXIT_SUCCESS, and the parent neither waits + // for it nor collects its result, it returns + // PROMISE_RESULT_SKIPPED for the backgrounded promise. } } else diff --git a/cf-agent/agent_operations.h b/cf-agent/agent_operations.h index b569e1f89c9..8604ba69513 100644 --- a/cf-agent/agent_operations.h +++ b/cf-agent/agent_operations.h @@ -27,8 +27,37 @@ #include +/** + * @brief Evaluates the promises of an agent or common bundle, like cf-agent + * does, using the evaluation order selected by + * EvalContextIsClassicOrder(). + * + * Used by cf-agent (bundlesequence and methods promises) and cf-reactor + * (bundle in the 'then' attribute of an events promise). Can be called + * recursively through methods promises. + * + * The caller is responsible for pushing the bundle frame + * (EvalContextStackPushBundleFrame()) before calling, and popping it after. + * + * @param ctx The evaluation context + * @param bp The bundle to evaluate + * @return The aggregated result of all promises evaluated in the bundle + */ PromiseResult ScheduleAgentOperations(EvalContext *ctx, const Bundle *bp); + +/** + * @brief Evaluates the promises of a bundle in normal order: by promise type, + * in AGENT_TYPESEQUENCE order, followed by custom promise types, in + * each evaluation pass. + * @see ScheduleAgentOperations() + */ PromiseResult ScheduleAgentOperationsNormalOrder(EvalContext *ctx, const Bundle *bp); + +/** + * @brief Evaluates the promises of a bundle in top-down order: in the order + * they are written in the policy, in each evaluation pass. + * @see ScheduleAgentOperations() + */ PromiseResult ScheduleAgentOperationsTopDownOrder(EvalContext *ctx, const Bundle *bp); extern int CFA_BACKGROUND; /* GLOBAL_X */ diff --git a/cf-agent/cf-agent.c b/cf-agent/cf-agent.c index 4a3a9612f44..a015224bffb 100644 --- a/cf-agent/cf-agent.c +++ b/cf-agent/cf-agent.c @@ -1030,7 +1030,7 @@ static void KeepControlPromises(EvalContext *ctx, const Policy *policy, GenericA for (const Rlist *rp = value; rp != NULL; rp = rp->next) { Log(LOG_LEVEL_VERBOSE, "%s", RlistScalarValue(rp)); - // TODO: why is this only done in verbose mode? + // TODO: ENT-14709 why is this only done in verbose mode? // original commit says 'optimization'. if (LogGetGlobalLevel() >= LOG_LEVEL_VERBOSE) { From 95fc5e60a6b698b4d064054d12a74afc3aadd00c Mon Sep 17 00:00:00 2001 From: Victor Moene Date: Wed, 30 Sep 2026 13:18:39 +0200 Subject: [PATCH 3/6] Run bundles on event Better concurrency for cf-reactor static state - No need for map anymore, since watcher keeps track of the promise - No need to keep key as variable anymore, since we do not use a map - No need to run bundle inside watcher.c - No risk of running a bundle from a wrong key between policy reads Signed-off-by: Victor Moene --- cf-reactor/Makefile.am | 3 +- cf-reactor/cf-reactor.c | 10 ++- cf-reactor/reactor_context.c | 5 +- cf-reactor/reactor_context.h | 3 +- cf-reactor/reactor_transform.c | 132 ++++++++++++++++++++++------- cf-reactor/reactor_transform.h | 17 +++- cf-reactor/watcher.c | 150 +++++++++++++++++++-------------- cf-reactor/watcher.h | 44 ++++++++-- 8 files changed, 258 insertions(+), 106 deletions(-) diff --git a/cf-reactor/Makefile.am b/cf-reactor/Makefile.am index 9d9a2e96aae..710befaef0c 100644 --- a/cf-reactor/Makefile.am +++ b/cf-reactor/Makefile.am @@ -26,6 +26,7 @@ noinst_LTLIBRARIES = libcf-reactor.la AM_CPPFLAGS = -I$(srcdir)/../libpromises -I$(srcdir)/../libntech/libutils \ -I$(srcdir)/../libcfecompat \ -I$(srcdir)/../libcfnet \ + -I$(srcdir)/../cf-agent \ $(OPENSSL_CPPFLAGS) \ $(PCRE2_CPPFLAGS) \ $(ENTERPRISE_CPPFLAGS) @@ -48,7 +49,7 @@ libcf_reactor_la_SOURCES = \ if !BUILTIN_EXTENSIONS bin_PROGRAMS = cf-reactor - cf_reactor_LDADD = libcf-reactor.la + cf_reactor_LDADD = libcf-reactor.la ../cf-agent/libcf-agent.la cf_reactor_SOURCES = endif diff --git a/cf-reactor/cf-reactor.c b/cf-reactor/cf-reactor.c index e6160e5d611..4a93b9d64cc 100644 --- a/cf-reactor/cf-reactor.c +++ b/cf-reactor/cf-reactor.c @@ -41,6 +41,7 @@ #include #include #include +#include /* WatcherRegistryClear */ /*****************************************************************************/ /* Globals */ @@ -237,6 +238,11 @@ static void CheckPolicyUpdates(EvalContext *ctx, Policy **policy, GenericAgentCo Log(LOG_LEVEL_NOTICE, "Rereading policy file '%s'", config->input_file); + /* Registered watchers refer to promises of the policy: handle the + * events detected so far with the current policy, and don't detect new + * ones until the watchers of the new policy are registered */ + EventWatcherPause(ctx, HandleReactorEvent); + WatcherRegistryClear(); EvalContextClear(ctx); PolicyDestroy(*policy); *policy = NULL; @@ -250,6 +256,8 @@ static void CheckPolicyUpdates(EvalContext *ctx, Policy **policy, GenericAgentCo { KeepReactorPromises(ctx, *policy); } + + EventWatcherResume(); } /*****************************************************************************/ @@ -377,7 +385,7 @@ int main(int argc, char *argv[]) } else { - ReactorContextHandleEvents(&reactor_ctx, &next_tick); + ReactorContextHandleEvents(ctx, &reactor_ctx, &next_tick); } diff --git a/cf-reactor/reactor_context.c b/cf-reactor/reactor_context.c index 907fe6d5974..e4d66f78ccf 100644 --- a/cf-reactor/reactor_context.c +++ b/cf-reactor/reactor_context.c @@ -26,6 +26,7 @@ #include /* ReactorNova*() */ #include /* GetSignalPipe() */ #include +#include /* HandleReactorEvent() */ #include #define INIT_FD_COUNT 8 @@ -186,7 +187,7 @@ static int GetWatcherFd(const ReactorContext *reactor_context) return 0; } -void ReactorContextHandleEvents(ReactorContext *reactor_context, time_t *next_tick) +void ReactorContextHandleEvents(EvalContext *ctx, ReactorContext *reactor_context, time_t *next_tick) { assert(reactor_context != NULL); @@ -214,7 +215,7 @@ void ReactorContextHandleEvents(ReactorContext *reactor_context, time_t *next_ti while (recv(GetSignalPipe(), &buf, 1, 0) > 0) { /* drain */ } } - EventWatcherHandleEvents(GetWatcherFd(reactor_context), &reactor_context->readfds); + EventWatcherHandleEvents(ctx, HandleReactorEvent, GetWatcherFd(reactor_context), &reactor_context->readfds); } void ReactorContextFinalize(ReactorContext *reactor_context) diff --git a/cf-reactor/reactor_context.h b/cf-reactor/reactor_context.h index 3950b94696a..fa8030239ff 100644 --- a/cf-reactor/reactor_context.h +++ b/cf-reactor/reactor_context.h @@ -26,6 +26,7 @@ #define CFENGINE_REACTOR_CONTEXT_H #include +#include #include typedef enum @@ -55,7 +56,7 @@ typedef struct bool ReactorContextInitialize(ReactorContext *reactor_context); int ReactorContextSetupFileDescriptors(ReactorContext *reactor_context); -void ReactorContextHandleEvents(ReactorContext *reactor_context, time_t *next_tick); +void ReactorContextHandleEvents(EvalContext *ctx, ReactorContext *reactor_context, time_t *next_tick); void ReactorContextFinalize(ReactorContext *reactor_context); #endif diff --git a/cf-reactor/reactor_transform.c b/cf-reactor/reactor_transform.c index 6c8e6c96022..4d8ed5825c9 100644 --- a/cf-reactor/reactor_transform.c +++ b/cf-reactor/reactor_transform.c @@ -34,6 +34,7 @@ #include #include #include +#include // ScheduleAgentOperations() /* Promise types evaluated within `bundle reactor NAME { ... }`. */ static const char *const REACTOR_TYPESEQUENCE[] = @@ -45,41 +46,51 @@ static const char *const REACTOR_TYPESEQUENCE[] = 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) +/* Resolves the agent or common bundle referred to by an events promise's + * 'then' attribute, either `then => "name"` or `then => name(args)`. If args is + * not NULL, it is set to the bundle arguments (or NULL if there are none). + * Logs an error and returns NULL if the bundle is not found. */ +static const Bundle *ResolveThenBundle( + const EvalContext *ctx, const Promise *pp, Rval then_rval, const Rlist **args) { - const char *bundle_name = NULL; + assert(pp != NULL); + + const char *name = NULL; + const Rlist *bundle_args = NULL; switch (then_rval.type) { case RVAL_TYPE_SCALAR: - bundle_name = RvalScalarValue(then_rval); + name = RvalScalarValue(then_rval); break; case RVAL_TYPE_FNCALL: - bundle_name = RvalFnCallValue(then_rval)->name; + name = RvalFnCallValue(then_rval)->name; + bundle_args = RvalFnCallValue(then_rval)->args; break; default: break; } - const Bundle *bundle = NULL; - if (bundle_name != NULL) + const Bundle *bp = NULL; + if (name != NULL) { - bundle = EvalContextResolveBundleExpression(ctx, PromiseGetPolicy(pp), bundle_name, "agent"); - if (bundle == NULL) + bp = EvalContextResolveBundleExpression(ctx, PromiseGetPolicy(pp), name, "agent"); + if (bp == NULL) { - bundle = EvalContextResolveBundleExpression(ctx, PromiseGetPolicy(pp), bundle_name, "common"); + bp = EvalContextResolveBundleExpression(ctx, PromiseGetPolicy(pp), name, "common"); } } - if (bundle == NULL) + if (bp == NULL) { - Log(LOG_LEVEL_ERR, "Reactor events promise '%s' refers to unknown bundle '%s', ignoring", - key, (bundle_name != NULL) ? bundle_name : "(invalid)"); - return false; + Log(LOG_LEVEL_ERR, "Reactor events promise '%s' refers to unknown bundle '%s'", + pp->promiser, (name != NULL) ? name : "(invalid)"); + return NULL; } - return true; + if (args != NULL) + { + *args = bundle_args; + } + return bp; } // Temporary limitations: @@ -90,37 +101,30 @@ 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); + pp->promiser); 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); + Log(LOG_LEVEL_ERR, "Reactor events promise '%s' does not specify a 'then' bundle, ignoring", pp->promiser); return PROMISE_RESULT_FAIL; } - if (!ThenBundleExists(ctx, pp, then_constraint->rval, key)) + if (ResolveThenBundle(ctx, pp, then_constraint->rval, NULL) == NULL) { - 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); + Log(LOG_LEVEL_INFO, "Registering a file_deleted watcher for events promise '%s', on file '%s'", pp->promiser, path); + bool kept = WatcherRegister(pp->promiser, EVENT_FILE_DELETED, FileWatcherStateNew(path), pp->org_pp, 1); return (kept) ? PROMISE_RESULT_NOOP : PROMISE_RESULT_FAIL; } @@ -176,10 +180,78 @@ static void EvaluateReactorBundle(EvalContext *ctx, const Bundle *bp) EvalContextStackPopFrame(ctx); } +/* Runs the bundle referred to by the 'then' attribute of the (expanded) + * events promise. */ +static PromiseResult RunThenBundle(EvalContext *ctx, const Promise *pp) +{ + assert(pp != NULL); + + 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", pp->promiser); + return PROMISE_RESULT_FAIL; + } + + const Rlist *args = NULL; + const Bundle *bp = ResolveThenBundle(ctx, pp, then_constraint->rval, &args); + if (bp == NULL) + { + return PROMISE_RESULT_FAIL; + } + + + BundleBanner(bp, args); + EvalContextSetBundleArgs(ctx, args); + EvalContextStackPushBundleFrame(ctx, bp, args, false, NULL); + + PromiseResult result = ScheduleAgentOperations(ctx, bp); + + EvalContextStackPopFrame(ctx); /* bundle */ + EvalContextSetBundleArgs(ctx, NULL); + EndBundleBanner(bp); + + return result; +} + +typedef struct +{ + const char *promiser; +} ReactorEventParam; + +/* Promise actuator for an events promise whose watcher fired. Only acts on + * the iteration the watcher was registered for. */ +static PromiseResult KeepEventsPromiseOnEvent(EvalContext *ctx, const Promise *pp, void *param) +{ + assert(pp != NULL); + assert(param != NULL); + const ReactorEventParam *event = param; + + if (!StringEqual(pp->promiser, event->promiser)) + { + return PROMISE_RESULT_SKIPPED; + } + + return RunThenBundle(ctx, pp); +} + +void HandleReactorEvent(EvalContext *ctx, const Promise *pp, const char *promiser) +{ + assert(pp != NULL); + assert(promiser != NULL); + + ReactorEventParam event = { .promiser = promiser }; + + EvalContextStackPushBundleFrame(ctx, PromiseGetBundle(pp), NULL, false, NULL); + EvalContextStackPushBundleSectionFrame(ctx, pp->parent_section); + ExpandPromise(ctx, pp, KeepEventsPromiseOnEvent, &event); + EvalContextStackPopFrame(ctx); /* bundle section */ + EvalContextStackPopFrame(ctx); /* bundle */ +} + void KeepReactorPromises(EvalContext *ctx, const Policy *policy) { assert(policy != NULL); - WatcherRegistryClear(); for (size_t i = 0; i < SeqLength(policy->bundles); i++) { diff --git a/cf-reactor/reactor_transform.h b/cf-reactor/reactor_transform.h index b089a9b53cb..5dfcb6bd4c4 100644 --- a/cf-reactor/reactor_transform.h +++ b/cf-reactor/reactor_transform.h @@ -33,10 +33,21 @@ * 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. + * Expects no watchers to be registered: on re-read of the policy, the + * watchers of the previous policy must be discarded with + * WatcherRegistryClear() before that policy is destroyed, as they refer to + * its promises. */ void KeepReactorPromises(EvalContext *ctx, const Policy *policy); +/** + * @brief Keep the events promise of a watcher that fired (see WatcherEventFn + * in watcher.h): run the bundle of its 'then' attribute. + * + * @param pp the unexpanded events promise the watcher was registered for + * @param promiser the expanded promiser the watcher was registered with, + * selecting the iteration of the events promise to keep + */ +void HandleReactorEvent(EvalContext *ctx, const Promise *pp, const char *promiser); + #endif diff --git a/cf-reactor/watcher.c b/cf-reactor/watcher.c index e0be04da879..bcf577d7776 100644 --- a/cf-reactor/watcher.c +++ b/cf-reactor/watcher.c @@ -29,13 +29,11 @@ #include // IsPendingTermination() #include #include -#include // StringHash_untyped(), StringEqual_untyped() -#include +#include // StringEqual() #include #include #include -#include -#include +#include // ThreadLock(), ThreadUnlock() /* 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, @@ -46,27 +44,31 @@ typedef struct { - char *key; + char *promiser; // expanded promiser of the iteration of `promise` WatcherCheckFn check_callback; // resolved from `type` at WatcherRegister() time WatcherStateDestroyFn destroy_state; void *state; + const Promise *promise; // the unexpanded events promise, owned by the policy time_t poll_interval_secs; time_t next_due; } Watcher; static void WatcherDestroy(void *item); -static void DestroyRval(void *item); -/* 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()). */ +/* Guards `watchers` and `paused` 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; +static bool paused = false; /* see EventWatcherPause() */ static WakeupChannel wakeup_channel = { .fds = { -1, -1 } }; + +/* Watchers that fired, queued by the watcher thread for the main thread to + * handle. Not owned: the watcher thread only queues watchers of the registry + * while holding watchers_mutex, and WatcherRegistryClear() requires the queue + * to be empty, so the queued watchers are always registered. */ static ThreadedQueue *event_queue = NULL; static StoppableThread *watcher_thread = NULL; @@ -75,57 +77,56 @@ static StoppableThread *watcher_thread = NULL; void WatcherRegistryInitialize(void) { assert(watchers == NULL); - assert(event_to_bundle == NULL); watchers = SeqNew(4, WatcherDestroy); - event_to_bundle = MapNew(StringHash_untyped, StringEqual_untyped, NULL, DestroyRval); } void WatcherRegistryFinalize(void) { SeqDestroy(watchers); watchers = NULL; - MapDestroy(event_to_bundle); - event_to_bundle = NULL; } void WatcherRegistryClear(void) { - assert(watchers != NULL && event_to_bundle != NULL); + assert(watchers != NULL); ThreadLock(&watchers_mutex); + if (event_queue != NULL && !ThreadedQueueIsEmpty(event_queue)) + { + ProgrammingError("Clearing the reactor watcher registry with events of its " + "watchers still queued, call EventWatcherPause() first"); + } + 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) +/* A watcher is identified by its events promise and the expanded promiser + * of the iteration: the iterations of a promise share the same (unexpanded) + * promise, and different promises can have the same promiser. */ +static const Watcher *FindWatcher(const Promise *pp, const char *promiser) { - Rval *new = xmalloc(sizeof(Rval)); - *new = RvalCopy(val); - return new; -} - -static void DestroyRval(void *item) -{ - Rval *rval = item; - if (rval != NULL) + for (size_t i = 0; i < SeqLength(watchers); i++) { - RvalDestroy(*rval); - free(rval); + const Watcher *w = SeqAt(watchers, i); + if (w->promise == pp && StringEqual(w->promiser, promiser)) + { + return w; + } } + return NULL; } // Expects interval to be strictly greater than 0, otherwise the watcher thread will busy spin -bool WatcherRegister(const char *key, EventType type, void *state, Rval val, time_t interval) +bool WatcherRegister(const char *promiser, EventType type, void *state, const Promise *pp, time_t interval) { - assert(key != NULL); - assert(watchers != NULL && event_to_bundle != NULL); + assert(promiser != NULL); + assert(pp != NULL); + assert(watchers != NULL); assert(interval > 0); WatcherCheckFn check_callback = NULL; @@ -140,14 +141,14 @@ bool WatcherRegister(const char *key, EventType type, void *state, Rval val, tim // TODO: add more event types default: - ProgrammingError("Unknown reactor event type %d for watcher '%s'", (int) type, key); + ProgrammingError("Unknown reactor event type %d for watcher '%s'", (int) type, promiser); } ThreadLock(&watchers_mutex); - if (MapHasKey(event_to_bundle, key)) + if (FindWatcher(pp, promiser) != NULL) { - Log(LOG_LEVEL_ERR, "Reactor watcher key '%s' is already registered, ignoring the duplicate", key); + Log(LOG_LEVEL_ERR, "Reactor watcher '%s' is already registered, ignoring the duplicate", promiser); ThreadUnlock(&watchers_mutex); if (destroy_state != NULL) { @@ -157,15 +158,15 @@ bool WatcherRegister(const char *key, EventType type, void *state, Rval val, tim } Watcher *w = xmalloc(sizeof(Watcher)); - w->key = xstrdup(key); + w->promiser = xstrdup(promiser); w->state = state; + w->promise = pp; w->poll_interval_secs = interval; w->next_due = 0; /* due immediately on the watcher thread's first pass */ w->check_callback = check_callback; w->destroy_state = destroy_state; SeqAppend(watchers, w); - MapInsert(event_to_bundle, w->key, AllocateRval(val)); ThreadUnlock(&watchers_mutex); return true; @@ -178,7 +179,7 @@ static void WatcherDestroy(void *item) { w->destroy_state(w->state); } - free(w->key); + free(w->promiser); free(w); } @@ -199,7 +200,7 @@ static void WatcherThreadMain(StoppableThread *thread, ARG_UNUSED void *unused) ThreadLock(&watchers_mutex); - for (size_t i = 0; i < SeqLength(watchers); i++) + for (size_t i = 0; !paused && i < SeqLength(watchers); i++) { Watcher *w = SeqAt(watchers, i); @@ -212,7 +213,7 @@ static void WatcherThreadMain(StoppableThread *thread, ARG_UNUSED void *unused) w->next_due = now + w->poll_interval_secs; if (fired) { - ThreadedQueuePush(event_queue, SafeStringDuplicate(w->key)); + ThreadedQueuePush(event_queue, w); any_event = true; } } @@ -235,7 +236,7 @@ static void WatcherThreadMain(StoppableThread *thread, ARG_UNUSED void *unused) bool EventWatcherInitialize(int *fd) { assert(fd != NULL); - assert(watchers != NULL && event_to_bundle != NULL); /* WatcherRegistryInitialize() first */ + assert(watchers != NULL); /* WatcherRegistryInitialize() first */ if (!WakeupChannelOpen(&wakeup_channel)) { @@ -243,7 +244,7 @@ bool EventWatcherInitialize(int *fd) return false; } - event_queue = ThreadedQueueNew(16, free); + event_queue = ThreadedQueueNew(16, NULL); watcher_thread = StoppableThreadStart(WatcherThreadMain, NULL); if (watcher_thread == NULL) @@ -261,8 +262,26 @@ bool EventWatcherInitialize(int *fd) return true; } -void EventWatcherHandleEvents(int fd, fd_set *readfds) +/* Watchers are only destroyed by the main thread (WatcherRegistryClear()), + * which is also the one handling events, and never while events are queued, + * so the queued watchers can be used without holding watchers_mutex. Their + * promiser and promise are never modified by the watcher thread. */ +static void HandleQueuedEvents(EvalContext *ctx, WatcherEventFn on_event) { + assert(on_event != NULL); + + void *item; + while (ThreadedQueuePop(event_queue, &item, 0)) + { + const Watcher *w = item; + Log(LOG_LEVEL_NOTICE, "Reactor watcher '%s' fired", w->promiser); + on_event(ctx, w->promise, w->promiser); + } +} + +void EventWatcherHandleEvents(EvalContext *ctx, WatcherEventFn on_event, int fd, fd_set *readfds) +{ + assert(on_event != NULL); assert(readfds != NULL); if (!FD_ISSET(fd, readfds)) @@ -270,31 +289,36 @@ void EventWatcherHandleEvents(int fd, fd_set *readfds) return; } + /* Drain before popping, so that a wakeup for an event queued after the + * last pop stays in the channel for the next select() */ WakeupChannelDrain(&wakeup_channel); - void *item; - while (ThreadedQueuePop(event_queue, &item, 0)) - { - const char *key = item; + HandleQueuedEvents(ctx, on_event); +} - ThreadLock(&watchers_mutex); - Rval *bundle = MapGet(event_to_bundle, key); +void EventWatcherPause(EvalContext *ctx, WatcherEventFn on_event) +{ + assert(on_event != NULL); - 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 - } + /* A pass of the watcher thread holds watchers_mutex, so once `paused` is + * set no further events are queued */ + ThreadLock(&watchers_mutex); + paused = true; + ThreadUnlock(&watchers_mutex); - ThreadUnlock(&watchers_mutex); - free(item); + if (event_queue != NULL) + { + HandleQueuedEvents(ctx, on_event); } } +void EventWatcherResume(void) +{ + ThreadLock(&watchers_mutex); + paused = false; + ThreadUnlock(&watchers_mutex); +} + void EventWatcherFinalize(void) { bool joined = StoppableThreadStop(watcher_thread, WATCHER_THREAD_EXIT_TIMEOUT_SECS); diff --git a/cf-reactor/watcher.h b/cf-reactor/watcher.h index 29a043b81d4..71448256348 100644 --- a/cf-reactor/watcher.h +++ b/cf-reactor/watcher.h @@ -35,27 +35,61 @@ typedef enum typedef bool (*WatcherCheckFn)(void *state); typedef void (*WatcherStateDestroyFn)(void *state); +/** + * @brief Called on the main thread for every event of a registered watcher. + * + * @param pp the (unexpanded) events promise the watcher was registered for + * @param promiser the expanded promiser the watcher was registered with, + * selecting the iteration of the events promise + * @note Called without the watcher registry locked, but it must not call + * WatcherRegistryClear(), as the events of the watchers are being + * handled. + */ +typedef void (*WatcherEventFn)(EvalContext *ctx, const Promise *pp, const char *promiser); + 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. + * + * @note No events of the watchers may be queued, i.e. once the event watcher + * is initialized, call EventWatcherPause() first. */ void WatcherRegistryClear(void); /** * @brief Register a specific watcher instance. - * - * @param key the events promise identifier + * + * @param promiser the expanded promiser of the iteration of the events + * promise. Together with pp, identifies the watcher. * @param type the type of watcher, defined in when bodies * @param state the data used for by the watcher, depending on the type - * @param val the rval holding the bundle to run on event + * @param pp the unexpanded events promise, passed back to the WatcherEventFn + * on event. Not owned, so the registry must be cleared before the + * policy holding it is destroyed. * @param interval interval between runs */ -bool WatcherRegister(const char *key, EventType type, void *state, Rval val, time_t interval); +bool WatcherRegister(const char *promiser, EventType type, void *state, const Promise *pp, time_t interval); bool EventWatcherInitialize(int *fd); -void EventWatcherHandleEvents(int fd, fd_set *readfds); +void EventWatcherHandleEvents(EvalContext *ctx, WatcherEventFn on_event, int fd, fd_set *readfds); + +/** + * @brief Stop checking for events, and handle the events already queued, so + * that the watchers can be cleared (see WatcherRegistryClear()) without + * losing any event. Events are checked for again after EventWatcherResume(). + * + * @note Only the events already detected are guaranteed to be handled: + * whether an event happening while paused is detected afterwards + * depends on the state the watcher is registered again with. + */ +void EventWatcherPause(EvalContext *ctx, WatcherEventFn on_event); + +/** + * @brief Check for events again, after EventWatcherPause(). + */ +void EventWatcherResume(void); void EventWatcherFinalize(void); #endif From ce5d017c7f9f9e932c807728f257377f29286ed8 Mon Sep 17 00:00:00 2001 From: Victor Moene Date: Wed, 30 Sep 2026 13:20:32 +0200 Subject: [PATCH 4/6] Fixed skipped promises on event Promises are skipped to prevent running them several times. However, on event, we want to run them every single time: - we clear the promise lock cache - we set the default if_elapsed time for bundles run from an events promise to be 0, so it doesn't skip the promises. Fixed also connection cache and custom promise prologue and epilogue. Clear function cache before "then" bundle run Signed-off-by: Victor Moene --- cf-reactor/reactor_transform.c | 18 ++++++++++++++++++ libpromises/attributes.c | 17 +++++++++++++++++ libpromises/attributes.h | 4 ++++ libpromises/eval_context.c | 12 ++++++++++++ libpromises/eval_context.h | 2 ++ libpromises/mod_custom.c | 1 + 6 files changed, 54 insertions(+) diff --git a/cf-reactor/reactor_transform.c b/cf-reactor/reactor_transform.c index 4d8ed5825c9..6660a8308aa 100644 --- a/cf-reactor/reactor_transform.c +++ b/cf-reactor/reactor_transform.c @@ -35,6 +35,9 @@ #include #include #include // ScheduleAgentOperations() +#include +#include // ConnCache_Init(), ConnCache_Destroy() +#include // Initialize/FinalizeCustomPromises() /* Promise types evaluated within `bundle reactor NAME { ... }`. */ static const char *const REACTOR_TYPESEQUENCE[] = @@ -200,12 +203,27 @@ static PromiseResult RunThenBundle(EvalContext *ctx, const Promise *pp) return PROMISE_RESULT_FAIL; } + /* The promise lock cache makes each promise act at most once per + * EvalContext, and the function cache keeps the results of functions like + * execresult(), but cf-reactor keeps the same EvalContext across events. + * Clear them so every event gets a fresh run of the bundle. */ + EvalContextPromiseLockCacheClear(ctx); + EvalContextFunctionCacheClear(ctx); BundleBanner(bp, args); EvalContextSetBundleArgs(ctx, args); EvalContextStackPushBundleFrame(ctx, bp, args, false, NULL); + /* Remote copy_from needs the connection cache and custom promise types + * need the promise modules map. Set them up per run like in cf-agent so + * that no connections or promise modules are kept between events. */ + ConnCache_Init(); + InitializeCustomPromises(); + int prev_ifelapsed = OverrideIfelapsed(0); PromiseResult result = ScheduleAgentOperations(ctx, bp); + RestoreIfelapsed(prev_ifelapsed); + FinalizeCustomPromises(); + ConnCache_Destroy(); EvalContextStackPopFrame(ctx); /* bundle */ EvalContextSetBundleArgs(ctx, NULL); diff --git a/libpromises/attributes.c b/libpromises/attributes.c index dba6e73b375..d54ec5b8ae0 100644 --- a/libpromises/attributes.c +++ b/libpromises/attributes.c @@ -666,6 +666,23 @@ LogLevel ActionAttributeLogLevelFromString(const char *log_level) } } +int OverrideIfelapsed(int minutes) +{ + assert(minutes >= 0); + + int prev_ifelapsed = VIFELAPSED; + VIFELAPSED = minutes; + + Log(LOG_LEVEL_DEBUG, "Default ifelapsed set to %d (was %d)", VIFELAPSED, prev_ifelapsed); + return prev_ifelapsed; +} + +void RestoreIfelapsed(int prev_ifelapsed) +{ + Log(LOG_LEVEL_DEBUG, "Default ifelapsed restored to %d (was %d)", prev_ifelapsed, VIFELAPSED); + VIFELAPSED = prev_ifelapsed; +} + static TransactionContext GetTransactionConstraints(const EvalContext *ctx, const Promise *pp) { TransactionContext t; diff --git a/libpromises/attributes.h b/libpromises/attributes.h index 02237f08b19..b836efb2af1 100644 --- a/libpromises/attributes.h +++ b/libpromises/attributes.h @@ -28,6 +28,10 @@ #include LogLevel ActionAttributeLogLevelFromString(const char *log_level); + +int OverrideIfelapsed(int minutes); +void RestoreIfelapsed(int prev_ifelapsed); + bool IsClassesBodyConstraint(const char *constraint); Attributes GetClassContextAttributes(const EvalContext *ctx, const Promise *pp); Attributes GetColumnAttributes(const EvalContext *ctx, const Promise *pp); diff --git a/libpromises/eval_context.c b/libpromises/eval_context.c index feb5766b012..2fbe3ffddb8 100644 --- a/libpromises/eval_context.c +++ b/libpromises/eval_context.c @@ -2929,6 +2929,18 @@ void EvalContextPromiseLockCacheRemove(EvalContext *ctx, const char *key) StringSetRemove(ctx->promise_lock_cache, key); } +void EvalContextPromiseLockCacheClear(EvalContext *ctx) +{ + assert(ctx != NULL); + StringSetClear(ctx->promise_lock_cache); +} + +void EvalContextFunctionCacheClear(EvalContext *ctx) +{ + assert(ctx != NULL); + FuncCacheMapClear(ctx->function_cache); +} + bool EvalContextFunctionCacheGet(const EvalContext *ctx, const FnCall *fp ARG_UNUSED, const Rlist *args, Rval *rval_out) diff --git a/libpromises/eval_context.h b/libpromises/eval_context.h index 2dd0f1df05d..2537c083f65 100644 --- a/libpromises/eval_context.h +++ b/libpromises/eval_context.h @@ -239,6 +239,8 @@ VariableTableIterator *EvalContextVariableTableFromRefIteratorNew(const EvalCont bool EvalContextPromiseLockCacheContains(const EvalContext *ctx, const char *key); void EvalContextPromiseLockCachePut(EvalContext *ctx, const char *key); void EvalContextPromiseLockCacheRemove(EvalContext *ctx, const char *key); +void EvalContextPromiseLockCacheClear(EvalContext *ctx); +void EvalContextFunctionCacheClear(EvalContext *ctx); bool EvalContextFunctionCacheGet(const EvalContext *ctx, const FnCall *fp, const Rlist *args, Rval *rval_out); void EvalContextFunctionCachePut(EvalContext *ctx, const FnCall *fp, const Rlist *args, const Rval *rval); diff --git a/libpromises/mod_custom.c b/libpromises/mod_custom.c index 31d6b985ac7..f56a691908c 100644 --- a/libpromises/mod_custom.c +++ b/libpromises/mod_custom.c @@ -1233,6 +1233,7 @@ bool InitializeCustomPromises() void FinalizeCustomPromises() { MapDestroy(custom_modules); + custom_modules = NULL; } PromiseResult EvaluateCustomPromise(EvalContext *ctx, const Promise *pp) From dae57f1919792aad338440d42dc5aa08deefa155 Mon Sep 17 00:00:00 2001 From: Victor Moene Date: Wed, 30 Sep 2026 17:05:42 +0200 Subject: [PATCH 5/6] Fix watcher test Signed-off-by: Victor Moene --- tests/unit/watcher_test.c | 115 ++++++++++++++++++++++++++++++++------ 1 file changed, 99 insertions(+), 16 deletions(-) diff --git a/tests/unit/watcher_test.c b/tests/unit/watcher_test.c index 507bfce37d9..9559d8ef89d 100644 --- a/tests/unit/watcher_test.c +++ b/tests/unit/watcher_test.c @@ -2,10 +2,14 @@ #include #include /* FileWatcherStateNew() */ -#include /* Rval */ +#include /* Promise */ #include /* LoggingPrivContext, LoggingPrivSetContext() */ #include /* strstr() */ +#include /* snprintf() */ +#include /* mkdtemp() */ +#include /* close(), unlink(), rmdir() */ +#include /* open() */ static int captured_err_count = 0; static char captured_err_message[256]; @@ -20,6 +24,18 @@ static char *CaptureErrorLogHook(ARG_UNUSED LoggingPrivContext *pctx, LogLevel l return (char *) message; } +/* The registry never dereferences the promise it is given, it only hands it + * back to the WatcherEventFn, so a dummy one is enough. */ +static const Promise dummy_promise; + +static int handled_event_count = 0; + +static void CountEvent(ARG_UNUSED EvalContext *ctx, ARG_UNUSED const Promise *pp, + ARG_UNUSED const char *promiser) +{ + handled_event_count++; +} + static void test_registry_initialize_finalize(void) { WatcherRegistryInitialize(); @@ -30,23 +46,16 @@ static void test_watcher_register_single(void) { WatcherRegistryInitialize(); - /* 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 }; - - assert_true(WatcherRegister("test-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/test-event"), bundle, 5)); + assert_true(WatcherRegister("test-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/test-event"), &dummy_promise, 5)); WatcherRegistryFinalize(); } -static void test_watcher_register_duplicate_key_ignored(void) +static void test_watcher_register_duplicate_ignored(void) { WatcherRegistryInitialize(); - const Rval bundle_a = { .item = (char *) "bundle_a", .type = RVAL_TYPE_SCALAR }; - const Rval bundle_b = { .item = (char *) "bundle_b", .type = RVAL_TYPE_SCALAR }; - - assert_true(WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-a"), bundle_a, 5)); + assert_true(WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-a"), &dummy_promise, 5)); captured_err_count = 0; captured_err_message[0] = '\0'; @@ -60,10 +69,10 @@ static void test_watcher_register_duplicate_key_ignored(void) const LogLevel old_level = LogGetGlobalLevel(); LogSetGlobalLevel(LOG_LEVEL_CRIT); - /* Registering the same key again must be rejected: an error is logged + /* Registering the same promise and promiser again must be rejected: an error is logged * (not silently swallowed) and the first registration is kept, not * replaced. */ - const bool registered = WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-b"), bundle_b, 5); + const bool registered = WatcherRegister("dup-event", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/dup-event-b"), &dummy_promise, 5); LogSetGlobalLevel(old_level); LoggingPrivSetContext(NULL); @@ -76,6 +85,21 @@ static void test_watcher_register_duplicate_key_ignored(void) WatcherRegistryFinalize(); } +static void test_watcher_register_same_promiser_of_other_promise(void) +{ + WatcherRegistryInitialize(); + + /* Watchers are identified by their promise and promiser: iterations of + * a promise (same promise, other promisers), and promises with the same + * promiser (other promises), are all different watchers */ + static const Promise other_promise; + assert_true(WatcherRegister("/a", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/a"), &dummy_promise, 5)); + assert_true(WatcherRegister("/b", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/b"), &dummy_promise, 5)); + assert_true(WatcherRegister("/a", EVENT_FILE_DELETED, FileWatcherStateNew("/nonexistent/a"), &other_promise, 5)); + + WatcherRegistryFinalize(); +} + static void test_event_watcher_lifecycle(void) { WatcherRegistryInitialize(); @@ -86,20 +110,77 @@ static void test_event_watcher_lifecycle(void) fd_set readfds; FD_ZERO(&readfds); + handled_event_count = 0; + /* fd not set: must return early without touching the wakeup channel or * the event queue. */ - EventWatcherHandleEvents(fd, &readfds); + EventWatcherHandleEvents(NULL, CountEvent, fd, &readfds); FD_SET(fd, &readfds); /* fd set but nothing queued: must drain the channel (no-op) and find * nothing to pop from the event queue. */ - EventWatcherHandleEvents(fd, &readfds); + EventWatcherHandleEvents(NULL, CountEvent, fd, &readfds); + assert_int_equal(handled_event_count, 0); /* Must wake up and stop the watcher thread by itself (IsPendingTermination() * is false here), join it, and finalize the watcher registry. */ EventWatcherFinalize(); } +static const Promise *last_event_promise = NULL; + +static void RecordEvent(ARG_UNUSED EvalContext *ctx, const Promise *pp, + ARG_UNUSED const char *promiser) +{ + handled_event_count++; + last_event_promise = pp; +} + +static void test_event_watcher_pause_handles_queued_events(void) +{ + char dir[] = "/tmp/watcher_test.XXXXXX"; + assert_true(mkdtemp(dir) != NULL); + char path[PATH_MAX]; + snprintf(path, sizeof(path), "%s/file", dir); + int file_fd = open(path, O_CREAT | O_WRONLY, 0600); + assert_true(file_fd >= 0); + close(file_fd); + + WatcherRegistryInitialize(); + assert_true(WatcherRegister("pause-event", EVENT_FILE_DELETED, FileWatcherStateNew(path), &dummy_promise, 1)); + + int fd = -1; + assert_true(EventWatcherInitialize(&fd)); + assert_int_equal(unlink(path), 0); + + /* Wait for the watcher thread to queue the event (it checks every + * second), without handling it */ + bool queued = false; + for (int i = 0; !queued && i < 50; i++) + { + fd_set readfds; + FD_ZERO(&readfds); + FD_SET(fd, &readfds); + struct timeval timeout = { .tv_sec = 0, .tv_usec = 100000 }; + queued = (select(fd + 1, &readfds, NULL, NULL, &timeout) > 0); + } + assert_true(queued); + + /* Pausing handles the queued event, with the promise of its watcher, so + * that the registry can be cleared without losing it */ + handled_event_count = 0; + last_event_promise = NULL; + EventWatcherPause(NULL, RecordEvent); + assert_int_equal(handled_event_count, 1); + assert_true(last_event_promise == &dummy_promise); + + WatcherRegistryClear(); + EventWatcherResume(); + + EventWatcherFinalize(); + assert_int_equal(rmdir(dir), 0); +} + int main() { PRINT_TEST_BANNER(); @@ -107,8 +188,10 @@ int main() { unit_test(test_registry_initialize_finalize), unit_test(test_watcher_register_single), - unit_test(test_watcher_register_duplicate_key_ignored), + unit_test(test_watcher_register_duplicate_ignored), + unit_test(test_watcher_register_same_promiser_of_other_promise), unit_test(test_event_watcher_lifecycle), + unit_test(test_event_watcher_pause_handles_queued_events), }; return run_tests(tests); From bcace8d033a6be207ee1a2b97e50e42f1627d1ad Mon Sep 17 00:00:00 2001 From: Victor Moene Date: Tue, 6 Oct 2026 17:23:44 +0200 Subject: [PATCH 6/6] Added syntax arity check for events promises' "then" - Moved mod_methods.c's MethodsParseTreeCheck implementation to PromiseCheckBundleCallArity in policy.c, and made it more generic - Now, MethodsParseTreeCheck and EventsParseTreeCheck both call this new function when calling bundles Signed-off-by: Victor Moene --- libpromises/mod_methods.c | 49 +------------------------------------ libpromises/mod_reactor.c | 8 +++++- libpromises/policy.c | 51 +++++++++++++++++++++++++++++++++++++++ libpromises/policy.h | 2 ++ 4 files changed, 61 insertions(+), 49 deletions(-) diff --git a/libpromises/mod_methods.c b/libpromises/mod_methods.c index a51e0a34365..a9716bda3cb 100644 --- a/libpromises/mod_methods.c +++ b/libpromises/mod_methods.c @@ -26,13 +26,6 @@ #include #include -#include -#include -#include -#include - -static const char *const POLICY_ERROR_METHODS_BUNDLE_ARITY = - "Conflicting arity in calling bundle %s, expected %d arguments, %d given"; static const ConstraintSyntax CF_METHOD_BODIES[] = { @@ -44,47 +37,7 @@ static const ConstraintSyntax CF_METHOD_BODIES[] = static bool MethodsParseTreeCheck(const Promise *pp, Seq *errors) { - bool success = true; - - for (size_t i = 0; i < SeqLength(pp->conlist); i++) - { - const Constraint *cp = SeqAt(pp->conlist, i); - - // ensure: if call and callee are resolved, then they have matching arity - if (StringEqual(cp->lval, "usebundle")) - { - if (cp->rval.type == RVAL_TYPE_FNCALL) - { - // HACK: exploiting the fact that class-references and call-references are similar - FnCall *call = RvalFnCallValue(cp->rval); - ClassRef ref = ClassRefParse(call->name); - if (!ClassRefIsQualified(ref)) - { - ClassRefQualify(&ref, PromiseGetNamespace(pp)); - } - - const Bundle *callee = PolicyGetBundle(PolicyFromPromise(pp), ref.ns, "agent", ref.name); - if (!callee) - { - callee = PolicyGetBundle(PolicyFromPromise(pp), ref.ns, "common", ref.name); - } - - ClassRefDestroy(ref); - - if (callee) - { - if (RlistLen(call->args) != RlistLen(callee->args)) - { - SeqAppend(errors, PolicyErrorNew(POLICY_ELEMENT_TYPE_CONSTRAINT, cp, - POLICY_ERROR_METHODS_BUNDLE_ARITY, - call->name, RlistLen(callee->args), RlistLen(call->args))); - success = false; - } - } - } - } - } - return success; + return PromiseCheckBundleCallArity(pp, "usebundle", errors); } const PromiseTypeSyntax CF_METHOD_PROMISE_TYPES[] = diff --git a/libpromises/mod_reactor.c b/libpromises/mod_reactor.c index 0bc16cd1fc8..d07b0bd46cd 100644 --- a/libpromises/mod_reactor.c +++ b/libpromises/mod_reactor.c @@ -25,6 +25,7 @@ #include #include +#include static const ConstraintSyntax when_constraints[] = { @@ -44,8 +45,13 @@ static const ConstraintSyntax CF_EVENT_BODIES[] = ConstraintSyntaxNewNull() }; +static bool EventsParseTreeCheck(const Promise *pp, Seq *errors) +{ + return PromiseCheckBundleCallArity(pp, "then", errors); +} + const PromiseTypeSyntax CF_REACTOR_PROMISE_TYPES[] = { - PromiseTypeSyntaxNew("reactor", "events", CF_EVENT_BODIES, NULL, SYNTAX_STATUS_NORMAL), + PromiseTypeSyntaxNew("reactor", "events", CF_EVENT_BODIES, &EventsParseTreeCheck, SYNTAX_STATUS_NORMAL), PromiseTypeSyntaxNewNull() }; diff --git a/libpromises/policy.c b/libpromises/policy.c index 0b9fc3ae496..0d72d75c3d1 100644 --- a/libpromises/policy.c +++ b/libpromises/policy.c @@ -33,6 +33,7 @@ #include #include #include +#include #include #include #include @@ -67,6 +68,8 @@ static const char *const POLICY_ERROR_PROMISE_ATTRIBUTE_NOT_SUPPORTED = "Common attribute '%s' not supported for custom promises, use '%s' instead (%s promises)"; static const char *const POLICY_ERROR_PROMISE_TYPE_UNSUPPORTED = "Promise type '%s' not supported by '%s' bundle type"; +static const char *const POLICY_ERROR_BUNDLE_CALL_ARITY = + "Conflicting arity in calling bundle %s, expected %d arguments, %d given"; static const char *const POLICY_ERROR_CONSTRAINT_TYPE_MISMATCH = "Type mismatch in constraint: %s"; @@ -2727,6 +2730,54 @@ static bool ValidateCustomPromise(const Promise *pp, Seq *errors) return valid; } +bool PromiseCheckBundleCallArity(const Promise *pp, const char *lval, Seq *errors) +{ + assert(pp != NULL); + assert(lval != NULL); + + bool success = true; + + for (size_t i = 0; i < SeqLength(pp->conlist); i++) + { + const Constraint *cp = SeqAt(pp->conlist, i); + + // ensure: if call and callee are resolved, then they have matching arity + if (StringEqual(cp->lval, lval)) + { + if (cp->rval.type == RVAL_TYPE_FNCALL) + { + // HACK: exploiting the fact that class-references and call-references are similar + FnCall *call = RvalFnCallValue(cp->rval); + ClassRef ref = ClassRefParse(call->name); + if (!ClassRefIsQualified(ref)) + { + ClassRefQualify(&ref, PromiseGetNamespace(pp)); + } + + const Bundle *callee = PolicyGetBundle(PolicyFromPromise(pp), ref.ns, "agent", ref.name); + if (!callee) + { + callee = PolicyGetBundle(PolicyFromPromise(pp), ref.ns, "common", ref.name); + } + + ClassRefDestroy(ref); + + if (callee) + { + if (RlistLen(call->args) != RlistLen(callee->args)) + { + SeqAppend(errors, PolicyErrorNew(POLICY_ELEMENT_TYPE_CONSTRAINT, cp, + POLICY_ERROR_BUNDLE_CALL_ARITY, + call->name, RlistLen(callee->args), RlistLen(call->args))); + success = false; + } + } + } + } + } + return success; +} + static bool PromiseCheck(const Promise *pp, Seq *errors) { assert(pp != NULL); diff --git a/libpromises/policy.h b/libpromises/policy.h index 24828fd5696..1defbd3a587 100644 --- a/libpromises/policy.h +++ b/libpromises/policy.h @@ -212,6 +212,8 @@ const char *PromiseGetNamespace(const Promise *pp); const Bundle *PromiseGetBundle(const Promise *pp); const Policy *PromiseGetPolicy(const Promise *pp); +bool PromiseCheckBundleCallArity(const Promise *pp, const char *lval, Seq *errors); + static inline const char *PromiseGetPromiseType(const Promise *pp) { assert(pp != NULL);