Skip to content
Merged
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
15 changes: 14 additions & 1 deletion google/cloud/storage/internal/async/open_object.cc
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,9 @@
// limitations under the License.

#include "google/cloud/storage/internal/async/open_object.h"
#include "google/cloud/storage/internal/grpc/scale_stall_timeout.h"
#include "google/cloud/storage/options.h"
#include "google/cloud/internal/async_read_write_stream_timeout.h"
#include "google/cloud/internal/make_status.h"
#include "absl/strings/str_cat.h"
#include <utility>
Expand Down Expand Up @@ -59,7 +62,17 @@ std::unique_ptr<OpenStream::StreamingRpc> OpenObject::CreateRpc(
google::storage::v2::BidiReadObjectRequest const& request) {
auto p = RequestParams(request);
if (!p.empty()) context->AddMetadata("x-goog-request-params", std::move(p));
return stub.AsyncBidiReadObject(cq, std::move(context), std::move(options));
std::chrono::milliseconds const timeout = ScaleStallTimeout(
options->get<storage::DownloadStallTimeoutOption>(),
options->get<storage::DownloadStallMinimumRateOption>(),
google::storage::v2::ServiceConstants::MAX_READ_CHUNK_BYTES);
std::unique_ptr<OpenStream::StreamingRpc> rpc =
stub.AsyncBidiReadObject(cq, std::move(context), std::move(options));
return std::make_unique<
google::cloud::internal::AsyncStreamingReadWriteRpcTimeout<
Comment thread
rajeevpodar marked this conversation as resolved.
google::storage::v2::BidiReadObjectRequest,
google::storage::v2::BidiReadObjectResponse>>(
cq, timeout, timeout, timeout, std::move(rpc));
}

void OpenObject::OnStart(bool ok) {
Expand Down
82 changes: 82 additions & 0 deletions google/cloud/storage/internal/async/open_object_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -14,14 +14,17 @@

#include "google/cloud/storage/internal/async/open_object.h"
#include "google/cloud/mocks/mock_async_streaming_read_write_rpc.h"
#include "google/cloud/storage/options.h"
#include "google/cloud/storage/testing/canonical_errors.h"
#include "google/cloud/storage/testing/mock_storage_stub.h"
#include "google/cloud/testing_util/async_sequencer.h"
#include "google/cloud/testing_util/is_proto_equal.h"
#include "google/cloud/testing_util/mock_completion_queue_impl.h"
#include "google/cloud/testing_util/status_matchers.h"
#include "google/cloud/testing_util/validate_metadata.h"
#include <google/protobuf/text_format.h>
#include <gmock/gmock.h>
#include <chrono>
#include <memory>

namespace google {
Expand All @@ -35,6 +38,7 @@ using ::google::cloud::storage::testing::canonical_errors::PermanentError;
using ::google::cloud::testing_util::AsyncSequencer;
using ::google::cloud::testing_util::IsOkAndHolds;
using ::google::cloud::testing_util::IsProtoEqual;
using ::google::cloud::testing_util::MockCompletionQueueImpl;
using ::google::cloud::testing_util::StatusIs;
using ::google::protobuf::TextFormat;
using ::testing::AllOf;
Expand All @@ -47,6 +51,16 @@ using MockStream = google::cloud::mocks::MockAsyncStreamingReadWriteRpc<
google::storage::v2::BidiReadObjectRequest,
google::storage::v2::BidiReadObjectResponse>;

StatusOr<std::chrono::system_clock::time_point> CancelledTimer() {
return internal::CancelledError("test-only", GCP_ERROR_INFO());
}

StatusOr<std::chrono::system_clock::time_point> MakeTimerStatus(
future<bool> f) {
if (!f.get()) return CancelledTimer();
return make_status_or(std::chrono::system_clock::now());
}

TEST(OpenImpl, RequestParams) {
auto constexpr kPlain = R"pb(
read_object_spec {
Expand Down Expand Up @@ -390,6 +404,74 @@ TEST(OpenImpl, UnexpectedFinish) {
EXPECT_THAT(response, StatusIs(StatusCode::kInternal));
}

TEST(OpenImpl, TimeoutCancellation) {
AsyncSequencer<bool> sequencer;
MockStorageStub mock;
EXPECT_CALL(mock, AsyncBidiReadObject).WillOnce([&sequencer]() {
auto stream = std::make_unique<MockStream>();
EXPECT_CALL(*stream, Start).WillOnce([&sequencer]() {
return sequencer.PushBack("Start").then([](auto f) { return f.get(); });
});
EXPECT_CALL(*stream, Write).WillOnce([&sequencer]() {
return sequencer.PushBack("Write").then([](auto f) { return f.get(); });
});
EXPECT_CALL(*stream, Read).WillOnce([&sequencer]() {
return sequencer.PushBack("Read").then([](auto) {
return std::optional<google::storage::v2::BidiReadObjectResponse>();
});
});
EXPECT_CALL(*stream, Cancel).Times(testing::AtLeast(1));
EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() {
return sequencer.PushBack("Finish").then([](auto) {
return Status(StatusCode::kCancelled, "Stream timeout");
});
});
return std::unique_ptr<OpenStream::StreamingRpc>(std::move(stream));
});

auto mock_cq = std::make_shared<MockCompletionQueueImpl>();
EXPECT_CALL(*mock_cq, MakeRelativeTimer).WillRepeatedly([&sequencer](auto) {
return sequencer.PushBack("MakeRelativeTimer").then(MakeTimerStatus);
});

CompletionQueue cq(mock_cq);
Options options;
options.set<storage::DownloadStallTimeoutOption>(std::chrono::seconds(1));
auto coro = std::make_shared<OpenObject>(
mock, cq, std::make_shared<grpc::ClientContext>(),
internal::MakeImmutableOptions(std::move(options)),
google::storage::v2::BidiReadObjectRequest{});
future<StatusOr<OpenStreamResult>> pending = coro->Call();

auto timer1 = sequencer.PopFrontWithName();
EXPECT_EQ(timer1.second, "MakeRelativeTimer");
auto start = sequencer.PopFrontWithName();
EXPECT_EQ(start.second, "Start");
start.first.set_value(true);
timer1.first.set_value(false);

auto timer2 = sequencer.PopFrontWithName();
EXPECT_EQ(timer2.second, "MakeRelativeTimer");
auto write = sequencer.PopFrontWithName();
EXPECT_EQ(write.second, "Write");
write.first.set_value(true);
timer2.first.set_value(false);

auto timer3 = sequencer.PopFrontWithName();
EXPECT_EQ(timer3.second, "MakeRelativeTimer");
auto read = sequencer.PopFrontWithName();
EXPECT_EQ(read.second, "Read");
timer3.first.set_value(true);
read.first.set_value(true);

auto finish = sequencer.PopFrontWithName();
EXPECT_EQ(finish.second, "Finish");
finish.first.set_value(true);

StatusOr<OpenStreamResult> response = pending.get();
EXPECT_THAT(response, StatusIs(StatusCode::kCancelled));
}

} // namespace
GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END
} // namespace storage_internal
Expand Down
Loading