From 36072174ff8d2e5fa174bb171b5ef088327586fb Mon Sep 17 00:00:00 2001 From: johha Date: Wed, 22 Jul 2026 14:15:54 +0200 Subject: [PATCH 1/4] Add async recursive delete for apps and service instances Behind `temporary_enable_async_recursive_delete flag` (default off). Recursive delete jobs re-enqueue instead of failing when sub-resources (service bindings) are deleted asynchronously, waiting for them to settle and surfacing the original broker error on failure. Adds root/sub job tracking via `root_job_guid` on the jobs table. Foundation for future recursive deletes (org, space). --- app/actions/app_delete.rb | 21 +- app/actions/app_stop.rb | 9 +- app/actions/mixins/bindings_delete.rb | 7 +- app/actions/service_instance_unshare.rb | 21 +- app/actions/v3/service_instance_delete.rb | 36 +- app/controllers/runtime/apps_controller.rb | 3 +- app/controllers/v3/apps_controller.rb | 27 +- .../v3/service_instances_controller.rb | 40 +- app/errors/sub_resource_error.rb | 41 ++ app/jobs/enqueuer.rb | 16 +- app/jobs/generic_enqueuer.rb | 8 + app/jobs/mixins/root_job_mixin.rb | 141 ++++++ app/jobs/pollable_job_wrapper.rb | 6 +- app/jobs/v3/recursive_delete_app_job.rb | 81 +++ .../recursive_delete_service_instance_job.rb | 103 ++++ app/models/runtime/pollable_job_model.rb | 5 + app/repositories/app_event_repository.rb | 5 +- config/cloud_controller.yml | 5 + ...0260723120000_add_root_job_guid_to_jobs.rb | 38 ++ .../config_schemas/api_schema.rb | 5 + .../config_schemas/worker_schema.rb | 5 + ...23120000_add_root_job_guid_to_jobs_spec.rb | 30 ++ spec/request/service_instances_spec.rb | 111 +++- spec/support/reset_generic_enqueuer.rb | 6 + .../jobs/recursive_delete_root_job.rb | 51 ++ spec/unit/actions/app_delete_spec.rb | 30 +- spec/unit/actions/app_stop_spec.rb | 13 +- .../actions/service_instance_unshare_spec.rb | 9 +- .../v3/service_instance_delete_spec.rb | 27 +- .../runtime/apps_controller_spec.rb | 2 +- .../controllers/v3/apps_controller_spec.rb | 98 ++-- spec/unit/jobs/deserialization_spec.rb | 1 + spec/unit/jobs/enqueuer_spec.rb | 70 +++ spec/unit/jobs/generic_enqueuer_spec.rb | 72 ++- spec/unit/jobs/mixins/root_job_mixin_spec.rb | 476 ++++++++++++++++++ .../mixins/root_job_orchestration_spec.rb | 240 +++++++++ spec/unit/jobs/pollable_job_wrapper_spec.rb | 24 + .../jobs/v3/recursive_delete_app_job_spec.rb | 197 ++++++++ ...ursive_delete_service_instance_job_spec.rb | 171 +++++++ .../models/runtime/pollable_job_model_spec.rb | 34 ++ .../repositories/app_event_repository_spec.rb | 16 + 41 files changed, 2140 insertions(+), 161 deletions(-) create mode 100644 app/errors/sub_resource_error.rb create mode 100644 app/jobs/mixins/root_job_mixin.rb create mode 100644 app/jobs/v3/recursive_delete_app_job.rb create mode 100644 app/jobs/v3/recursive_delete_service_instance_job.rb create mode 100644 db/migrations/20260723120000_add_root_job_guid_to_jobs.rb create mode 100644 spec/migrations/20260723120000_add_root_job_guid_to_jobs_spec.rb create mode 100644 spec/support/reset_generic_enqueuer.rb create mode 100644 spec/support/shared_examples/jobs/recursive_delete_root_job.rb create mode 100644 spec/unit/jobs/mixins/root_job_mixin_spec.rb create mode 100644 spec/unit/jobs/mixins/root_job_orchestration_spec.rb create mode 100644 spec/unit/jobs/v3/recursive_delete_app_job_spec.rb create mode 100644 spec/unit/jobs/v3/recursive_delete_service_instance_job_spec.rb diff --git a/app/actions/app_delete.rb b/app/actions/app_delete.rb index 1d7cd5a25b4..739350909d9 100644 --- a/app/actions/app_delete.rb +++ b/app/actions/app_delete.rb @@ -13,23 +13,12 @@ require 'actions/route_mapping_delete' require 'actions/staging_cancel' require 'actions/mixins/bindings_delete' +require 'errors/sub_resource_error' module VCAP::CloudController class AppDelete include V3::BindingsDeleteMixin - class AsyncBindingDeletionsTriggered < StandardError; end - - class SubResourceError < StandardError - def initialize(errors) - @errors = errors - end - - def underlying_errors - @errors - end - end - def initialize(user_audit_info) @user_audit_info = user_audit_info end @@ -86,7 +75,7 @@ def delete_subresources(app) def delete_non_transactional_subresources(app) errors = delete_bindings(app.service_bindings, user_audit_info: @user_audit_info) - raise SubResourceError.new(errors) if errors.any? + SubResourceError.raise_from(errors) end def stagers @@ -106,10 +95,8 @@ def logger @logger ||= Steno.logger('cc.action.app_delete') end - def unbinding_operation_in_progress!(binding) - raise AsyncBindingDeletionsTriggered.new( - "An operation for the service binding between app #{binding.app.name} and service instance #{binding.service_instance.name} is in progress." - ) + def unbinding_in_progress_message(binding) + "An operation for the service binding between app #{binding.app.name} and service instance #{binding.service_instance.name} is in progress." end end end diff --git a/app/actions/app_stop.rb b/app/actions/app_stop.rb index 199224ecb57..080ad3b7911 100644 --- a/app/actions/app_stop.rb +++ b/app/actions/app_stop.rb @@ -3,7 +3,7 @@ class AppStop class InvalidApp < StandardError; end class << self - def stop(app:, user_audit_info:, record_event: true) + def stop(app:, user_audit_info:, record_event: true, delete_triggered: false) app.db.transaction do app.lock! @@ -13,7 +13,7 @@ def stop(app:, user_audit_info:, record_event: true) process.update(state: ProcessModel::STOPPED) end - record_audit_event(app, user_audit_info) if record_event + record_audit_event(app, user_audit_info, delete_triggered:) if record_event end rescue Sequel::ValidationFailed => e raise InvalidApp.new(e.message) @@ -25,10 +25,11 @@ def stop_without_event(app) private - def record_audit_event(app, user_audit_info) + def record_audit_event(app, user_audit_info, delete_triggered: false) Repositories::AppEventRepository.new.record_app_stop( app, - user_audit_info + user_audit_info, + delete_triggered: ) end end diff --git a/app/actions/mixins/bindings_delete.rb b/app/actions/mixins/bindings_delete.rb index b6750daa073..49ef748d4e3 100644 --- a/app/actions/mixins/bindings_delete.rb +++ b/app/actions/mixins/bindings_delete.rb @@ -3,6 +3,7 @@ require 'jobs/generic_enqueuer' require 'jobs/v3/delete_binding_job' require 'jobs/v3/delete_service_binding_job_factory' +require 'errors/sub_resource_error' module VCAP::CloudController module V3 @@ -21,12 +22,16 @@ def delete_bindings(bindings, user_audit_info:) unless result[:finished] polling_job = DeleteBindingJob.new(type, binding.guid, user_audit_info:) Jobs::GenericEnqueuer.shared.enqueue_pollable(polling_job) - unbinding_operation_in_progress!(binding) + raise AsyncOperationInProgress.new(unbinding_in_progress_message(binding)) end rescue StandardError => e errors << e end end + + def unbinding_in_progress_message(binding) + "An operation for service binding #{binding.guid} is in progress." + end end end end diff --git a/app/actions/service_instance_unshare.rb b/app/actions/service_instance_unshare.rb index b83d7f22500..5099cd498d1 100644 --- a/app/actions/service_instance_unshare.rb +++ b/app/actions/service_instance_unshare.rb @@ -1,5 +1,6 @@ require 'actions/service_credential_binding_delete' require 'actions/mixins/bindings_delete' +require 'errors/sub_resource_error' module VCAP::CloudController class ServiceInstanceUnshare @@ -8,11 +9,15 @@ class ServiceInstanceUnshare class Error < ::StandardError end - def unshare(service_instance, target_space, user_audit_info) + def unshare(service_instance, target_space, user_audit_info, fail_if_in_progress: true) errors = delete_bindings_in_target_space!(service_instance, target_space, user_audit_info) if errors.any? - error!("Unshare of service instance failed because one or more bindings could not be deleted.\n\n " \ - "#{errors.map { |err| "\t#{err.message}" }.join("\n\n")}") + if fail_if_in_progress + raise Error.new("Unshare of service instance failed because one or more bindings could not be deleted.\n\n " \ + "#{errors.map { |err| "\t#{err.message}" }.join("\n\n")}") + else + SubResourceError.raise_from(errors) + end end service_instance.remove_shared_space(target_space) @@ -24,19 +29,15 @@ def unshare(service_instance, target_space, user_audit_info) private - def error!(message) - raise Error.new(message) - end - def delete_bindings_in_target_space!(service_instance, target_space, user_audit_info) active_bindings = ServiceBinding.where(service_instance_guid: service_instance.guid) bindings_in_target_space = active_bindings.all.select { |b| b.app.space_guid == target_space.guid } delete_bindings(bindings_in_target_space, user_audit_info:) end - def unbinding_operation_in_progress!(binding) - raise Error.new("The binding between an application and service instance #{binding.service_instance.name} " \ - "in space #{binding.app.space.name} is being deleted asynchronously.") + def unbinding_in_progress_message(binding) + "The binding between an application and service instance #{binding.service_instance.name} " \ + "in space #{binding.app.space.name} is being deleted asynchronously." end end end diff --git a/app/actions/v3/service_instance_delete.rb b/app/actions/v3/service_instance_delete.rb index 0b4f839c739..4be55120ba3 100644 --- a/app/actions/v3/service_instance_delete.rb +++ b/app/actions/v3/service_instance_delete.rb @@ -4,6 +4,7 @@ require 'actions/service_instance_unshare' require 'cloud_controller/errors/api_error' require 'actions/mixins/bindings_delete' +require 'errors/sub_resource_error' module VCAP::CloudController module V3 @@ -13,9 +14,6 @@ class ServiceInstanceDelete class DeleteFailed < StandardError end - class UnbindingOperatationInProgress < StandardError - end - DeleteStatus = Struct.new(:finished, :operation).freeze DeleteStarted = ->(operation) { DeleteStatus.new(false, operation) } DeleteComplete = DeleteStatus.new(true, nil).freeze @@ -24,9 +22,10 @@ class UnbindingOperatationInProgress < StandardError PollingFinished = PollingStatus.new(true, nil).freeze ContinuePolling = ->(retry_after) { PollingStatus.new(false, retry_after) } - def initialize(service_instance, event_repo) + def initialize(service_instance, event_repo, fail_if_in_progress: true) @service_instance = service_instance @service_event_repository = event_repo + @fail_if_in_progress = fail_if_in_progress end def blocking_operation_in_progress? @@ -38,7 +37,11 @@ def delete operation_in_progress! if blocking_operation_in_progress? errors = remove_associations - raise errors.first if errors.any? + if errors.any? + raise errors.first if @fail_if_in_progress # Single-shot callers fail on the first error + + SubResourceError.raise_from(errors) + end result = send_deprovison_to_broker if result[:finished] @@ -49,6 +52,11 @@ def delete end result + rescue SubResourceError => e + raise if !@fail_if_in_progress && e.any_in_progress? # re-raise SubResourceError so that root job continues to run + + update_last_operation_with_failure(e.message) unless service_instance.operation_in_progress? + raise e rescue StandardError => e update_last_operation_with_failure(e.message) unless service_instance.operation_in_progress? raise e @@ -154,7 +162,9 @@ def unshare_all_spaces unshare_action = ServiceInstanceUnshare.new space_guids.each_with_object([]) do |space_guid, errors| - unshare_action.unshare(service_instance, Space.first(guid: space_guid), service_event_repository.user_audit_info) + unshare_action.unshare(service_instance, Space.first(guid: space_guid), service_event_repository.user_audit_info, fail_if_in_progress: @fail_if_in_progress) + rescue SubResourceError => e + errors.concat(e.underlying_errors) rescue StandardError => e errors << e end @@ -182,14 +192,12 @@ def operation_in_progress! raise CloudController::Errors::ApiError.new_from_details('AsyncServiceInstanceOperationInProgress', service_instance.name) end - def unbinding_operation_in_progress!(binding) - raise UnbindingOperatationInProgress.new( - if binding.is_a?(VCAP::CloudController::ServiceBinding) - "An operation for the service binding between app #{binding.app.name} and service instance #{service_instance.name} is in progress." - else - "An operation for a service binding of service instance #{service_instance.name} is in progress." - end - ) + def unbinding_in_progress_message(binding) + if binding.is_a?(VCAP::CloudController::ServiceBinding) + "An operation for the service binding between app #{binding.app.name} and service instance #{service_instance.name} is in progress." + else + "An operation for a service binding of service instance #{service_instance.name} is in progress." + end end def delete_failed!(message) diff --git a/app/controllers/runtime/apps_controller.rb b/app/controllers/runtime/apps_controller.rb index d57f39b5775..40f9d3a29c0 100644 --- a/app/controllers/runtime/apps_controller.rb +++ b/app/controllers/runtime/apps_controller.rb @@ -6,6 +6,7 @@ require 'actions/v2/route_mapping_create' require 'models/helpers/process_types' require 'controllers/runtime/mixins/find_process_through_app' +require 'errors/sub_resource_error' module VCAP::CloudController class AppsController < RestController::ModelController @@ -148,7 +149,7 @@ def delete(guid) AppDelete.new(UserAuditInfo.from_context(SecurityContext)).delete_without_event([process.app]) rescue Sequel::NoExistingObject raise self.class.not_found_exception(guid, AppModel) - rescue AppDelete::SubResourceError => e + rescue SubResourceError => e error_message = e.underlying_errors.map { |err| "\t" + err.message }.join("\n") raise CloudController::Errors::ApiError.new_from_details('AppRecursiveDeleteFailed', process.app.name, error_message) end diff --git a/app/controllers/v3/apps_controller.rb b/app/controllers/v3/apps_controller.rb index 69e19dbcb93..500d354d022 100644 --- a/app/controllers/v3/apps_controller.rb +++ b/app/controllers/v3/apps_controller.rb @@ -5,11 +5,13 @@ require 'actions/app_update' require 'actions/app_patch_environment_variables' require 'actions/app_delete' +require 'errors/sub_resource_error' require 'actions/app_restart' require 'actions/app_apply_manifest' require 'actions/app_start' require 'actions/app_stop' require 'actions/app_assign_droplet' +require 'jobs/v3/recursive_delete_app_job' require 'decorators/include_space_decorator' require 'decorators/include_organization_decorator' require 'decorators/include_space_organization_decorator' @@ -154,12 +156,23 @@ def destroy unauthorized! unless permission_queryer.can_write_to_active_space?(space.id) require_writable_space!(space) - delete_action = AppDelete.new(user_audit_info) - deletion_job = VCAP::CloudController::Jobs::DeleteActionJob.new(AppModel, app.guid, delete_action) + if async_recursive_delete_enabled? + job = Jobs::Enqueuer.new(queue: Jobs::Queues.generic).enqueue_or_find_active_pollable( + resource_model: AppModel, + resource_guid: app.guid, + operation: 'app.delete' + ) { VCAP::CloudController::V3::RecursiveDeleteAppJob.new(app.guid, user_audit_info) } - job = Jobs::Enqueuer.new(queue: Jobs::Queues.generic).enqueue_pollable(deletion_job) do |pollable_job| - DeleteAppErrorTranslatorJob.new(pollable_job) + app_not_found! unless job + else + delete_action = AppDelete.new(user_audit_info) + deletion_job = VCAP::CloudController::Jobs::DeleteActionJob.new(AppModel, app.guid, delete_action) + + job = Jobs::Enqueuer.new(queue: Jobs::Queues.generic).enqueue_pollable(deletion_job) do |pollable_job| + DeleteAppErrorTranslatorJob.new(pollable_job) + end end + VCAP::AppLogEmitter.emit(app.guid, "Enqueued job to delete app with guid #{app.guid}") head HTTP::ACCEPTED, 'Location' => url_builder.build_url(path: "/v3/jobs/#{job.guid}") end @@ -372,7 +385,7 @@ class DeleteAppErrorTranslatorJob < VCAP::CloudController::Jobs::ErrorTranslator include V3ErrorsHelper def translate_error(e) - if e.instance_of?(VCAP::CloudController::AppDelete::SubResourceError) + if e.instance_of?(VCAP::CloudController::SubResourceError) underlying_errors = e.underlying_errors.map { |err| unprocessable(err.message) } e = CloudController::Errors::CompoundError.new(underlying_errors) end @@ -382,6 +395,10 @@ def translate_error(e) private + def async_recursive_delete_enabled? + !!VCAP::CloudController::Config.config.get(:temporary_enable_async_recursive_delete, :apps) + end + def handle_order_by_presented_value(page_results) return unless page_results.try(:pagination_options).try(:order_by) == 'desired_state' diff --git a/app/controllers/v3/service_instances_controller.rb b/app/controllers/v3/service_instances_controller.rb index e5663120484..a605ce3b3b8 100644 --- a/app/controllers/v3/service_instances_controller.rb +++ b/app/controllers/v3/service_instances_controller.rb @@ -29,6 +29,7 @@ require 'decorators/field_service_instance_plan_decorator' require 'jobs/v3/create_service_instance_job' require 'jobs/v3/update_service_instance_job' +require 'jobs/v3/recursive_delete_service_instance_job' class ServiceInstancesV3Controller < ApplicationController include ServicePermissions @@ -115,13 +116,33 @@ def destroy end delete_action = V3::ServiceInstanceDelete.new(service_instance, service_event_repository) - operation_in_progress! if delete_action.blocking_operation_in_progress? case service_instance when VCAP::CloudController::ManagedServiceInstance - job_guid = enqueue_delete_job(service_instance) - head :accepted, 'Location' => url_builder.build_url(path: "/v3/jobs/#{job_guid}") + if async_recursive_delete_enabled? + # Idempotent retry: return the active delete job instead of 422-ing on its delete-in-progress last_operation. + active_delete = PollableJobModel.find_active_delete(resource_guid: service_instance.guid, operation: 'service_instance.delete') + operation_in_progress! if active_delete.nil? && delete_action.blocking_operation_in_progress? + + job = Jobs::Enqueuer.new(queue: Jobs::Queues.generic).enqueue_or_find_active_pollable( + resource_model: ManagedServiceInstance, + resource_guid: service_instance.guid, + operation: 'service_instance.delete' + ) do |locked_instance| + log_service_instance_deletion(locked_instance) + V3::RecursiveDeleteServiceInstanceJob.new(locked_instance.guid, user_audit_info) + end + + service_instance_not_found! unless job + + head :accepted, 'Location' => url_builder.build_url(path: "/v3/jobs/#{job.guid}") + else + operation_in_progress! if delete_action.blocking_operation_in_progress? + job_guid = enqueue_delete_job(service_instance) + head :accepted, 'Location' => url_builder.build_url(path: "/v3/jobs/#{job_guid}") + end when VCAP::CloudController::UserProvidedServiceInstance + operation_in_progress! if delete_action.blocking_operation_in_progress? delete_action.delete head :no_content end @@ -391,9 +412,7 @@ def fetch_writable_service_instance(guid) service_instance end - def enqueue_delete_job(service_instance) - delete_job = V3::DeleteServiceInstanceJob.new(service_instance.guid, user_audit_info) - + def log_service_instance_deletion(service_instance) plan = service_instance.service_plan service = plan.service broker = service.service_broker @@ -404,11 +423,20 @@ def enqueue_delete_job(service_instance) "from service offering '#{service.label}' " \ "provided by broker '#{broker.name}'." ) + end + # Legacy enqueue path used when the async-recursive-delete rollout flag is off. Preserves main's behaviour. + def enqueue_delete_job(service_instance) + log_service_instance_deletion(service_instance) + delete_job = V3::DeleteServiceInstanceJob.new(service_instance.guid, user_audit_info) pollable_job = Jobs::Enqueuer.new(queue: Jobs::Queues.generic).enqueue_pollable(delete_job) pollable_job.guid end + def async_recursive_delete_enabled? + !!VCAP::CloudController::Config.config.get(:temporary_enable_async_recursive_delete, :service_instances) + end + def unreadable_error_message(service_instance_name, unreadable_space_guids) return unless unreadable_space_guids.any? diff --git a/app/errors/sub_resource_error.rb b/app/errors/sub_resource_error.rb new file mode 100644 index 00000000000..3584cc6a002 --- /dev/null +++ b/app/errors/sub_resource_error.rb @@ -0,0 +1,41 @@ +require 'cloud_controller/errors/api_error' +require 'cloud_controller/errors/compound_error' + +module VCAP::CloudController + class AsyncOperationInProgress < StandardError; end + + class SubResourceError < StandardError + attr_reader :errors + + def initialize(errors) + super() + @errors = errors + end + + def underlying_errors + @errors + end + + def failures + @errors.reject { |e| e.is_a?(AsyncOperationInProgress) } + end + + def in_progress_operations + @errors.select { |e| e.is_a?(AsyncOperationInProgress) } + end + + def any_in_progress? + in_progress_operations.any? + end + + def message + @errors.map(&:message).join("\n") + end + + def self.raise_from(errors) + return if errors.empty? + + raise new(errors) + end + end +end diff --git a/app/jobs/enqueuer.rb b/app/jobs/enqueuer.rb index e52b6a613f6..d4afe2d328f 100644 --- a/app/jobs/enqueuer.rb +++ b/app/jobs/enqueuer.rb @@ -9,6 +9,8 @@ module VCAP::CloudController module Jobs class Enqueuer + attr_accessor :root_job_guid + def initialize(opts={}) @opts = opts @timeout_calculator = JobTimeoutCalculator.new(VCAP::CloudController::Config.config) @@ -21,7 +23,7 @@ def enqueue(job, run_at: nil, priority_increment: nil) end def enqueue_pollable(job, existing_guid: nil, run_at: nil, priority_increment: nil, preserve_priority: false) - wrapped_job = PollableJobWrapper.new(job, existing_guid:) + wrapped_job = PollableJobWrapper.new(job, existing_guid:, root_job_guid:) wrapped_job = yield wrapped_job if block_given? @@ -29,6 +31,18 @@ def enqueue_pollable(job, existing_guid: nil, run_at: nil, priority_increment: n PollableJobModel.find_by_delayed_job(delayed_job) end + def enqueue_or_find_active_pollable(resource_model:, resource_guid:, operation:) + resource_model.db.transaction do + resource = resource_model.where(guid: resource_guid).for_update.first + return nil unless resource + + existing = PollableJobModel.find_active_delete(resource_guid:, operation:) + return existing if existing + + enqueue_pollable(yield(resource)) + end + end + def self.unwrap_job(job) job.is_a?(WrappingJob) ? unwrap_job(job.handler) : job end diff --git a/app/jobs/generic_enqueuer.rb b/app/jobs/generic_enqueuer.rb index 5cd8ba408dd..053a8185c80 100644 --- a/app/jobs/generic_enqueuer.rb +++ b/app/jobs/generic_enqueuer.rb @@ -15,6 +15,14 @@ def self.shared(priority: nil) def self.reset! Thread.current[:generic_enqueuer] = nil end + + def activate_root_context(root_job_guid:) + self.root_job_guid = root_job_guid + end + + def deactivate_root_context + self.root_job_guid = nil + end end end end diff --git a/app/jobs/mixins/root_job_mixin.rb b/app/jobs/mixins/root_job_mixin.rb new file mode 100644 index 00000000000..d515a66e80c --- /dev/null +++ b/app/jobs/mixins/root_job_mixin.rb @@ -0,0 +1,141 @@ +require 'errors/sub_resource_error' +require 'cloud_controller/errors/api_error' +require 'cloud_controller/errors/compound_error' + +module VCAP::CloudController + module Jobs + module RootJobMixin + # Buffer added on top of the sub-jobs' next run_at so the root wakes just after them, never before. + ROOT_JOB_BUFFER_SECONDS = 5 + + private + + def perform_with_root_job_handling + activate_root_job_context + yield + rescue SubResourceError => e + return if e.any_in_progress? + + raise compound_error_for(e.failures) + rescue CloudController::Errors::ApiError, CloudController::Errors::CompoundError + raise + rescue StandardError => e + raise CloudController::Errors::ApiError.new_from_details('UnableToPerform', 'delete', e.message) + ensure + deactivate_root_job_context + end + + attr_reader :root_job, :sub_jobs + + # One shared tag for the whole recursive-delete subsystem so ops can grep it by a single logger name. + def logger + @logger ||= Steno.logger('cc.jobs.v3.recursive_delete') + end + + # Resolves in any state (unlike root_job), so log lines stay searchable by this guid after the job settles. + def pollable_job_guid + @pollable_job_guid ||= PollableJobModel.first(resource_guid: resource_guid, operation: display_name)&.guid + end + + def activate_root_job_context + fetch_root_context + Jobs::GenericEnqueuer.shared.activate_root_context(root_job_guid: root_job&.guid) + end + + def deactivate_root_job_context + # Must clear in perform's ensure: the job YAML-serialises itself on reschedule, reviving a stale cache otherwise. + @root_job = nil + @sub_jobs = nil + Jobs::GenericEnqueuer.shared.deactivate_root_context + end + + def fetch_root_context + @root_job = PollableJobModel.find_active_delete(resource_guid: resource_guid, operation: display_name) + @sub_jobs = @root_job ? @root_job.sub_jobs : [] + end + + def active_sub_jobs + sub_jobs.select { |s| [PollableJobModel::PROCESSING_STATE, PollableJobModel::POLLING_STATE].include?(s.state) } + end + + # Pace off the slowest active sub-job's next run (else the normal interval) so the root never re-runs early. + def next_execution_in + interval = (seconds_until_slowest_sub_job || super) + ROOT_JOB_BUFFER_SECONDS + [interval, Config.config.get(:broker_client_max_async_poll_interval_seconds)].min + end + + def seconds_until_slowest_sub_job + job = PollableJobModel.find_active_delete(resource_guid: resource_guid, operation: display_name) + return nil unless job + + active_guids = job.sub_jobs_dataset.where(state: [PollableJobModel::PROCESSING_STATE, PollableJobModel::POLLING_STATE]).select_map(:delayed_job_guid) + return nil if active_guids.empty? + + latest = Delayed::Job.where(guid: active_guids).max(:run_at) + now = Delayed::Job.db_time_now + return nil unless latest && latest > now + + (latest - now).ceil + end + + def sub_jobs_in_flight? + return false if active_sub_jobs.empty? + + add_in_progress_warning(root_job) + true + end + + def raise_if_sub_jobs_failed + return if sub_job_errors.empty? + + raise CloudController::Errors::CompoundError.new(all_failure_errors) + end + + def add_in_progress_warning(job) + return if job.warnings_dataset.any? + + JobWarningModel.create(job: job, detail: in_progress_warning_detail) + rescue Sequel::Error => e + logger.warn("could not add in-progress warning for #{resource_type} #{resource_guid} (job #{job.guid}): #{e.message}") + end + + def in_progress_warning_detail + 'This operation is still in progress: it is waiting for one or more dependent operations to finish.' + end + + def compound_error_for(raised_failures) + errors = all_failure_errors + errors = raised_failures.map { |e| CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', e.message) } if errors.empty? + CloudController::Errors::CompoundError.new(errors) + end + + def all_failure_errors + by_guid = {} + sub_resource_errors.each { |guid, err| by_guid[guid] = err } + sub_job_errors.each { |guid, err| by_guid[guid] ||= err } + by_guid.values + end + + def sub_job_errors + sub_jobs.select { |s| s.state == PollableJobModel::FAILED_STATE }.map do |sub_job| + [sub_job.resource_guid, CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', sub_job_error_detail(sub_job))] + end + end + + def sub_resource_errors + [] + end + + def sub_job_error_detail(sub_job) + fallback = "#{sub_job.resource_type} #{sub_job.resource_guid}" + return fallback if sub_job.cf_api_error.nil? + + parsed = Psych.safe_load(sub_job.cf_api_error, strict_integer: true) + detail = parsed && parsed['errors']&.first&.fetch('detail', nil) + detail.presence || fallback + rescue Psych::Exception + fallback + end + end + end +end diff --git a/app/jobs/pollable_job_wrapper.rb b/app/jobs/pollable_job_wrapper.rb index 6803ba49133..b1ea5c482b1 100644 --- a/app/jobs/pollable_job_wrapper.rb +++ b/app/jobs/pollable_job_wrapper.rb @@ -7,8 +7,9 @@ module Jobs class PollableJobWrapper < WrappingJob attr_reader :existing_guid - def initialize(handler, existing_guid: nil) + def initialize(handler, existing_guid: nil, root_job_guid: nil) @existing_guid = existing_guid + @root_job_guid = root_job_guid super(handler) end @@ -36,7 +37,8 @@ def before_enqueue(job) operation: @handler.display_name, resource_guid: @handler.resource_guid, resource_type: @handler.resource_type, - user_guid: user_guid + user_guid: user_guid, + root_job_guid: @root_job_guid ) end end diff --git a/app/jobs/v3/recursive_delete_app_job.rb b/app/jobs/v3/recursive_delete_app_job.rb new file mode 100644 index 00000000000..fadd9533aaa --- /dev/null +++ b/app/jobs/v3/recursive_delete_app_job.rb @@ -0,0 +1,81 @@ +require 'jobs/reoccurring_job' +require 'jobs/mixins/root_job_mixin' +require 'actions/app_delete' +require 'actions/app_stop' + +module VCAP::CloudController + module V3 + class RecursiveDeleteAppJob < Jobs::ReoccurringJob + include Jobs::RootJobMixin + + attr_reader :app_guid + + def initialize(app_guid, user_audit_info) + super() + @app_guid = app_guid + @user_audit_info = user_audit_info + end + + def perform + perform_with_root_job_handling do + if sub_jobs_in_flight? + logger.info("app delete #{app_guid} (job #{pollable_job_guid}) waiting on in-progress service binding deletions") + return + end + + log_failed_bindings + raise_if_sub_jobs_failed + + app = AppModel.first(guid: app_guid) + return finish unless app + + AppStop.stop(app: app, user_audit_info: @user_audit_info, delete_triggered: true) if app.desired_state != ProcessModel::STOPPED + AppDelete.new(@user_audit_info).delete([app]) + finish + end + end + + def handle_timeout; end + + def resource_guid + app_guid + end + + def resource_type + 'app' + end + + def display_name + 'app.delete' + end + + def max_attempts + 1 + end + + private + + attr_reader :user_audit_info + + def log_failed_bindings + sub_resource_errors.each do |guid, error| + logger.warn("app delete #{app_guid} (job #{pollable_job_guid}): service binding #{guid} deletion failed: #{error.message}") + end + end + + def in_progress_warning_detail + 'Deletion of the app is still in progress: one or more service bindings are still being deleted. ' \ + 'It will complete once those operations finish.' + end + + def sub_resource_errors + app = AppModel.first(guid: app_guid) + return [] unless app + + app.service_bindings.select(&:delete_failed?).map do |binding| + [binding.guid, CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', binding.last_operation.description)] + end + end + end + end +end diff --git a/app/jobs/v3/recursive_delete_service_instance_job.rb b/app/jobs/v3/recursive_delete_service_instance_job.rb new file mode 100644 index 00000000000..9b2aa7d9460 --- /dev/null +++ b/app/jobs/v3/recursive_delete_service_instance_job.rb @@ -0,0 +1,103 @@ +require 'jobs/reoccurring_job' +require 'jobs/mixins/root_job_mixin' +require 'actions/v3/service_instance_delete' + +module VCAP::CloudController + module V3 + class RecursiveDeleteServiceInstanceJob < VCAP::CloudController::Jobs::ReoccurringJob + include Jobs::RootJobMixin + + attr_reader :resource_guid + + def initialize(guid, user_audit_info) + super() + @resource_guid = guid + @user_audit_info = user_audit_info + end + + def perform + perform_with_root_job_handling do + if sub_jobs_in_flight? + logger.info("service instance delete #{resource_guid} (job #{pollable_job_guid}) waiting on in-progress binding deletions") + return + end + + log_failed_children + raise_if_sub_jobs_failed + + return finish unless service_instance + + self.maximum_duration_seconds = service_instance.service_plan.try(:maximum_polling_duration) + + unless delete_in_progress? + result = action.delete + return finish if result[:finished] + end + + result = action.poll + return finish if result[:finished] + + self.polling_interval_seconds = result[:retry_after].to_i if result[:retry_after] + end + end + + def handle_timeout + action.update_last_operation_with_failure("Service Broker failed to #{operation} within the required time.") + end + + def operation + :deprovision + end + + def operation_type + 'delete' + end + + def resource_type + 'service_instance' + end + + def display_name + "#{resource_type}.#{operation_type}" + end + + private + + attr_reader :user_audit_info + + def in_progress_warning_detail + 'Deletion of the service instance is still in progress: one or more bindings are still being ' \ + 'deleted. It will complete once those operations finish.' + end + + def service_instance + ManagedServiceInstance.first(guid: resource_guid) + end + + def delete_in_progress? + service_instance.last_operation&.type == 'delete' && + service_instance.last_operation&.state == 'in progress' + end + + def action + ServiceInstanceDelete.new(service_instance, Repositories::ServiceEventRepository.new(user_audit_info), fail_if_in_progress: false) + end + + def log_failed_children + sub_resource_errors.each do |guid, error| + logger.warn("service instance delete #{resource_guid} (job #{pollable_job_guid}): binding #{guid} deletion failed: #{error.message}") + end + end + + def sub_resource_errors + si = service_instance + return [] unless si + + children = si.service_bindings + si.service_keys + RouteBinding.where(service_instance: si).all + children.select(&:delete_failed?).map do |child| + [child.guid, CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', child.last_operation.description)] + end + end + end + end +end diff --git a/app/models/runtime/pollable_job_model.rb b/app/models/runtime/pollable_job_model.rb index 29dbf4a44aa..d252de16d0a 100644 --- a/app/models/runtime/pollable_job_model.rb +++ b/app/models/runtime/pollable_job_model.rb @@ -6,6 +6,7 @@ class PollableJobModel < Sequel::Model(:jobs) POLLING_STATE = 'POLLING'.freeze one_to_many :warnings, class: 'VCAP::CloudController::JobWarningModel', key: :job_id + one_to_many :sub_jobs, class: 'VCAP::CloudController::PollableJobModel', key: :root_job_guid, primary_key: :guid plugin :serialization add_association_dependencies warnings: :destroy @@ -47,6 +48,10 @@ def self.find_by_delayed_job_guid(delayed_job_guid) pollable_job end + def self.find_active_delete(resource_guid:, operation:) + PollableJobModel.first(resource_guid: resource_guid, operation: operation, state: [PROCESSING_STATE, POLLING_STATE]) + end + def self.number_of_active_jobs_by_user(user_guid) PollableJobModel.where(state: %w[PROCESSING POLLING], user_guid: user_guid).count end diff --git a/app/repositories/app_event_repository.rb b/app/repositories/app_event_repository.rb index 4b1c99c14a7..947f8008e19 100644 --- a/app/repositories/app_event_repository.rb +++ b/app/repositories/app_event_repository.rb @@ -75,11 +75,12 @@ def record_app_restart(app, user_audit_info) create_app_audit_event(EventTypes::APP_RESTART, app, app.space, actor, nil) end - def record_app_stop(app, user_audit_info) + def record_app_stop(app, user_audit_info, delete_triggered: false) VCAP::AppLogEmitter.emit(app.guid, "Stopping app with guid #{app.guid}") actor = { name: user_audit_info.user_email, guid: user_audit_info.user_guid, user_name: user_audit_info.user_name, type: 'user' } - create_app_audit_event(EventTypes::APP_STOP, app, app.space, actor, nil) + metadata = delete_triggered ? { delete_triggered: true } : nil + create_app_audit_event(EventTypes::APP_STOP, app, app.space, actor, metadata) end def record_app_delete_request(app, space, user_audit_info, recursive=nil) diff --git a/config/cloud_controller.yml b/config/cloud_controller.yml index c2e7d248ab8..2dc252a8322 100644 --- a/config/cloud_controller.yml +++ b/config/cloud_controller.yml @@ -285,6 +285,11 @@ rate_limiter_v2_api: temporary_enable_v2: true +# Rollout gate: when enabled, DELETE of apps/service_instances uses the recursive root-job pipeline. +temporary_enable_async_recursive_delete: + apps: false + service_instances: false + max_concurrent_service_broker_requests: 0 diego: diff --git a/db/migrations/20260723120000_add_root_job_guid_to_jobs.rb b/db/migrations/20260723120000_add_root_job_guid_to_jobs.rb new file mode 100644 index 00000000000..e5f0db7a975 --- /dev/null +++ b/db/migrations/20260723120000_add_root_job_guid_to_jobs.rb @@ -0,0 +1,38 @@ +Sequel.migration do + no_transaction + + up do + if database_type == :postgres + add_column :jobs, :root_job_guid, String, size: 255, if_not_exists: true + VCAP::Migration.with_concurrent_timeout(self) do + add_index :jobs, :root_job_guid, name: :jobs_root_job_guid_index, if_not_exists: true, concurrently: true + end + + elsif database_type == :mysql + alter_table :jobs do + add_column :root_job_guid, String, size: 255 unless @db.schema(:jobs).map(&:first).include?(:root_job_guid) + # rubocop:disable Sequel/ConcurrentIndex + add_index :root_job_guid, name: :jobs_root_job_guid_index unless @db.indexes(:jobs).include?(:jobs_root_job_guid_index) + # rubocop:enable Sequel/ConcurrentIndex + end + end + end + + down do + if database_type == :postgres + VCAP::Migration.with_concurrent_timeout(self) do + drop_index :jobs, :root_job_guid, name: :jobs_root_job_guid_index, if_exists: true, concurrently: true + end + drop_column :jobs, :root_job_guid, if_exists: true + end + + if database_type == :mysql + alter_table :jobs do + # rubocop:disable Sequel/ConcurrentIndex + drop_index :root_job_guid, name: :jobs_root_job_guid_index if @db.indexes(:jobs).include?(:jobs_root_job_guid_index) + # rubocop:enable Sequel/ConcurrentIndex + drop_column :root_job_guid if @db.schema(:jobs).map(&:first).include?(:root_job_guid) + end + end + end +end diff --git a/lib/cloud_controller/config_schemas/api_schema.rb b/lib/cloud_controller/config_schemas/api_schema.rb index ad57647eea9..172cd6b51a4 100644 --- a/lib/cloud_controller/config_schemas/api_schema.rb +++ b/lib/cloud_controller/config_schemas/api_schema.rb @@ -401,6 +401,11 @@ class ApiSchema < VCAP::Config optional(:temporary_enable_v2) => bool, + optional(:temporary_enable_async_recursive_delete) => { + optional(:apps) => bool, + optional(:service_instances) => bool + }, + allow_app_ssh_access: bool, optional(:external_host) => String, diff --git a/lib/cloud_controller/config_schemas/worker_schema.rb b/lib/cloud_controller/config_schemas/worker_schema.rb index f4f7e753964..ac3c3e6eaa7 100644 --- a/lib/cloud_controller/config_schemas/worker_schema.rb +++ b/lib/cloud_controller/config_schemas/worker_schema.rb @@ -23,6 +23,11 @@ class WorkerSchema < VCAP::Config optional(:temporary_enable_v2) => bool, + optional(:temporary_enable_async_recursive_delete) => { + optional(:apps) => bool, + optional(:service_instances) => bool + }, + uaa: { internal_url: String, optional(:ca_file) => String, diff --git a/spec/migrations/20260723120000_add_root_job_guid_to_jobs_spec.rb b/spec/migrations/20260723120000_add_root_job_guid_to_jobs_spec.rb new file mode 100644 index 00000000000..8dd63190f22 --- /dev/null +++ b/spec/migrations/20260723120000_add_root_job_guid_to_jobs_spec.rb @@ -0,0 +1,30 @@ +require 'spec_helper' +require 'migrations/helpers/migration_shared_context' + +RSpec.describe 'migration to add root_job_guid column to jobs table', isolation: :truncation, type: :migration do + include_context 'migration' do + let(:migration_filename) { '20260723120000_add_root_job_guid_to_jobs.rb' } + end + + describe 'jobs table' do + it 'adds column and index, and handles idempotency gracefully' do + expect(db[:jobs].columns).not_to include(:root_job_guid) + expect(db.indexes(:jobs)).not_to be_key(:jobs_root_job_guid_index) + expect { Sequel::Migrator.run(db, migrations_path, target: current_migration_index, allow_missing_migration_files: true) }.not_to raise_error + expect(db[:jobs].columns).to include(:root_job_guid) + expect(db.indexes(:jobs)).to be_key(:jobs_root_job_guid_index) + + expect { Sequel::Migrator.run(db, migrations_path, target: current_migration_index, allow_missing_migration_files: true) }.not_to raise_error + expect(db[:jobs].columns).to include(:root_job_guid) + expect(db.indexes(:jobs)).to be_key(:jobs_root_job_guid_index) + + expect { Sequel::Migrator.run(db, migrations_path, target: current_migration_index - 1, allow_missing_migration_files: true) }.not_to raise_error + expect(db[:jobs].columns).not_to include(:root_job_guid) + expect(db.indexes(:jobs)).not_to be_key(:jobs_root_job_guid_index) + + expect { Sequel::Migrator.run(db, migrations_path, target: current_migration_index - 1, allow_missing_migration_files: true) }.not_to raise_error + expect(db[:jobs].columns).not_to include(:root_job_guid) + expect(db.indexes(:jobs)).not_to be_key(:jobs_root_job_guid_index) + end + end +end diff --git a/spec/request/service_instances_spec.rb b/spec/request/service_instances_spec.rb index 0bbab9e9c12..6621adefdb6 100644 --- a/spec/request/service_instances_spec.rb +++ b/spec/request/service_instances_spec.rb @@ -3086,6 +3086,7 @@ def check_filtered_instances(*instances) before do allow(Steno).to receive(:logger).and_call_original allow(Steno).to receive(:logger).with('cc.api').and_return(mock_logger) + TestConfig.override(temporary_enable_async_recursive_delete: { apps: true, service_instances: true }) end it 'responds with job resource' do @@ -3101,6 +3102,89 @@ def check_filtered_instances(*instances) expect(job.resource_type).to eq('service_instance') end + context 'when temporary_enable_async_recursive_delete.service_instances is disabled (legacy path)' do + before { TestConfig.override(temporary_enable_async_recursive_delete: { apps: false, service_instances: false }) } + + it 'still enqueues a pollable service_instance.delete job (legacy shape)' do + api_call.call(admin_headers) + expect(last_response).to have_status_code(202) + + job = VCAP::CloudController::PollableJobModel.last + expect(job.operation).to eq('service_instance.delete') + expect(job.resource_guid).to eq(instance.guid) + expect(last_response.headers['Location']).to end_with("/v3/jobs/#{job.guid}") + end + end + + context 'when a delete is already in flight for the instance' do + let!(:existing_job) do + create(:pollable_job_model, + state: VCAP::CloudController::PollableJobModel::PROCESSING_STATE, + resource_guid: instance.guid, + operation: 'service_instance.delete') + end + + it 'returns 202 pointing at the existing job and does not enqueue a second one' do + expect do + api_call.call(admin_headers) + end.not_to(change(VCAP::CloudController::PollableJobModel, :count)) + + expect(last_response).to have_status_code(202) + expect(last_response.headers['Location']).to end_with("/v3/jobs/#{existing_job.guid}") + end + + context 'and last_operation is delete/in progress (retry during async deprovision)' do + before do + instance.save_with_new_operation({}, { type: 'delete', state: 'in progress', description: 'draining bindings' }) + end + + it 'returns 202 pointing at the existing job rather than 422' do + expect do + api_call.call(admin_headers) + end.not_to(change(VCAP::CloudController::PollableJobModel, :count)) + + expect(last_response).to have_status_code(202) + expect(last_response.headers['Location']).to end_with("/v3/jobs/#{existing_job.guid}") + end + end + end + + context 'when a create operation is initial (no delete job)' do + before do + instance.save_with_new_operation({}, { type: 'create', state: 'initial' }) + end + + it 'returns 422 with an operation-in-progress error' do + api_call.call(admin_headers) + expect(last_response).to have_status_code(422) + expect(parsed_response['errors'].first['detail']).to include('operation in progress') + end + end + + context 'when an update operation is in progress (no delete job)' do + before do + instance.save_with_new_operation({}, { type: 'update', state: 'in progress' }) + end + + it 'returns 422 with an operation-in-progress error' do + api_call.call(admin_headers) + expect(last_response).to have_status_code(422) + expect(parsed_response['errors'].first['detail']).to include('operation in progress') + end + end + + context 'when last_operation is delete/in progress but there is no active job (stale state)' do + before do + instance.save_with_new_operation({}, { type: 'delete', state: 'in progress' }) + end + + it 'returns 422 with an operation-in-progress error (preserves safety net for orphaned state)' do + api_call.call(admin_headers) + expect(last_response).to have_status_code(422) + expect(parsed_response['errors'].first['detail']).to include('operation in progress') + end + end + it 'logs the correct names when deleting a managed service instance' do api_call.call(admin_headers) @@ -3433,14 +3517,12 @@ def check_filtered_instances(*instances) to_return(status: 202, body: '{}', headers: {}) end - it 'fails when the unbind is async' do + it 'enqueues an async unbind sub-job and defers the SI delete' do api_call.call(admin_headers) - execute_all_jobs(expected_successes: 0, expected_failures: 1, jobs_to_execute: 1) + # The root job enqueues an unbind sub-job then defers (re-enqueues itself), counting as 1 success. + execute_all_jobs(expected_successes: 1, expected_failures: 0, jobs_to_execute: 1) - lo = instance.last_operation - expect(lo.type).to eq('delete') - expect(lo.state).to eq('failed') - expect(lo.description).to eq("An operation for the service binding between app #{application.name} and service instance #{instance.name} is in progress.") + expect(VCAP::CloudController::ServiceInstance.first(guid: instance.guid)).not_to be_nil expect( stub_request(:delete, "#{instance.service_broker.broker_url}/v2/service_instances/#{instance.guid}/service_bindings/#{service_binding.guid}"). @@ -3517,14 +3599,14 @@ def check_filtered_instances(*instances) end end - it 'fails and starts the delete operation on the bindings' do + it 'defers and starts the delete operation on the bindings' do api_call.call(admin_headers) - execute_all_jobs(expected_successes: 0, expected_failures: 1, jobs_to_execute: 1) + # The root job enqueues 3 unbind sub-jobs then defers (re-enqueues itself), counting as 1 success. + execute_all_jobs(expected_successes: 1, expected_failures: 0, jobs_to_execute: 1) - lo = VCAP::CloudController::ServiceInstance.first.last_operation - expect(lo.type).to eq('delete') - expect(lo.state).to eq('failed') - expect(lo.description).to eq("An operation for a service binding of service instance #{instance.name} is in progress.") + instance.reload + expect(VCAP::CloudController::ServiceInstance.first(guid: instance.guid)).not_to be_nil + expect(instance.last_operation&.state).not_to eq('failed') lo = VCAP::CloudController::RouteBinding.first.last_operation expect(lo.type).to eq('delete') @@ -3541,7 +3623,8 @@ def check_filtered_instances(*instances) it 'continues to poll the last operation for the bindings' do api_call.call(admin_headers) - execute_all_jobs(expected_successes: 3, expected_failures: 1) + # 4 successes: the 3 binding-poll sub-jobs each poll once; the deferred root re-enqueues for later. + execute_all_jobs(expected_successes: 4, expected_failures: 0) [route_binding, service_binding, service_key].each do |binding| expect( @@ -3556,7 +3639,7 @@ def check_filtered_instances(*instances) it 'eventually removes the bindings' do api_call.call(admin_headers) - execute_all_jobs(expected_successes: 3, expected_failures: 1) + execute_all_jobs(expected_successes: 4, expected_failures: 0) expect(VCAP::CloudController::RouteBinding.all).to be_empty expect(VCAP::CloudController::ServiceBinding.all).to be_empty diff --git a/spec/support/reset_generic_enqueuer.rb b/spec/support/reset_generic_enqueuer.rb new file mode 100644 index 00000000000..83c68476ced --- /dev/null +++ b/spec/support/reset_generic_enqueuer.rb @@ -0,0 +1,6 @@ +# Reset the GenericEnqueuer thread-local after every example so a root_job_guid can't leak across parallel-worker specs. +RSpec.configure do |config| + config.after do + VCAP::CloudController::Jobs::GenericEnqueuer.reset! + end +end diff --git a/spec/support/shared_examples/jobs/recursive_delete_root_job.rb b/spec/support/shared_examples/jobs/recursive_delete_root_job.rb new file mode 100644 index 00000000000..78bd463cb1d --- /dev/null +++ b/spec/support/shared_examples/jobs/recursive_delete_root_job.rb @@ -0,0 +1,51 @@ +# Asserts a RootJobMixin-based delete job runs the mixin guards before its own delete action (mixin +# behaviour itself is covered in root_job_mixin_spec). The host spec must define: +## subject(:job) - the job instance under test +## root_operation - the root pollable operation string, e.g. 'app.delete' +## resource_guid_for_job - the guid used on the root pollable's resource_guid +## expect_no_delete_attempt { } - wraps a block, asserting the host action's delete is NOT invoked +## destroy_resource - removes the underlying resource (app/instance) from the db +RSpec.shared_examples 'a recursive delete root job' do + let!(:root_pollable_job) do + create(:pollable_job_model, + state: VCAP::CloudController::PollableJobModel::PROCESSING_STATE, + resource_guid: resource_guid_for_job, + operation: root_operation) + end + + def make_failed_sub_job(resource_guid: 'binding-guid') + create(:pollable_job_model, + root_job_guid: root_pollable_job.guid, + state: VCAP::CloudController::PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', + resource_guid: resource_guid, + cf_api_error: YAML.dump({ 'errors' => [{ 'detail' => 'unbind could not be completed: broker down' }] })) + end + + context 'when a sub-job is still in flight' do + let!(:pending_sub_job) do + create(:pollable_job_model, root_job_guid: root_pollable_job.guid, state: VCAP::CloudController::PollableJobModel::PROCESSING_STATE) + end + + it 'defers: does not attempt the delete and does not finish' do + expect_no_delete_attempt { job.perform } + expect(job.finished).to be_falsey + end + + it 'still defers when the resource is already gone (guard runs before the resource lookup)' do + destroy_resource + expect_no_delete_attempt { job.perform } + expect(job.finished).to be_falsey + end + end + + context 'when a sub-job has failed' do + let!(:failed_sub_job) { make_failed_sub_job } + + it 'surfaces the failure and does not attempt the delete' do + expect_no_delete_attempt do + expect { job.perform }.to raise_error(CloudController::Errors::CompoundError) + end + end + end +end diff --git a/spec/unit/actions/app_delete_spec.rb b/spec/unit/actions/app_delete_spec.rb index 44e7b0b0c77..32689de52f7 100644 --- a/spec/unit/actions/app_delete_spec.rb +++ b/spec/unit/actions/app_delete_spec.rb @@ -237,25 +237,24 @@ module VCAP::CloudController it 'alwayses call the broker with accepts_incomplete true' do expect(client).to receive(:unbind).with(binding1, user_guid: user_audit_info.user_guid, accepts_incomplete: true) - expect { app_delete.delete(app_dataset) }.to raise_error(AppDelete::SubResourceError) + expect { app_delete.delete(app_dataset) }.to raise_error(SubResourceError) end - it 'return an error that a service binding is being deleted asynchronously' do - expect { app_delete.delete(app_dataset) }.to raise_error(AppDelete::SubResourceError) do |err| - expect(err.underlying_errors.map(&:message)).to contain_exactly( - "An operation for the service binding between app #{binding1.app.name} and service instance #{binding1.service_instance.name} is in progress." - ) - end + it 'raises SubResourceError referencing the app and service instance' do + expect { app_delete.delete(app_dataset) }.to raise_error( + SubResourceError, + "An operation for the service binding between app #{binding1.app.name} and service instance #{binding1.service_instance.name} is in progress." + ) end it 'does not delete the app' do - expect { app_delete.delete(app_dataset) }.to raise_error(AppDelete::SubResourceError) + expect { app_delete.delete(app_dataset) }.to raise_error(SubResourceError) expect(binding1).to exist expect(app).to exist end it 'does not rollback the enqueuing of a job to delete the service binding' do - expect { app_delete.delete(app_dataset) }.to raise_error(AppDelete::SubResourceError) + expect { app_delete.delete(app_dataset) }.to raise_error(SubResourceError) expect(Delayed::Job.count).to eq 1 end end @@ -264,13 +263,14 @@ module VCAP::CloudController let!(:binding1) { create(:service_binding, app: app, service_instance: create(:managed_service_instance, space: app.space)) } let!(:binding2) { create(:service_binding, app: app, service_instance: create(:managed_service_instance, space: app.space)) } - it 'returns some errors describing that the service bindings are being deleted asynchronously' do - expect { app_delete.delete(app_dataset) }.to raise_error(AppDelete::SubResourceError) do |err| - expect(err.underlying_errors.map(&:message)).to contain_exactly( - "An operation for the service binding between app #{binding2.app.name} and service instance #{binding2.service_instance.name} is in progress.", - "An operation for the service binding between app #{binding1.app.name} and service instance #{binding1.service_instance.name} is in progress." + it 'raises SubResourceError carrying every binding\'s message after enqueueing a polling job each' do + expect { app_delete.delete(app_dataset) }.to raise_error(SubResourceError) do |err| + expect(err.message).to include( + "An operation for the service binding between app #{binding1.app.name} and service instance #{binding1.service_instance.name} is in progress.", + "An operation for the service binding between app #{binding2.app.name} and service instance #{binding2.service_instance.name} is in progress." ) end + expect(Delayed::Job.count).to eq 2 end end end @@ -291,7 +291,7 @@ module VCAP::CloudController it 'raises the errors wrapped into a SubResourceError' do expect do app_delete.delete(app_dataset) - end.to raise_error(AppDelete::SubResourceError) do |err| + end.to raise_error(SubResourceError) do |err| expect(err.underlying_errors).to have(2).items expect(err.underlying_errors).to all(be_a(StandardError)) expect(err.underlying_errors.map(&:message)).to eq(['error 1', 'error 2']) diff --git a/spec/unit/actions/app_stop_spec.rb b/spec/unit/actions/app_stop_spec.rb index 7b34f069d8f..9aacab7054a 100644 --- a/spec/unit/actions/app_stop_spec.rb +++ b/spec/unit/actions/app_stop_spec.rb @@ -20,12 +20,23 @@ module VCAP::CloudController it 'creates an audit event' do expect_any_instance_of(Repositories::AppEventRepository).to receive(:record_app_stop).with( app, - user_audit_info + user_audit_info, + delete_triggered: false ) AppStop.stop(app:, user_audit_info:) end + it 'passes delete_triggered through to the audit event when set' do + expect_any_instance_of(Repositories::AppEventRepository).to receive(:record_app_stop).with( + app, + user_audit_info, + delete_triggered: true + ) + + AppStop.stop(app: app, user_audit_info: user_audit_info, delete_triggered: true) + end + it 'prepares the sub-processes of the app' do AppStop.stop(app:, user_audit_info:) app.processes.each do |process| diff --git a/spec/unit/actions/service_instance_unshare_spec.rb b/spec/unit/actions/service_instance_unshare_spec.rb index 79140a0d900..e02774e5ae5 100644 --- a/spec/unit/actions/service_instance_unshare_spec.rb +++ b/spec/unit/actions/service_instance_unshare_spec.rb @@ -80,7 +80,7 @@ module VCAP::CloudController it 'fails to unshare' do expect { service_instance_unshare.unshare(service_instance, target_space, user_audit_info) }. to raise_error(VCAP::CloudController::ServiceInstanceUnshare::Error) do |err| - expect(err.message).to include("\n\tThe binding between an application and service instance #{service_instance.name} " \ + expect(err.message).to include("The binding between an application and service instance #{service_instance.name} " \ "in space #{target_space.name} is being deleted asynchronously.") end @@ -95,6 +95,13 @@ module VCAP::CloudController jobs = VCAP::CloudController::PollableJobModel.where(resource_guid: binding_guids) expect(jobs.count).to eq(3) end + + context 'when called with fail_if_in_progress: false' do + it 'raises SubResourceError instead of Error' do + expect { service_instance_unshare.unshare(service_instance, target_space, user_audit_info, fail_if_in_progress: false) }. + to raise_error(VCAP::CloudController::SubResourceError) + end + end end context 'when a binding has operation in progress and the delete action fails' do diff --git a/spec/unit/actions/v3/service_instance_delete_spec.rb b/spec/unit/actions/v3/service_instance_delete_spec.rb index 55ad84a30eb..bb3f27475f5 100644 --- a/spec/unit/actions/v3/service_instance_delete_spec.rb +++ b/spec/unit/actions/v3/service_instance_delete_spec.rb @@ -436,9 +436,9 @@ module V3 expect(delete_service_key_action).to have_received(:delete).with(service_key_3) expect(ServiceInstanceUnshare).to have_received(:new) - expect(unshare_action).to have_received(:unshare).with(service_instance, shared_space_1, event_repository.user_audit_info) - expect(unshare_action).to have_received(:unshare).with(service_instance, shared_space_2, event_repository.user_audit_info) - expect(unshare_action).to have_received(:unshare).with(service_instance, shared_space_3, event_repository.user_audit_info) + expect(unshare_action).to have_received(:unshare).with(service_instance, shared_space_1, event_repository.user_audit_info, fail_if_in_progress: true) + expect(unshare_action).to have_received(:unshare).with(service_instance, shared_space_2, event_repository.user_audit_info, fail_if_in_progress: true) + expect(unshare_action).to have_received(:unshare).with(service_instance, shared_space_3, event_repository.user_audit_info, fail_if_in_progress: true) end context 'when deleting bindings or unsharing spaces raises' do @@ -504,12 +504,7 @@ module V3 end it 'fails and schedules a polling job' do - expect do - action.delete - end.to raise_error( - ServiceInstanceDelete::UnbindingOperatationInProgress, - "An operation for a service binding of service instance #{service_instance.name} is in progress." - ) + expect { action.delete }.to raise_error(VCAP::CloudController::AsyncOperationInProgress) expect(Delayed::Job.all).to have(1).job end @@ -528,12 +523,7 @@ module V3 end it 'fails and schedules a polling job' do - expect do - action.delete - end.to raise_error( - ServiceInstanceDelete::UnbindingOperatationInProgress, - "An operation for the service binding between app #{service_binding.app.name} and service instance #{service_instance.name} is in progress." - ) + expect { action.delete }.to raise_error(VCAP::CloudController::AsyncOperationInProgress) expect(Delayed::Job.all).to have(1).job end @@ -552,12 +542,7 @@ module V3 end it 'fails and schedules a polling job' do - expect do - action.delete - end.to raise_error( - ServiceInstanceDelete::UnbindingOperatationInProgress, - "An operation for a service binding of service instance #{service_instance.name} is in progress." - ) + expect { action.delete }.to raise_error(VCAP::CloudController::AsyncOperationInProgress) expect(Delayed::Job.all).to have(1).job end diff --git a/spec/unit/controllers/runtime/apps_controller_spec.rb b/spec/unit/controllers/runtime/apps_controller_spec.rb index 96be7475581..28a89371696 100644 --- a/spec/unit/controllers/runtime/apps_controller_spec.rb +++ b/spec/unit/controllers/runtime/apps_controller_spec.rb @@ -1309,7 +1309,7 @@ def delete_app context 'when the error is a SubResource error' do before do errs = [StandardError.new('oops-1'), StandardError.new('oops-2')] - allow_any_instance_of(AppDelete).to receive(:delete_without_event).and_raise(VCAP::CloudController::AppDelete::SubResourceError.new(errs)) + allow_any_instance_of(AppDelete).to receive(:delete_without_event).and_raise(VCAP::CloudController::SubResourceError.new(errs)) end it 'returns all errors contained within it from the action' do diff --git a/spec/unit/controllers/v3/apps_controller_spec.rb b/spec/unit/controllers/v3/apps_controller_spec.rb index fa55b83a9eb..ecf6941b334 100644 --- a/spec/unit/controllers/v3/apps_controller_spec.rb +++ b/spec/unit/controllers/v3/apps_controller_spec.rb @@ -1033,15 +1033,13 @@ let(:space) { app_model.space } let(:org) { space.organization } let(:user) { set_current_user(create(:user)) } - let(:app_delete_stub) { instance_double(VCAP::CloudController::AppDelete) } before do allow_user_read_access_for(user, spaces: [space]) allow_user_write_access(user, space:) create(:buildpack_lifecycle_data_model, app: app_model, buildpacks: nil, stack: VCAP::CloudController::Stack.default.name) - allow(VCAP::CloudController::Jobs::DeleteActionJob).to receive(:new).and_call_original - allow(VCAP::CloudController::AppDelete).to receive(:new).and_return(app_delete_stub) - allow(AppsV3Controller::DeleteAppErrorTranslatorJob).to receive(:new).and_call_original + allow(VCAP::CloudController::V3::RecursiveDeleteAppJob).to receive(:new).and_call_original + TestConfig.override(temporary_enable_async_recursive_delete: { apps: true, service_instances: true }) end context 'when the app does not exist' do @@ -1053,23 +1051,17 @@ end end - it 'successfully deletes the app in a background job' do + it 'enqueues a RecursiveDeleteAppJob for the app' do delete :destroy, params: { guid: app_model.guid } - app_delete_jobs = Delayed::Job.where(Sequel.lit("handler like '%AppDelete%'")) - expect(app_delete_jobs.count).to eq 1 - app_delete_jobs.first - - expect(VCAP::CloudController::AppModel.find(guid: app_model.guid)).not_to be_nil - expect(VCAP::CloudController::Jobs::DeleteActionJob).to have_received(:new).with( - VCAP::CloudController::AppModel, - app_model.guid, - app_delete_stub + expect(VCAP::CloudController::V3::RecursiveDeleteAppJob).to have_received(:new).with( + app_model.guid, instance_of(VCAP::CloudController::UserAuditInfo) ) - expect(AppsV3Controller::DeleteAppErrorTranslatorJob).to have_received(:new) + expect(VCAP::CloudController::AppModel.find(guid: app_model.guid)).not_to be_nil + expect(Delayed::Job.count).to eq 1 end - it 'creates a job to track the deletion and returns it in the location header' do + it 'creates a pollable job to track the deletion and returns it in the location header' do expect do delete :destroy, params: { guid: app_model.guid } end.to change(VCAP::CloudController::PollableJobModel, :count).by(1) @@ -1085,6 +1077,52 @@ expect(response).to have_http_status(:accepted) expect(response.headers['Location']).to include "#{link_prefix}/v3/jobs/#{job.guid}" end + + context 'when a delete is already in flight for the app' do + let!(:existing_job) do + create(:pollable_job_model, + state: VCAP::CloudController::PollableJobModel::PROCESSING_STATE, + resource_guid: app_model.guid, + operation: 'app.delete') + end + + it 'returns 202 pointing at the existing job and does not enqueue a second one' do + expect do + delete :destroy, params: { guid: app_model.guid } + end.not_to(change(VCAP::CloudController::PollableJobModel, :count)) + + expect(VCAP::CloudController::V3::RecursiveDeleteAppJob).not_to have_received(:new) + expect(response).to have_http_status(:accepted) + expect(response.headers['Location']).to include "#{link_prefix}/v3/jobs/#{existing_job.guid}" + end + end + + context 'when temporary_enable_async_recursive_delete.apps is disabled (legacy path)' do + before { TestConfig.override(temporary_enable_async_recursive_delete: { apps: false, service_instances: false }) } + + it 'enqueues a DeleteActionJob (main behaviour) rather than a RecursiveDeleteAppJob' do + delete :destroy, params: { guid: app_model.guid } + + expect(response).to have_http_status(:accepted) + expect(VCAP::CloudController::V3::RecursiveDeleteAppJob).not_to have_received(:new) + expect(Delayed::Job.count).to eq(1) + + job = VCAP::CloudController::PollableJobModel.last + expect(job.resource_guid).to eq(app_model.guid) + expect(response.headers['Location']).to include "#{link_prefix}/v3/jobs/#{job.guid}" + end + + it 'does not consult the row-lock idempotency guard when a stale active pollable exists' do + create(:pollable_job_model, + state: VCAP::CloudController::PollableJobModel::PROCESSING_STATE, + resource_guid: app_model.guid, + operation: 'app.delete') + + expect do + delete :destroy, params: { guid: app_model.guid } + end.to change(VCAP::CloudController::PollableJobModel, :count).by(1) + end + end end describe '#start' do @@ -2316,32 +2354,4 @@ end end end - - describe 'DeleteAppErrorTranslatorJob' do - let(:error_translator) { AppsV3Controller::DeleteAppErrorTranslatorJob.new(job) } - let(:job) {} - - context 'when the error is a SubResourceError' do - it 'translates it to CompoundError with underlying API errors' do - translated_error = error_translator.translate_error(VCAP::CloudController::AppDelete::SubResourceError.new([ - StandardError.new('oops-1'), - StandardError.new('oops-2') - ])) - - expect(translated_error).to be_a(CloudController::Errors::CompoundError) - expect(translated_error.underlying_errors).to contain_exactly(CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', 'oops-1'), - CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', 'oops-2')) - end - end - - context 'when the error is not a SubResourceError' do - it 'justs return it' do - err = StandardError.new('oops') - - translated_error = error_translator.translate_error(err) - - expect(translated_error).to eq(err) - end - end - end end diff --git a/spec/unit/jobs/deserialization_spec.rb b/spec/unit/jobs/deserialization_spec.rb index 5cc0ab38d7f..3c8c622062e 100644 --- a/spec/unit/jobs/deserialization_spec.rb +++ b/spec/unit/jobs/deserialization_spec.rb @@ -49,6 +49,7 @@ module Jobs handler: !ruby/object:VCAP::CloudController::Jobs::TimeoutJob handler: !ruby/object:VCAP::CloudController::Jobs::PollableJobWrapper existing_guid: null + root_job_guid: null handler: !ruby/object:VCAP::CloudController::V3::CreateServiceInstanceJob start_time: #{job.start_time} finished: false diff --git a/spec/unit/jobs/enqueuer_spec.rb b/spec/unit/jobs/enqueuer_spec.rb index 1c3f274c285..86d6536abbb 100644 --- a/spec/unit/jobs/enqueuer_spec.rb +++ b/spec/unit/jobs/enqueuer_spec.rb @@ -4,6 +4,7 @@ require 'jobs/delete_action_job' require 'jobs/runtime/model_deletion' require 'jobs/error_translator_job' +require 'jobs/v3/recursive_delete_app_job' module VCAP::CloudController::Jobs RSpec.describe Enqueuer, job_context: :api do @@ -290,5 +291,74 @@ module VCAP::CloudController::Jobs end end end + + describe '#enqueue_or_find_active_pollable' do + let(:app_model) { create(:app_model) } + let(:user_audit_info) { VCAP::CloudController::UserAuditInfo.new(user_guid: create(:user).guid, user_email: 'test@example.com') } + let(:job_factory) { ->(_resource) { VCAP::CloudController::V3::RecursiveDeleteAppJob.new(app_model.guid, user_audit_info) } } + + def enqueue + Enqueuer.new(queue: Queues.generic).enqueue_or_find_active_pollable( + resource_model: VCAP::CloudController::AppModel, resource_guid: app_model.guid, operation: 'app.delete', &job_factory + ) + end + + context 'when no active delete job exists for the resource' do + it 'creates a new pollable job and returns it' do + job = nil + expect { job = enqueue }.to change(VCAP::CloudController::PollableJobModel, :count).by(1) + + expect(job).to be_a(VCAP::CloudController::PollableJobModel) + expect(job.state).to eq(VCAP::CloudController::PollableJobModel::PROCESSING_STATE) + expect(job.operation).to eq('app.delete') + expect(job.resource_guid).to eq(app_model.guid) + end + end + + context 'when an active delete job already exists for the resource' do + let!(:existing) do + create(:pollable_job_model, + state: VCAP::CloudController::PollableJobModel::PROCESSING_STATE, + resource_guid: app_model.guid, + operation: 'app.delete') + end + + it 'returns the existing pollable job without enqueueing a new one' do + result = nil + expect { result = enqueue }.not_to change(VCAP::CloudController::PollableJobModel, :count) + + expect(result.guid).to eq(existing.guid) + end + + it 'does not invoke the job factory block' do + invoked = false + Enqueuer.new(queue: Queues.generic).enqueue_or_find_active_pollable( + resource_model: VCAP::CloudController::AppModel, resource_guid: app_model.guid, operation: 'app.delete' + ) do |_resource| + invoked = true + VCAP::CloudController::V3::RecursiveDeleteAppJob.new(app_model.guid, user_audit_info) + end + + expect(invoked).to be(false) + end + end + + context 'when the resource no longer exists' do + before { app_model.destroy } + + it 'returns nil and does not enqueue a job' do + result = nil + expect { result = enqueue }.not_to change(VCAP::CloudController::PollableJobModel, :count) + expect(result).to be_nil + end + end + + context 'row lock' do + it 'builds a SELECT ... FOR UPDATE query on the resource row' do + sql = VCAP::CloudController::AppModel.where(guid: app_model.guid).for_update.sql + expect(sql).to match(/FOR UPDATE/i) + end + end + end end end diff --git a/spec/unit/jobs/generic_enqueuer_spec.rb b/spec/unit/jobs/generic_enqueuer_spec.rb index 31b8ef90cb6..c8b464f31fc 100644 --- a/spec/unit/jobs/generic_enqueuer_spec.rb +++ b/spec/unit/jobs/generic_enqueuer_spec.rb @@ -20,7 +20,6 @@ def perform end before do - # Reset singleton instance to ensure clean tests GenericEnqueuer.reset! end @@ -81,5 +80,76 @@ def perform expect(Delayed::Job.first.priority).to eq(7) end end + + describe 'root context' do + let(:job) { DummyPerformJob.new } + + describe '#activate_root_context' do + it 'stamps root_job_guid onto pollable rows created while active' do + enqueuer = generic_enqueuer.shared + enqueuer.activate_root_context(root_job_guid: 'root-guid-1') + + pollable_job = enqueuer.enqueue_pollable( + VCAP::CloudController::Jobs::DeleteActionJob.new(VCAP::CloudController::DropletModel, 'fake', + VCAP::CloudController::DropletDelete.new('fake')) + ) + + expect(pollable_job.root_job_guid).to eq('root-guid-1') + end + + it 'stamps root_job_guid onto every pollable row created while active' do + enqueuer = generic_enqueuer.shared + enqueuer.activate_root_context(root_job_guid: 'root-guid-1') + + first = enqueuer.enqueue_pollable( + VCAP::CloudController::Jobs::DeleteActionJob.new(VCAP::CloudController::DropletModel, 'fake-1', + VCAP::CloudController::DropletDelete.new('fake-1')) + ) + second = enqueuer.enqueue_pollable( + VCAP::CloudController::Jobs::DeleteActionJob.new(VCAP::CloudController::DropletModel, 'fake-2', + VCAP::CloudController::DropletDelete.new('fake-2')) + ) + + expect(first.root_job_guid).to eq('root-guid-1') + expect(second.root_job_guid).to eq('root-guid-1') + end + end + + describe '#deactivate_root_context' do + it 'clears the root_job_guid' do + enqueuer = generic_enqueuer.shared + enqueuer.activate_root_context(root_job_guid: 'root-guid-1') + enqueuer.deactivate_root_context + + expect(enqueuer.root_job_guid).to be_nil + end + + it 'subsequent enqueues no longer carry the root_job_guid' do + enqueuer = generic_enqueuer.shared + enqueuer.activate_root_context(root_job_guid: 'root-guid-1') + enqueuer.deactivate_root_context + + pollable_job = enqueuer.enqueue_pollable( + VCAP::CloudController::Jobs::DeleteActionJob.new(VCAP::CloudController::DropletModel, 'fake', + VCAP::CloudController::DropletDelete.new('fake')) + ) + + expect(pollable_job.root_job_guid).to be_nil + end + end + + describe 'without an active root context (default)' do + it 'enqueues pollable jobs with root_job_guid nil' do + enqueuer = generic_enqueuer.shared + + pollable_job = enqueuer.enqueue_pollable( + VCAP::CloudController::Jobs::DeleteActionJob.new(VCAP::CloudController::DropletModel, 'fake', + VCAP::CloudController::DropletDelete.new('fake')) + ) + + expect(pollable_job.root_job_guid).to be_nil + end + end + end end end diff --git a/spec/unit/jobs/mixins/root_job_mixin_spec.rb b/spec/unit/jobs/mixins/root_job_mixin_spec.rb new file mode 100644 index 00000000000..6692587049d --- /dev/null +++ b/spec/unit/jobs/mixins/root_job_mixin_spec.rb @@ -0,0 +1,476 @@ +require 'spec_helper' +require 'jobs/mixins/root_job_mixin' +require 'jobs/reoccurring_job' + +module VCAP::CloudController + module Jobs + RSpec.describe RootJobMixin do + let(:test_job_class) do + Class.new(ReoccurringJob) do + include RootJobMixin + + attr_reader :resource_guid + + def initialize(resource_guid) + super() + @resource_guid = resource_guid + end + + def perform; end + + def display_name + 'test.delete' + end + + def resource_type + 'test' + end + + def max_attempts + 1 + end + + def logger + Steno.logger('cc.jobs.test') + end + end + end + + let(:job) { test_job_class.new('resource-guid-1') } + + before { Jobs::GenericEnqueuer.reset! } + after { Jobs::GenericEnqueuer.reset! } + + def make_root(state: PollableJobModel::PROCESSING_STATE) + create(:pollable_job_model, state: state, resource_guid: 'resource-guid-1', operation: 'test.delete') + end + + def make_sub_job(state:, **attrs) + create(:pollable_job_model, root_job_guid: root_pollable_job.guid, state: state, **attrs) + end + + describe '#next_execution_in' do + let(:max_interval) { 100 } + + before do + allow(Config.config).to receive(:get).and_call_original + allow(Config.config).to receive(:get).with(:broker_client_max_async_poll_interval_seconds).and_return(max_interval) + end + + context 'when an active sub-job is due later' do + before { allow(job).to receive(:seconds_until_slowest_sub_job).and_return(30) } + + it 'wakes after that many seconds, plus the buffer' do + expect(job.send(:next_execution_in)).to eq(30 + RootJobMixin::ROOT_JOB_BUFFER_SECONDS) + end + + it 'caps the interval at the max async poll interval' do + allow(job).to receive(:seconds_until_slowest_sub_job).and_return(9999) + expect(job.send(:next_execution_in)).to eq(max_interval) + end + end + + context 'when no sub-job is due in the future' do + before { allow(job).to receive(:seconds_until_slowest_sub_job).and_return(nil) } + + it 'falls back to the ReoccurringJob interval plus the buffer' do + allow(job).to receive(:polling_interval_seconds).and_return(80) + expect(job.send(:next_execution_in)).to eq(80 + RootJobMixin::ROOT_JOB_BUFFER_SECONDS) + end + end + + context 'across occurrences with real sub-job rows' do + let!(:root_pollable_job) { make_root } + let(:now) { Delayed::Job.db_time_now } + + def add_active_sub_job(run_at:, state: PollableJobModel::POLLING_STATE) + dj = Delayed::Job.create!(guid: SecureRandom.uuid, handler: 'fake', run_at: run_at, queue: 'cc-generic') + create(:pollable_job_model, root_job_guid: root_pollable_job.guid, state: state, delayed_job_guid: dj.guid) + dj + end + + it 'derives the interval from the latest run_at among multiple active sub-jobs' do + add_active_sub_job(run_at: now + 10) + add_active_sub_job(run_at: now + 40) + add_active_sub_job(run_at: now + 25) + + expect(job.send(:next_execution_in)).to be_within(1).of(40 + RootJobMixin::ROOT_JOB_BUFFER_SECONDS) + end + + it 'ignores completed and failed sub-job rows, pacing only off the active ones' do + add_active_sub_job(run_at: now + 20) + add_active_sub_job(run_at: now + 90, state: PollableJobModel::COMPLETE_STATE) + add_active_sub_job(run_at: now + 90, state: PollableJobModel::FAILED_STATE) + + expect(job.send(:next_execution_in)).to be_within(1).of(20 + RootJobMixin::ROOT_JOB_BUFFER_SECONDS) + end + + it 're-derives from fresh rows when a sub-job re-enqueues' do + first = add_active_sub_job(run_at: now + 15) + expect(job.send(:next_execution_in)).to be_within(1).of(15 + RootJobMixin::ROOT_JOB_BUFFER_SECONDS) + + # next_execution_in reads sub-jobs directly, so it tracks the fresh row, not the destroyed one. + first.destroy + PollableJobModel.where(delayed_job_guid: first.guid).delete + add_active_sub_job(run_at: now + 55) + + expect(job.send(:next_execution_in)).to be_within(1).of(55 + RootJobMixin::ROOT_JOB_BUFFER_SECONDS) + end + end + end + + describe '#active_sub_jobs' do + let!(:root_pollable_job) { make_root } + + it 'returns only the active (processing/polling) sub-jobs' do + make_sub_job(state: PollableJobModel::PROCESSING_STATE, delayed_job_guid: 'active-1') + make_sub_job(state: PollableJobModel::POLLING_STATE, delayed_job_guid: 'active-2') + make_sub_job(state: PollableJobModel::COMPLETE_STATE, delayed_job_guid: 'done-1') + make_sub_job(state: PollableJobModel::FAILED_STATE, delayed_job_guid: 'failed-1') + job.send(:fetch_root_context) + + expect(job.send(:active_sub_jobs).map(&:delayed_job_guid)).to contain_exactly('active-1', 'active-2') + end + + it 'returns empty when there is no active root job' do + root_pollable_job.update(state: PollableJobModel::COMPLETE_STATE) + job.send(:fetch_root_context) + expect(job.send(:active_sub_jobs)).to eq([]) + end + end + + describe '#fetch_root_context' do + it 'loads the active pollable job for this resource and operation into root_job' do + pollable_job = make_root + + job.send(:fetch_root_context) + expect(job.send(:root_job)).to eq(pollable_job) + end + + it 'leaves root_job nil and sub_jobs empty when no active job exists' do + make_root(state: PollableJobModel::COMPLETE_STATE) + + job.send(:fetch_root_context) + expect(job.send(:root_job)).to be_nil + expect(job.send(:sub_jobs)).to eq([]) + end + end + + describe '#activate_root_job_context' do + let!(:root_pollable_job) { make_root } + + it 'installs the root_job_guid on the shared enqueuer so sub-jobs are linked' do + job.send(:activate_root_job_context) + expect(Jobs::GenericEnqueuer.shared.root_job_guid).to eq(root_pollable_job.guid) + ensure + job.send(:deactivate_root_job_context) + end + + it 'is a no-op when no active pollable job exists yet' do + root_pollable_job.update(state: PollableJobModel::COMPLETE_STATE) + + job.send(:activate_root_job_context) + expect(Jobs::GenericEnqueuer.shared.root_job_guid).to be_nil + ensure + job.send(:deactivate_root_job_context) + end + end + + describe '#deactivate_root_job_context' do + it 'clears the root_job_guid on the shared enqueuer' do + make_root + job.send(:activate_root_job_context) + job.send(:deactivate_root_job_context) + + expect(Jobs::GenericEnqueuer.shared.root_job_guid).to be_nil + end + end + + describe '#sub_jobs_in_flight?' do + let!(:root_pollable_job) { make_root } + + it 'returns false when there are no sub-jobs' do + job.send(:fetch_root_context) + expect(job.send(:sub_jobs_in_flight?)).to be(false) + end + + it 'returns true when any sub-job is PROCESSING' do + make_sub_job(state: PollableJobModel::PROCESSING_STATE) + job.send(:fetch_root_context) + expect(job.send(:sub_jobs_in_flight?)).to be(true) + end + + it 'returns true when any sub-job is POLLING' do + make_sub_job(state: PollableJobModel::POLLING_STATE) + job.send(:fetch_root_context) + expect(job.send(:sub_jobs_in_flight?)).to be(true) + end + + it 'returns false when no root pollable job is registered yet' do + root_pollable_job.update(state: PollableJobModel::COMPLETE_STATE) + job.send(:fetch_root_context) + expect(job.send(:sub_jobs_in_flight?)).to be(false) + end + + it 'does not raise even when a settled sub-job has failed' do + make_sub_job(state: PollableJobModel::FAILED_STATE) + job.send(:fetch_root_context) + expect(job.send(:sub_jobs_in_flight?)).to be(false) + end + + it 'persists a user-facing in-progress warning (without the word "async") while deferring' do + make_sub_job(state: PollableJobModel::PROCESSING_STATE) + job.send(:fetch_root_context) + + expect(job.send(:sub_jobs_in_flight?)).to be(true) + + warnings = root_pollable_job.reload.warnings + expect(warnings.map(&:detail)).to contain_exactly(job.send(:in_progress_warning_detail)) + expect(warnings.first.detail).not_to match(/async/i) + end + + it 'uses the job-provided warning text when the job overrides it' do + overriding_job = Class.new(test_job_class) do + def in_progress_warning_detail + 'custom in-progress message for this resource' + end + end.new('resource-guid-1') + make_sub_job(state: PollableJobModel::PROCESSING_STATE) + overriding_job.send(:fetch_root_context) + + overriding_job.send(:sub_jobs_in_flight?) + + expect(root_pollable_job.reload.warnings.map(&:detail)).to contain_exactly('custom in-progress message for this resource') + end + + it 'persists the in-progress warning only once across reoccurring runs' do + make_sub_job(state: PollableJobModel::PROCESSING_STATE) + job.send(:fetch_root_context) + + job.send(:sub_jobs_in_flight?) + job.send(:sub_jobs_in_flight?) + + expect(root_pollable_job.reload.warnings.count).to eq(1) + end + + it 'logs and does not raise when persisting the warning fails' do + make_sub_job(state: PollableJobModel::PROCESSING_STATE) + job.send(:fetch_root_context) + logger = instance_double(Steno::Logger, info: nil, warn: nil, error: nil) + allow(job).to receive(:logger).and_return(logger) + allow(JobWarningModel).to receive(:create).and_raise(Sequel::DatabaseError.new('warning insert failed')) + + expect { job.send(:sub_jobs_in_flight?) }.not_to raise_error + expect(logger).to have_received(:warn).with(/could not add in-progress warning/) + end + + context 'when some sub-jobs have failed but others are still active' do + before do + make_sub_job(state: PollableJobModel::FAILED_STATE) + make_sub_job(state: PollableJobModel::PROCESSING_STATE) + end + + it 'returns true (still waiting on the active sub-job)' do + job.send(:fetch_root_context) + expect(job.send(:sub_jobs_in_flight?)).to be(true) + end + end + end + + describe '#raise_if_sub_jobs_failed' do + let!(:root_pollable_job) { make_root } + + it 'does nothing when there are no failed sub-jobs' do + make_sub_job(state: PollableJobModel::COMPLETE_STATE) + job.send(:fetch_root_context) + expect { job.send(:raise_if_sub_jobs_failed) }.not_to raise_error + end + + it 'does nothing when no active root job exists' do + root_pollable_job.update(state: PollableJobModel::COMPLETE_STATE) + job.send(:fetch_root_context) + expect { job.send(:raise_if_sub_jobs_failed) }.not_to raise_error + end + + context 'when a settled sub-job has failed' do + before do + make_sub_job(state: PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', resource_guid: 'binding-1', + cf_api_error: YAML.dump({ 'errors' => [{ 'title' => 'CF-UnableToPerform', 'code' => 10_009, + 'detail' => 'unbind could not be completed: broker exploded' }] })) + make_sub_job(state: PollableJobModel::COMPLETE_STATE) + end + + it 'raises a CompoundError of UnprocessableEntity carrying the failed sub-job detail' do + job.send(:fetch_root_context) + expect { job.send(:raise_if_sub_jobs_failed) }.to raise_error(CloudController::Errors::CompoundError) do |err| + expect(err.underlying_errors.map(&:name)).to eq(%w[UnprocessableEntity]) + expect(err.underlying_errors.first.message).to include('unbind could not be completed: broker exploded') + end + end + + context 'when the failed sub-job has no stored error detail' do + before do + make_sub_job(state: PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', resource_guid: 'binding-2') + end + + it 'falls back to a resource reference for that entry' do + job.send(:fetch_root_context) + expect { job.send(:raise_if_sub_jobs_failed) }.to raise_error(CloudController::Errors::CompoundError) do |err| + expect(err.underlying_errors.map(&:message)).to include(a_string_including('service_credential_binding binding-2')) + end + end + end + end + end + + describe 'sub_resource_errors (durable sync-failure hook)' do + let!(:root_pollable_job) { make_root } + + it 'defaults to none so jobs without sub-resources are unaffected' do + make_sub_job(state: PollableJobModel::COMPLETE_STATE) + job.send(:fetch_root_context) + expect { job.send(:raise_if_sub_jobs_failed) }.not_to raise_error + end + + context 'when a subclass reports a failed sub-resource with no matching sub-job' do + let(:job_with_sync_failure) do + klass = Class.new(test_job_class) do + def sub_resource_errors + [['sync-binding', CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', 'sync unbind failed')]] + end + end + klass.new('resource-guid-1') + end + + it 'does NOT halt when there is no failed sub-job, so the action can re-run and retry it' do + job_with_sync_failure.send(:fetch_root_context) + expect { job_with_sync_failure.send(:raise_if_sub_jobs_failed) }.not_to raise_error + end + + it 'is reported once a sub-job has terminally failed (action is then skipped), merged with that sub-job' do + make_sub_job(state: PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', resource_guid: 'async-binding', + cf_api_error: YAML.dump({ 'errors' => [{ 'detail' => 'async unbind failed' }] })) + job_with_sync_failure.send(:fetch_root_context) + + expect { job_with_sync_failure.send(:raise_if_sub_jobs_failed) }.to raise_error(CloudController::Errors::CompoundError) do |err| + expect(err.underlying_errors.map(&:message)).to include(a_string_including('sync unbind failed'), a_string_including('async unbind failed')) + end + end + end + + context 'when a failed sub-resource shares its guid with a failed sub-job' do + let(:job_with_dup) do + klass = Class.new(test_job_class) do + def sub_resource_errors + [['shared-guid', CloudController::Errors::ApiError.new_from_details('UnprocessableEntity', 'unbind failed once')]] + end + end + klass.new('resource-guid-1') + end + + it 'reports the resource only once' do + make_sub_job(state: PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', resource_guid: 'shared-guid', + cf_api_error: YAML.dump({ 'errors' => [{ 'detail' => 'unbind failed once' }] })) + job_with_dup.send(:fetch_root_context) + + expect { job_with_dup.send(:raise_if_sub_jobs_failed) }.to raise_error(CloudController::Errors::CompoundError) do |err| + expect(err.underlying_errors.size).to eq(1) + end + end + end + end + + describe '#perform_with_root_job_handling' do + let!(:root_pollable_job) { make_root } + + it 'activates the root context for the duration of the block' do + observed = nil + job.send(:perform_with_root_job_handling) do + observed = Jobs::GenericEnqueuer.shared.root_job_guid + end + + expect(observed).to eq(root_pollable_job.guid) + end + + it 'fetches the root job from the DB only once, even when the block hits several sub-job helpers' do + allow(PollableJobModel).to receive(:find_active_delete).and_call_original + + job.send(:perform_with_root_job_handling) do + job.send(:sub_jobs_in_flight?) + job.send(:raise_if_sub_jobs_failed) + job.send(:active_sub_jobs) + end + + expect(PollableJobModel).to have_received(:find_active_delete).once + end + + it 'always deactivates the root context, even when the block raises' do + expect do + job.send(:perform_with_root_job_handling) { raise StandardError.new('boom') } + end.to raise_error(CloudController::Errors::ApiError) + + expect(Jobs::GenericEnqueuer.shared.root_job_guid).to be_nil + end + + it 'swallows SubResourceError carrying only in-progress signals' do + expect do + job.send(:perform_with_root_job_handling) do + raise SubResourceError.new([AsyncOperationInProgress.new('async')]) + end + end.not_to raise_error + end + + it 'translates SubResourceError with real failures to a CompoundError of UnprocessableEntity' do + expect do + job.send(:perform_with_root_job_handling) do + raise SubResourceError.new([StandardError.new('one broke'), StandardError.new('two broke')]) + end + end.to raise_error(CloudController::Errors::CompoundError) do |err| + expect(err.underlying_errors).to all(be_a(CloudController::Errors::ApiError)) + expect(err.underlying_errors.map(&:name)).to eq(%w[UnprocessableEntity UnprocessableEntity]) + expect(err.underlying_errors.map(&:message)).to include(match(/one broke/), match(/two broke/)) + end + end + + it 'merges current-tick sync failures with settled failed sub-jobs into one CompoundError' do + make_sub_job(state: PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', resource_guid: 'async-binding', + cf_api_error: YAML.dump({ 'errors' => [{ 'title' => 'CF-UnableToPerform', 'code' => 10_009, + 'detail' => 'async unbind failed' }] })) + + expect do + job.send(:perform_with_root_job_handling) do + raise SubResourceError.new([StandardError.new('sync unbind failed')]) + end + end.to raise_error(CloudController::Errors::CompoundError) do |err| + expect(err.underlying_errors.map(&:name)).to all(eq('UnprocessableEntity')) + expect(err.underlying_errors.map(&:message)).to include(match(/sync unbind failed/), match(/async unbind failed/)) + end + end + + it 'passes ApiErrors through unchanged' do + original = CloudController::Errors::ApiError.new_from_details('UnableToPerform', 'delete', 'broker said no') + + expect do + job.send(:perform_with_root_job_handling) { raise original } + end.to raise_error(CloudController::Errors::ApiError) do |err| + expect(err.name).to eq('UnableToPerform') + end + end + + it 'surfaces unexpected StandardErrors as UnableToPerform' do + expect do + job.send(:perform_with_root_job_handling) { raise StandardError.new('uncategorised') } + end.to raise_error(CloudController::Errors::ApiError) do |err| + expect(err.name).to eq('UnableToPerform') + expect(err.message).to include('uncategorised') + end + end + end + end + end +end diff --git a/spec/unit/jobs/mixins/root_job_orchestration_spec.rb b/spec/unit/jobs/mixins/root_job_orchestration_spec.rb new file mode 100644 index 00000000000..ec3fe5f2604 --- /dev/null +++ b/spec/unit/jobs/mixins/root_job_orchestration_spec.rb @@ -0,0 +1,240 @@ +require 'spec_helper' +require 'jobs/mixins/root_job_mixin' +require 'jobs/reoccurring_job' + +# Integration-style: drives the real DelayedJob worker loop (+ Timecop) to prove root + sub-jobs converge. +module VCAP::CloudController + module Jobs + RSpec.describe 'RootJobMixin orchestration', isolation: :truncation do + # Per-resource run counter, shared across ReoccurringJob's YAML reschedule and terminal-row deletion. + RUN_COUNTS = Hash.new(0) + + # A sub-job that completes after `iterations` runs, or fails on run number `fail_on_run` (1-based). + class CountingSubJob < ReoccurringJob + def initialize(resource_guid, iterations: 1, fail_on_run: nil) + super() + @resource_guid = resource_guid + @iterations = iterations + @fail_on_run = fail_on_run + end + + attr_reader :resource_guid + + def perform + RUN_COUNTS[resource_guid] += 1 + run = RUN_COUNTS[resource_guid] + raise CloudController::Errors::ApiError.new_from_details('UnableToPerform', 'delete', "sub-job #{resource_guid} failed on run #{run}") if @fail_on_run == run + + finish if run >= @iterations + end + + def display_name = "#{resource_guid}.delete" + def resource_type = 'sub-resource' + def max_attempts = 1 + def handle_timeout; end + end + + # A root job mirroring the production body: defer while sub-jobs are in flight, else surface failure or finish. + class CountingRootJob < ReoccurringJob + include RootJobMixin + + attr_reader :resource_guid + + def initialize(resource_guid) + super() + @resource_guid = resource_guid + end + + def perform + perform_with_root_job_handling do + RUN_COUNTS[resource_guid] += 1 + return if sub_jobs_in_flight? + + raise_if_sub_jobs_failed + finish + end + end + + def display_name = 'root.delete' + def resource_type = 'root-resource' + def max_attempts = 1 + def handle_timeout; end + + def logger = Steno.logger('cc.jobs.test.root') + end + + # Flat poll interval for every job so a single Timecop jump makes all rescheduled jobs due at once. + let(:interval) { 60 } + + before do + RUN_COUNTS.clear + TestConfig.override( + broker_client_default_async_poll_interval_seconds: interval, + broker_client_async_poll_exponential_backoff_rate: 1.0, + broker_client_max_async_poll_interval_seconds: interval + RootJobMixin::ROOT_JOB_BUFFER_SECONDS, + broker_client_max_async_poll_duration_minutes: 24 * 60 + ) + Jobs::GenericEnqueuer.reset! + end + + after do + Jobs::GenericEnqueuer.reset! + Timecop.return + end + + # Enqueue the root, then each sub-job under the root's context (stamping root_job_guid on each row). + def enqueue_root_with_sub_jobs(root, sub_jobs) + Timecop.freeze do + root_pollable = Jobs::Enqueuer.new(queue: Jobs::Queues.generic).enqueue_pollable(root) + + enqueuer = Jobs::GenericEnqueuer.shared + enqueuer.activate_root_context(root_job_guid: root_pollable.guid) + sub_jobs.each { |sj| enqueuer.enqueue_pollable(sj) } + enqueuer.deactivate_root_context + + root_pollable + end + end + + # Runs the worker until nothing is due, advancing the clock one interval per pass; yields the root state each pass. + def drain_until_settled(root_pollable, max_passes: 20) + passes = 0 + max_passes.times do + successes, failures = Delayed::Worker.new.work_off(100) + break if successes.zero? && failures.zero? + + passes += 1 + yield(root_pollable.reload.state) if block_given? + Timecop.freeze(Time.now + interval + RootJobMixin::ROOT_JOB_BUFFER_SECONDS + 1) + end + passes + end + + def state_of(resource_guid) + PollableJobModel.where(resource_guid:).first.state + end + + context 'when sub-jobs take multiple runs and the root waits for all of them' do + it 'reruns each occurrence until every sub-job completes, then completes the root, and never settles early' do + root = CountingRootJob.new('root-1') + root_pollable = enqueue_root_with_sub_jobs(root, [ + CountingSubJob.new('fast', iterations: 1), + CountingSubJob.new('slow', iterations: 3) + ]) + + root_states = [] + passes = drain_until_settled(root_pollable) { |state| root_states << state } + + expect(state_of('fast')).to eq(PollableJobModel::COMPLETE_STATE) + expect(state_of('slow')).to eq(PollableJobModel::COMPLETE_STATE) + expect(root_pollable.reload.state).to eq(PollableJobModel::COMPLETE_STATE) + + # Root runs exactly as many times as the slowest sub-job (3) — no wasted trailing run. + expect(RUN_COUNTS['fast']).to eq(1) + expect(RUN_COUNTS['slow']).to eq(3) + expect(passes).to eq(3) + expect(RUN_COUNTS['root-1']).to eq(3) + + # Invariant: the root only reached COMPLETE on the final pass, never while a sub-job was active. + expect(root_states[0...-1]).to all(eq(PollableJobModel::POLLING_STATE)) + expect(root_states.last).to eq(PollableJobModel::COMPLETE_STATE) + end + end + + context 'when one sub-job fails while another is still running' do + it 'keeps deferring while the other sub-job is in flight, then fails the root once all have settled' do + root = CountingRootJob.new('root-2') + root_pollable = enqueue_root_with_sub_jobs(root, [ + CountingSubJob.new('failing', iterations: 5, fail_on_run: 1), + CountingSubJob.new('running', iterations: 3) + ]) + + root_states = [] + passes = drain_until_settled(root_pollable) { |state| root_states << state } + + expect(state_of('failing')).to eq(PollableJobModel::FAILED_STATE) + expect(state_of('running')).to eq(PollableJobModel::COMPLETE_STATE) + expect(root_pollable.reload.state).to eq(PollableJobModel::FAILED_STATE) + + # Early failure adds no extra wakes: the root paces off the still-running sub-job (3 runs). + expect(RUN_COUNTS['failing']).to eq(1) + expect(RUN_COUNTS['running']).to eq(3) + expect(passes).to eq(3) + expect(RUN_COUNTS['root-2']).to eq(3) + + # The root kept polling until the still-running sub-job also settled, only then failing. + expect(root_states[0...-1]).to all(eq(PollableJobModel::POLLING_STATE)) + expect(root_states.last).to eq(PollableJobModel::FAILED_STATE) + end + end + + context 'when a long-running sub-job fails only after several runs' do + it 'defers across every occurrence until the sub-job eventually fails, then fails the root' do + root = CountingRootJob.new('root-3') + root_pollable = enqueue_root_with_sub_jobs(root, [ + CountingSubJob.new('eventual', iterations: 10, fail_on_run: 3) + ]) + + root_states = [] + passes = drain_until_settled(root_pollable) { |state| root_states << state } + + expect(state_of('eventual')).to eq(PollableJobModel::FAILED_STATE) + expect(root_pollable.reload.state).to eq(PollableJobModel::FAILED_STATE) + + # Root runs exactly 3 times, no wasted trailing run after the failure settles. + expect(RUN_COUNTS['eventual']).to eq(3) + expect(passes).to eq(3) + expect(RUN_COUNTS['root-3']).to eq(3) + + expect(root_states[0...-1]).to all(eq(PollableJobModel::POLLING_STATE)) + expect(root_states.last).to eq(PollableJobModel::FAILED_STATE) + end + end + + # next_execution_in makes the root pace off its sub-jobs' run_at instead of its own backoff. + # First-tick caveat: the sub-job's run_at isn't committed yet, so the root takes one extra backoff wake. + context 'when a sub-job polls slower than the root default backoff' do + let(:sub_poll) { 120 } # sub-job re-polls every 120s (e.g. a broker Retry-After) + let(:root_backoff) { 30 } # the root's default backoff would otherwise fire every 30s + + before do + TestConfig.override( + broker_client_default_async_poll_interval_seconds: root_backoff, + broker_client_async_poll_exponential_backoff_rate: 1.0, + broker_client_max_async_poll_interval_seconds: sub_poll + RootJobMixin::ROOT_JOB_BUFFER_SECONDS, + broker_client_max_async_poll_duration_minutes: 24 * 60 + ) + end + + it 'collapses wasted wakes: runs far fewer times than the default backoff would' do + root = CountingRootJob.new('root-4') + slow_poller = CountingSubJob.new('slow-poller', iterations: 2) + slow_poller.polling_interval_seconds = sub_poll + + root_pollable = enqueue_root_with_sub_jobs(root, [slow_poller]) + start = nil + Timecop.freeze { start = Time.now } + + # Step in fine root_backoff increments so a root ignoring its sub-jobs' schedule could wake at every step. + Timecop.freeze(start) + total_span = (sub_poll * 2) + root_backoff + steps = total_span / root_backoff + steps.times do + Delayed::Worker.new.work_off(100) + break if root_pollable.reload.state != PollableJobModel::POLLING_STATE + + Timecop.freeze(Time.now + root_backoff) + end + + expect(state_of('slow-poller')).to eq(PollableJobModel::COMPLETE_STATE) + expect(root_pollable.reload.state).to eq(PollableJobModel::COMPLETE_STATE) + expect(RUN_COUNTS['slow-poller']).to eq(2) + + # A bound, not an exact count: the precise number jitters by ±1 with DB-clock/worker interleaving. + expect(RUN_COUNTS['root-4']).to be <= RUN_COUNTS['slow-poller'] + 1 + expect(RUN_COUNTS['root-4']).to be < steps / 2 + end + end + end + end +end diff --git a/spec/unit/jobs/pollable_job_wrapper_spec.rb b/spec/unit/jobs/pollable_job_wrapper_spec.rb index 14c1683f6c7..e96b7547f2b 100644 --- a/spec/unit/jobs/pollable_job_wrapper_spec.rb +++ b/spec/unit/jobs/pollable_job_wrapper_spec.rb @@ -169,6 +169,30 @@ class BigException < StandardError end end + context 'when a root_job_guid is provided' do + let(:root_pollable_job) { create(:pollable_job_model) } + let(:pollable_job) { PollableJobWrapper.new(job, root_job_guid: root_pollable_job.guid) } + + it 'persists the root_job_guid on the new PollableJobModel row' do + enqueued_job = VCAP::CloudController::Jobs::Enqueuer.new.enqueue(pollable_job) + job_record = VCAP::CloudController::PollableJobModel.find(delayed_job_guid: enqueued_job.guid) + expect(job_record.root_job_guid).to eq(root_pollable_job.guid) + end + + it 'makes the new row queryable as a sub-job of the root' do + VCAP::CloudController::Jobs::Enqueuer.new.enqueue(pollable_job) + expect(root_pollable_job.sub_jobs.length).to eq(1) + end + end + + context 'when no root_job_guid is provided (default)' do + it 'creates a row with root_job_guid nil' do + enqueued_job = VCAP::CloudController::Jobs::Enqueuer.new.enqueue(pollable_job) + job_record = VCAP::CloudController::PollableJobModel.find(delayed_job_guid: enqueued_job.guid) + expect(job_record.root_job_guid).to be_nil + end + end + context 'when the job fails' do before do allow_any_instance_of(VCAP::CloudController::Jobs::DeleteActionJob). diff --git a/spec/unit/jobs/v3/recursive_delete_app_job_spec.rb b/spec/unit/jobs/v3/recursive_delete_app_job_spec.rb new file mode 100644 index 00000000000..4eb7885d1b7 --- /dev/null +++ b/spec/unit/jobs/v3/recursive_delete_app_job_spec.rb @@ -0,0 +1,197 @@ +require 'spec_helper' +require 'support/shared_examples/jobs/delayed_job' +require 'support/shared_examples/jobs/recursive_delete_root_job' +require 'jobs/v3/recursive_delete_app_job' + +module VCAP::CloudController + module V3 + RSpec.describe RecursiveDeleteAppJob do + let(:user_audit_info) { UserAuditInfo.new(user_guid: create(:user).guid, user_email: 'test@example.com') } + let(:org) { create(:organization) } + let(:space) { create(:space, organization: org) } + let(:app_model) { create(:app_model, space: space, name: 'my-app') } + + subject(:job) { described_class.new(app_model.guid, user_audit_info) } + + before { Jobs::GenericEnqueuer.reset! } + after { Jobs::GenericEnqueuer.reset! } + + it_behaves_like 'delayed job', described_class + + describe '#perform' do + def root_pollable(state: PollableJobModel::PROCESSING_STATE) + create(:pollable_job_model, state: state, resource_guid: app_model.guid, operation: 'app.delete') + end + + def make_failed_binding(desc:) + service_instance = create(:managed_service_instance, space:) + create(:service_binding, app: app_model, service_instance: service_instance).tap do |b| + b.save_with_attributes_and_new_operation({}, { type: 'delete', state: 'failed', description: desc }) + end + end + + context 'when the app does not exist' do + before { app_model.destroy } + + it 'finishes the job' do + job.perform + expect(job.finished).to be(true) + end + end + + context 'when the app has no service bindings' do + it 'deletes the app and finishes' do + job.perform + expect(job.finished).to be(true) + expect(AppModel.find(guid: app_model.guid)).to be_nil + end + end + + context 'when the app has synchronous service bindings' do + let(:service_instance) { create(:managed_service_instance, space:) } + let!(:binding) { create(:service_binding, app: app_model, service_instance: service_instance) } + + before do + stub_unbind(binding, accepts_incomplete: true) + end + + it 'deletes the app and bindings, and finishes' do + job.perform + expect(job.finished).to be(true) + expect(AppModel.find(guid: app_model.guid)).to be_nil + expect(ServiceBinding.where(app_guid: app_model.guid).count).to eq(0) + end + end + + describe 'stopping the app before deletion' do + let(:process) { create(:process_model, app: app_model, state: ProcessModel::STARTED) } + + before do + app_model.update(desired_state: ProcessModel::STARTED) + process + end + + it 'stops the app so it is not running during async unbinding' do + allow_any_instance_of(AppDelete).to receive(:delete).and_raise( + VCAP::CloudController::SubResourceError.new([VCAP::CloudController::AsyncOperationInProgress.new('async binding')]) + ) + + job.perform + + expect(app_model.reload.desired_state).to eq(ProcessModel::STOPPED) + expect(process.reload.state).to eq(ProcessModel::STOPPED) + end + + it 'records an audit.app.stop event tagged with delete_triggered so the cascade is transparent' do + expect { job.perform }.to change { Event.where(type: 'audit.app.stop', actee: app_model.guid).count }.by(1) + + stop_event = Event.where(type: 'audit.app.stop', actee: app_model.guid).last + expect(stop_event.metadata['delete_triggered']).to be(true) + end + + context 'when the app is already stopped' do + before { app_model.update(desired_state: ProcessModel::STOPPED) } + + it 'does not record a redundant audit.app.stop event' do + expect { job.perform }.not_to(change { Event.where(type: 'audit.app.stop', actee: app_model.guid).count }) + end + end + end + + context 'when AppDelete raises SubResourceError carrying only async-in-progress signals' do + before do + allow_any_instance_of(AppDelete).to receive(:delete).and_raise( + VCAP::CloudController::SubResourceError.new([VCAP::CloudController::AsyncOperationInProgress.new('async binding')]) + ) + end + + it 'does not finish or destroy the app' do + job.perform + expect(job.finished).to be(false) + expect(AppModel.find(guid: app_model.guid)).not_to be_nil + end + end + + context 'shared root-job behaviour' do + let(:resource_guid_for_job) { app_model.guid } + let(:root_operation) { 'app.delete' } + + def expect_no_delete_attempt + expect_any_instance_of(AppDelete).not_to receive(:delete) + yield + end + + def destroy_resource + app_model.destroy + end + + it_behaves_like 'a recursive delete root job' + + it 'leaves the app intact on failure so the DELETE can be retried' do + root = root_pollable + create(:pollable_job_model, root_job_guid: root.guid, state: PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', resource_guid: 'binding-guid', + cf_api_error: YAML.dump({ 'errors' => [{ 'detail' => 'broker down' }] })) + expect { job.perform }.to raise_error(CloudController::Errors::CompoundError) + expect(AppModel.find(guid: app_model.guid)).not_to be_nil + end + end + + context 'when a binding is left in delete/failed but no sub-job failed (e.g. a sporadic error)' do + let!(:sporadically_failed_binding) { make_failed_binding(desc: 'sporadic db error') } + let!(:root_pollable_job) { root_pollable } + + before { stub_unbind(sporadically_failed_binding, accepts_incomplete: true) } + + it 're-runs the action so the binding is retried and can self-heal' do + expect_any_instance_of(AppDelete).to receive(:delete).and_call_original + job.perform + expect(job.finished).to be(true) + expect(AppModel.find(guid: app_model.guid)).to be_nil + end + + it 'logs the failed binding (with detail and job guid) for operator visibility' do + logger = instance_double(Steno::Logger, info: nil, warn: nil, error: nil) + allow(job).to receive(:logger).and_return(logger) + + job.perform + + expect(logger).to have_received(:warn).with( + a_string_including(sporadically_failed_binding.guid).and(including('sporadic db error')).and(including(root_pollable_job.guid)) + ) + end + end + end + + describe '#resource_guid' do + it 'returns the app guid' do + expect(job.resource_guid).to eq(app_model.guid) + end + end + + describe '#resource_type' do + it 'returns "app"' do + expect(job.resource_type).to eq('app') + end + end + + describe '#display_name' do + it 'returns "app.delete"' do + expect(job.display_name).to eq('app.delete') + end + end + + describe '#max_attempts' do + it 'returns 1' do + expect(job.max_attempts).to eq(1) + end + end + + describe '#handle_timeout' do + it 'is a no-op' do + expect { job.handle_timeout }.not_to raise_error + end + end + end + end +end diff --git a/spec/unit/jobs/v3/recursive_delete_service_instance_job_spec.rb b/spec/unit/jobs/v3/recursive_delete_service_instance_job_spec.rb new file mode 100644 index 00000000000..9f833a46f44 --- /dev/null +++ b/spec/unit/jobs/v3/recursive_delete_service_instance_job_spec.rb @@ -0,0 +1,171 @@ +require 'spec_helper' +require 'support/shared_examples/jobs/delayed_job' +require 'support/shared_examples/jobs/recursive_delete_root_job' +require 'jobs/v3/recursive_delete_service_instance_job' +require 'cloud_controller/errors/api_error' +require 'cloud_controller/user_audit_info' +require 'actions/v3/service_instance_delete' + +module VCAP::CloudController + module V3 + RSpec.describe RecursiveDeleteServiceInstanceJob do + let(:user_audit_info) { UserAuditInfo.new(user_guid: create(:user).guid, user_email: 'foo@example.com') } + let(:service_instance) { create(:managed_service_instance, service_plan:) } + let(:service_plan) { create(:service_plan, service: service_offering) } + let(:service_offering) { create(:service) } + + subject(:job) { described_class.new(service_instance.guid, user_audit_info) } + + before { Jobs::GenericEnqueuer.reset! } + after { Jobs::GenericEnqueuer.reset! } + + it_behaves_like 'delayed job', described_class + + describe '#perform' do + let(:delete_response) { { finished: false, operation: 'test-operation' } } + let(:poll_response) { { finished: false } } + let(:action) do + double(VCAP::CloudController::V3::ServiceInstanceDelete, { delete: delete_response, poll: poll_response }) + end + + before do + allow(VCAP::CloudController::V3::ServiceInstanceDelete).to receive(:new).and_return(action) + end + + def root_pollable(state: PollableJobModel::PROCESSING_STATE) + create(:pollable_job_model, state: state, resource_guid: service_instance.guid, operation: 'service_instance.delete') + end + + def make_failed_binding(desc:) + create(:service_binding, service_instance: service_instance).tap do |b| + b.save_with_attributes_and_new_operation({}, { type: 'delete', state: 'failed', description: desc }) + end + end + + it 'passes fail_if_in_progress: false to the action' do + job.perform + + expect(VCAP::CloudController::V3::ServiceInstanceDelete).to have_received(:new).with( + service_instance, + an_instance_of(VCAP::CloudController::Repositories::ServiceEventRepository), + fail_if_in_progress: false + ).at_least(:once) + end + + context 'when a binding-delete sub-job has been enqueued (async)' do + before do + allow(action).to receive(:delete).and_raise(VCAP::CloudController::SubResourceError.new([VCAP::CloudController::AsyncOperationInProgress.new('async')])) + end + + it 'swallows the error, does not finish, and does not poll the broker this cycle' do + expect { job.perform }.not_to raise_error + expect(job.finished).to be_falsey + expect(action).not_to have_received(:poll) + end + end + + context 'shared root-job behaviour' do + let(:resource_guid_for_job) { service_instance.guid } + let(:root_operation) { 'service_instance.delete' } + + def expect_no_delete_attempt + expect(action).not_to receive(:delete) + yield + end + + def destroy_resource + service_instance.destroy + end + + it_behaves_like 'a recursive delete root job' + end + + context 'when a binding is left in delete/failed but no sub-job failed (e.g. a sporadic error)' do + let!(:sporadically_failed_binding) { make_failed_binding(desc: 'sporadic db error') } + let!(:root_pollable_job) { root_pollable } + + it 're-runs the action so the binding is retried and can self-heal' do + expect(action).to receive(:delete).and_return({ finished: true }) + job.perform + expect(job.finished).to be(true) + end + + it 'logs the failed binding (with detail and job guid) for operator visibility' do + logger = instance_double(Steno::Logger, info: nil, warn: nil, error: nil) + allow(job).to receive(:logger).and_return(logger) + + job.perform + + expect(logger).to have_received(:warn).with( + a_string_including(sporadically_failed_binding.guid).and(including('sporadic db error')).and(including(root_pollable_job.guid)) + ) + end + end + + context 'when a previous delete attempt for this instance failed and left stale pollable rows' do + let!(:previous_failed_root) { root_pollable(state: PollableJobModel::FAILED_STATE) } + + let!(:previous_failed_sub) do + create(:pollable_job_model, + root_job_guid: previous_failed_root.guid, + state: PollableJobModel::FAILED_STATE, + resource_type: 'service_credential_binding', + resource_guid: 'binding-guid') + end + + it 'ignores the stale rows and runs the delete this cycle' do + expect(action).to receive(:delete).and_return({ finished: true }) + job.perform + expect(job.finished).to be(true) + end + end + end + + describe 'handle timeout' do + let(:action) do + double(VCAP::CloudController::V3::ServiceInstanceDelete, { update_last_operation_with_failure: nil }) + end + + before do + allow(VCAP::CloudController::V3::ServiceInstanceDelete).to receive(:new).and_return(action) + end + + it 'asks the action to update the last operation' do + job.handle_timeout + + expect(action).to have_received(:update_last_operation_with_failure).with('Service Broker failed to deprovision within the required time.') + end + end + + describe '#operation' do + it 'returns "deprovision"' do + expect(job.operation).to eq(:deprovision) + end + end + + describe '#operation_type' do + it 'returns "delete"' do + expect(job.operation_type).to eq('delete') + end + end + + describe '#resource_type' do + it 'returns "service_instance"' do + expect(job.resource_type).to eq('service_instance') + end + end + + describe '#resource_guid' do + it 'returns the service instance guid' do + expect(job.resource_guid).to eq(service_instance.guid) + end + end + + describe '#display_name' do + it 'returns the display name' do + expect(job.display_name).to eq('service_instance.delete') + end + end + end + end +end diff --git a/spec/unit/models/runtime/pollable_job_model_spec.rb b/spec/unit/models/runtime/pollable_job_model_spec.rb index cb1d8bd91a5..7616c0045ec 100644 --- a/spec/unit/models/runtime/pollable_job_model_spec.rb +++ b/spec/unit/models/runtime/pollable_job_model_spec.rb @@ -167,5 +167,39 @@ module VCAP::CloudController expect(JobWarningModel.all).to be_empty end end + + describe '#sub_jobs' do + let(:root) { create(:pollable_job_model) } + let!(:sub1) { create(:pollable_job_model, root_job_guid: root.guid) } + let!(:sub2) { create(:pollable_job_model, root_job_guid: root.guid) } + let!(:unrelated) { create(:pollable_job_model) } + + it 'returns jobs whose root_job_guid matches this job\'s guid' do + expect(root.sub_jobs).to contain_exactly(sub1, sub2) + end + end + + describe '.find_active_delete' do + it 'returns a PROCESSING job for the given resource and operation' do + job = create(:pollable_job_model, state: 'PROCESSING', resource_guid: 'guid-1', operation: 'app.delete') + expect(PollableJobModel.find_active_delete(resource_guid: 'guid-1', operation: 'app.delete')).to eq(job) + end + + it 'returns a POLLING job for the given resource and operation' do + job = create(:pollable_job_model, state: 'POLLING', resource_guid: 'guid-1', operation: 'app.delete') + expect(PollableJobModel.find_active_delete(resource_guid: 'guid-1', operation: 'app.delete')).to eq(job) + end + + it 'ignores terminal jobs' do + create(:pollable_job_model, state: 'COMPLETE', resource_guid: 'guid-1', operation: 'app.delete') + create(:pollable_job_model, state: 'FAILED', resource_guid: 'guid-1', operation: 'app.delete') + expect(PollableJobModel.find_active_delete(resource_guid: 'guid-1', operation: 'app.delete')).to be_nil + end + + it 'ignores jobs for a different operation' do + create(:pollable_job_model, state: 'PROCESSING', resource_guid: 'guid-1', operation: 'something.else') + expect(PollableJobModel.find_active_delete(resource_guid: 'guid-1', operation: 'app.delete')).to be_nil + end + end end end diff --git a/spec/unit/repositories/app_event_repository_spec.rb b/spec/unit/repositories/app_event_repository_spec.rb index cb1f2762d7c..6b9f4972e8a 100644 --- a/spec/unit/repositories/app_event_repository_spec.rb +++ b/spec/unit/repositories/app_event_repository_spec.rb @@ -596,6 +596,22 @@ module Repositories expect(event.space).to eq(app.space) expect(event.space_guid).to eq(app.space.guid) end + + context 'when delete_triggered is true' do + it 'tags the audit event metadata so it is clear the stop was caused by an app delete' do + event = app_event_repository.record_app_stop(app, user_audit_info, delete_triggered: true) + + expect(event.metadata.fetch(:delete_triggered)).to be(true) + end + end + + context 'when delete_triggered is false (default)' do + it 'leaves the audit event metadata untouched' do + event = app_event_repository.record_app_stop(app, user_audit_info) + + expect(event.metadata).to be_nil.or(be_empty) + end + end end describe '#record_app_restart' do From 05ca1aaf87a23f30c433f61f00a15e1ad8de3d60 Mon Sep 17 00:00:00 2001 From: johha Date: Wed, 22 Jul 2026 17:45:46 +0200 Subject: [PATCH 2/4] temp logging --- app/jobs/enqueuer.rb | 10 +++++++++- app/jobs/logging_context_job.rb | 9 +++++++++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/app/jobs/enqueuer.rb b/app/jobs/enqueuer.rb index d4afe2d328f..6220a7c87b3 100644 --- a/app/jobs/enqueuer.rb +++ b/app/jobs/enqueuer.rb @@ -68,7 +68,15 @@ def enqueue_job(job, run_at: nil, priority_increment: nil, preserve_priority: fa local_opts[:priority] = final_priority if final_priority > 0 local_opts[:run_at] = run_at if run_at - Delayed::Job.enqueue(logging_context_job, @opts.merge(local_opts)) + delayed_job = Delayed::Job.enqueue(logging_context_job, @opts.merge(local_opts)) + + # TEMP recursive-delete debug (remove before PR ready): log every enqueue with its run_at. + Steno.logger('cc.temp.recursive_delete').info( + "ENQUEUE job=#{self.class.unwrap_job(job).class.name} dj_id=#{delayed_job.id} dj_guid=#{delayed_job.guid} " \ + "run_at=#{delayed_job.run_at.utc.iso8601} priority=#{delayed_job.priority} root_job_guid=#{root_job_guid || '-'}" + ) + + delayed_job end def load_delayed_job_plugins diff --git a/app/jobs/logging_context_job.rb b/app/jobs/logging_context_job.rb index 41e837c5713..e60a7e903b6 100644 --- a/app/jobs/logging_context_job.rb +++ b/app/jobs/logging_context_job.rb @@ -11,6 +11,15 @@ def initialize(handler, request_id) @request_id = request_id end + def before(job) + # TEMP recursive-delete debug (remove before PR ready): log every delayed job as it starts executing. + Steno.logger('cc.temp.recursive_delete').info( + "EXECUTE job=#{wrapped_handler.class.name} dj_id=#{job.id} dj_guid=#{job.guid} " \ + "run_at=#{job.run_at&.utc&.iso8601} attempts=#{job.attempts} queue=#{job.queue}" + ) + super + end + def perform with_request_id_set do logger.info("about to run job #{wrapped_handler.class.name}") From 52c1b1f202456639cb2d75946b7bfa64de228175 Mon Sep 17 00:00:00 2001 From: johha Date: Thu, 23 Jul 2026 13:29:01 +0200 Subject: [PATCH 3/4] fix logger --- app/jobs/mixins/root_job_mixin.rb | 4 +++- spec/unit/jobs/mixins/root_job_mixin_spec.rb | 12 ++++++++++++ 2 files changed, 15 insertions(+), 1 deletion(-) diff --git a/app/jobs/mixins/root_job_mixin.rb b/app/jobs/mixins/root_job_mixin.rb index d515a66e80c..7c9a9e8ee70 100644 --- a/app/jobs/mixins/root_job_mixin.rb +++ b/app/jobs/mixins/root_job_mixin.rb @@ -28,8 +28,10 @@ def perform_with_root_job_handling attr_reader :root_job, :sub_jobs # One shared tag for the whole recursive-delete subsystem so ops can grep it by a single logger name. + # NOT memoized: this job YAML-serialises itself on reschedule, and a cached Steno::Logger would drag + # its file-sink IO into the dump, reviving as an "uninitialized stream" that raises on the next write. def logger - @logger ||= Steno.logger('cc.jobs.v3.recursive_delete') + Steno.logger('cc.jobs.v3.recursive_delete') end # Resolves in any state (unlike root_job), so log lines stay searchable by this guid after the job settles. diff --git a/spec/unit/jobs/mixins/root_job_mixin_spec.rb b/spec/unit/jobs/mixins/root_job_mixin_spec.rb index 6692587049d..a263ed96db4 100644 --- a/spec/unit/jobs/mixins/root_job_mixin_spec.rb +++ b/spec/unit/jobs/mixins/root_job_mixin_spec.rb @@ -1,6 +1,7 @@ require 'spec_helper' require 'jobs/mixins/root_job_mixin' require 'jobs/reoccurring_job' +require 'jobs/v3/recursive_delete_app_job' module VCAP::CloudController module Jobs @@ -471,6 +472,17 @@ def sub_resource_errors end end end + + # Regression: the job YAML-serialises itself on reschedule. A memoized Steno::Logger dragged its + # file-sink IO into the dump, reviving as an "uninitialized stream" that raised on the next log write. + describe 'serialisation safety across reschedule' do + it 'does not carry the logger into its YAML dump' do + job = VCAP::CloudController::V3::RecursiveDeleteAppJob.new('app-guid', nil) + job.send(:logger).info('warm up the logger so a memoised one would be cached before dumping') + + expect(YAML.dump(job)).not_to include('Steno::Logger') + end + end end end end From 10296462636550f062ec804ab135654d28393d2c Mon Sep 17 00:00:00 2001 From: johha Date: Thu, 23 Jul 2026 13:46:23 +0200 Subject: [PATCH 4/4] improve audit log so it shows in CF CLI --- app/repositories/app_event_repository.rb | 3 ++- spec/unit/repositories/app_event_repository_spec.rb | 2 ++ 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/app/repositories/app_event_repository.rb b/app/repositories/app_event_repository.rb index 947f8008e19..64347ba6690 100644 --- a/app/repositories/app_event_repository.rb +++ b/app/repositories/app_event_repository.rb @@ -79,7 +79,8 @@ def record_app_stop(app, user_audit_info, delete_triggered: false) VCAP::AppLogEmitter.emit(app.guid, "Stopping app with guid #{app.guid}") actor = { name: user_audit_info.user_email, guid: user_audit_info.user_guid, user_name: user_audit_info.user_name, type: 'user' } - metadata = delete_triggered ? { delete_triggered: true } : nil + # 'reason' is on the cf CLI's known-metadata allowlist, so it surfaces in `cf events`; delete_triggered stays for API consumers. + metadata = delete_triggered ? { delete_triggered: true, reason: 'stopped as part of app deletion' } : nil create_app_audit_event(EventTypes::APP_STOP, app, app.space, actor, metadata) end diff --git a/spec/unit/repositories/app_event_repository_spec.rb b/spec/unit/repositories/app_event_repository_spec.rb index 6b9f4972e8a..e2115cfdb06 100644 --- a/spec/unit/repositories/app_event_repository_spec.rb +++ b/spec/unit/repositories/app_event_repository_spec.rb @@ -602,6 +602,8 @@ module Repositories event = app_event_repository.record_app_stop(app, user_audit_info, delete_triggered: true) expect(event.metadata.fetch(:delete_triggered)).to be(true) + # 'reason' surfaces in `cf events` so operators see why the app was stopped. + expect(event.metadata.fetch(:reason)).to eq('stopped as part of app deletion') end end