fix: race on warned flag in versioncheck pkg, make consistent use of client/server terms

This commit is contained in:
Pasha Sviderski
2026-04-08 20:42:27 +10:00
parent 883f184e37
commit ab644f6a47
8 changed files with 122 additions and 138 deletions
@@ -1,12 +1,14 @@
package versioncheck package grpcversion
import ( import (
"context" "context"
"fmt" "fmt"
"io"
"os" "os"
"sync/atomic"
"github.com/Masterminds/semver" "github.com/Masterminds/semver"
internalVersion "github.com/psviderski/uncloud/internal/version" "github.com/psviderski/uncloud/internal/version"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/codes" "google.golang.org/grpc/codes"
"google.golang.org/grpc/metadata" "google.golang.org/grpc/metadata"
@@ -14,39 +16,43 @@ import (
) )
const ( const (
MetadataKeyCLIVersion = "uncloud-client-version" MetadataKeyClientVersion = "uncloud-client-version"
MetadataKeyMinDaemonVersion = "uncloud-min-server-version" MetadataKeyMinServerVersion = "uncloud-min-server-version"
MetadataKeyDaemonVersion = "uncloud-server-version" MetadataKeyServerVersion = "uncloud-server-version"
// MinCLIVersion is the minimum client version the daemon accepts. The daemon // MinClientVersion is the minimum client version the daemon accepts. The daemon
// rejects requests from older clients, forcing them to upgrade. This provides // rejects requests from older clients, forcing them to upgrade. This provides
// a clean cut-off for dropping support for old clients. // a clean cut-off for dropping support for old clients.
// //
// MinDaemonVersion is the minimum daemon version the client requires. The client // MinServerVersion is the minimum daemon version the client requires. The client
// sends this with each request so the daemon can immediately reject if it's too old, // sends this with each request so the daemon can immediately reject if it's too old,
// avoiding the need for a preflight request. This is useful when a new client feature // avoiding the need for a preflight request. This is useful when a new client feature
// requires daemon capabilities that didn't exist in older versions. // requires daemon capabilities that didn't exist in older versions.
// //
// The two minimums are independent: a client might require a newer daemon for new // The two minimums are independent: a client might require a newer daemon for new
// features, while that same daemon could still handle requests from older clients. // features, while that same daemon could still handle requests from older clients.
MinCLIVersion = "0.0.0" MinClientVersion = "0.0.0"
MinDaemonVersion = "0.0.0" MinServerVersion = "0.0.0"
ReleaseURL = "https://github.com/psviderski/uncloud/releases/latest" ReleaseURL = "https://github.com/psviderski/uncloud/releases/latest"
) )
var ( var (
// currentVersion is the version of this binary (CLI or daemon) // currentVersion is the version of this binary (CLI or daemon).
currentVersion = semver.MustParse(internalVersion.String()) currentVersion = semver.MustParse(version.String())
// zeroVersion is used when no version is specified (treated as 0.0.0) // zeroVersion is used when no version is specified (treated as 0.0.0).
zeroVersion = semver.MustParse("0.0.0") zeroVersion = semver.MustParse("0.0.0")
// Pre-parsed minimum versions for comparison // Pre-parsed minimum versions for comparison.
minCLIVersion = semver.MustParse(MinCLIVersion) minClientVersion = semver.MustParse(MinClientVersion)
minDaemonVersion = semver.MustParse(MinDaemonVersion) minServerVersion = semver.MustParse(MinServerVersion)
// warned tracks if we've already printed the daemon version warning // warned tracks if we've already printed the daemon version warning.
// TODO: remove when checkDaemonVersionInResponse is no longer needed (see below) // TODO: Remove when checkServerVersionInResponse is no longer needed (see below).
warned bool warned atomic.Bool
// WarnWriter is the writer used for version mismatch warnings. Defaults to os.Stderr.
// Tests can override this to capture warning output.
WarnWriter io.Writer = os.Stderr
) )
func extractVersion(md metadata.MD, key string) *semver.Version { func extractVersion(md metadata.MD, key string) *semver.Version {
@@ -67,18 +73,18 @@ func extractVersion(md metadata.MD, key string) *semver.Version {
func checkClientVersionHeaders(ctx context.Context) error { func checkClientVersionHeaders(ctx context.Context) error {
md, _ := metadata.FromIncomingContext(ctx) md, _ := metadata.FromIncomingContext(ctx)
actualCLIVersion := extractVersion(md, MetadataKeyCLIVersion) actualClientVersion := extractVersion(md, MetadataKeyClientVersion)
if actualCLIVersion.LessThan(minCLIVersion) { if actualClientVersion.LessThan(minClientVersion) {
return status.Errorf(codes.FailedPrecondition, return status.Errorf(codes.FailedPrecondition,
"version check failed: client version is below minimum %s. Please upgrade: %s", "version check failed: client version is below minimum %s. Please upgrade: %s",
minCLIVersion, ReleaseURL) minClientVersion, ReleaseURL)
} }
requiredMinDaemon := extractVersion(md, MetadataKeyMinDaemonVersion) requiredMinServer := extractVersion(md, MetadataKeyMinServerVersion)
if currentVersion.LessThan(requiredMinDaemon) { if currentVersion.LessThan(requiredMinServer) {
return status.Errorf(codes.FailedPrecondition, return status.Errorf(codes.FailedPrecondition,
"version check failed: daemon version %s is below client's minimum required version %s. Please upgrade the daemon: %s", "version check failed: daemon version %s is below client's minimum required version %s. Please upgrade the daemon: %s",
currentVersion, requiredMinDaemon, ReleaseURL) currentVersion, requiredMinServer, ReleaseURL)
} }
return nil return nil
@@ -88,7 +94,7 @@ func ServerUnaryInterceptor(ctx context.Context, req any, info *grpc.UnaryServer
if err := checkClientVersionHeaders(ctx); err != nil { if err := checkClientVersionHeaders(ctx); err != nil {
return nil, err return nil, err
} }
if err := grpc.SetHeader(ctx, metadata.Pairs(MetadataKeyDaemonVersion, currentVersion.String())); err != nil { if err := grpc.SetHeader(ctx, metadata.Pairs(MetadataKeyServerVersion, currentVersion.String())); err != nil {
return nil, err return nil, err
} }
return handler(ctx, req) return handler(ctx, req)
@@ -98,7 +104,7 @@ func ServerStreamInterceptor(srv any, ss grpc.ServerStream, info *grpc.StreamSer
if err := checkClientVersionHeaders(ss.Context()); err != nil { if err := checkClientVersionHeaders(ss.Context()); err != nil {
return err return err
} }
if err := ss.SetHeader(metadata.Pairs(MetadataKeyDaemonVersion, currentVersion.String())); err != nil { if err := ss.SetHeader(metadata.Pairs(MetadataKeyServerVersion, currentVersion.String())); err != nil {
return err return err
} }
return handler(srv, ss) return handler(srv, ss)
@@ -106,11 +112,11 @@ func ServerStreamInterceptor(srv any, ss grpc.ServerStream, info *grpc.StreamSer
func ClientUnaryInterceptor(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error { func ClientUnaryInterceptor(ctx context.Context, method string, req, reply any, cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption) error {
ctx = metadata.AppendToOutgoingContext(ctx, ctx = metadata.AppendToOutgoingContext(ctx,
MetadataKeyCLIVersion, currentVersion.String(), MetadataKeyClientVersion, currentVersion.String(),
MetadataKeyMinDaemonVersion, MinDaemonVersion, MetadataKeyMinServerVersion, MinServerVersion,
) )
// TODO: remove when checkDaemonVersionInResponse is no longer needed, // TODO: Remove when checkServerVersionInResponse is no longer needed,
// as we'll no longer need to extract headers from the response here. // as we'll no longer need to extract headers from the response here.
var respMD metadata.MD var respMD metadata.MD
opts = append(opts, grpc.Header(&respMD)) opts = append(opts, grpc.Header(&respMD))
@@ -120,34 +126,33 @@ func ClientUnaryInterceptor(ctx context.Context, method string, req, reply any,
return err return err
} }
// TODO: Remove eventually (see note on method below) // TODO: Remove eventually (see note on method below).
checkDaemonVersionInResponse(respMD) checkServerVersionInResponse(respMD)
return nil return nil
} }
// This is just needed as a warning during the transition to version checking // checkServerVersionInResponse warns the user when they communicated with a daemon that
// releases. It warns the user when they just communicated with a daemon that did // did not check the version requirements. This is only needed during the transition to
// not check the version requirements. // version-checking releases.
// TODO: Remove this in some later release, after users have upgraded. // TODO: Remove this in some later release, after users have upgraded.
func checkDaemonVersionInResponse(md metadata.MD) { func checkServerVersionInResponse(md metadata.MD) {
daemonVersion := extractVersion(md, MetadataKeyDaemonVersion) serverVersion := extractVersion(md, MetadataKeyServerVersion)
if daemonVersion.LessThan(minDaemonVersion) { if serverVersion.LessThan(minServerVersion) {
if warned { if warned.Swap(true) {
return return
} }
warned = true
msg := fmt.Sprintf("daemon version is below minimum required version %s. The daemon did not verify this CLI's minimum version requirement, so the operation may not have behaved as intended. Please upgrade the daemon: %s", msg := fmt.Sprintf("daemon version is below minimum required version %s. The daemon did not verify this CLI's minimum version requirement, so the operation may not have behaved as intended. Please upgrade the daemon: %s",
minDaemonVersion, ReleaseURL) minServerVersion, ReleaseURL)
fmt.Fprintf(os.Stderr, "WARNING: %s\n", msg) fmt.Fprintf(WarnWriter, "WARNING: %s\n", msg)
} }
} }
func ClientStreamInterceptor(ctx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn, method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) { func ClientStreamInterceptor(ctx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn, method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) {
ctx = metadata.AppendToOutgoingContext(ctx, ctx = metadata.AppendToOutgoingContext(ctx,
MetadataKeyCLIVersion, currentVersion.String(), MetadataKeyClientVersion, currentVersion.String(),
MetadataKeyMinDaemonVersion, MinDaemonVersion, MetadataKeyMinServerVersion, MinServerVersion,
) )
stream, err := streamer(ctx, desc, cc, method, opts...) stream, err := streamer(ctx, desc, cc, method, opts...)
@@ -157,23 +162,23 @@ func ClientStreamInterceptor(ctx context.Context, desc *grpc.StreamDesc, cc *grp
// TODO: Wrapping the stream in versionedClientStream will no longer // TODO: Wrapping the stream in versionedClientStream will no longer
// be necessary when we are ready to remove the temporary, transition // be necessary when we are ready to remove the temporary, transition
// safety check checkDaemonVersionInResponse (see note on method above) // safety check checkServerVersionInResponse (see note on method above).
return &versionedClientStream{ClientStream: stream}, nil return &versionedClientStream{ClientStream: stream}, nil
} }
// TODO: remove when checkDaemonVersionInResponse is no longer needed // TODO: Remove when checkServerVersionInResponse is no longer needed.
type versionedClientStream struct { type versionedClientStream struct {
grpc.ClientStream grpc.ClientStream
} }
// TODO: remove when checkDaemonVersionInResponse is no longer needed // TODO: Remove when checkServerVersionInResponse is no longer needed.
func (s *versionedClientStream) Header() (metadata.MD, error) { func (s *versionedClientStream) Header() (metadata.MD, error) {
md, err := s.ClientStream.Header() md, err := s.ClientStream.Header()
if err != nil { if err != nil {
return nil, err return nil, err
} }
checkDaemonVersionInResponse(md) checkServerVersionInResponse(md)
return md, nil return md, nil
} }
@@ -1,11 +1,8 @@
package versioncheck package grpcversion
import ( import (
"bytes" "bytes"
"context" "context"
"io"
"os"
"strings"
"testing" "testing"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
@@ -25,37 +22,37 @@ func TestExtractVersion(t *testing.T) {
{ {
name: "nil metadata", name: "nil metadata",
md: nil, md: nil,
key: MetadataKeyCLIVersion, key: MetadataKeyClientVersion,
expected: "0.0.0", expected: "0.0.0",
}, },
{ {
name: "missing key", name: "missing key",
md: metadata.MD{}, md: metadata.MD{},
key: MetadataKeyCLIVersion, key: MetadataKeyClientVersion,
expected: "0.0.0", expected: "0.0.0",
}, },
{ {
name: "empty value", name: "empty value",
md: metadata.Pairs(MetadataKeyCLIVersion, ""), md: metadata.Pairs(MetadataKeyClientVersion, ""),
key: MetadataKeyCLIVersion, key: MetadataKeyClientVersion,
expected: "0.0.0", expected: "0.0.0",
}, },
{ {
name: "invalid version", name: "invalid version",
md: metadata.Pairs(MetadataKeyCLIVersion, "not-a-version"), md: metadata.Pairs(MetadataKeyClientVersion, "not-a-version"),
key: MetadataKeyCLIVersion, key: MetadataKeyClientVersion,
expected: "0.0.0", expected: "0.0.0",
}, },
{ {
name: "valid version", name: "valid version",
md: metadata.Pairs(MetadataKeyCLIVersion, "1.2.3"), md: metadata.Pairs(MetadataKeyClientVersion, "1.2.3"),
key: MetadataKeyCLIVersion, key: MetadataKeyClientVersion,
expected: "1.2.3", expected: "1.2.3",
}, },
{ {
name: "version with prerelease", name: "version with prerelease",
md: metadata.Pairs(MetadataKeyCLIVersion, "0.0.0-dev"), md: metadata.Pairs(MetadataKeyClientVersion, "0.0.0-dev"),
key: MetadataKeyCLIVersion, key: MetadataKeyClientVersion,
expected: "0.0.0-dev", expected: "0.0.0-dev",
}, },
} }
@@ -79,7 +76,7 @@ func TestCheckClientVersionHeaders(t *testing.T) {
{ {
name: "cli version below minimum", name: "cli version below minimum",
md: metadata.Pairs( md: metadata.Pairs(
MetadataKeyCLIVersion, "0.0.0-dev", MetadataKeyClientVersion, "0.0.0-dev",
), ),
wantErr: true, wantErr: true,
errCode: codes.FailedPrecondition, errCode: codes.FailedPrecondition,
@@ -88,15 +85,15 @@ func TestCheckClientVersionHeaders(t *testing.T) {
{ {
name: "cli version above minimum", name: "cli version above minimum",
md: metadata.Pairs( md: metadata.Pairs(
MetadataKeyCLIVersion, "999.0.0", MetadataKeyClientVersion, "999.0.0",
), ),
wantErr: false, wantErr: false,
}, },
{ {
name: "min daemon version above current daemon", name: "min daemon version above current daemon",
md: metadata.Pairs( md: metadata.Pairs(
MetadataKeyCLIVersion, "999.0.0", MetadataKeyClientVersion, "999.0.0",
MetadataKeyMinDaemonVersion, "999.0.0", MetadataKeyMinServerVersion, "999.0.0",
), ),
wantErr: true, wantErr: true,
errCode: codes.FailedPrecondition, errCode: codes.FailedPrecondition,
@@ -105,8 +102,8 @@ func TestCheckClientVersionHeaders(t *testing.T) {
{ {
name: "min daemon version below current daemon", name: "min daemon version below current daemon",
md: metadata.Pairs( md: metadata.Pairs(
MetadataKeyCLIVersion, "999.0.0", MetadataKeyClientVersion, "999.0.0",
MetadataKeyMinDaemonVersion, "0.0.1", MetadataKeyMinServerVersion, "0.0.1",
), ),
wantErr: false, wantErr: false,
}, },
@@ -126,8 +123,7 @@ func TestCheckClientVersionHeaders(t *testing.T) {
st, ok := status.FromError(err) st, ok := status.FromError(err)
require.True(t, ok, "expected gRPC status error, got %T", err) require.True(t, ok, "expected gRPC status error, got %T", err)
assert.Equal(t, tt.errCode, st.Code()) assert.Equal(t, tt.errCode, st.Code())
assert.True(t, strings.Contains(st.Message(), tt.errContain), assert.Contains(t, st.Message(), tt.errContain)
"error message = %q, want to contain %q", st.Message(), tt.errContain)
} else { } else {
assert.NoError(t, err) assert.NoError(t, err)
} }
@@ -135,7 +131,17 @@ func TestCheckClientVersionHeaders(t *testing.T) {
} }
} }
func TestCheckDaemonVersionInResponse(t *testing.T) { func captureWarnings(t *testing.T, fn func()) string {
t.Helper()
var buf bytes.Buffer
old := WarnWriter
WarnWriter = &buf
t.Cleanup(func() { WarnWriter = old })
fn()
return buf.String()
}
func TestCheckServerVersionInResponse(t *testing.T) {
tests := []struct { tests := []struct {
name string name string
md metadata.MD md metadata.MD
@@ -143,74 +149,47 @@ func TestCheckDaemonVersionInResponse(t *testing.T) {
}{ }{
{ {
name: "daemon version below minimum", name: "daemon version below minimum",
md: metadata.Pairs(MetadataKeyDaemonVersion, "0.0.0-dev"), md: metadata.Pairs(MetadataKeyServerVersion, "0.0.0-dev"),
wantWarning: true, wantWarning: true,
}, },
{ {
name: "daemon version above minimum", name: "daemon version above minimum",
md: metadata.Pairs(MetadataKeyDaemonVersion, "999.0.0"), md: metadata.Pairs(MetadataKeyServerVersion, "999.0.0"),
wantWarning: false, wantWarning: false,
}, },
} }
for _, tt := range tests { for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
// Reset warned flag for each test warned.Store(false)
warned = false
// Capture stderr output := captureWarnings(t, func() {
old := os.Stderr checkServerVersionInResponse(tt.md)
r, w, _ := os.Pipe() })
os.Stderr = w
checkDaemonVersionInResponse(tt.md)
w.Close()
var buf bytes.Buffer
io.Copy(&buf, r)
os.Stderr = old
output := buf.String()
if tt.wantWarning { if tt.wantWarning {
assert.True(t, strings.Contains(output, "WARNING"), "expected warning output, got none") assert.Contains(t, output, "WARNING")
} else { } else {
assert.Equal(t, "", output) assert.Empty(t, output)
} }
}) })
} }
} }
func TestCheckDaemonVersionInResponse_WarnOnce(t *testing.T) { func TestCheckServerVersionInResponse_WarnOnce(t *testing.T) {
// Reset warned flag warned.Store(false)
warned = false
md := metadata.Pairs(MetadataKeyDaemonVersion, "0.0.0-dev") md := metadata.Pairs(MetadataKeyServerVersion, "0.0.0-dev")
// First call - should warn // First call should warn.
old := os.Stderr output1 := captureWarnings(t, func() {
r, w, _ := os.Pipe() checkServerVersionInResponse(md)
os.Stderr = w })
assert.Contains(t, output1, "WARNING", "first call should warn")
checkDaemonVersionInResponse(md) // Second call should not warn (warned flag is now true).
output2 := captureWarnings(t, func() {
w.Close() checkServerVersionInResponse(md)
var buf bytes.Buffer })
io.Copy(&buf, r) assert.Empty(t, output2, "second call should not warn")
os.Stderr = old
assert.True(t, strings.Contains(buf.String(), "WARNING"), "first call should warn")
// Second call - should NOT warn (warned flag is now true)
r2, w2, _ := os.Pipe()
os.Stderr = w2
checkDaemonVersionInResponse(md)
w2.Close()
var buf2 bytes.Buffer
io.Copy(&buf2, r2)
os.Stderr = old
assert.Equal(t, "", buf2.String(), "second call should not warn")
} }
+5 -5
View File
@@ -22,6 +22,7 @@ import (
"github.com/psviderski/uncloud/internal/corrosion" "github.com/psviderski/uncloud/internal/corrosion"
"github.com/psviderski/uncloud/internal/docker" "github.com/psviderski/uncloud/internal/docker"
"github.com/psviderski/uncloud/internal/fs" "github.com/psviderski/uncloud/internal/fs"
"github.com/psviderski/uncloud/internal/grpcversion"
"github.com/psviderski/uncloud/internal/journal" "github.com/psviderski/uncloud/internal/journal"
"github.com/psviderski/uncloud/internal/machine/api/pb" "github.com/psviderski/uncloud/internal/machine/api/pb"
apiproxy "github.com/psviderski/uncloud/internal/machine/api/proxy" apiproxy "github.com/psviderski/uncloud/internal/machine/api/proxy"
@@ -34,7 +35,6 @@ import (
"github.com/psviderski/uncloud/internal/machine/network" "github.com/psviderski/uncloud/internal/machine/network"
"github.com/psviderski/uncloud/internal/machine/store" "github.com/psviderski/uncloud/internal/machine/store"
"github.com/psviderski/uncloud/pkg/api" "github.com/psviderski/uncloud/pkg/api"
versionpkg "github.com/psviderski/uncloud/pkg/versioncheck"
"github.com/psviderski/unregistry" "github.com/psviderski/unregistry"
"github.com/siderolabs/grpc-proxy/proxy" "github.com/siderolabs/grpc-proxy/proxy"
"golang.org/x/sync/errgroup" "golang.org/x/sync/errgroup"
@@ -273,8 +273,8 @@ func NewMachine(config *Config) (*Machine, error) {
proxyDirector := apiproxy.NewDirector(config.MachineSockPath, constants.MachineAPIPort) proxyDirector := apiproxy.NewDirector(config.MachineSockPath, constants.MachineAPIPort)
localProxyServer := grpc.NewServer( localProxyServer := grpc.NewServer(
grpc.ForceServerCodecV2(proxy.Codec()), grpc.ForceServerCodecV2(proxy.Codec()),
grpc.UnaryInterceptor(versionpkg.ServerUnaryInterceptor), grpc.UnaryInterceptor(grpcversion.ServerUnaryInterceptor),
grpc.StreamInterceptor(versionpkg.ServerStreamInterceptor), grpc.StreamInterceptor(grpcversion.ServerStreamInterceptor),
grpc.UnknownServiceHandler( grpc.UnknownServiceHandler(
proxy.TransparentHandler(proxyDirector.Director), proxy.TransparentHandler(proxyDirector.Director),
), ),
@@ -428,8 +428,8 @@ func (m *Machine) Run(ctx context.Context) error {
m.proxyDirector.UpdateLocalAddress(m.state.Network.ManagementIP.String()) m.proxyDirector.UpdateLocalAddress(m.state.Network.ManagementIP.String())
proxyServer := grpc.NewServer( proxyServer := grpc.NewServer(
grpc.ForceServerCodecV2(proxy.Codec()), grpc.ForceServerCodecV2(proxy.Codec()),
grpc.UnaryInterceptor(versionpkg.ServerUnaryInterceptor), grpc.UnaryInterceptor(grpcversion.ServerUnaryInterceptor),
grpc.StreamInterceptor(versionpkg.ServerStreamInterceptor), grpc.StreamInterceptor(grpcversion.ServerStreamInterceptor),
grpc.UnknownServiceHandler( grpc.UnknownServiceHandler(
proxy.TransparentHandler(m.proxyDirector.Director), proxy.TransparentHandler(m.proxyDirector.Director),
), ),
+3 -3
View File
@@ -7,9 +7,9 @@ import (
"net" "net"
"strings" "strings"
"github.com/psviderski/uncloud/internal/grpcversion"
"github.com/psviderski/uncloud/internal/machine" "github.com/psviderski/uncloud/internal/machine"
"github.com/psviderski/uncloud/internal/sshexec" "github.com/psviderski/uncloud/internal/sshexec"
"github.com/psviderski/uncloud/pkg/versioncheck"
"golang.org/x/crypto/ssh" "golang.org/x/crypto/ssh"
"golang.org/x/net/proxy" "golang.org/x/net/proxy"
"google.golang.org/grpc" "google.golang.org/grpc"
@@ -75,8 +75,8 @@ func (c *SSHConnector) Connect(ctx context.Context) (*grpc.ClientConn, error) {
"unix://"+sockPath, "unix://"+sockPath,
grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(defaultServiceConfig), grpc.WithDefaultServiceConfig(defaultServiceConfig),
grpc.WithUnaryInterceptor(versioncheck.ClientUnaryInterceptor), grpc.WithUnaryInterceptor(grpcversion.ClientUnaryInterceptor),
grpc.WithStreamInterceptor(versioncheck.ClientStreamInterceptor), grpc.WithStreamInterceptor(grpcversion.ClientStreamInterceptor),
grpc.WithContextDialer( grpc.WithContextDialer(
func(ctx context.Context, addr string) (net.Conn, error) { func(ctx context.Context, addr string) (net.Conn, error) {
addr = strings.TrimPrefix(addr, "unix://") addr = strings.TrimPrefix(addr, "unix://")
+3 -3
View File
@@ -11,7 +11,7 @@ import (
"strings" "strings"
"github.com/docker/cli/cli/connhelper/commandconn" "github.com/docker/cli/cli/connhelper/commandconn"
"github.com/psviderski/uncloud/pkg/versioncheck" "github.com/psviderski/uncloud/internal/grpcversion"
"golang.org/x/net/proxy" "golang.org/x/net/proxy"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
@@ -78,8 +78,8 @@ func (c *SSHCLIConnector) Connect(ctx context.Context) (*grpc.ClientConn, error)
"passthrough:///", // Dummy target since we're using a custom dialer. "passthrough:///", // Dummy target since we're using a custom dialer.
grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(defaultServiceConfig), grpc.WithDefaultServiceConfig(defaultServiceConfig),
grpc.WithUnaryInterceptor(versioncheck.ClientUnaryInterceptor), grpc.WithUnaryInterceptor(grpcversion.ClientUnaryInterceptor),
grpc.WithStreamInterceptor(versioncheck.ClientStreamInterceptor), grpc.WithStreamInterceptor(grpcversion.ClientStreamInterceptor),
grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) { grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) {
dialArgs := append(c.buildSSHArgs(), "uncloudd", "dial-stdio") dialArgs := append(c.buildSSHArgs(), "uncloudd", "dial-stdio")
if c.config.SockPath != "" { if c.config.SockPath != "" {
+3 -3
View File
@@ -5,7 +5,7 @@ import (
"fmt" "fmt"
"net/netip" "net/netip"
"github.com/psviderski/uncloud/pkg/versioncheck" "github.com/psviderski/uncloud/internal/grpcversion"
"golang.org/x/net/proxy" "golang.org/x/net/proxy"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
@@ -25,8 +25,8 @@ func (c *TCPConnector) Connect(_ context.Context) (*grpc.ClientConn, error) {
c.apiAddr.String(), c.apiAddr.String(),
grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(defaultServiceConfig), grpc.WithDefaultServiceConfig(defaultServiceConfig),
grpc.WithUnaryInterceptor(versioncheck.ClientUnaryInterceptor), grpc.WithUnaryInterceptor(grpcversion.ClientUnaryInterceptor),
grpc.WithStreamInterceptor(versioncheck.ClientStreamInterceptor), grpc.WithStreamInterceptor(grpcversion.ClientStreamInterceptor),
) )
if err != nil { if err != nil {
return nil, fmt.Errorf("create machine API client: %w", err) return nil, fmt.Errorf("create machine API client: %w", err)
+3 -3
View File
@@ -4,7 +4,7 @@ import (
"context" "context"
"fmt" "fmt"
"github.com/psviderski/uncloud/pkg/versioncheck" "github.com/psviderski/uncloud/internal/grpcversion"
"golang.org/x/net/proxy" "golang.org/x/net/proxy"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
@@ -27,8 +27,8 @@ func (c *UnixConnector) Connect(_ context.Context) (*grpc.ClientConn, error) {
target, target,
grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultServiceConfig(defaultServiceConfig), grpc.WithDefaultServiceConfig(defaultServiceConfig),
grpc.WithUnaryInterceptor(versioncheck.ClientUnaryInterceptor), grpc.WithUnaryInterceptor(grpcversion.ClientUnaryInterceptor),
grpc.WithStreamInterceptor(versioncheck.ClientStreamInterceptor), grpc.WithStreamInterceptor(grpcversion.ClientStreamInterceptor),
) )
if err != nil { if err != nil {
return nil, fmt.Errorf("create machine API client: %w", err) return nil, fmt.Errorf("create machine API client: %w", err)
+3 -3
View File
@@ -8,11 +8,11 @@ import (
"strconv" "strconv"
"github.com/psviderski/uncloud/internal/cli/config" "github.com/psviderski/uncloud/internal/cli/config"
"github.com/psviderski/uncloud/internal/grpcversion"
"github.com/psviderski/uncloud/internal/machine/constants" "github.com/psviderski/uncloud/internal/machine/constants"
"github.com/psviderski/uncloud/internal/machine/network" "github.com/psviderski/uncloud/internal/machine/network"
"github.com/psviderski/uncloud/internal/machine/network/tunnel" "github.com/psviderski/uncloud/internal/machine/network/tunnel"
"github.com/psviderski/uncloud/pkg/client" "github.com/psviderski/uncloud/pkg/client"
"github.com/psviderski/uncloud/pkg/versioncheck"
"golang.org/x/net/proxy" "golang.org/x/net/proxy"
"google.golang.org/grpc" "google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/credentials/insecure"
@@ -71,8 +71,8 @@ func (c *WireGuardConnector) Connect(ctx context.Context) (*grpc.ClientConn, err
grpc.WithContextDialer(func(ctx context.Context, addr string) (net.Conn, error) { grpc.WithContextDialer(func(ctx context.Context, addr string) (net.Conn, error) {
return c.tun.DialContext(ctx, "tcp", addr) return c.tun.DialContext(ctx, "tcp", addr)
}), }),
grpc.WithUnaryInterceptor(versioncheck.ClientUnaryInterceptor), grpc.WithUnaryInterceptor(grpcversion.ClientUnaryInterceptor),
grpc.WithStreamInterceptor(versioncheck.ClientStreamInterceptor), grpc.WithStreamInterceptor(grpcversion.ClientStreamInterceptor),
) )
if err != nil { if err != nil {
return nil, fmt.Errorf("connect to machine API through WireGuard tunnel: %w", err) return nil, fmt.Errorf("connect to machine API through WireGuard tunnel: %w", err)