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
27 changes: 25 additions & 2 deletions 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 All @@ -26,14 +29,30 @@ func prepareModelVerbosity(ctx context.Context, binary string, plan *SessionPlan
}
cmd.Dir = plan.Cwd
cmd.Env = append(os.Environ(), plan.Env...)
commandStarted := time.Now()
catalog, err := cmd.Output()
commandEnded := time.Now()
// The launcher can exit before its children, ending the context watcher.
cleanupStarted := time.Now()
var cleanupErr error
if cmd.Process != nil {
_ = cmd.Cancel()
cleanupErr = cmd.Cancel()
}
cleanupEnded := time.Now()
observePreparationInterval(ctx, "model_catalog_command", commandStarted, commandEnded, err)
// This is the existing best-effort process-group signal, not a second Wait.
if cmd.Process != nil {
observePreparationInterval(ctx, "model_catalog_cleanup", cleanupStarted, cleanupEnded, cleanupErr)
}
if err != nil {
return fmt.Errorf("codex: cannot verify model verbosity support: %w", err)
}
validationStarted := time.Now()
defer func() {
if !validationStarted.IsZero() {
observePreparationStage(ctx, "model_catalog_validation", validationStarted, resultErr)
}
}()
supported, err := catalogSupportsVerbosity(catalog, plan.Model)
if err != nil {
return fmt.Errorf("codex: cannot read model verbosity support: %w", err)
Expand All @@ -49,6 +68,10 @@ func prepareModelVerbosity(ctx context.Context, binary string, plan *SessionPlan
if !filepath.IsAbs(codexHome) {
return fmt.Errorf("codex: missing managed home for model catalog")
}
observePreparationStage(ctx, "model_catalog_validation", validationStarted, nil)
validationStarted = time.Time{}
snapshotStarted := time.Now()
defer func() { observePreparationStage(ctx, "model_catalog_snapshot", snapshotStarted, resultErr) }()
file, err := os.CreateTemp(codexHome, "model-catalog-*.json")
if err != nil {
return err
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,12 +6,13 @@ import (
"fmt"
"os"
"sync"
"time"

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

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 All @@ -34,7 +35,9 @@ func newPreparation(parent context.Context, req proto.PromptRequestPayload, cfg
if err != nil {
return nil, err
}
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 @@ -84,6 +87,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
19 changes: 19 additions & 0 deletions apps/daemon/internal/agent/codex/preparation_observation.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
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) {
observePreparationInterval(ctx, stage, started, time.Now(), err)
}

func observePreparationInterval(ctx context.Context, stage string, started, ended time.Time, err error) {
obslog.Ctx(ctx).Info("codex preparation stage", "stage", stage,
"duration_ms", float64(ended.Sub(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")
}
}
18 changes: 13 additions & 5 deletions apps/daemon/internal/agent/codex/rpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -127,10 +127,11 @@ type deferReplySentinel struct{}
var DeferReply = deferReplySentinel{}

type pendingRequest struct {
method string
resp chan rpcResponse
timer *time.Timer
onResult func(json.RawMessage) error
method string
resp chan rpcResponse
timer *time.Timer
onResult func(json.RawMessage) error
onResponse func(bool)
}

type rpcResponse struct {
Expand Down Expand Up @@ -170,7 +171,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 +181,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, 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 +204,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 Expand Up @@ -370,6 +375,9 @@ func (c *JSONRPCClient) handleResponse(rawID, rawResult, rawError json.RawMessag
if p.timer != nil {
p.timer.Stop()
}
if p.onResponse != nil {
p.onResponse(len(rawError) == 0 || string(rawError) == "null")
}
if len(rawError) > 0 && string(rawError) != "null" {
var errBody JsonRpcError
if err := json.Unmarshal(rawError, &errBody); err != nil {
Expand Down
30 changes: 29 additions & 1 deletion apps/daemon/internal/agent/codex/rpc_request.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ func (c *JSONRPCClient) requestWithResult(ctx context.Context, method string, pa
}

func (c *JSONRPCClient) requestWithTimeout(ctx context.Context, method string, params any, write func(any) error, timeout time.Duration, onResult func(json.RawMessage) error) (json.RawMessage, error) {
started := time.Now()
if !c.Alive() {
return nil, errors.New("codex rpc: client not alive")
}
Expand All @@ -37,12 +38,39 @@ func (c *JSONRPCClient) requestWithTimeout(ctx context.Context, method string, p
resp: make(chan rpcResponse, 1),
onResult: onResult,
}
if method == "initialize" {
type responseObservation struct {
at time.Time
accepted bool
}
observed := make(chan responseObservation, 1)
pending.onResponse = func(success bool) {
// Only capture on the reader: logging must not delay response delivery.
observed <- responseObservation{at: time.Now(), accepted: success}
}
defer func() {
select {
case response := <-observed:
var responseErr error
if !response.accepted {
responseErr = errors.New("native initialize rejection")
}
observePreparationInterval(ctx, "rpc_initialize_response", started, response.at, responseErr)
default:
}
}()
}
c.pendingMu.Lock()
c.pending[id] = pending
c.pendingMu.Unlock()

frame := JsonRpcRequest{JsonRpc: JsonRpcVersion, ID: id, Method: method, Params: params}
if err := write(frame); err != nil {
err = write(frame)
if method == "initialize" {
ended, writeErr := time.Now(), err
defer func() { observePreparationInterval(ctx, "rpc_initialize_write", started, ended, writeErr) }()
}
if err != nil {
c.pendingMu.Lock()
delete(c.pending, id)
c.pendingMu.Unlock()
Expand Down
Loading
Loading