From 958ac113d3fa7112b37f1a01cea0b0c24700eede Mon Sep 17 00:00:00 2001 From: Fredrik Fornwall Date: Sat, 18 Jul 2026 13:20:43 +0200 Subject: [PATCH] fix(spanner): recover inline-begin read-only transaction after failed first statement The transactionIdFuture of a MultiUseReadOnlyTransaction with BeginTransactionOption.INLINE was created once and never reset. If the statement carrying the inline BeginTransaction failed (or its ResultSet was closed before a transaction id was returned), every subsequent statement on the transaction rethrew the stale error of that first statement forever, and concurrently waiting statements failed with an error belonging to an unrelated statement. Reset transactionIdFuture when it is failed, and retry the selector loop in getTransactionSelector, so the next (or a concurrently waiting) statement attempts a new inline BeginTransaction, matching the recovery behavior of the explicit BeginTransaction path. Signed-off-by: Fredrik Fornwall --- .../cloud/spanner/AbstractReadContext.java | 69 +++++---- .../google/cloud/spanner/SessionImplTest.java | 134 ++++++++++++++++-- 2 files changed, 166 insertions(+), 37 deletions(-) diff --git a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java index a4d7614ebd2b..4df8d54c2385 100644 --- a/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java +++ b/java-spanner/google-cloud-spanner/src/main/java/com/google/cloud/spanner/AbstractReadContext.java @@ -433,38 +433,47 @@ TransactionSelector getTransactionSelector() { return selector; } - ApiFuture futureToWaitFor = null; - txnLock.lock(); - try { - if (transactionId != null) { - return TransactionSelector.newBuilder().setId(transactionId).build(); + final long deadlineNanos = + System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(WAIT_FOR_INLINE_BEGIN_TIMEOUT_MILLIS); + while (true) { + ApiFuture futureToWaitFor; + txnLock.lock(); + try { + if (transactionId != null) { + return TransactionSelector.newBuilder().setId(transactionId).build(); + } + if (transactionIdFuture == null) { + transactionIdFuture = SettableApiFuture.create(); + return TransactionSelector.newBuilder() + .setBegin(createReadOnlyTransactionOptions()) + .build(); + } + futureToWaitFor = transactionIdFuture; + } finally { + txnLock.unlock(); } - if (transactionIdFuture == null) { - transactionIdFuture = SettableApiFuture.create(); + + try { return TransactionSelector.newBuilder() - .setBegin(createReadOnlyTransactionOptions()) + .setId(futureToWaitFor.get(deadlineNanos - System.nanoTime(), TimeUnit.NANOSECONDS)) .build(); + } catch (ExecutionException e) { + // The statement that carried the inline BeginTransaction failed before a transaction was + // returned, and failTransactionIdFuture has reset transactionIdFuture. Retry the loop so + // this statement either takes over the BeginTransaction itself, or waits for another + // statement that already has. The error is propagated to the failed statement itself and + // should not fail unrelated statements. Read-only transactions take no locks, so unlike + // the read/write inline-begin path there is no need to abort the entire transaction. + } catch (TimeoutException e) { + throw SpannerExceptionFactory.newSpannerException( + ErrorCode.DEADLINE_EXCEEDED, + "Timeout while waiting for an inlined read-only transaction to be returned by another" + + " statement.", + e); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw SpannerExceptionFactory.newSpannerExceptionForCancellation(null, e); } - futureToWaitFor = transactionIdFuture; - } finally { - txnLock.unlock(); - } - - try { - return TransactionSelector.newBuilder() - .setId(futureToWaitFor.get(WAIT_FOR_INLINE_BEGIN_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS)) - .build(); - } catch (ExecutionException e) { - throw SpannerExceptionFactory.asSpannerException(e.getCause()); - } catch (TimeoutException e) { - throw SpannerExceptionFactory.newSpannerException( - ErrorCode.DEADLINE_EXCEEDED, - "Timeout while waiting for an inlined read-only transaction to be returned by another" - + " statement.", - e); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw SpannerExceptionFactory.newSpannerExceptionForCancellation(null, e); } } @@ -624,6 +633,10 @@ private void failTransactionIdFuture(Throwable t) { try { if (transactionIdFuture != null && !transactionIdFuture.isDone()) { transactionIdFuture.setException(t); + // Reset the future so the next statement attempts a new inline BeginTransaction instead + // of failing on the error of an earlier, unrelated statement. This mirrors the explicit + // BeginTransaction path, where a failed begin is retried by the next statement. + transactionIdFuture = null; } } finally { txnLock.unlock(); diff --git a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java index d1779b459e4c..50cb2cc4a842 100644 --- a/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java +++ b/java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/SessionImplTest.java @@ -67,6 +67,7 @@ import java.util.Map; import java.util.TimeZone; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; @@ -725,18 +726,32 @@ public void multiUseReadOnlyTransactionCanUseInlineBeginForQuery() throws ParseE } @Test - public void multiUseReadOnlyTransactionInlineBeginFirstQueryErrorPropagates() { + public void multiUseReadOnlyTransactionInlineBeginRetriesBeginAfterFailedFirstStatement() + throws ParseException { SpannerException error = SpannerExceptionFactory.newSpannerException(ErrorCode.INVALID_ARGUMENT, "bad query"); - final ArgumentCaptor consumer = + final ArgumentCaptor queryConsumer = ArgumentCaptor.forClass(SpannerRpc.ResultStreamConsumer.class); - final ArgumentCaptor request = + final ArgumentCaptor queryRequest = ArgumentCaptor.forClass(ExecuteSqlRequest.class); Mockito.when( - rpc.executeQuery(request.capture(), consumer.capture(), anyMap(), any(), eq(false))) + rpc.executeQuery( + queryRequest.capture(), queryConsumer.capture(), anyMap(), any(), eq(false))) .then( invocation -> { - consumer.getValue().onError(error); + queryConsumer.getValue().onError(error); + return new NoOpStreamingCall(); + }); + PartialResultSet readResultSet = inlineBeginResultSet("recovered-tx"); + final ArgumentCaptor readConsumer = + ArgumentCaptor.forClass(SpannerRpc.ResultStreamConsumer.class); + final ArgumentCaptor readRequest = ArgumentCaptor.forClass(ReadRequest.class); + Mockito.when( + rpc.read(readRequest.capture(), readConsumer.capture(), anyMap(), any(), eq(false))) + .then( + invocation -> { + readConsumer.getValue().onPartialResultSet(readResultSet); + readConsumer.getValue().onCompleted(); return new NoOpStreamingCall(); }); @@ -745,18 +760,119 @@ public void multiUseReadOnlyTransactionInlineBeginFirstQueryErrorPropagates() { SpannerException e = assertThrows(SpannerException.class, () -> rs.next()); assertEquals(ErrorCode.INVALID_ARGUMENT, e.getErrorCode()); } + // The transaction should not be poisoned by the failed first statement: the next statement + // should attempt a new inline BeginTransaction. + txn.readRow("Dummy", Key.of(), Collections.singletonList("C")); + txn.readRow("Dummy", Key.of(), Collections.singletonList("C")); + assertEquals( + Timestamp.fromProto(Timestamps.parse("2015-10-01T10:54:20.021Z")), + txn.getReadTimestamp()); + } + + Mockito.verify(rpc, Mockito.never()).beginTransaction(Mockito.any(), anyMap(), eq(false)); + assertEquals(1, queryRequest.getAllValues().size()); + assertThat(queryRequest.getAllValues().get(0).getTransaction().hasBegin()).isTrue(); + assertEquals(2, readRequest.getAllValues().size()); + assertThat(readRequest.getAllValues().get(0).getTransaction().hasBegin()).isTrue(); + assertEquals( + ByteString.copyFromUtf8("recovered-tx"), + readRequest.getAllValues().get(1).getTransaction().getId()); + } + + @Test + public void multiUseReadOnlyTransactionInlineBeginRecoversWhenNoTransactionIsReturned() + throws ParseException { + PartialResultSet resultSetWithoutTransaction = resultSetWithoutTransaction(); + PartialResultSet resultSetWithTransaction = inlineBeginResultSet("second-attempt-tx"); + final ArgumentCaptor request = ArgumentCaptor.forClass(ReadRequest.class); + final AtomicInteger callCount = new AtomicInteger(); + Mockito.when(rpc.read(request.capture(), Mockito.any(), anyMap(), any(), eq(false))) + .then( + invocation -> { + SpannerRpc.ResultStreamConsumer consumer = invocation.getArgument(1); + consumer.onPartialResultSet( + callCount.incrementAndGet() == 1 + ? resultSetWithoutTransaction + : resultSetWithTransaction); + consumer.onCompleted(); + return new NoOpStreamingCall(); + }); + + try (ReadOnlyTransaction txn = inlineReadOnlyTransaction()) { SpannerException e = assertThrows( SpannerException.class, () -> txn.readRow("Dummy", Key.of(), Collections.singletonList("C"))); - assertEquals(ErrorCode.INVALID_ARGUMENT, e.getErrorCode()); + assertEquals(ErrorCode.FAILED_PRECONDITION, e.getErrorCode()); + // The next statement should attempt a new inline BeginTransaction instead of failing on the + // error of the first statement. + txn.readRow("Dummy", Key.of(), Collections.singletonList("C")); } Mockito.verify(rpc, Mockito.never()).beginTransaction(Mockito.any(), anyMap(), eq(false)); - assertEquals(1, request.getAllValues().size()); + assertEquals(2, request.getAllValues().size()); assertThat(request.getAllValues().get(0).getTransaction().hasBegin()).isTrue(); - Mockito.verify(rpc, Mockito.never()) - .read(Mockito.any(), Mockito.any(), anyMap(), any(), eq(false)); + assertThat(request.getAllValues().get(1).getTransaction().hasBegin()).isTrue(); + } + + @Test + public void multiUseReadOnlyTransactionInlineBeginConcurrentStatementTakesOverFailedBegin() + throws Exception { + PartialResultSet takeOverResultSet = inlineBeginResultSet("take-over-tx"); + SpannerException error = + SpannerExceptionFactory.newSpannerException(ErrorCode.INVALID_ARGUMENT, "bad read"); + final List requests = Collections.synchronizedList(new ArrayList<>()); + final List consumers = + Collections.synchronizedList(new ArrayList<>()); + final AtomicInteger callCount = new AtomicInteger(); + final CountDownLatch firstRpcStarted = new CountDownLatch(1); + Mockito.when(rpc.read(Mockito.any(), Mockito.any(), anyMap(), any(), eq(false))) + .then( + invocation -> { + int call = callCount.incrementAndGet(); + ReadRequest readRequest = invocation.getArgument(0); + SpannerRpc.ResultStreamConsumer consumer = invocation.getArgument(1); + requests.add(readRequest); + consumers.add(consumer); + if (call == 1) { + firstRpcStarted.countDown(); + } else { + consumer.onPartialResultSet(takeOverResultSet); + consumer.onCompleted(); + } + return new NoOpStreamingCall(); + }); + + ExecutorService executor = Executors.newFixedThreadPool(2); + try (ReadOnlyTransaction txn = inlineReadOnlyTransaction()) { + Future first = + executor.submit(() -> txn.readRow("Dummy", Key.of(), Collections.singletonList("C"))); + assertThat(firstRpcStarted.await(5, TimeUnit.SECONDS)).isTrue(); + + Future second = + executor.submit(() -> txn.readRow("Dummy", Key.of(), Collections.singletonList("C"))); + Thread.sleep(100L); + assertThat(callCount.get()).isEqualTo(1); + assertThat(second.isDone()).isFalse(); + + // Fail the statement that carried the inline BeginTransaction. The concurrently waiting + // statement should take over the BeginTransaction instead of failing with the error of the + // first statement. + consumers.get(0).onError(error); + + ExecutionException e = + assertThrows(ExecutionException.class, () -> first.get(5, TimeUnit.SECONDS)); + assertThat(e.getCause()).isInstanceOf(SpannerException.class); + assertEquals(ErrorCode.INVALID_ARGUMENT, ((SpannerException) e.getCause()).getErrorCode()); + assertThat(second.get(5, TimeUnit.SECONDS)).isNull(); + } finally { + executor.shutdownNow(); + } + + Mockito.verify(rpc, Mockito.never()).beginTransaction(Mockito.any(), anyMap(), eq(false)); + assertEquals(2, requests.size()); + assertThat(requests.get(0).getTransaction().hasBegin()).isTrue(); + assertThat(requests.get(1).getTransaction().hasBegin()).isTrue(); } @Test