Class: S3DirectMultipartUpload::UploadSession

Inherits:
ApplicationRecord
  • Object
show all
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

Instance Method Summary collapse

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(message)
  RuntimeError.new(message)
end

.default_bucketObject



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_sizeObject



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_urlObject



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

Returns:



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

Returns:



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

Raises:



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_clientObject



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
    client_options = {
      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(**client_options)
  end
end

.start!(filename:, byte_size:, content_type:, chunk_size: nil, metadata: {}) ⇒ Object

Raises:



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_configObject



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_pathObject



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_configurationsObject



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_namespaceObject



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_sizeObject



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_partsObject



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_partsObject



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_keyObject



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_limitObject



139
140
141
# File 'app/models/s3_direct_multipart_upload/upload_session.rb', line 139

def part_limit
  MAX_PARTS
end

#parts_uploaded?Boolean

Returns:



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

Raises:



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

Raises:



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