Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
22 changes: 22 additions & 0 deletions google/cloud/storage/async/client.cc
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
#include "google/cloud/storage/internal/async/connection_impl.h"
#include "google/cloud/storage/internal/async/connection_tracing.h"
#include "google/cloud/storage/internal/async/default_options.h"
#include "google/cloud/storage/internal/async/options.h"
#include "google/cloud/storage/internal/grpc/stub.h"
#include "google/cloud/grpc_options.h"
#include <memory>
Expand Down Expand Up @@ -84,6 +85,27 @@ future<StatusOr<ObjectDescriptor>> AsyncClient::Open(
});
}

future<StatusOr<ObjectDescriptor>> AsyncClient::Open(
BucketName const& bucket_name, std::string object_name,
InitialReadRanges const& config, Options opts) {
auto spec = google::storage::v2::BidiReadObjectSpec{};
spec.set_bucket(bucket_name.FullName());
spec.set_object(std::move(object_name));

// Convert the user-facing `InitialReadRanges` to the internal
// `ReadRangesOption` so it can be propagated down to the connection
// implementation.
if (!config.initial_ranges.empty()) {
std::vector<storage_internal::ReadRangeConfig> internal_ranges;
internal_ranges.reserve(config.initial_ranges.size());
for (auto const& r : config.initial_ranges) {
internal_ranges.push_back({r.offset, r.length});
}
opts.set<storage_internal::ReadRangesOption>(std::move(internal_ranges));
}
return Open(std::move(spec), std::move(opts));
}

future<StatusOr<std::pair<AsyncReader, AsyncToken>>> AsyncClient::ReadObject(
BucketName const& bucket_name, std::string object_name, Options opts) {
auto request = google::storage::v2::ReadObjectRequest{};
Expand Down
44 changes: 44 additions & 0 deletions google/cloud/storage/async/client.h
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,27 @@ GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN
*/
class AsyncClient {
public:
/**
* Specifies a byte range for a read request.
*/
struct ByteRange {
std::int64_t offset = 0;
std::int64_t length = 0;
};

/**
* Specifies initial byte ranges to request concurrently when opening an
* object.
*
* Passing initial read ranges allows the client to begin fetching expected
* byte ranges during connection setup, which may improve first-byte retrieval
* times for known access patterns.
*/
struct InitialReadRanges {
// The initial read ranges.
std::vector<ByteRange> initial_ranges;
};

/// Create a new client configured with @p options.
explicit AsyncClient(Options options = {});
/// Create a new client using @p connection. This is often used for mocking.
Expand Down Expand Up @@ -273,6 +294,29 @@ class AsyncClient {
std::string object_name,
Options opts = {});

/**
* Open an object descriptor, requesting specified initial read ranges
* concurrently.
*
* @par Example
* @snippet storage_async_samples.cc open-object-initial-read-ranges
*
* @par Idempotency
* This is a read-only operation and is always idempotent. The operation will
* retry until the descriptor is successfully created. The descriptor itself
* will resume any incomplete ranged reads if the connection(s) are
* interrupted. Use `ResumePolicyOption` and `ResumePolicy` to control this.
*
* @param bucket_name the name of the bucket that contains the object.
* @param object_name the name of the object to be read.
* @param config initial byte ranges to request during connection setup.
* @param opts options controlling the behavior of this RPC.
*/
future<StatusOr<ObjectDescriptor>> Open(BucketName const& bucket_name,
std::string object_name,
InitialReadRanges const& config,
Options opts = {});

/**
* Open an object descriptor to perform one or more ranged reads.
*
Expand Down
38 changes: 38 additions & 0 deletions google/cloud/storage/async/client_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
// limitations under the License.

#include "google/cloud/storage/async/client.h"
#include "google/cloud/storage/internal/async/options.h"
#include "google/cloud/storage/mocks/mock_async_connection.h"
#include "google/cloud/storage/mocks/mock_async_object_descriptor_connection.h"
#include "google/cloud/storage/mocks/mock_async_reader_connection.h"
Expand Down Expand Up @@ -290,6 +291,43 @@ TEST(AsyncClient, Open) {
"empty response", [](auto const& p) { return p.size(); }, 0)));
}

TEST(AsyncClient, OpenWithInitialReadRanges) {
auto constexpr kExpectedRequest = R"pb(
bucket: "projects/_/buckets/test-bucket"
object: "test-object"
)pb";
auto mock = std::make_shared<MockAsyncConnection>();
EXPECT_CALL(*mock, options).WillRepeatedly(Return(Options{}));

EXPECT_CALL(*mock, Open).WillOnce([&](AsyncConnection::OpenParams const& p) {
EXPECT_TRUE(p.options.has<storage_internal::ReadRangesOption>());
auto const& ranges = p.options.get<storage_internal::ReadRangesOption>();
EXPECT_EQ(ranges.size(), 2);
if (ranges.size() >= 2) {
EXPECT_EQ(ranges[0].offset, 0);
EXPECT_EQ(ranges[0].length, 100);
EXPECT_EQ(ranges[1].offset, 1000);
EXPECT_EQ(ranges[1].length, 200);
}

auto expected = google::storage::v2::BidiReadObjectSpec{};
EXPECT_TRUE(TextFormat::ParseFromString(kExpectedRequest, &expected));
EXPECT_THAT(p.read_spec, IsProtoEqual(expected));

auto descriptor = std::make_shared<MockAsyncObjectDescriptorConnection>();
return make_ready_future(make_status_or(
std::shared_ptr<ObjectDescriptorConnection>(std::move(descriptor))));
});

auto client = AsyncClient(mock);
AsyncClient::InitialReadRanges config;
config.initial_ranges = {{0, 100}, {1000, 200}};
auto descriptor =
client.Open(BucketName("test-bucket"), "test-object", std::move(config))
.get();
ASSERT_STATUS_OK(descriptor);
}

TEST(AsyncClient, OpenWithInvalidBucket) {
auto constexpr kExpectedRequest = R"pb(
bucket: "test-only-invalid"
Expand Down
60 changes: 60 additions & 0 deletions google/cloud/storage/examples/storage_async_samples.cc
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,55 @@ void OpenObjectMultipleRangedRead(google::cloud::storage::AsyncClient& client,
std::cout << "The ranges contain " << count << " newlines\n";
}

void OpenObjectWithInitialReadRanges(
google::cloud::storage::AsyncClient& client,
std::vector<std::string> const& argv) {
//! [open-object-initial-read-ranges]
// [START storage_open_object_initial_read_ranges]
namespace gcs = google::cloud::storage;

// Helper coroutine to count newlines returned by an AsyncReader.
auto count_newlines =
[](gcs::AsyncReader reader,
gcs::AsyncToken token) -> google::cloud::future<std::uint64_t> {
std::uint64_t count = 0;
while (token.valid()) {
auto [payload, t] = (co_await reader.Read(std::move(token))).value();
token = std::move(t);
for (auto const& buffer : payload.contents()) {
count += std::count(buffer.begin(), buffer.end(), '\n');
}
}
co_return count;
};

auto coro =
[&count_newlines](
gcs::AsyncClient& client, std::string bucket_name,
std::string object_name) -> google::cloud::future<std::uint64_t> {
gcs::AsyncClient::InitialReadRanges config;
config.initial_ranges = {{0, 1024}, {1024, 1024}};

auto descriptor =
(co_await client.Open(gcs::BucketName(std::move(bucket_name)),
std::move(object_name), std::move(config)))
.value();

auto [r1, t1] = descriptor.Read(0, 1024);
auto [r2, t2] = descriptor.Read(1024, 1024);

auto c1 = count_newlines(std::move(r1), std::move(t1));
auto c2 = count_newlines(std::move(r2), std::move(t2));
co_return (co_await std::move(c1)) + (co_await std::move(c2));
};
// [END storage_open_object_initial_read_ranges]
//! [open-object-initial-read-ranges]
// The example is easier to test and run if we call the coroutine and block
// until it completes.
auto const count = coro(client, argv.at(0), argv.at(1)).get();
std::cout << "The pre-warmed ranges contain " << count << " newlines\n";
}

void OpenObjectReadFullObject(google::cloud::storage::AsyncClient& client,
std::vector<std::string> const& argv) {
//! [open-object-read-full-object]
Expand Down Expand Up @@ -997,6 +1046,11 @@ void OpenObjectMultipleRangedRead(google::cloud::storage::AsyncClient&,
std::cerr << "AsyncClient::Open() example requires coroutines\n";
}

void OpenObjectWithInitialReadRanges(google::cloud::storage::AsyncClient&,
std::vector<std::string> const&) {
std::cerr << "AsyncClient::Open() example requires coroutines\n";
}

void OpenMultipleObjectsRangedRead(google::cloud::storage::AsyncClient&,
std::vector<std::string> const&) {
std::cerr << "AsyncClient::Open() example requires coroutines\n";
Expand Down Expand Up @@ -1306,6 +1360,10 @@ void AutoRun(std::vector<std::string> const& argv) {
<< std::endl;
OpenObjectMultipleRangedRead(client, {bucket_name, composed_name});

std::cout << "Running the OpenObjectWithInitialReadRanges() example"
<< std::endl;
OpenObjectWithInitialReadRanges(client, {bucket_name, composed_name});

std::cout << "Running the OpenMultipleObjectsRangedRead() example"
<< std::endl;
auto const multi_read_o1 =
Expand Down Expand Up @@ -1563,6 +1621,8 @@ int main(int argc, char* argv[]) try {
OpenObjectSingleRangedRead),
make_entry("open-object-multiple-ranged-read", {},
OpenObjectMultipleRangedRead),
make_entry("open-object-initial-read-ranges", {},
OpenObjectWithInitialReadRanges),
make_entry("open-object-read-full-object", {}, OpenObjectReadFullObject),
make_entry("open-multiple-objects-ranged-read",
{"<object-name-1>", "<object-name-2>", "<object-name-3>"},
Expand Down
1 change: 1 addition & 0 deletions google/cloud/storage/google_cloud_cpp_storage_grpc.bzl
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ google_cloud_cpp_storage_grpc_hdrs = [
"internal/async/open_object.h",
"internal/async/open_object_telemetry.h",
"internal/async/open_stream.h",
"internal/async/options.h",
"internal/async/partial_upload.h",
"internal/async/read_payload_fwd.h",
"internal/async/read_payload_impl.h",
Expand Down
1 change: 1 addition & 0 deletions google/cloud/storage/google_cloud_cpp_storage_grpc.cmake
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,7 @@ add_library(
internal/async/open_object_telemetry.h
internal/async/open_stream.cc
internal/async/open_stream.h
internal/async/options.h
internal/async/partial_upload.cc
internal/async/partial_upload.h
internal/async/read_payload_fwd.h
Expand Down
15 changes: 15 additions & 0 deletions google/cloud/storage/internal/async/connection_impl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,9 @@
#include "google/cloud/storage/internal/async/object_descriptor_impl.h"
#include "google/cloud/storage/internal/async/open_object.h"
#include "google/cloud/storage/internal/async/open_stream.h"
#include "google/cloud/storage/internal/async/options.h"
#include "google/cloud/storage/internal/async/read_payload_impl.h"
#include "google/cloud/storage/internal/async/read_range.h"
#include "google/cloud/storage/internal/async/reader_connection_impl.h"
#include "google/cloud/storage/internal/async/reader_connection_resume.h"
#include "google/cloud/storage/internal/async/rewriter_connection_impl.h"
Expand Down Expand Up @@ -58,6 +60,7 @@
#include "google/cloud/internal/make_status.h"
#include <grpcpp/grpcpp.h>
#include <memory>
#include <set>
#include <utility>

namespace google {
Expand Down Expand Up @@ -220,6 +223,18 @@ AsyncConnectionImpl::Open(OpenParams p) {
auto initial_request = google::storage::v2::BidiReadObjectRequest{};
*initial_request.mutable_read_object_spec() = p.read_spec;
auto current = internal::MakeImmutableOptions(p.options);
// If pre-warmed ranges are configured, populate the initial request
// with these ranges to start downloading them as soon as the stream opens.
if (current->has<ReadRangesOption>()) {
auto const& ranges = current->get<ReadRangesOption>();

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

medium

When a function returns a reference to a container or a large object, use auto const& instead of auto to avoid unnecessary and inefficient copies.

    auto const& ranges =
        current->get<ReadRangesOption>();
References
  1. When a function returns a reference to a container or a large object (such as std::multimap), use 'auto const&' instead of 'auto' to avoid unnecessary and inefficient copies.

for (auto const& r : DeduplicateRanges(ranges)) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

medium

Avoid using auto when it hides domain objects. Explicitly declare the loop variable type.

    for (DedupedReadRange const& r : DeduplicateRanges(ranges)) {
References
  1. Reject auto when it hides domain objects or function return types. (link)

auto* proto_range = initial_request.add_read_ranges();

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

medium

Avoid using auto when it hides protobuf messages or fields. Use the explicit protobuf type instead.

      google::storage::v2::ReadRange* proto_range =
          initial_request.add_read_ranges();
References
  1. Reject auto when it hides protobuf messages/fields. (link)

proto_range->set_read_offset(r.config.offset);
proto_range->set_read_length(r.config.length);
// Generate sequential IDs starting at 1. The receiver must match these.
proto_range->set_read_id(r.read_id);
}
}
// Get the policy factory and immediately create a policy.
auto resume_policy = current->get<storage::ResumePolicyOption>()();

Expand Down
Loading
Loading