From de4a6626514e73dc4b702f631ba4797725df6fdb Mon Sep 17 00:00:00 2001 From: Sam Cartwright Date: Thu, 20 Feb 2020 10:28:40 +0000 Subject: [PATCH 1/4] add kill flag to aws pipelines --- README.md | 13 +++++++ lib/eventq/eventq_aws/aws_queue_worker.rb | 8 ++++ lib/eventq/eventq_base/message_args.rb | 2 + lib/eventq/queue_worker.rb | 14 ++++++- .../integration/aws_queue_worker_v2_spec.rb | 37 +++++++++++++++++++ .../rabbitmq_queue_worker_spec.rb | 25 +++++++++++++ 6 files changed, 98 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index b74d2b3..e49c889 100644 --- a/README.md +++ b/README.md @@ -133,6 +133,19 @@ The on_retry_exceeded method allows you to specify a block that should execute w .... end +#### #on_killed + +The on_killed method allows you to specify a block that should execute whenever an event kills itself. The event object passed to the block is a **[QueueMessage]** object. + +**Example** + + worker.on_killed do |event| + .... + #Do something with the failed event + .... + end + + #### #on_retry The on_retry method allows you to specify a block that should execute whenever an event fails to process and is retried. The event object passed to the block is a **[QueueMessage]** object, and the abort arg is a Boolean that specifies if the message was aborted (true or false). diff --git a/lib/eventq/eventq_aws/aws_queue_worker.rb b/lib/eventq/eventq_aws/aws_queue_worker.rb index 5e9c4d2..b8382fb 100644 --- a/lib/eventq/eventq_aws/aws_queue_worker.rb +++ b/lib/eventq/eventq_aws/aws_queue_worker.rb @@ -113,6 +113,14 @@ def reject_message(queue, poller, msg, retry_attempts, message, args) EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Retry attempt limit exceeded.") context.call_on_retry_exceeded_block(message) end + elsif args.kill + EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Rejected without retry. Message: #{serialize_message(message)}") + + # remove the message from the queue so that it does not get retried again + poller.delete_message(msg) + + EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Message killed.") + context.call_on_killed_block(message) elsif queue.allow_retry retry_attempts += 1 diff --git a/lib/eventq/eventq_base/message_args.rb b/lib/eventq/eventq_base/message_args.rb index 854743a..40374ff 100644 --- a/lib/eventq/eventq_base/message_args.rb +++ b/lib/eventq/eventq_base/message_args.rb @@ -5,6 +5,7 @@ class MessageArgs attr_reader :retry_attempts attr_accessor :abort attr_accessor :drop + attr_accessor :kill attr_reader :context attr_reader :id attr_reader :sent @@ -14,6 +15,7 @@ def initialize(type:, retry_attempts:, context: {}, content_type:, id: nil, sent @retry_attempts = retry_attempts @abort = false @drop = false + @kill = false @context = context @content_type = content_type @id = id diff --git a/lib/eventq/queue_worker.rb b/lib/eventq/queue_worker.rb index a61b730..874fbc5 100644 --- a/lib/eventq/queue_worker.rb +++ b/lib/eventq/queue_worker.rb @@ -113,6 +113,7 @@ def start_thread(queue, options, block) # @return [Symbol, MessageArgs] :accepted, :duplicate, :reject def process_message(block, message, retry_attempts, acceptance_args) abort = false + kill = false error = false status = nil @@ -140,6 +141,9 @@ def process_message(block, message, retry_attempts, acceptance_args) if message_args.abort == true abort = true EventQ.logger.debug("[#{self.class}] - Message aborted. Id: #{message.id}.") + elsif message_args.kill == true + kill = true + EventQ.logger.debug("[#{self.class}] - Message killed. Id: #{message.id}.") else # accept the message as processed status = :accepted @@ -156,7 +160,7 @@ def process_message(block, message, retry_attempts, acceptance_args) call_on_error_block(error: e, message: message) end - if error || abort + if error || abort || kill EventQ::NonceManager.failed(message.id) status = :reject else @@ -248,6 +252,10 @@ def on_retry_exceeded(&block) @on_retry_exceeded_block = block end + def on_killed(&block) + @on_killed_block = block + end + def on_retry(&block) @on_retry_block = block end @@ -264,6 +272,10 @@ def call_on_retry_exceeded_block(message) call_block(:on_retry_exceeded_block, message) end + def call_on_killed_block(message) + call_block(:on_killed_block, message) + end + def call_on_retry_block(message) call_block(:on_retry_block, message) end diff --git a/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb b/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb index a17d1c5..7dc1ab8 100644 --- a/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb +++ b/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb @@ -200,6 +200,43 @@ expect(queue_worker.running?).to eq(false) end + it 'should receive an event from the subscriber queue and not retry it (kill).' do + + subscriber_queue.retry_delay = 1000 + subscriber_queue.allow_retry = true + + subscription_manager.subscribe(event_type, subscriber_queue) + eventq_client.raise_event(event_type, message) + + received = false + received_count = 0 + received_attribute = 0; + + # wait 1 second to allow the message to be sent and broadcast to the queue + sleep(1) + + queue_worker.start(subscriber_queue, { worker_adapter: subject, wait: false, block_process: false, client: queue_client }) do |event, args| + expect(event).to eq(message) + expect(args).to be_a(EventQ::MessageArgs) + received = true + received_count += 1 + received_attribute = args.retry_attempts + EventQ.logger.debug { "Message Received: #{event}" } + if received_count > 1 + args.kill = true + end + end + + sleep(4) + + queue_worker.stop + + expect(received).to eq(true) + expect(received_count).to eq(2) + expect(received_attribute).to eq(1) + expect(queue_worker.running?).to eq(false) + end + it 'should receive multiple events from the subscriber queue' do subscription_manager.subscribe(event_type2, subscriber_queue) diff --git a/spec/eventq_rabbitmq/rabbitmq_queue_worker_spec.rb b/spec/eventq_rabbitmq/rabbitmq_queue_worker_spec.rb index e518dfe..278ee4c 100644 --- a/spec/eventq_rabbitmq/rabbitmq_queue_worker_spec.rb +++ b/spec/eventq_rabbitmq/rabbitmq_queue_worker_spec.rb @@ -235,6 +235,31 @@ end end + describe '#call_on_killed_block' do + let(:message) { double } + context 'when a block is specified' do + let(:block) { Proc.new { } } + before do + queue_worker.on_killed &block + allow(block).to receive(:call) + end + it 'should execute the block' do + expect(block).to receive(:call).with(message).once + queue_worker.call_on_killed_block(message) + end + end + context 'when a block is NOT specified' do + let(:block) { nil } + before do + queue_worker.on_killed &block + end + it 'should NOT execute the block' do + expect(block).not_to receive(:call) + queue_worker.call_on_killed_block(message) + end + end + end + it 'should receive an event from the subscriber queue' do event_type = 'queue.worker.event1' From cf27e49914142d1849e4863cd3427a0b6e245ea6 Mon Sep 17 00:00:00 2001 From: Sam Cartwright Date: Thu, 20 Feb 2020 10:37:39 +0000 Subject: [PATCH 2/4] refactor reject method --- lib/eventq/eventq_aws/aws_queue_worker.rb | 72 +++++++++++-------- .../integration/aws_queue_worker_v2_spec.rb | 10 +-- 2 files changed, 47 insertions(+), 35 deletions(-) diff --git a/lib/eventq/eventq_aws/aws_queue_worker.rb b/lib/eventq/eventq_aws/aws_queue_worker.rb index b8382fb..6364997 100644 --- a/lib/eventq/eventq_aws/aws_queue_worker.rb +++ b/lib/eventq/eventq_aws/aws_queue_worker.rb @@ -104,44 +104,56 @@ def process_message(msg, poller, queue, block) def reject_message(queue, poller, msg, retry_attempts, message, args) if !queue.allow_retry || retry_attempts >= queue.max_retry_attempts - EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Rejected removing from queue. Message: #{serialize_message(message)}") + queue_will_not_retry_message(queue, poller, msg, message) + elsif args.kill + queue_will_kill_message(poller, msg, message) + elsif queue.allow_retry + queue_will_retry_message(queue, poller, msg, retry_attempts, message, args) + end + end - # remove the message from the queue so that it does not get retried again - poller.delete_message(msg) + def queue_will_not_retry_message(queue, poller, msg, message) + EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Rejected removing from queue. Message: #{serialize_message(message)}") - if retry_attempts >= queue.max_retry_attempts - EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Retry attempt limit exceeded.") - context.call_on_retry_exceeded_block(message) - end - elsif args.kill - EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Rejected without retry. Message: #{serialize_message(message)}") + # remove the message from the queue so that it does not get retried again + poller.delete_message(msg) - # remove the message from the queue so that it does not get retried again - poller.delete_message(msg) + if retry_attempts >= queue.max_retry_attempts + EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Retry attempt limit exceeded.") + context.call_on_retry_exceeded_block(message) + end + end - EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Message killed.") - context.call_on_killed_block(message) - elsif queue.allow_retry - retry_attempts += 1 + def queue_will_kill_message(poller, msg, message) + EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Rejected without retry. Message: #{serialize_message(message)}") - EventQ.logger.warn("[#{self.class}] - Message Id: #{args.id}. Rejected requesting retry. Attempts: #{retry_attempts}") + # remove the message from the queue so that it does not get retried again + poller.delete_message(msg) - visibility_timeout = @calculate_visibility_timeout.call( - retry_attempts: retry_attempts, - queue_settings: { - allow_retry_back_off: queue.allow_retry_back_off, - max_retry_delay: queue.max_retry_delay, - retry_back_off_grace: queue.retry_back_off_grace, - retry_back_off_weight: queue.retry_back_off_weight, - retry_delay: queue.retry_delay - } - ) + EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Message killed.") + context.call_on_killed_block(message) + end - EventQ.logger.debug { "[#{self.class}] - Sending message for retry. Message TTL: #{visibility_timeout}" } - poller.change_message_visibility_timeout(msg, visibility_timeout) + def queue_will_retry_message(queue, poller, msg, retry_attempts, message, args) + retry_attempts += 1 - context.call_on_retry_block(message) - end + EventQ.logger.warn("[#{self.class}] - Message Id: #{args.id}. Rejected requesting retry. Attempts: #{retry_attempts}") + + visibility_timeout = @calculate_visibility_timeout.call( + retry_attempts: retry_attempts, + queue_settings: { + allow_retry_back_off: queue.allow_retry_back_off, + max_retry_delay: queue.max_retry_delay, + retry_back_off_grace: queue.retry_back_off_grace, + retry_back_off_weight: queue.retry_back_off_weight, + retry_delay: queue.retry_delay + } + ) + + EventQ.logger.debug { "[#{self.class}] - Sending message for retry. Message TTL: #{visibility_timeout}" } + poller.change_message_visibility_timeout(msg, visibility_timeout) + + context.call_on_retry_block(message) end end end diff --git a/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb b/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb index 7dc1ab8..f5fed03 100644 --- a/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb +++ b/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb @@ -173,7 +173,7 @@ received = false received_count = 0 - received_attribute = 0; + received_attribute = 0 # wait 1 second to allow the message to be sent and broadcast to the queue sleep(1) @@ -210,7 +210,7 @@ received = false received_count = 0 - received_attribute = 0; + received_attribute = 0 # wait 1 second to allow the message to be sent and broadcast to the queue sleep(1) @@ -222,7 +222,7 @@ received_count += 1 received_attribute = args.retry_attempts EventQ.logger.debug { "Message Received: #{event}" } - if received_count > 1 + if received_count == 3 args.kill = true end end @@ -232,8 +232,8 @@ queue_worker.stop expect(received).to eq(true) - expect(received_count).to eq(2) - expect(received_attribute).to eq(1) + expect(received_count).to eq(1) + expect(received_attribute).to eq(3) expect(queue_worker.running?).to eq(false) end From 70288b1f77d4afb4d1b75d6ae71510a5619f1362 Mon Sep 17 00:00:00 2001 From: Sam Cartwright Date: Thu, 20 Feb 2020 11:58:06 +0000 Subject: [PATCH 3/4] correct test logic --- spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb b/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb index f5fed03..86e7a4d 100644 --- a/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb +++ b/spec/eventq_aws/integration/aws_queue_worker_v2_spec.rb @@ -224,6 +224,8 @@ EventQ.logger.debug { "Message Received: #{event}" } if received_count == 3 args.kill = true + else + args.abort = true end end @@ -232,8 +234,8 @@ queue_worker.stop expect(received).to eq(true) - expect(received_count).to eq(1) - expect(received_attribute).to eq(3) + expect(received_count).to eq(3) + expect(received_attribute).to eq(2) expect(queue_worker.running?).to eq(false) end From 635d03b73a6c82dd11b8f3d27ed4812cea219b14 Mon Sep 17 00:00:00 2001 From: Sam Cartwright Date: Thu, 20 Feb 2020 13:39:04 +0000 Subject: [PATCH 4/4] unfactor to avoid linter whinge --- lib/eventq/eventq_aws/aws_queue_worker.rb | 42 ++++++++++------------- 1 file changed, 19 insertions(+), 23 deletions(-) diff --git a/lib/eventq/eventq_aws/aws_queue_worker.rb b/lib/eventq/eventq_aws/aws_queue_worker.rb index 6364997..a8c53fd 100644 --- a/lib/eventq/eventq_aws/aws_queue_worker.rb +++ b/lib/eventq/eventq_aws/aws_queue_worker.rb @@ -108,7 +108,25 @@ def reject_message(queue, poller, msg, retry_attempts, message, args) elsif args.kill queue_will_kill_message(poller, msg, message) elsif queue.allow_retry - queue_will_retry_message(queue, poller, msg, retry_attempts, message, args) + retry_attempts += 1 + + EventQ.logger.warn("[#{self.class}] - Message Id: #{args.id}. Rejected requesting retry. Attempts: #{retry_attempts}") + + visibility_timeout = @calculate_visibility_timeout.call( + retry_attempts: retry_attempts, + queue_settings: { + allow_retry_back_off: queue.allow_retry_back_off, + max_retry_delay: queue.max_retry_delay, + retry_back_off_grace: queue.retry_back_off_grace, + retry_back_off_weight: queue.retry_back_off_weight, + retry_delay: queue.retry_delay + } + ) + + EventQ.logger.debug { "[#{self.class}] - Sending message for retry. Message TTL: #{visibility_timeout}" } + poller.change_message_visibility_timeout(msg, visibility_timeout) + + context.call_on_retry_block(message) end end @@ -133,28 +151,6 @@ def queue_will_kill_message(poller, msg, message) EventQ.logger.error("[#{self.class}] - Message Id: #{args.id}. Message killed.") context.call_on_killed_block(message) end - - def queue_will_retry_message(queue, poller, msg, retry_attempts, message, args) - retry_attempts += 1 - - EventQ.logger.warn("[#{self.class}] - Message Id: #{args.id}. Rejected requesting retry. Attempts: #{retry_attempts}") - - visibility_timeout = @calculate_visibility_timeout.call( - retry_attempts: retry_attempts, - queue_settings: { - allow_retry_back_off: queue.allow_retry_back_off, - max_retry_delay: queue.max_retry_delay, - retry_back_off_grace: queue.retry_back_off_grace, - retry_back_off_weight: queue.retry_back_off_weight, - retry_delay: queue.retry_delay - } - ) - - EventQ.logger.debug { "[#{self.class}] - Sending message for retry. Message TTL: #{visibility_timeout}" } - poller.change_message_visibility_timeout(msg, visibility_timeout) - - context.call_on_retry_block(message) - end end end end