iop/apps/node/internal/adapters/cli/codex_app_server.go
toki 24735a5914 feat: openai workspace agent execution contract implementation
- Update agent contract documentation for openai-compatible API
- Update automation-runtime-bridge milestone tracking
- Add edge smoke tests for openai CLI workspace
- Add node CLI adapters (codex, opencode, oneshot, persistent)
- Add e2e openai-cli-workspace script
- Add agent task tracking for execution contract
2026-06-13 23:12:07 +09:00

643 lines
18 KiB
Go

package cli
import (
"bufio"
"context"
"encoding/json"
"fmt"
"io"
"os/exec"
"sync"
"sync/atomic"
"time"
"iop/apps/node/internal/runtime"
"iop/packages/go/config"
)
// codexAppServerSession tracks a single long-lived `codex app-server --stdio` process
// bound to a (target, sessionID) key. The process is started lazily on first use and
// reused across Execute calls until it exits or is terminated.
type codexAppServerSession struct {
key sessionKey
threadID string // Codex thread id assigned by thread/start
mu sync.Mutex
proc *codexAppServerProc
closed bool
}
// codexAppServerProc wraps the OS process and its stdio pipes.
type codexAppServerProc struct {
cmd *exec.Cmd
stdin io.WriteCloser
stdout *bufio.Scanner
nextID atomic.Int64
pending sync.Map // int64 -> chan appServerResponse
notifCh chan appServerNotification
done chan struct{}
closeOnce sync.Once
}
type appServerRequest struct {
ID int64 `json:"id"`
Method string `json:"method"`
Params map[string]any `json:"params"`
}
type appServerResponse struct {
ID int64 `json:"id"`
Result map[string]any `json:"result,omitempty"`
Error *appServerErr `json:"error,omitempty"`
}
type appServerErr struct {
Code int `json:"code"`
Message string `json:"message"`
}
type appServerNotification struct {
Method string `json:"method"`
Params map[string]any `json:"params,omitempty"`
}
// codexAppServerExecutor is the sessionReporter + executor for modeCodexAppServer profiles.
type codexAppServerExecutor struct {
cli *CLI
mu sync.Mutex
sessions map[sessionKey]*codexAppServerSession
}
func (e *codexAppServerExecutor) Execute(ctx context.Context, spec runtime.ExecutionSpec, profile config.CLIProfileConf, sink runtime.EventSink) error {
sess, err := e.resolveCodexAppServerSession(ctx, spec, profile)
if err != nil {
return emitReturnedError(ctx, sink, spec.RunID, err)
}
sess.mu.Lock()
defer sess.mu.Unlock()
if err := sink.Emit(ctx, runtime.RuntimeEvent{
RunID: spec.RunID,
Type: runtime.EventTypeStart,
Timestamp: time.Now(),
}); err != nil {
return err
}
prompt := extractPrompt(spec.Input)
turnID, err := codexAppServerTurnStart(ctx, sess.proc, sess.threadID, prompt)
if err != nil {
return err
}
return codexAppServerDrainNotifications(ctx, sess.proc, spec.RunID, sess.threadID, turnID, sink)
}
// codexAppServerTurnStart sends turn/start with the prompt and returns the turn id.
func codexAppServerTurnStart(ctx context.Context, p *codexAppServerProc, threadID, prompt string) (string, error) {
// TurnStartParams.input is required; UserInput text variant requires text_elements.
params := map[string]any{
"threadId": threadID,
"input": []map[string]any{
{"type": "text", "text": prompt, "text_elements": []any{}},
},
}
resp, err := p.send(ctx, "turn/start", params)
if err != nil {
return "", fmt.Errorf("turn/start: %w", err)
}
// Actual shape: result.turn.id (TurnStartResponse)
turnID := nestedString(resp.Result, "turn", "id")
// Legacy/flat fallbacks.
if turnID == "" {
turnID, _ = resp.Result["turnId"].(string)
}
if turnID == "" {
turnID, _ = resp.Result["id"].(string)
}
if turnID == "" {
return "", fmt.Errorf("turn/start: server did not return a turn id")
}
return turnID, nil
}
// codexAppServerDrainNotifications consumes notifications until turn/completed or error.
// When the process exits (p.done closes), we drain any remaining buffered notifications
// before surfacing the exit error, so that a terminal event sent just before exit is
// not lost.
func codexAppServerDrainNotifications(ctx context.Context, p *codexAppServerProc, runID, threadID, turnID string, sink runtime.EventSink) error {
for {
select {
case <-ctx.Done():
return ctx.Err()
case notif, ok := <-p.notifCh:
if !ok {
return fmt.Errorf("codex app-server notification channel closed")
}
events, done, err := decodeAppServerNotification(notif, runID, threadID, turnID)
for _, ev := range events {
if emitErr := sink.Emit(ctx, ev); emitErr != nil {
return emitErr
}
}
if err != nil {
return err
}
if done {
return nil
}
case <-p.done:
// Process exited. Drain any notifications already in the buffer before
// declaring an error — the terminal event may have arrived with the exit.
for {
select {
case notif, ok := <-p.notifCh:
if !ok {
return fmt.Errorf("codex app-server process exited during turn")
}
events, done, err := decodeAppServerNotification(notif, runID, threadID, turnID)
for _, ev := range events {
if emitErr := sink.Emit(ctx, ev); emitErr != nil {
return emitErr
}
}
if err != nil {
return err
}
if done {
return nil
}
default:
return fmt.Errorf("codex app-server process exited during turn")
}
}
}
}
}
// decodeAppServerNotification maps a single app-server notification to RuntimeEvents.
// threadID and turnID are the execution-context values used as fallbacks when the
// notification params do not carry them. Returns (events, turnDone, error).
func decodeAppServerNotification(notif appServerNotification, runID, threadID, turnID string) ([]runtime.RuntimeEvent, bool, error) {
// Prefer IDs from the notification params; fall back to execution-context values.
effectiveThread := notif.Params["threadId"]
if s, ok := effectiveThread.(string); ok && s != "" {
threadID = s
}
effectiveTurn := notif.Params["turnId"]
if s, ok := effectiveTurn.(string); ok && s != "" {
turnID = s
}
meta := map[string]string{"source": "codex-app-server"}
if threadID != "" {
meta["thread_id"] = threadID
}
if turnID != "" {
meta["turn_id"] = turnID
}
switch notif.Method {
case "item/agentMessage/delta":
// AgentMessageDeltaNotification: { threadId, turnId, itemId, delta }
delta := nestedString(notif.Params, "delta")
if delta == "" {
return nil, false, nil
}
return []runtime.RuntimeEvent{{
RunID: runID,
Type: runtime.EventTypeDelta,
Delta: delta,
Metadata: meta,
Timestamp: time.Now(),
}}, false, nil
case "item/completed":
// Suppress duplicate final delta — already streamed via delta notifications.
return nil, false, nil
case "turn/completed":
// TurnCompletedNotification: { threadId, turn: { id, status, error, ... } }
// status may be "completed" | "interrupted" | "failed" | "inProgress"
status := nestedString(notif.Params, "turn", "status")
if status == "failed" {
msg := nestedString(notif.Params, "turn", "error", "message")
if msg == "" {
msg = "turn failed"
}
return []runtime.RuntimeEvent{{
RunID: runID,
Type: runtime.EventTypeError,
Error: msg,
Metadata: meta,
Timestamp: time.Now(),
}}, false, fmt.Errorf("codex app-server turn failed: %s", msg)
}
return []runtime.RuntimeEvent{{
RunID: runID,
Type: runtime.EventTypeComplete,
Metadata: meta,
Timestamp: time.Now(),
}}, true, nil
case "error":
// ErrorNotification: { error: { message, ... }, willRetry, threadId, turnId }
msg := nestedString(notif.Params, "error", "message")
if msg == "" {
// flat fallback for non-standard shapes
msg = nestedString(notif.Params, "message")
}
if msg == "" {
msg = "codex app-server error"
}
return []runtime.RuntimeEvent{{
RunID: runID,
Type: runtime.EventTypeError,
Error: msg,
Metadata: meta,
Timestamp: time.Now(),
}}, false, fmt.Errorf("codex app-server: %s", msg)
case "turn/failed":
// Legacy method name kept for compatibility.
msg := nestedString(notif.Params, "error", "message")
if msg == "" {
msg = nestedString(notif.Params, "message")
}
if msg == "" {
msg = "turn failed"
}
return []runtime.RuntimeEvent{{
RunID: runID,
Type: runtime.EventTypeError,
Error: msg,
Metadata: meta,
Timestamp: time.Now(),
}}, false, fmt.Errorf("codex app-server turn failed: %s", msg)
}
return nil, false, nil
}
// nestedString walks a params map by keys and returns the string value at the path.
func nestedString(params map[string]any, keys ...string) string {
var cur any = params
for i, k := range keys {
m, ok := cur.(map[string]any)
if !ok {
return ""
}
cur = m[k]
if i == len(keys)-1 {
s, _ := cur.(string)
return s
}
}
return ""
}
func (e *codexAppServerExecutor) resolveCodexAppServerSession(ctx context.Context, spec runtime.ExecutionSpec, profile config.CLIProfileConf) (*codexAppServerSession, error) {
target := cliTargetName(spec)
key := newSessionKey(spec)
e.mu.Lock()
if sess, ok := e.sessions[key]; ok {
e.mu.Unlock()
return sess, nil
}
if spec.SessionMode == runtime.SessionModeRequireExisting {
e.mu.Unlock()
return nil, fmt.Errorf("cli adapter: no codex app-server session for target %q session %q", target, key.sessionID)
}
e.mu.Unlock()
proc, err := startCodexAppServerProc(ctx, profile, spec.Workspace)
if err != nil {
return nil, fmt.Errorf("cli adapter: start codex app-server: %w", err)
}
threadID, err := codexAppServerInit(ctx, proc)
if err != nil {
proc.close()
return nil, fmt.Errorf("cli adapter: codex app-server init: %w", err)
}
sess := &codexAppServerSession{
key: key,
threadID: threadID,
proc: proc,
}
e.mu.Lock()
if existing, ok := e.sessions[key]; ok {
e.mu.Unlock()
proc.close()
return existing, nil
}
e.sessions[key] = sess
e.mu.Unlock()
return sess, nil
}
func (e *codexAppServerExecutor) Sessions() []sessionListEntry {
e.mu.Lock()
defer e.mu.Unlock()
snaps := make([]sessionListEntry, 0, len(e.sessions))
for k := range e.sessions {
snaps = append(snaps, sessionListEntry{modeCodexAppServer, k.target, k.sessionID, k.workspace})
}
return snaps
}
// Terminate closes every workspace variant of the given target/sessionID.
func (e *codexAppServerExecutor) Terminate(_ context.Context, target, sessionID string) (bool, error) {
sid := normalizeSessionID(sessionID)
e.mu.Lock()
var matched []*codexAppServerSession
for k, sess := range e.sessions {
if k.target == target && k.sessionID == sid {
matched = append(matched, sess)
delete(e.sessions, k)
}
}
e.mu.Unlock()
if len(matched) == 0 {
return false, nil
}
for _, sess := range matched {
closeCodexAppServerSession(sess)
}
return true, nil
}
func (e *codexAppServerExecutor) Stop(_ context.Context) error {
e.mu.Lock()
cp := make(map[sessionKey]*codexAppServerSession, len(e.sessions))
for k, v := range e.sessions {
cp[k] = v
}
e.sessions = make(map[sessionKey]*codexAppServerSession)
e.mu.Unlock()
for _, sess := range cp {
closeCodexAppServerSession(sess)
}
return nil
}
// --- process lifecycle ---
func startCodexAppServerProc(ctx context.Context, profile config.CLIProfileConf, workspace string) (*codexAppServerProc, error) {
args := append(append([]string{}, profile.Args...), "--stdio")
cmd := exec.Command(profile.Command, args...)
if len(profile.Env) > 0 {
cmd.Env = append(cmd.Environ(), profile.Env...)
}
dir, err := prepareWorkspaceDir(workspace)
if err != nil {
return nil, err
}
if dir != "" {
cmd.Dir = dir
}
stdinPipe, err := cmd.StdinPipe()
if err != nil {
return nil, fmt.Errorf("stdin pipe: %w", err)
}
stdoutPipe, err := cmd.StdoutPipe()
if err != nil {
return nil, fmt.Errorf("stdout pipe: %w", err)
}
// stderr is intentionally discarded; process exit errors surface via done channel.
cmd.Stderr = nil
if err := cmd.Start(); err != nil {
return nil, fmt.Errorf("start: %w", err)
}
p := &codexAppServerProc{
cmd: cmd,
stdin: stdinPipe,
stdout: bufio.NewScanner(stdoutPipe),
notifCh: make(chan appServerNotification, 32),
done: make(chan struct{}),
}
p.stdout.Buffer(make([]byte, 64*1024), 4*1024*1024)
go p.readLoop()
go func() {
_ = cmd.Wait()
p.closeOnce.Do(func() { close(p.done) })
}()
return p, nil
}
// readLoop dispatches incoming newline-delimited JSON to pending response channels
// or the notification channel.
func (p *codexAppServerProc) readLoop() {
for p.stdout.Scan() {
line := p.stdout.Bytes()
if len(line) == 0 {
continue
}
// Try response first (has "id" field as number).
var resp appServerResponse
if err := json.Unmarshal(line, &resp); err == nil && resp.ID != 0 {
if ch, ok := p.pending.Load(resp.ID); ok {
select {
case ch.(chan appServerResponse) <- resp:
default:
}
}
continue
}
// Treat as notification. Blocking send ensures no notification is dropped;
// the drain goroutine or a closed done channel will unblock this if needed.
var notif appServerNotification
if err := json.Unmarshal(line, &notif); err == nil && notif.Method != "" {
select {
case p.notifCh <- notif:
case <-p.done:
return
}
}
}
p.closeOnce.Do(func() { close(p.done) })
}
// recvResponse waits for a response on ch, giving priority to an already-arrived
// response over p.done or ctx cancellation. This prevents the case where the
// process sends its response and then immediately exits, and Go's random select
// picks p.done instead of ch.
func (p *codexAppServerProc) recvResponse(ctx context.Context, ch chan appServerResponse) (appServerResponse, error) {
// Fast path: response already in buffer.
select {
case resp := <-ch:
return resp, nil
default:
}
// Slow path: wait for any of the three signals.
select {
case resp := <-ch:
return resp, nil
case <-p.done:
// One last drain: response may have arrived between the default check and here.
select {
case resp := <-ch:
return resp, nil
default:
}
return appServerResponse{}, fmt.Errorf("codex app-server process exited")
case <-ctx.Done():
return appServerResponse{}, ctx.Err()
}
}
func (p *codexAppServerProc) send(ctx context.Context, method string, params map[string]any) (appServerResponse, error) {
id := p.nextID.Add(1)
if params == nil {
params = map[string]any{}
}
req := appServerRequest{ID: id, Method: method, Params: params}
raw, err := json.Marshal(req)
if err != nil {
return appServerResponse{}, fmt.Errorf("marshal request: %w", err)
}
raw = append(raw, '\n')
ch := make(chan appServerResponse, 1)
p.pending.Store(id, ch)
defer p.pending.Delete(id)
// Write stdin in a goroutine so that a blocked pipe write does not prevent
// us from observing process exit or an early response via p.done / ctx.
writeErr := make(chan error, 1)
go func() {
_, err := p.stdin.Write(raw)
writeErr <- err
}()
// Wait for write completion, but also watch for an early response or exit.
// A fast server may respond and exit before the write goroutine reports back.
select {
case resp, ok := <-ch:
// Response arrived before write confirmed — still valid.
if ok {
if resp.Error != nil {
return resp, fmt.Errorf("rpc error %d: %s", resp.Error.Code, resp.Error.Message)
}
return resp, nil
}
default:
}
select {
case err := <-writeErr:
if err != nil {
return appServerResponse{}, fmt.Errorf("write stdin: %w", err)
}
case resp := <-ch:
if resp.Error != nil {
return resp, fmt.Errorf("rpc error %d: %s", resp.Error.Code, resp.Error.Message)
}
return resp, nil
case <-p.done:
// Check for response that arrived simultaneously with exit.
select {
case resp := <-ch:
if resp.Error != nil {
return resp, fmt.Errorf("rpc error %d: %s", resp.Error.Code, resp.Error.Message)
}
return resp, nil
default:
}
return appServerResponse{}, fmt.Errorf("codex app-server process exited")
case <-ctx.Done():
return appServerResponse{}, ctx.Err()
}
resp, err := p.recvResponse(ctx, ch)
if err != nil {
return appServerResponse{}, err
}
if resp.Error != nil {
return resp, fmt.Errorf("rpc error %d: %s", resp.Error.Code, resp.Error.Message)
}
return resp, nil
}
// sendNotification sends a JSON-RPC notification (no id, no response expected).
func (p *codexAppServerProc) sendNotification(method string, params map[string]any) error {
type notifMsg struct {
Method string `json:"method"`
Params map[string]any `json:"params,omitempty"`
}
raw, err := json.Marshal(notifMsg{Method: method, Params: params})
if err != nil {
return fmt.Errorf("marshal notification: %w", err)
}
raw = append(raw, '\n')
_, err = p.stdin.Write(raw)
return err
}
func (p *codexAppServerProc) close() {
_ = p.stdin.Close()
if p.cmd != nil && p.cmd.Process != nil {
_ = p.cmd.Process.Kill()
}
p.closeOnce.Do(func() { close(p.done) })
}
// codexAppServerInit runs initialize → initialized → thread/start and returns the thread id.
func codexAppServerInit(ctx context.Context, p *codexAppServerProc) (string, error) {
initParams := map[string]any{
"protocolVersion": "2024-11-05",
"clientInfo": map[string]any{"name": "iop-node", "version": "0"},
"capabilities": map[string]any{},
}
if _, err := p.send(ctx, "initialize", initParams); err != nil {
return "", fmt.Errorf("initialize: %w", err)
}
if err := p.sendNotification("initialized", nil); err != nil {
return "", fmt.Errorf("initialized notification: %w", err)
}
resp, err := p.send(ctx, "thread/start", map[string]any{})
if err != nil {
return "", fmt.Errorf("thread/start: %w", err)
}
// Actual shape: result.thread.id (ThreadStartResponse)
threadID := nestedString(resp.Result, "thread", "id")
// Legacy/flat fallbacks for compatibility with older or fake servers.
if threadID == "" {
threadID, _ = resp.Result["threadId"].(string)
}
if threadID == "" {
threadID, _ = resp.Result["id"].(string)
}
if threadID == "" {
return "", fmt.Errorf("thread/start: server did not return a thread id")
}
return threadID, nil
}
func closeCodexAppServerSession(sess *codexAppServerSession) {
if sess == nil {
return
}
sess.mu.Lock()
defer sess.mu.Unlock()
if sess.closed {
return
}
sess.closed = true
if sess.proc != nil {
sess.proc.close()
}
}