fix(instance): stabilize desktop access and stream defaults

- Reuse desktop direct upstream resolution and include upstream in short external access tokens
- Default and normalize Selkies stream settings for desktop runtimes
- Remove legacy deployments with the same instance-id before ensuring the current deployment
- Preserve the unsaved desktop stream profile state in the frontend
- Add tests for stream profiles, pod envs, and deployment cleanup
This commit is contained in:
litiantian03
2026-06-18 11:20:13 +08:00
parent b8a7e8c67a
commit 3c35b26b74
10 changed files with 255 additions and 41 deletions
+38 -27
View File
@@ -55,6 +55,39 @@ func desktopProxyMode(directEnabled bool, upstream string) string {
return "direct"
}
func (h *InstanceHandler) desktopAccessUpstream(c *gin.Context, instance *models.Instance, targetPort int32) (string, bool) {
directProxyEnabled := desktopDirectProxyEnabled()
if !directProxyEnabled {
return "", false
}
if !h.proxyService.IsWebtopInstanceType(instance.Type) {
fmt.Printf("Desktop direct proxy fallback: unsupported desktop instance type instance=%d user=%d type=%s target_port=%d\n",
instance.ID, instance.UserID, instance.Type, targetPort)
return "", true
}
resolved, resolveErr := h.proxyService.ResolveUpstreamHostPort(
c.Request.Context(),
instance.UserID,
instance.ID,
targetPort,
)
if resolveErr != nil {
fmt.Printf("Desktop direct proxy fallback: failed to resolve upstream instance=%d user=%d type=%s target_port=%d error=%v\n",
instance.ID, instance.UserID, instance.Type, targetPort, resolveErr)
return "", true
}
if strings.TrimSpace(resolved) == "" {
fmt.Printf("Desktop direct proxy fallback: resolved empty upstream instance=%d user=%d type=%s target_port=%d\n",
instance.ID, instance.UserID, instance.Type, targetPort)
return "", true
}
fmt.Printf("Desktop direct proxy resolved: instance=%d user=%d type=%s target_port=%d upstream=%s\n",
instance.ID, instance.UserID, instance.Type, targetPort, resolved)
return resolved, true
}
func workspaceArchiveMaxMiB() int64 {
value := strings.TrimSpace(os.Getenv(workspaceArchiveMaxMiBEnv))
if value == "" {
@@ -880,32 +913,7 @@ func (h *InstanceHandler) GenerateAccessToken(c *gin.Context) {
// "host:port" into the token so the edge gateway can dial the instance
// directly. On any failure we fall back to an empty upstream, which keeps
// the request flowing through the in-process control-plane proxy.
upstream := ""
directProxyEnabled := desktopDirectProxyEnabled()
if directProxyEnabled {
if h.proxyService.IsWebtopInstanceType(instance.Type) {
resolved, resolveErr := h.proxyService.ResolveUpstreamHostPort(
c.Request.Context(),
instance.UserID,
instance.ID,
targetPort,
)
if resolveErr != nil {
fmt.Printf("Desktop direct proxy fallback: failed to resolve upstream instance=%d user=%d type=%s target_port=%d error=%v\n",
instance.ID, instance.UserID, instance.Type, targetPort, resolveErr)
} else if strings.TrimSpace(resolved) == "" {
fmt.Printf("Desktop direct proxy fallback: resolved empty upstream instance=%d user=%d type=%s target_port=%d\n",
instance.ID, instance.UserID, instance.Type, targetPort)
} else {
upstream = resolved
fmt.Printf("Desktop direct proxy resolved: instance=%d user=%d type=%s target_port=%d upstream=%s\n",
instance.ID, instance.UserID, instance.Type, targetPort, upstream)
}
} else {
fmt.Printf("Desktop direct proxy fallback: unsupported desktop instance type instance=%d user=%d type=%s target_port=%d\n",
instance.ID, instance.UserID, instance.Type, targetPort)
}
}
upstream, directProxyEnabled := h.desktopAccessUpstream(c, instance, targetPort)
// Generate access token (valid for 1 hour)
maxAgeSeconds := int(time.Hour.Seconds())
@@ -1605,12 +1613,15 @@ func (h *InstanceHandler) issueShortExternalAccessToken(c *gin.Context, instance
utils.Error(c, http.StatusServiceUnavailable, "Unable to generate access URL")
return nil, false
}
targetPort := h.proxyService.GetTargetPortForInstance(instance)
upstream, _ := h.desktopAccessUpstream(c, instance, targetPort)
instanceToken, err := h.accessService.GenerateToken(
instance.UserID,
instance.ID,
instance.Type,
accessURL,
h.proxyService.GetTargetPortForInstance(instance),
upstream,
targetPort,
1*time.Hour,
)
if err != nil {
@@ -189,7 +189,7 @@ func TestProxyAccessTokenPrefersCookieOverRuntimeQueryToken(t *testing.T) {
gin.SetMode(gin.TestMode)
accessService := services.NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(1, 76, "hermes", "/api/v1/instances/76/proxy/chat/", 3000, time.Hour)
token, err := accessService.GenerateToken(1, 76, "hermes", "/api/v1/instances/76/proxy/chat/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -13,6 +13,7 @@ const (
var selkiesDesktopStreamEnvKeys = []string{
desktopStreamProfileEnvKey,
"SELKIES_ENCODER",
"SELKIES_USE_CSS_SCALING",
"SELKIES_FRAMERATE",
"SELKIES_H264_CRF",
"SELKIES_AUDIO_ENABLED",
@@ -49,7 +50,8 @@ func applyDesktopStreamProfileEnv(overrides map[string]string, profile string) m
}
overrides[desktopStreamProfileEnvKey] = normalized
overrides["SELKIES_ENCODER"] = "x264enc"
overrides["SELKIES_ENCODER"] = "x264enc,jpeg"
overrides["SELKIES_USE_CSS_SCALING"] = "true"
overrides["SELKIES_SECOND_SCREEN"] = "false"
overrides["SELKIES_AUDIO_ENABLED"] = "false"
@@ -68,6 +70,25 @@ func applyDesktopStreamProfileEnv(overrides map[string]string, profile string) m
return overrides
}
func ensureDesktopStreamProfileEnv(overrides map[string]string, runtimeType string) map[string]string {
if normalizeInstanceRuntimeType(runtimeType) != RuntimeBackendDesktop {
return overrides
}
profile := desktopStreamProfileFromEnv(overrides)
if profile != "" {
return applyDesktopStreamProfileEnv(overrides, profile)
}
for _, key := range selkiesDesktopStreamEnvKeys {
if _, ok := overrides[key]; ok {
return overrides
}
}
return applyDesktopStreamProfileEnv(overrides, DesktopStreamProfileStandard)
}
func desktopStreamProfileFromEnv(overrides map[string]string) string {
if profile, ok := overrides[desktopStreamProfileEnvKey]; ok {
if normalized, valid := normalizeDesktopStreamProfile(profile); valid && normalized != "" {
@@ -0,0 +1,76 @@
package services
import "testing"
func TestApplyDesktopStreamProfileEnvStandard(t *testing.T) {
overrides := applyDesktopStreamProfileEnv(map[string]string{
"SELKIES_ENCODER": "x264enc",
"SELKIES_USE_CSS_SCALING": "false",
"SELKIES_FRAMERATE": "10",
"SELKIES_H264_CRF": "20",
}, DesktopStreamProfileStandard)
want := map[string]string{
"CLAWMANAGER_DESKTOP_STREAM_PROFILE": "standard",
"SELKIES_ENCODER": "x264enc,jpeg",
"SELKIES_USE_CSS_SCALING": "true",
"SELKIES_FRAMERATE": "35",
"SELKIES_H264_CRF": "34",
"SELKIES_SECOND_SCREEN": "false",
"SELKIES_AUDIO_ENABLED": "false",
}
for key, value := range want {
if got := overrides[key]; got != value {
t.Fatalf("%s = %q, want %q", key, got, value)
}
}
}
func TestDesktopStreamProfileFromEnvIgnoresEncoderDetails(t *testing.T) {
overrides := map[string]string{
"SELKIES_ENCODER": "x264enc,jpeg",
"SELKIES_USE_CSS_SCALING": "true",
"SELKIES_FRAMERATE": "40",
"SELKIES_H264_CRF": "24",
}
if got := desktopStreamProfileFromEnv(overrides); got != DesktopStreamProfileHigh {
t.Fatalf("desktopStreamProfileFromEnv() = %q, want %q", got, DesktopStreamProfileHigh)
}
}
func TestEnsureDesktopStreamProfileEnvUpgradesSavedProfile(t *testing.T) {
overrides := ensureDesktopStreamProfileEnv(map[string]string{
"CLAWMANAGER_DESKTOP_STREAM_PROFILE": "standard",
"SELKIES_ENCODER": "x264enc",
"SELKIES_FRAMERATE": "35",
"SELKIES_H264_CRF": "34",
}, RuntimeBackendDesktop)
if got := overrides["SELKIES_ENCODER"]; got != "x264enc,jpeg" {
t.Fatalf("SELKIES_ENCODER = %q, want x264enc,jpeg", got)
}
if got := overrides["SELKIES_USE_CSS_SCALING"]; got != "true" {
t.Fatalf("SELKIES_USE_CSS_SCALING = %q, want true", got)
}
}
func TestEnsureDesktopStreamProfileEnvDefaultsDesktopOnlyWhenUnset(t *testing.T) {
desktop := ensureDesktopStreamProfileEnv(nil, RuntimeBackendDesktop)
if got := desktop["CLAWMANAGER_DESKTOP_STREAM_PROFILE"]; got != DesktopStreamProfileStandard {
t.Fatalf("desktop profile = %q, want %q", got, DesktopStreamProfileStandard)
}
custom := ensureDesktopStreamProfileEnv(map[string]string{
"SELKIES_ENCODER": "custom",
}, RuntimeBackendDesktop)
if got := custom["SELKIES_ENCODER"]; got != "custom" {
t.Fatalf("custom encoder = %q, want custom", got)
}
shell := ensureDesktopStreamProfileEnv(nil, RuntimeBackendShell)
if shell != nil {
t.Fatalf("shell runtime should not receive desktop stream env")
}
}
@@ -117,6 +117,7 @@ func buildInstancePodEnv(instance *models.Instance, runtimeEnv, gatewayEnv, agen
if err != nil {
return nil, err
}
overrides = ensureDesktopStreamProfileEnv(overrides, instance.RuntimeType)
resolved := mergeEnvMaps(runtimeEnv, mergeEnvMaps(gatewayEnv, agentEnv))
resolved = withInstanceProxyEnv(instance.Type, instance.ID, resolved)
@@ -73,6 +73,35 @@ func TestBuildInstancePodEnvAppliesOverridesAfterDefaults(t *testing.T) {
}
}
func TestBuildInstancePodEnvNormalizesDesktopStreamProfile(t *testing.T) {
raw, err := marshalEnvironmentOverrides(map[string]string{
"CLAWMANAGER_DESKTOP_STREAM_PROFILE": "standard",
"SELKIES_ENCODER": "x264enc",
"SELKIES_FRAMERATE": "35",
"SELKIES_H264_CRF": "34",
})
if err != nil {
t.Fatalf("marshalEnvironmentOverrides returned error: %v", err)
}
env, err := buildInstancePodEnv(&models.Instance{
ID: 42,
Type: "openclaw",
RuntimeType: RuntimeBackendDesktop,
EnvironmentOverridesJSON: raw,
}, nil, nil, nil)
if err != nil {
t.Fatalf("buildInstancePodEnv returned error: %v", err)
}
if got := env["SELKIES_ENCODER"]; got != "x264enc,jpeg" {
t.Fatalf("SELKIES_ENCODER = %q, want x264enc,jpeg", got)
}
if got := env["SELKIES_USE_CSS_SCALING"]; got != "true" {
t.Fatalf("SELKIES_USE_CSS_SCALING = %q, want true", got)
}
}
func TestPopSHMSizeGB(t *testing.T) {
tests := []struct {
name string
@@ -60,7 +60,7 @@ func TestInstanceProxyServiceUsesRuntimeBindingForV2(t *testing.T) {
}
accessService := NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(45, 123, "openclaw", "/api/v1/instances/123/proxy/", 3000, time.Hour)
token, err := accessService.GenerateToken(45, 123, "openclaw", "/api/v1/instances/123/proxy/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -128,7 +128,7 @@ func TestInstanceProxyServiceInjectsInstanceTokenForHermesLite(t *testing.T) {
}
accessService := NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(45, 127, "hermes", "/api/v1/instances/127/proxy/", 3000, time.Hour)
token, err := accessService.GenerateToken(45, 127, "hermes", "/api/v1/instances/127/proxy/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -193,7 +193,7 @@ func TestInstanceProxyServiceUsesHermesLiteAccessURLForRootEntry(t *testing.T) {
}
accessService := NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(45, 131, "hermes", "/api/v1/instances/131/proxy/chat/", 3000, time.Hour)
token, err := accessService.GenerateToken(45, 131, "hermes", "/api/v1/instances/131/proxy/chat/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -263,7 +263,7 @@ func TestInstanceProxyServicePreservesHermesRuntimeQueryToken(t *testing.T) {
}
accessService := NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(45, 130, "hermes", "/api/v1/instances/130/proxy/chat/", 3000, time.Hour)
token, err := accessService.GenerateToken(45, 130, "hermes", "/api/v1/instances/130/proxy/chat/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -328,7 +328,7 @@ func TestInstanceProxyServiceRewritesHermesLiteHTMLBase(t *testing.T) {
}
accessService := NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(45, 128, "hermes", "/api/v1/instances/128/proxy/", 3000, time.Hour)
token, err := accessService.GenerateToken(45, 128, "hermes", "/api/v1/instances/128/proxy/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -424,7 +424,7 @@ func TestInstanceProxyServiceProxiesHermesLiteWebSocket(t *testing.T) {
}
accessService := NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(45, 129, "hermes", "/api/v1/instances/129/proxy/", 3000, time.Hour)
token, err := accessService.GenerateToken(45, 129, "hermes", "/api/v1/instances/129/proxy/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -641,7 +641,7 @@ func TestInstanceProxyServiceReturnsUnavailableWhenV2BindingMissing(t *testing.T
}
accessService := NewInstanceAccessService()
defer accessService.Stop()
token, err := accessService.GenerateToken(45, 124, "hermes", "/api/v1/instances/124/proxy/", 3000, time.Hour)
token, err := accessService.GenerateToken(45, 124, "hermes", "/api/v1/instances/124/proxy/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -663,7 +663,7 @@ func newV2ProxyTestService(t *testing.T, instanceRepo repository.InstanceReposit
t.Helper()
accessService := NewInstanceAccessService()
t.Cleanup(accessService.Stop)
token, err := accessService.GenerateToken(userID, instanceID, instanceType, "/api/v1/instances/"+strconv.Itoa(instanceID)+"/proxy/", 3000, time.Hour)
token, err := accessService.GenerateToken(userID, instanceID, instanceType, "/api/v1/instances/"+strconv.Itoa(instanceID)+"/proxy/", "", 3000, time.Hour)
if err != nil {
t.Fatalf("GenerateToken returned error: %v", err)
}
@@ -43,6 +43,9 @@ func (s *InstanceDeploymentService) EnsureDeployment(ctx context.Context, config
desired := BuildInstanceDeployment(s.client, config, replicas)
deployments := s.client.Clientset.AppsV1().Deployments(desired.Namespace)
if err := s.deleteLegacyDeployments(ctx, desired.Namespace, config.InstanceID, desired.Name); err != nil {
return nil, err
}
existing, err := deployments.Get(ctx, desired.Name, metav1.GetOptions{})
if errors.IsNotFound(err) {
created, createErr := deployments.Create(ctx, desired, metav1.CreateOptions{})
@@ -75,6 +78,27 @@ func (s *InstanceDeploymentService) EnsureDeployment(ctx context.Context, config
return result, nil
}
func (s *InstanceDeploymentService) deleteLegacyDeployments(ctx context.Context, namespace string, instanceID int, expectedName string) error {
deployments := s.client.Clientset.AppsV1().Deployments(namespace)
items, err := deployments.List(ctx, metav1.ListOptions{
LabelSelector: fmt.Sprintf("instance-id=%d", instanceID),
})
if err != nil {
return fmt.Errorf("failed to list instance deployments %s/instance-id=%d: %w", namespace, instanceID, err)
}
for _, item := range items.Items {
if item.Name == expectedName {
continue
}
propagation := metav1.DeletePropagationForeground
if err := deployments.Delete(ctx, item.Name, metav1.DeleteOptions{PropagationPolicy: &propagation}); err != nil && !errors.IsNotFound(err) {
return fmt.Errorf("failed to delete legacy instance deployment %s/%s: %w", namespace, item.Name, err)
}
}
return nil
}
func (s *InstanceDeploymentService) ScaleDeployment(ctx context.Context, userID, instanceID int, replicas int32) error {
if s.client == nil {
return fmt.Errorf("k8s client not initialized")
@@ -4,6 +4,7 @@ import (
"context"
"testing"
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes/fake"
@@ -120,3 +121,54 @@ func TestInstanceDeploymentServiceEnsureAndScale(t *testing.T) {
t.Fatalf("scaled replicas = %#v, want 0", scaled.Spec.Replicas)
}
}
func TestInstanceDeploymentServiceEnsureDeletesLegacyDeployments(t *testing.T) {
client := &Client{Clientset: fake.NewSimpleClientset(), Namespace: "clawreef"}
service := &InstanceDeploymentService{
client: client,
namespaceService: &NamespaceService{client: client},
}
ctx := context.Background()
namespace := "clawreef-user-9"
if _, err := client.Clientset.AppsV1().Deployments(namespace).Create(ctx, &appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: "clawreef-44-old-name",
Namespace: namespace,
Labels: map[string]string{
"app": "clawreef",
"instance-id": "44",
"managed-by": "clawreef",
},
},
}, metav1.CreateOptions{}); err != nil {
t.Fatalf("create legacy deployment: %v", err)
}
if _, err := service.EnsureDeployment(ctx, PodConfig{
InstanceID: 44,
InstanceName: "Pro Desktop",
UserID: 9,
Type: "openclaw",
RuntimeType: "desktop",
CPUCores: 2,
MemoryGB: 4,
Image: "registry/openclaw:v2",
MountPath: "/config",
ContainerPort: 3001,
}, 1); err != nil {
t.Fatalf("EnsureDeployment returned error: %v", err)
}
deployments, err := client.Clientset.AppsV1().Deployments(namespace).List(ctx, metav1.ListOptions{
LabelSelector: "instance-id=44",
})
if err != nil {
t.Fatalf("list deployments: %v", err)
}
if len(deployments.Items) != 1 {
t.Fatalf("deployment count = %d, want 1: %#v", len(deployments.Items), deployments.Items)
}
if got := deployments.Items[0].Name; got != "clawreef-44-deployment" {
t.Fatalf("remaining deployment = %q, want stable deployment", got)
}
}
@@ -317,7 +317,7 @@ const InstanceDetailPage: React.FC = () => {
const [desktopStreamProfile, setDesktopStreamProfile] =
useState<DesktopStreamProfile>("standard");
const [desktopStreamSavedProfile, setDesktopStreamSavedProfile] =
useState<DesktopStreamProfile>("standard");
useState<DesktopStreamProfile | "">("");
const [desktopStreamMessage, setDesktopStreamMessage] = useState<string | null>(null);
const fetchMeta = useCallback(
@@ -395,9 +395,9 @@ const InstanceDetailPage: React.FC = () => {
}, []);
useEffect(() => {
const profile = instance?.desktop_stream_profile || "standard";
setDesktopStreamProfile(profile);
setDesktopStreamSavedProfile(profile);
const savedProfile = instance?.desktop_stream_profile || "";
setDesktopStreamProfile(savedProfile || "standard");
setDesktopStreamSavedProfile(savedProfile);
setDesktopStreamMessage(null);
}, [instance?.desktop_stream_profile, instance?.id]);