From d7d5f10825e8a817049b28b239e91bc29466e9a0 Mon Sep 17 00:00:00 2001 From: CarterPerez-dev Date: Wed, 29 Apr 2026 01:05:44 -0400 Subject: [PATCH] feat(engine): Scheduler fiber publishes SchedulerTick on a fixed interval Initial tick fires immediately on start so first policy evaluation doesn't wait for the full interval on boot. The PolicyEvaluator subscribes to SchedulerTick events and invokes evaluate_all -> overdue discovery -> RotationScheduled fan-out. 3 specs verify boot-tick, periodic cadence, and clean stop(). --- .../spec/unit/engine/scheduler_spec.cr | 76 +++++++++++++++++++ .../src/cre/engine/scheduler.cr | 52 +++++++++++++ 2 files changed, 128 insertions(+) create mode 100644 PROJECTS/intermediate/credential-rotation-enforcer/spec/unit/engine/scheduler_spec.cr create mode 100644 PROJECTS/intermediate/credential-rotation-enforcer/src/cre/engine/scheduler.cr diff --git a/PROJECTS/intermediate/credential-rotation-enforcer/spec/unit/engine/scheduler_spec.cr b/PROJECTS/intermediate/credential-rotation-enforcer/spec/unit/engine/scheduler_spec.cr new file mode 100644 index 00000000..7e197951 --- /dev/null +++ b/PROJECTS/intermediate/credential-rotation-enforcer/spec/unit/engine/scheduler_spec.cr @@ -0,0 +1,76 @@ +# =================== +# ©AngelaMos | 2026 +# scheduler_spec.cr +# =================== + +require "../../spec_helper" +require "../../../src/cre/engine/scheduler" +require "../../../src/cre/engine/event_bus" + +private def drain(ch : ::Channel(CRE::Events::Event)) : Array(CRE::Events::Event) + out = [] of CRE::Events::Event + loop do + select + when ev = ch.receive + out << ev + else + break + end + end + out +end + +describe CRE::Engine::Scheduler do + it "publishes a tick immediately on start" do + bus = CRE::Engine::EventBus.new + ch = bus.subscribe + bus.run + + scheduler = CRE::Engine::Scheduler.new(bus, interval: 5.seconds) + scheduler.start + sleep 0.1.seconds + scheduler.stop + + ticks = drain(ch).count(&.is_a?(CRE::Events::SchedulerTick)) + ticks.should be >= 1 + ensure + bus.try(&.stop) + end + + it "publishes ticks at the configured interval" do + bus = CRE::Engine::EventBus.new + ch = bus.subscribe(buffer: 64) + bus.run + + scheduler = CRE::Engine::Scheduler.new(bus, interval: 0.05.seconds) + scheduler.start + sleep 0.18.seconds + scheduler.stop + sleep 0.05.seconds # let final tick land + + ticks = drain(ch).count(&.is_a?(CRE::Events::SchedulerTick)) + # Initial + ~3 interval ticks = 3-5 expected; allow some scheduling variance + ticks.should be >= 2 + ticks.should be <= 6 + ensure + bus.try(&.stop) + end + + it "stop halts publication" do + bus = CRE::Engine::EventBus.new + ch = bus.subscribe + bus.run + + scheduler = CRE::Engine::Scheduler.new(bus, interval: 0.05.seconds) + scheduler.start + sleep 0.06.seconds + scheduler.stop + drain(ch) # drain whatever was already published + sleep 0.2.seconds + + later = drain(ch).count(&.is_a?(CRE::Events::SchedulerTick)) + later.should be <= 1 # at most one in-flight tick from the last sleep cycle + ensure + bus.try(&.stop) + end +end diff --git a/PROJECTS/intermediate/credential-rotation-enforcer/src/cre/engine/scheduler.cr b/PROJECTS/intermediate/credential-rotation-enforcer/src/cre/engine/scheduler.cr new file mode 100644 index 00000000..3b232b2b --- /dev/null +++ b/PROJECTS/intermediate/credential-rotation-enforcer/src/cre/engine/scheduler.cr @@ -0,0 +1,52 @@ +# =================== +# ©AngelaMos | 2026 +# scheduler.cr +# =================== + +require "log" +require "./event_bus" +require "../events/system_events" + +module CRE::Engine + # Scheduler is a fiber that publishes SchedulerTick events at a fixed interval. + # The PolicyEvaluator subscriber listens for these and runs evaluate_all, + # which discovers overdue credentials and publishes RotationScheduled etc. + # + # The fiber owns its own lifecycle: start() spawns it; stop() flips a flag and + # the next tick's check exits cleanly. + class Scheduler + Log = ::Log.for("cre.scheduler") + + @running : Bool + + def initialize(@bus : EventBus, @interval : Time::Span = 60.seconds) + @running = false + end + + def start : Nil + @running = true + spawn(name: "scheduler") do + # Fire one tick immediately so the first evaluation doesn't wait + # for the full interval on boot. + @bus.publish Events::SchedulerTick.new + while @running + sleep @interval + break unless @running + begin + @bus.publish Events::SchedulerTick.new + rescue ex + Log.error(exception: ex) { "scheduler tick publish failed" } + end + end + end + end + + def stop : Nil + @running = false + end + + def running? : Bool + @running + end + end +end