diff --git a/historyserver/pkg/historyserver/enter_cluster_test.go b/historyserver/pkg/historyserver/enter_cluster_test.go new file mode 100644 index 00000000000..7bfd9090b04 --- /dev/null +++ b/historyserver/pkg/historyserver/enter_cluster_test.go @@ -0,0 +1,175 @@ +package historyserver + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "sort" + "testing" + + "github.com/emicklei/go-restful/v3" + "github.com/ray-project/kuberay/historyserver/pkg/utils" +) + +func TestEnterCluster(t *testing.T) { + // Reset default container to avoid polluting or using duplicate services across runs + restful.DefaultContainer = restful.NewContainer() + + // Create ServerHandler with fake sessionLoader + handler := &ServerHandler{ + maxClusters: 100, + clustersMap: make(map[utils.ClusterKey][]utils.ClusterInfo), + } + + // Setup fake processor for sessionLoader + fp := &fakeProcessor{ + fn: func(ctx context.Context, info utils.ClusterInfo) (SessionStatus, error) { + if info.SessionName == "session_2026-04-22_10-00-00_000000_1" { + return SessionStatusProcessed, nil + } + if info.SessionName == "session_2026-04-22_10-00-00_000000_2_live" { + return SessionStatusLive, nil + } + return SessionStatusEventsErr, fmt.Errorf("unknown session") + }, + } + handler.sessionLoader = NewSessionLoader(fp, context.Background(), DefaultSessionProcessTimeout) + + // Single session cluster + keyA := utils.ClusterKey{ + Namespace: "default", + Name: "cluster-a", + } + handler.clustersMap[keyA] = []utils.ClusterInfo{ + { + Namespace: "default", + Name: "cluster-a", + SessionName: "session_2026-04-22_10-00-00_000000_1", + OwnerKind: "RayJob", + OwnerName: "job-a", + }, + } + + // Multi-session cluster (past session AND live session) + keyB := utils.ClusterKey{ + Namespace: "default", + Name: "cluster-b", + } + handler.clustersMap[keyB] = []utils.ClusterInfo{ + { + Namespace: "default", + Name: "cluster-b", + SessionName: "session_2026-04-22_10-00-00_000000_1", + OwnerKind: "RayService", + OwnerName: "svc-b", + CreateTimeStamp: 1000, // Older + }, + { + Namespace: "default", + Name: "cluster-b", + SessionName: "live", + OwnerKind: "RayService", + OwnerName: "svc-b", + CreateTimeStamp: 2000, // Newer (latest) + }, + } + + // Cluster with session that is resolved but will trigger live resolution + keyC := utils.ClusterKey{ + Namespace: "default", + Name: "cluster-c", + } + handler.clustersMap[keyC] = []utils.ClusterInfo{ + { + Namespace: "default", + Name: "cluster-c", + SessionName: "session_2026-04-22_10-00-00_000000_2_live", + OwnerKind: "RayJob", + OwnerName: "job-c", + }, + } + + // Cluster with invalid session format + keyD := utils.ClusterKey{ + Namespace: "default", + Name: "cluster-d", + } + handler.clustersMap[keyD] = []utils.ClusterInfo{ + { + Namespace: "default", + Name: "cluster-d", + SessionName: "invalid-session-name", + }, + } + + // Explicitly sort `Multi-session cluster` to simulate listClusters post-sorting logic + sort.Sort(utils.ClusterInfoList(handler.clustersMap[keyB])) + + // Register actual router + routerRayClusterSet(handler) + + container := restful.DefaultContainer + + t.Run("Verify sorted order of multi-session slice puts latest first", func(t *testing.T) { + sessions := handler.clustersMap[keyB] + if len(sessions) != 2 { + t.Fatalf("Expected 2 sessions, got %d", len(sessions)) + } + // Index 0 must be the newer session (live, timestamp 2000) + if sessions[0].SessionName != "live" { + t.Errorf("Expected latest session 'live' at index 0, got %s", sessions[0].SessionName) + } + }) + + t.Run("Enter existing single-session cluster with explicit session (Successful Dead Session Loading)", func(t *testing.T) { + req := httptest.NewRequest("GET", "/enter_cluster/default/cluster-a/session_2026-04-22_10-00-00_000000_1", nil) + resp := httptest.NewRecorder() + container.ServeHTTP(resp, req) + + if resp.Code != http.StatusOK { + t.Fatalf("Expected status 200, got %d: %s", resp.Code, resp.Body.String()) + } + + cookies := resp.Result().Cookies() + cookieMap := make(map[string]*http.Cookie) + for _, cookie := range cookies { + cookieMap[cookie.Name] = cookie + } + + if c, ok := cookieMap[COOKIE_SESSION_NAME_KEY]; !ok || c.Value != "session_2026-04-22_10-00-00_000000_1" { + t.Errorf("Expected cookie %s to be 'session_2026-04-22_10-00-00_000000_1', got %v", COOKIE_SESSION_NAME_KEY, c) + } + }) + + t.Run("Enter cluster with session that is resolved but triggers live resolution (Line 326 path)", func(t *testing.T) { + req := httptest.NewRequest("GET", "/enter_cluster/default/cluster-c/session_2026-04-22_10-00-00_000000_2_live", nil) + resp := httptest.NewRecorder() + container.ServeHTTP(resp, req) + + if resp.Code != http.StatusOK { + t.Fatalf("Expected status 200, got %d: %s", resp.Code, resp.Body.String()) + } + + cookies := resp.Result().Cookies() + cookieMap := make(map[string]*http.Cookie) + for _, cookie := range cookies { + cookieMap[cookie.Name] = cookie + } + + // Since the fake processor returns SessionStatusLive for this session, resolvedSession is set to "live" + if c, ok := cookieMap[COOKIE_SESSION_NAME_KEY]; !ok || c.Value != "live" { + t.Errorf("Expected cookie %s to be 'live' (resolved from timestamp because cluster is live), got %v", COOKIE_SESSION_NAME_KEY, c) + } + }) + + t.Run("Enter cluster with invalid session name format (Line 310 path)", func(t *testing.T) { + req := httptest.NewRequest("GET", "/enter_cluster/default/cluster-d/invalid-session-name", nil) + resp := httptest.NewRecorder() + container.ServeHTTP(resp, req) + + if resp.Code != http.StatusBadRequest { + t.Fatalf("Expected status 400 (BadRequest), got %d", resp.Code) + } + }) +} diff --git a/historyserver/pkg/historyserver/reader.go b/historyserver/pkg/historyserver/reader.go index 89e64b62875..db866417497 100644 --- a/historyserver/pkg/historyserver/reader.go +++ b/historyserver/pkg/historyserver/reader.go @@ -67,9 +67,47 @@ func (s *ServerHandler) listClusters(limit int) []utils.ClusterInfo { clusters = clusters[:limit] } clusters = append(liveClusterInfos, clusters...) + + clustersMap := make(map[utils.ClusterKey][]utils.ClusterInfo) + for _, c := range clusters { + key := utils.ClusterKey{ + Namespace: c.Namespace, + Name: c.Name, + } + clustersMap[key] = append(clustersMap[key], c) + } + + for key := range clustersMap { + sort.Sort(utils.ClusterInfoList(clustersMap[key])) + } + + s.mu.Lock() + s.clustersMap = clustersMap + s.mu.Unlock() + return clusters } +func (s *ServerHandler) findSessionInMap(namespace, name, session string) (string, bool) { + s.mu.RLock() + defer s.mu.RUnlock() + key := utils.ClusterKey{ + Namespace: namespace, + Name: name, + } + if list, ok := s.clustersMap[key]; ok { + if len(list) == 0 { + return "", false + } + for _, c := range list { + if c.SessionName == session { + return c.SessionName, true + } + } + } + return "", false +} + func (s *ServerHandler) _getNodeLogs(rayClusterNameNamespace, sessionId, nodeId, folder, glob string) ([]byte, error) { logPath := path.Join(sessionId, utils.RAY_SESSIONDIR_LOGDIR_NAME, nodeId) if folder != "" { diff --git a/historyserver/pkg/historyserver/router.go b/historyserver/pkg/historyserver/router.go index 2f074b5b1e1..b7870fd9f6d 100644 --- a/historyserver/pkg/historyserver/router.go +++ b/historyserver/pkg/historyserver/router.go @@ -294,21 +294,30 @@ func routerRayClusterSet(s *ServerHandler) { defer restful.Add(ws) ws.Path("/enter_cluster").Consumes(restful.MIME_JSON).Produces(restful.MIME_JSON).Filter(RequestLogFilter) - ws.Route(ws.GET("/{namespace}/{name}/{session}").To(func(r1 *restful.Request, r2 *restful.Response) { - name := r1.PathParameter("name") - namespace := r1.PathParameter("namespace") - session := r1.PathParameter("session") + enterHandler := func(r1 *restful.Request, r2 *restful.Response, namespace, name, session string) { + resolvedSession, found := s.findSessionInMap(namespace, name, session) + if !found { + if s.clientManager != nil && s.reader != nil { + s.listClusters(s.maxClusters) + } + resolvedSession, found = s.findSessionInMap(namespace, name, session) + } - if session != "live" { - if ParseSessionTimestamp(session).IsZero() { - logrus.Warnf("Rejecting invalid session name: %s/%s/%s", namespace, name, session) - r2.WriteErrorString(http.StatusBadRequest, fmt.Sprintf("invalid session name: %q", session)) + if !found { + r2.WriteErrorString(http.StatusNotFound, fmt.Sprintf("cluster %s/%s with session %s not found", namespace, name, session)) + return + } + + if resolvedSession != "live" { + if ParseSessionTimestamp(resolvedSession).IsZero() { + logrus.Warnf("Rejecting invalid session name: %s/%s/%s", namespace, name, resolvedSession) + r2.WriteErrorString(http.StatusBadRequest, fmt.Sprintf("invalid session name: %q", resolvedSession)) return } - info := utils.ClusterInfo{Name: name, Namespace: namespace, SessionName: session} + info := utils.ClusterInfo{Name: name, Namespace: namespace, SessionName: resolvedSession} live, err := s.sessionLoader.LoadSession(r1.Request.Context(), info) if err != nil { - logrus.Errorf("Failed to load session %s/%s/%s: %v", namespace, name, session, err) + logrus.Errorf("Failed to load session %s/%s/%s: %v", namespace, name, resolvedSession, err) r2.WriteErrorString(http.StatusInternalServerError, err.Error()) return } @@ -316,21 +325,27 @@ func routerRayClusterSet(s *ServerHandler) { // Users might use the complete session name to enter a live cluster, // so "live" sentinel is set to avoid querying empty in-memory state. if live { - session = "live" + resolvedSession = "live" } } - // Set cookies only after a successful load (dead) or a skip (live). http.SetCookie(r2, &http.Cookie{MaxAge: 600, Path: "/", Name: COOKIE_CLUSTER_NAME_KEY, Value: name}) http.SetCookie(r2, &http.Cookie{MaxAge: 600, Path: "/", Name: COOKIE_CLUSTER_NAMESPACE_KEY, Value: namespace}) - http.SetCookie(r2, &http.Cookie{MaxAge: 600, Path: "/", Name: COOKIE_SESSION_NAME_KEY, Value: session}) + http.SetCookie(r2, &http.Cookie{MaxAge: 600, Path: "/", Name: COOKIE_SESSION_NAME_KEY, Value: resolvedSession}) r2.WriteJson(map[string]interface{}{ "result": "success", "name": name, "namespace": namespace, - "session": session, + "session": resolvedSession, }, "application/json") + } + + ws.Route(ws.GET("/{namespace}/{name}/{session}").To(func(r1 *restful.Request, r2 *restful.Response) { + name := r1.PathParameter("name") + namespace := r1.PathParameter("namespace") + session := r1.PathParameter("session") + enterHandler(r1, r2, namespace, name, session) }). Doc("set cookie for cluster"). Param(ws.PathParameter("namespace", "namespace")). diff --git a/historyserver/pkg/historyserver/server.go b/historyserver/pkg/historyserver/server.go index bf5aa22e2b3..8fa7d18ba52 100644 --- a/historyserver/pkg/historyserver/server.go +++ b/historyserver/pkg/historyserver/server.go @@ -5,11 +5,13 @@ import ( "fmt" "log" "net/http" + "sync" "time" "github.com/ray-project/kuberay/historyserver/pkg/collector/types" "github.com/ray-project/kuberay/historyserver/pkg/eventserver" "github.com/ray-project/kuberay/historyserver/pkg/storage" + "github.com/ray-project/kuberay/historyserver/pkg/utils" "github.com/sirupsen/logrus" "k8s.io/client-go/transport" ) @@ -26,6 +28,9 @@ type ServerHandler struct { httpClient *http.Client useKubernetesProxy bool + + mu sync.RWMutex + clustersMap map[utils.ClusterKey][]utils.ClusterInfo } func NewServerHandler(c *types.RayHistoryServerConfig, dashboardDir string, reader storage.StorageReader, clientManager *ClientManager, eventHandler *eventserver.EventHandler, sessionLoader *SessionLoader, useKubernetesProxy bool) (*ServerHandler, error) { diff --git a/historyserver/pkg/utils/types.go b/historyserver/pkg/utils/types.go index d979ba0e7a2..03223c67fd2 100644 --- a/historyserver/pkg/utils/types.go +++ b/historyserver/pkg/utils/types.go @@ -10,6 +10,11 @@ type ClusterInfo struct { OwnerName string `json:"ownerName,omitempty"` } +type ClusterKey struct { + Namespace string + Name string +} + type ClusterInfoList []ClusterInfo func (a ClusterInfoList) Len() int { return len(a) }