diff --git a/services/core/IMPLEMENTATION.md b/services/core/IMPLEMENTATION.md index 433cf678f..d410e2568 100644 --- a/services/core/IMPLEMENTATION.md +++ b/services/core/IMPLEMENTATION.md @@ -94,7 +94,7 @@ The Worker scans pending inputs with the same scheduling slots, Session locks, d `runtimegateway.LinkAuthority` is Core's [Link](../../docs/sandbox-link-protocol.md) `Authority`. `cmd/server` builds it over `sessionpg.Store` and gives it to one `relay.New`, which the API serves at `/api/v1/sandbox-link`; the Worker revokes through that relay as `execution.Dispatcher.Links`. Every Hello, Open and renewal rereads the database. A Serve credential authenticates only its own resource while the `sandbox_resources` view marks it live, at that resource's current generation: an allocation in `creating` or `running` by its `serve_credential_hash`, or a `sandbox_enrollments` row by its executor key while the key would still authenticate for the Environment. Rotating the key advances the generation of each of its enrollments and the epoch of their Sessions' bound assignments in the same statement, so Opens and renewals for the old generation are refused, its attachments end within one lease and the next bind supersedes the old epoch. A key without an Environment restriction may hold several enrollments, and its holder is trusted for every Environment the key authenticates for: it may Serve any of them, replacing that enrollment's serve peer. Every device is a deployment-scoped agent host and may Attach; its `credential_revision` is the peer's `Revision`. `cmd/server` registers it at startup from `OAC_AGENT_HOST_IDENTITY_FILE`, which advances the revision only for a new credential and never lifts a revocation. An attach grant is the assignment ID, epoch and resource generation followed by their keyed digest under the credential key, so Core stores none. It opens a service only while its assignment is the Session's bound assignment at that epoch, held by the peer's Runtime, and its generation is current. File gets the `world` export, Network gets every destination while the Session's network access is enabled and none otherwise, and each lease lasts one minute. Cleanup of an allocation and a release first commit, which withdraws their authority, then revoke the resource at the relay, and only then destroy the compute or send the release. Every other end of Serve authority, such as a credential revocation or rotation, an Environment expiry or a Session deletion, reaches the relay through the Worker's connection pass: it keeps each live resource's last generation in memory and revokes only on a transition, the previous generation when the generation advances and the last one when the resource leaves the view, because each revocation advances the relay's epoch. The relay and that memory are process-local, so a restarted Core starts both empty. -A self-hosted machine enrolls as its Environment's Link resource: enrollment authorizes the executor key and, under the Session lock, inserts the `sandbox_enrollments` row, where the first key wins; it creates no device and binds nothing. Its connection status is that resource Serving while the enrolled key keeps its authority (`runtimeenrollment.RuntimeConnected`), the same rule as a hosted Environment's. +A self-hosted machine enrolls as its Environment's Link resource: enrollment authorizes the executor key and, under the Session lock, inserts the `sandbox_enrollments` row, where the first key wins; it creates no device and binds nothing. Its connection status is that resource Serving while the enrolled key keeps its authority (`runtimeenrollment.Connections.ExecutorConnected`), the same rule as a hosted Environment's. Connection observations use the execution lease and the Session lock. A separate `environment_connections` row holds the current generation and revision, and `environments.status` commits together with its Session Environment event. The producer serializes replacements and numbers socket observations within each generation; duplicate or older revisions and superseded generations are inert, and a replacement retires the previous connected observation before registering the new one. Registration alone creates no `connected` event. Event payloads carry only public Environment identity, type, status and nullable error, never configuration, credentials, registration IDs or revisions, and have no Turn association. `connected` and `disconnected` are distinct from native readiness: never cast resource `expired` into this vocabulary or emit `ready` for a self-hosted connection. On restart the Worker reconciles old observations before admitting new ones, and a failed or stale observation never establishes a connection. diff --git a/services/core/cmd/server/executor_connections.go b/services/core/cmd/server/executor_connections.go deleted file mode 100644 index 3a0b12b7c..000000000 --- a/services/core/cmd/server/executor_connections.go +++ /dev/null @@ -1,19 +0,0 @@ -package main - -import ( - "context" - - "github.com/MiniMax-AI/OpenAgentCore/internal/sandboxlink/relay" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeenrollment" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" -) - -// executorConnections observes enrolled sandboxes through the Link relay. -type executorConnections struct { - sessions sessions.Reader - links *relay.Relay -} - -func (c executorConnections) ExecutorConnected(ctx context.Context, environment, digest string) (bool, error) { - return runtimeenrollment.RuntimeConnected(ctx, c.sessions, c.links, environment, digest) -} diff --git a/services/core/cmd/server/main.go b/services/core/cmd/server/main.go index 062ae3564..12b6db7ef 100644 --- a/services/core/cmd/server/main.go +++ b/services/core/cmd/server/main.go @@ -58,7 +58,6 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/vaultpg" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/processconfig" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/projects" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimeenrollment" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimehistory" @@ -68,6 +67,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/skills" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/vaults" "github.com/MiniMax-AI/OpenAgentCore/services/core/migrations" + "github.com/go-chi/chi/v5" "github.com/jackc/pgx/v5/pgxpool" ) @@ -211,13 +211,17 @@ func run(config processconfig.Config) error { } executorURL := config.PublicOrigin.DaemonWebSocket() links := runtimegateway.NewLinkAuthority(sessionStore) - daemonHandler, registry, err := runtime.NewGateway(sessionStore, sessionService, links, executorURL) - if err != nil { - return err - } - defer runtime.CloseConnections(registry) + registry := runtimegateway.NewRegistry() + gateway := runtimegateway.NewHandler(runtimegateway.HandlerConfig{ + Authenticator: runtimegateway.NewAuthenticator(sessionStore), Registry: registry, + Heartbeat: sessionService, Links: links, PublicWSURL: executorURL, + }) + daemonHandler := chi.NewRouter() + daemonHandler.Route("/api/v1", func(r chi.Router) { runtimegateway.RegisterRoutes(r, gateway) }) + defer registry.CloseConnections() linkRelay := relay.New(links) defer linkRelay.Close() + connections := &runtimeenrollment.Connections{Store: sessionStore, Links: linkRelay} var catalog *nativeinstaller.Catalog if config.NativeInstallers != "" { catalog, err = nativeinstaller.Load(config.NativeInstallers, buildRevision) @@ -324,7 +328,7 @@ func run(config processconfig.Config) error { Artifacts: sessionService, ArtifactsReader: sessionStore, SessionAdmin: sessionStore, - Environments: sessionService, EnvironmentsReader: sessionStore, ExecutorConnections: executorConnections{sessions: sessionStore, links: linkRelay}, + Environments: sessionService, EnvironmentsReader: sessionStore, ExecutorConnections: connections, Admin: sessionStore, AdminAudit: auditStore, WriteAudit: auditStore, Metrics: metrics, RuntimeObservations: observationService, RuntimeHistory: historyService, Execution: api.Execution{ @@ -350,7 +354,7 @@ func run(config processconfig.Config) error { } handler := serverHandler(apiHandler, &daemonRoutes{gateway: daemonHandler, enrollment: runtimeenrollment.EnrollmentHandler(sessionService, config.PublicOrigin), - connection: runtimeenrollment.ConnectionHandler(sessionStore, linkRelay), + connection: connections, nodeConnect: managedNodes.hub}) server := &http.Server{Addr: config.Addr, Handler: handler, ReadHeaderTimeout: 10 * time.Second, ReadTimeout: 30 * time.Second, WriteTimeout: 30 * time.Second, IdleTimeout: 60 * time.Second} done := make(chan error, 1) diff --git a/services/core/internal/runtime/gateway.go b/services/core/internal/runtime/gateway.go deleted file mode 100644 index 7058db9ac..000000000 --- a/services/core/internal/runtime/gateway.go +++ /dev/null @@ -1,38 +0,0 @@ -// Package runtime connects execution devices without product dependencies. -package runtime - -import ( - "errors" - "net/http" - - "github.com/go-chi/chi/v5" - - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" -) - -// NewGateway serves the V1 agent host transport. It authenticates devices with -// credentials, records their heartbeats and gives binds their Link fields from -// links. Its credentials never grant public Session API access. -func NewGateway(credentials runtimegateway.RuntimeStore, heartbeat runtimegateway.HeartbeatTouch, links *runtimegateway.LinkAuthority, publicWSURL string) (http.Handler, *runtimegateway.Registry, error) { - if credentials == nil || heartbeat == nil || links == nil { - return nil, nil, errors.New("daemon gateway dependencies are required") - } - registry := runtimegateway.NewRegistry() - h := runtimegateway.NewHandler(runtimegateway.HandlerConfig{ - Authenticator: runtimegateway.NewAuthenticator(credentials), Registry: registry, - Heartbeat: heartbeat, Links: links, PublicWSURL: publicWSURL, - }) - r := chi.NewRouter() - r.Route("/api/v1", func(r chi.Router) { runtimegateway.RegisterRoutes(r, h) }) - return r, registry, nil -} - -// CloseConnections releases upgraded WebSockets, which http.Server.Shutdown -// does not close. Call after stopping new HTTP upgrades. -func CloseConnections(registry *runtimegateway.Registry) { - for _, id := range registry.Devices() { - if session, err := registry.LookupDevice(id); err == nil { - session.Close("execution service shutting down") - } - } -} diff --git a/services/core/internal/runtimeenrollment/connection.go b/services/core/internal/runtimeenrollment/connection.go index 1250acb0c..fb016fec3 100644 --- a/services/core/internal/runtimeenrollment/connection.go +++ b/services/core/internal/runtimeenrollment/connection.go @@ -21,89 +21,93 @@ type ConnectionStore interface { GetEnvironmentResource(context.Context, string, string) (runtimedevice.ServeAuthority, error) } -// ConnectionHandler observes an existing enrollment without enrollment or -// execution. Executor authority never grants access to the public Session API. -func ConnectionHandler(s ConnectionStore, links *relay.Relay) http.Handler { - return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - w.Header().Set("Cache-Control", "no-store") - fail := func(status int) { http.Error(w, http.StatusText(status), status) } - if r.Method != http.MethodGet { - w.Header().Set("Allow", http.MethodGet) - fail(http.StatusMethodNotAllowed) - return - } - authorization := strings.Fields(r.Header.Get("Authorization")) - if len(authorization) != 2 || !strings.EqualFold(authorization[0], "Bearer") { - fail(http.StatusUnauthorized) - return - } - query, err := url.ParseQuery(r.URL.RawQuery) - if err != nil || len(query) != 1 || len(query["environment_id"]) != 1 || query.Get("environment_id") == "" { - fail(http.StatusBadRequest) - return - } - environment := query.Get("environment_id") - digest := runtimedevice.HashCredential(authorization[1]) - ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) - defer cancel() - connected, err := RuntimeConnected(ctx, s, links, environment, digest) - switch { - case errors.Is(err, sessions.ErrNotFound): - fail(http.StatusUnauthorized) - case errors.Is(err, sessions.ErrDeviceBindingConflict): - fail(http.StatusConflict) - case err != nil: - fail(http.StatusServiceUnavailable) - default: - status := "disconnected" - if connected { - status = "connected" - } - w.Header().Set("Content-Type", "application/json") - _ = json.NewEncoder(w).Encode(struct { - EnvironmentID string `json:"environment_id"` - Status string `json:"status"` - }{environment, status}) +// Connections observes enrolled sandboxes through their live Link resources. +type Connections struct { + Store ConnectionStore + Links *relay.Relay +} + +// ServeHTTP observes an existing enrollment without enrollment or execution. +// Executor authority never grants access to the public Session API. +func (c *Connections) ServeHTTP(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Cache-Control", "no-store") + fail := func(status int) { http.Error(w, http.StatusText(status), status) } + if r.Method != http.MethodGet { + w.Header().Set("Allow", http.MethodGet) + fail(http.StatusMethodNotAllowed) + return + } + authorization := strings.Fields(r.Header.Get("Authorization")) + if len(authorization) != 2 || !strings.EqualFold(authorization[0], "Bearer") { + fail(http.StatusUnauthorized) + return + } + query, err := url.ParseQuery(r.URL.RawQuery) + if err != nil || len(query) != 1 || len(query["environment_id"]) != 1 || query.Get("environment_id") == "" { + fail(http.StatusBadRequest) + return + } + environment := query.Get("environment_id") + digest := runtimedevice.HashCredential(authorization[1]) + ctx, cancel := context.WithTimeout(r.Context(), 10*time.Second) + defer cancel() + connected, err := c.ExecutorConnected(ctx, environment, digest) + switch { + case errors.Is(err, sessions.ErrNotFound): + fail(http.StatusUnauthorized) + case errors.Is(err, sessions.ErrDeviceBindingConflict): + fail(http.StatusConflict) + case err != nil: + fail(http.StatusServiceUnavailable) + default: + status := "disconnected" + if connected { + status = "connected" } - }) + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(struct { + EnvironmentID string `json:"environment_id"` + Status string `json:"status"` + }{environment, status}) + } } -// RuntimeConnected reports whether the sandbox the executor credential +// ExecutorConnected reports whether the sandbox the executor credential // enrolled serves the Environment: the credential authenticates for the // Environment, the Environment's live Link resource is that enrollment with // the same credential, the relay holds its serve peer, and the credential // still has that authority afterwards. A live resource of another credential // is ErrDeviceBindingConflict. -func RuntimeConnected(ctx context.Context, s ConnectionStore, links *relay.Relay, environment, digest string) (bool, error) { - tenant, err := s.AuthenticateEnvironmentExecutor(ctx, environment, digest) +func (c *Connections) ExecutorConnected(ctx context.Context, environment, digest string) (bool, error) { + tenant, err := c.Store.AuthenticateEnvironmentExecutor(ctx, environment, digest) if err != nil { return false, err } - current, err := s.GetEnvironment(ctx, tenant, environment) + current, err := c.Store.GetEnvironment(ctx, tenant, environment) if err != nil { return false, err } if current.Status == "failed" || current.Status == "expired" { return false, sessions.ErrNotFound } - resource, found, err := enrolledResource(ctx, s, tenant, environment, digest) - if err != nil || !found || !links.Serving(resource.Ref()) { + resource, found, err := c.enrolledResource(ctx, tenant, environment, digest) + if err != nil || !found || !c.Links.Serving(resource.Ref()) { return false, err } // Recheck authority after the relay; rotation or revocation never // inherits the serve peer of the former credential. - if _, err = s.AuthenticateEnvironmentExecutor(ctx, environment, digest); err != nil { + if _, err = c.Store.AuthenticateEnvironmentExecutor(ctx, environment, digest); err != nil { return false, err } - again, found, err := enrolledResource(ctx, s, tenant, environment, digest) - return err == nil && found && again == resource && links.Serving(resource.Ref()), err + again, found, err := c.enrolledResource(ctx, tenant, environment, digest) + return err == nil && found && again == resource && c.Links.Serving(resource.Ref()), err } // enrolledResource reads the Environment's live Link resource and reports // whether it has one; the resource must be an enrollment served with the // credential. -func enrolledResource(ctx context.Context, s ConnectionStore, tenant, environment, digest string) (sandboxbootstrap.Resource, bool, error) { - authority, err := s.GetEnvironmentResource(ctx, tenant, environment) +func (c *Connections) enrolledResource(ctx context.Context, tenant, environment, digest string) (sandboxbootstrap.Resource, bool, error) { + authority, err := c.Store.GetEnvironmentResource(ctx, tenant, environment) if errors.Is(err, sessions.ErrNotFound) { return sandboxbootstrap.Resource{}, false, nil } diff --git a/services/core/internal/runtimeenrollment/connection_test.go b/services/core/internal/runtimeenrollment/connection_test.go index eb36aa14f..cec00bb22 100644 --- a/services/core/internal/runtimeenrollment/connection_test.go +++ b/services/core/internal/runtimeenrollment/connection_test.go @@ -57,7 +57,7 @@ func TestConnectionReadContract(t *testing.T) { req := httptest.NewRequest(tc.method, "/api/v1/agent-daemon/connection?"+tc.query, nil) req.Header.Set("Authorization", tc.bearer) res := httptest.NewRecorder() - ConnectionHandler(s, relay.New(sandboxlinktest.NewAuthority())).ServeHTTP(res, req) + (&Connections{Store: s, Links: relay.New(sandboxlinktest.NewAuthority())}).ServeHTTP(res, req) if res.Code != tc.code || s.calls != tc.calls || res.Header().Get("Cache-Control") != "no-store" { t.Fatalf("%s %s: %d, %d calls", tc.method, tc.query, res.Code, s.calls) } @@ -147,7 +147,7 @@ func TestRuntimeConnectedFollowsServeAndCurrentAuthority(t *testing.T) { } { t.Run(tc.name, func(t *testing.T) { s := tc.store - got, err := RuntimeConnected(t.Context(), &s, srv.Relay, resource.EnvironmentID, digest) + got, err := (&Connections{Store: &s, Links: srv.Relay}).ExecutorConnected(t.Context(), resource.EnvironmentID, digest) if got != tc.want || !errors.Is(err, tc.err) { t.Fatalf("connected=%v err=%v", got, err) } diff --git a/services/core/internal/runtimegateway/registry.go b/services/core/internal/runtimegateway/registry.go index b5984821e..91522a6d2 100644 --- a/services/core/internal/runtimegateway/registry.go +++ b/services/core/internal/runtimegateway/registry.go @@ -237,3 +237,13 @@ func (r *Registry) removeWaiter(deviceID string, ch chan *Session) { r.waiters[deviceID] = filtered } } + +// CloseConnections releases upgraded WebSockets, which http.Server.Shutdown +// does not close. Call after stopping new HTTP upgrades. +func (r *Registry) CloseConnections() { + for _, id := range r.Devices() { + if session, err := r.LookupDevice(id); err == nil { + session.Close("execution service shutting down") + } + } +} diff --git a/services/core/tests/integration/credential_matrix_http_test.go b/services/core/tests/integration/credential_matrix_http_test.go index ad5ce68bd..071fbcb10 100644 --- a/services/core/tests/integration/credential_matrix_http_test.go +++ b/services/core/tests/integration/credential_matrix_http_test.go @@ -48,7 +48,7 @@ func TestCredentialNamespaceMatrix(t *testing.T) { t.Fatal(err) } mux.Handle("/api/v1/agent-daemon/enroll", runtimeenrollment.EnrollmentHandler(sessionService(t, s), origin)) - mux.Handle("/api/v1/agent-daemon/connection", runtimeenrollment.ConnectionHandler(sessionAdapter(s), relay.New(runtimegateway.NewLinkAuthority(sessionAdapter(s))))) + mux.Handle("/api/v1/agent-daemon/connection", &runtimeenrollment.Connections{Store: sessionAdapter(s), Links: relay.New(runtimegateway.NewLinkAuthority(sessionAdapter(s)))}) mux.Handle("/", handler) server := api.CanonicalPaths(mux) call := func(method, path, token, body string) *httptest.ResponseRecorder { diff --git a/services/core/tests/integration/devices_test.go b/services/core/tests/integration/devices_test.go index fc4f498a1..23b74ab20 100644 --- a/services/core/tests/integration/devices_test.go +++ b/services/core/tests/integration/devices_test.go @@ -17,7 +17,6 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/sessionpg" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" ) @@ -109,13 +108,10 @@ func TestStandaloneGatewayUsesExecutionCredentials(t *testing.T) { foreignSecret := registerAgentHost(t, s).Credential server := httptest.NewUnstartedServer(nil) wsURL := "ws://" + server.Listener.Addr().String() + "/api/v1/agent-daemon/ws" - handler, registry, err := runtime.NewGateway(sessionAdapter(s), sessionService(t, s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), wsURL) - if err != nil { - t.Fatal(err) - } + handler, registry := fixtureGateway(t, s, wsURL) server.Config.Handler = handler server.Start() - t.Cleanup(func() { server.Close(); runtime.CloseConnections(registry) }) + t.Cleanup(func() { server.Close(); registry.CloseConnections() }) for _, token := range []string{"", "session-api-key", foreignSecret, secret} { body, _ := json.Marshal(map[string]string{"device_id": a.ID}) req, _ := http.NewRequest(http.MethodPost, server.URL+"/api/v1/agent-daemon/bootstrap", bytes.NewReader(body)) diff --git a/services/core/tests/integration/dispatch_test.go b/services/core/tests/integration/dispatch_test.go index c9eb6eca6..477a62439 100644 --- a/services/core/tests/integration/dispatch_test.go +++ b/services/core/tests/integration/dispatch_test.go @@ -19,13 +19,26 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/modelconfigurationpg" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgtest" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/persistence/postgres/pgunit" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" + "github.com/go-chi/chi/v5" "github.com/google/uuid" "github.com/gorilla/websocket" ) +// fixtureGateway composes the daemon transport as cmd/server does. +func fixtureGateway(t *testing.T, s *Store, publicWSURL string) (http.Handler, *runtimegateway.Registry) { + t.Helper() + registry := runtimegateway.NewRegistry() + h := runtimegateway.NewHandler(runtimegateway.HandlerConfig{ + Authenticator: runtimegateway.NewAuthenticator(sessionAdapter(s)), Registry: registry, + Heartbeat: sessionService(t, s), Links: runtimegateway.NewLinkAuthority(sessionAdapter(s)), PublicWSURL: publicWSURL, + }) + r := chi.NewRouter() + r.Route("/api/v1", func(r chi.Router) { runtimegateway.RegisterRoutes(r, h) }) + return r, registry +} + type dispatchHarness struct { writeMu sync.Mutex assignments map[string]proto.AssignmentRef // by frame and Run ID; writeMu guards it @@ -94,13 +107,10 @@ func newDispatchHarnessForSession(t *testing.T, configuration []byte) *dispatchH } server := httptest.NewUnstartedServer(nil) wsURL := "ws://" + server.Listener.Addr().String() + "/api/v1/agent-daemon/ws" - server.Config.Handler, h.registry, err = runtime.NewGateway(sessionAdapter(s), sessionService(t, s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), wsURL) - if err != nil { - t.Fatal(err) - } + server.Config.Handler, h.registry = fixtureGateway(t, s, wsURL) server.Start() h.url = server.URL - t.Cleanup(func() { server.Close(); runtime.CloseConnections(h.registry) }) + t.Cleanup(func() { server.Close(); h.registry.CloseConnections() }) u, _ := url.Parse(wsURL) u.RawQuery = url.Values{"device_id": {h.device.ID}, "version": {proto.Version}}.Encode() h.conn, _, err = websocket.DefaultDialer.Dial(u.String(), http.Header{"Authorization": {"Bearer " + secret}}) diff --git a/services/core/tests/integration/link_authority_test.go b/services/core/tests/integration/link_authority_test.go index c2c149444..f9a6af576 100644 --- a/services/core/tests/integration/link_authority_test.go +++ b/services/core/tests/integration/link_authority_test.go @@ -31,7 +31,6 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/environmentconfig" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/processconfig" - "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtime" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimedevice" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/runtimegateway" "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sandbox" @@ -541,13 +540,10 @@ func TestInitializationBindsAgentHost(t *testing.T) { resource, serve := fixtureLinkResource(t, s, tenant, session) server := httptest.NewUnstartedServer(nil) endpoint := "ws://" + server.Listener.Addr().String() + "/api/v1/agent-daemon/ws" - handler, registry, err := runtime.NewGateway(sessionAdapter(s), sessionService(t, s), runtimegateway.NewLinkAuthority(sessionAdapter(s)), endpoint) - if err != nil { - t.Fatal(err) - } + handler, registry := fixtureGateway(t, s, endpoint) server.Config.Handler = handler server.Start() - t.Cleanup(func() { server.Close(); runtime.CloseConnections(registry) }) + t.Cleanup(func() { server.Close(); registry.CloseConnections() }) link := startLinkRoute(t, s) runWorker(t, startWorker(t, t.Context(), s, &execution.Dispatcher{Registry: registry, Links: link.Relay})) within(t, startLinkServe(t, link, serve, resource.Ref()).connected) diff --git a/services/core/tests/integration/runtime_enrollment_connection_test.go b/services/core/tests/integration/runtime_enrollment_connection_test.go index 9114aa906..fad330a55 100644 --- a/services/core/tests/integration/runtime_enrollment_connection_test.go +++ b/services/core/tests/integration/runtime_enrollment_connection_test.go @@ -46,7 +46,7 @@ func TestEnrolledSandboxConnectionRevocationAndRestart(t *testing.T) { t.Fatal(err) } link := sandboxlinktest.StartRelay(t, runtimegateway.NewLinkAuthority(sessionAdapter(s))) - connection := runtimeenrollment.ConnectionHandler(sessionAdapter(s), link.Relay) + connection := &runtimeenrollment.Connections{Store: sessionAdapter(s), Links: link.Relay} assertConnection := func(target, token, status string, code int) { t.Helper() request := httptest.NewRequest("GET", "/api/v1/agent-daemon/connection?environment_id="+target, nil)