| 47 | } |
| 48 | |
| 49 | void DeleteS3Object::onTrigger(const std::shared_ptr<core::ProcessContext> &context, const std::shared_ptr<core::ProcessSession> &session) { |
| 50 | logger_->log_debug("DeleteS3Object onTrigger"); |
| 51 | std::shared_ptr<core::FlowFile> flow_file = session->get(); |
| 52 | if (!flow_file) { |
| 53 | context->yield(); |
| 54 | return; |
| 55 | } |
| 56 | |
| 57 | auto common_properties = getCommonELSupportedProperties(context, flow_file); |
| 58 | if (!common_properties) { |
| 59 | session->transfer(flow_file, Failure); |
| 60 | return; |
| 61 | } |
| 62 | |
| 63 | std::string version; |
| 64 | context->getProperty(Version, version, flow_file); |
| 65 | logger_->log_debug("DeleteS3Object: Version [%s]", version); |
| 66 | |
| 67 | bool delete_succeeded = false; |
| 68 | { |
| 69 | std::lock_guard<std::mutex> lock(s3_wrapper_mutex_); |
| 70 | configureS3Wrapper(common_properties.value()); |
| 71 | delete_succeeded = s3_wrapper_->deleteObject(common_properties->bucket, common_properties->object_key, version); |
| 72 | } |
| 73 | |
| 74 | if (delete_succeeded) { |
| 75 | logger_->log_debug("Successfully deleted S3 object '%s' from bucket '%s'", common_properties->object_key, common_properties->bucket); |
| 76 | session->transfer(flow_file, Success); |
| 77 | } else { |
| 78 | logger_->log_error("Failed to delete S3 object '%s' from bucket '%s'", common_properties->object_key, common_properties->bucket); |
| 79 | session->transfer(flow_file, Failure); |
| 80 | } |
| 81 | } |
| 82 | |
| 83 | } // namespace processors |
| 84 | } // namespace aws |
nothing calls this directly
no test coverage detected