Skip to content
Draft
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
42 changes: 20 additions & 22 deletions apps/daemon/internal/cli/connect.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ func runConnect(ctx *runContext, args []string) error {
if *remote != "" || *environment != "" || *credentialFile != "" || fs.NArg() != 0 {
return errors.New("connect: bootstrap input cannot be combined with enrollment options")
}
bootstrapped, err = bootstrapProfile(*bootstrapFile)
bootstrapped, err = bootstrapProfile(*bootstrapFile, ctx)
if err != nil {
return err
}
Expand Down Expand Up @@ -98,13 +98,6 @@ func runConnect(ctx *runContext, args []string) error {
return spawnBackground(context.Background(), ctx, *profile, os.Args)
}

// Self-check before loading credentials so a machine with no
// supported agent CLI fails fast.
agentCLIs, err := preflightAgentCLIs(context.Background(), ctx, *profile)
if err != nil {
return err
}

var prof auth.Profile
if bootstrapped != nil {
prof = *bootstrapped
Expand All @@ -115,7 +108,7 @@ func runConnect(ctx *runContext, args []string) error {
}
}

return mainLoop(ctx, *profile, prof, agentCLIs)
return mainLoop(ctx, *profile, prof)
}

// spawnBackground forks the daemon into the background. Parent
Expand Down Expand Up @@ -186,11 +179,16 @@ func spawnBackground(ctx context.Context, rc *runContext, profile string, argv [
// background process. SIGINT / SIGTERM cancels the root context, which
// unblocks the read pump and any in-flight Send so the daemon exits
// without orphaning agent subprocesses.
func mainLoop(rc *runContext, profile string, prof auth.Profile, agentCLIs agentCLIDiscovery) error {
return mainLoopRemote(context.Background(), rc, profile, prof, agentCLIs, "")
func mainLoop(rc *runContext, profile string, prof auth.Profile) error {
return mainLoopRemote(context.Background(), rc, profile, prof, "")
}

func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof auth.Profile, agentCLIs agentCLIDiscovery, remote string) error {
func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof auth.Profile, remote string) error {
return mainLoopRemoteWithDiscovery(parent, rc, profile, prof, remote, func(ctx context.Context) (agentCLIDiscovery, error) {
return preflightAgentCLIs(ctx, rc, profile)
})
}
func mainLoopRemoteWithDiscovery(parent context.Context, rc *runContext, profile string, prof auth.Profile, remote string, discover func(context.Context) (agentCLIDiscovery, error)) error {
// Route through obs/log so daemon log lines pick up the same
// trace_id / span_id auto-injection as the server side — when the
// daemon adopts an envelope's trace, every log call under that ctx
Expand All @@ -205,18 +203,18 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof
rootCtx, cancel := daemonize.NotifyContext(parent)
defer cancel()

bootCtx, bootCancel := context.WithTimeout(rootCtx, bootstrapTimeout)
var boot *transport.BootstrapResponse
var err error
if remote == "" {
boot, err = transport.Bootstrap(bootCtx, prof.ServerURL, prof.RuntimeID, prof.RunnerCredential, Version)
} else {
boot, err = environmentBootstrap(bootCtx, prof, remote)
}
bootCancel()
boot, agentCLIs, err := prepareConnection(rootCtx, discover, func(ctx context.Context) (*transport.BootstrapResponse, error) {
bootCtx, stop := context.WithTimeout(ctx, bootstrapTimeout)
defer stop()
if remote != "" {
return environmentBootstrap(bootCtx, prof, remote)
}
return transport.Bootstrap(bootCtx, prof.ServerURL, prof.RuntimeID, prof.RunnerCredential, Version)
})
if err != nil {
return fmt.Errorf("connect: bootstrap: %w", err)
return err
}

wsURL, err := transport.DeriveWSURL(*boot, prof.ServerURL)
if err != nil {
return fmt.Errorf("connect: derive ws url: %w", err)
Expand Down
3 changes: 2 additions & 1 deletion apps/daemon/internal/cli/connect_bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import (

// The launch file is the sole credential source for this connection. Reopening
// it on process restart never reads or overwrites an auth profile.
func bootstrapProfile(path string) (*auth.Profile, error) {
func bootstrapProfile(path string, rc *runContext) (*auth.Profile, error) {
raw, err := runtimefs.ReadPrivatePath(path, runtimebootstrap.MaxBytes)
if err != nil {
return nil, errors.New("connect: Runtime bootstrap file unavailable")
Expand All @@ -18,5 +18,6 @@ func bootstrapProfile(path string) (*auth.Profile, error) {
if err != nil {
return nil, err
}
rc.installedKinds = map[string]bool{input.Harness: true}
return &auth.Profile{ServerURL: input.CoreURL, RuntimeID: input.DeviceID, RunnerCredential: input.Credential}, nil
}
12 changes: 8 additions & 4 deletions apps/daemon/internal/cli/connect_bootstrap_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ func TestBootstrapConnectionDoesNotReadOrOverwritePrivateProfile(t *testing.T) {
if err = os.WriteFile(filepath.Join(profileDir, "auth.json"), priorRaw, 0600); err != nil {
t.Fatal(err)
}
input := runtimebootstrap.Connection{Version: runtimebootstrap.Version, CoreURL: "https://core.example/api/v1", DeviceID: "da912024-1543-4242-a2c1-5f4f7ebbc6c7", Credential: "bootstrap-secret"}
input := runtimebootstrap.Connection{Version: runtimebootstrap.Version, CoreURL: "https://core.example/api/v1", DeviceID: "da912024-1543-4242-a2c1-5f4f7ebbc6c7", Credential: "bootstrap-secret", Harness: "codex"}
raw, err := input.Marshal()
if err != nil {
t.Fatal(err)
Expand All @@ -41,7 +41,11 @@ func TestBootstrapConnectionDoesNotReadOrOverwritePrivateProfile(t *testing.T) {
t.Fatal(err)
}
for range 2 {
p, err := bootstrapProfile(path)
rc := &runContext{}
p, err := bootstrapProfile(path, rc)
if !rc.installedKinds["codex"] || len(rc.installedKinds) != 1 {
t.Fatal("bootstrap did not scope discovery")
}
if err != nil || p.ServerURL != input.CoreURL || p.RuntimeID != input.DeviceID || p.RunnerCredential != input.Credential {
t.Fatal("failed bootstrap/restart", err)
}
Expand All @@ -53,13 +57,13 @@ func TestBootstrapConnectionDoesNotReadOrOverwritePrivateProfile(t *testing.T) {
if err = os.WriteFile(path, []byte("bootstrap-secret"), 0600); err != nil {
t.Fatal(err)
}
if _, err = bootstrapProfile(path); err == nil || strings.Contains(err.Error(), input.Credential) {
if _, err = bootstrapProfile(path, &runContext{}); err == nil || strings.Contains(err.Error(), input.Credential) {
t.Fatal("invalid input fell back or leaked")
}
if err = os.Remove(path); err != nil {
t.Fatal(err)
}
if _, err = bootstrapProfile(path); err == nil {
if _, err = bootstrapProfile(path, &runContext{}); err == nil {
t.Fatal("missing input fell back to private auth")
}
}
Expand Down
7 changes: 1 addition & 6 deletions apps/daemon/internal/cli/connect_environment.go
Original file line number Diff line number Diff line change
Expand Up @@ -163,13 +163,8 @@ func runEnvironmentConnect(parent context.Context, rc *runContext, profile strin
if background && !daemonize.IsBackgroundChild() {
return spawnBackground(parent, rc, profile, os.Args)
}
// Discovery consumes the immutable Runtime binding; it must follow enrollment.
discovery, err := preflightAgentCLIs(parent, rc, profile)
if err != nil {
return err
}
prof := auth.Profile{ServerURL: base, RuntimeID: bound.DeviceID, RunnerCredential: credential}
return rejected(mainLoopRemote(parent, rc, profile, prof, discovery, remote))
return rejected(mainLoopRemote(parent, rc, profile, prof, remote))
}

var (
Expand Down
50 changes: 50 additions & 0 deletions apps/daemon/internal/cli/connect_startup.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
package cli

import (
"context"
"errors"
"fmt"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport"
)

// prepareConnection overlaps independent preflight work after local credentials
// and enrollment are resolved. Both workers are joined before returning, including
// on failure. Their temporary context never owns the live connection or executor.
func prepareConnection(parent context.Context, discover func(context.Context) (agentCLIDiscovery, error), bootstrap func(context.Context) (*transport.BootstrapResponse, error)) (*transport.BootstrapResponse, agentCLIDiscovery, error) {
if err := parent.Err(); err != nil {
return nil, nil, err
}
ctx, cancel := context.WithCancel(parent)
defer cancel()
var agents agentCLIDiscovery
var boot *transport.BootstrapResponse
done := make(chan error, 2)
go func() {
var err error
agents, err = discover(ctx)
done <- err
}()
go func() {
var err error
boot, err = bootstrap(ctx)
if err != nil {
err = fmt.Errorf("connect: bootstrap: %w", err)
}
done <- err
}()
var result error
for range 2 {
if err := <-done; err != nil {
result = errors.Join(result, err)
cancel()
}
}
if result != nil {
return nil, nil, result
}
if err := parent.Err(); err != nil {
return nil, nil, err
}
return boot, agents, nil
}
Loading
Loading