mirror of
https://github.com/google/cdc-file-transfer.git
synced 2026-09-13 09:32:30 +03:00
[cdc_stream] Add a CLI client to start/stop asset streaming sessions (#4)
Implements the cdc_stream client and adjusts asset streaming in various places to work better outside of a GGP environment. This CL tries to get quoting for SSH commands right. It also brings back the ability to start a streaming session from asset_stream_manager. Also cleans up Bazel targets setup. Since the sln file is now in root, it is no longer necessary to prepend ../ to relative filenames to make clicking on errors work.
This commit is contained in:
@@ -17,6 +17,7 @@
|
||||
#include <iomanip>
|
||||
|
||||
#include "absl/strings/str_format.h"
|
||||
#include "absl/strings/str_replace.h"
|
||||
#include "absl/strings/str_split.h"
|
||||
#include "asset_stream_manager/multi_session.h"
|
||||
#include "asset_stream_manager/session_manager.h"
|
||||
@@ -26,11 +27,21 @@
|
||||
#include "common/process.h"
|
||||
#include "common/sdk_util.h"
|
||||
#include "common/status.h"
|
||||
#include "google/protobuf/text_format.h"
|
||||
#include "manifest/manifest_updater.h"
|
||||
|
||||
using TextFormat = google::protobuf::TextFormat;
|
||||
|
||||
namespace cdc_ft {
|
||||
namespace {
|
||||
|
||||
std::string RequestToString(const google::protobuf::Message& request) {
|
||||
std::string str;
|
||||
google::protobuf::TextFormat::PrintToString(request, &str);
|
||||
if (!str.empty() && str.back() == '\n') str.pop_back();
|
||||
return absl::StrReplaceAll(str, {{"\n", ", "}});
|
||||
}
|
||||
|
||||
// Parses |instance_name| of the form
|
||||
// "organizations/{org-id}/projects/{proj-id}/pools/{pool-id}/gamelets/{gamelet-id}"
|
||||
// into parts. The pool id is not returned.
|
||||
@@ -102,48 +113,18 @@ LocalAssetsStreamManagerServiceImpl::~LocalAssetsStreamManagerServiceImpl() =
|
||||
grpc::Status LocalAssetsStreamManagerServiceImpl::StartSession(
|
||||
grpc::ServerContext* /*context*/, const StartSessionRequest* request,
|
||||
StartSessionResponse* /*response*/) {
|
||||
LOG_INFO("RPC:StartSession(gamelet_name='%s', workstation_directory='%s'",
|
||||
request->gamelet_name(), request->workstation_directory());
|
||||
LOG_INFO("RPC:StartSession(%s)", RequestToString(*request));
|
||||
|
||||
metrics::DeveloperLogEvent evt;
|
||||
evt.as_manager_data = std::make_unique<metrics::AssetStreamingManagerData>();
|
||||
evt.as_manager_data->session_start_data =
|
||||
std::make_unique<metrics::SessionStartData>();
|
||||
evt.as_manager_data->session_start_data->absl_status = absl::StatusCode::kOk;
|
||||
evt.as_manager_data->session_start_data->status =
|
||||
metrics::SessionStartStatus::kOk;
|
||||
evt.as_manager_data->session_start_data->origin =
|
||||
ConvertOrigin(request->origin());
|
||||
|
||||
// Parse instance/project/org id.
|
||||
absl::Status status;
|
||||
MultiSession* ms = nullptr;
|
||||
std::string instance_id, project_id, organization_id, instance_ip;
|
||||
uint16_t instance_port = 0;
|
||||
if (!ParseInstanceName(request->gamelet_name(), &instance_id, &project_id,
|
||||
&organization_id)) {
|
||||
status = absl::InvalidArgumentError(absl::StrFormat(
|
||||
"Failed to parse instance name '%s'", request->gamelet_name()));
|
||||
} else {
|
||||
evt.project_id = project_id;
|
||||
evt.organization_id = organization_id;
|
||||
|
||||
status = InitSsh(instance_id, project_id, organization_id, &instance_ip,
|
||||
&instance_port);
|
||||
|
||||
if (status.ok()) {
|
||||
status = session_manager_->StartSession(
|
||||
instance_id, project_id, organization_id, instance_ip, instance_port,
|
||||
request->workstation_directory(), &ms,
|
||||
&evt.as_manager_data->session_start_data->status);
|
||||
}
|
||||
}
|
||||
metrics::DeveloperLogEvent evt;
|
||||
std::string instance_id;
|
||||
absl::Status status = StartSessionInternal(request, &instance_id, &ms, &evt);
|
||||
|
||||
evt.as_manager_data->session_start_data->absl_status = status.code();
|
||||
if (ms) {
|
||||
evt.as_manager_data->session_start_data->concurrent_session_count =
|
||||
ms->GetSessionCount();
|
||||
if (!instance_id.empty() && ms->HasSessionForInstance(instance_id)) {
|
||||
if (!instance_id.empty() && ms->HasSession(instance_id)) {
|
||||
ms->RecordSessionEvent(std::move(evt), metrics::EventType::kSessionStart,
|
||||
instance_id);
|
||||
} else {
|
||||
@@ -166,9 +147,14 @@ grpc::Status LocalAssetsStreamManagerServiceImpl::StartSession(
|
||||
grpc::Status LocalAssetsStreamManagerServiceImpl::StopSession(
|
||||
grpc::ServerContext* /*context*/, const StopSessionRequest* request,
|
||||
StopSessionResponse* /*response*/) {
|
||||
LOG_INFO("RPC:StopSession(gamelet_id='%s')", request->gamelet_id());
|
||||
LOG_INFO("RPC:StopSession(%s)", RequestToString(*request));
|
||||
|
||||
absl::Status status = session_manager_->StopSession(request->gamelet_id());
|
||||
std::string instance_id =
|
||||
!request->gamelet_id().empty() // Stadia use case
|
||||
? request->gamelet_id()
|
||||
: absl::StrCat(request->user_host(), ":", request->mount_dir());
|
||||
|
||||
absl::Status status = session_manager_->StopSession(instance_id);
|
||||
if (status.ok()) {
|
||||
LOG_INFO("StopSession() succeeded");
|
||||
} else {
|
||||
@@ -177,6 +163,86 @@ grpc::Status LocalAssetsStreamManagerServiceImpl::StopSession(
|
||||
return ToGrpcStatus(status);
|
||||
}
|
||||
|
||||
absl::Status LocalAssetsStreamManagerServiceImpl::StartSessionInternal(
|
||||
const StartSessionRequest* request, std::string* instance_id,
|
||||
MultiSession** ms, metrics::DeveloperLogEvent* evt) {
|
||||
instance_id->clear();
|
||||
*ms = nullptr;
|
||||
evt->as_manager_data = std::make_unique<metrics::AssetStreamingManagerData>();
|
||||
evt->as_manager_data->session_start_data =
|
||||
std::make_unique<metrics::SessionStartData>();
|
||||
evt->as_manager_data->session_start_data->absl_status = absl::StatusCode::kOk;
|
||||
evt->as_manager_data->session_start_data->status =
|
||||
metrics::SessionStartStatus::kOk;
|
||||
evt->as_manager_data->session_start_data->origin =
|
||||
ConvertOrigin(request->origin());
|
||||
|
||||
if (!(request->gamelet_name().empty() ^ request->user_host().empty())) {
|
||||
return absl::InvalidArgumentError(
|
||||
"Must set either gamelet_name or user_host.");
|
||||
}
|
||||
|
||||
if (request->mount_dir().empty()) {
|
||||
return absl::InvalidArgumentError("mount_dir cannot be empty.");
|
||||
}
|
||||
|
||||
SessionTarget target;
|
||||
if (!request->gamelet_name().empty()) {
|
||||
ASSIGN_OR_RETURN(target,
|
||||
GetTargetForStadia(*request, instance_id, &evt->project_id,
|
||||
&evt->organization_id));
|
||||
} else {
|
||||
target = GetTarget(*request, instance_id);
|
||||
}
|
||||
|
||||
return session_manager_->StartSession(
|
||||
*instance_id, request->workstation_directory(), target, evt->project_id,
|
||||
evt->organization_id, ms,
|
||||
&evt->as_manager_data->session_start_data->status);
|
||||
}
|
||||
|
||||
absl::StatusOr<SessionTarget>
|
||||
LocalAssetsStreamManagerServiceImpl::GetTargetForStadia(
|
||||
const StartSessionRequest& request, std::string* instance_id,
|
||||
std::string* project_id, std::string* organization_id) {
|
||||
SessionTarget target;
|
||||
target.mount_dir = request.mount_dir();
|
||||
target.ssh_command = request.ssh_command();
|
||||
target.scp_command = request.scp_command();
|
||||
|
||||
// Parse instance/project/org id.
|
||||
if (!ParseInstanceName(request.gamelet_name(), instance_id, project_id,
|
||||
organization_id)) {
|
||||
return absl::InvalidArgumentError(absl::StrFormat(
|
||||
"Failed to parse instance name '%s'", request.gamelet_name()));
|
||||
}
|
||||
|
||||
// Run 'ggp ssh init' to determine IP (host) and port.
|
||||
std::string instance_ip;
|
||||
uint16_t instance_port = 0;
|
||||
RETURN_IF_ERROR(InitSsh(*instance_id, *project_id, *organization_id,
|
||||
&instance_ip, &instance_port));
|
||||
|
||||
target.user_host = "cloudcast@" + instance_ip;
|
||||
target.ssh_port = instance_port;
|
||||
return target;
|
||||
}
|
||||
|
||||
SessionTarget LocalAssetsStreamManagerServiceImpl::GetTarget(
|
||||
const StartSessionRequest& request, std::string* instance_id) {
|
||||
SessionTarget target;
|
||||
target.user_host = request.user_host();
|
||||
target.mount_dir = request.mount_dir();
|
||||
target.ssh_command = request.ssh_command();
|
||||
target.scp_command = request.scp_command();
|
||||
target.ssh_port = request.port() > 0 && request.port() <= UINT16_MAX
|
||||
? static_cast<uint16_t>(request.port())
|
||||
: RemoteUtil::kDefaultSshPort;
|
||||
|
||||
*instance_id = absl::StrCat(target.user_host, ":", target.mount_dir);
|
||||
return target;
|
||||
}
|
||||
|
||||
metrics::RequestOrigin LocalAssetsStreamManagerServiceImpl::ConvertOrigin(
|
||||
StartSessionRequestOrigin origin) const {
|
||||
switch (origin) {
|
||||
|
||||
Reference in New Issue
Block a user