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
5 changes: 4 additions & 1 deletion apps/daemon/internal/agent/codex/model_verbosity.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,11 +8,14 @@ import (
"path/filepath"
"slices"
"strings"
"time"
)

// Validate against the binary's active catalog and use that same snapshot for
// execution. A CLI override alone is silently ignored for unsupported models.
func prepareModelVerbosity(ctx context.Context, binary string, plan *SessionPlan) error {
func prepareModelVerbosity(ctx context.Context, binary string, plan *SessionPlan) (resultErr error) {
started := time.Now()
defer func() { observePreparationStage(ctx, "model_catalog", started, resultErr) }()
args := []string{}
for _, kv := range plan.ExtraConfig {
args = append(args, "-c", kv[0]+"="+kv[1])
Expand Down
7 changes: 6 additions & 1 deletion apps/daemon/internal/agent/codex/preparation.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"fmt"
"os"
"sync"
"time"

"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
Expand Down Expand Up @@ -35,7 +36,7 @@ func newSession(parent context.Context, req proto.PromptRequestPayload, out chan
return prepared.start(parent, runID, prompt, out)
}

func newPreparation(parent context.Context, req proto.PromptRequestPayload, cfg sessionConfig) (*Prepared, error) {
func newPreparation(parent context.Context, req proto.PromptRequestPayload, cfg sessionConfig) (_ *Prepared, resultErr error) {
if req.ExecutionControls != nil && req.ExecutionControls.OutputFormat != nil {
return nil, errors.New("codex: structured output is not qualified")
}
Expand Down Expand Up @@ -64,7 +65,9 @@ func newPreparation(parent context.Context, req proto.PromptRequestPayload, cfg
req.AgentStateKey = effectiveAgentStateKey(req)

req.AgentOptions = executionOptions(req)
planStarted := time.Now()
plan, skillRoots, err := prepareSessionPlan(parent, req, cfg)
observePreparationStage(parent, "session_plan", planStarted, err)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -117,6 +120,8 @@ func newPreparation(parent context.Context, req proto.PromptRequestPayload, cfg
if _, err := rpc.Start(cancelCtx, initParams); err != nil {
return p.preparationFailed(fmt.Errorf("codex: rpc start: %w", err))
}
verificationStarted := time.Now()
defer func() { observePreparationStage(parent, "verification", verificationStarted, resultErr) }()
if req.ExecutionControls != nil && req.ExecutionControls.DisableProgrammaticToolCalling {
if err := verifyProgrammaticToolsDisabled(cancelCtx, rpc); err != nil {
return p.preparationFailed(err)
Expand Down
15 changes: 15 additions & 0 deletions apps/daemon/internal/agent/codex/preparation_observation.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
package codex

import (
"context"
"time"

obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
)

// Callers supply fixed stage names. Never include native error text, environment,
// catalog contents or command output; the owner context carries trace correlation.
func observePreparationStage(ctx context.Context, stage string, started time.Time, err error) {
obslog.Info(ctx, "codex preparation stage", "stage", stage,
"duration_ms", float64(time.Since(started))/float64(time.Millisecond), "success", err == nil)
}
54 changes: 54 additions & 0 deletions apps/daemon/internal/agent/codex/preparation_observation_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package codex

import (
"bytes"
"encoding/json"
"errors"
"log/slog"
"strings"
"testing"
"time"

obslog "github.com/MiniMax-AI/OpenAgentCore/internal/obs/log"
)

func TestPreparationObservationIsCorrelatedAndSecretSafe(t *testing.T) {
var output bytes.Buffer
previous := slog.Default()
slog.SetDefault(slog.New(obslog.NewContextHandler(slog.NewJSONHandler(&output, nil))))
defer slog.SetDefault(previous)
carrier, err := obslog.ParseTraceparent("00-12345678901234567890123456789012-1234567890123456-01")
if err != nil {
t.Fatal(err)
}
ctx := obslog.WithTrace(t.Context(), carrier)
for _, err := range []error{nil, errors.New("fixture-secret-command-output")} {
observePreparationStage(ctx, "model_catalog", time.Now(), err)
}
lines := strings.Split(strings.TrimSpace(output.String()), "\n")
if len(lines) != 2 {
t.Fatal(output.String())
}
for i, line := range lines {
var entry map[string]any
if err := json.Unmarshal([]byte(line), &entry); err != nil {
t.Fatal(err)
}
if entry["trace_id"] != "12345678901234567890123456789012" || entry["stage"] != "model_catalog" || entry["success"] != (i == 0) {
t.Fatal(entry)
}
if entry["duration_ms"].(float64) < 0 {
t.Fatal(entry)
}
for key := range entry {
switch key {
case "time", "level", "msg", "trace_id", "span_id", "stage", "duration_ms", "success":
default:
t.Fatalf("unexpected field %s", key)
}
}
}
if strings.Contains(output.String(), "fixture-secret") {
t.Fatal("native error leaked")
}
}
6 changes: 5 additions & 1 deletion apps/daemon/internal/agent/codex/rpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -170,7 +170,7 @@ func NewJSONRPCClient(cfg JSONRPCConfig) *JSONRPCClient {
// - exec.LookPath / Start failure → returns the spawn error verbatim
// - initialize timeout → kills the child, returns context.DeadlineExceeded
// - JSON-RPC error on initialize → kills the child, returns the error
func (c *JSONRPCClient) Start(ctx context.Context, init InitializeParams) (InitializeResult, error) {
func (c *JSONRPCClient) Start(ctx context.Context, init InitializeParams) (_ InitializeResult, resultErr error) {
args := append([]string{}, c.cfg.ExtraArgs...)
args = append(args, "app-server", "--stdio")
for _, f := range c.cfg.EnableFeatures {
Expand All @@ -180,10 +180,12 @@ func (c *JSONRPCClient) Start(ctx context.Context, init InitializeParams) (Initi
args = append(args, "--disable", f)
}

spawnStarted := time.Now()
process, err := clirunner.Start(clirunner.StartOptions{
Parent: ctx, Binary: c.cfg.Binary, Args: args, Dir: c.cfg.Cwd, Env: c.cfg.Env,
NeedStdin: true, OwnProcessGroup: true, KillTimeout: 250 * time.Millisecond,
})
observePreparationStage(ctx, "process_spawn", spawnStarted, err)
if err != nil {
return InitializeResult{}, fmt.Errorf("codex rpc: spawn %q: %w", c.cfg.Binary, err)
}
Expand All @@ -201,6 +203,8 @@ func (c *JSONRPCClient) Start(ctx context.Context, init InitializeParams) (Initi

initCtx, cancel := context.WithTimeout(ctx, rpcInitTimeout)
defer cancel()
initializeStarted := time.Now()
defer func() { observePreparationStage(ctx, "rpc_initialize", initializeStarted, resultErr) }()
rawResult, err := c.Request(initCtx, "initialize", init)
if err != nil {
_ = c.Close()
Expand Down
53 changes: 32 additions & 21 deletions apps/daemon/internal/cli/connect.go
Original file line number Diff line number Diff line change
Expand Up @@ -136,12 +136,17 @@ func runConnect(ctx *runContext, args []string) error {
return spawnBackground(context.Background(), ctx, *profile, argv, extraEnv)
}

// Self-check before pairing/loading credentials so a machine with
// no supported agent CLI fails before consuming a one-shot token.
initializeRuntimeObservations(ctx)
agentCLIs, err := preflightAgentCLIs(context.Background(), ctx, *profile)
if err != nil {
return err
discover := func(parent context.Context) (agentCLIDiscovery, error) {
return preflightAgentCLIs(parent, ctx, *profile)
}
// Validate before consuming a one-shot pairing token, then reuse that result.
if inlinePair {
discovery, err := discover(context.Background())
if err != nil {
return err
}
discover = func(context.Context) (agentCLIDiscovery, error) { return discovery, nil }
}

var prof auth.Profile
Expand All @@ -154,7 +159,7 @@ func runConnect(ctx *runContext, args []string) error {
}
}

return mainLoop(ctx, *profile, prof, agentCLIs)
return mainLoopRemoteWithDiscovery(context.Background(), ctx, *profile, prof, "", discover)
}

func loadInlineConnectEnv(serverURL, token, deviceName *string) {
Expand Down Expand Up @@ -290,11 +295,17 @@ 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 @@ -309,19 +320,19 @@ func mainLoopRemote(parent context.Context, rc *runContext, profile string, prof
rootCtx, cancel := daemonize.NotifyContext(parent)
defer cancel()

bootstrapStarted := time.Now()
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()
observeRuntimeStartup(rootCtx, "bootstrap", bootstrapStarted, err)
boot, agentCLIs, err := prepareConnection(rootCtx, discover,
func(ctx context.Context) (boot *transport.BootstrapResponse, err error) {
started := time.Now()
defer func() { observeRuntimeStartup(ctx, "bootstrap", started, err) }()
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 {
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 @@ -164,13 +164,8 @@ func runEnvironmentConnect(parent context.Context, rc *runContext, profile strin
if background && !daemonize.IsBackgroundChild() {
return spawnBackground(parent, rc, profile, os.Args, nil)
}
// 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