Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -326,6 +326,9 @@ class HealthCheckMetrics {
@GuardedBy("lock")
private long windowedQueuedRetriesMax;

@GuardedBy("lock")
private long windowedInflightBytesMax;

private long windowedConnectionAttemptCount;
private long windowedConnectionClosedCount;

Expand All @@ -338,6 +341,12 @@ void updateWindowedQueuedRequestsMax(long currentQueueLength, long currentRetryC
}
}

void updateInflightBytesMax(long currentInflightBytes) {
if (currentInflightBytes > windowedInflightBytesMax) {
windowedInflightBytesMax = currentInflightBytes;
}
}
Comment on lines +344 to +348

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.

high

The field windowedInflightBytesMax is annotated with @GuardedBy("lock"), but updateInflightBytesMax accesses and modifies it without acquiring the lock. This can lead to race conditions and thread-safety issues. In performance-sensitive code, prefer using explicit locks over the 'synchronized' keyword to protect shared state while ensuring thread safety and visibility.

    void updateInflightBytesMax(long currentInflightBytes) {
      lock.lock();
      try {
        if (currentInflightBytes > windowedInflightBytesMax) {
          windowedInflightBytesMax = currentInflightBytes;
        }
      } finally {
        lock.unlock();
      }
    }
References
  1. In performance-sensitive code, prefer using explicit locks over the 'synchronized' keyword to protect shared state while ensuring thread safety and visibility.


void updateResponseWait(Instant sendInstant) {
long currentWaitTime = Duration.between(sendInstant, Instant.now()).toMillis();
if (currentWaitTime > windowedMilliResponseWaitTimeMax) {
Expand Down Expand Up @@ -380,7 +389,7 @@ class HealthCheckFields {
long responseCount;
long queuedRequestCountMax;
long queuedRetryCountMax; // How many active waiting or inflight requests are retries
long inflightBytes;
long inflightBytesMax;
long connectionAttemptCount;
long connectionClosedCount;
boolean isConnected; // snapshot at instant metrics are gathered
Expand All @@ -399,7 +408,7 @@ private void gatherHealthCheckMetrics(HealthCheckFields healthCheckFields) {
healthCheckFields.queuedRequestCountMax = windowedQueuedRequestsMax;
healthCheckFields.queuedRetryCountMax = windowedQueuedRetriesMax;
healthCheckFields.msecLongestResponseWaitTime = windowedMilliResponseWaitTimeMax;
healthCheckFields.inflightBytes = inflightBytes.get();
healthCheckFields.inflightBytesMax = windowedInflightBytesMax;
healthCheckFields.requestsSentCount = windowedRequestsSent;
healthCheckFields.responseCount = windowedResponsesAcked;
if (HEALTH_CHECK_INTERVAL.toMillis() > 0) {
Expand All @@ -425,7 +434,7 @@ private void gatherHealthCheckMetrics(HealthCheckFields healthCheckFields) {
*/
private boolean checkThresholds(HealthCheckFields healthCheckFields) {
if ((healthCheckFields.queuedRequestCountMax >= queuedRequestsThreshold)
|| (healthCheckFields.inflightBytes >= queuedBytesThreshold)
|| (healthCheckFields.inflightBytesMax >= queuedBytesThreshold)
|| (healthCheckFields.msecLongestResponseWaitTime >= responseWaitTimeThreshold.toMillis())
|| (healthCheckFields.msecMaxLatency >= latencyThreshold.toMillis())
|| (healthCheckFields.connectionAttemptCount >= connectionAttemptThreshold)
Expand Down Expand Up @@ -479,6 +488,7 @@ private void resetWindowedMetrics() {
windowedConnectionClosedCount = 0;
windowedQueuedRequestsMax = 0;
windowedQueuedRetriesMax = 0;
windowedInflightBytesMax = 0;
}

/*
Expand Down Expand Up @@ -785,6 +795,7 @@ private void addMessageToWaitingQueue(
AppendRequestAndResponse requestWrapper, boolean addToFront) {
this.inflightRequests.incrementAndGet();
this.inflightBytes.addAndGet(requestWrapper.messageSize);
healthCheckMetrics.updateInflightBytesMax(this.inflightBytes);

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.

critical

The method updateInflightBytesMax expects a long primitive argument, but this.inflightBytes is an AtomicLong object. Passing it directly will cause a compilation error. Use this.inflightBytes.get() to retrieve the primitive value.

Suggested change
healthCheckMetrics.updateInflightBytesMax(this.inflightBytes);
healthCheckMetrics.updateInflightBytesMax(this.inflightBytes.get());

hasMessageInWaitingQueue.signal();
requestProfilerHook.startOperation(
RequestProfiler.OperationName.WAIT_QUEUE, requestWrapper.requestUniqueId);
Expand Down Expand Up @@ -932,6 +943,7 @@ private ApiFuture<AppendRowsResponse> appendInternal(
requestProfilerHook.startOperation(RequestProfiler.OperationName.WAIT_QUEUE, requestUniqueId);
this.inflightRequests.incrementAndGet();
this.inflightBytes.addAndGet(requestWrapper.messageSize);
healthCheckMetrics.updateInflightBytesMax(this.inflightBytes);

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.

critical

The method updateInflightBytesMax expects a long primitive argument, but this.inflightBytes is an AtomicLong object. Passing it directly will cause a compilation error. Use this.inflightBytes.get() to retrieve the primitive value.

Suggested change
healthCheckMetrics.updateInflightBytesMax(this.inflightBytes);
healthCheckMetrics.updateInflightBytesMax(this.inflightBytes.get());

requestWrapper.placedInWaitingQueueTime = Instant.now();
waitingRequestQueue.addLast(requestWrapper);
healthCheckMetrics.updateWindowedQueuedRequestsMax(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1296,7 +1296,7 @@ void testHealthCheck() throws Exception {
healthCheckFields.msecLongestResponseWaitTime > 1
&& healthCheckFields.msecLongestResponseWaitTime < msecResponseDelay);
assertEquals(appendCount, healthCheckFields.queuedRequestCountMax);
assertEquals(appendCount * sizePerRequest, healthCheckFields.inflightBytes);
assertEquals(appendCount * sizePerRequest, healthCheckFields.inflightBytesMax);
assertEquals("MULTIPLEXING", healthCheckFields.streamName);
assertEquals(connectionWorker.getWriterId(), healthCheckFields.writerId);

Expand Down Expand Up @@ -1363,7 +1363,7 @@ void testHealthCheckThresholds() throws Exception {
fields.responseCount = 0;
fields.queuedRequestCountMax = 0;
fields.queuedRetryCountMax = 0;
fields.inflightBytes = 0;
fields.inflightBytesMax = 0;
fields.connectionAttemptCount = 0;
fields.connectionClosedCount = 0;
fields.isConnected = false;
Expand Down Expand Up @@ -1415,9 +1415,9 @@ void testHealthCheckThresholds() throws Exception {
fields.queuedRetryCountMax = 0;

// inflightBytes >= queuedBytesThreshold (52428800)
fields.inflightBytes = 50 * 1024 * 1024;
fields.inflightBytesMax = 50 * 1024 * 1024;
assertEquals(true, connectionWorker.checkTestOnlyHealthCheckThresholds(fields));
fields.inflightBytes = 0;
fields.inflightBytesMax = 0;

// connectionAttemptCount >= connectionAttemptThreshold (1)
fields.connectionAttemptCount = 1;
Expand Down
Loading