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
365 changes: 311 additions & 54 deletions cmd/autobahn-e2e/aws.go

Large diffs are not rendered by default.

213 changes: 201 additions & 12 deletions cmd/autobahn-e2e/command_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"bytes"
"context"
"errors"
"fmt"
"os"
"path/filepath"
"strings"
Expand Down Expand Up @@ -78,6 +79,15 @@ func TestClusterNodesMatchDockerComposePorts(t *testing.T) {
}, clusterNodes(4))
}

func TestNativeClusterNodesUseOneRPCPortPerHost(t *testing.T) {
require.Equal(t, []node{
{Index: 0, Name: "node-0", EVMHostPort: 8545},
{Index: 1, Name: "node-1", EVMHostPort: 8545},
{Index: 2, Name: "node-2", EVMHostPort: 8545},
{Index: 3, Name: "node-3", EVMHostPort: 8545},
}, nativeClusterNodes(4))
}

func TestFindNodeAcceptsNamesAndIndexes(t *testing.T) {
nodes := clusterNodes(4)
for _, selector := range []string{"2", "node-2", "sei-node-2"} {
Expand All @@ -104,9 +114,16 @@ func TestAWSDeployCreatesManagedResourcesAndReadyState(t *testing.T) {
case strings.Contains(joined, "create-key-pair"):
return "-----BEGIN OPENSSH PRIVATE KEY-----\ntest\n-----END OPENSSH PRIVATE KEY-----\n", nil
case strings.Contains(joined, "run-instances"):
return "i-123\n", nil
return "i-100\ti-101\ti-102\ti-103\n", nil
case strings.Contains(joined, "describe-instances"):
return "203.0.113.10\n", nil
return `[
{"InstanceID":"i-100","PublicIP":"203.0.113.10","PrivateIP":"10.0.0.10"},
{"InstanceID":"i-101","PublicIP":"203.0.113.11","PrivateIP":"10.0.0.11"},
{"InstanceID":"i-102","PublicIP":"203.0.113.12","PrivateIP":"10.0.0.12"},
{"InstanceID":"i-103","PublicIP":"203.0.113.13","PrivateIP":"10.0.0.13"}
]`, nil
case spec.name == "ssh" && strings.Contains(joined, nativeBuildFailedFile):
return nativeBuildStatusReady, nil
case spec.name == "ssh":
return "", nil
default:
Expand All @@ -133,8 +150,13 @@ func TestAWSDeployCreatesManagedResourcesAndReadyState(t *testing.T) {
state, err := app.store().load(options.name)
require.NoError(t, err)
require.Equal(t, "ready", state.Status)
require.Equal(t, "i-123", state.AWS.InstanceID)
require.Equal(t, "203.0.113.10", state.AWS.PublicIP)
require.Equal(t, []awsInstanceState{
{NodeIndex: 0, InstanceID: "i-100", PublicIP: "203.0.113.10", PrivateIP: "10.0.0.10"},
{NodeIndex: 1, InstanceID: "i-101", PublicIP: "203.0.113.11", PrivateIP: "10.0.0.11"},
{NodeIndex: 2, InstanceID: "i-102", PublicIP: "203.0.113.12", PrivateIP: "10.0.0.12"},
{NodeIndex: 3, InstanceID: "i-103", PublicIP: "203.0.113.13", PrivateIP: "10.0.0.13"},
}, state.AWS.Instances)
require.Equal(t, nativeClusterNodes(4), state.Nodes)
require.True(t, state.AWS.ManagedKey)
require.FileExists(t, state.AWS.SSHKeyPath)
keyInfo, err := os.Stat(state.AWS.SSHKeyPath)
Expand All @@ -145,7 +167,12 @@ func TestAWSDeployCreatesManagedResourcesAndReadyState(t *testing.T) {
commands := joinedCommands(runner.commands)
require.Contains(t, commands, "authorize-security-group-ingress")
require.Contains(t, commands, "--cidr 198.51.100.4/32")
require.Contains(t, commands, "AUTOBAHN_EVMONLY_IN_MEMORY=true")
require.Contains(t, commands, "--count 4")
require.Contains(t, commands, "UserIdGroupPairs=[{GroupId=sg-123}]")
require.Contains(t, commands, "prepare_native_cluster.sh")
require.Contains(t, commands, "install_native_node.sh")
require.Contains(t, commands, "systemctl start seid.service")
require.NotContains(t, commands, "docker")
require.Contains(t, commands, "-o StrictHostKeyChecking=accept-new")
}

Expand Down Expand Up @@ -190,18 +217,85 @@ func TestAWSDeployRetainsFailedState(t *testing.T) {
require.Equal(t, "sg-123", state.AWS.SecurityGroupID)
}

func TestAWSRuntimeOptionValidation(t *testing.T) {
parameter, err := ubuntuAMIParameter("amd64")
require.NoError(t, err)
require.Equal(t, ubuntuAMD64AMIParameter, parameter)
parameter, err = ubuntuAMIParameter("arm64")
require.NoError(t, err)
require.Equal(t, ubuntuARM64AMIParameter, parameter)
_, err = ubuntuAMIParameter("riscv64")
require.Error(t, err)

for _, value := range []string{"off", "0", "100", "200"} {
require.NoError(t, validateGoGC(value))
}
for _, value := range []string{"", "-1", "invalid"} {
require.Error(t, validateGoGC(value))
}
}

func TestDescribeAWSInstancesHandlesMissingPublicIP(t *testing.T) {
runner := &fakeRunner{outputFn: func(commandSpec) (string, error) {
return `[{"InstanceID":"i-100","PublicIP":null,"PrivateIP":"10.0.0.10"}]`, nil
}}
instances, err := describeAWSInstances(t.Context(), awsClient{runner: runner}, []string{"i-100"})
require.NoError(t, err)
require.Equal(t, awsInstanceState{InstanceID: "i-100", PrivateIP: "10.0.0.10"}, instances["i-100"])
}

func TestWaitForNativeBuildReturnsRemoteLogOnFailure(t *testing.T) {
pendingInstance := awsInstanceState{NodeIndex: 1, PublicIP: "203.0.113.11"}
instance := awsInstanceState{NodeIndex: 2, PublicIP: "203.0.113.12"}
state := clusterState{AWS: &awsState{
SSHUser: "ubuntu",
SSHKeyPath: "/tmp/test.pem",
Instances: []awsInstanceState{pendingInstance, instance},
}}
runner := &fakeRunner{outputFn: func(spec commandSpec) (string, error) {
joined := strings.Join(spec.args, " ")
if strings.Contains(joined, "tail -n 200") {
return "compile failed: missing library\n", nil
}
if strings.Contains(joined, pendingInstance.PublicIP) {
return "", nil
}
return nativeBuildStatusFailed, nil
}}
app := &application{runner: runner, stdout: &bytes.Buffer{}, stderr: &bytes.Buffer{}}

err := app.waitForNativeBuilds(
t.Context(),
state,
"/remote/build/"+nativeBuildReadyFile,
"/remote/build/"+nativeBuildFailedFile,
"/remote/build/"+nativeBuildLogFile,
)

require.ErrorContains(t, err, "native build failed")
require.ErrorContains(t, err, "node-2")
require.ErrorContains(t, err, "compile failed: missing library")
require.Len(t, runner.commands, 3)
require.Contains(t, strings.Join(runner.commands[2].args, " "), "tail -n 200")
}

func TestAWSForwardUsesChosenNodePort(t *testing.T) {
stateDir := t.TempDir()
state := clusterState{
Version: stateVersion,
Name: "forward-test",
Target: "aws",
Status: "ready",
Nodes: clusterNodes(4),
Nodes: nativeClusterNodes(4),
AWS: &awsState{
PublicIP: "203.0.113.10",
SSHUser: "ubuntu",
SSHKeyPath: "/tmp/test.pem",
Instances: []awsInstanceState{
{NodeIndex: 0, InstanceID: "i-100", PublicIP: "203.0.113.10", PrivateIP: "10.0.0.10"},
{NodeIndex: 1, InstanceID: "i-101", PublicIP: "203.0.113.11", PrivateIP: "10.0.0.11"},
{NodeIndex: 2, InstanceID: "i-102", PublicIP: "203.0.113.12", PrivateIP: "10.0.0.12"},
{NodeIndex: 3, InstanceID: "i-103", PublicIP: "203.0.113.13", PrivateIP: "10.0.0.13"},
},
},
}
require.NoError(t, newStateStore(stateDir).save(state))
Expand All @@ -218,8 +312,8 @@ func TestAWSForwardUsesChosenNodePort(t *testing.T) {
require.Len(t, runner.commands, 1)
require.Equal(t, "ssh", runner.commands[0].name)
joined := strings.Join(runner.commands[0].args, " ")
require.Contains(t, joined, "-L 127.0.0.1:18545:127.0.0.1:8551")
require.True(t, strings.HasSuffix(joined, "ubuntu@203.0.113.10"))
require.Contains(t, joined, "-L 127.0.0.1:18545:127.0.0.1:8545")
require.True(t, strings.HasSuffix(joined, "ubuntu@203.0.113.13"))
}

func TestListShowsPartialAWSDeploymentWithoutCredentials(t *testing.T) {
Expand Down Expand Up @@ -247,6 +341,55 @@ func TestListShowsPartialAWSDeploymentWithoutCredentials(t *testing.T) {
require.Empty(t, runner.commands)
}

func TestListInspectsEachNativeAWSInstance(t *testing.T) {
stateDir := t.TempDir()
state := clusterState{
Version: stateVersion,
Name: "native-aws",
Target: targetAWS,
Status: "ready",
Nodes: nativeClusterNodes(4),
AWS: &awsState{
Region: "us-west-2",
SSHUser: "ubuntu",
SSHKeyPath: "/tmp/test.pem",
Instances: []awsInstanceState{
{NodeIndex: 0, InstanceID: "i-100", PublicIP: "203.0.113.10", PrivateIP: "10.0.0.10"},
{NodeIndex: 1, InstanceID: "i-101", PublicIP: "203.0.113.11", PrivateIP: "10.0.0.11"},
{NodeIndex: 2, InstanceID: "i-102", PublicIP: "203.0.113.12", PrivateIP: "10.0.0.12"},
{NodeIndex: 3, InstanceID: "i-103", PublicIP: "203.0.113.13", PrivateIP: "10.0.0.13"},
},
},
}
require.NoError(t, newStateStore(stateDir).save(state))
runner := &fakeRunner{outputFn: func(spec commandSpec) (string, error) {
joined := strings.Join(spec.args, " ")
switch {
case strings.Contains(joined, "sts get-caller-identity"):
return `{}`, nil
case strings.Contains(joined, "describe-instances"):
return "running\n", nil
case spec.name == "ssh" && strings.Contains(joined, "systemctl is-active"):
return "active\n", nil
case spec.name == "ssh" && strings.Contains(joined, "26660/metrics"):
return "tendermint_internal_autobahn_data_next_block{stage=\"execute\"} 43\n", nil
default:
return "", nil
}
}}
var stdout bytes.Buffer
app := &application{runner: runner, stdout: &stdout, stderr: &bytes.Buffer{}, stateDir: stateDir}

require.NoError(t, app.list(context.Background(), listOptions{name: state.Name}))
for nodeIndex := range 4 {
require.Contains(t, stdout.String(), "node-"+fmt.Sprint(nodeIndex))
require.Contains(t, stdout.String(), "i-10"+fmt.Sprint(nodeIndex))
require.Contains(t, stdout.String(), "203.0.113.1"+fmt.Sprint(nodeIndex))
}
require.Equal(t, 4, strings.Count(stdout.String(), "active"))
require.Equal(t, 4, strings.Count(stdout.String(), "42"))
}

func TestAWSTeardownToleratesAlreadyDeletedManagedResources(t *testing.T) {
stateDir := t.TempDir()
keyPath := filepath.Join(stateDir, "managed.pem")
Expand Down Expand Up @@ -281,6 +424,51 @@ func TestAWSTeardownToleratesAlreadyDeletedManagedResources(t *testing.T) {
require.NoFileExists(t, store.path(state.Name))
}

func TestAWSTeardownStopsAndTerminatesNativeInstances(t *testing.T) {
stateDir := t.TempDir()
keyPath := filepath.Join(stateDir, "managed.pem")
require.NoError(t, os.WriteFile(keyPath, []byte("key"), 0o600))
state := clusterState{
Version: stateVersion,
Name: "native-teardown",
Target: targetAWS,
Status: "ready",
Nodes: nativeClusterNodes(4),
AWS: &awsState{
Region: "us-west-2",
SSHUser: "ubuntu",
SSHKeyPath: keyPath,
SecurityGroupID: "sg-native",
KeyName: "key-native",
ManagedKey: true,
Instances: []awsInstanceState{
{NodeIndex: 0, InstanceID: "i-100", PublicIP: "203.0.113.10"},
{NodeIndex: 1, InstanceID: "i-101", PublicIP: "203.0.113.11"},
{NodeIndex: 2, InstanceID: "i-102", PublicIP: "203.0.113.12"},
{NodeIndex: 3, InstanceID: "i-103", PublicIP: "203.0.113.13"},
},
},
}
store := newStateStore(stateDir)
require.NoError(t, store.save(state))
runner := &fakeRunner{outputFn: func(spec commandSpec) (string, error) {
if strings.Contains(strings.Join(spec.args, " "), "sts get-caller-identity") {
return `{}`, nil
}
return "", nil
}}
app := &application{runner: runner, stdout: &bytes.Buffer{}, stderr: &bytes.Buffer{}, stateDir: stateDir}

require.NoError(t, app.teardown(context.Background(), teardownOptions{name: state.Name}))
commands := joinedCommands(runner.commands)
require.Equal(t, 4, strings.Count(commands, "systemctl stop seid.service"))
require.Contains(t, commands, "terminate-instances --instance-ids i-100 i-101 i-102 i-103")
require.Contains(t, commands, "instance-terminated --instance-ids i-100 i-101 i-102 i-103")
require.NotContains(t, commands, "docker-cluster-stop")
require.NoFileExists(t, keyPath)
require.NoFileExists(t, store.path(state.Name))
}

func TestLocalTeardownRemovesState(t *testing.T) {
stateDir := t.TempDir()
state := clusterState{
Expand Down Expand Up @@ -311,14 +499,15 @@ func TestParseAutobahnExecutedHeight(t *testing.T) {
require.Equal(t, "-", parseAutobahnExecutedHeight("not-prometheus"))
}

func TestWriteUserDataUsesSelectedSSHUser(t *testing.T) {
path, err := writeUserData(t.TempDir(), "test", "ec2-user")
func TestWriteUserDataInstallsNativeBuildDependencies(t *testing.T) {
path, err := writeUserData(t.TempDir(), "test")
require.NoError(t, err)
data, err := os.ReadFile(path)
require.NoError(t, err)
require.Contains(t, string(data), "usermod -aG docker ec2-user")
require.Contains(t, string(data), "build-essential python3")
require.Contains(t, string(data), "go1.25.6")
require.Contains(t, string(data), "/var/lib/autobahn-e2e-ready")
require.NotContains(t, string(data), "docker")
}

func TestShellQuote(t *testing.T) {
Expand Down
11 changes: 9 additions & 2 deletions cmd/autobahn-e2e/deploy.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
)

const dockerClusterSize = 4
const awsClusterSize = 4

type deployOptions struct {
name string
Expand All @@ -21,13 +22,16 @@ type deployOptions struct {
region string
profile string
instanceType string
architecture string
amiID string
subnetID string
sshCIDR string
sshUser string
keyName string
sshKeyPath string
volumeSize int
goMaxProcs int
goGC string
repoURL string
ref string
}
Expand All @@ -44,17 +48,20 @@ func (a *application) newDeployCommand() *cobra.Command {
flags := cmd.Flags()
flags.StringVar(&options.name, "name", defaultClusterName, "cluster name")
flags.StringVar(&options.target, "target", targetLocal, "deployment target: local or aws")
flags.DurationVar(&options.timeout, "timeout", 20*time.Minute, "deployment readiness timeout")
flags.DurationVar(&options.timeout, "timeout", 45*time.Minute, "deployment readiness timeout")
flags.StringVar(&options.region, "region", "us-west-2", "AWS region")
flags.StringVar(&options.profile, "profile", "", "AWS CLI profile")
flags.StringVar(&options.instanceType, "instance-type", "c7g.2xlarge", "EC2 instance type")
flags.StringVar(&options.amiID, "ami-id", "", "EC2 AMI ID; defaults to Ubuntu 24.04 ARM64")
flags.StringVar(&options.architecture, "architecture", "arm64", "EC2 architecture used to resolve the default AMI: arm64 or amd64")
flags.StringVar(&options.amiID, "ami-id", "", "EC2 AMI ID; defaults to Ubuntu 24.04 for --architecture")
flags.StringVar(&options.subnetID, "subnet-id", "", "EC2 subnet; defaults to a default VPC subnet")
flags.StringVar(&options.sshCIDR, "ssh-cidr", "", "CIDR allowed to SSH; defaults to the caller's public IP")
flags.StringVar(&options.sshUser, "ssh-user", "ubuntu", "EC2 SSH user")
flags.StringVar(&options.keyName, "key-name", "", "existing EC2 key pair name; omitted creates a managed key")
flags.StringVar(&options.sshKeyPath, "ssh-key", "", "private key for --key-name")
flags.IntVar(&options.volumeSize, "volume-size", 100, "EC2 root volume size in GiB")
flags.IntVar(&options.goMaxProcs, "gomaxprocs", 0, "GOMAXPROCS for each validator; 0 uses all instance CPUs")
flags.StringVar(&options.goGC, "gogc", "200", "GOGC for each validator, or off")
flags.StringVar(&options.repoURL, "repo-url", "", "Git repository cloned on EC2; defaults to origin")
flags.StringVar(&options.ref, "ref", "", "Git ref deployed on EC2; defaults to the current commit")
return cmd
Expand Down
8 changes: 6 additions & 2 deletions cmd/autobahn-e2e/forward.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,8 +60,12 @@ func (a *application) forward(ctx context.Context, options forwardOptions) error
if state.AWS == nil {
return fmt.Errorf("aws metadata is missing")
}
_, _ = fmt.Fprintf(a.stdout, "Forwarding %s to %s:8545 through %s. Press Ctrl-C to stop.\n", localAddress, node.Name, state.AWS.PublicIP)
baseArgs := sshBaseArgs(state)
instance, err := awsInstanceForNode(state, node)
if err != nil {
return err
}
_, _ = fmt.Fprintf(a.stdout, "Forwarding %s to %s:8545 through %s. Press Ctrl-C to stop.\n", localAddress, node.Name, instance.PublicIP)
baseArgs := sshBaseArgs(state, instance)
destination := baseArgs[len(baseArgs)-1]
args := append(baseArgs[:len(baseArgs)-1],
"-o", "ExitOnForwardFailure=yes",
Expand Down
Loading
Loading