Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
make seekable s3 stream generic
Signed-off-by: Robin Appelman <[email protected]>
  • Loading branch information
icewind1991 committed Apr 1, 2020
commit 7b07e7251c8a92e95da922f34dde158ddffbeeee
36 changes: 25 additions & 11 deletions lib/private/Files/ObjectStore/S3ObjectTrait.php
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
use Aws\S3\ObjectUploader;
use Aws\S3\S3Client;
use Icewind\Streams\CallbackWrapper;
use OC\Files\Stream\SeekableHttpStream;

const S3_UPLOAD_PART_SIZE = 524288000; // 500MB

Expand All @@ -49,16 +50,29 @@ abstract protected function getConnection();
* @since 7.0.0
*/
function readObject($urn) {
$context = stream_context_create([
's3seek' => [
'client' => $this->getConnection(),
'bucket' => $this->bucket,
'urn' => $urn,
],
]);
return SeekableHttpStream::open(function ($range) use ($urn) {
$command = $this->getConnection()->getCommand('GetObject', [
'Bucket' => $this->bucket,
'Key' => $urn,
'Range' => 'bytes=' . $range,
]);
$request = \Aws\serialize($command);
$headers = [];
foreach ($request->getHeaders() as $key => $values) {
foreach ($values as $value) {
$headers[] = "$key: $value";
}
}
$opts = [
'http' => [
'protocol_version' => 1.1,
'header' => $headers,
],
];

S3SeekableReadStream::registerIfNeeded();
return fopen('s3seek://', 'r', false, $context);
$context = stream_context_create($opts);
return fopen($request->getUri(), 'r', false, $context);
});
}

/**
Expand All @@ -76,7 +90,7 @@ function writeObject($urn, $stream) {
$uploader = new MultipartUploader($this->getConnection(), $countStream, [
'bucket' => $this->bucket,
'key' => $urn,
'part_size' => S3_UPLOAD_PART_SIZE
'part_size' => S3_UPLOAD_PART_SIZE,
]);

try {
Expand All @@ -103,7 +117,7 @@ function writeObject($urn, $stream) {
function deleteObject($urn) {
$this->getConnection()->deleteObject([
'Bucket' => $this->bucket,
'Key' => $urn
'Key' => $urn,
]);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,64 +20,72 @@
*
*/

namespace OC\Files\ObjectStore;
namespace OC\Files\Stream;

use Icewind\Streams\File;

/**
* A stream wrapper that uses http range requests to provide a seekable
* stream of a file in S3 storage.
* A stream wrapper that uses http range requests to provide a seekable stream for http reading
*/
class S3SeekableReadStream {
class SeekableHttpStream implements File {
private const PROTOCOL = 'httpseek';

private static $registered = false;

/**
* Registers the stream wrapper using the `s3seek://` url scheme
* Registers the stream wrapper using the `httpseek://` url scheme
* $return void
*/
public static function registerIfNeeded() {
private static function registerIfNeeded() {
if (!self::$registered) {
stream_wrapper_register(
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

stream_wrapper_register() will return FALSE if the protocol already has a handler.

https://www.php.net/manual/en/function.stream-wrapper-register.php

Do we really need a state (registered y/no)? I guess stream_wrapper_register will not register more than one handler for one protocol.

's3seek',
'OC\Files\ObjectStore\S3SeekableReadStream'
self::PROTOCOL,
self::class
);
self::$registered = true;
}
}

private $client;
private $bucket;
private $urn;
/**
* Open a readonly-seekable http stream
*
* The provided callback will be called with byte range and should return an http stream for the requested range
*
* @param callable $callback
* @return false|resource
*/
public static function open(callable $callback) {
$context = stream_context_create([
SeekableHttpStream::PROTOCOL => [
'callback' => $callback
],
]);

SeekableHttpStream::registerIfNeeded();
return fopen(SeekableHttpStream::PROTOCOL . '://', 'r', false, $context);
}

/** @var resource */
public $context;

/** @var callable */
private $openCallback;

/** @var resource */
private $current;
/** @var int */
private $offset = 0;

private function reconnect($range) {
if ($this->current != null) {
fclose($this->current);
}

$command = $this->client->getCommand('GetObject', [
'Bucket' => $this->bucket,
'Key' => $this->urn,
'Range' => 'bytes=' . $range,
]);
$request = \Aws\serialize($command);
$headers = [];
foreach ($request->getHeaders() as $key => $values) {
foreach ($values as $value) {
$headers[] = "$key: $value";
}
}
$opts = [
'http' => [
'protocol_version' => 1.1,
'header' => $headers,
],
];

$context = stream_context_create($opts);
$this->current = fopen($request->getUri(), 'r', false, $context);
$this->current = ($this->openCallback)($range);

if ($this->current === false) {return false;}
if ($this->current === false) {
return false;
}

$responseHead = stream_get_meta_data($this->current)['wrapper_data'];
$contentRange = array_values(array_filter($responseHead, function ($v) {
Expand All @@ -93,30 +101,35 @@ private function reconnect($range) {
}

function stream_open($path, $mode, $options, &$opened_path) {
$o = stream_context_get_options($this->context)['s3seek'];
$this->bucket = $o['bucket'];
$this->urn = $o['urn'];
$this->client = $o['client'];
$options = stream_context_get_options($this->context)[self::PROTOCOL];
$this->openCallback = $options['callback'];

return $this->reconnect('0-');
}

function stream_read($count) {
if (!$this->current) {
return false;
}
$ret = fread($this->current, $count);
$this->offset += strlen($ret);
return $ret;
}

function stream_seek($offset, $whence) {
function stream_seek($offset, $whence = SEEK_SET) {
switch ($whence) {
case SEEK_SET:
if ($offset === $this->offset) {return true;}
return $this->reconnect($offset . '-');
case SEEK_CUR:
if ($offset === 0) {return true;}
return $this->reconnect(($this->offset + $offset) . '-');
case SEEK_END:
return false;
case SEEK_SET:
if ($offset === $this->offset) {
return true;
}
return $this->reconnect($offset . '-');
case SEEK_CUR:
if ($offset === 0) {
return true;
}
return $this->reconnect(($this->offset + $offset) . '-');
case SEEK_END:
return false;
}
return false;
}
Expand All @@ -136,4 +149,26 @@ function stream_eof() {
function stream_close() {
fclose($this->current);
}

public function stream_write($data) {
return false;
}

public function stream_set_option($option, $arg1, $arg2) {
return false;
}

public function stream_truncate($size) {
return false;
}

public function stream_lock($operation) {
return false;
}

public function stream_flush() {
return; //noop because readonly stream
}


}