Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Gemfile.lock
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
PATH
remote: .
specs:
entitlements-app (1.2.4)
entitlements-app (1.2.5)
concurrent-ruby (~> 1.3, >= 1.3.1)
dogstatsd-ruby (~> 5.7)
faraday (~> 2.0)
Expand Down
35 changes: 25 additions & 10 deletions lib/entitlements.rb
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
require "ostruct"
require "resolv"
require "stringio"
require "thread"
require "uri"
require "yaml"

Expand Down Expand Up @@ -457,19 +458,33 @@ def self.calculate_actions
# Calculate old and new membership in each group.
thread_pool = Concurrent::FixedThreadPool.new(max_parallelism)
logger.debug("Begin prefetch and validate for all groups")

prep_start = Time.now
futures = Entitlements.child_classes.map do |group_name, obj|
Concurrent::Future.execute({ executor: thread_pool }) do
group_start = Time.now
logger.debug("Begin prefetch and validate for #{group_name}")
provider = Entitlements.config["groups"].fetch(group_name).fetch("type")
timed_operation(phase: "prefetch", provider: provider, target: group_name, concurrent: true) { obj.prefetch }
timed_operation(phase: "validate", provider: provider, target: group_name, concurrent: true) { obj.validate }
logger.debug("Finished prefetch and validate for #{group_name} in #{Time.now - group_start}")
jobs = Entitlements.child_classes
completions = Queue.new

begin
jobs.each do |group_name, obj|
thread_pool.post do
group_start = Time.now
logger.debug("Begin prefetch and validate for #{group_name}")
provider = Entitlements.config["groups"].fetch(group_name).fetch("type")
timed_operation(phase: "prefetch", provider: provider, target: group_name, concurrent: true) { obj.prefetch }
timed_operation(phase: "validate", provider: provider, target: group_name, concurrent: true) { obj.validate }
logger.debug("Finished prefetch and validate for #{group_name} in #{Time.now - group_start}")
completions << nil
rescue => e
completions << e
end
end
end

futures.each(&:value!)
jobs.size.times do
exception = completions.pop
raise exception if exception
end
ensure
thread_pool.kill
end
logger.debug("Finished all prefetch and validate in #{Time.now - prep_start}")

logger.debug("Begin all calculations")
Expand Down
2 changes: 1 addition & 1 deletion lib/version.rb
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,6 @@

module Entitlements
module Version
VERSION = "1.2.4"
VERSION = "1.2.5"
end
end
21 changes: 21 additions & 0 deletions spec/unit/entitlements_spec.rb
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
# frozen_string_literal: true

require_relative "spec_helper"
require "timeout"

describe Entitlements do
let(:subject) { Entitlements }
Expand Down Expand Up @@ -346,6 +347,26 @@

expect(cache[:change_count]).to eq(3)
end

it "raises a worker failure without waiting for an earlier job" do
Entitlements.config["max_parallelism"] = 2
started = Concurrent::Event.new
blocker = Concurrent::Event.new
allow(Entitlements).to receive(:child_classes)
.and_return("ldap-dir" => ldap_controller, "other-ldap-dir" => other_controller)
allow(ldap_controller).to receive(:prefetch) do
started.set
blocker.wait
end
allow(other_controller).to receive(:prefetch) do
started.wait
raise "Boom"
end

expect do
Timeout.timeout(1) { described_class.calculate }
end.to raise_error(RuntimeError, "Boom")
end
end

describe "#execute" do
Expand Down
Loading