2727import java .util .List ;
2828import java .util .Map ;
2929import java .util .concurrent .ConcurrentHashMap ;
30+ import java .util .concurrent .CopyOnWriteArrayList ;
3031import java .util .concurrent .LinkedBlockingQueue ;
3132import java .util .concurrent .ScheduledExecutorService ;
3233import java .util .concurrent .Semaphore ;
3334import java .util .concurrent .TimeUnit ;
3435import java .util .concurrent .atomic .AtomicInteger ;
36+ import java .util .concurrent .atomic .AtomicLong ;
3537import java .util .function .Supplier ;
3638import java .util .logging .Logger ;
3739
@@ -59,11 +61,11 @@ class FakeBigQueryWriteImpl extends BigQueryWriteGrpc.BigQueryWriteImplBase {
5961
6062 private long numberTimesToClose = 0 ;
6163 private long closeAfter = 0 ;
62- private long recordCount = 0 ;
63- private long connectionCount = 0 ;
64+ private final AtomicLong recordCount = new AtomicLong ( 0 ) ;
65+ private final AtomicLong connectionCount = new AtomicLong ( 0 ) ;
6466 private long closeForeverAfter = 0 ;
65- private int responseIndex = 0 ;
66- private long expectedOffset = 0 ;
67+ private final AtomicInteger responseIndex = new AtomicInteger ( 0 ) ;
68+ private final AtomicLong expectedOffset = new AtomicLong ( 0 ) ;
6769 private boolean verifyOffset = false ;
6870 private boolean returnErrorDuringExclusiveStreamRetry = false ;
6971 private boolean returnErrorUntilRetrySuccess = false ;
@@ -74,7 +76,8 @@ class FakeBigQueryWriteImpl extends BigQueryWriteGrpc.BigQueryWriteImplBase {
7476 private final Map <StreamObserver <AppendRowsResponse >, Boolean > connectionToFirstRequest =
7577 new ConcurrentHashMap <>();
7678 private Status failedStatus = Status .ABORTED ;
77- private ArrayList <Instant > requestReceivedInstants = new ArrayList <>();
79+ private final CopyOnWriteArrayList <Instant > requestReceivedInstants =
80+ new CopyOnWriteArrayList <>();
7881
7982 /** Class used to save the state of a possible response. */
8083 public static class Response {
@@ -114,7 +117,7 @@ public String toString() {
114117 }
115118
116119 public ArrayList <Instant > getLatestRequestReceivedInstants () {
117- return requestReceivedInstants ;
120+ return new ArrayList <>( requestReceivedInstants ) ;
118121 }
119122
120123 @ Override
@@ -153,7 +156,7 @@ void waitForResponseScheduled() throws InterruptedException {
153156
154157 /* Return the number of times the stream was connected. */
155158 public long getConnectionCount () {
156- return connectionCount ;
159+ return connectionCount . get () ;
157160 }
158161
159162 void setFailedStatus (Status failedStatus ) {
@@ -197,20 +200,23 @@ private Response determineResponse(long offset) {
197200 @ Override
198201 public StreamObserver <AppendRowsRequest > appendRows (
199202 final StreamObserver <AppendRowsResponse > responseObserver ) {
200- this . connectionCount ++ ;
203+ connectionCount . incrementAndGet () ;
201204 connectionToFirstRequest .put (responseObserver , true );
202205 StreamObserver <AppendRowsRequest > requestObserver =
203206 new StreamObserver <AppendRowsRequest >() {
204207 @ Override
205208 public void onNext (AppendRowsRequest value ) {
206209 requestReceivedInstants .add (Instant .now ());
207- recordCount ++;
210+ long currentRecordCount = recordCount .incrementAndGet ();
211+ int currentResponseIndex = responseIndex .getAndIncrement ();
212+ int responseIndexAfterIncrement = currentResponseIndex + 1 ;
213+ long currentConnectionCount = connectionCount .get ();
214+
208215 requests .add (value );
209216 long offset = value .getOffset ().getValue ();
210217 if (offset == -1 || !value .hasOffset ()) {
211- offset = responseIndex ;
218+ offset = currentResponseIndex ;
212219 }
213- responseIndex ++;
214220 if (responseSleep .compareTo (Duration .ZERO ) > 0 ) {
215221 LOG .info ("Sleeping before response for " + responseSleep .toString ());
216222 Uninterruptibles .sleepUninterruptibly (
@@ -233,33 +239,38 @@ public void onNext(AppendRowsRequest value) {
233239 }
234240 connectionToFirstRequest .put (responseObserver , false );
235241 if (closeAfter > 0
236- && responseIndex % closeAfter == 0
237- && recordCount % closeAfter == 0
238- && (numberTimesToClose == 0 || connectionCount <= numberTimesToClose )) {
242+ && responseIndexAfterIncrement % closeAfter == 0
243+ && currentRecordCount % closeAfter == 0
244+ && (numberTimesToClose == 0 || currentConnectionCount <= numberTimesToClose )) {
239245 LOG .info ("Shutting down connection from test..." );
240246 responseObserver .onError (failedStatus .asException ());
241- } else if (closeForeverAfter > 0 && recordCount > closeForeverAfter ) {
247+ } else if (closeForeverAfter > 0 && currentRecordCount > closeForeverAfter ) {
242248 LOG .info ("Shutting down connection from test..." );
243249 responseObserver .onError (failedStatus .asException ());
244250 } else {
245251 Response response = determineResponse (offset );
246252 if (verifyOffset
247253 && !response .getResponse ().hasError ()
248254 && response .getResponse ().getAppendResult ().getOffset ().getValue () > -1 ) {
249- // No error and offset is present; verify order
250- if (response .getResponse ().getAppendResult ().getOffset ().getValue ()
251- != expectedOffset ) {
255+ long responseOffset =
256+ response .getResponse ().getAppendResult ().getOffset ().getValue ();
257+ // Atomically verify that the response offset matches the expected offset
258+ // and increment the expected offset for the next request. This avoids
259+ // using a synchronized block while ensuring thread safety across concurrent
260+ // streams.
261+ if (!expectedOffset .compareAndSet (responseOffset , responseOffset + 1 )) {
262+ LOG .info (
263+ String .format (
264+ "Offset mismatch: expected %s, got %s" ,
265+ expectedOffset .get (), responseOffset ));
252266 com .google .rpc .Status status =
253267 com .google .rpc .Status .newBuilder ().setCode (Code .INTERNAL_VALUE ).build ();
254268 response = new Response (AppendRowsResponse .newBuilder ().setError (status ).build ());
255269 } else {
256270 LOG .info (
257271 String .format (
258- "asserted offset: %s expected: %s" ,
259- response .getResponse ().getAppendResult ().getOffset ().getValue (),
260- expectedOffset ));
272+ "asserted offset: %s expected: %s" , responseOffset , responseOffset ));
261273 LOG .info (String .format ("sending response: %s" , response .getResponse ()));
262- expectedOffset ++;
263274 }
264275 }
265276 sendResponse (response , responseObserver );
0 commit comments