Skip to content

Commit 67fcc6e

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 67fcc6e

7 files changed

Lines changed: 311 additions & 0 deletions

File tree

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,108 @@
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+
# Find stuck service instance 'in progress' update operations where the broker is still working
20+
# but CC's polling job has permanently failed due to a transient error (e.g. brief db connection flip).
21+
# Join path: service_instance_operations → service_instances → jobs → delayed_jobs.
22+
#
23+
# Filters:
24+
# - service_instance_operations.state='in progress': the broker has not yet reported a final state
25+
# (succeeded or failed) that CC could successfully persist; if CC had received and saved a final
26+
# state from the broker, this column would already be 'succeeded' or 'failed' — not 'in progress'
27+
# - service_instance_operations.type='update': scope to update operations only
28+
# - service_instance_operations.created_at > CURRENT_TIMESTAMP - max_duration: operations beyond the max async polling window
29+
# are intentionally excluded — the broker has given up on them too, so they are out of scope for this cleanup
30+
# - jobs.state IN (POLLING, FAILED): the pollable job has not reached COMPLETE (a successful job
31+
# would already be done and is out of scope); POLLING covers the case where the failure hook
32+
# itself couldn't write FAILED due to the DB flip
33+
# - jobs.operation='service_instance.update': prevents matching create/delete jobs for the same
34+
# service instance that happen to share the same resource_guid
35+
# - delayed_jobs.failed_at IS NOT NULL: the delayed job permanently failed (exhausted max_attempts);
36+
# jobs still alive or locked have failed_at=NULL and must not be touched
37+
#
38+
# Unlike the create cleanup, no orphan mitigation is triggered: an update targets a resource that
39+
# already exists, so we must not deprovision it. We only give the client a definitive final state
40+
# by marking the operation and its pollable job as failed.
41+
operation_table = operation_model.table_name
42+
instance_table = instance_model.table_name
43+
44+
stuck = operation_model.
45+
join(instance_table, id: Sequel[operation_table][foreign_key]).
46+
join(:jobs, resource_guid: Sequel[instance_table][:guid]).
47+
join(:delayed_jobs, guid: Sequel[:jobs][:delayed_job_guid]).
48+
where(Sequel[operation_table][:state] => 'in progress').
49+
where(Sequel[operation_table][:type] => 'update').
50+
where(Sequel.lit("#{operation_table}.created_at > CURRENT_TIMESTAMP - INTERVAL '?' SECOND", default_maximum_duration_seconds.to_i)).
51+
where(Sequel[:jobs][:state] => [PollableJobModel::POLLING_STATE, PollableJobModel::FAILED_STATE]).
52+
where(Sequel[:jobs][:operation] => jobs_operation).
53+
exclude(Sequel[:delayed_jobs][:failed_at] => nil).
54+
select(
55+
Sequel[:jobs][:guid].as(:pollable_guid),
56+
Sequel[operation_table][:id].as(:op_id),
57+
Sequel[operation_table][foreign_key].as(:resource_id)
58+
).
59+
order(Sequel[operation_table][:created_at]).
60+
limit(BATCH_SIZE)
61+
62+
stuck.each do |row|
63+
resolve_stuck(operation_model, instance_model, row[:op_id], row[:resource_id], row[:pollable_guid])
64+
end
65+
end
66+
67+
def resolve_stuck(operation_model, instance_model, op_id, resource_id, pollable_guid)
68+
# Mark the stuck update operation as failed and mark its pollable job as failed,
69+
# giving the client a definitive final state. No broker-side changes are reverted.
70+
operation_model.db.transaction do
71+
operation = operation_model.where(id: op_id, state: 'in progress').for_update.skip_locked.first
72+
return unless operation
73+
74+
instance = instance_model.first(id: resource_id)
75+
return unless instance
76+
77+
instance_type = instance_model.to_s.split('::').last
78+
79+
logger.info(
80+
"#{instance_type} #{instance.guid} update operation is stuck in 'in progress'. " \
81+
"Setting operation's state to 'failed' and pollable job's state to 'FAILED'.",
82+
instance_type: instance_type,
83+
instance_guid: instance.guid,
84+
operation_id: op_id,
85+
pollable_job_guid: pollable_guid
86+
)
87+
88+
operation.update(state: 'failed',
89+
description: "Operation was stuck in 'in progress' state. Set to 'failed' by cleanup job.")
90+
PollableJobModel.where(guid: pollable_guid).update(state: PollableJobModel::FAILED_STATE)
91+
end
92+
end
93+
94+
def default_maximum_duration_seconds
95+
Config.config.get(:broker_client_max_async_poll_duration_minutes).minutes
96+
end
97+
98+
def logger
99+
@logger ||= Steno.logger('cc.background.service-operations-update-in-progress-cleanup')
100+
end
101+
102+
def job_name_in_configuration
103+
:service_operations_update_in_progress_cleanup
104+
end
105+
end
106+
end
107+
end
108+
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: 188 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,188 @@
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+
# Builds a fully stuck scenario for ServiceInstance update that the job should pick up and resolve.
17+
# All filter conditions are satisfied: sio is in progress/update/within cutoff,
18+
# pjob is FAILED with operation=service_instance.update, delayed_job has failed_at set.
19+
# Override individual parameters to break a single filter and test exclusion.
20+
def prepare_stuck_service_instance(
21+
service_instance_state: 'in progress',
22+
service_instance_type: 'update',
23+
service_instance_created_at: Time.now,
24+
pollable_job_state: PollableJobModel::FAILED_STATE,
25+
pollable_job_operation: 'service_instance.update',
26+
delayed_job_failed_at: Time.now
27+
)
28+
service_instance = create(:managed_service_instance)
29+
30+
create(:service_instance_operation,
31+
service_instance_id: service_instance.id,
32+
type: service_instance_type,
33+
state: service_instance_state,
34+
created_at: service_instance_created_at)
35+
36+
dj = Delayed::Job.create!(
37+
guid: SecureRandom.uuid,
38+
handler: 'fake',
39+
run_at: Time.now,
40+
failed_at: delayed_job_failed_at,
41+
queue: 'cc-generic'
42+
)
43+
44+
pjob = create(:pollable_job_model,
45+
state: pollable_job_state,
46+
operation: pollable_job_operation,
47+
resource_guid: service_instance.guid,
48+
resource_type: 'service_instances',
49+
delayed_job_guid: dj.guid)
50+
51+
{ service_instance: service_instance, pjob: pjob, delayed_job: dj }
52+
end
53+
54+
shared_examples 'does not resolve the operation' do
55+
it 'leaves the operation in progress and the pollable job untouched' do
56+
scenario = subject_scenario
57+
job.perform
58+
expect(scenario[:service_instance].last_operation.reload.state).to eq('in progress')
59+
expect(scenario[:pjob].reload.state).to eq(scenario[:pjob].state)
60+
end
61+
end
62+
63+
it { is_expected.to be_a_valid_job }
64+
65+
describe '#perform' do
66+
context 'when sio state is not in progress' do
67+
it 'does not resolve when state is succeeded' do
68+
scenario = prepare_stuck_service_instance(service_instance_state: 'succeeded')
69+
job.perform
70+
expect(scenario[:service_instance].last_operation.reload.state).to eq('succeeded')
71+
end
72+
73+
it 'does not resolve when state is failed' do
74+
scenario = prepare_stuck_service_instance(service_instance_state: 'failed')
75+
job.perform
76+
expect(scenario[:service_instance].last_operation.reload.state).to eq('failed')
77+
end
78+
end
79+
80+
context 'when sio type is not update' do
81+
let(:subject_scenario) { prepare_stuck_service_instance(service_instance_type: 'create') }
82+
83+
it_behaves_like 'does not resolve the operation'
84+
end
85+
86+
context 'when sio created_at is beyond the max polling window' do
87+
let(:subject_scenario) { prepare_stuck_service_instance(service_instance_created_at: Time.now - (max_poll_duration_minutes + 1).minutes) }
88+
89+
it_behaves_like 'does not resolve the operation'
90+
end
91+
92+
context 'when delayed_job.failed_at is nil (job still running or locked)' do
93+
let(:subject_scenario) { prepare_stuck_service_instance(delayed_job_failed_at: nil) }
94+
95+
it_behaves_like 'does not resolve the operation'
96+
end
97+
98+
context 'when pollable job state is COMPLETE' do
99+
let(:subject_scenario) { prepare_stuck_service_instance(pollable_job_state: PollableJobModel::COMPLETE_STATE) }
100+
101+
it_behaves_like 'does not resolve the operation'
102+
end
103+
104+
context 'when pollable job state is PROCESSING' do
105+
let(:subject_scenario) { prepare_stuck_service_instance(pollable_job_state: PollableJobModel::PROCESSING_STATE) }
106+
107+
it_behaves_like 'does not resolve the operation'
108+
end
109+
110+
context 'when pollable job operation is not service_instance.update' do
111+
let(:subject_scenario) { prepare_stuck_service_instance(pollable_job_operation: 'service_instance.create') }
112+
113+
it_behaves_like 'does not resolve the operation'
114+
end
115+
116+
context 'when a service instance update job is stuck with state FAILED' do
117+
it 'sets operation to failed and pollable job to FAILED' do
118+
scenario = prepare_stuck_service_instance
119+
job.perform
120+
expect(scenario[:service_instance].last_operation.reload.state).to eq('failed')
121+
expect(scenario[:pjob].reload.state).to eq(PollableJobModel::FAILED_STATE)
122+
end
123+
end
124+
125+
context 'when a service instance update job is stuck with state POLLING (DB flip before failure hook)' do
126+
it 'sets operation to failed and pollable job to FAILED' do
127+
scenario = prepare_stuck_service_instance(pollable_job_state: PollableJobModel::POLLING_STATE)
128+
job.perform
129+
expect(scenario[:service_instance].last_operation.reload.state).to eq('failed')
130+
expect(scenario[:pjob].reload.state).to eq(PollableJobModel::FAILED_STATE)
131+
end
132+
end
133+
134+
context 'when there are multiple stuck jobs within the batch size' do
135+
it 'resolves each one' do
136+
3.times { prepare_stuck_service_instance }
137+
job.perform
138+
expect(ServiceInstanceOperation.where(state: 'failed').count).to eq(3)
139+
end
140+
end
141+
142+
context 'when there are more stuck jobs than the batch size' do
143+
it 'processes only up to BATCH_SIZE jobs per run' do
144+
(ServiceOperationsUpdateInProgressCleanup::BATCH_SIZE + 1).times { prepare_stuck_service_instance }
145+
job.perform
146+
expect(ServiceInstanceOperation.where(state: 'failed').count).to eq(ServiceOperationsUpdateInProgressCleanup::BATCH_SIZE)
147+
end
148+
end
149+
end
150+
151+
describe '#resolve_stuck' do
152+
context 'when another process already resolved it (skip_locked returns nil)' do
153+
it 'does nothing' do
154+
scenario = prepare_stuck_service_instance
155+
156+
expect do
157+
job.send(:resolve_stuck, ServiceInstanceOperation, ServiceInstance,
158+
-1, scenario[:service_instance].id, scenario[:pjob].guid)
159+
end.not_to raise_error
160+
expect(scenario[:service_instance].last_operation.reload.state).to eq('in progress')
161+
end
162+
end
163+
164+
context 'when the operation is stuck in progress' do
165+
it 'sets the operation state from in progress to failed' do
166+
scenario = prepare_stuck_service_instance
167+
op = scenario[:service_instance].last_operation
168+
169+
expect do
170+
job.send(:resolve_stuck, ServiceInstanceOperation, ServiceInstance,
171+
op.id, scenario[:service_instance].id, scenario[:pjob].guid)
172+
end.to change { op.reload.state }.from('in progress').to('failed')
173+
end
174+
175+
it 'sets the pollable job state to FAILED' do
176+
scenario = prepare_stuck_service_instance(pollable_job_state: PollableJobModel::POLLING_STATE)
177+
op = scenario[:service_instance].last_operation
178+
179+
expect do
180+
job.send(:resolve_stuck, ServiceInstanceOperation, ServiceInstance,
181+
op.id, scenario[:service_instance].id, scenario[:pjob].guid)
182+
end.to change { scenario[:pjob].reload.state }.from(PollableJobModel::POLLING_STATE).to(PollableJobModel::FAILED_STATE)
183+
end
184+
end
185+
end
186+
end
187+
end
188+
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)