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
16 changes: 8 additions & 8 deletions .kokoro/system.sh
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,14 @@ run_package_test() {
trap 'rm -rf "$gcloud_config_dir"' EXIT

case "${package_name}" in
"gapic-generator")
source "${KOKORO_GFILE_DIR}/gapic-generator-env.sh"
PROJECT_ID=$(cat "${KOKORO_GFILE_DIR}/gapic-generator-project-id.json")
GOOGLE_APPLICATION_CREDENTIALS="${KOKORO_GFILE_DIR}/gapic-generator-service-account.json"
export GOOGLE_ADS_VIDEO_PATH="${KOKORO_GFILE_DIR}/gapic-generator-video.mp4"
NOX_FILE="noxfile.py"
NOX_SESSION="${NOX_SESSION:-system-3.12}"
;;
"google-auth")
# Copy files needed for google-auth system tests
mkdir -p "${package_path}/system_tests/data"
Expand All @@ -86,14 +94,6 @@ run_package_test() {
NOX_FILE="noxfile.py"
fi
;;
"google-cloud-dns")
# EXPERIMENTAL: Force running all system sessions to test mixed results. This will be reverted
# before merge. You can safely ignore it.
PROJECT_ID=$(cat "${KOKORO_GFILE_DIR}/project-id.json")
GOOGLE_APPLICATION_CREDENTIALS="${KOKORO_GFILE_DIR}/service-account.json"
NOX_FILE="noxfile.py"
NOX_SESSION="system"
;;
*)
PROJECT_ID=$(cat "${KOKORO_GFILE_DIR}/project-id.json")
GOOGLE_APPLICATION_CREDENTIALS="${KOKORO_GFILE_DIR}/service-account.json"
Expand Down
116 changes: 112 additions & 4 deletions packages/gapic-generator/noxfile.py
Original file line number Diff line number Diff line change
Expand Up @@ -978,12 +978,120 @@ def format(session):
)


@contextmanager
def google_ads_library(session):
"""Generate and install the Google Ads GAPIC library for live system tests."""
session.install("-e", ".")
session.install("grpcio-tools", "pyYAML", "pypandoc-binary==1.16.2")

with tempfile.TemporaryDirectory() as tmp_dir:
googleapis_dir = path.join(tmp_dir, "googleapis")
sdk_dir = path.join(tmp_dir, "sdk")
os.makedirs(sdk_dir, exist_ok=True)

session.run(
"git",
"clone",
"--depth",
"1",
"--filter=blob:none",
"--sparse",
"https://github.com/googleapis/googleapis.git",
googleapis_dir,
external=True,
silent=True,
)
session.run(
"git",
"-C",
googleapis_dir,
"sparse-checkout",
"set",
"google/ads/googleads/v23",
"google/api",
"google/rpc",
"google/longrunning",
"google/type",
external=True,
silent=True,
)

src_yaml = path.join(
googleapis_dir, "google", "ads", "googleads", "v23", "googleads_v23.yaml"
)
dst_yaml = path.join(sdk_dir, "googleads_v23.yaml")
shutil.copyfile(src_yaml, dst_yaml)
session.run(
"python",
"-c",
(
"import yaml; "
f"p = {dst_yaml!r}; "
"data = yaml.safe_load(open(p, encoding='utf-8')); "
"data['apis'] = [{'name': 'google.ads.googleads.v23.services.YouTubeVideoUploadService'}]; "
"data.setdefault('publishing', {})['library_settings'] = ["
"{'version': 'google.ads.googleads.v23', "
"'python_settings': {'experimental_features': {'rest_async_io_enabled': True}}}"
"]; "
"yaml.safe_dump(data, open(p, 'w', encoding='utf-8'), default_flow_style=False, sort_keys=False)"
),
)

ads_protos = (
"google/ads/googleads/v23/services/youtube_video_upload_service.proto",
"google/ads/googleads/v23/resources/youtube_video_upload.proto",
"google/ads/googleads/v23/enums/youtube_video_privacy.proto",
)
desc_path = path.join(sdk_dir, "googleads.desc")
session.run(
"python",
"-m",
"grpc_tools.protoc",
"--experimental_allow_proto3_optional",
f"--proto_path={googleapis_dir}",
"--include_imports",
"--include_source_info",
f"-o{desc_path}",
*(path.join(googleapis_dir, p) for p in ads_protos),
external=True,
)

retry_config = path.join(
googleapis_dir,
"google",
"ads",
"googleads",
"v23",
"googleads_grpc_service_config.json",
)
session.run(
"python",
"-m",
"grpc_tools.protoc",
"--experimental_allow_proto3_optional",
f"--descriptor_set_in={desc_path}",
(
"--python_gapic_opt="
f"transport=grpc+rest,autogen-snippets=False,service-yaml={dst_yaml},retry-config={retry_config}"
),
f"--python_gapic_out={sdk_dir}",
*ads_protos,
external=True,
)

session.install("-e", f"{sdk_dir}[async_rest]")
yield sdk_dir


@nox.session(python=ALL_PYTHON)
def system(session):
# TODO(https://github.com/googleapis/google-cloud-python/issues/16190):
# Implement system test session.
"""Run the system test suite (skipped for migration)."""
session.skip(f"system session is not yet implemented for gapic-generator-python.")
"""Run the system test suite."""
with google_ads_library(session):
session.install("pytest", "pytest-asyncio")
session.run(
"py.test",
*(session.posargs or ["-vv", path.join("tests", "system_live")]),
)


@nox.session(python=NEWEST_PYTHON)
Expand Down
200 changes: 199 additions & 1 deletion packages/gapic-generator/tests/system/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,16 @@
HAS_GOOGLE_AUTH_AIO = False
import google.auth
from google.auth import credentials as ga_credentials

from google.showcase import EchoClient, IdentityClient, MessagingClient
try:
from google.showcase import ResumableUploadServiceClient

HAS_RESUMABLE_UPLOAD_CLIENT = True
except ImportError:
HAS_RESUMABLE_UPLOAD_CLIENT = False

HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = False
if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true":
import asyncio

Expand All @@ -63,6 +71,16 @@
HAS_ASYNC_REST_IDENTITY_TRANSPORT = True
except:
HAS_ASYNC_REST_IDENTITY_TRANSPORT = False
try:
from google.showcase import ResumableUploadServiceAsyncClient
from google.showcase_v1beta1.services.resumable_upload_service.transports.rest_asyncio import (
AsyncResumableUploadServiceRestTransport,
AsyncResumableUploadServiceRestInterceptor,
)

HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = True
except:
HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = False

_GRPC_VERSION = grpc.__version__

Expand All @@ -83,7 +101,9 @@ def async_anonymous_credentials():

@pytest.fixture
def event_loop():
return asyncio.get_event_loop()
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
return loop

@pytest_asyncio.fixture(params=["grpc_asyncio", "rest_asyncio"])
def async_echo(use_mtls, request, event_loop):
Expand Down Expand Up @@ -296,6 +316,33 @@ def post_expand_with_metadata(self, request, metadata):
return request, metadata


if HAS_RESUMABLE_UPLOAD_CLIENT:
try:
from google.showcase_v1beta1.services.resumable_upload_service.transports import (
ResumableUploadServiceRestInterceptor,
)

class ResumableUploadMetadataClientRestInterceptor(
ResumableUploadServiceRestInterceptor
):
request_metadata: Sequence[Tuple[str, str]] = []
response_metadata: Sequence[Tuple[str, str]] = []

def pre_upload_media(self, request, metadata):
self.request_metadata = metadata
return request, metadata

def post_upload_media_with_metadata(self, request, metadata):
self.response_metadata = metadata
return request, metadata

HAS_RESUMABLE_UPLOAD_INTERCEPTOR = True
except ImportError:
HAS_RESUMABLE_UPLOAD_INTERCEPTOR = False
else:
HAS_RESUMABLE_UPLOAD_INTERCEPTOR = False


if HAS_ASYNC_REST_ECHO_TRANSPORT:

class EchoMetadataClientRestAsyncInterceptor(AsyncEchoRestInterceptor):
Expand All @@ -319,6 +366,19 @@ async def post_expand_with_metadata(self, request, metadata):
return request, metadata


if HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT:

class ResumableUploadMetadataClientRestAsyncInterceptor(
AsyncResumableUploadServiceRestInterceptor
):
request_metadata: Sequence[Tuple[str, str]] = []
response_metadata: Sequence[Tuple[str, str]] = []

async def pre_upload_media(self, request, metadata):
self.request_metadata = metadata
return request, metadata


class EchoMetadataClientGrpcInterceptor(
grpc.UnaryUnaryClientInterceptor,
grpc.UnaryStreamClientInterceptor,
Expand Down Expand Up @@ -535,6 +595,144 @@ def intercepted_echo_rest_async():
return EchoAsyncClient(transport=transport), interceptor


@pytest.fixture
def intercepted_resumable_upload_rest(use_mtls, use_tls):
if not HAS_RESUMABLE_UPLOAD_CLIENT or not HAS_RESUMABLE_UPLOAD_INTERCEPTOR:
pytest.skip("ResumableUploadServiceClient not available.")

transport_name = "rest"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

We need tests starting from clients created with grpc transports too, because those are the ones I'm more concerned about

I worry some of these fixtures are obscuring important details here Do we have any tests that go from creating a client to reading the result, exactly we expect end-users would?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I created a live test which doesn't use the conftest.py fixtures: packages/gapic-generator/tests/system_live/test_google_ads_resumable_upload.py . This now runs as a system test

transport_cls = ResumableUploadServiceClient.get_transport_class(transport_name)
interceptor = ResumableUploadMetadataClientRestInterceptor()

url_scheme = "https" if (use_mtls or use_tls) else "http"
transport = transport_cls(
credentials=ga_credentials.AnonymousCredentials(),
host="localhost:7469",
url_scheme=url_scheme,
interceptor=interceptor,
)
if use_mtls or use_tls:
transport._session.verify = CERT_PATH
transport._session.mount("https://", HostNameIgnoringAdapter())
if use_mtls:
transport._session.cert = (CERT_PATH, KEY_PATH)

return ResumableUploadServiceClient(transport=transport), interceptor


@pytest.fixture
def intercepted_resumable_upload_rest_async():
if not HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT:
pytest.skip("Skipping test with async rest.")

interceptor = ResumableUploadMetadataClientRestAsyncInterceptor()

transport = AsyncResumableUploadServiceRestTransport(
credentials=async_anonymous_credentials(),
host="localhost:7469",
url_scheme="http",
interceptor=interceptor,
)

return ResumableUploadServiceAsyncClient(transport=transport), interceptor


try:
from google.api_core.resumable_transfer import (
ResumableUploadConfig,
ResumableUploadSession,
)
except ImportError:
ResumableUploadConfig = None
ResumableUploadSession = None


def make_resumable_upload(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

It looks like the tests still exercise this code in a way different than the user would. Is it possible to create a client, and start the stream using the client apis, instead of creating the session directly?

If this is the main system test, we should make sure we have good end-to-end coverage

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I created a live test which doesn't use conftest.py : packages/gapic-generator/tests/system_live/test_google_ads_resumable_upload.py

packages/gapic-generator/tests/system/conftest.py is used only for showcase integration tests

transport,
request_body,
stream,
upload_url,
size=None,
config=None,
**kwargs,
):
content_type = kwargs.pop("content_type", "application/octet-stream")
response_type = kwargs.pop("response_type", None)
retry = kwargs.pop("retry", None)
timeout = kwargs.pop("timeout", None)

if config is None:
config = ResumableUploadConfig(**kwargs)
elif kwargs:
for k, v in kwargs.items():
if hasattr(config, k):
setattr(config, k, v)

# ``make_resumable_upload`` instantiates ``ResumableUploadSession`` directly
# rather than calling the generated GAPIC client method (``client.upload_media``).
# Unlike the generated GAPIC REST transport, ``ResumableUploadSession`` in
# ``google-api-core`` is payload-format agnostic and does not set
# ``Content-Type: application/json`` on the start request automatically.
headers = dict(config.start_headers or [])
if request_body and "Content-Type" not in headers:
headers["Content-Type"] = "application/json"
config.headers = headers

session = ResumableUploadSession(
upload_url=upload_url,
config=config,
content_type=content_type,
response_type=response_type,
transport=transport,
)
return session.upload(
stream=stream,
request_body=request_body,
content_type=content_type,
size=size,
transport=transport,
retry=retry,
timeout=timeout,
)


def resume_resumable_upload(
transport,
upload_url,
stream,
size=None,
config=None,
**kwargs,
):
content_type = kwargs.pop("content_type", None)
response_type = kwargs.pop("response_type", None)
retry = kwargs.pop("retry", None)
timeout = kwargs.pop("timeout", None)

if config is None:
config = ResumableUploadConfig(**kwargs)
elif kwargs:
for k, v in kwargs.items():
if hasattr(config, k):
setattr(config, k, v)

session = ResumableUploadSession(
upload_url=upload_url,
config=config,
transport=transport,
content_type=content_type,
response_type=response_type,
)
return session.resume(
upload_url=upload_url,
stream=stream,
size=size,
transport=transport,
retry=retry,
timeout=timeout,
)


def pytest_terminal_summary(terminalreporter, exitstatus, config):
"""Prints a Telemetry Span Compliance summary to the console.

Expand Down
Loading
Loading