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..58a9a094db9 --- /dev/null +++ b/cf-agent/agent_operations.c @@ -0,0 +1,819 @@ +/* + 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. +*/ + +/* 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 + +#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; + +/* 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[] = +{ + "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) +// NOTE: 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(); + + /* 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(); + } + + 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. + // 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) + { + 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(); + + /* 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(); + } + + 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++) + { + 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: 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"); + 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: + /* 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; + + 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: 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 + { + 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..8604ba69513 --- /dev/null +++ b/cf-agent/agent_operations.h @@ -0,0 +1,67 @@ +/* + 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 + +/** + * @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 */ +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..a015224bffb 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(); @@ -1077,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) { @@ -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/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..6660a8308aa 100644 --- a/cf-reactor/reactor_transform.c +++ b/cf-reactor/reactor_transform.c @@ -34,6 +34,10 @@ #include #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[] = @@ -45,41 +49,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 +104,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 +183,93 @@ 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; + } + + /* 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); + 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 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) 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); 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); 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);