PluginProbe
Media Cloud Sync / 1.2.13
Media Cloud Sync v1.2.13
1.4.2 1.4.1 1.4.0 1.3.12 1.3.11 1.3.10 trunk 1.0.0 1.0.1 1.0.2 1.0.3 1.1.0 1.1.1 1.2.0 1.2.10 1.2.11 1.2.12 1.2.13 1.2.2 1.2.3 1.2.4 1.2.5 1.2.6 1.2.7 1.2.8 All 36 releases
media-cloud-sync / includes / sdk / s3 / Aws / S3 / Transfer.php

Transfer.php in Media Cloud Sync 1.2.13, at includes/sdk/s3/Aws/S3/Transfer.php

341 lines 14.6 KB
No matching file
Up and down to move Enter to open Esc to close
Raw Download Zip
1 <?php
2
3 namespace Dudlewebs\WPMCS\s3\Aws\S3;
4
5 use Dudlewebs\WPMCS\s3\Aws;
6 use Dudlewebs\WPMCS\s3\Aws\CommandInterface;
7 use Dudlewebs\WPMCS\s3\Aws\Exception\AwsException;
8 use Dudlewebs\WPMCS\s3\GuzzleHttp\Promise;
9 use Dudlewebs\WPMCS\s3\GuzzleHttp\Promise\PromiseInterface;
10 use Dudlewebs\WPMCS\s3\GuzzleHttp\Promise\PromisorInterface;
11 use Iterator;
12 /**
13 * Transfers files from the local filesystem to S3 or from S3 to the local
14 * filesystem.
15 *
16 * This class does not support copying from the local filesystem to somewhere
17 * else on the local filesystem or from one S3 bucket to another.
18 */
19 class Transfer implements PromisorInterface
20 {
21 private $client;
22 private $promise;
23 private $source;
24 private $sourceMetadata;
25 private $destination;
26 private $concurrency;
27 private $mupThreshold;
28 private $before;
29 private $s3Args = [];
30 private $addContentMD5;
31 /**
32 * When providing the $source argument, you may provide a string referencing
33 * the path to a directory on disk to upload, an s3 scheme URI that contains
34 * the bucket and key (e.g., "s3://bucket/key"), or an \Iterator object
35 * that yields strings containing filenames that are the path to a file on
36 * disk or an s3 scheme URI. The bucket portion of the s3 URI may be an S3
37 * access point ARN. The "/key" portion of an s3 URI is optional.
38 *
39 * When providing an iterator for the $source argument, you must also
40 * provide a 'base_dir' key value pair in the $options argument.
41 *
42 * The $dest argument can be the path to a directory on disk or an s3
43 * scheme URI (e.g., "s3://bucket/key").
44 *
45 * The options array can contain the following key value pairs:
46 *
47 * - base_dir: (string) Base directory of the source, if $source is an
48 * iterator. If the $source option is not an array, then this option is
49 * ignored.
50 * - before: (callable) A callback to invoke before each transfer. The
51 * callback accepts a single argument: Aws\CommandInterface $command.
52 * The provided command will be either a GetObject, PutObject,
53 * InitiateMultipartUpload, or UploadPart command.
54 * - mup_threshold: (int) Size in bytes in which a multipart upload should
55 * be used instead of PutObject. Defaults to 20971520 (20 MB).
56 * - concurrency: (int, default=5) Number of files to upload concurrently.
57 * The ideal concurrency value will vary based on the number of files
58 * being uploaded and the average size of each file. Generally speaking,
59 * smaller files benefit from a higher concurrency while larger files
60 * will not.
61 * - debug: (bool) Set to true to print out debug information for
62 * transfers. Set to an fopen() resource to write to a specific stream
63 * rather than writing to STDOUT.
64 *
65 * @param S3ClientInterface $client Client used for transfers.
66 * @param string|Iterator $source Where the files are transferred from.
67 * @param string $dest Where the files are transferred to.
68 * @param array $options Hash of options.
69 */
70 public function __construct(S3ClientInterface $client, $source, $dest, array $options = [])
71 {
72 $this->client = $client;
73 // Prepare the destination.
74 $this->destination = $this->prepareTarget($dest);
75 if ($this->destination['scheme'] === 's3') {
76 $this->s3Args = $this->getS3Args($this->destination['path']);
77 }
78 // Prepare the source.
79 if (\is_string($source)) {
80 $this->sourceMetadata = $this->prepareTarget($source);
81 $this->source = $source;
82 } elseif ($source instanceof Iterator) {
83 if (empty($options['base_dir'])) {
84 throw new \InvalidArgumentException('You must provide the source' . ' argument as a string or provide the "base_dir" option.');
85 }
86 $this->sourceMetadata = $this->prepareTarget($options['base_dir']);
87 $this->source = $source;
88 } else {
89 throw new \InvalidArgumentException('source must be the path to a ' . 'directory or an iterator that yields file names.');
90 }
91 // Validate schemes.
92 if ($this->sourceMetadata['scheme'] === $this->destination['scheme']) {
93 throw new \InvalidArgumentException("You cannot copy from" . " {$this->sourceMetadata['scheme']} to" . " {$this->destination['scheme']}.");
94 }
95 // Handle multipart-related options.
96 $this->concurrency = isset($options['concurrency']) ? $options['concurrency'] : MultipartUploader::DEFAULT_CONCURRENCY;
97 $this->mupThreshold = isset($options['mup_threshold']) ? $options['mup_threshold'] : 16777216;
98 if ($this->mupThreshold < MultipartUploader::PART_MIN_SIZE) {
99 throw new \InvalidArgumentException('mup_threshold must be >= 5MB');
100 }
101 // Handle "before" callback option.
102 if (isset($options['before'])) {
103 $this->before = $options['before'];
104 if (!\is_callable($this->before)) {
105 throw new \InvalidArgumentException('before must be a callable.');
106 }
107 }
108 // Handle "debug" option.
109 if (isset($options['debug'])) {
110 if ($options['debug'] === \true) {
111 $options['debug'] = \fopen('php://output', 'w');
112 }
113 if (\is_resource($options['debug'])) {
114 $this->addDebugToBefore($options['debug']);
115 }
116 }
117 // Handle "add_content_md5" option.
118 $this->addContentMD5 = isset($options['add_content_md5']) && $options['add_content_md5'] === \true;
119 }
120 /**
121 * Transfers the files.
122 *
123 * @return PromiseInterface
124 */
125 public function promise() : PromiseInterface
126 {
127 // If the promise has been created, just return it.
128 if (!$this->promise) {
129 // Create an upload/download promise for the transfer.
130 $this->promise = $this->sourceMetadata['scheme'] === 'file' ? $this->createUploadPromise() : $this->createDownloadPromise();
131 }
132 return $this->promise;
133 }
134 /**
135 * Transfers the files synchronously.
136 */
137 public function transfer()
138 {
139 $this->promise()->wait();
140 }
141 private function prepareTarget($targetPath)
142 {
143 $target = ['path' => $this->normalizePath($targetPath), 'scheme' => $this->determineScheme($targetPath)];
144 if ($target['scheme'] !== 's3' && $target['scheme'] !== 'file') {
145 throw new \InvalidArgumentException('Scheme must be "s3" or "file".');
146 }
147 return $target;
148 }
149 /**
150 * Creates an array that contains Bucket and Key by parsing the filename.
151 *
152 * @param string $path Path to parse.
153 *
154 * @return array
155 */
156 private function getS3Args($path)
157 {
158 $parts = \explode('/', \str_replace('s3://', '', $path), 2);
159 $args = ['Bucket' => $parts[0]];
160 if (isset($parts[1])) {
161 $args['Key'] = $parts[1];
162 }
163 return $args;
164 }
165 /**
166 * Parses the scheme from a filename.
167 *
168 * @param string $path Path to parse.
169 *
170 * @return string
171 */
172 private function determineScheme($path)
173 {
174 return !\strpos($path, '://') ? 'file' : \explode('://', $path)[0];
175 }
176 /**
177 * Normalize a path so that it has UNIX-style directory separators and no trailing /
178 *
179 * @param string $path
180 *
181 * @return string
182 */
183 private function normalizePath($path)
184 {
185 return \rtrim(\str_replace('\\', '/', $path), '/');
186 }
187 private function resolvesOutsideTargetDirectory($sink, $objectKey)
188 {
189 $resolved = [];
190 $sections = \explode('/', $sink);
191 $targetSectionsLength = \count(\explode('/', $objectKey));
192 $targetSections = \array_slice($sections, -($targetSectionsLength + 1));
193 $targetDirectory = $targetSections[0];
194 foreach ($targetSections as $section) {
195 if ($section === '.' || $section === '') {
196 continue;
197 }
198 if ($section === '..') {
199 \array_pop($resolved);
200 if (empty($resolved) || $resolved[0] !== $targetDirectory) {
201 return \true;
202 }
203 } else {
204 $resolved[] = $section;
205 }
206 }
207 return \false;
208 }
209 private function createDownloadPromise()
210 {
211 $parts = $this->getS3Args($this->sourceMetadata['path']);
212 $prefix = "s3://{$parts['Bucket']}/" . (isset($parts['Key']) ? $parts['Key'] . '/' : '');
213 $commands = [];
214 foreach ($this->getDownloadsIterator() as $object) {
215 // Prepare the sink.
216 $objectKey = \preg_replace('/^' . \preg_quote($prefix, '/') . '/', '', $object);
217 $sink = $this->destination['path'] . '/' . $objectKey;
218 $command = $this->client->getCommand('GetObject', $this->getS3Args($object) + ['@http' => ['sink' => $sink]]);
219 if ($this->resolvesOutsideTargetDirectory($sink, $objectKey)) {
220 throw new AwsException('Cannot download key ' . $objectKey . ', its relative path resolves outside the' . ' parent directory', $command);
221 }
222 // Create the directory if needed.
223 $dir = \dirname($sink);
224 if (!\is_dir($dir) && !\mkdir($dir, 0777, \true)) {
225 throw new \RuntimeException("Could not create dir: {$dir}");
226 }
227 // Create the command.
228 $commands[] = $command;
229 }
230 // Create a GetObject command pool and return the promise.
231 return (new Aws\CommandPool($this->client, $commands, ['concurrency' => $this->concurrency, 'before' => $this->before, 'rejected' => function ($reason, $idx, Promise\PromiseInterface $p) {
232 $p->reject($reason);
233 }]))->promise();
234 }
235 private function createUploadPromise()
236 {
237 // Map each file into a promise that performs the actual transfer.
238 $files = \Dudlewebs\WPMCS\s3\Aws\map($this->getUploadsIterator(), function ($file) {
239 return \filesize($file) >= $this->mupThreshold ? $this->uploadMultipart($file) : $this->upload($file);
240 });
241 // Create an EachPromise, that will concurrently handle the upload
242 // operations' yielded promises from the iterator.
243 return Promise\Each::ofLimitAll($files, $this->concurrency);
244 }
245 /** @return Iterator */
246 private function getUploadsIterator()
247 {
248 if (\is_string($this->source)) {
249 return Aws\filter(Aws\recursive_dir_iterator($this->sourceMetadata['path']), function ($file) {
250 return !\is_dir($file);
251 });
252 }
253 return $this->source;
254 }
255 /** @return Iterator */
256 private function getDownloadsIterator()
257 {
258 if (\is_string($this->source)) {
259 $listArgs = $this->getS3Args($this->sourceMetadata['path']);
260 if (isset($listArgs['Key'])) {
261 $listArgs['Prefix'] = $listArgs['Key'] . '/';
262 unset($listArgs['Key']);
263 }
264 $files = $this->client->getPaginator('ListObjects', $listArgs)->search('Contents[].Key');
265 $files = Aws\map($files, function ($key) use($listArgs) {
266 return "s3://{$listArgs['Bucket']}/{$key}";
267 });
268 return Aws\filter($files, function ($key) {
269 return \substr($key, -1, 1) !== '/';
270 });
271 }
272 return $this->source;
273 }
274 private function upload($filename)
275 {
276 $args = $this->s3Args;
277 $args['SourceFile'] = $filename;
278 $args['Key'] = $this->createS3Key($filename);
279 $args['AddContentMD5'] = $this->addContentMD5;
280 $command = $this->client->getCommand('PutObject', $args);
281 $this->before and \call_user_func($this->before, $command);
282 return $this->client->executeAsync($command);
283 }
284 private function uploadMultipart($filename)
285 {
286 $args = $this->s3Args;
287 $args['Key'] = $this->createS3Key($filename);
288 $filename = $filename instanceof \SplFileInfo ? $filename->getPathname() : $filename;
289 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();
290 }
291 private function createS3Key($filename)
292 {
293 $filename = $this->normalizePath($filename);
294 $relative_file_path = \ltrim(\preg_replace('#^' . \preg_quote($this->sourceMetadata['path']) . '#', '', $filename), '/\\');
295 if (isset($this->s3Args['Key'])) {
296 return \rtrim($this->s3Args['Key'], '/') . '/' . $relative_file_path;
297 }
298 return $relative_file_path;
299 }
300 private function addDebugToBefore($debug)
301 {
302 $before = $this->before;
303 $sourcePath = $this->sourceMetadata['path'];
304 $s3Args = $this->s3Args;
305 $this->before = static function (CommandInterface $command) use($before, $debug, $sourcePath, $s3Args) {
306 // Call the composed before function.
307 $before and $before($command);
308 // Determine the source and dest values based on operation.
309 switch ($operation = $command->getName()) {
310 case 'GetObject':
311 $source = "s3://{$command['Bucket']}/{$command['Key']}";
312 $dest = $command['@http']['sink'];
313 break;
314 case 'PutObject':
315 $source = $command['SourceFile'];
316 $dest = "s3://{$command['Bucket']}/{$command['Key']}";
317 break;
318 case 'UploadPart':
319 $part = $command['PartNumber'];
320 case 'CreateMultipartUpload':
321 case 'CompleteMultipartUpload':
322 $sourceKey = $command['Key'];
323 if (isset($s3Args['Key']) && \strpos($sourceKey, $s3Args['Key']) === 0) {
324 $sourceKey = \substr($sourceKey, \strlen($s3Args['Key']) + 1);
325 }
326 $source = "{$sourcePath}/{$sourceKey}";
327 $dest = "s3://{$command['Bucket']}/{$command['Key']}";
328 break;
329 default:
330 throw new \UnexpectedValueException("Transfer encountered an unexpected operation: {$operation}.");
331 }
332 // Print the debugging message.
333 $context = \sprintf('%s -> %s (%s)', $source, $dest, $operation);
334 if (isset($part)) {
335 $context .= " : Part={$part}";
336 }
337 \fwrite($debug, "Transferring {$context}\n");
338 };
339 }
340 }
341