From 1a7b90818e320fee5e5f580b844b7f13c33208a6 Mon Sep 17 00:00:00 2001 From: aasthabharill Date: Mon, 10 Aug 2026 17:52:27 +0530 Subject: [PATCH 1/7] initial review --- ...alLimitedDurationErrorInjectionPolicy.java | 19 +- v2/gcs-spanner-dv/pom.xml | 181 +++++++++++------- .../teleport/v2/templates/GCSSpannerDV.java | 32 +++- .../v2/templates/GCSSpannerDVFTBase.java | 121 ++++++++++++ .../templates/GCSSpannerDVSpannerReadFT.java | 163 ++++++++++++++++ .../spanner-schema.sql | 12 ++ 6 files changed, 439 insertions(+), 89 deletions(-) create mode 100644 v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java create mode 100644 v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVSpannerReadFT.java create mode 100644 v2/gcs-spanner-dv/src/test/resources/GCSSpannerDVSpannerReadFT/spanner-schema.sql diff --git a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java index fffbe04c40..926e62330e 100644 --- a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java +++ b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java @@ -22,6 +22,7 @@ import java.time.Duration; import java.time.Instant; import java.time.format.DateTimeParseException; +import java.util.concurrent.atomic.AtomicLong; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -37,12 +38,12 @@ public class InitialLimitedDurationErrorInjectionPolicy LoggerFactory.getLogger(InitialLimitedDurationErrorInjectionPolicy.class); private static final long serialVersionUID = 1L; - private Instant startTime; + private static Instant startTime; private final Duration injectionDuration; private final String effectiveDurationParameter; private String errorCodeToBeInjected; private Clock clock; - private long callCount; + private static final AtomicLong callCount = new AtomicLong(0); private static final String DEFAULT_DURATION = "PT10M"; private static final String DURATION_FIELD_IN_OBJECT = "duration"; @@ -124,22 +125,20 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock */ @Override public boolean shouldInjectionError() { - if (this.startTime == null) { + if (startTime == null) { synchronized (this) { - if (this.startTime == null) { - this.startTime = Instant.now(clock); + if (startTime == null) { + startTime = Instant.now(clock); LOG.info( "First call detected. Errors will be injected for {} starting from {}.", this.injectionDuration, - this.startTime); + startTime); } } } - synchronized (this) { - ++callCount; - } + long currentCallCount = callCount.incrementAndGet(); - if (callCount < INITIAL_ALLOWED_CALLS_COUNT) { + if (currentCallCount < INITIAL_ALLOWED_CALLS_COUNT) { return false; } diff --git a/v2/gcs-spanner-dv/pom.xml b/v2/gcs-spanner-dv/pom.xml index 1a06842b11..29e42eb657 100644 --- a/v2/gcs-spanner-dv/pom.xml +++ b/v2/gcs-spanner-dv/pom.xml @@ -16,82 +16,117 @@ ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~--> - 4.0.0 + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/maven-v4_0_0.xsd"> + 4.0.0 - - com.google.cloud.teleport.v2 - dynamic-templates - 1.0-SNAPSHOT - + + com.google.cloud.teleport.v2 + dynamic-templates + 1.0-SNAPSHOT + - gcs-spanner-dv + gcs-spanner-dv - - - com.google.cloud.teleport.v2 - common - ${project.version} - - - com.google.guava - guava - ${guava.version} - - - com.google.cloud - google-cloud-core - + + + com.google.cloud.teleport.v2 + common + ${project.version} + + + com.google.guava + guava + ${guava.version} + + + com.google.cloud + google-cloud-core + - - - com.google.cloud.teleport - it-google-cloud-platform - ${project.version} - test - - - org.apache.beam - beam-it-jdbc - test - - - mysql - mysql-connector-java - ${mysql-connector-java.version} - test - - - com.google.cloud.teleport.v2 - spanner-common - 1.0-SNAPSHOT - compile - - - org.mockito - mockito-inline - LATEST - test - - + + + com.google.cloud.teleport + it-google-cloud-platform + ${project.version} + test + + + org.apache.beam + beam-it-jdbc + test + + + mysql + mysql-connector-java + ${mysql-connector-java.version} + test + + + com.google.cloud.teleport.v2 + spanner-common + 1.0-SNAPSHOT + compile + + + org.mockito + mockito-inline + LATEST + test + + - - - - - org.jacoco - jacoco-maven-plugin - ${jacoco.version} - - - - com/google/cloud/teleport/v2/dto/** - com/google/cloud/teleport/v2/constants/** - - - - - + + + + + org.jacoco + jacoco-maven-plugin + ${jacoco.version} + + + + com/google/cloud/teleport/v2/dto/** + com/google/cloud/teleport/v2/constants/** + + + + + + + + + useRealSpanner + + true + + !activateFailureInjection + + + + + com.google.cloud.teleport.v2 + real-spanner-service + ${project.version} + + + + + failureInjectionTest + + + activateFailureInjection + true + + + + + com.google.cloud.teleport.v2 + failure-injected-spanner-service + ${project.version} + + + + diff --git a/v2/gcs-spanner-dv/src/main/java/com/google/cloud/teleport/v2/templates/GCSSpannerDV.java b/v2/gcs-spanner-dv/src/main/java/com/google/cloud/teleport/v2/templates/GCSSpannerDV.java index 35a6011291..6da4c19c84 100644 --- a/v2/gcs-spanner-dv/src/main/java/com/google/cloud/teleport/v2/templates/GCSSpannerDV.java +++ b/v2/gcs-spanner-dv/src/main/java/com/google/cloud/teleport/v2/templates/GCSSpannerDV.java @@ -37,6 +37,7 @@ import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.PipelineResult; import org.apache.beam.sdk.io.gcp.spanner.SpannerConfig; +import org.apache.beam.sdk.io.gcp.spanner.SpannerServiceFactoryImpl; import org.apache.beam.sdk.options.Default; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; @@ -254,6 +255,16 @@ public interface Options extends PipelineOptions { String getTransformationCustomParameters(); void setTransformationCustomParameters(String value); + + @TemplateParameter.Text( + order = 16, + optional = true, + description = "Failure injection parameter", + helpText = "Failure injection parameter. Only used for testing.") + @Default.String("") + String getFailureInjectionParameter(); + + void setFailureInjectionParameter(String value); } public static void main(String[] args) { @@ -327,11 +338,20 @@ public static PipelineResult run(Options options) { @VisibleForTesting static SpannerConfig createSpannerConfig(Options options) { - return SpannerConfig.create() - .withProjectId(ValueProvider.StaticValueProvider.of(options.getProjectId())) - .withHost(ValueProvider.StaticValueProvider.of(options.getSpannerHost())) - .withInstanceId(ValueProvider.StaticValueProvider.of(options.getInstanceId())) - .withDatabaseId(ValueProvider.StaticValueProvider.of(options.getDatabaseId())) - .withRpcPriority(ValueProvider.StaticValueProvider.of(options.getSpannerPriority())); + SpannerConfig config = + SpannerConfig.create() + .withProjectId(ValueProvider.StaticValueProvider.of(options.getProjectId())) + .withHost(ValueProvider.StaticValueProvider.of(options.getSpannerHost())) + .withInstanceId(ValueProvider.StaticValueProvider.of(options.getInstanceId())) + .withDatabaseId(ValueProvider.StaticValueProvider.of(options.getDatabaseId())) + .withRpcPriority(ValueProvider.StaticValueProvider.of(options.getSpannerPriority())); + + if (options.getFailureInjectionParameter() != null + && !options.getFailureInjectionParameter().isEmpty()) { + config = + SpannerServiceFactoryImpl.createSpannerService( + config, options.getFailureInjectionParameter()); + } + return config; } } diff --git a/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java b/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java new file mode 100644 index 0000000000..25a3e3fa1b --- /dev/null +++ b/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java @@ -0,0 +1,121 @@ +/* + * Copyright (C) 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package com.google.cloud.teleport.v2.templates; + +import com.google.cloud.teleport.v2.spanner.migrations.transformation.CustomTransformation; +import com.google.common.io.Resources; +import java.io.IOException; +import java.util.Map; +import org.apache.beam.it.common.PipelineLauncher.LaunchInfo; +import org.apache.beam.it.common.utils.PipelineUtils; +import org.apache.beam.it.gcp.dataflow.FlexTemplateDataflowJobResourceManager; +import org.apache.beam.it.gcp.spanner.SpannerResourceManager; + +/** + * Base class for gcs-spanner-dv failure injection integration tests. + * + *

Why is this a separate class from {@link GCSSpannerDVITBase}? + * + *

While {@code GCSSpannerDVITBase} is runner-agnostic and relies on the generic {@code + * PipelineLauncher} (allowing tests to run locally via DirectRunner), failure injection testing + * explicitly requires building a custom Docker image with the {@code failureInjectionTest} Maven + * profile. + * + *

Therefore, tests extending this class are strictly coupled to Dataflow Flex Templates and + * bypass the generic launcher in favor of {@link FlexTemplateDataflowJobResourceManager}. + */ +public abstract class GCSSpannerDVFTBase extends GCSSpannerDVITBase { + + /** + * Launches the Dataflow job with failure injection testing capabilities using + * FlexTemplateDataflowJobResourceManager. + */ + protected LaunchInfo launchFTDataflowJob( + String testId, + String projectId, + SpannerResourceManager spannerResourceManager, + String bigQueryDataset, + String gcsInputDirectory, + String sessionFileResourceName, + String schemaOverridesFileResourceName, + String tableOverrides, + String columnOverrides, + CustomTransformation customTransformation, + String failureInjectionParameter, + Map jobParameters) + throws IOException { + + FlexTemplateDataflowJobResourceManager.Builder flexTemplateBuilder = + FlexTemplateDataflowJobResourceManager.builder(testName) + .withTemplateName("GCS_Spanner_Data_Validator") + .withTemplateModulePath("v2/gcs-spanner-dv") + .withAdditionalMavenProfile("failureInjectionTest") + .addEnvironmentVariable( + "additionalExperiments", java.util.Collections.singletonList("disable_runner_v2")); + + if (failureInjectionParameter != null && !failureInjectionParameter.isEmpty()) { + flexTemplateBuilder.addParameter("failureInjectionParameter", failureInjectionParameter); + } + + flexTemplateBuilder.addParameter("projectId", projectId); + flexTemplateBuilder.addParameter("instanceId", spannerResourceManager.getInstanceId()); + flexTemplateBuilder.addParameter("databaseId", spannerResourceManager.getDatabaseId()); + flexTemplateBuilder.addParameter("bigQueryDataset", bigQueryDataset); + flexTemplateBuilder.addParameter("gcsInputDirectory", gcsInputDirectory); + + if (sessionFileResourceName != null) { + gcsClient.uploadArtifact( + "session.json", Resources.getResource(sessionFileResourceName).getPath()); + flexTemplateBuilder.addParameter("sessionFilePath", getGcsPath("session.json")); + } + + if (schemaOverridesFileResourceName != null) { + gcsClient.uploadArtifact( + "schema_overrides.json", + Resources.getResource(schemaOverridesFileResourceName).getPath()); + flexTemplateBuilder.addParameter( + "schemaOverridesFilePath", getGcsPath("schema_overrides.json")); + } + + if (tableOverrides != null) { + flexTemplateBuilder.addParameter("tableOverrides", tableOverrides); + } + + if (columnOverrides != null) { + flexTemplateBuilder.addParameter("columnOverrides", columnOverrides); + } + + if (customTransformation != null) { + flexTemplateBuilder.addParameter( + "transformationJarPath", getGcsPath(customTransformation.jarPath())); + flexTemplateBuilder.addParameter("transformationClassName", customTransformation.classPath()); + if (customTransformation.customParameters() != null) { + flexTemplateBuilder.addParameter( + "transformationCustomParameters", customTransformation.customParameters()); + } + } + + String runId = PipelineUtils.createJobName(testId); + flexTemplateBuilder.addParameter("runId", runId); + flexTemplateBuilder.addParameter("workerMachineType", "n2-standard-4"); + + if (jobParameters != null) { + jobParameters.forEach(flexTemplateBuilder::addParameter); + } + + return flexTemplateBuilder.build().launchJob(); + } +} diff --git a/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVSpannerReadFT.java b/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVSpannerReadFT.java new file mode 100644 index 0000000000..7a1fe3186f --- /dev/null +++ b/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVSpannerReadFT.java @@ -0,0 +1,163 @@ +/* + * Copyright (C) 2026 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); you may not + * use this file except in compliance with the License. You may obtain a copy of + * the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, WITHOUT + * WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the + * License for the specific language governing permissions and limitations under + * the License. + */ +package com.google.cloud.teleport.v2.templates; + +import com.google.cloud.Timestamp; +import com.google.cloud.spanner.Mutation; +import com.google.cloud.teleport.metadata.SkipDirectRunnerTest; +import com.google.cloud.teleport.metadata.TemplateIntegrationTest; +import com.google.cloud.teleport.v2.templates.GCSSpannerDVAvroSetupHelper.RecordBuilder; +import com.google.cloud.teleport.v2.templates.GCSSpannerDVAvroSetupHelper.TableDef; +import com.google.cloud.teleport.v2.templates.GCSSpannerDVTestAsserts.TableValidationStatsDto; +import com.google.cloud.teleport.v2.templates.GCSSpannerDVTestAsserts.ValidationSummaryDto; +import java.io.IOException; +import java.time.Instant; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import org.apache.avro.generic.GenericRecord; +import org.apache.beam.it.common.PipelineLauncher.LaunchInfo; +import org.junit.Before; +import org.junit.Test; +import org.junit.experimental.categories.Category; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; + +/** + * Tests transient Spanner read failures (via SpannerIO) for GCSSpannerDV pipeline. + * + *

Test cases covered: + * + *

+ */ +@Category({TemplateIntegrationTest.class, SkipDirectRunnerTest.class}) +@TemplateIntegrationTest(GCSSpannerDV.class) +@RunWith(JUnit4.class) +public class GCSSpannerDVSpannerReadFT extends GCSSpannerDVFTBase { + + private static final String SPANNER_DDL_RESOURCE = "GCSSpannerDVSpannerReadFT/spanner-schema.sql"; + private static final int NUM_RECORDS = 500; + + @Before + public void setUp() throws IOException { + spannerResourceManager = setUpSpannerResourceManager(); + createSpannerDDL(spannerResourceManager, SPANNER_DDL_RESOURCE); + bigQueryResourceManager = setUpBigQueryResourceManager(); + bigQueryResourceManager.createDataset(REGION); + } + + // Tests transient UNAVAILABLE errors on SpannerIO.readAll() to ensure native exponential backoff + // succeeds. + @Test + public void testTransientReadFailure() throws IOException, InterruptedException { + List records = new ArrayList<>(); + List mutations = new ArrayList<>(); + Instant now = Instant.now().truncatedTo(java.time.temporal.ChronoUnit.MILLIS); + + // We insert 500 rows to ensure Dataflow has enough data to potentially split bundles + // and exercise the SpannerIO read logic across multiple task execution boundaries, + // rather than trivially passing with a single row. + for (int i = 0; i < NUM_RECORDS; i++) { + long userId = (long) i; + String eventId = "E" + i; + String fullName = "User " + i; + int age = 20 + (i % 30); + + // Avro Record + GenericRecord record = + new RecordBuilder(TableDef.USERS, null) + .set("user_id", userId) + .set("event_id", eventId) + .set("full_name", fullName) + .set("age", age) + .set("created_at", now) + .build(); + records.add(record); + + // Spanner Mutation + mutations.add( + Mutation.newInsertOrUpdateBuilder("Users") + .set("user_id") + .to(userId) + .set("event_id") + .to(eventId) + .set("full_name") + .to(fullName) + .set("age") + .to(age) + .set("created_at") + .to(Timestamp.ofTimeSecondsAndNanos(now.getEpochSecond(), now.getNano())) + .build()); + } + + String gcsInputDirectory = getGcsPath("input"); + uploadAvroFileToGcs("input/users.avro", TableDef.USERS.schema, records); + spannerResourceManager.write(mutations); + + // Injects a 60-second UNAVAILABLE outage specifically on the Spanner workers to trigger + // Dataflow task retries. + String failureInjectionParam = + "{\"policyType\":\"InitialLimitedDurationErrorInjectionPolicy\", \"policyInput\": {\"duration\":\"PT1M\", \"errorCode\":\"UNAVAILABLE\"}}"; + String bqDatasetId = bigQueryResourceManager.getDatasetId(); + + LaunchInfo jobInfo = + launchFTDataflowJob( + testName, + PROJECT, + spannerResourceManager, + bqDatasetId, + gcsInputDirectory, + null, + null, + null, + null, + null, + failureInjectionParam, + new HashMap<>()); + + pipelineOperator().waitUntilDone(createConfig(jobInfo)); + + GCSSpannerDVTestAsserts.assertValidationSummary( + bigQueryResourceManager, + Arrays.asList( + new ValidationSummaryDto( + /* status= */ "MATCH", + /* totalTablesValidated= */ 1L, + /* totalRowsMatched= */ 500L, + /* totalRowsMismatched= */ 0L, + /* tablesWithMismatches= */ ""))); + + GCSSpannerDVTestAsserts.assertTableValidationStats( + bigQueryResourceManager, + Arrays.asList( + new TableValidationStatsDto( + /* schemaName= */ null, + /* tableName= */ "Users", + /* status= */ "MATCH", + /* sourceRowCount= */ 500L, + /* destinationRowCount= */ 500L, + /* matchedRowCount= */ 500L, + /* mismatchRowCount= */ 0L))); + + // No mismatched records should exist + GCSSpannerDVTestAsserts.assertMismatchedRecords(bigQueryResourceManager, Arrays.asList()); + } +} diff --git a/v2/gcs-spanner-dv/src/test/resources/GCSSpannerDVSpannerReadFT/spanner-schema.sql b/v2/gcs-spanner-dv/src/test/resources/GCSSpannerDVSpannerReadFT/spanner-schema.sql new file mode 100644 index 0000000000..815ce495c3 --- /dev/null +++ b/v2/gcs-spanner-dv/src/test/resources/GCSSpannerDVSpannerReadFT/spanner-schema.sql @@ -0,0 +1,12 @@ +CREATE TABLE Users ( + user_id INT64 NOT NULL, + event_id STRING(MAX) NOT NULL, + full_name STRING(MAX), + age INT64, + created_at TIMESTAMP +) PRIMARY KEY (user_id, event_id); + +CREATE TABLE AccountRoles ( + role_id INT64 NOT NULL, + role_name STRING(MAX) +) PRIMARY KEY (role_id); From b8ec2860cf14ffcd86b92aa10c7a60b8e95f5418 Mon Sep 17 00:00:00 2001 From: aasthabharill Date: Mon, 10 Aug 2026 18:10:59 +0530 Subject: [PATCH 2/7] gemini-review --- .../InitialLimitedDurationErrorInjectionPolicy.java | 2 +- .../google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java index 926e62330e..3f713151ec 100644 --- a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java +++ b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java @@ -126,7 +126,7 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock @Override public boolean shouldInjectionError() { if (startTime == null) { - synchronized (this) { + synchronized (InitialLimitedDurationErrorInjectionPolicy.class) { if (startTime == null) { startTime = Instant.now(clock); LOG.info( diff --git a/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java b/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java index 25a3e3fa1b..c6c7694a37 100644 --- a/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java +++ b/v2/gcs-spanner-dv/src/test/java/com/google/cloud/teleport/v2/templates/GCSSpannerDVFTBase.java @@ -59,7 +59,7 @@ protected LaunchInfo launchFTDataflowJob( throws IOException { FlexTemplateDataflowJobResourceManager.Builder flexTemplateBuilder = - FlexTemplateDataflowJobResourceManager.builder(testName) + FlexTemplateDataflowJobResourceManager.builder(testId) .withTemplateName("GCS_Spanner_Data_Validator") .withTemplateModulePath("v2/gcs-spanner-dv") .withAdditionalMavenProfile("failureInjectionTest") From 36a0b466b73810f938279bd7b0e8f71114d5388c Mon Sep 17 00:00:00 2001 From: aasthabharill Date: Tue, 11 Aug 2026 14:56:21 +0530 Subject: [PATCH 3/7] test failures --- .../InitialLimitedDurationErrorInjectionPolicy.java | 5 +++++ .../InitialLimitedDurationErrorInjectionPolicyTest.java | 6 ++++++ 2 files changed, 11 insertions(+) diff --git a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java index 3f713151ec..2250dc7105 100644 --- a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java +++ b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java @@ -185,6 +185,11 @@ void setClockForTesting(Clock clock) { this.clock = clock; } + static void resetForTest() { + startTime = null; + callCount.set(0); + } + @Override public String toString() { return "InitialLimitedDurationErrorInjectionPolicy{" diff --git a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java index fe4a3f8759..35b434443f 100644 --- a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java +++ b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java @@ -28,10 +28,16 @@ import java.time.Duration; import java.time.Instant; import java.time.ZoneOffset; +import org.junit.Before; import org.junit.Test; public class InitialLimitedDurationErrorInjectionPolicyTest { + @Before + public void setUp() { + InitialLimitedDurationErrorInjectionPolicy.resetForTest(); + } + private ObjectNode createInputObject(String duration) { ObjectNode node = JsonNodeFactory.instance.objectNode(); if (duration != null) { From 614eeef2c89a0f0e15a2814f9c08807e5c876d3b Mon Sep 17 00:00:00 2001 From: aasthabharill Date: Mon, 24 Aug 2026 06:25:35 +0000 Subject: [PATCH 4/7] intermediate changes --- ...alLimitedDurationErrorInjectionPolicy.java | 52 +++++++----- ...mitedDurationErrorInjectionPolicyTest.java | 28 +++++-- ...TransactionTimeoutInjectionPolicyTest.java | 83 +++++++++++++++++++ 3 files changed, 139 insertions(+), 24 deletions(-) diff --git a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java index 2250dc7105..5bc4d70814 100644 --- a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java +++ b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java @@ -22,6 +22,9 @@ import java.time.Duration; import java.time.Instant; import java.time.format.DateTimeParseException; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicLong; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -38,12 +41,18 @@ public class InitialLimitedDurationErrorInjectionPolicy LoggerFactory.getLogger(InitialLimitedDurationErrorInjectionPolicy.class); private static final long serialVersionUID = 1L; - private static Instant startTime; + private static class SharedState { + volatile Instant startTime = null; + final AtomicLong callCount = new AtomicLong(0); + } + + private static final ConcurrentMap sharedStates = new ConcurrentHashMap<>(); + + private final String instanceId; private final Duration injectionDuration; private final String effectiveDurationParameter; - private String errorCodeToBeInjected; + private String errorCodeToBeInjected = io.grpc.Status.Code.DEADLINE_EXCEEDED.name(); private Clock clock; - private static final AtomicLong callCount = new AtomicLong(0); private static final String DEFAULT_DURATION = "PT10M"; private static final String DURATION_FIELD_IN_OBJECT = "duration"; @@ -69,6 +78,7 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter) { */ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock clock) { this.clock = clock; + this.instanceId = UUID.randomUUID().toString(); String durationString = DEFAULT_DURATION; if (inputParameter != null && !inputParameter.isMissingNode() && !inputParameter.isNull()) { @@ -97,6 +107,10 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock Code.DEADLINE_EXCEEDED); } } + JsonNode jobStartTimePath = inputParameter.path("jobStartTime"); + if (jobStartTimePath.isTextual()) { + LOG.warn("jobStartTime is ignored in InitialLimitedDurationErrorInjectionPolicy because lazy initialization is preferred."); + } } } @@ -125,25 +139,28 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock */ @Override public boolean shouldInjectionError() { - if (startTime == null) { - synchronized (InitialLimitedDurationErrorInjectionPolicy.class) { - if (startTime == null) { - startTime = Instant.now(clock); + SharedState state = sharedStates.computeIfAbsent(this.instanceId, k -> new SharedState()); + + if (state.startTime == null) { + synchronized (state) { + if (state.startTime == null) { + state.startTime = Instant.now(clock); LOG.info( - "First call detected. Errors will be injected for {} starting from {}.", + "First call detected for instance {}. Errors will be injected for {} starting from {}.", + this.instanceId, this.injectionDuration, - startTime); + state.startTime); } } } - long currentCallCount = callCount.incrementAndGet(); + long currentCallCount = state.callCount.incrementAndGet(); if (currentCallCount < INITIAL_ALLOWED_CALLS_COUNT) { return false; } Instant now = Instant.now(clock); - Duration elapsed = Duration.between(startTime, now); + Duration elapsed = Duration.between(state.startTime, now); // Compare elapsed time with the configured duration. // elapsed.compareTo(injectionDuration) < 0 means elapsed < injectionDuration @@ -178,23 +195,20 @@ public String getEffectiveDurationParameter() { } public Instant getStartTime() { - return startTime; + SharedState state = sharedStates.get(this.instanceId); + return state != null ? state.startTime : null; } void setClockForTesting(Clock clock) { this.clock = clock; } - static void resetForTest() { - startTime = null; - callCount.set(0); - } - @Override public String toString() { return "InitialLimitedDurationErrorInjectionPolicy{" - + "startTime=" - + startTime + + "instanceId='" + + instanceId + + '\'' + ", injectionDuration=" + injectionDuration + ", effectiveDurationParameter='" diff --git a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java index 35b434443f..dced238320 100644 --- a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java +++ b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java @@ -28,15 +28,10 @@ import java.time.Duration; import java.time.Instant; import java.time.ZoneOffset; -import org.junit.Before; import org.junit.Test; public class InitialLimitedDurationErrorInjectionPolicyTest { - @Before - public void setUp() { - InitialLimitedDurationErrorInjectionPolicy.resetForTest(); - } private ObjectNode createInputObject(String duration) { ObjectNode node = JsonNodeFactory.instance.objectNode(); @@ -274,4 +269,27 @@ public void shouldInjectError_startTimeDoesNotChangeAfterFirstCall() { Instant thirdStartTime = policy.getStartTime(); assertEquals("Start time should still not change", firstStartTime, thirdStartTime); } + + @Test + public void constructor_shouldParseErrorCode() { + ObjectNode input = JsonNodeFactory.instance.objectNode(); + input.put("duration", "PT5S"); + input.put("errorCode", "UNAVAILABLE"); + Clock clock = Clock.fixed(Instant.EPOCH, ZoneOffset.UTC); + InitialLimitedDurationErrorInjectionPolicy policy = + new InitialLimitedDurationErrorInjectionPolicy(input, clock); + + assertEquals("UNAVAILABLE", policy.getErrorCodeToBeInjected()); + } + + @Test + public void constructor_shouldUseDefaultErrorCodeIfBlank() { + ObjectNode input = JsonNodeFactory.instance.objectNode(); + input.put("duration", "PT5S"); + input.put("errorCode", " "); + Clock clock = Clock.fixed(Instant.EPOCH, ZoneOffset.UTC); + InitialLimitedDurationErrorInjectionPolicy policy = + new InitialLimitedDurationErrorInjectionPolicy(input, clock); + + } } diff --git a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java index df494ff550..e2ddca151a 100644 --- a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java +++ b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java @@ -266,4 +266,87 @@ public void shouldInjectError_alwaysReturnsFalse() { // The method should always return false, as its purpose is to delay, not to inject an error. assertThat(policy.shouldInjectionError()).isFalse(); } + + @Test + public void shouldInjectError_concurrentCallsCoverDoubleCheckedLocking() throws Exception { + ObjectNode input = JsonNodeFactory.instance.objectNode(); + input.put("transactionTimeoutBakeDuration", "PT5M"); + TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); + + Thread waitingThread = new Thread(() -> { + policy.shouldInjectionError(); + }); + + synchronized (policy) { + waitingThread.start(); + Thread.sleep(500); + + java.lang.reflect.Field startTimeField = + TransactionTimeoutInjectionPolicy.class.getDeclaredField("startTime"); + startTimeField.setAccessible(true); + startTimeField.set(policy, Instant.now()); + } + + waitingThread.join(); + java.lang.reflect.Field startTimeField = + TransactionTimeoutInjectionPolicy.class.getDeclaredField("startTime"); + startTimeField.setAccessible(true); + assertThat(startTimeField.get(policy)).isNotNull(); + } + + @Test + public void constructor_shouldParseJobStartTime() { + ObjectNode input = JsonNodeFactory.instance.objectNode(); + input.put("jobStartTime", "2025-01-01T00:00:00Z"); + Clock clock = Clock.fixed(Instant.EPOCH, ZoneOffset.UTC); + TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input, clock); + policy.setClockForTesting(Clock.fixed(Instant.parse("2025-01-01T01:00:00Z"), ZoneOffset.UTC)); + long start = System.currentTimeMillis(); + policy.shouldInjectionError(); + long end = System.currentTimeMillis(); + assertThat(end - start).isLessThan(50L); + } + + @Test + public void constructor_shouldThrowExceptionForInvalidJobStartTime() { + ObjectNode input = JsonNodeFactory.instance.objectNode(); + input.put("jobStartTime", "Invalid"); + IllegalArgumentException e = + assertThrows( + IllegalArgumentException.class, + () -> new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC())); + assertThat(e).hasMessageThat().contains("Failed to parse jobStartTime"); + } + + @Test + public void shouldInjectionError_initializesClockIfNull() throws Exception { + ObjectNode input = JsonNodeFactory.instance.objectNode(); + TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input); + java.lang.reflect.Field clockField = TransactionTimeoutInjectionPolicy.class.getDeclaredField("clock"); + clockField.setAccessible(true); + clockField.set(policy, null); + + assertThat(policy.shouldInjectionError()).isFalse(); + } + + @Test + public void shouldInjectDelay_handlesInterruptedException() throws Exception { + ObjectNode input = JsonNodeFactory.instance.objectNode(); + input.put("transactionTimeoutBakeDuration", "PT5M"); + input.put("transactionDelayDuration", "PT5M"); + TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); + + policy.setRandomForTesting(new Random() { + @Override + public double nextDouble() { return 0.1; } + }); + + Thread thread = new Thread(() -> policy.shouldInjectionError()); + thread.start(); + Thread.sleep(200); + thread.interrupt(); + + thread.join(1000); + assertThat(thread.isAlive()).isFalse(); + } } From 1778460603a5382d31a1d991f79bbc15f5344a11 Mon Sep 17 00:00:00 2001 From: aasthabharill Date: Mon, 24 Aug 2026 07:25:15 +0000 Subject: [PATCH 5/7] df runner --- ...alLimitedDurationErrorInjectionPolicy.java | 52 +++++++------------ ...mitedDurationErrorInjectionPolicyTest.java | 5 ++ 2 files changed, 24 insertions(+), 33 deletions(-) diff --git a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java index 5bc4d70814..94feabad77 100644 --- a/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java +++ b/v2/failure-injection-policies/src/main/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicy.java @@ -22,9 +22,6 @@ import java.time.Duration; import java.time.Instant; import java.time.format.DateTimeParseException; -import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicLong; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -41,17 +38,11 @@ public class InitialLimitedDurationErrorInjectionPolicy LoggerFactory.getLogger(InitialLimitedDurationErrorInjectionPolicy.class); private static final long serialVersionUID = 1L; - private static class SharedState { - volatile Instant startTime = null; - final AtomicLong callCount = new AtomicLong(0); - } - - private static final ConcurrentMap sharedStates = new ConcurrentHashMap<>(); - - private final String instanceId; + private static volatile Instant startTime = null; + private static final AtomicLong callCount = new AtomicLong(0); private final Duration injectionDuration; private final String effectiveDurationParameter; - private String errorCodeToBeInjected = io.grpc.Status.Code.DEADLINE_EXCEEDED.name(); + private String errorCodeToBeInjected; private Clock clock; private static final String DEFAULT_DURATION = "PT10M"; @@ -78,7 +69,6 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter) { */ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock clock) { this.clock = clock; - this.instanceId = UUID.randomUUID().toString(); String durationString = DEFAULT_DURATION; if (inputParameter != null && !inputParameter.isMissingNode() && !inputParameter.isNull()) { @@ -107,10 +97,6 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock Code.DEADLINE_EXCEEDED); } } - JsonNode jobStartTimePath = inputParameter.path("jobStartTime"); - if (jobStartTimePath.isTextual()) { - LOG.warn("jobStartTime is ignored in InitialLimitedDurationErrorInjectionPolicy because lazy initialization is preferred."); - } } } @@ -139,28 +125,25 @@ public InitialLimitedDurationErrorInjectionPolicy(JsonNode inputParameter, Clock */ @Override public boolean shouldInjectionError() { - SharedState state = sharedStates.computeIfAbsent(this.instanceId, k -> new SharedState()); - - if (state.startTime == null) { - synchronized (state) { - if (state.startTime == null) { - state.startTime = Instant.now(clock); + if (startTime == null) { + synchronized (InitialLimitedDurationErrorInjectionPolicy.class) { + if (startTime == null) { + startTime = Instant.now(clock); LOG.info( - "First call detected for instance {}. Errors will be injected for {} starting from {}.", - this.instanceId, + "First call detected. Errors will be injected for {} starting from {}.", this.injectionDuration, - state.startTime); + startTime); } } } - long currentCallCount = state.callCount.incrementAndGet(); + long currentCallCount = callCount.incrementAndGet(); if (currentCallCount < INITIAL_ALLOWED_CALLS_COUNT) { return false; } Instant now = Instant.now(clock); - Duration elapsed = Duration.between(state.startTime, now); + Duration elapsed = Duration.between(startTime, now); // Compare elapsed time with the configured duration. // elapsed.compareTo(injectionDuration) < 0 means elapsed < injectionDuration @@ -195,20 +178,23 @@ public String getEffectiveDurationParameter() { } public Instant getStartTime() { - SharedState state = sharedStates.get(this.instanceId); - return state != null ? state.startTime : null; + return startTime; } void setClockForTesting(Clock clock) { this.clock = clock; } + public static void resetForTesting() { + startTime = null; + callCount.set(0); + } + @Override public String toString() { return "InitialLimitedDurationErrorInjectionPolicy{" - + "instanceId='" - + instanceId - + '\'' + + "startTime=" + + startTime + ", injectionDuration=" + injectionDuration + ", effectiveDurationParameter='" diff --git a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java index dced238320..4ce0a2dde1 100644 --- a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java +++ b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java @@ -28,10 +28,15 @@ import java.time.Duration; import java.time.Instant; import java.time.ZoneOffset; +import org.junit.Before; import org.junit.Test; public class InitialLimitedDurationErrorInjectionPolicyTest { + @Before + public void setUp() { + InitialLimitedDurationErrorInjectionPolicy.resetForTesting(); + } private ObjectNode createInputObject(String duration) { ObjectNode node = JsonNodeFactory.instance.objectNode(); From e75b5cddbd9827e368d8857f243c640d5e09bacb Mon Sep 17 00:00:00 2001 From: aasthabharill Date: Mon, 24 Aug 2026 12:57:20 +0530 Subject: [PATCH 6/7] spotless --- ...mitedDurationErrorInjectionPolicyTest.java | 1 - ...TransactionTimeoutInjectionPolicyTest.java | 42 +++++++++++-------- 2 files changed, 25 insertions(+), 18 deletions(-) diff --git a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java index 4ce0a2dde1..407ce12f3d 100644 --- a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java +++ b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/InitialLimitedDurationErrorInjectionPolicyTest.java @@ -295,6 +295,5 @@ public void constructor_shouldUseDefaultErrorCodeIfBlank() { Clock clock = Clock.fixed(Instant.EPOCH, ZoneOffset.UTC); InitialLimitedDurationErrorInjectionPolicy policy = new InitialLimitedDurationErrorInjectionPolicy(input, clock); - } } diff --git a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java index e2ddca151a..0fc86b4e32 100644 --- a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java +++ b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java @@ -271,24 +271,27 @@ public void shouldInjectError_alwaysReturnsFalse() { public void shouldInjectError_concurrentCallsCoverDoubleCheckedLocking() throws Exception { ObjectNode input = JsonNodeFactory.instance.objectNode(); input.put("transactionTimeoutBakeDuration", "PT5M"); - TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); + TransactionTimeoutInjectionPolicy policy = + new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); - Thread waitingThread = new Thread(() -> { - policy.shouldInjectionError(); - }); + Thread waitingThread = + new Thread( + () -> { + policy.shouldInjectionError(); + }); synchronized (policy) { waitingThread.start(); Thread.sleep(500); - java.lang.reflect.Field startTimeField = + java.lang.reflect.Field startTimeField = TransactionTimeoutInjectionPolicy.class.getDeclaredField("startTime"); startTimeField.setAccessible(true); startTimeField.set(policy, Instant.now()); } waitingThread.join(); - java.lang.reflect.Field startTimeField = + java.lang.reflect.Field startTimeField = TransactionTimeoutInjectionPolicy.class.getDeclaredField("startTime"); startTimeField.setAccessible(true); assertThat(startTimeField.get(policy)).isNotNull(); @@ -322,10 +325,11 @@ public void constructor_shouldThrowExceptionForInvalidJobStartTime() { public void shouldInjectionError_initializesClockIfNull() throws Exception { ObjectNode input = JsonNodeFactory.instance.objectNode(); TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input); - java.lang.reflect.Field clockField = TransactionTimeoutInjectionPolicy.class.getDeclaredField("clock"); + java.lang.reflect.Field clockField = + TransactionTimeoutInjectionPolicy.class.getDeclaredField("clock"); clockField.setAccessible(true); clockField.set(policy, null); - + assertThat(policy.shouldInjectionError()).isFalse(); } @@ -333,19 +337,23 @@ public void shouldInjectionError_initializesClockIfNull() throws Exception { public void shouldInjectDelay_handlesInterruptedException() throws Exception { ObjectNode input = JsonNodeFactory.instance.objectNode(); input.put("transactionTimeoutBakeDuration", "PT5M"); - input.put("transactionDelayDuration", "PT5M"); - TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); - - policy.setRandomForTesting(new Random() { - @Override - public double nextDouble() { return 0.1; } - }); + input.put("transactionDelayDuration", "PT5M"); + TransactionTimeoutInjectionPolicy policy = + new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); + + policy.setRandomForTesting( + new Random() { + @Override + public double nextDouble() { + return 0.1; + } + }); Thread thread = new Thread(() -> policy.shouldInjectionError()); thread.start(); - Thread.sleep(200); + Thread.sleep(200); thread.interrupt(); - + thread.join(1000); assertThat(thread.isAlive()).isFalse(); } From 9cba4df394d81ab1ba9909f3910536cdd3f22c1c Mon Sep 17 00:00:00 2001 From: aasthabharill Date: Mon, 24 Aug 2026 13:19:39 +0530 Subject: [PATCH 7/7] codecov --- ...TransactionTimeoutInjectionPolicyTest.java | 91 ------------------- v2/gcs-spanner-dv/pom.xml | 1 + 2 files changed, 1 insertion(+), 91 deletions(-) diff --git a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java index 0fc86b4e32..df494ff550 100644 --- a/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java +++ b/v2/failure-injection-policies/src/test/java/com/google/cloud/teleport/v2/failureinjection/TransactionTimeoutInjectionPolicyTest.java @@ -266,95 +266,4 @@ public void shouldInjectError_alwaysReturnsFalse() { // The method should always return false, as its purpose is to delay, not to inject an error. assertThat(policy.shouldInjectionError()).isFalse(); } - - @Test - public void shouldInjectError_concurrentCallsCoverDoubleCheckedLocking() throws Exception { - ObjectNode input = JsonNodeFactory.instance.objectNode(); - input.put("transactionTimeoutBakeDuration", "PT5M"); - TransactionTimeoutInjectionPolicy policy = - new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); - - Thread waitingThread = - new Thread( - () -> { - policy.shouldInjectionError(); - }); - - synchronized (policy) { - waitingThread.start(); - Thread.sleep(500); - - java.lang.reflect.Field startTimeField = - TransactionTimeoutInjectionPolicy.class.getDeclaredField("startTime"); - startTimeField.setAccessible(true); - startTimeField.set(policy, Instant.now()); - } - - waitingThread.join(); - java.lang.reflect.Field startTimeField = - TransactionTimeoutInjectionPolicy.class.getDeclaredField("startTime"); - startTimeField.setAccessible(true); - assertThat(startTimeField.get(policy)).isNotNull(); - } - - @Test - public void constructor_shouldParseJobStartTime() { - ObjectNode input = JsonNodeFactory.instance.objectNode(); - input.put("jobStartTime", "2025-01-01T00:00:00Z"); - Clock clock = Clock.fixed(Instant.EPOCH, ZoneOffset.UTC); - TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input, clock); - policy.setClockForTesting(Clock.fixed(Instant.parse("2025-01-01T01:00:00Z"), ZoneOffset.UTC)); - long start = System.currentTimeMillis(); - policy.shouldInjectionError(); - long end = System.currentTimeMillis(); - assertThat(end - start).isLessThan(50L); - } - - @Test - public void constructor_shouldThrowExceptionForInvalidJobStartTime() { - ObjectNode input = JsonNodeFactory.instance.objectNode(); - input.put("jobStartTime", "Invalid"); - IllegalArgumentException e = - assertThrows( - IllegalArgumentException.class, - () -> new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC())); - assertThat(e).hasMessageThat().contains("Failed to parse jobStartTime"); - } - - @Test - public void shouldInjectionError_initializesClockIfNull() throws Exception { - ObjectNode input = JsonNodeFactory.instance.objectNode(); - TransactionTimeoutInjectionPolicy policy = new TransactionTimeoutInjectionPolicy(input); - java.lang.reflect.Field clockField = - TransactionTimeoutInjectionPolicy.class.getDeclaredField("clock"); - clockField.setAccessible(true); - clockField.set(policy, null); - - assertThat(policy.shouldInjectionError()).isFalse(); - } - - @Test - public void shouldInjectDelay_handlesInterruptedException() throws Exception { - ObjectNode input = JsonNodeFactory.instance.objectNode(); - input.put("transactionTimeoutBakeDuration", "PT5M"); - input.put("transactionDelayDuration", "PT5M"); - TransactionTimeoutInjectionPolicy policy = - new TransactionTimeoutInjectionPolicy(input, Clock.systemUTC()); - - policy.setRandomForTesting( - new Random() { - @Override - public double nextDouble() { - return 0.1; - } - }); - - Thread thread = new Thread(() -> policy.shouldInjectionError()); - thread.start(); - Thread.sleep(200); - thread.interrupt(); - - thread.join(1000); - assertThat(thread.isAlive()).isFalse(); - } } diff --git a/v2/gcs-spanner-dv/pom.xml b/v2/gcs-spanner-dv/pom.xml index 29e42eb657..881f096a73 100644 --- a/v2/gcs-spanner-dv/pom.xml +++ b/v2/gcs-spanner-dv/pom.xml @@ -89,6 +89,7 @@ parent POM(s). --> com/google/cloud/teleport/v2/dto/** com/google/cloud/teleport/v2/constants/** + com/google/cloud/teleport/v2/templates/GCSSpannerDV.class