Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
3 changes: 1 addition & 2 deletions internal/storage/storage_handle.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import (
"os"
"strconv"
"strings"

"time"

"cloud.google.com/go/storage"
Expand Down Expand Up @@ -416,8 +417,6 @@ func NewStorageHandle(ctx context.Context, clientConfig storageutil.StorageClien
// Wrap the control client with retry-on-stall logic.
// This will retry on only on GetStorageLayout call for all buckets.
controlClient = withRetryOnStorageLayout(controlClientWithBillingProject, &clientConfig)
} else {
logger.Infof("Skipping storage control client creation because custom-endpoint %q was passed, which is assumed to be a storage testbench server because of 'localhost' in it.", clientConfig.CustomEndpoint)
}

sh = &storageClient{
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
proxyType: grpc
targetHost: "localhost:8888"
retryConfig:
- method: CreateFolder
retryInstruction: "stall-for-40s"
retryCount: 1
skipCount: 2

Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
proxyType: grpc
targetHost: "localhost:8888"
retryConfig:
- method: DeleteFolder
retryInstruction: "stall-for-40s"
retryCount: 1
skipCount: 0

Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
proxyType: grpc
targetHost: "localhost:8888"
retryConfig:
- method: GetFolder
retryInstruction: "stall-for-40s"
retryCount: 1
skipCount: 0

Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
proxyType: grpc
targetHost: "localhost:8888"
retryConfig:
- method: GetStorageLayout
retryInstruction: "stall-for-60s"
retryCount: 1
skipCount: 0

Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
proxyType: grpc
targetHost: "localhost:8888"
retryConfig:
- method: RenameFolder
retryInstruction: "stall-for-40s"
retryCount: 1
skipCount: 0

Original file line number Diff line number Diff line change
@@ -0,0 +1,248 @@
// Copyright 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 control_client_stall

import (
"fmt"
"log"
"os"
"path"
"testing"
"time"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/stretchr/testify/suite"

emulator_tests "github.com/googlecloudplatform/gcsfuse/v3/tools/integration_tests/emulator_tests/util"
"github.com/googlecloudplatform/gcsfuse/v3/tools/integration_tests/util/mounting/static_mounting"
"github.com/googlecloudplatform/gcsfuse/v3/tools/integration_tests/util/setup"
)

type controlClientStallBase struct {
port int
proxyProcessId int
proxyServerLogFile string
flags []string
configFileName string
maxMountDurationSecs int
testDirPath string
suite.Suite
}

func (c *controlClientStallBase) SetupTest() {
c.testDirPath = setup.SetupTestDirectory(c.T().Name())
c.proxyServerLogFile = setup.CreateProxyServerLogFile(c.T())
var err error
c.port, c.proxyProcessId, err = emulator_tests.StartProxyServer(c.configFileName, c.proxyServerLogFile)
require.NoError(c.T(), err)

// Add proxy server endpoint and configure gRPC testing
// Copy flags to avoid mutating the original slice across suites
mountFlags := append([]string(nil), c.flags...)
mountFlags = append(mountFlags, fmt.Sprintf("--custom-endpoint=localhost:%d", c.port))
mountFlags = append(mountFlags, "--anonymous-access") // Required for gRPC localhost endpoint

startTime := time.Now()
setup.MountGCSFuseWithGivenMountFunc(mountFlags, static_mounting.MountGcsfuseWithStaticMounting)
mountTime := time.Since(startTime)
if c.maxMountDurationSecs != -1 {
assert.True(c.T(), mountTime < time.Duration(c.maxMountDurationSecs)*time.Second, "Mount time %v should be less than %ds", mountTime, c.maxMountDurationSecs)
}
}

func (c *controlClientStallBase) TearDownTest() {
setup.UnmountGCSFuse(setup.MntDir())
assert.NoError(c.T(), emulator_tests.KillProxyServerProcess(c.proxyProcessId))
Comment thread
raj-prince marked this conversation as resolved.
Outdated
setup.SaveGCSFuseLogFileInCaseOfFailure(c.T())
setup.SaveProxyServerLogFileInCaseOfFailure(c.proxyServerLogFile, c.T())
}

// --- Test Suites ---

type createFolderStallSuite struct{ controlClientStallBase }

// TestCreateFolderStallInducedShouldComplete verifies that creating a folder via
// os.MkdirAll completes successfully even when a stall is induced,
// proving that the experimental retry logic correctly handles the delayed gRPC response.
func (c *createFolderStallSuite) TestCreateFolderStallInducedShouldComplete() {
folderPath := path.Join(c.testDirPath, "stalled_folder_for_create")

startTime := time.Now()
err := os.MkdirAll(folderPath, 0755)
elapsedTime := time.Since(startTime)

assert.NoError(c.T(), err)
// Ensure the CreateFolder operation took less than 32 seconds (30 seconds first stalled-call timeout + 2 second wiggle-room) because of the internal abortion and retry.
assert.True(c.T(), elapsedTime < 32*time.Second, "Elapsed time %v should be less than 32s", elapsedTime)
// Validate the directory actually exists
info, err := os.Stat(folderPath)
assert.NoError(c.T(), err)
assert.True(c.T(), info.IsDir())
}

type getFolderStallSuite struct{ controlClientStallBase }

// TestGetFolderStallInducedShouldComplete verifies that stating a pre-created folder via
// os.Stat completes successfully even when a stall is induced,
// proving that the experimental retry logic correctly handles the delayed gRPC response.
func (c *getFolderStallSuite) TestGetFolderStallInducedShouldComplete() {
folderPath := path.Join(c.testDirPath, "stalled_folder_for_get")
err := os.MkdirAll(folderPath, 0755)
assert.NoError(c.T(), err)

startTime := time.Now()
_, err = os.Stat(folderPath)
elapsedTime := time.Since(startTime)

assert.NoError(c.T(), err)
// Ensure the GetFolder operation took less than 32 seconds (30 seconds first stalled-call timeout + 2 second wiggle-room) because of the internal abortion and retry.
assert.True(c.T(), elapsedTime < 32*time.Second, "Elapsed time %v should be less than 32s", elapsedTime)
}

type deleteFolderStallSuite struct{ controlClientStallBase }

// TestDeleteFolderStallInducedShouldComplete verifies that deleting a pre-created empty folder via
// os.Remove completes successfully even when a stall is induced,
// proving that the experimental retry logic correctly handles the delayed gRPC response.
func (c *deleteFolderStallSuite) TestDeleteFolderStallInducedShouldComplete() {
folderPath := path.Join(c.testDirPath, "stalled_folder_for_delete")
err := os.MkdirAll(folderPath, 0755)
assert.NoError(c.T(), err)

startTime := time.Now()
err = os.Remove(folderPath)
elapsedTime := time.Since(startTime)

assert.NoError(c.T(), err)
// Ensure the DeleteFolder operation took less than 32 seconds (30 seconds first stalled-call timeout + 2 second wiggle-room) because of the internal abortion and retry.
assert.True(c.T(), elapsedTime < 32*time.Second, "Elapsed time %v should be less than 32s", elapsedTime)
// Validate the directory actually deleted
_, err = os.Stat(folderPath)
assert.True(c.T(), os.IsNotExist(err))
}

type renameFolderStallSuite struct{ controlClientStallBase }

// TestRenameFolderStallInducedShouldComplete verifies that renaming a pre-created folder via
// os.Rename completes successfully even when a stall is induced,
// proving that the experimental retry logic correctly handles the delayed gRPC response.
func (c *renameFolderStallSuite) TestRenameFolderStallInducedShouldComplete() {
folderPath := path.Join(c.testDirPath, "stalled_folder_for_rename")
err := os.MkdirAll(folderPath, 0755)
assert.NoError(c.T(), err)
destPath := path.Join(c.testDirPath, "stalled_folder_renamed")

startTime := time.Now()
err = os.Rename(folderPath, destPath)
elapsedTime := time.Since(startTime)

assert.NoError(c.T(), err)
// Ensure the RenameFolder operation took less than 32 seconds (30 seconds first stalled-call timeout + 2 second wiggle-room) because of the internal abortion and retry.
assert.True(c.T(), elapsedTime < 32*time.Second, "Elapsed time %v should be less than 32s", elapsedTime)
// Validate the directory actually renamed
_, err = os.Stat(destPath)
assert.NoError(c.T(), err)
}

type getStorageLayoutStallSuite struct{ controlClientStallBase }

// TestGetStorageLayoutStallInducedShouldComplete verifies that a GCSFuse mount
// completes successfully even when a stall is induced in the very first GetStorageLayout call,
// proving that the experimental retry logic correctly handles the delayed gRPC response.
func (c *getStorageLayoutStallSuite) TestGetStorageLayoutStallInducedShouldComplete() {
// The stall for GetStorageLayout is actually triggered during the bucket mount
// phase within SetupTest() when GCSFuse queries the bucket layout to determine
// if Hierarchical Namespace (HNS) is enabled.
// The fact that this test function is executing means the mount succeeded,
// effectively proving that GetStorageLayout correctly timed out after 30s,
// retried successfully, and allowed the mount to complete!

// Just perform a basic verification to ensure the mount is functional.
folderPath := path.Join(c.testDirPath, "stalled_folder_for_layout")
err := os.MkdirAll(folderPath, 0755)
assert.NoError(c.T(), err)
// Validate the directory actually exists
info, err := os.Stat(folderPath)
assert.NoError(c.T(), err)
assert.True(c.T(), info.IsDir())
}

func TestControlClientStall(t *testing.T) {
// Test matrix
flagsSet := [][]string{
{
"--client-protocol=grpc",
"--rename-dir-limit=0", // Force use of Control API for folders instead of legacy workarounds
"--experimental-nonrapid-folder-api-stall-retry=true", // Enable the new flag we are testing!
"--metadata-cache-ttl-secs=0", // Disable cache so calls definitely hit the backend
},
}

for _, flags := range flagsSet {
suites := []suite.TestingSuite{
&createFolderStallSuite{
controlClientStallBase: controlClientStallBase{
flags: flags,
// Artificially stall the very first CreateFolder call by 40 seconds.
configFileName: "../configs/control_client_stall_create_40s.yaml",
// No check on how long mount takes.
maxMountDurationSecs: -1,
},
},
&getFolderStallSuite{
controlClientStallBase: controlClientStallBase{
flags: flags,
// Artificially stall the very first GetFolder call by 40 seconds.
configFileName: "../configs/control_client_stall_get_40s.yaml",
// No check on how long mount takes.
maxMountDurationSecs: -1,
},
},
&deleteFolderStallSuite{
controlClientStallBase: controlClientStallBase{
flags: flags,
// Artificially stall the very first DeleteFolder call by 40 seconds.
configFileName: "../configs/control_client_stall_delete_40s.yaml",
// No check on how long mount takes.
maxMountDurationSecs: -1,
},
},
&renameFolderStallSuite{
controlClientStallBase: controlClientStallBase{
flags: flags,
// Artificially stall the very first RenameFolder call by 40 seconds.
configFileName: "../configs/control_client_stall_rename_40s.yaml",
// No check on how long mount takes.
maxMountDurationSecs: -1,
},
},
&getStorageLayoutStallSuite{
controlClientStallBase: controlClientStallBase{
flags: flags,
// Artificially stall the very first GetStorageLayout call by 60 seconds.
configFileName: "../configs/control_client_stall_layout_60s.yaml",
// Ensure that the whole mount operation including a stalled GetStorageLayout call took less than 40 seconds (30 seconds first stalled-call GetStorageLayout timeout + 10 second wiggle-room for 2nd attempt and other steps for mount to go through) because of the internal abortion and retry.
maxMountDurationSecs: 40,
},
},
}

for _, s := range suites {
log.Printf("Running Control Client Stall test suite %T with flags: %s", s, flags)
suite.Run(t, s)
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
// Copyright 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 control_client_stall

import (
"log"
"os"
"testing"

"github.com/googlecloudplatform/gcsfuse/v3/tools/integration_tests/util/setup"
)

func TestMain(m *testing.M) {
setup.ParseSetUpFlags()

if setup.MountedDirectory() != "" {
log.Printf("These tests will not run with mounted directory.")
return
}

// Set up test directory.
setup.SetUpTestDirForTestBucketFlag()

log.Println("Running control client stall tests...")
successCode := m.Run()
os.Exit(successCode)
}
24 changes: 21 additions & 3 deletions tools/integration_tests/emulator_tests/emulator_tests.sh
Original file line number Diff line number Diff line change
Expand Up @@ -118,8 +118,9 @@ sudo docker pull $DOCKER_IMAGE
CONTAINER_ID=$(sudo docker ps -aqf "name=$CONTAINER_NAME")
if [[ -n "$CONTAINER_ID" ]]; then
log_info "Container with ID:[$CONTAINER_ID] is already running with name:[$CONTAINER_NAME]"
log_info "Stopping container...."
sudo docker stop $CONTAINER_ID
log_info "Stopping and removing container...."
docker stop $CONTAINER_ID || true
docker rm $CONTAINER_ID || true
Comment thread
raj-prince marked this conversation as resolved.
Outdated
fi

wait_for_emulator() {
Expand Down Expand Up @@ -175,6 +176,17 @@ if ! curl -X POST --data-binary @test.json \
fi
rm test.json

# Create an HNS bucket for control client tests
cat << EOF > test_hns.json
{"name":"test-hns-bucket", "hierarchicalNamespace": {"enabled": true}}
EOF
if ! curl -X POST --data-binary @test_hns.json -H "Content-Type: application/json" "$STORAGE_EMULATOR_HOST/storage/v1/b?project=test-project"; then
log_error "Failed to create bucket test-hns-bucket"
exit 1
fi
rm test_hns.json


Comment on lines +179 to +189

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Nit: can we make move this code to a method and use that in both places here and just above it?

# Start the gRPC server on port 8888.
log_info "Starting the gRPC server on port 8888"
response=$(curl -w "%{http_code}\n" --retry 5 --retry-max-time 40 -o /dev/null "$STORAGE_EMULATOR_HOST/start_grpc?port=8888")
Expand All @@ -192,4 +204,10 @@ if [[ -n "$GCSFUSE_PREBUILT_DIR" ]]; then
fi

# Run all emulator test packages in sequence to avoid high cpu usage.
go test -v -p 1 -timeout 10m ./tools/integration_tests/emulator_tests/... --integrationTest --testbucket=test-bucket "${args[@]}"
TEST_TARGET=${TEST_TARGET:-"./tools/integration_tests/emulator_tests/..."}
# Run all emulator test packages in sequence to avoid high cpu usage.
# Run all other emulator tests with standard bucket
go test -v -p 1 -timeout 10m $(go list ${TEST_TARGET} | grep -v control_client_stall) --integrationTest --testbucket=${TEST_BUCKET:-test-bucket} "${args[@]}"
Comment on lines +207 to +210

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.

medium

Under set -euo pipefail, if TEST_TARGET is set to only run the control_client_stall tests (e.g., ./tools/integration_tests/emulator_tests/control_client_stall/...), grep -v control_client_stall will find no matches and exit with status 1. This will cause the entire script to exit immediately with an error. We should handle this gracefully by using || true and checking if the package list is non-empty before running go test.

Suggested change
TEST_TARGET=${TEST_TARGET:-"./tools/integration_tests/emulator_tests/..."}
# Run all emulator test packages in sequence to avoid high cpu usage.
# Run all other emulator tests with standard bucket
go test -v -p 1 -timeout 10m $(go list ${TEST_TARGET} | grep -v control_client_stall) --integrationTest --testbucket=${TEST_BUCKET:-test-bucket} "${args[@]}"
TEST_TARGET=${TEST_TARGET:-"./tools/integration_tests/emulator_tests/..."}
# Run all emulator test packages in sequence to avoid high cpu usage.
# Run all other emulator tests with standard bucket
TEST_PKGS=$(go list ${TEST_TARGET} | grep -v control_client_stall || true)
if [[ -n "$TEST_PKGS" ]]; then
go test -v -p 1 -timeout 10m $TEST_PKGS --integrationTest --testbucket=${TEST_BUCKET:-test-bucket} "${args[@]}"
fi


# Run control_client_stall with HNS bucket
go test -v -p 1 -timeout 10m ./tools/integration_tests/emulator_tests/control_client_stall/... --integrationTest --testbucket=test-hns-bucket "${args[@]}"
Loading
Loading