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
2 changes: 1 addition & 1 deletion services/core/IMPLEMENTATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
19 changes: 0 additions & 19 deletions services/core/cmd/server/executor_connections.go

This file was deleted.

20 changes: 12 additions & 8 deletions services/core/cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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"
)

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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{
Expand All @@ -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)
Expand Down
38 changes: 0 additions & 38 deletions services/core/internal/runtime/gateway.go

This file was deleted.

114 changes: 59 additions & 55 deletions services/core/internal/runtimeenrollment/connection.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
4 changes: 2 additions & 2 deletions services/core/internal/runtimeenrollment/connection_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down
10 changes: 10 additions & 0 deletions services/core/internal/runtimegateway/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
8 changes: 2 additions & 6 deletions services/core/tests/integration/devices_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down Expand Up @@ -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))
Expand Down
Loading