diff --git a/internal/app/bootstrap_test.go b/internal/app/bootstrap_test.go index b63339b..d5cfac6 100644 --- a/internal/app/bootstrap_test.go +++ b/internal/app/bootstrap_test.go @@ -133,6 +133,8 @@ func TestBindFromEnvIfNeededReconnectsExistingAgentTokenEvenWhenSessionAlreadyBo } func TestBindFromEnvIfNeededAcceptsColonSeparatedContainerEnv(t *testing.T) { + t.Setenv(moltenHubTokenEnvVar, "") + t.Setenv(moltenHubRegionEnvVar, "") t.Setenv(moltenHubTokenEnvVar+":t_env-agent-123", "") t.Setenv(moltenHubRegionEnvVar+":"+HubRegionNA, "") @@ -205,6 +207,7 @@ func TestBindFromEnvIfNeededReportsFailure(t *testing.T) { func TestBindFromEnvIfNeededRequiresRegion(t *testing.T) { t.Setenv(moltenHubTokenEnvVar, "b_bind-123") + t.Setenv(moltenHubRegionEnvVar, "") store, err := NewStore(t.TempDir()+"/config.json", DefaultSettings()) if err != nil { diff --git a/internal/app/coverage_additional_test.go b/internal/app/coverage_additional_test.go index 9f73067..3c1cab1 100644 --- a/internal/app/coverage_additional_test.go +++ b/internal/app/coverage_additional_test.go @@ -76,6 +76,41 @@ func TestAdditionalPrimitiveHelpers(t *testing.T) { } } +func TestAdditionalEnvRuntimeHelpers(t *testing.T) { + t.Setenv("APP_DATA_DIR", "") + t.Setenv(moltenHubTokenEnvVar, "") + t.Setenv(moltenHubRegionEnvVar, "") + t.Setenv(moltenHubLocalModeEnvVar, "") + t.Setenv(moltenHubURLEnvVar, "") + t.Setenv(moltenHubAPIBaseEnvVar, "") + + t.Setenv("APP_DATA_DIR:/var/lib/molten", "") + if got, ok := envValue("APP_DATA_DIR"); !ok || got != "/var/lib/molten" { + t.Fatalf("colon APP_DATA_DIR = %q, %v; want /var/lib/molten, true", got, ok) + } + + t.Setenv(moltenHubLocalModeEnvVar, "true") + t.Setenv(moltenHubURLEnvVar, "http://127.0.0.1:8080/root/") + runtime, err, ok := runtimeFromEnv() + if err != nil || !ok { + t.Fatalf("runtimeFromEnv local = %#v, %v, %v; want configured local runtime", runtime, err, ok) + } + if runtime.ID != HubRegionLocal || runtime.HubURL != "http://127.0.0.1:8080" { + t.Fatalf("local runtime = %#v", runtime) + } + + t.Setenv(moltenHubAPIBaseEnvVar, "http://127.0.0.1:9090/api/") + if got := configuredAPIBaseForHub(runtime.HubURL); got != "http://127.0.0.1:9090/api" { + t.Fatalf("configuredAPIBaseForHub local override = %q", got) + } + + t.Setenv(moltenHubLocalModeEnvVar, "false") + t.Setenv(moltenHubRegionEnvVar, "bad-region") + if _, err, ok := runtimeFromEnv(); !ok || err == nil { + t.Fatalf("runtimeFromEnv bad region err=%v ok=%v, want error and configured", err, ok) + } +} + func TestAdditionalNormalizeOnboardingTokens(t *testing.T) { cases := []struct { name string @@ -100,6 +135,34 @@ func TestAdditionalNormalizeOnboardingTokens(t *testing.T) { } } +func TestAdditionalDispatchTextHelpers(t *testing.T) { + if got := textMessageDetail(hub.OpenClawMessage{Payload: map[string]any{"message": " hi "}}); got != "hi" { + t.Fatalf("textMessageDetail map payload = %q, want hi", got) + } + if got := textMessageDetail(hub.OpenClawMessage{Payload: nil, Input: map[string]any{"content": []string{"a", "b"}}}); got != `["a","b"]` { + t.Fatalf("textMessageDetail JSON input = %q", got) + } + if got := textMessageDetail(hub.OpenClawMessage{}); got != "Received text message." { + t.Fatalf("textMessageDetail default = %q", got) + } + if got := textMessageValue(make(chan int)); got == "" { + t.Fatalf("textMessageValue unmarshalable = %q, want fmt fallback", got) + } + + uuid, uri := callerTargetFromMessage(hub.PullResponse{FromAgentUUID: " uuid ", FromAgentURI: " uri "}) + if uuid != "uuid" || uri != "uri" { + t.Fatalf("callerTargetFromMessage explicit = %q, %q", uuid, uri) + } + uuid, uri = callerTargetFromMessage(hub.PullResponse{OpenClawMessage: hub.OpenClawMessage{ReplyTarget: "molten://agent/1"}}) + if uuid != "" || uri != "molten://agent/1" { + t.Fatalf("callerTargetFromMessage URI reply = %q, %q", uuid, uri) + } + uuid, uri = callerTargetFromMessage(hub.PullResponse{OpenClawMessage: hub.OpenClawMessage{ReplyTarget: "uuid-1"}}) + if uuid != "uuid-1" || uri != "" { + t.Fatalf("callerTargetFromMessage UUID reply = %q, %q", uuid, uri) + } +} + func TestAdditionalConnectedAgentHelpers(t *testing.T) { service, _ := newTestService(t) if err := service.AddConnectedAgent(ConnectedAgent{}); err == nil { diff --git a/internal/hub/realtime_test.go b/internal/hub/realtime_test.go index 866a6cf..2f636e3 100644 --- a/internal/hub/realtime_test.go +++ b/internal/hub/realtime_test.go @@ -532,6 +532,96 @@ func TestConnectRuntimeMessagesKeepsIdleConnectionReadingUntilDelivery(t *testin } } +func TestWebsocketSessionAckAndNackResponses(t *testing.T) { + t.Parallel() + + actions := make(chan map[string]any, 2) + server := httptest.NewServer(websocket.Handler(func(conn *websocket.Conn) { + defer conn.Close() + if err := websocket.JSON.Send(conn, map[string]any{"type": "session_ready"}); err != nil { + return + } + for i := 0; i < 2; i++ { + var payload map[string]any + if err := websocket.JSON.Receive(conn, &payload); err != nil { + return + } + actions <- payload + action, _ := payload["type"].(string) + requestID, _ := payload["request_id"].(string) + response := map[string]any{ + "type": "response", + "request_id": requestID, + "ok": action == "ack", + } + if action != "ack" { + response["error"] = map[string]any{"code": "nack_failed", "message": "cannot nack"} + } + if err := websocket.JSON.Send(conn, response); err != nil { + return + } + } + })) + defer server.Close() + + client := NewClient(server.URL) + session, err := client.ConnectRuntimeMessages(context.Background(), "agent-token", "main") + if err != nil { + t.Fatalf("ConnectRuntimeMessages() error = %v", err) + } + defer session.Close() + + if err := session.Ack(context.Background(), "delivery-1"); err != nil { + t.Fatalf("Ack() error = %v", err) + } + if err := session.Nack(context.Background(), "delivery-2"); err == nil || !strings.Contains(err.Error(), "nack_failed") { + t.Fatalf("Nack() error = %v, want nack_failed", err) + } + + first := <-actions + second := <-actions + if first["type"] != "ack" || first["delivery_id"] != "delivery-1" { + t.Fatalf("first action = %#v", first) + } + if second["type"] != "nack" || second["delivery_id"] != "delivery-2" { + t.Fatalf("second action = %#v", second) + } +} + +func TestWebsocketSessionResponseWithoutPendingRequestIsIgnored(t *testing.T) { + t.Parallel() + + server := httptest.NewServer(websocket.Handler(func(conn *websocket.Conn) { + defer conn.Close() + if err := websocket.JSON.Send(conn, map[string]any{"type": "session_ready"}); err != nil { + return + } + _ = websocket.JSON.Send(conn, map[string]any{"type": "response", "request_id": "missing", "ok": true}) + _ = websocket.JSON.Send(conn, map[string]any{"type": "delivery", "result": map[string]any{ + "delivery_id": "delivery-after-response", + "message": map[string]any{"kind": "ack"}, + }}) + })) + defer server.Close() + + client := NewClient(server.URL) + session, err := client.ConnectRuntimeMessages(context.Background(), "agent-token", "main") + if err != nil { + t.Fatalf("ConnectRuntimeMessages() error = %v", err) + } + defer session.Close() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + message, err := session.Receive(ctx) + if err != nil { + t.Fatalf("Receive() error = %v", err) + } + if message.DeliveryID != "delivery-after-response" { + t.Fatalf("delivery_id = %q, want delivery-after-response", message.DeliveryID) + } +} + func websocketAcceptKey(key string) string { sum := sha1.Sum([]byte(strings.TrimSpace(key) + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11")) return base64.StdEncoding.EncodeToString(sum[:])