Class: Gapic::ResumableUpload
- Inherits:
-
Object
- Object
- Gapic::ResumableUpload
- Defined in:
- lib/gapic/resumable_upload.rb
Overview
Coordinates resumable uploads for a client method that performs them.
A client method that uploads media returns one of these handles instead of a response. No request is sent and no byte is read from the stream until #start or #resume is called on it. Both are synchronous: they block the calling thread for the whole upload and return the decoded response message.
Reusable
A handle is reusable, and a failed run is resumed on the same object:
The stream handed to #resume must be positioned at byte 0 of the whole object, not at the server's acknowledged offset; the upload fast-forwards on its own, by seeking on a seekable stream or by reading and discarding on an unseekable one. An unseekable stream therefore has to be freshly opened rather than rewound.
A run that failed in a way the protocol can recover from leaves a Gapic::Rest::ResumableUpload::ResumeHandle behind, readable from #resume_handle and also carried on the error. Persisting that handle lets a later process resume the same upload:
A completed upload is finalized: #resume_handle returns nil and #resumable? returns false, so
there is no handle to resume from. Calling #start again is permitted and begins a second, unrelated
upload.
The Two Timeouts
An upload is bounded by two independent budgets, and they are three orders of magnitude apart:
| Budget | Set by | Covers |
|---|---|---|
| whole upload | upload_timeout: on #start and #resume |
every request, retry and byte of the run |
| initiation request | per-call timeout |
creating the session |
The per-call timeout a client method takes reaches only the initiation request. An upload still
transferring bytes an hour later has long outlived it, and that is expected. To bound the run as a
whole, pass upload_timeout:.
Threading
#start and #resume block the calling thread, and the on_progress callback runs on that same
thread. The readers (#resume_handle, #resumable?, #running?) are guarded by an internal mutex and
may be called from another thread mid-run; values read that way are a best-effort snapshot of a state
the upload thread is still advancing.
Defaults
chunk_sizedefaults to 8 MB, then rounds down to a multiple of any chunk granularity the server requires.upload_timeoutdefaults toupload_size / 1 MB per secondwhenupload_sizeis known, floored at one hour, and to one hour flat when it is not.
Retry Policies
Retry behavior is partitioned across three policies. Only the initiation policy is caller-supplied; the other two are the protocol's own and govern the requests no call option describes.
| Policy | Governs | Default retry_codes |
|---|---|---|
| initiation | session initiation | the 4xx and 5xx sets below |
| control plane | query and cancel |
the 4xx and 5xx sets below |
| data plane | upload and finalize |
the 5xx set below only |
- 4xx set:
ALREADY_EXISTS(HTTP409),RESOURCE_EXHAUSTED(429),CANCELLED(499). - 5xx set:
INTERNAL(HTTP500),UNAVAILABLE(503),DEADLINE_EXCEEDED(504).
None of the defaults carries a retry_predicate, so one supplied by the caller is consulted as-is,
ahead of retry_codes. All three share the same backoff: initial_delay 1.0 s, max_delay
15.0 s, multiplier 1.3.
Some decisions belong to the protocol and are made before a policy is asked:
- Initiation and control plane requests are re-sent, within the policy's deadline, after a connection
or TLS failure, and after a
200response missingX-Goog-Upload-Status. - Data plane requests are never re-sent after an outcome that leaves the server offset unknown — a
timeout, a connection failure, a missing status header, or any
4xx. The upload re-queries the session and resumes from the offset the server reports instead. - Consecutive recovery attempts that make no progress back off on a single schedule, with the same delays as above; the schedule starts over once the server confirms new bytes.
- A response carrying
X-Goog-Upload-Status: finalis never retried.
Instance Method Summary collapse
-
#resumable? ⇒ Boolean
Returns whether the last run left an upload that can be resumed.
-
#resume(stream:, resume_handle: nil, content_type: nil, upload_size: nil, upload_timeout: nil, on_progress: nil) ⇒ Object
Resumes an upload session the server has already created, transferring whatever it has not yet acknowledged.
-
#resume_handle ⇒ Gapic::Rest::ResumableUpload::ResumeHandle?
Returns the handle needed to resume the last run, carrying its upload URL and resolved chunk size.
-
#running? ⇒ Boolean
Returns whether a run is currently executing.
-
#start(stream:, content_type: nil, upload_size: nil, chunk_size: nil, upload_timeout: nil, on_progress: nil) ⇒ Object
Creates an upload session on the server and transfers the stream into it.
Instance Method Details
#resumable? ⇒ Boolean
Returns whether the last run left an upload that can be resumed.
323 324 325 |
# File 'lib/gapic/resumable_upload.rb', line 323 def resumable? !resume_handle.nil? end |
#resume(stream:, resume_handle: nil, content_type: nil, upload_size: nil, upload_timeout: nil, on_progress: nil) ⇒ Object
Resumes an upload session the server has already created, transferring whatever it has not yet acknowledged.
Blocks the calling thread until the upload completes or fails. A resumed run sends no initiation request, so it takes no initiation arguments and needs no request message.
The target is a Gapic::Rest::ResumableUpload::ResumeHandle: either the one passed in, or — when
resume_handle is omitted — the one left behind by this handle's last run. The bare form raises
ArgumentError when there is none, which covers a handle that has never run, a run that finished
successfully, and a run that failed in a way the protocol considers unresumable.
Resuming against a finalized upload URL is undefined behavior: it queries the server and might return the response body or raise an error, depending on the server response.
283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 |
# File 'lib/gapic/resumable_upload.rb', line 283 def resume stream:, resume_handle: nil, content_type: nil, upload_size: nil, upload_timeout: nil, on_progress: nil execute_run do verify_stream_at_origin stream handle = resolve_resume_handle resume_handle client_stub = @client_stub_proc.call config = ::Gapic::Rest::ResumableUpload::ResumeUploadConfig.new( upload_url: handle.upload_url, chunk_size: handle.chunk_size, **run_config_args(stream: stream, content_type: content_type, upload_size: upload_size, upload_timeout: upload_timeout, on_progress: on_progress) ) build_driver client_stub, config end end |
#resume_handle ⇒ Gapic::Rest::ResumableUpload::ResumeHandle?
Returns the handle needed to resume the last run, carrying its upload URL and resolved chunk size.
nil before the first run, and after any run that left nothing to resume: a completed upload is
finalized, and rejected and cancelled uploads are too.
Each run replaces this value rather than accumulating handles, so it always describes the most recent one. Starting a second upload therefore discards whatever the previous run left behind — persist the handle first if the earlier upload still matters. A call that fails while building its configuration does not count as a run: it never reaches a driver, so the earlier value survives.
315 316 317 |
# File 'lib/gapic/resumable_upload.rb', line 315 def resume_handle @mutex.synchronize { @driver&.resume_handle } end |
#running? ⇒ Boolean
Returns whether a run is currently executing.
331 332 333 |
# File 'lib/gapic/resumable_upload.rb', line 331 def running? @mutex.synchronize { @running } end |
#start(stream:, content_type: nil, upload_size: nil, chunk_size: nil, upload_timeout: nil, on_progress: nil) ⇒ Object
Creates an upload session on the server and transfers the stream into it.
Blocks the calling thread until the upload completes or fails. The stream is assumed to be positioned at byte 0 (it is not rewound before reading) and is not closed after use.
219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 |
# File 'lib/gapic/resumable_upload.rb', line 219 def start stream:, content_type: nil, upload_size: nil, chunk_size: nil, upload_timeout: nil, on_progress: nil execute_run do client_stub = @client_stub_proc.call initial_url, initial_body = @initial_request_proc.call config = ::Gapic::Rest::ResumableUpload::StartUploadConfig.new( initial_url: initial_url, initial_body: initial_body, initial_headers: @initial_headers, chunk_size: chunk_size, start_retry_policy: @start_retry_policy, **run_config_args(stream: stream, content_type: content_type, upload_size: upload_size, upload_timeout: upload_timeout, on_progress: on_progress) ) build_driver client_stub, config end end |