Class: S3DirectMultipartUpload::UploadSession
- Inherits:
-
ApplicationRecord
- Object
- ApplicationRecord
- S3DirectMultipartUpload::UploadSession
- Defined in:
- app/models/s3_direct_multipart_upload/upload_session.rb
Constant Summary collapse
- MAX_PARTS =
10_000- MIN_MULTIPART_CHUNK_SIZE =
5.megabytes
- MAX_CHUNK_SIZE =
5.gigabytes
- SINGLE_PART_THRESHOLD =
MIN_MULTIPART_CHUNK_SIZE- PRESIGNED_URL_TTL =
24.hours
Class Method Summary collapse
- .build_object_key(key_prefix:, filename:) ⇒ Object
- .configuration_error(message) ⇒ Object
- .default_bucket ⇒ Object
- .default_chunk_size ⇒ Object
- .default_key_prefix(session_id) ⇒ Object
- .dev_endpoint_base_url ⇒ Object
- .dev_mode_enabled? ⇒ Boolean
- .dev_presigned_url(object_key:, upload_id:, part_number:) ⇒ Object
- .dev_signature(path:, expires_at:, upload_id:, part_number:) ⇒ Object
- .disk_service? ⇒ Boolean
- .expected_parts_for(byte_size:, chunk_size:) ⇒ Object
- .normalize_config(config) ⇒ Object
- .resolve_chunk_size(byte_size:, chunk_size:) ⇒ Object
- .s3_client ⇒ Object
- .start!(filename:, byte_size:, content_type:, chunk_size: nil, metadata: {}) ⇒ Object
- .storage_config ⇒ Object
- .storage_config_path ⇒ Object
- .storage_configurations ⇒ Object
- .stub_responses_option(config) ⇒ Object
- .upload_namespace ⇒ Object
- .validate_config!(config) ⇒ Object
Instance Method Summary collapse
- #complete!(checksum: nil) ⇒ Object
- #default_chunk_size ⇒ Object
- #expected_parts ⇒ Object
- #missing_parts ⇒ Object
- #object_key ⇒ Object
- #part_limit ⇒ Object
- #parts_uploaded? ⇒ Boolean
- #presigned_parts(part_numbers:) ⇒ Object
- #report_part!(part_number:, etag:) ⇒ Object
- #schedule_expiration! ⇒ Object
Class Method Details
.build_object_key(key_prefix:, filename:) ⇒ Object
278 279 280 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 278 def build_object_key(key_prefix:, filename:) File.join(key_prefix, filename) end |
.configuration_error(message) ⇒ Object
389 390 391 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 389 def configuration_error() RuntimeError.new() end |
.default_bucket ⇒ Object
233 234 235 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 233 def default_bucket storage_config.fetch(:bucket) end |
.default_chunk_size ⇒ Object
246 247 248 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 246 def default_chunk_size 128.megabytes end |
.default_key_prefix(session_id) ⇒ Object
237 238 239 240 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 237 def default_key_prefix(session_id) date_prefix = Time.current.strftime("%Y%m%d") File.join(upload_namespace, date_prefix, session_id) end |
.dev_endpoint_base_url ⇒ Object
367 368 369 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 367 def dev_endpoint_base_url ENV.fetch("S3_DIRECT_DEV_ENDPOINT", "http://localhost:3000") end |
.dev_mode_enabled? ⇒ Boolean
385 386 387 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 385 def dev_mode_enabled? Rails.env.development? || Rails.env.test? end |
.dev_presigned_url(object_key:, upload_id:, part_number:) ⇒ Object
337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 337 def dev_presigned_url(object_key:, upload_id:, part_number:) raise configuration_error("dev presigned URLs are only available for Disk service") unless disk_service? expires_at = PRESIGNED_URL_TTL.from_now signature = dev_signature( path: object_key, expires_at:, upload_id:, part_number: ) path = S3DirectMultipartUpload::Engine.routes.url_helpers.dev_storage_upload_path(path: object_key) url = URI.join(dev_endpoint_base_url, path) query = Rack::Utils.build_query( uploadId: upload_id, partNumber: part_number, expires_at: expires_at.to_i, signature: ) url.query = query url.to_s end |
.dev_signature(path:, expires_at:, upload_id:, part_number:) ⇒ Object
360 361 362 363 364 365 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 360 def dev_signature(path:, expires_at:, upload_id:, part_number:) secret = ENV.fetch("S3_DIRECT_DEV_SECRET", "dev-secret") normalized_path = path.to_s.sub(/\A\//, "") normalized_part = part_number.presence&.to_s Digest::SHA256.hexdigest([ normalized_path, expires_at.to_i, upload_id.to_s, normalized_part, secret ].compact.join(":")) end |
.disk_service? ⇒ Boolean
333 334 335 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 333 def disk_service? storage_config[:service].to_s.casecmp("disk").zero? end |
.expected_parts_for(byte_size:, chunk_size:) ⇒ Object
261 262 263 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 261 def expected_parts_for(byte_size:, chunk_size:) (byte_size.to_f / chunk_size).ceil end |
.normalize_config(config) ⇒ Object
321 322 323 324 325 326 327 328 329 330 331 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 321 def normalize_config(config) service = config[:service].to_s.downcase return config if service == "s3" config.merge( service: "Disk", bucket: config[:bucket] || "s3_direct_disk", region: config[:region] || "us-east-1", stub_responses: true ) end |
.resolve_chunk_size(byte_size:, chunk_size:) ⇒ Object
250 251 252 253 254 255 256 257 258 259 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 250 def resolve_chunk_size(byte_size:, chunk_size:) return byte_size if byte_size < SINGLE_PART_THRESHOLD requested = chunk_size.presence&.to_i || default_chunk_size raise ArgumentError, "chunk_size must be positive" if requested <= 0 raise ArgumentError, "chunk_size must be at least #{MIN_MULTIPART_CHUNK_SIZE}" if requested < MIN_MULTIPART_CHUNK_SIZE raise ArgumentError, "chunk_size must be at most #{MAX_CHUNK_SIZE}" if requested > MAX_CHUNK_SIZE requested end |
.s3_client ⇒ Object
265 266 267 268 269 270 271 272 273 274 275 276 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 265 def s3_client @s3_client ||= begin config = storage_config = { region: config.fetch(:region), endpoint: config[:endpoint], force_path_style: config[:force_path_style], stub_responses: stub_responses_option(config) }.compact Aws::S3::Client.new(**) end end |
.start!(filename:, byte_size:, content_type:, chunk_size: nil, metadata: {}) ⇒ Object
55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 55 def self.start!(filename:, byte_size:, content_type:, chunk_size: nil, metadata: {}) byte_size = byte_size.to_i raise ArgumentError, "byte_size must be positive" if byte_size <= 0 resolved_chunk_size = resolve_chunk_size(byte_size:, chunk_size:) expected_parts = expected_parts_for(byte_size:, chunk_size: resolved_chunk_size) raise ArgumentError, "part number limit exceeded" if expected_parts > MAX_PARTS session_id = SecureRandom.uuid key_prefix = default_key_prefix(session_id) object_key = build_object_key(key_prefix:, filename:) response = s3_client.create_multipart_upload( bucket: default_bucket, key: object_key, content_type: ) create!( session_id:, upload_id: response.upload_id, bucket: response.bucket || default_bucket, key_prefix:, filename:, byte_size:, content_type:, chunk_size: resolved_chunk_size, metadata: ).tap(&:schedule_expiration!) end |
.storage_config ⇒ Object
282 283 284 285 286 287 288 289 290 291 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 282 def storage_config @storage_config ||= begin config = storage_configurations["s3_direct_multipart"] || storage_configurations[:s3_direct_multipart] raise configuration_error("ActiveStorage service `s3_direct_multipart` is not configured") unless config config = normalize_config(config.deep_symbolize_keys) validate_config!(config) config end end |
.storage_config_path ⇒ Object
306 307 308 309 310 311 312 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 306 def storage_config_path env_path = Rails.root.join("config/storage/#{Rails.env}.yml") return env_path if env_path.exist? default_path = Rails.root.join("config/storage.yml") default_path if default_path.exist? end |
.storage_configurations ⇒ Object
293 294 295 296 297 298 299 300 301 302 303 304 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 293 def storage_configurations explicit = Rails.application.config.respond_to?(:active_storage) && Rails.application.config.active_storage.respond_to?(:service_configurations) && Rails.application.config.active_storage.service_configurations.present? return Rails.application.config.active_storage.service_configurations if explicit path = storage_config_path raise configuration_error("ActiveStorage service configurations are empty; add config/storage.yml or config/storage/<env>.yml") unless path erb = ERB.new(path.read).result YAML.safe_load(erb, aliases: true) || {} end |
.stub_responses_option(config) ⇒ Object
371 372 373 374 375 376 377 378 379 380 381 382 383 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 371 def stub_responses_option(config) return config[:stub_responses] unless disk_service? && config[:stub_responses] return config[:stub_responses] if config[:stub_responses].respond_to?(:to_hash) { create_multipart_upload: ->(*) { { upload_id: SecureRandom.uuid, bucket: config[:bucket] || default_bucket } } } end |
.upload_namespace ⇒ Object
242 243 244 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 242 def upload_namespace "uploads/s3_direct" end |
.validate_config!(config) ⇒ Object
314 315 316 317 318 319 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 314 def validate_config!(config) service = config[:service].to_s.downcase raise configuration_error("ActiveStorage service `s3_direct_multipart` must set service: S3 or Disk") unless %w[s3 disk].include?(service) raise configuration_error("ActiveStorage service `s3_direct_multipart` requires bucket") if config[:bucket].blank? raise configuration_error("ActiveStorage service `s3_direct_multipart` requires region") if config[:region].blank? end |
Instance Method Details
#complete!(checksum: nil) ⇒ Object
107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 107 def complete!(checksum: nil) missing = missing_parts raise "missing parts: #{missing}" if missing.present? self.class.s3_client.complete_multipart_upload( bucket:, key: object_key, upload_id:, multipart_upload: { parts: parts.order(:part_number).map { { part_number: _1.part_number, etag: _1.etag } } } ) completed! update!(metadata: .merge("checksum" => checksum).compact) combine_dev_parts_if_needed self end |
#default_chunk_size ⇒ Object
155 156 157 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 155 def default_chunk_size self.class.default_chunk_size end |
#expected_parts ⇒ Object
133 134 135 136 137 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 133 def expected_parts return 0 if chunk_size.blank? || byte_size.blank? self.class.expected_parts_for(byte_size:, chunk_size:) end |
#missing_parts ⇒ Object
143 144 145 146 147 148 149 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 143 def missing_parts return [] unless expected_parts.positive? uploaded = parts.pluck(:part_number).uniq.sort expected = (1..expected_parts).to_a expected - uploaded end |
#object_key ⇒ Object
125 126 127 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 125 def object_key @object_key ||= self.class.build_object_key(key_prefix:, filename:) end |
#part_limit ⇒ Object
139 140 141 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 139 def part_limit MAX_PARTS end |
#parts_uploaded? ⇒ Boolean
129 130 131 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 129 def parts_uploaded? expected_parts.positive? && missing_parts.empty? end |
#presigned_parts(part_numbers:) ⇒ Object
86 87 88 89 90 91 92 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 86 def presigned_parts(part_numbers:) raise ArgumentError, "part number limit exceeded" if part_numbers.size > MAX_PARTS part_numbers.map do |number| disk_service? ? dev_presigned_part(number) : aws_presigned_part(number) end end |
#report_part!(part_number:, etag:) ⇒ Object
94 95 96 97 98 99 100 101 102 103 104 105 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 94 def report_part!(part_number:, etag:) part_number = part_number.to_i raise ArgumentError, "part_number is invalid" if part_number.zero? || part_number.negative? raise ArgumentError, "part_number exceeds expected parts" if expected_parts.positive? && part_number > expected_parts raise ArgumentError, "upload already completed or aborted" if completed? || aborted? transaction do record = parts.find_or_initialize_by(part_number:) record.update!(etag:, uploaded_at: Time.current) uploading! if pending? end end |
#schedule_expiration! ⇒ Object
151 152 153 |
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 151 def schedule_expiration! update!(expires_at: PRESIGNED_URL_TTL.from_now) end |