|
@@ -50,6 +50,39 @@ void WorkflowScheduler::stop() {
|
|
|
spdlog::info("Workflow scheduler stopped");
|
|
spdlog::info("Workflow scheduler stopped");
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+namespace {
|
|
|
|
|
+
|
|
|
|
|
+// When this entry should next run, as a steady_clock point.
|
|
|
|
|
+//
|
|
|
|
|
+// For a cron entry the answer comes from wall-clock time in the workflow's
|
|
|
|
|
+// timezone and is then expressed as a delay from `from`. Recomputed after every
|
|
|
|
|
+// run rather than advanced by a fixed step, so a schedule stays pinned to the
|
|
|
|
|
+// times it names even across a daylight-saving change.
|
|
|
|
|
+std::chrono::steady_clock::time_point nextRunAfter(const ScheduledWorkflow& entry,
|
|
|
|
|
+ std::chrono::steady_clock::time_point from) {
|
|
|
|
|
+ if (!entry.cron) {
|
|
|
|
|
+ return from + std::chrono::minutes(entry.interval_minutes);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ const auto wall_now = std::chrono::system_clock::now();
|
|
|
|
|
+ auto wall_next = entry.cron->nextAfter(wall_now, entry.timezone);
|
|
|
|
|
+ if (!wall_next) {
|
|
|
|
|
+ // Nothing this expression can ever match - "30 4 31 2 *" and the like.
|
|
|
|
|
+ // Parked rather than retried every tick, and said out loud once here.
|
|
|
|
|
+ spdlog::error("Workflow '{}' has cron \"{}\" which can never match a real date; "
|
|
|
|
|
+ "it will not run", entry.workflow_name, entry.cron->expression());
|
|
|
|
|
+ return from + std::chrono::hours(24 * 365);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ auto delay = *wall_next - wall_now;
|
|
|
|
|
+ if (delay < std::chrono::seconds(0)) {
|
|
|
|
|
+ delay = std::chrono::seconds(0);
|
|
|
|
|
+ }
|
|
|
|
|
+ return from + std::chrono::duration_cast<std::chrono::steady_clock::duration>(delay);
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+} // namespace
|
|
|
|
|
+
|
|
|
void WorkflowScheduler::registerWorkflow(
|
|
void WorkflowScheduler::registerWorkflow(
|
|
|
const std::string& workflow_id,
|
|
const std::string& workflow_id,
|
|
|
const std::string& workflow_name,
|
|
const std::string& workflow_name,
|
|
@@ -58,9 +91,26 @@ void WorkflowScheduler::registerWorkflow(
|
|
|
int interval_minutes,
|
|
int interval_minutes,
|
|
|
OverlapPolicy overlap_policy,
|
|
OverlapPolicy overlap_policy,
|
|
|
int max_concurrent,
|
|
int max_concurrent,
|
|
|
- int max_run_minutes
|
|
|
|
|
|
|
+ int max_run_minutes,
|
|
|
|
|
+ const std::string& cron_expression,
|
|
|
|
|
+ const std::string& timezone
|
|
|
) {
|
|
) {
|
|
|
- if (interval_minutes <= 0) {
|
|
|
|
|
|
|
+ std::optional<common::CronSchedule> cron;
|
|
|
|
|
+ if (!cron_expression.empty()) {
|
|
|
|
|
+ auto parsed = common::CronSchedule::parse(cron_expression);
|
|
|
|
|
+ if (parsed.failed()) {
|
|
|
|
|
+ // Not scheduled at all. Falling back to the interval would be the
|
|
|
|
|
+ // very failure this replaced: a workflow running on a schedule
|
|
|
|
|
+ // nobody asked for, with nothing saying so.
|
|
|
|
|
+ spdlog::error("Workflow '{}' ({}) has a cron expression that cannot be used, so it "
|
|
|
|
|
+ "is not scheduled: {}", workflow_name, workflow_id,
|
|
|
|
|
+ parsed.error().message());
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ cron = parsed.value();
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!cron && interval_minutes <= 0) {
|
|
|
spdlog::debug("Workflow {} has interval 0, not scheduling", workflow_id);
|
|
spdlog::debug("Workflow {} has interval 0, not scheduling", workflow_id);
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
@@ -75,8 +125,10 @@ void WorkflowScheduler::registerWorkflow(
|
|
|
entry.trigger_node_id = trigger_node_id;
|
|
entry.trigger_node_id = trigger_node_id;
|
|
|
entry.trigger_type = trigger_type;
|
|
entry.trigger_type = trigger_type;
|
|
|
entry.interval_minutes = interval_minutes;
|
|
entry.interval_minutes = interval_minutes;
|
|
|
|
|
+ entry.cron = cron;
|
|
|
|
|
+ entry.timezone = timezone;
|
|
|
entry.last_run = now; // Consider it just ran to avoid immediate execution
|
|
entry.last_run = now; // Consider it just ran to avoid immediate execution
|
|
|
- entry.next_run = now + std::chrono::minutes(interval_minutes);
|
|
|
|
|
|
|
+ entry.next_run = nextRunAfter(entry, now);
|
|
|
entry.overlap_policy = overlap_policy;
|
|
entry.overlap_policy = overlap_policy;
|
|
|
entry.max_concurrent = max_concurrent > 0 ? max_concurrent : 1;
|
|
entry.max_concurrent = max_concurrent > 0 ? max_concurrent : 1;
|
|
|
entry.max_run_minutes = max_run_minutes;
|
|
entry.max_run_minutes = max_run_minutes;
|
|
@@ -87,14 +139,32 @@ void WorkflowScheduler::registerWorkflow(
|
|
|
if (existing != workflows_.end()) {
|
|
if (existing != workflows_.end()) {
|
|
|
entry.active_runs = existing->second.active_runs;
|
|
entry.active_runs = existing->second.active_runs;
|
|
|
entry.pending_start = existing->second.pending_start;
|
|
entry.pending_start = existing->second.pending_start;
|
|
|
- entry.next_run = existing->second.next_run;
|
|
|
|
|
|
|
+
|
|
|
|
|
+ // The next run is only carried over when the schedule itself has not
|
|
|
|
|
+ // changed. Keeping it unconditionally means a schedule edit does not
|
|
|
|
|
+ // take effect until the workflow next fires on the OLD schedule -
|
|
|
|
|
+ // switching an hourly trigger to "Monday 02:00" left it due in 54
|
|
|
|
|
+ // minutes, and changing an interval from an hour to five minutes left
|
|
|
|
|
+ // it waiting the old hour out.
|
|
|
|
|
+ const auto& before = existing->second;
|
|
|
|
|
+ const std::string before_cron = before.cron ? before.cron->expression() : std::string();
|
|
|
|
|
+ const std::string after_cron = entry.cron ? entry.cron->expression() : std::string();
|
|
|
|
|
+ const bool same_schedule = before_cron == after_cron &&
|
|
|
|
|
+ before.timezone == entry.timezone &&
|
|
|
|
|
+ before.interval_minutes == entry.interval_minutes;
|
|
|
|
|
+ if (same_schedule) {
|
|
|
|
|
+ entry.next_run = before.next_run;
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
workflows_[workflow_id] = entry;
|
|
workflows_[workflow_id] = entry;
|
|
|
|
|
|
|
|
- spdlog::info("Scheduled workflow '{}' ({}) with {} trigger, interval: {} minutes, "
|
|
|
|
|
- "overlap: {}, maxConcurrent: {}",
|
|
|
|
|
- workflow_name, workflow_id, trigger_type, interval_minutes,
|
|
|
|
|
|
|
+ const std::string when = entry.cron
|
|
|
|
|
+ ? ("cron \"" + entry.cron->expression() + "\"" +
|
|
|
|
|
+ (entry.timezone.empty() ? std::string(" (UTC)") : " (" + entry.timezone + ")"))
|
|
|
|
|
+ : ("every " + std::to_string(interval_minutes) + " minutes");
|
|
|
|
|
+ spdlog::info("Scheduled workflow '{}' ({}) with {} trigger, {}, overlap: {}, maxConcurrent: {}",
|
|
|
|
|
+ workflow_name, workflow_id, trigger_type, when,
|
|
|
overlapPolicyToString(entry.overlap_policy), entry.max_concurrent);
|
|
overlapPolicyToString(entry.overlap_policy), entry.max_concurrent);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -184,6 +254,14 @@ void WorkflowScheduler::notifyExecutionStarted(const std::string& workflow_id,
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
int minutes = it->second.max_run_minutes;
|
|
int minutes = it->second.max_run_minutes;
|
|
|
|
|
+ if (minutes <= 0 && it->second.cron) {
|
|
|
|
|
+ // Deriving the stuck-run timeout from interval_minutes would be wrong
|
|
|
|
|
+ // here: on a cron entry that field is the trigger's poll interval, not
|
|
|
|
|
+ // its schedule. The gap to the next occurrence is the real horizon -
|
|
|
|
|
+ // past it, a run that has not finished is holding up the next one.
|
|
|
|
|
+ const auto gap = it->second.next_run - std::chrono::steady_clock::now();
|
|
|
|
|
+ minutes = static_cast<int>(std::chrono::duration_cast<std::chrono::minutes>(gap).count());
|
|
|
|
|
+ }
|
|
|
if (minutes <= 0) {
|
|
if (minutes <= 0) {
|
|
|
minutes = it->second.interval_minutes * 5;
|
|
minutes = it->second.interval_minutes * 5;
|
|
|
}
|
|
}
|
|
@@ -332,7 +410,7 @@ void WorkflowScheduler::checkAndExecute() {
|
|
|
if (active >= entry.max_concurrent) {
|
|
if (active >= entry.max_concurrent) {
|
|
|
spdlog::info("Workflow '{}' at concurrency limit {}, skipping tick",
|
|
spdlog::info("Workflow '{}' at concurrency limit {}, skipping tick",
|
|
|
entry.workflow_name, entry.max_concurrent);
|
|
entry.workflow_name, entry.max_concurrent);
|
|
|
- entry.next_run = now + std::chrono::minutes(entry.interval_minutes);
|
|
|
|
|
|
|
+ entry.next_run = nextRunAfter(entry, now);
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
|
break;
|
|
break;
|
|
@@ -344,7 +422,7 @@ void WorkflowScheduler::checkAndExecute() {
|
|
|
entry.workflow_name);
|
|
entry.workflow_name);
|
|
|
entry.pending_start = true;
|
|
entry.pending_start = true;
|
|
|
}
|
|
}
|
|
|
- entry.next_run = now + std::chrono::minutes(entry.interval_minutes);
|
|
|
|
|
|
|
+ entry.next_run = nextRunAfter(entry, now);
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
|
break;
|
|
break;
|
|
@@ -355,7 +433,7 @@ void WorkflowScheduler::checkAndExecute() {
|
|
|
spdlog::info("Workflow '{}' still running ({} in flight), skipping this "
|
|
spdlog::info("Workflow '{}' still running ({} in flight), skipping this "
|
|
|
"tick",
|
|
"tick",
|
|
|
entry.workflow_name, active);
|
|
entry.workflow_name, active);
|
|
|
- entry.next_run = now + std::chrono::minutes(entry.interval_minutes);
|
|
|
|
|
|
|
+ entry.next_run = nextRunAfter(entry, now);
|
|
|
continue;
|
|
continue;
|
|
|
}
|
|
}
|
|
|
break;
|
|
break;
|
|
@@ -389,9 +467,12 @@ void WorkflowScheduler::checkAndExecute() {
|
|
|
// and the following tick fires immediately.
|
|
// and the following tick fires immediately.
|
|
|
auto dispatched_at = std::chrono::steady_clock::now();
|
|
auto dispatched_at = std::chrono::steady_clock::now();
|
|
|
it->second.last_run = dispatched_at;
|
|
it->second.last_run = dispatched_at;
|
|
|
- it->second.next_run = dispatched_at + std::chrono::minutes(it->second.interval_minutes);
|
|
|
|
|
|
|
+ it->second.next_run = nextRunAfter(it->second, dispatched_at);
|
|
|
|
|
|
|
|
- auto next_in_minutes = it->second.interval_minutes;
|
|
|
|
|
|
|
+ // Read back from next_run rather than the interval, which a
|
|
|
|
|
+ // cron entry does not run on.
|
|
|
|
|
+ const auto next_in_minutes = std::chrono::duration_cast<std::chrono::minutes>(
|
|
|
|
|
+ it->second.next_run - dispatched_at).count();
|
|
|
spdlog::debug("Workflow '{}' next run in {} minutes",
|
|
spdlog::debug("Workflow '{}' next run in {} minutes",
|
|
|
entry.workflow_name, next_in_minutes);
|
|
entry.workflow_name, next_in_minutes);
|
|
|
}
|
|
}
|