From cfc7ccb5f1ee6a934659252ef5bad2920d345006 Mon Sep 17 00:00:00 2001 From: sam Date: Thu, 8 Oct 2026 21:40:32 +0800 Subject: [PATCH] perf(runtime): integrate concurrent startup into beta --- .../internal/agent/codex/model_verbosity.go | 5 +- .../internal/agent/codex/preparation.go | 7 +- .../agent/codex/preparation_observation.go | 15 + .../codex/preparation_observation_test.go | 54 +++ apps/daemon/internal/agent/codex/rpc.go | 6 +- apps/daemon/internal/cli/connect.go | 53 +-- .../internal/cli/connect_environment.go | 7 +- apps/daemon/internal/cli/connect_startup.go | 50 +++ .../internal/cli/connect_startup_test.go | 316 ++++++++++++++++++ apps/daemon/internal/dispatch/executor.go | 2 + deploy/kubernetes/BETA_CHANGELOG.md | 14 + docs/getting-started/operations.md | 9 + docs/zh/getting-started/operations.md | 13 +- 13 files changed, 520 insertions(+), 31 deletions(-) create mode 100644 apps/daemon/internal/agent/codex/preparation_observation.go create mode 100644 apps/daemon/internal/agent/codex/preparation_observation_test.go create mode 100644 apps/daemon/internal/cli/connect_startup.go create mode 100644 apps/daemon/internal/cli/connect_startup_test.go diff --git a/apps/daemon/internal/agent/codex/model_verbosity.go b/apps/daemon/internal/agent/codex/model_verbosity.go index 6c9d4f28e..01e0d1e82 100644 --- a/apps/daemon/internal/agent/codex/model_verbosity.go +++ b/apps/daemon/internal/agent/codex/model_verbosity.go @@ -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]) diff --git a/apps/daemon/internal/agent/codex/preparation.go b/apps/daemon/internal/agent/codex/preparation.go index 44d750c04..1b4589222 100644 --- a/apps/daemon/internal/agent/codex/preparation.go +++ b/apps/daemon/internal/agent/codex/preparation.go @@ -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" @@ -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") } @@ -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 } @@ -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) diff --git a/apps/daemon/internal/agent/codex/preparation_observation.go b/apps/daemon/internal/agent/codex/preparation_observation.go new file mode 100644 index 000000000..1fdbbba62 --- /dev/null +++ b/apps/daemon/internal/agent/codex/preparation_observation.go @@ -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) +} diff --git a/apps/daemon/internal/agent/codex/preparation_observation_test.go b/apps/daemon/internal/agent/codex/preparation_observation_test.go new file mode 100644 index 000000000..36b9e7063 --- /dev/null +++ b/apps/daemon/internal/agent/codex/preparation_observation_test.go @@ -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") + } +} diff --git a/apps/daemon/internal/agent/codex/rpc.go b/apps/daemon/internal/agent/codex/rpc.go index 2102eda9e..8e99bdfbf 100644 --- a/apps/daemon/internal/agent/codex/rpc.go +++ b/apps/daemon/internal/agent/codex/rpc.go @@ -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 { @@ -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) } @@ -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() diff --git a/apps/daemon/internal/cli/connect.go b/apps/daemon/internal/cli/connect.go index cd990ce30..c9e65aef0 100644 --- a/apps/daemon/internal/cli/connect.go +++ b/apps/daemon/internal/cli/connect.go @@ -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 @@ -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) { @@ -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 @@ -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 { diff --git a/apps/daemon/internal/cli/connect_environment.go b/apps/daemon/internal/cli/connect_environment.go index 55147716f..ac6f9b82d 100644 --- a/apps/daemon/internal/cli/connect_environment.go +++ b/apps/daemon/internal/cli/connect_environment.go @@ -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 ( diff --git a/apps/daemon/internal/cli/connect_startup.go b/apps/daemon/internal/cli/connect_startup.go new file mode 100644 index 000000000..68258f297 --- /dev/null +++ b/apps/daemon/internal/cli/connect_startup.go @@ -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 +} diff --git a/apps/daemon/internal/cli/connect_startup_test.go b/apps/daemon/internal/cli/connect_startup_test.go new file mode 100644 index 000000000..6c26a21f9 --- /dev/null +++ b/apps/daemon/internal/cli/connect_startup_test.go @@ -0,0 +1,316 @@ +package cli + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/auth" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "io" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/transport" +) + +func awaitStartup(t *testing.T, ch <-chan struct{}) { + t.Helper() + select { + case <-ch: + case <-time.After(3 * time.Second): + t.Fatal("startup worker did not reach barrier") + } +} + +func TestConnectionPreflightOverlapsAndJoins(t *testing.T) { + for _, first := range []string{"bootstrap", "discovery"} { + t.Run(first, func(t *testing.T) { + discoveryStarted, bootstrapStarted := make(chan struct{}), make(chan struct{}) + releaseDiscovery, releaseBootstrap := make(chan struct{}), make(chan struct{}) + done := make(chan error, 1) + boot := &transport.BootstrapResponse{DeviceID: "device"} + go func() { + result, _, err := prepareConnection(t.Context(), func(context.Context) (agentCLIDiscovery, error) { + close(discoveryStarted) + <-releaseDiscovery + return agentCLIDiscovery{}, nil + }, func(context.Context) (*transport.BootstrapResponse, error) { + close(bootstrapStarted) + <-releaseBootstrap + return boot, nil + }) + if err == nil && result != boot { + err = errors.New("bootstrap result lost") + } + done <- err + }() + awaitStartup(t, discoveryStarted) + awaitStartup(t, bootstrapStarted) + if first == "bootstrap" { + close(releaseBootstrap) + } else { + close(releaseDiscovery) + } + select { + case <-done: + t.Fatal("returned before both prerequisites completed") + default: + } + if first == "bootstrap" { + close(releaseDiscovery) + } else { + close(releaseBootstrap) + } + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(3 * time.Second): + t.Fatal("join blocked") + } + }) + } +} + +func TestConnectionPreflightFailureCancelsAndJoinsSibling(t *testing.T) { + for _, failing := range []string{"bootstrap", "discovery"} { + t.Run(failing, func(t *testing.T) { + cause := errors.New("rejected") + started, canceled, release := make(chan struct{}), make(chan struct{}), make(chan struct{}) + done := make(chan error, 1) + wait := func(ctx context.Context) error { + close(started) + <-ctx.Done() + close(canceled) + <-release + return ctx.Err() + } + fail := func() error { <-started; return cause } + go func() { + boot, agents, err := prepareConnection(t.Context(), func(ctx context.Context) (agentCLIDiscovery, error) { + if failing == "discovery" { + return nil, fail() + } + return nil, wait(ctx) + }, func(ctx context.Context) (*transport.BootstrapResponse, error) { + if failing == "bootstrap" { + return nil, fail() + } + return nil, wait(ctx) + }) + if boot != nil || agents != nil { + err = errors.New("partial successful result escaped") + } + done <- err + }() + awaitStartup(t, canceled) + select { + case <-done: + t.Fatal("returned before sibling cleanup") + default: + } + close(release) + select { + case err := <-done: + if !errors.Is(err, cause) { + t.Fatal(err) + } + case <-time.After(3 * time.Second): + t.Fatal("join blocked") + } + }) + } +} + +func TestConnectionPreflightParentCancellation(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + started := make(chan struct{}, 2) + done := make(chan error, 1) + go func() { + _, _, err := prepareConnection(ctx, func(ctx context.Context) (agentCLIDiscovery, error) { + started <- struct{}{} + <-ctx.Done() + return nil, ctx.Err() + }, func(ctx context.Context) (*transport.BootstrapResponse, error) { + started <- struct{}{} + <-ctx.Done() + return nil, ctx.Err() + }) + done <- err + }() + awaitStartup(t, started) + awaitStartup(t, started) + cancel() + select { + case err := <-done: + if !errors.Is(err, context.Canceled) { + t.Fatal(err) + } + case <-time.After(3 * time.Second): + t.Fatal("cancellation blocked") + } + _, _, err := prepareConnection(ctx, func(context.Context) (agentCLIDiscovery, error) { + t.Error("discovery started after cancellation") + return nil, nil + }, func(context.Context) (*transport.BootstrapResponse, error) { + t.Error("bootstrap started after cancellation") + return nil, nil + }) + if !errors.Is(err, context.Canceled) { + t.Fatal(err) + } +} + +// Controlled waits model independent work, not production latency. Compare the +// same two operations to demonstrate the removed dependency edge. +func BenchmarkConnectionPreflight(b *testing.B) { + for _, parallel := range []bool{false, true} { + name := "serial" + if parallel { + name = "overlap" + } + b.Run(name, func(b *testing.B) { + discover := func(context.Context) (agentCLIDiscovery, error) { + time.Sleep(5 * time.Millisecond) + return agentCLIDiscovery{}, nil + } + bootstrap := func(context.Context) (*transport.BootstrapResponse, error) { + time.Sleep(5 * time.Millisecond) + return &transport.BootstrapResponse{}, nil + } + for b.Loop() { + if parallel { + _, _, _ = prepareConnection(b.Context(), discover, bootstrap) + } else { + _, _ = discover(b.Context()) + _, _ = bootstrap(b.Context()) + } + } + }) + } +} + +func TestMainLoopOverlapsPreflightWithoutEarlyRegistration(t *testing.T) { + for _, failure := range []string{"", "discovery", "bootstrap"} { + t.Run(failure, func(t *testing.T) { + t.Setenv("OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE", "") + original := harnessDeclarations + defer func() { harnessDeclarations = original }() + started, release, completed := make(chan struct{}), make(chan struct{}), make(chan struct{}) + var probes, boots, dials atomic.Int32 + declaration := original[0] + declaration.Discover = func(ctx context.Context, _ agent.DiscoveryOptions, info proto.SupportedAgentKind) *agent.Runtime { + probes.Add(1) + close(started) + select { + case <-release: + case <-ctx.Done(): + } + defer close(completed) + info.Available = failure != "discovery" && ctx.Err() == nil + return &agent.Runtime{Info: info, Session: func(context.Context, proto.PromptRequestPayload, chan<- proto.Envelope) (agent.Session, error) { + return nil, errors.New("unexpected session") + }} + } + harnessDeclarations = []agent.Declaration{declaration} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/agent-daemon/bootstrap": + boots.Add(1) + select { + case <-started: + case <-r.Context().Done(): + return + } + if dials.Load() != 0 { + t.Error("dial before discovery completed") + } + close(release) + if failure == "bootstrap" { + http.Error(w, "rejected", http.StatusUnauthorized) + return + } + _ = json.NewEncoder(w).Encode(transport.BootstrapResponse{DeviceID: "device"}) + case "/agent-daemon/ws": + select { + case <-completed: + default: + t.Error("registration before discovery joined") + } + dials.Add(1) + http.Error(w, "end test", http.StatusUnauthorized) + default: + t.Errorf("unexpected route %s", r.URL.Path) + http.NotFound(w, r) + } + })) + defer server.Close() + ctx, cancel := context.WithTimeout(t.Context(), 3*time.Second) + defer cancel() + err := mainLoopRemote(ctx, &runContext{stdout: io.Discard, stderr: io.Discard}, "default", auth.Profile{ServerURL: server.URL, RuntimeID: "device", RunnerCredential: "fixture"}, "") + if err == nil || ctx.Err() != nil { + t.Fatalf("unexpected completion: %v, context %v", err, ctx.Err()) + } + expected := int32(0) + if failure == "" { + expected = 1 + } + if probes.Load() != 1 || boots.Load() != 1 || dials.Load() != expected { + t.Fatalf("calls discovery=%d bootstrap=%d dial=%d", probes.Load(), boots.Load(), dials.Load()) + } + }) + } +} + +func TestInlinePairDiscoveryBeforeTokenAndOnlyOnce(t *testing.T) { + for _, available := range []bool{false, true} { + t.Run(fmt.Sprint(available), func(t *testing.T) { + t.Setenv("OAC_RUNTIME_HOME", t.TempDir()) + t.Setenv("OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE", "") + original := harnessDeclarations + defer func() { harnessDeclarations = original }() + var probes, pairs, boots atomic.Int32 + declaration := original[0] + declaration.Discover = func(_ context.Context, _ agent.DiscoveryOptions, info proto.SupportedAgentKind) *agent.Runtime { + probes.Add(1) + info.Available = available + return &agent.Runtime{Info: info} + } + harnessDeclarations = []agent.Declaration{declaration} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/api/v1/runtimes/pair": + if probes.Load() != 1 { + t.Error("pair before discovery") + } + pairs.Add(1) + _ = json.NewEncoder(w).Encode(pairResponse{Runtime: pairRuntime{ID: "device"}, RunnerCredential: "fixture"}) + case "/agent-daemon/bootstrap": + boots.Add(1) + http.Error(w, "end test", http.StatusUnauthorized) + default: + t.Errorf("unexpected request: %s", r.URL.Path) + http.NotFound(w, r) + } + })) + defer server.Close() + err := runConnect(&runContext{stdout: io.Discard, stderr: io.Discard}, []string{"--url", server.URL, "--token", "fixture"}) + if err == nil { + t.Fatal("expected failure") + } + expected := int32(0) + if available { + expected = 1 + } + if probes.Load() != 1 || pairs.Load() != expected || boots.Load() != expected { + t.Fatalf("probes=%d pairs=%d bootstrap=%d", probes.Load(), pairs.Load(), boots.Load()) + } + }) + } +} diff --git a/apps/daemon/internal/dispatch/executor.go b/apps/daemon/internal/dispatch/executor.go index 90544c9a7..b79fa9675 100644 --- a/apps/daemon/internal/dispatch/executor.go +++ b/apps/daemon/internal/dispatch/executor.go @@ -181,7 +181,9 @@ func (r *Router) prepareExecutor(p *preparationState, req proto.PromptRequestPay var native agent.Executor var err error if owner.ctx.Err() == nil { + workspaceStarted := time.Now() req, err = r.localWorkspace.Prepare(owner.ctx, req) + r.log.InfoContext(owner.ctx, "executor preparation stage", "stage", "workspace", "executor_id", owner.id, "session_id", owner.sessionID, "duration_ms", float64(time.Since(workspaceStarted))/float64(time.Millisecond), "success", err == nil) if err == nil && owner.ctx.Err() == nil { native, err = factory(owner.ctx, req) } diff --git a/deploy/kubernetes/BETA_CHANGELOG.md b/deploy/kubernetes/BETA_CHANGELOG.md index 2abb166bd..022eaa9de 100644 --- a/deploy/kubernetes/BETA_CHANGELOG.md +++ b/deploy/kubernetes/BETA_CHANGELOG.md @@ -121,3 +121,17 @@ Validation includes database-backed ordinary/stream creation and input recovery Create validates the managed entry point through SDK file information and a streaming open instead of starting a probe process. The initialization command returns its bounded durable receipt, removing the separate ready-file read on the successful path. Ownership and configuration checks remain mandatory; uncertain responses recover through inspection without replaying initialization. The [helper guide](../../services/core/tools/e2b-provider/README.md#create) owns these behaviors. Validation includes generated-contract checks, 11 template tests, 261 helper tests, SDK transport fixtures, real local subprocess failure/recovery checks and independent review. These checks do not establish production latency savings; compare Create stages and the full cold Session path after rollout. + +## 2026-10-08 — Concurrent Runtime startup and preparation timings + +| Item | Value | +| --- | --- | +| Fork beta baseline | `d5ea46f1e7f931e6f92d91ebb25b37b52e5c5d86` | +| Source PR | [Runtime startup optimization #19](https://github.com/sandbaseai/OpenAgentCore/pull/19) | +| Integration branch | `codex/runtime-startup-beta-qa` | +| Schema migrations / SQL / public protocol / native package pins | None | +| Deployment requirements | Build and qualify a new combined Runtime template; activate its immutable build and deploy Core/Web from the integrated beta revision | + +Harness discovery and authenticated bootstrap overlap while registration requires both to succeed. Failure cancels and joins the sibling. Inline pairing retains validation before consuming its one-shot token and reuses the discovery result. The selected-Harness filter, suspension/wake behavior, native process groups and PR18 E2B receipts remain unchanged. [Operations](../../docs/getting-started/operations.md#runtime-startup-latency) owns the new preparation timings and interpretation. + +The beta adapter retains its existing build-script layout; the main-only builder-path correction is not included. Local race, Runtime contract, translation and Linux compilation checks passed, with independent review. Synthetic wait overlap is not production acceptance. Retain the previous Core/Web images and immutable template selection for rollback; existing sandboxes retain their Runtime until their normal lifecycle ends. diff --git a/docs/getting-started/operations.md b/docs/getting-started/operations.md index 3192faccf..60b456cda 100644 --- a/docs/getting-started/operations.md +++ b/docs/getting-started/operations.md @@ -24,6 +24,15 @@ docker compose -f ~/.oac/core/compose.yaml ps For a second installation, use its directory, such as `~/.oac/second`. +## Runtime startup latency + +Inline pairing validates Harness discovery before consuming its one-shot token and reuses that result. For other connections, after local credentials and enrollment are resolved, Runtime Harness discovery and the authenticated bootstrap HTTP request run concurrently. Both must succeed before the Runtime opens its connection or publishes capabilities. Failure cancels the sibling operation and waits for its cleanup. Reconnect and suspension retain their existing lifecycle; concurrency does not skip executable or credential validation. + +The daemon logs `executor preparation stage` with `stage=workspace`, the executor and Session IDs, duration in milliseconds and `success`. Codex logs `codex preparation stage` for `session_plan`, `model_catalog`, `process_spawn`, `rpc_initialize` and `verification`, with the owner trace, duration and `success`. These records contain no native error text, credentials, configuration, catalog contents or command output. An omitted conditional stage is unobserved, not zero. `session_plan` contains `model_catalog`; the executor readiness interval contains workspace preparation, the adapter stages and transport overhead. Do not add nested intervals together. A failed `rpc_initialize` includes its required child cleanup. + +The model catalog remains validated and pinned before app-server initialization. These timings distinguish catalog preparation from native process initialization; they do not establish a latency improvement. Compare fresh and reused Sessions on the same Runtime template, model and provider, and verify persisted replies and usage in addition to first-text latency. A Runtime startup change needs a rebuilt, qualified Runtime template; replacing Core alone does not update existing sandboxes. + + ## Service health Use these observations for different questions: diff --git a/docs/zh/getting-started/operations.md b/docs/zh/getting-started/operations.md index e5ded119c..efbfb0a03 100644 --- a/docs/zh/getting-started/operations.md +++ b/docs/zh/getting-started/operations.md @@ -1,7 +1,7 @@ --- title: "管理你的安装" source: docs/getting-started/operations.md -source_hash: 9012f887f10171fb2abe27dc49770d429485673559c0d9661bd2100a34adce5f +source_hash: 7e90cbe03d7b6ac8edeb5ea24177b0cb2abc5e79958157e6ed51d8b38b69a8a4 --- 安装运维人员负责 Core 主机、存储和可用性。节点主机运行各自的服务;参阅[节点](nodes.md)。设置见[配置参考](../configuration.md)。 @@ -26,6 +26,17 @@ docker compose -f ~/.oac/core/compose.yaml ps 第二个安装使用自己的目录,例如 `~/.oac/second`。 +## Runtime 启动延迟 {#runtime-startup-latency} + +内联配对在消耗一次性 token 前验证 Harness,并复用发现结果。 + +本地凭据解析和注册绑定完成后,Runtime 的 Harness 探测与带认证的 bootstrap HTTP 请求并行执行。两项都成功后,Runtime 才建立连接并公布能力;任一失败都会取消另一项并等待其清理完成。重连和挂起继续使用原有生命周期,并行执行不会跳过可执行程序或凭据校验。 + +daemon 通过 `executor preparation stage` 记录 `stage=workspace`、executor 和 Session ID、毫秒耗时及 `success`。Codex 通过 `codex preparation stage` 记录 `session_plan`、`model_catalog`、`process_spawn`、`rpc_initialize` 和 `verification`,携带所属请求的 trace、耗时及 `success`。这些记录不包含原生错误文本、凭据、配置、模型目录内容或命令输出。未执行的条件阶段表示未观测,不能按零计算。`session_plan` 包含 `model_catalog`;executor 就绪耗时包含工作区准备、adapter 各阶段和传输开销,不应重复相加。失败的 `rpc_initialize` 包含必要的子进程清理。 + +app-server 初始化前仍会校验并固定模型目录。阶段计时用于区分目录准备与原生进程初始化,本身不证明提速。应在相同 Runtime 模板、模型和 Provider 下对比全新及复用 Session,并同时验证持久化回复、用量和首字延迟。Runtime 启动逻辑变更需要重新构建和验收 Runtime 模板;仅替换 Core 不会更新既有沙箱。 + + ## 服务健康状态 {#service-health} 根据不同问题使用这些观察: