Skip to content
Merged
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
3 changes: 3 additions & 0 deletions internal/app/bootstrap_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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, "")

Expand Down Expand Up @@ -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 {
Expand Down
63 changes: 63 additions & 0 deletions internal/app/coverage_additional_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 {
Expand Down
90 changes: 90 additions & 0 deletions internal/hub/realtime_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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[:])
Expand Down
Loading