Skip to content

Commit 0d5cd44

Browse files
committed
Fail stuck service instance 'update' operations
When CCDB is briefly unavailable during a broker polling cycle, the CC polling job can fail permanently (max_attempts=1) while the broker is still processing. This leaves an update stuck: last_operation.state stays 'in progress' with no delayed job working on it, requiring operator intervention. Add ServiceOperationsUpdateInProgressCleanup, a periodic job that detects stuck 'update' operations whose polling job has permanently failed and marks the operation and its pollable job as 'failed', giving clients a definitive final state. A FOR UPDATE SKIP LOCKED guard prevents double processing across concurrent CC instances. Unlike the create cleanup, no orphan mitigation is triggered: an update targets a resource that already exists, so it must not be deprovisioned.
1 parent 846c257 commit 0d5cd44

7 files changed

Lines changed: 283 additions & 0 deletions

File tree

Lines changed: 84 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,84 @@
1+
module VCAP::CloudController
2+
module Jobs
3+
module Runtime
4+
class ServiceOperationsUpdateInProgressCleanup < VCAP::CloudController::Jobs::CCJob
5+
BATCH_SIZE = 10
6+
7+
def perform
8+
logger.info("Cleaning up service 'update' operations stuck in 'in progress'")
9+
cleanup_operations(ServiceInstanceOperation, ServiceInstance, :service_instance_id, 'service_instance.update')
10+
end
11+
12+
def max_attempts
13+
1
14+
end
15+
16+
private
17+
18+
def cleanup_operations(operation_model, instance_model, foreign_key, jobs_operation)
19+
operation_table = operation_model.table_name
20+
instance_table = instance_model.table_name
21+
22+
stuck = operation_model.
23+
join(instance_table, id: Sequel[operation_table][foreign_key]).
24+
join(:jobs, resource_guid: Sequel[instance_table][:guid]).
25+
join(:delayed_jobs, guid: Sequel[:jobs][:delayed_job_guid]).
26+
where(Sequel[operation_table][:state] => 'in progress').
27+
where(Sequel[operation_table][:type] => 'update').
28+
where(Sequel.lit("#{operation_table}.created_at > CURRENT_TIMESTAMP - INTERVAL '?' SECOND", default_maximum_duration_seconds.to_i)).
29+
where(Sequel[:jobs][:state] => [PollableJobModel::POLLING_STATE, PollableJobModel::FAILED_STATE]).
30+
where(Sequel[:jobs][:operation] => jobs_operation).
31+
exclude(Sequel[:delayed_jobs][:failed_at] => nil).
32+
select(
33+
Sequel[:jobs][:guid].as(:pollable_guid),
34+
Sequel[operation_table][:id].as(:op_id),
35+
Sequel[operation_table][foreign_key].as(:resource_id)
36+
).
37+
order(Sequel[operation_table][:created_at]).
38+
limit(BATCH_SIZE)
39+
40+
stuck.each do |row|
41+
resolve_stuck(operation_model, instance_model, row[:op_id], row[:resource_id], row[:pollable_guid])
42+
end
43+
end
44+
45+
def resolve_stuck(operation_model, instance_model, op_id, resource_id, pollable_guid)
46+
operation_model.db.transaction do
47+
operation = operation_model.where(id: op_id, state: 'in progress').for_update.skip_locked.first
48+
return unless operation
49+
50+
instance = instance_model.first(id: resource_id)
51+
return unless instance
52+
53+
instance_type = instance_model.to_s.split('::').last
54+
55+
logger.info(
56+
"#{instance_type} #{instance.guid} update operation is stuck in 'in progress'. " \
57+
"Setting operation's state to 'failed' and pollable job's state to 'FAILED'.",
58+
instance_type: instance_type,
59+
instance_guid: instance.guid,
60+
operation_id: op_id,
61+
pollable_job_guid: pollable_guid
62+
)
63+
64+
operation.update(state: 'failed',
65+
description: "Operation was stuck in 'in progress' state. Set to 'failed' by cleanup job.")
66+
PollableJobModel.where(guid: pollable_guid).update(state: PollableJobModel::FAILED_STATE)
67+
end
68+
end
69+
70+
def default_maximum_duration_seconds
71+
Config.config.get(:broker_client_max_async_poll_duration_minutes).minutes
72+
end
73+
74+
def logger
75+
@logger ||= Steno.logger('cc.background.service-operations-update-in-progress-cleanup')
76+
end
77+
78+
def job_name_in_configuration
79+
:service_operations_update_in_progress_cleanup
80+
end
81+
end
82+
end
83+
end
84+
end

config/cloud_controller.yml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,9 @@ service_operations_initial_cleanup:
5858
service_operations_create_in_progress_cleanup:
5959
frequency_in_seconds: 3600 #1h
6060

61+
service_operations_update_in_progress_cleanup:
62+
frequency_in_seconds: 3600 #1h
63+
6164
# One-off backfill - to be removed in a future version.
6265
lifecycle_type_backfill:
6366
frequency_in_seconds: 3600 #1h

lib/cloud_controller/clock/scheduler.rb

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ class Scheduler
2626
{ name: 'failed_jobs', class: Jobs::Runtime::FailedJobsCleanup },
2727
{ name: 'service_operations_initial_cleanup', class: Jobs::Runtime::ServiceOperationsInitialCleanup },
2828
{ name: 'service_operations_create_in_progress_cleanup', class: Jobs::Runtime::ServiceOperationsCreateInProgressCleanup },
29+
{ name: 'service_operations_update_in_progress_cleanup', class: Jobs::Runtime::ServiceOperationsUpdateInProgressCleanup },
2930
# One-off backfill - to be removed in a future version.
3031
{ name: 'lifecycle_type_backfill', class: Jobs::Runtime::LifecycleTypeBackfill }
3132
].freeze

lib/cloud_controller/config_schemas/clock_schema.rb

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,9 @@ class ClockSchema < VCAP::Config
3737
service_operations_create_in_progress_cleanup: {
3838
frequency_in_seconds: Integer
3939
},
40+
service_operations_update_in_progress_cleanup: {
41+
frequency_in_seconds: Integer
42+
},
4043
# One-off backfill - to be removed in a future version.
4144
lifecycle_type_backfill: {
4245
frequency_in_seconds: Integer

lib/cloud_controller/jobs.rb

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
require 'jobs/runtime/expired_orphaned_blob_cleanup'
2727
require 'jobs/runtime/expired_resource_cleanup'
2828
require 'jobs/runtime/service_operations_create_in_progress_cleanup'
29+
require 'jobs/runtime/service_operations_update_in_progress_cleanup'
2930
require 'jobs/runtime/failed_jobs_cleanup'
3031
require 'jobs/runtime/service_operations_initial_cleanup'
3132
require 'jobs/runtime/legacy_jobs'
Lines changed: 184 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,184 @@
1+
require 'spec_helper'
2+
3+
module VCAP::CloudController
4+
module Jobs::Runtime
5+
RSpec.describe ServiceOperationsUpdateInProgressCleanup, job_context: :worker do
6+
subject(:job) { ServiceOperationsUpdateInProgressCleanup.new }
7+
8+
let(:fake_logger) { instance_double(Steno::Logger, info: nil, warn: nil) }
9+
let(:max_poll_duration_minutes) { 60 }
10+
11+
before do
12+
allow(Steno).to receive(:logger).and_return(fake_logger)
13+
TestConfig.override(broker_client_max_async_poll_duration_minutes: max_poll_duration_minutes)
14+
end
15+
16+
def prepare_stuck_service_instance(
17+
service_instance_state: 'in progress',
18+
service_instance_type: 'update',
19+
service_instance_created_at: Time.now,
20+
pollable_job_state: PollableJobModel::FAILED_STATE,
21+
pollable_job_operation: 'service_instance.update',
22+
delayed_job_failed_at: Time.now
23+
)
24+
service_instance = create(:managed_service_instance)
25+
26+
create(:service_instance_operation,
27+
service_instance_id: service_instance.id,
28+
type: service_instance_type,
29+
state: service_instance_state,
30+
created_at: service_instance_created_at)
31+
32+
dj = Delayed::Job.create!(
33+
guid: SecureRandom.uuid,
34+
handler: 'fake',
35+
run_at: Time.now,
36+
failed_at: delayed_job_failed_at,
37+
queue: 'cc-generic'
38+
)
39+
40+
pjob = create(:pollable_job_model,
41+
state: pollable_job_state,
42+
operation: pollable_job_operation,
43+
resource_guid: service_instance.guid,
44+
resource_type: 'service_instances',
45+
delayed_job_guid: dj.guid)
46+
47+
{ service_instance: service_instance, pjob: pjob, delayed_job: dj }
48+
end
49+
50+
shared_examples 'does not resolve the operation' do
51+
it 'leaves the operation in progress and the pollable job untouched' do
52+
scenario = subject_scenario
53+
job.perform
54+
expect(scenario[:service_instance].last_operation.reload.state).to eq('in progress')
55+
expect(scenario[:pjob].reload.state).to eq(scenario[:pjob].state)
56+
end
57+
end
58+
59+
it { is_expected.to be_a_valid_job }
60+
61+
describe '#perform' do
62+
context 'when sio state is not in progress' do
63+
it 'does not resolve when state is succeeded' do
64+
scenario = prepare_stuck_service_instance(service_instance_state: 'succeeded')
65+
job.perform
66+
expect(scenario[:service_instance].last_operation.reload.state).to eq('succeeded')
67+
end
68+
69+
it 'does not resolve when state is failed' do
70+
scenario = prepare_stuck_service_instance(service_instance_state: 'failed')
71+
job.perform
72+
expect(scenario[:service_instance].last_operation.reload.state).to eq('failed')
73+
end
74+
end
75+
76+
context 'when sio type is not update' do
77+
let(:subject_scenario) { prepare_stuck_service_instance(service_instance_type: 'create') }
78+
79+
it_behaves_like 'does not resolve the operation'
80+
end
81+
82+
context 'when sio created_at is beyond the max polling window' do
83+
let(:subject_scenario) { prepare_stuck_service_instance(service_instance_created_at: Time.now - (max_poll_duration_minutes + 1).minutes) }
84+
85+
it_behaves_like 'does not resolve the operation'
86+
end
87+
88+
context 'when delayed_job.failed_at is nil (job still running or locked)' do
89+
let(:subject_scenario) { prepare_stuck_service_instance(delayed_job_failed_at: nil) }
90+
91+
it_behaves_like 'does not resolve the operation'
92+
end
93+
94+
context 'when pollable job state is COMPLETE' do
95+
let(:subject_scenario) { prepare_stuck_service_instance(pollable_job_state: PollableJobModel::COMPLETE_STATE) }
96+
97+
it_behaves_like 'does not resolve the operation'
98+
end
99+
100+
context 'when pollable job state is PROCESSING' do
101+
let(:subject_scenario) { prepare_stuck_service_instance(pollable_job_state: PollableJobModel::PROCESSING_STATE) }
102+
103+
it_behaves_like 'does not resolve the operation'
104+
end
105+
106+
context 'when pollable job operation is not service_instance.update' do
107+
let(:subject_scenario) { prepare_stuck_service_instance(pollable_job_operation: 'service_instance.create') }
108+
109+
it_behaves_like 'does not resolve the operation'
110+
end
111+
112+
context 'when a service instance update job is stuck with state FAILED' do
113+
it 'sets operation to failed and pollable job to FAILED' do
114+
scenario = prepare_stuck_service_instance
115+
job.perform
116+
expect(scenario[:service_instance].last_operation.reload.state).to eq('failed')
117+
expect(scenario[:pjob].reload.state).to eq(PollableJobModel::FAILED_STATE)
118+
end
119+
end
120+
121+
context 'when a service instance update job is stuck with state POLLING (DB flip before failure hook)' do
122+
it 'sets operation to failed and pollable job to FAILED' do
123+
scenario = prepare_stuck_service_instance(pollable_job_state: PollableJobModel::POLLING_STATE)
124+
job.perform
125+
expect(scenario[:service_instance].last_operation.reload.state).to eq('failed')
126+
expect(scenario[:pjob].reload.state).to eq(PollableJobModel::FAILED_STATE)
127+
end
128+
end
129+
130+
context 'when there are multiple stuck jobs within the batch size' do
131+
it 'resolves each one' do
132+
3.times { prepare_stuck_service_instance }
133+
job.perform
134+
expect(ServiceInstanceOperation.where(state: 'failed').count).to eq(3)
135+
end
136+
end
137+
138+
context 'when there are more stuck jobs than the batch size' do
139+
it 'processes only up to BATCH_SIZE jobs per run' do
140+
(ServiceOperationsUpdateInProgressCleanup::BATCH_SIZE + 1).times { prepare_stuck_service_instance }
141+
job.perform
142+
expect(ServiceInstanceOperation.where(state: 'failed').count).to eq(ServiceOperationsUpdateInProgressCleanup::BATCH_SIZE)
143+
end
144+
end
145+
end
146+
147+
describe '#resolve_stuck' do
148+
context 'when another process already resolved it (skip_locked returns nil)' do
149+
it 'does nothing' do
150+
scenario = prepare_stuck_service_instance
151+
152+
expect do
153+
job.send(:resolve_stuck, ServiceInstanceOperation, ServiceInstance,
154+
-1, scenario[:service_instance].id, scenario[:pjob].guid)
155+
end.not_to raise_error
156+
expect(scenario[:service_instance].last_operation.reload.state).to eq('in progress')
157+
end
158+
end
159+
160+
context 'when the operation is stuck in progress' do
161+
it 'sets the operation state from in progress to failed' do
162+
scenario = prepare_stuck_service_instance
163+
op = scenario[:service_instance].last_operation
164+
165+
expect do
166+
job.send(:resolve_stuck, ServiceInstanceOperation, ServiceInstance,
167+
op.id, scenario[:service_instance].id, scenario[:pjob].guid)
168+
end.to change { op.reload.state }.from('in progress').to('failed')
169+
end
170+
171+
it 'sets the pollable job state to FAILED' do
172+
scenario = prepare_stuck_service_instance(pollable_job_state: PollableJobModel::POLLING_STATE)
173+
op = scenario[:service_instance].last_operation
174+
175+
expect do
176+
job.send(:resolve_stuck, ServiceInstanceOperation, ServiceInstance,
177+
op.id, scenario[:service_instance].id, scenario[:pjob].guid)
178+
end.to change { scenario[:pjob].reload.state }.from(PollableJobModel::POLLING_STATE).to(PollableJobModel::FAILED_STATE)
179+
end
180+
end
181+
end
182+
end
183+
end
184+
end

spec/unit/lib/cloud_controller/clock/scheduler_spec.rb

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ module VCAP::CloudController
2222
pollable_jobs: { cutoff_age_in_days: 2 },
2323
service_operations_initial_cleanup: { frequency_in_seconds: 600 },
2424
service_operations_create_in_progress_cleanup: { frequency_in_seconds: 600 },
25+
service_operations_update_in_progress_cleanup: { frequency_in_seconds: 600 },
2526
lifecycle_type_backfill: { frequency_in_seconds: 500 },
2627
service_usage_events: { cutoff_age_in_days: 5 },
2728
completed_tasks: { cutoff_age_in_days: 6 },
@@ -169,6 +170,12 @@ module VCAP::CloudController
169170
expect(block.call).to be_instance_of(Jobs::Runtime::ServiceOperationsCreateInProgressCleanup)
170171
end
171172

173+
expect(clock).to receive(:schedule_frequent_worker_job) do |args, &block|
174+
expect(args).to eql(name: 'service_operations_update_in_progress_cleanup', interval: 600)
175+
expect(Jobs::Runtime::ServiceOperationsUpdateInProgressCleanup).to receive(:new).and_call_original
176+
expect(block.call).to be_instance_of(Jobs::Runtime::ServiceOperationsUpdateInProgressCleanup)
177+
end
178+
172179
expect(clock).to receive(:schedule_frequent_worker_job) do |args, &block|
173180
expect(args).to eql(name: 'lifecycle_type_backfill', interval: 500)
174181
expect(Jobs::Runtime::LifecycleTypeBackfill).to receive(:new).and_call_original

0 commit comments

Comments
 (0)