| @@ -4,8 +4,9 @@ | ||
| 4 | 4 | |
| 5 | 5 | use Dudlewebs\WPMCS\s3\Aws; |
| 6 | 6 | use Dudlewebs\WPMCS\s3\Aws\CommandInterface; |
| 7 | 7 | use Dudlewebs\WPMCS\s3\Aws\Exception\AwsException; |
| 8 | +use Dudlewebs\WPMCS\s3\Aws\MetricsBuilder; | |
| 8 | 9 | use Dudlewebs\WPMCS\s3\GuzzleHttp\Promise; |
| 9 | 10 | use Dudlewebs\WPMCS\s3\GuzzleHttp\Promise\PromiseInterface; |
| 10 | 11 | use Dudlewebs\WPMCS\s3\GuzzleHttp\Promise\PromisorInterface; |
| 11 | 12 | use Iterator; |
| @@ -25,9 +26,11 @@ | ||
| 25 | 26 | private $destination; |
| 26 | 27 | private $concurrency; |
| 27 | 28 | private $mupThreshold; |
| 28 | 29 | private $before; |
| 30 | + private $after; | |
| 29 | 31 | private $s3Args = []; |
| 32 | + private $addContentMD5; | |
| 30 | 33 | /** |
| 31 | 34 | * When providing the $source argument, you may provide a string referencing |
| 32 | 35 | * the path to a directory on disk to upload, an s3 scheme URI that contains |
| 33 | 36 | * the bucket and key (e.g., "s3://bucket/key"), or an \Iterator object |
| @@ -49,8 +52,13 @@ | ||
| 49 | 52 | * - before: (callable) A callback to invoke before each transfer. The |
| 50 | 53 | * callback accepts a single argument: Aws\CommandInterface $command. |
| 51 | 54 | * The provided command will be either a GetObject, PutObject, |
| 52 | 55 | * InitiateMultipartUpload, or UploadPart command. |
| 56 | + * - after: (callable) A callback to invoke after each transfer promise is fulfilled. | |
| 57 | + * The function is invoked with three arguments: the fulfillment value, the index | |
| 58 | + * position from the iterable list of the promise, and the aggregate | |
| 59 | + * promise that manages all the promises. The aggregate promise may | |
| 60 | + * be resolved from within the callback to short-circuit the promise. | |
| 53 | 61 | * - mup_threshold: (int) Size in bytes in which a multipart upload should |
| 54 | 62 | * be used instead of PutObject. Defaults to 20971520 (20 MB). |
| 55 | 63 | * - concurrency: (int, default=5) Number of files to upload concurrently. |
| 56 | 64 | * The ideal concurrency value will vary based on the number of files |
| @@ -103,8 +111,15 @@ | ||
| 103 | 111 | if (!\is_callable($this->before)) { |
| 104 | 112 | throw new \InvalidArgumentException('before must be a callable.'); |
| 105 | 113 | } |
| 106 | 114 | } |
| 115 | + // Handle "after" callback option. | |
| 116 | + if (isset($options['after'])) { | |
| 117 | + $this->after = $options['after']; | |
| 118 | + if (!\is_callable($this->after)) { | |
| 119 | + throw new \InvalidArgumentException('after must be a callable.'); | |
| 120 | + } | |
| 121 | + } | |
| 107 | 122 | // Handle "debug" option. |
| 108 | 123 | if (isset($options['debug'])) { |
| 109 | 124 | if ($options['debug'] === \true) { |
| 110 | 125 | $options['debug'] = \fopen('php://output', 'w'); |
| @@ -112,8 +127,11 @@ | ||
| 112 | 127 | if (\is_resource($options['debug'])) { |
| 113 | 128 | $this->addDebugToBefore($options['debug']); |
| 114 | 129 | } |
| 115 | 130 | } |
| 131 | + // Handle "add_content_md5" option. | |
| 132 | + $this->addContentMD5 = isset($options['add_content_md5']) && $options['add_content_md5'] === \true; | |
| 133 | + MetricsBuilder::appendMetricsCaptureMiddleware($this->client->getHandlerList(), MetricsBuilder::S3_TRANSFER); | |
| 116 | 134 | } |
| 117 | 135 | /** |
| 118 | 136 | * Transfers the files. |
| 119 | 137 | * |
| @@ -118,9 +136,9 @@ | ||
| 118 | 136 | * Transfers the files. |
| 119 | 137 | * |
| 120 | 138 | * @return PromiseInterface |
| 121 | 139 | */ |
| 122 | - public function promise() | |
| 140 | + public function promise() : PromiseInterface | |
| 123 | 141 | { |
| 124 | 142 | // If the promise has been created, just return it. |
| 125 | 143 | if (!$this->promise) { |
| 126 | 144 | // Create an upload/download promise for the transfer. |
| @@ -224,9 +242,9 @@ | ||
| 224 | 242 | // Create the command. |
| 225 | 243 | $commands[] = $command; |
| 226 | 244 | } |
| 227 | 245 | // Create a GetObject command pool and return the promise. |
| 228 | - return (new Aws\CommandPool($this->client, $commands, ['concurrency' => $this->concurrency, 'before' => $this->before, 'rejected' => function ($reason, $idx, Promise\PromiseInterface $p) { | |
| 246 | + return (new Aws\CommandPool($this->client, $commands, ['concurrency' => $this->concurrency, 'before' => $this->before, 'fulfill' => $this->after, 'rejected' => function ($reason, $idx, Promise\PromiseInterface $p) { | |
| 229 | 247 | $p->reject($reason); |
| 230 | 248 | }]))->promise(); |
| 231 | 249 | } |
| 232 | 250 | private function createUploadPromise() |
| @@ -236,9 +254,9 @@ | ||
| 236 | 254 | return \filesize($file) >= $this->mupThreshold ? $this->uploadMultipart($file) : $this->upload($file); |
| 237 | 255 | }); |
| 238 | 256 | // Create an EachPromise, that will concurrently handle the upload |
| 239 | 257 | // operations' yielded promises from the iterator. |
| 240 | - return Promise\Each::ofLimitAll($files, $this->concurrency); | |
| 258 | + return Promise\Each::ofLimitAll($files, $this->concurrency, $this->after); | |
| 241 | 259 | } |
| 242 | 260 | /** @return Iterator */ |
| 243 | 261 | private function getUploadsIterator() |
| 244 | 262 | { |
| @@ -272,8 +290,9 @@ | ||
| 272 | 290 | { |
| 273 | 291 | $args = $this->s3Args; |
| 274 | 292 | $args['SourceFile'] = $filename; |
| 275 | 293 | $args['Key'] = $this->createS3Key($filename); |
| 294 | + $args['AddContentMD5'] = $this->addContentMD5; | |
| 276 | 295 | $command = $this->client->getCommand('PutObject', $args); |
| 277 | 296 | $this->before and \call_user_func($this->before, $command); |
| 278 | 297 | return $this->client->executeAsync($command); |
| 279 | 298 | } |
| @@ -281,9 +300,9 @@ | ||
| 281 | 300 | { |
| 282 | 301 | $args = $this->s3Args; |
| 283 | 302 | $args['Key'] = $this->createS3Key($filename); |
| 284 | 303 | $filename = $filename instanceof \SplFileInfo ? $filename->getPathname() : $filename; |
| 285 | - return (new MultipartUploader($this->client, $filename, ['bucket' => $args['Bucket'], 'key' => $args['Key'], 'before_initiate' => $this->before, 'before_upload' => $this->before, 'before_complete' => $this->before, 'concurrency' => $this->concurrency]))->promise(); | |
| 304 | + return (new MultipartUploader($this->client, $filename, ['bucket' => $args['Bucket'], 'key' => $args['Key'], 'before_initiate' => $this->before, 'before_upload' => $this->before, 'before_complete' => $this->before, 'concurrency' => $this->concurrency, 'add_content_md5' => $this->addContentMD5]))->promise(); | |
| 286 | 305 | } |
| 287 | 306 | private function createS3Key($filename) |
| 288 | 307 | { |
| 289 | 308 | $filename = $this->normalizePath($filename); |