diff --git a/lib/floe.rb b/lib/floe.rb index 6cfe559e..c5d90d34 100644 --- a/lib/floe.rb +++ b/lib/floe.rb @@ -21,6 +21,7 @@ require_relative "floe/workflow/choice_rule/data" require_relative "floe/workflow/context" require_relative "floe/workflow/item_processor" +require_relative "floe/workflow/item_reader" require_relative "floe/workflow/intrinsic_function" require_relative "floe/workflow/intrinsic_function/parser" require_relative "floe/workflow/intrinsic_function/transformer" diff --git a/lib/floe/workflow/item_reader.rb b/lib/floe/workflow/item_reader.rb new file mode 100644 index 00000000..d78b0a19 --- /dev/null +++ b/lib/floe/workflow/item_reader.rb @@ -0,0 +1,24 @@ +# frozen_string_literal: true + +module Floe + class Workflow + class ItemReader + include ValidationMixin + + attr_reader :name, :resource, :parameters, :reader_config, :max_items, :runner + + def initialize(payload, name) + @name = name + @resource = payload["Resource"] + @parameters = PayloadTemplate.new(payload["Parameters"]) if payload["Parameters"] + @reader_config = payload["ReaderConfig"] || {} + @max_items = reader_config["MaxItems"] + + missing_field_error!("Resource") unless @resource.kind_of?(String) + invalid_field_error!("ReaderConfig.MaxItems", @max_items, "must be positive") if @max_items && @max_items < 0 + + @runner = wrap_parser_error("Resource", @resource) { Floe::Runner.for_resource(@resource) } + end + end + end +end diff --git a/lib/floe/workflow/states/map.rb b/lib/floe/workflow/states/map.rb index e1af12cd..4c1618b4 100644 --- a/lib/floe/workflow/states/map.rb +++ b/lib/floe/workflow/states/map.rb @@ -30,7 +30,7 @@ def initialize(workflow, name, payload) @catch = payload["Catch"].to_a.map { |catcher| Catcher.new(catcher) } @item_processor = ItemProcessor.new(payload["ItemProcessor"], name) @items_path = ReferencePath.new(payload.fetch("ItemsPath", "$")) - @item_reader = payload["ItemReader"] + @item_reader = ItemReader.new(payload["ItemReader"], name + ["ItemReader"]) if payload["ItemReader"] @item_selector = payload["ItemSelector"] @item_batcher = payload["ItemBatcher"] @result_writer = payload["ResultWriter"] diff --git a/spec/workflow/item_reader_spec.rb b/spec/workflow/item_reader_spec.rb new file mode 100644 index 00000000..bb20a99c --- /dev/null +++ b/spec/workflow/item_reader_spec.rb @@ -0,0 +1,27 @@ +RSpec.describe Floe::Workflow::ItemReader do + let(:subject) { described_class.new(payload, ["Map", "ItemReader"]) } + + describe "#initialize" do + let(:payload) { {"Resource" => "docker://item_reader:latest"} } + + it "returns an ItemReader instance" do + expect(subject).to be_kind_of(described_class) + end + + context "Missing a \"Resource\" field" do + let(:payload) { {} } + + it "raises an exception" do + expect { subject }.to raise_error(Floe::InvalidWorkflowError, "Map.ItemReader does not have required field \"Resource\"") + end + end + + context "with an invalid ReaderConfig" do + let(:payload) { {"Resource" => "docker://item_reader:latest", "ReaderConfig" => {"MaxItems" => -1}} } + + it "raises an exception" do + expect { subject }.to raise_error(Floe::InvalidWorkflowError, "Map.ItemReader field \"ReaderConfig.MaxItems\" value \"-1\" must be positive") + end + end + end +end