From 46d1d430a986ce53f4d93164014293b7c14f4dff Mon Sep 17 00:00:00 2001 From: zewelor Date: Thu, 19 Mar 2026 20:41:52 +0000 Subject: [PATCH 1/4] Add change-detection and trigger registry - Add ChangeDetectionJob and TriggerState model with migration - Introduce TriggerCollection and stable trigger `unique_key` values - Add ChangeDetecting concern for detector triggers and scheduler support - Have RecurringTasksConfig schedule ChangeDetectionJob for detectors - Make RunWorkflowJob accept `trigger_key` and payload; pass payload to workflow - Add tests and fake trigger; update schema, docs, HTTP client, and add pg gem --- AGENTS.md | 22 ++-- Gemfile | 2 + Gemfile.lock | 19 ++- app/jobs/r3x/change_detection_job.rb | 69 +++++++++++ app/jobs/r3x/run_workflow_job.rb | 21 +++- app/lib/r3x/client/http.rb | 6 +- app/models/r3x/trigger_state.rb | 26 ++++ .../20260318100000_create_trigger_states.rb | 20 +++ db/schema.rb | 18 ++- lib/r3x/recurring_tasks_config.rb | 22 +++- lib/r3x/trigger_collection.rb | 35 ++++++ lib/r3x/trigger_execution.rb | 5 +- lib/r3x/triggers/base.rb | 14 ++- lib/r3x/triggers/concerns/change_detecting.rb | 15 +++ lib/r3x/workflow.rb | 20 +-- test/jobs/r3x/change_detection_job_test.rb | 115 ++++++++++++++++++ test/jobs/r3x/run_workflow_job_test.rb | 61 +++++++--- test/lib/r3x/recurring_tasks_config_test.rb | 33 ++++- test/lib/r3x/workflow_context_test.rb | 11 ++ test/lib/r3x/workflow_test.rb | 65 +++++++--- test/support/fake_change_detecting_trigger.rb | 30 +++++ 21 files changed, 553 insertions(+), 76 deletions(-) create mode 100644 app/jobs/r3x/change_detection_job.rb create mode 100644 app/models/r3x/trigger_state.rb create mode 100644 db/migrate/20260318100000_create_trigger_states.rb create mode 100644 lib/r3x/trigger_collection.rb create mode 100644 lib/r3x/triggers/concerns/change_detecting.rb create mode 100644 test/jobs/r3x/change_detection_job_test.rb create mode 100644 test/support/fake_change_detecting_trigger.rb diff --git a/AGENTS.md b/AGENTS.md index 04b3ffe..596a646 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -11,10 +11,12 @@ This Rails app uses a small set of preferred libraries for common integration wo ## Codebase Map -- `lib/r3x/`: core framework code for the workflow DSL, trigger types, workflow loading, registry, execution context, and recurring-task config. +- `lib/r3x/`: core framework code for the workflow DSL, trigger types, workflow loading, registry, execution context, recurring-task config, and shared DSL helpers. - `lib/r3x/dsl/`: shared DSL infrastructure, especially validation concerns and configuration errors used by workflow-declared objects. -- `app/lib/r3x/`: runtime support code such as outputs, service clients, and shared concerns. -- `app/jobs/r3x/`: job entrypoints, especially `R3x::RunWorkflowJob`, which resolves and executes workflows. +- `lib/r3x/trigger_collection.rb`: internal collection class that manages workflow triggers as a hash keyed by `unique_key`. +- `app/lib/r3x/`: runtime support code such as outputs, client wrappers, and shared concerns. +- `app/jobs/r3x/`: job entrypoints, especially `R3x::RunWorkflowJob`, which resolves and executes workflows, and `R3x::ChangeDetectionJob`, which evaluates change-detecting triggers before enqueueing workflow runs. +- `app/models/r3x/`: runtime support models such as `R3x::TriggerState` for per-trigger change-detection state. - `workflows/`: user workflow packs. These are not the framework itself; they are loaded by the framework. - `config/initializers/r3x_workflow_loader.rb`: boot-time workflow loading hook. - `test/fixtures/workflows/`: fixture workflows for framework tests. Prefer these over hardcoding real workflows in tests. @@ -24,8 +26,10 @@ This Rails app uses a small set of preferred libraries for common integration wo - Workflows subclass `R3x::Workflow`, declare triggers via the DSL, and implement `#run(ctx)`. - Workflow-declared DSL objects must validate themselves before being registered; invalid DSL configuration should raise `R3x::ConfigurationError` with collected validation errors. - `R3x::WorkflowPackLoader` discovers `workflow.rb` entrypoints from directories listed in `R3X_WORKFLOW_PATHS`, loads them, and registers their classes in `R3x::WorkflowRegistry`. -- `R3x::RecurringTasksConfig` turns schedulable workflow triggers into Solid Queue recurring-task definitions. -- `R3x::RunWorkflowJob` fetches the workflow from the registry, builds a `WorkflowContext`, and calls `workflow_class.new.run(ctx)`. +- `R3x::RecurringTasksConfig` turns schedulable workflow triggers into Solid Queue recurring-task definitions. All triggers have a `unique_key` (based on type + options hash) used for identification and duplicate detection. +- Change-detecting triggers are file-defined trigger objects that provide `cron`, `unique_key`, and `detect_changes(workflow_key:, state:)`. Their durable runtime state lives in `R3x::TriggerState`. +- `R3x::ChangeDetectionJob` loads the trigger, fetches/updates `R3x::TriggerState`, and only enqueues `R3x::RunWorkflowJob` when the trigger reports a change. +- `R3x::RunWorkflowJob` fetches the workflow from the registry, resolves the trigger by `trigger_key`, builds a `WorkflowContext`, and calls `workflow_class.new.run(ctx)`. - Trigger discovery is filesystem-backed through `lib/r3x/triggers/*.rb`, so trigger file names, constants, and supported types must stay aligned. ## Maintenance Warning @@ -108,16 +112,19 @@ This Rails app uses a small set of preferred libraries for common integration wo ## Control Flow - `case` statements that dispatch on configuration values (e.g., ENV modes) must either exhaustively list all supported values or raise an exception in the `else` branch for unsupported values. -- **Good**: +- **Good**: + ```ruby case mode when "real" then # handle real - when "test" then # handle test + when "test" then # handle test else raise ArgumentError, "Unsupported mode: #{mode}" end ``` + - **Bad**: + ```ruby case mode when "real" then # handle real @@ -125,6 +132,7 @@ This Rails app uses a small set of preferred libraries for common integration wo # silently assumes test mode, hides typos in configuration end ``` + - Reasoning: Failing fast with a clear error prevents silent misconfiguration and makes debugging easier when an invalid mode is accidentally provided. ## Environment Variables diff --git a/Gemfile b/Gemfile index 5d3a0c0..260d6fc 100644 --- a/Gemfile +++ b/Gemfile @@ -4,6 +4,8 @@ source "https://rubygems.org" gem "rails", "~> 8.1.2" # Use sqlite3 as the database for Active Record gem "sqlite3", ">= 2.1" +# Use PostgreSQL as an alternative database adapter +gem "pg", "~> 1.1" # Use the Puma web server [https://github.com/puma/puma] gem "puma", ">= 5.0" diff --git a/Gemfile.lock b/Gemfile.lock index b8fd65a..2b434f0 100644 --- a/Gemfile.lock +++ b/Gemfile.lock @@ -260,7 +260,7 @@ GEM language_server-protocol (3.17.0.5) lint_roller (1.1.0) logger (1.7.0) - loofah (2.25.0) + loofah (2.25.1) crass (~> 1.0.2) nokogiri (>= 1.12.0) mail (2.9.0) @@ -322,6 +322,13 @@ GEM parser (3.3.10.2) ast (~> 2.4.1) racc + pg (1.6.3) + pg (1.6.3-aarch64-linux) + pg (1.6.3-aarch64-linux-musl) + pg (1.6.3-arm64-darwin) + pg (1.6.3-x86_64-darwin) + pg (1.6.3-x86_64-linux) + pg (1.6.3-x86_64-linux-musl) pp (0.6.3) prettyprint prettyprint (0.2.0) @@ -509,6 +516,7 @@ DEPENDENCIES mission_control-jobs multi_json nokogiri + pg (~> 1.1) propshaft puma (>= 5.0) rails (~> 8.1.2) @@ -603,7 +611,7 @@ CHECKSUMS language_server-protocol (3.17.0.5) sha256=fd1e39a51a28bf3eec959379985a72e296e9f9acfce46f6a79d31ca8760803cc lint_roller (1.1.0) sha256=2c0c845b632a7d172cb849cc90c1bce937a28c5c8ccccb50dfd46a485003cc87 logger (1.7.0) sha256=196edec7cc44b66cfb40f9755ce11b392f21f7967696af15d274dde7edff0203 - loofah (2.25.0) sha256=df5ed7ac3bac6a4ec802df3877ee5cc86d027299f8952e6243b3dac446b060e6 + loofah (2.25.1) sha256=d436c73dbd0c1147b16c4a41db097942d217303e1f7728704b37e4df9f6d2e04 mail (2.9.0) sha256=6fa6673ecd71c60c2d996260f9ee3dd387d4673b8169b502134659ece6d34941 marcel (1.1.0) sha256=fdcfcfa33cc52e93c4308d40e4090a5d4ea279e160a7f6af988260fa970e0bee mcp (0.8.0) sha256=ae8bd146bb8e168852866fd26f805f52744f6326afb3211e073f78a95e0c34fb @@ -630,6 +638,13 @@ CHECKSUMS os (1.1.4) sha256=57816d6a334e7bd6aed048f4b0308226c5fb027433b67d90a9ab435f35108d3f parallel (1.27.0) sha256=4ac151e1806b755fb4e2dc2332cbf0e54f2e24ba821ff2d3dcf86bf6dc4ae130 parser (3.3.10.2) sha256=6f60c84aa4bdcedb6d1a2434b738fe8a8136807b6adc8f7f53b97da9bc4e9357 + pg (1.6.3) sha256=1388d0563e13d2758c1089e35e973a3249e955c659592d10e5b77c468f628a99 + pg (1.6.3-aarch64-linux) sha256=0698ad563e02383c27510b76bf7d4cd2de19cd1d16a5013f375dd473e4be72ea + pg (1.6.3-aarch64-linux-musl) sha256=06a75f4ea04b05140146f2a10550b8e0d9f006a79cdaf8b5b130cde40e3ecc2c + pg (1.6.3-arm64-darwin) sha256=7240330b572e6355d7c75a7de535edb5dfcbd6295d9c7777df4d9dddfb8c0e5f + pg (1.6.3-x86_64-darwin) sha256=ee2e04a17c0627225054ffeb43e31a95be9d7e93abda2737ea3ce4a62f2729d6 + pg (1.6.3-x86_64-linux) sha256=5d9e188c8f7a0295d162b7b88a768d8452a899977d44f3274d1946d67920ae8d + pg (1.6.3-x86_64-linux-musl) sha256=9c9c90d98c72f78eb04c0f55e9618fe55d1512128e411035fe229ff427864009 pp (0.6.3) sha256=2951d514450b93ccfeb1df7d021cae0da16e0a7f95ee1e2273719669d0ab9df6 prettyprint (0.2.0) sha256=2bc9e15581a94742064a3cc8b0fb9d45aae3d03a1baa6ef80922627a0766f193 prism (1.9.0) sha256=7b530c6a9f92c24300014919c9dcbc055bf4cdf51ec30aed099b06cd6674ef85 diff --git a/app/jobs/r3x/change_detection_job.rb b/app/jobs/r3x/change_detection_job.rb new file mode 100644 index 0000000..818e350 --- /dev/null +++ b/app/jobs/r3x/change_detection_job.rb @@ -0,0 +1,69 @@ +module R3x + class ChangeDetectionJob < ApplicationJob + queue_as :default + + def perform(workflow_key, trigger_key:) + R3x::WorkflowPackLoader.load! + workflow_class = R3x::WorkflowRegistry.fetch(workflow_key) + trigger = find_trigger(workflow_class: workflow_class, trigger_key: trigger_key) + trigger_state = load_trigger_state(workflow_key: workflow_key, trigger_key: trigger_key, trigger_type: trigger.type) + result = normalize_result( + trigger.detect_changes( + workflow_key: workflow_key, + state: trigger_state.state.deep_symbolize_keys + ) + ) + + trigger_state.record_check!(result) + + return unless result[:changed] + + R3x::RunWorkflowJob.perform_later( + workflow_key, + trigger_key: trigger_key, + trigger_payload: result[:payload] + ) + rescue StandardError => e + trigger_state.record_error!(e) if defined?(trigger_state) && trigger_state&.persisted? + raise + end + + private + + def find_trigger(workflow_class:, trigger_key:) + trigger = workflow_class.triggers_by_key[trigger_key] + + if trigger.nil? + raise ArgumentError, "Unknown trigger key '#{trigger_key}' for workflow '#{workflow_class.workflow_key}'" + end + + unless trigger.change_detecting? + raise ArgumentError, "Trigger '#{trigger_key}' is not change-detecting" + end + + trigger + end + + def load_trigger_state(workflow_key:, trigger_key:, trigger_type:) + TriggerState.find_or_create_by!( + workflow_key: workflow_key, + trigger_key: trigger_key + ) do |state| + state.trigger_type = trigger_type.to_s + state.state = {} + end + end + + def normalize_result(result) + normalized = result.deep_symbolize_keys + + unless normalized.key?(:changed) && normalized.key?(:state) + raise ArgumentError, "Change-detecting trigger must return a hash with :changed and :state" + end + + normalized[:state] ||= {} + normalized[:payload] ||= nil + normalized + end + end +end diff --git a/app/jobs/r3x/run_workflow_job.rb b/app/jobs/r3x/run_workflow_job.rb index 1c25339..2575760 100644 --- a/app/jobs/r3x/run_workflow_job.rb +++ b/app/jobs/r3x/run_workflow_job.rb @@ -2,16 +2,15 @@ module R3x class RunWorkflowJob < ApplicationJob queue_as :default - def perform(workflow_key, trigger_type: "manual") + def perform(workflow_key, trigger_key:, trigger_payload: nil) R3x::WorkflowPackLoader.load! workflow_class = R3x::WorkflowRegistry.fetch(workflow_key) - - trigger = workflow_class.triggers.find { |t| t.type.to_s == trigger_type } - trigger ||= Triggers::Manual.new + trigger = find_trigger(workflow_class: workflow_class, trigger_key: trigger_key) execution = TriggerExecution.new( trigger: trigger, - workflow_key: workflow_key + workflow_key: workflow_key, + payload: trigger_payload ) ctx = WorkflowContext.new( @@ -20,5 +19,17 @@ def perform(workflow_key, trigger_type: "manual") ) workflow_class.new.run(ctx) end + + private + + def find_trigger(workflow_class:, trigger_key:) + trigger = workflow_class.triggers_by_key[trigger_key] + + if trigger.nil? + raise ArgumentError, "Unknown trigger key '#{trigger_key}' for workflow '#{workflow_class.workflow_key}'" + end + + trigger + end end end diff --git a/app/lib/r3x/client/http.rb b/app/lib/r3x/client/http.rb index 1b57116..173fabf 100644 --- a/app/lib/r3x/client/http.rb +++ b/app/lib/r3x/client/http.rb @@ -6,7 +6,11 @@ def initialize end def get(url) - connection.get(url).body + connection.get(url) + end + + def post(url, payload) + connection.post(url, payload) end private diff --git a/app/models/r3x/trigger_state.rb b/app/models/r3x/trigger_state.rb new file mode 100644 index 0000000..4e94b3e --- /dev/null +++ b/app/models/r3x/trigger_state.rb @@ -0,0 +1,26 @@ +module R3x + class TriggerState < ApplicationRecord + serialize :state, coder: MultiJson if connection.adapter_name.downcase == "sqlite" + + validates :workflow_key, presence: true + validates :trigger_type, presence: true + validates :trigger_key, presence: true, uniqueness: { scope: :workflow_key } + + def record_check!(result) + update!( + state: result.fetch(:state), + last_checked_at: Time.current, + last_error_at: nil, + last_error_message: nil, + last_triggered_at: result[:changed] ? Time.current : last_triggered_at + ) + end + + def record_error!(error) + update!( + last_error_at: Time.current, + last_error_message: error.message + ) + end + end +end diff --git a/db/migrate/20260318100000_create_trigger_states.rb b/db/migrate/20260318100000_create_trigger_states.rb new file mode 100644 index 0000000..835a098 --- /dev/null +++ b/db/migrate/20260318100000_create_trigger_states.rb @@ -0,0 +1,20 @@ +class CreateTriggerStates < ActiveRecord::Migration[8.1] + def change + create_table :trigger_states do |t| + t.string :workflow_key, null: false + t.string :trigger_key, null: false + t.string :trigger_type, null: false + t.json :state, null: false, default: {} + t.datetime :last_checked_at + t.datetime :last_triggered_at + t.datetime :last_error_at + t.text :last_error_message + + t.timestamps + end + + add_index :trigger_states, [ :workflow_key, :trigger_key ], unique: true + add_index :trigger_states, :workflow_key + add_index :trigger_states, :trigger_type + end +end diff --git a/db/schema.rb b/db/schema.rb index d083e34..b05276a 100644 --- a/db/schema.rb +++ b/db/schema.rb @@ -10,7 +10,7 @@ # # It's strongly recommended that you check this file into your version control system. -ActiveRecord::Schema[8.1].define(version: 2026_03_15_120100) do +ActiveRecord::Schema[8.1].define(version: 2026_03_18_100000) do create_table "solid_queue_blocked_executions", force: :cascade do |t| t.string "concurrency_key", null: false t.datetime "created_at", null: false @@ -132,6 +132,22 @@ t.index ["key"], name: "index_solid_queue_semaphores_on_key", unique: true end + create_table "trigger_states", force: :cascade do |t| + t.datetime "created_at", null: false + t.datetime "last_checked_at" + t.datetime "last_error_at" + t.text "last_error_message" + t.datetime "last_triggered_at" + t.json "state", default: {}, null: false + t.string "trigger_key", null: false + t.string "trigger_type", null: false + t.datetime "updated_at", null: false + t.string "workflow_key", null: false + t.index ["trigger_type"], name: "index_trigger_states_on_trigger_type" + t.index ["workflow_key", "trigger_key"], name: "index_trigger_states_on_workflow_key_and_trigger_key", unique: true + t.index ["workflow_key"], name: "index_trigger_states_on_workflow_key" + end + add_foreign_key "solid_queue_blocked_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_claimed_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade add_foreign_key "solid_queue_failed_executions", "solid_queue_jobs", column: "job_id", on_delete: :cascade diff --git a/lib/r3x/recurring_tasks_config.rb b/lib/r3x/recurring_tasks_config.rb index 3f75acb..720b610 100644 --- a/lib/r3x/recurring_tasks_config.rb +++ b/lib/r3x/recurring_tasks_config.rb @@ -11,17 +11,27 @@ def to_h workflow_key = workflow_class.workflow_key triggers.each do |trigger| - result[workflow_key] = { - "class" => "R3x::RunWorkflowJob", - "args" => [ workflow_key, { "triggered_by" => trigger.type.to_s } ], - "schedule" => trigger.cron, - "queue" => "default" - } + result_key = "#{workflow_key}:#{trigger.unique_key}" + result[result_key] = task_definition_for( + workflow_key: workflow_key, + trigger: trigger + ) end end result end + + private + + def task_definition_for(workflow_key:, trigger:) + { + "class" => trigger.change_detecting? ? "R3x::ChangeDetectionJob" : "R3x::RunWorkflowJob", + "args" => [ workflow_key, { "trigger_key" => trigger.unique_key } ], + "schedule" => trigger.cron, + "queue" => "default" + } + end end end end diff --git a/lib/r3x/trigger_collection.rb b/lib/r3x/trigger_collection.rb new file mode 100644 index 0000000..09d33da --- /dev/null +++ b/lib/r3x/trigger_collection.rb @@ -0,0 +1,35 @@ +module R3x + class TriggerCollection + include Enumerable + + def initialize + @by_key = {} + end + + def add(trigger) + key = trigger.unique_key + raise ArgumentError, "Trigger with key '#{key}' already exists" if @by_key.key?(key) + @by_key[key] = trigger + end + + def each(&block) + @by_key.values.each(&block) + end + + def by_key + @by_key.dup + end + + def select(&block) + @by_key.values.select(&block) + end + + def to_a + @by_key.values + end + + def size + @by_key.size + end + end +end diff --git a/lib/r3x/trigger_execution.rb b/lib/r3x/trigger_execution.rb index af0ca12..1a49455 100644 --- a/lib/r3x/trigger_execution.rb +++ b/lib/r3x/trigger_execution.rb @@ -1,10 +1,11 @@ module R3x class TriggerExecution - attr_reader :trigger, :workflow_key + attr_reader :trigger, :workflow_key, :payload - def initialize(trigger:, workflow_key:) + def initialize(trigger:, workflow_key:, payload: nil) @trigger = trigger @workflow_key = workflow_key + @payload = payload end def type diff --git a/lib/r3x/triggers/base.rb b/lib/r3x/triggers/base.rb index 0c7feb3..901f2ed 100644 --- a/lib/r3x/triggers/base.rb +++ b/lib/r3x/triggers/base.rb @@ -10,12 +10,18 @@ def initialize(type, **options) @options = options end - def to_h - options.dup + def change_detecting? + false end - def validation_subject - "trigger :#{type}" + def cron_schedulable? + false + end + + def unique_key + # Generate unique key from type + sorted options hash + key_json = MultiJson.dump(options.sort.to_h) + "#{type}:#{Digest::SHA256.hexdigest(key_json)[0..15]}" end end end diff --git a/lib/r3x/triggers/concerns/change_detecting.rb b/lib/r3x/triggers/concerns/change_detecting.rb new file mode 100644 index 0000000..13fb95b --- /dev/null +++ b/lib/r3x/triggers/concerns/change_detecting.rb @@ -0,0 +1,15 @@ +module R3x + module Triggers + module Concerns + module ChangeDetecting + def change_detecting? + true + end + + def detect_changes(workflow_key:, state:) + raise NotImplementedError, "#{self.class.name} must implement #detect_changes" + end + end + end + end +end diff --git a/lib/r3x/workflow.rb b/lib/r3x/workflow.rb index 294ea43..8fa33e1 100644 --- a/lib/r3x/workflow.rb +++ b/lib/r3x/workflow.rb @@ -2,7 +2,8 @@ module R3x class Workflow class << self def inherited(subclass) - subclass.instance_variable_set(:@_triggers, []) + super + subclass._triggers = TriggerCollection.new end def workflow_key @@ -10,21 +11,24 @@ def workflow_key end def trigger(type, **options) - trigger_class = Triggers.resolve(type) - trigger_instance = trigger_class.new(**options) - + trigger_instance = Triggers.resolve(type).new(**options) trigger_instance.validate!(message_prefix: "Invalid trigger :#{type} for #{name}") - @_triggers << trigger_instance + _triggers.add(trigger_instance) end def triggers - @_triggers ||= [] - @_triggers.dup + _triggers.to_a end def schedulable_triggers - triggers.select { |t| t.respond_to?(:cron_schedulable?) && t.cron_schedulable? } + _triggers.select(&:cron_schedulable?) + end + + def triggers_by_key + _triggers.by_key end + + attr_accessor :_triggers end def run(ctx) diff --git a/test/jobs/r3x/change_detection_job_test.rb b/test/jobs/r3x/change_detection_job_test.rb new file mode 100644 index 0000000..ad9d56a --- /dev/null +++ b/test/jobs/r3x/change_detection_job_test.rb @@ -0,0 +1,115 @@ +require "test_helper" +require_relative "../../support/fake_change_detecting_trigger" + +module R3x + class ChangeDetectionJobTest < ActiveSupport::TestCase + include ActiveJob::TestHelper + + setup do + @original_workflow_paths = ENV["R3X_WORKFLOW_PATHS"] + ENV["R3X_WORKFLOW_PATHS"] = Rails.root.join("test/fixtures/workflows").to_s + WorkflowPackLoader.load!(force: true) + clear_enqueued_jobs + R3x::TriggerState.delete_all + end + + teardown do + clear_enqueued_jobs + R3x::TriggerState.delete_all + ENV["R3X_WORKFLOW_PATHS"] = @original_workflow_paths + WorkflowRegistry.reset! + WorkflowPackLoader.load!(force: true) + end + + test "creates trigger state and does not enqueue workflow when unchanged" do + fake_trigger = R3x::TestSupport::FakeChangeDetectingTrigger.new( + identity: "feed", + detector: ->(workflow_key:, state:) do + assert_equal "test_change_detecting_feed", workflow_key + assert_equal({}, state) + { changed: false, state: { cursor: "v1" }, payload: nil } + end + ) + + register_change_detecting_workflow(fake_trigger) + + assert_no_enqueued_jobs do + ChangeDetectionJob.perform_now("test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + end + + state = R3x::TriggerState.find_by!(workflow_key: "test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + assert_equal "fake_change_detecting", state.trigger_type + assert_equal({ "cursor" => "v1" }, state.state) + assert state.last_checked_at.present? + assert_nil state.last_triggered_at + assert_nil state.last_error_at + end + + test "enqueues workflow and records last_triggered_at when changed" do + fake_trigger = R3x::TestSupport::FakeChangeDetectingTrigger.new( + identity: "feed", + detector: ->(workflow_key:, state:) do + { changed: true, state: state.merge(cursor: "v2"), payload: { "entries" => [ { "title" => "Hello" } ] } } + end + ) + + register_change_detecting_workflow(fake_trigger) + + assert_enqueued_jobs 1, only: R3x::RunWorkflowJob do + ChangeDetectionJob.perform_now("test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + end + + enqueued_job = enqueued_jobs.last + assert_equal R3x::RunWorkflowJob, enqueued_job[:job] + assert_equal "test_change_detecting_feed", enqueued_job[:args][0] + assert_equal fake_trigger.unique_key, enqueued_job[:args][1]["trigger_key"] + assert_equal( + { "entries" => [ { "title" => "Hello", "_aj_symbol_keys" => [ "title" ] } ], "_aj_symbol_keys" => [ "entries" ] }, + enqueued_job[:args][1]["trigger_payload"] + ) + + state = R3x::TriggerState.find_by!(workflow_key: "test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + assert_equal({ "cursor" => "v2" }, state.state) + assert state.last_triggered_at.present? + end + + test "persists last error details when detection fails" do + fake_trigger = R3x::TestSupport::FakeChangeDetectingTrigger.new( + identity: "feed", + detector: ->(workflow_key:, state:) do + raise ArgumentError, "detection failed" + end + ) + + register_change_detecting_workflow(fake_trigger) + + error = assert_raises(ArgumentError) do + ChangeDetectionJob.perform_now("test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + end + + assert_equal "detection failed", error.message + + state = R3x::TriggerState.find_by!(workflow_key: "test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + assert state.last_error_at.present? + assert_equal "detection failed", state.last_error_message + end + + private + + def register_change_detecting_workflow(fake_trigger) + workflow_class = Class.new(R3x::Workflow) do + def self.name + "TestChangeDetectingFeed" + end + + define_singleton_method(:triggers_by_key) { { fake_trigger.unique_key => fake_trigger } } + + def run(ctx) + ctx.trigger.payload + end + end + + WorkflowRegistry.register(workflow_class) + end + end +end diff --git a/test/jobs/r3x/run_workflow_job_test.rb b/test/jobs/r3x/run_workflow_job_test.rb index e85d08b..8e62638 100644 --- a/test/jobs/r3x/run_workflow_job_test.rb +++ b/test/jobs/r3x/run_workflow_job_test.rb @@ -1,10 +1,12 @@ require "test_helper" +require_relative "../../support/fake_change_detecting_trigger" module R3x class RunWorkflowJobTest < ActiveSupport::TestCase setup do @original_workflow_paths = ENV["R3X_WORKFLOW_PATHS"] ENV["R3X_WORKFLOW_PATHS"] = Rails.root.join("test/fixtures/workflows").to_s + WorkflowPackLoader.load!(force: true) end @@ -12,23 +14,42 @@ class RunWorkflowJobTest < ActiveSupport::TestCase ENV["R3X_WORKFLOW_PATHS"] = @original_workflow_paths end - test "performs workflow with manual trigger when no triggered_by provided" do + test "performs workflow with manual trigger" do job = RunWorkflowJob.new - # The test_workflow is already loaded from fixtures and has a simple run method - result = job.perform("test_workflow") + test_workflow_class = Class.new(R3x::Workflow) do + def self.name + "TestManual" + end + + trigger :manual + + def run(ctx) + { + "trigger_type" => ctx.trigger.type.to_s, + "schedule?" => ctx.trigger.schedule? + } + end + end + + WorkflowRegistry.register(test_workflow_class) + manual_trigger = test_workflow_class.triggers.first + + result = job.perform("test_manual", trigger_key: manual_trigger.unique_key) - assert_equal true, result["test"] - assert_equal "Test workflow executed successfully", result["message"] + assert_equal "manual", result["trigger_type"] + refute result["schedule?"] + ensure + WorkflowRegistry.reset! + WorkflowPackLoader.load!(force: true) end test "performs workflow with schedule trigger" do job = RunWorkflowJob.new - # Create a test workflow class that checks trigger test_workflow_class = Class.new(R3x::Workflow) do def self.name - "TestTriggerType" + "TestSchedule" end trigger :schedule, cron: "0 * * * *" @@ -42,8 +63,9 @@ def run(ctx) end WorkflowRegistry.register(test_workflow_class) + schedule_trigger = test_workflow_class.triggers.first - result = job.perform("test_trigger_type", trigger_type: "schedule") + result = job.perform("test_schedule", trigger_key: schedule_trigger.unique_key) assert_equal "schedule", result["trigger_type"] assert result["schedule?"] @@ -52,30 +74,35 @@ def run(ctx) WorkflowPackLoader.load!(force: true) end - test "performs workflow with manual trigger explicitly" do + test "performs workflow with change-detecting trigger and payload" do job = RunWorkflowJob.new + fake_trigger = R3x::TestSupport::FakeChangeDetectingTrigger.new(identity: "feed") - test_workflow_class = Class.new(R3x::Workflow) do + workflow_class = Class.new(R3x::Workflow) do def self.name - "TestManual" + "TestChangeDetecting" end - trigger :manual + define_singleton_method(:triggers_by_key) { { fake_trigger.unique_key => fake_trigger } } def run(ctx) { "trigger_type" => ctx.trigger.type.to_s, - "schedule?" => ctx.trigger.schedule? + "payload" => ctx.trigger.payload } end end - WorkflowRegistry.register(test_workflow_class) + WorkflowRegistry.register(workflow_class) - result = job.perform("test_manual", trigger_type: "manual") + result = job.perform( + "test_change_detecting", + trigger_key: fake_trigger.unique_key, + trigger_payload: { "entries" => [ { "title" => "Hello" } ] } + ) - assert_equal "manual", result["trigger_type"] - refute result["schedule?"] + assert_equal "fake_change_detecting", result["trigger_type"] + assert_equal({ "entries" => [ { "title" => "Hello" } ] }, result["payload"]) ensure WorkflowRegistry.reset! WorkflowPackLoader.load!(force: true) diff --git a/test/lib/r3x/recurring_tasks_config_test.rb b/test/lib/r3x/recurring_tasks_config_test.rb index d611962..6bdfc64 100644 --- a/test/lib/r3x/recurring_tasks_config_test.rb +++ b/test/lib/r3x/recurring_tasks_config_test.rb @@ -1,4 +1,5 @@ require "test_helper" +require_relative "../../support/fake_change_detecting_trigger" module R3x class RecurringTasksConfigTest < ActiveSupport::TestCase @@ -15,11 +16,11 @@ class RecurringTasksConfigTest < ActiveSupport::TestCase test "generates recurring tasks from workflow DSL" do tasks = RecurringTasksConfig.to_h - assert tasks.key?("test_workflow") - task = tasks["test_workflow"] + expected_key = "test_workflow:schedule:" + task = tasks.find { |k, _| k.start_with?(expected_key) }&.last + assert task, "Expected task with key starting with #{expected_key}" assert_equal "R3x::RunWorkflowJob", task["class"] - assert_equal [ "test_workflow", { "triggered_by" => "schedule" } ], task["args"] assert_equal "0 * * * *", task["schedule"] assert_equal "default", task["queue"] end @@ -41,5 +42,31 @@ def self.name WorkflowRegistry.reset! WorkflowPackLoader.load!(force: true) end + + test "generates change detection tasks for change-detecting triggers" do + fake_trigger = R3x::TestSupport::FakeChangeDetectingTrigger.new(identity: "feed") + workflow_class = Class.new(R3x::Workflow) do + def self.name + "Workflows::ChangeDetectingFeed" + end + + define_singleton_method(:triggers) { [ fake_trigger ] } + define_singleton_method(:schedulable_triggers) { [ fake_trigger ] } + end + + WorkflowRegistry.register(workflow_class) + + tasks = RecurringTasksConfig.to_h + expected_key = "change_detecting_feed:#{fake_trigger.unique_key}" + task = tasks.fetch(expected_key) + + assert_equal "R3x::ChangeDetectionJob", task["class"] + assert_equal [ "change_detecting_feed", { "trigger_key" => fake_trigger.unique_key } ], task["args"] + assert_equal "every 15 minutes", task["schedule"] + assert_equal "default", task["queue"] + ensure + WorkflowRegistry.reset! + WorkflowPackLoader.load!(force: true) + end end end diff --git a/test/lib/r3x/workflow_context_test.rb b/test/lib/r3x/workflow_context_test.rb index 08f4064..c84a318 100644 --- a/test/lib/r3x/workflow_context_test.rb +++ b/test/lib/r3x/workflow_context_test.rb @@ -27,6 +27,17 @@ class TriggerExecutionTest < ActiveSupport::TestCase execution = TriggerExecution.new(trigger: trigger, workflow_key: "test") assert_equal({ cron: "0 13 * * *" }, execution.options) end + + test "exposes runtime payload" do + trigger = R3x::Triggers::Manual.new + execution = TriggerExecution.new( + trigger: trigger, + workflow_key: "test", + payload: { "entries" => [ { "title" => "Hello" } ] } + ) + + assert_equal({ "entries" => [ { "title" => "Hello" } ] }, execution.payload) + end end class WorkflowExecutionTest < ActiveSupport::TestCase diff --git a/test/lib/r3x/workflow_test.rb b/test/lib/r3x/workflow_test.rb index 7bd38fe..6483be5 100644 --- a/test/lib/r3x/workflow_test.rb +++ b/test/lib/r3x/workflow_test.rb @@ -1,4 +1,5 @@ require "test_helper" +require_relative "../../support/fake_change_detecting_trigger" module R3x class WorkflowTest < ActiveSupport::TestCase @@ -100,30 +101,22 @@ def self.name assert_equal [ :schedule ], triggers.map(&:type) end - test "trigger :schedule rejects blank cron" do - error = assert_raises(ConfigurationError) do - Class.new(R3x::Workflow) do - def self.name - "Test" + test "trigger :schedule rejects blank cron (empty string and whitespace)" do + [ + "", + " " + ].each do |cron| + error = assert_raises(ConfigurationError) do + Class.new(R3x::Workflow) do + def self.name + "Test" + end + trigger :schedule, cron: cron end - trigger :schedule, cron: "" end - end - - assert_includes error.message, "Cron can't be blank" - end - test "trigger :schedule rejects whitespace-only cron" do - error = assert_raises(ConfigurationError) do - Class.new(R3x::Workflow) do - def self.name - "Test" - end - trigger :schedule, cron: " " - end + assert_includes error.message, "Cron can't be blank" end - - assert_includes error.message, "Cron can't be blank" end test "supported_types returns list of available trigger files" do @@ -145,5 +138,37 @@ def self.name assert_match(/Unknown trigger type: nonexistent/, error.message) assert_match(/Supported types:.*:schedule/, error.message) end + + test "rejects duplicate change-detecting trigger keys in one workflow" do + original_resolve = R3x::Triggers.method(:resolve) + + R3x::Triggers.define_singleton_method(:resolve) do |_type| + R3x::TestSupport::FakeChangeDetectingTrigger + end + + error = begin + assert_raises(ArgumentError) do + Class.new(R3x::Workflow) do + def self.name + "Workflows::DuplicateChangeDetecting" + end + + trigger :fake_change_detecting, identity: "same" + trigger :fake_change_detecting, identity: "same", cron: "every hour" + end + end + end + + assert_match(/Trigger with key .* already exists/, error.message) + ensure + R3x::Triggers.define_singleton_method(:resolve, original_resolve) + end + + test "change-detecting trigger key does not change when only cron changes" do + trigger_one = R3x::TestSupport::FakeChangeDetectingTrigger.new(identity: "feed", cron: "every 15 minutes") + trigger_two = R3x::TestSupport::FakeChangeDetectingTrigger.new(identity: "feed", cron: "every hour") + + assert_equal trigger_one.unique_key, trigger_two.unique_key + end end end diff --git a/test/support/fake_change_detecting_trigger.rb b/test/support/fake_change_detecting_trigger.rb new file mode 100644 index 0000000..9262838 --- /dev/null +++ b/test/support/fake_change_detecting_trigger.rb @@ -0,0 +1,30 @@ +module R3x + module TestSupport + class FakeChangeDetectingTrigger < R3x::Triggers::Base + include R3x::Triggers::Concerns::CronSchedulable + include R3x::Triggers::Concerns::ChangeDetecting + + def initialize(identity:, cron: "every 15 minutes", detector: nil, **options) + @detector = detector || ->(workflow_key:, state:) { { changed: false, state: state, payload: nil } } + super(:fake_change_detecting, identity: identity, cron: cron, **options) + end + + def validate!(**) + true + end + + def cron + options[:cron] + end + + def unique_key + # Identity-based key - doesn't change when cron changes + "fake_change_detecting:#{options[:identity]}" + end + + def detect_changes(workflow_key:, state:) + @detector.call(workflow_key:, state:) + end + end + end +end From add9c0b1e100534a767e0efe2ab8e51ab90c4294 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Ciemi=C4=99ga?= Date: Thu, 19 Mar 2026 21:03:13 +0000 Subject: [PATCH 2/4] Update app/models/r3x/trigger_state.rb Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- app/models/r3x/trigger_state.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/models/r3x/trigger_state.rb b/app/models/r3x/trigger_state.rb index 4e94b3e..9c3f623 100644 --- a/app/models/r3x/trigger_state.rb +++ b/app/models/r3x/trigger_state.rb @@ -1,6 +1,6 @@ module R3x class TriggerState < ApplicationRecord - serialize :state, coder: MultiJson if connection.adapter_name.downcase == "sqlite" + serialize :state, coder: MultiJson if ActiveRecord::Base.connection_db_config.adapter.to_s.downcase == "sqlite" validates :workflow_key, presence: true validates :trigger_type, presence: true From a8af286693ed75071ecb6733c966d89c92df59b8 Mon Sep 17 00:00:00 2001 From: zewelor Date: Thu, 19 Mar 2026 21:05:15 +0000 Subject: [PATCH 3/4] Remove HTTP/Discord helpers from WorkflowContext - Remove fetch_body, discord_output, and http_client methods. - Change made in `lib/r3x/workflow_context.rb`. - Reduce hidden dependencies and side effects in the context. - Encourage injection of HTTP clients/outputs for testability. --- lib/r3x/workflow_context.rb | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/lib/r3x/workflow_context.rb b/lib/r3x/workflow_context.rb index 1c52869..8aab33a 100644 --- a/lib/r3x/workflow_context.rb +++ b/lib/r3x/workflow_context.rb @@ -8,19 +8,5 @@ def initialize(trigger:, workflow_key:) @trigger = trigger @execution = WorkflowExecution.new(workflow_key: workflow_key) end - - def fetch_body(url) - http_client.get(url) - end - - def discord_output - @discord_output ||= R3x::Outputs::Discord.new - end - - private - - def http_client - @http_client ||= R3x::Client::Http.new - end end end From 714a8682ee25d507a7e3b2dc9050d02e37428df9 Mon Sep 17 00:00:00 2001 From: zewelor Date: Fri, 20 Mar 2026 07:39:20 +0000 Subject: [PATCH 4/4] Make enqueueing atomic for change detection - Make trigger state update and job enqueue atomic in ChangeDetectionJob - Add argument normalization to RunWorkflowJob and ChangeDetectionJob - Add test that simulates enqueue failure and preserves trigger state - Adjust test payload expectation to use symbolized keys - Update AGENTS.md to document Solid Queue transactional behavior --- AGENTS.md | 5 +++ app/jobs/r3x/change_detection_job.rb | 46 +++++++++++++++++----- app/jobs/r3x/run_workflow_job.rb | 32 ++++++++++++++- test/jobs/r3x/change_detection_job_test.rb | 36 ++++++++++++++++- test/jobs/r3x/run_workflow_job_test.rb | 2 +- 5 files changed, 108 insertions(+), 13 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 596a646..5ff4540 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -7,6 +7,9 @@ This Rails app uses a small set of preferred libraries for common integration wo - `r3x` is a Rails API app that acts as a Ruby-native workflow executor and automation engine. - The high-level split is: framework/runtime code lives in the app and `lib/r3x/`, while user-defined workflows live under `workflows/`. - Workflows are file-based, Git-friendly, and loaded into a database-backed runtime that uses Active Job + Solid Queue for execution and recurring scheduling. +- `Solid Queue` is the active job backend for app/runtime execution. Treat queueing semantics as database-backed, not Redis-backed. +- In the current app configuration, `Solid Queue` is not wired through `config.solid_queue.connects_to`, so queue records use the same Active Record database connection as the app in the environments configured here. That means queue inserts can participate in the same database transaction as app writes. +- If `Solid Queue` is ever moved to a separate database, or replaced with a non-database backend, revisit any code that relies on transactional integrity between app writes and job enqueueing. In that setup, `enqueue_after_transaction_commit` and related tests become important again. - The default local UI surface is Mission Control Jobs mounted at `/jobs`; the root route redirects there. ## Codebase Map @@ -29,6 +32,7 @@ This Rails app uses a small set of preferred libraries for common integration wo - `R3x::RecurringTasksConfig` turns schedulable workflow triggers into Solid Queue recurring-task definitions. All triggers have a `unique_key` (based on type + options hash) used for identification and duplicate detection. - Change-detecting triggers are file-defined trigger objects that provide `cron`, `unique_key`, and `detect_changes(workflow_key:, state:)`. Their durable runtime state lives in `R3x::TriggerState`. - `R3x::ChangeDetectionJob` loads the trigger, fetches/updates `R3x::TriggerState`, and only enqueues `R3x::RunWorkflowJob` when the trigger reports a change. +- Because the app currently uses `Solid Queue` as a database-backed backend on the same Active Record database connection, code may intentionally rely on a database transaction covering both `TriggerState` updates and `perform_later`. Do not assume those guarantees survive a future backend or database split. - `R3x::RunWorkflowJob` fetches the workflow from the registry, resolves the trigger by `trigger_key`, builds a `WorkflowContext`, and calls `workflow_class.new.run(ctx)`. - Trigger discovery is filesystem-backed through `lib/r3x/triggers/*.rb`, so trigger file names, constants, and supported types must stay aligned. @@ -37,6 +41,7 @@ This Rails app uses a small set of preferred libraries for common integration wo - Keep this file synchronized with the real codebase. If you change workflow loading, trigger discovery, scheduling flow, top-level directory structure, namespaces, or the framework/user-workflow boundary, update the relevant `AGENTS.md` sections in the same change. - In particular, update examples and notes here when changing files such as `lib/r3x/workflow.rb`, `lib/r3x/workflow_pack_loader.rb`, `lib/r3x/workflow_registry.rb`, `lib/r3x/recurring_tasks_config.rb`, `lib/r3x/triggers.rb`, `app/jobs/r3x/run_workflow_job.rb`, or `config/initializers/r3x_workflow_loader.rb`. - Also update this file when changing the shared DSL validation contract in files such as `lib/r3x/dsl/validatable.rb`, `lib/r3x/configuration_error.rb`, or the base classes for workflow-declared objects. +- Also update this file when changing Active Job backend semantics, `Solid Queue` database wiring, or any logic that depends on enqueueing being inside the same database transaction as app writes. - When adding a new subsystem or moving code between `lib/r3x/`, `app/lib/r3x/`, `app/jobs/r3x/`, or `workflows/`, refresh the project overview and codebase map so future agents can still orient themselves quickly. ## JSON diff --git a/app/jobs/r3x/change_detection_job.rb b/app/jobs/r3x/change_detection_job.rb index 818e350..c758aaa 100644 --- a/app/jobs/r3x/change_detection_job.rb +++ b/app/jobs/r3x/change_detection_job.rb @@ -2,7 +2,10 @@ module R3x class ChangeDetectionJob < ApplicationJob queue_as :default - def perform(workflow_key, trigger_key:) + def perform(workflow_key, options = nil) + workflow_key, options = normalize_arguments(workflow_key, options) + trigger_key = options.fetch(:trigger_key) + R3x::WorkflowPackLoader.load! workflow_class = R3x::WorkflowRegistry.fetch(workflow_key) trigger = find_trigger(workflow_class: workflow_class, trigger_key: trigger_key) @@ -14,15 +17,17 @@ def perform(workflow_key, trigger_key:) ) ) - trigger_state.record_check!(result) - - return unless result[:changed] + TriggerState.transaction do + if result[:changed] + R3x::RunWorkflowJob.perform_later( + workflow_key, + trigger_key: trigger_key, + trigger_payload: result[:payload] + ) + end - R3x::RunWorkflowJob.perform_later( - workflow_key, - trigger_key: trigger_key, - trigger_payload: result[:payload] - ) + trigger_state.record_check!(result) + end rescue StandardError => e trigger_state.record_error!(e) if defined?(trigger_state) && trigger_state&.persisted? raise @@ -54,6 +59,29 @@ def load_trigger_state(workflow_key:, trigger_key:, trigger_type:) end end + def normalize_arguments(workflow_key, options) + if workflow_key.is_a?(Hash) && options.nil? + options = workflow_key + workflow_key = nil + end + + options = normalize_options_hash(options) + workflow_key ||= options[:workflow_key] + + [ workflow_key.presence || raise(ArgumentError, "Missing workflow_key"), options ] + end + + def normalize_options_hash(options) + case options + when nil + {} + when Hash + options.deep_symbolize_keys + else + raise ArgumentError, "Expected options hash, got #{options.class.name}" + end + end + def normalize_result(result) normalized = result.deep_symbolize_keys diff --git a/app/jobs/r3x/run_workflow_job.rb b/app/jobs/r3x/run_workflow_job.rb index 2575760..f6b4146 100644 --- a/app/jobs/r3x/run_workflow_job.rb +++ b/app/jobs/r3x/run_workflow_job.rb @@ -2,7 +2,11 @@ module R3x class RunWorkflowJob < ApplicationJob queue_as :default - def perform(workflow_key, trigger_key:, trigger_payload: nil) + def perform(workflow_key, options = nil) + workflow_key, options = normalize_arguments(workflow_key, options) + trigger_key = options.fetch(:trigger_key) + trigger_payload = options[:trigger_payload] + R3x::WorkflowPackLoader.load! workflow_class = R3x::WorkflowRegistry.fetch(workflow_key) trigger = find_trigger(workflow_class: workflow_class, trigger_key: trigger_key) @@ -17,11 +21,37 @@ def perform(workflow_key, trigger_key:, trigger_payload: nil) trigger: execution, workflow_key: workflow_key ) + workflow_class.new.run(ctx) end private + def normalize_arguments(workflow_key, options) + if workflow_key.is_a?(Hash) && options.nil? + options = workflow_key + workflow_key = nil + end + + options = normalize_options_hash(options) + workflow_key ||= options[:workflow_key] + + raise ArgumentError, "Missing workflow_key" if workflow_key.blank? + + [ workflow_key, options ] + end + + def normalize_options_hash(options) + case options + when nil + {} + when Hash + options.deep_symbolize_keys + else + raise ArgumentError, "Expected options hash, got #{options.class.name}" + end + end + def find_trigger(workflow_class:, trigger_key:) trigger = workflow_class.triggers_by_key[trigger_key] diff --git a/test/jobs/r3x/change_detection_job_test.rb b/test/jobs/r3x/change_detection_job_test.rb index ad9d56a..bc81619 100644 --- a/test/jobs/r3x/change_detection_job_test.rb +++ b/test/jobs/r3x/change_detection_job_test.rb @@ -34,7 +34,7 @@ class ChangeDetectionJobTest < ActiveSupport::TestCase register_change_detecting_workflow(fake_trigger) assert_no_enqueued_jobs do - ChangeDetectionJob.perform_now("test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + ChangeDetectionJob.perform_now("test_change_detecting_feed", { "trigger_key" => fake_trigger.unique_key }) end state = R3x::TriggerState.find_by!(workflow_key: "test_change_detecting_feed", trigger_key: fake_trigger.unique_key) @@ -56,7 +56,7 @@ class ChangeDetectionJobTest < ActiveSupport::TestCase register_change_detecting_workflow(fake_trigger) assert_enqueued_jobs 1, only: R3x::RunWorkflowJob do - ChangeDetectionJob.perform_now("test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + ChangeDetectionJob.perform_now("test_change_detecting_feed", { "trigger_key" => fake_trigger.unique_key }) end enqueued_job = enqueued_jobs.last @@ -73,6 +73,38 @@ class ChangeDetectionJobTest < ActiveSupport::TestCase assert state.last_triggered_at.present? end + test "does not advance trigger state when enqueue fails" do + fake_trigger = R3x::TestSupport::FakeChangeDetectingTrigger.new( + identity: "feed", + detector: ->(workflow_key:, state:) do + { changed: true, state: state.merge(cursor: "v2"), payload: { "entries" => [ { "title" => "Hello" } ] } } + end + ) + + register_change_detecting_workflow(fake_trigger) + + original_perform_later = R3x::RunWorkflowJob.method(:perform_later) + R3x::RunWorkflowJob.singleton_class.send(:define_method, :perform_later) do |*| + raise ActiveJob::EnqueueError, "enqueue failed" + end + + error = assert_raises(ActiveJob::EnqueueError) do + ChangeDetectionJob.perform_now("test_change_detecting_feed", { "trigger_key" => fake_trigger.unique_key }) + ensure + R3x::RunWorkflowJob.singleton_class.send(:define_method, :perform_later, original_perform_later) + end + + assert_equal "enqueue failed", error.message + + state = R3x::TriggerState.find_by!(workflow_key: "test_change_detecting_feed", trigger_key: fake_trigger.unique_key) + + assert_equal({}, state.state) + assert_nil state.last_checked_at + assert_nil state.last_triggered_at + assert state.last_error_at.present? + assert_equal "enqueue failed", state.last_error_message + end + test "persists last error details when detection fails" do fake_trigger = R3x::TestSupport::FakeChangeDetectingTrigger.new( identity: "feed", diff --git a/test/jobs/r3x/run_workflow_job_test.rb b/test/jobs/r3x/run_workflow_job_test.rb index 8e62638..5400105 100644 --- a/test/jobs/r3x/run_workflow_job_test.rb +++ b/test/jobs/r3x/run_workflow_job_test.rb @@ -102,7 +102,7 @@ def run(ctx) ) assert_equal "fake_change_detecting", result["trigger_type"] - assert_equal({ "entries" => [ { "title" => "Hello" } ] }, result["payload"]) + assert_equal({ entries: [ { title: "Hello" } ] }, result["payload"]) ensure WorkflowRegistry.reset! WorkflowPackLoader.load!(force: true)