diff --git a/AGENTS.md b/AGENTS.md index 04b3ffe4..5ff45406 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -7,14 +7,19 @@ 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 -- `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 +29,11 @@ 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. +- 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. ## Maintenance Warning @@ -33,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 @@ -108,16 +117,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 +137,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 5d3a0c07..260d6fc8 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 b8fd65a9..2b434f00 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 00000000..c758aaaf --- /dev/null +++ b/app/jobs/r3x/change_detection_job.rb @@ -0,0 +1,97 @@ +module R3x + class ChangeDetectionJob < ApplicationJob + queue_as :default + + 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) + 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 + ) + ) + + TriggerState.transaction do + if result[:changed] + R3x::RunWorkflowJob.perform_later( + workflow_key, + trigger_key: trigger_key, + trigger_payload: result[:payload] + ) + end + + trigger_state.record_check!(result) + end + 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_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 + + 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 1c25339b..f6b41467 100644 --- a/app/jobs/r3x/run_workflow_job.rb +++ b/app/jobs/r3x/run_workflow_job.rb @@ -2,23 +2,64 @@ module R3x class RunWorkflowJob < ApplicationJob queue_as :default - def perform(workflow_key, trigger_type: "manual") + 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 = 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( 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] + + 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 1b57116e..173fabfe 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 00000000..9c3f6236 --- /dev/null +++ b/app/models/r3x/trigger_state.rb @@ -0,0 +1,26 @@ +module R3x + class TriggerState < ApplicationRecord + 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 + 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 00000000..835a0989 --- /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 d083e34c..b05276a9 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 3f75acbb..720b610f 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 00000000..09d33da8 --- /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 af0ca123..1a494552 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 0c7feb38..901f2ede 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 00000000..13fb95b9 --- /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 294ea439..8fa33e14 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/lib/r3x/workflow_context.rb b/lib/r3x/workflow_context.rb index 1c528693..8aab33a5 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 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 00000000..bc81619d --- /dev/null +++ b/test/jobs/r3x/change_detection_job_test.rb @@ -0,0 +1,147 @@ +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 "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", + 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 e85d08be..5400105b 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 d611962b..6bdfc646 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 08f40643..c84a318b 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 7bd38fe0..6483be5c 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 00000000..92628384 --- /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