RSpec.shared_examples 'a swf-feed-ingestion flow' do |flow, date, timeout_seconds, licensor|
  ENVied.require

  def method_missing(method_name)
    raise NotImplementedError, "Please define a #{method_name} method for your caller flow."
  end

  def snowflake_conn
    @snowflake_conn ||= Sequel.odbc(ENVied.SNOWFLAKE_DSN, user: ENVied.SNOWFLAKE_USER, password: ENVied.SNOWFLAKE_PASSWORD)
  end

  def pg
    tries ||= 5
    @pg ||= PG.connect(host: ENVied.RS_HOST,
                       user: ENVied.RS_USER,
                       password: ENVied.RS_PASS,
                       port: ENVied.RS_PORT,
                       dbname: ENVied.RS_DBNAME)
  rescue PG::Error => e
    sleep 1
    retry unless (tries -= 1).zero?
    raise e
  end

  decider_pid, worker_pid = nil

  before(:all) do
    # truncate rows, aka clean up before company comes over
    truncate_rows

    # run prerequisites, if available
    setup_flow if defined?(setup_flow)

    # run the flow to completion
    # spawn a Decider
    puts '*** Starting Decider... ***'
    decider_pid = spawn("garcon decider #{flow}", chdir: "#{Dir.pwd}/../")

    # spawn a Worker
    puts '*** Starting Worker... ***'
    worker_pid = spawn("garcon worker #{flow}", chdir: "#{Dir.pwd}/../")

    # run our the setup script (copies sample data to FTP drop)
    puts '*** Running stage script... ***'
    `#{Dir.pwd}/stage_#{flow}.sh`
    raise "#{Dir.pwd}/stage_#{flow}.sh failed with exit code: #{$?.exitstatus}" unless $?.success?

    # have Executor kick off the flow
    puts '*** Firing up Executor... ***'
    executor_pid = spawn(
      "garcon exec #{flow} -o #{Dir.pwd}/#{flow}.out -c '{\"context_date\":\"#{date}\", \"reload\":\"True\", \"licensor\":\"#{licensor}\"}'")
    Process.wait executor_pid

    # read in Executor details
    puts '*** Reading Flow output file... ***'
    flow_details = JSON.parse(File.read("#{Dir.pwd}/#{flow}.out"))

    # poll until Workflow is complete
    puts '*** Polling for Workflow completion... ***'
    swf = Aws::SWF::Client.new
    max_time = Time.now + timeout_seconds
    while Time.now < max_time
      exec = swf.describe_workflow_execution(
        domain: ENVied.WORKFLOW_DOMAIN,
        execution: { workflow_id: flow_details['workflow_id'],
                     run_id: flow_details['run_id'] })
      unless exec.successful?
        puts ">>>>>> Execution for #{flow} FAILED: <<<<<<<"
        puts exec.data.execution_info.inspect if exec.data
        exit(-1)
      end
      break if exec.data && exec.data.execution_info.execution_status == 'CLOSED'
      sleep 10
    end
  end

  after(:all) do
    # truncate rows, aka clean up after company leaves
    truncate_rows

    # shutdown flow execution
    # shut down Worker
    puts '*** Shutting down Worker... ***'
    Process.kill('TERM', worker_pid)

    # shut down Decider
    puts '*** Shutting down Decider... ***'
    Process.kill('TERM', decider_pid)

    # remove out file
    puts '*** Clean up Flow output file... ***'
    `rm #{Dir.pwd}/#{flow}.out`
  end
end
