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
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ func TestLiveClaudeSDKCancelResume(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
defer cancel()
out := make(chan proto.Envelope, 64)
request := proto.PromptRequestPayload{RunID: uuid.NewString(), Input: proto.TextInput(prompt), AgentSessionID: resume, ObserveMessages: true, DisableExecutionEnvironment: true, DisableSubagents: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: "MiniMax-M3", ModelProvider: provider, SystemPrompt: "Follow the user's requested format. Preserve the exact verification value in conversation history. Use no tools."}
request := proto.PromptRequestPayload{RunID: uuid.NewString(), Input: proto.TextInput(prompt), AgentSessionID: resume, DisableExecutionEnvironment: true, DisableSubagents: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: "MiniMax-M3", ModelProvider: provider, SystemPrompt: "Follow the user's requested format. Preserve the exact verification value in conversation history. Use no tools."}
running, err := startSingleTurn(ctx, config, request, out)
if err != nil {
t.Fatal(err)
Expand Down Expand Up @@ -148,11 +148,11 @@ func TestLiveClaudeSDKCancelResume(t *testing.T) {
nonce := "cancel-history-" + uuid.NewString()
first := run("Remember this exact verification value: "+nonce+". First repeat it, then write two hundred numbered sentences about trees. Do not use tools.", "", true)
id, _ := first.Outcome.Metadata[proto.DoneMetaAgentSessionID].(string)
if id == "" || first.Outcome.Content == "" || first.Failure == "" {
if id == "" || messageText(first.Events) == "" || first.Failure == "" {
t.Fatal("live cancellation lost identity, partial output or interruption evidence")
}
second := run("Return only the exact cancel-history verification value in the earlier user request. Ignore the earlier request for numbered sentences.", id, false)
if second.Failure != "" || second.Outcome.Metadata[proto.DoneMetaAgentSessionID] != id || !strings.Contains(second.Outcome.Content, nonce) || first.NodePID == second.NodePID {
if second.Failure != "" || second.Outcome.Metadata[proto.DoneMetaAgentSessionID] != id || !strings.Contains(messageText(second.Events), nonce) || first.NodePID == second.NodePID {
t.Fatalf("cold continuation did not preserve identity/history; evidence %s", root)
}
data, _ := json.MarshalIndent(map[string]any{"scope": "private Go factory -> maintained SDK/native -> real MiniMax cancellation and cold continuation; public admission remains separate", "verification_value": nonce, "executions": []evidence{first, second}}, "", " ")
Expand Down
12 changes: 6 additions & 6 deletions apps/daemon/internal/agent/claudesdk/cancellation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ func TestCancellationWaitsForDrainAndPublishesOutcome(t *testing.T) {
t.Fatal("successful cancellation preceded owned process release")
}
got := running.CancellationOutcome()
if got.Content != "partialtaildrained" || got.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || got.Usage.Raw["claude_sdk_result"] == nil || got.Usage.Tokens != nil {
if got.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || got.Usage.Raw["claude_sdk_result"] == nil || got.Usage.Tokens != nil {
t.Fatalf("lost drained cancellation outcome: %+v", got)
}
var done proto.DonePayload
Expand Down Expand Up @@ -108,7 +108,7 @@ func TestFailureKeepsOnlyVerifiedNativeIdentity(t *testing.T) {
t.Fatal("native failure was not reported")
}
if mode == "failure" {
if done.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || done.Content != "partial" || done.Usage.Raw["claude_sdk_result"] == nil {
if done.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || done.Usage.Raw["claude_sdk_result"] == nil {
t.Fatal("verified failure outcome was lost", done)
}
} else if done.Metadata[proto.DoneMetaAgentSessionID] != nil {
Expand Down Expand Up @@ -168,7 +168,7 @@ func cancellationRequest() proto.PromptRequestPayload {

func runCancellationHelper(request startRequest, mode string, scanner *bufio.Scanner, emit func(bridgeEvent)) {
if mode == "cancellation-wait" {
emit(bridgeEvent{Type: "delta", Delta: "partial"})
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "partial"})
if !scanner.Scan() {
return
}
Expand All @@ -179,15 +179,15 @@ func runCancellationHelper(request startRequest, mode string, scanner *bufio.Sca
// These valid observations were in flight when cancellation started.
emit(bridgeEvent{Type: "input_ready", SessionID: request.Resume})
emit(bridgeEvent{Type: "usage", ResultID: "native-result", SessionID: request.Resume, Usage: json.RawMessage(usageFixture)})
emit(bridgeEvent{Type: "delta", Delta: "tail"})
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "tail"})
for {
if _, err := os.Stat(filepath.Join(os.Getenv("CLAUDE_CONFIG_DIR"), "release")); err == nil {
break
}
time.Sleep(time.Millisecond)
}
_, _ = os.Stderr.WriteString(strings.Repeat("x", 2*1024*1024))
emit(bridgeEvent{Type: "delta", Delta: "drained"})
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "drained"})
emit(bridgeEvent{Type: "input_closed", SessionID: request.Resume})
emit(bridgeEvent{Type: "error", Code: "cancelled"})
return
Expand All @@ -199,7 +199,7 @@ func runCancellationHelper(request startRequest, mode string, scanner *bufio.Sca
}
emit(bridgeEvent{Type: "input_ready", SessionID: id})
emit(bridgeEvent{Type: "usage", ResultID: "native-result", SessionID: id, Usage: json.RawMessage(usageFixture)})
emit(bridgeEvent{Type: "delta", Delta: "partial"})
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "partial"})
}
emit(bridgeEvent{Type: "error", Code: "execution_failed"})
}
Original file line number Diff line number Diff line change
Expand Up @@ -175,5 +175,5 @@ func runCommandsHelper(request startRequest, mode string, scanner *bufio.Scanner
return
}
emit(bridgeEvent{Type: "input_closed", SessionID: request.Resume})
emit(bridgeEvent{Type: "result", SessionID: request.Resume, Text: "completed"})
emit(bridgeEvent{Type: "result", SessionID: request.Resume})
}
1 change: 0 additions & 1 deletion apps/daemon/internal/agent/claudesdk/declaration.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@ const claudeSDKNodeEnv = "OAC_RUNTIME_CLAUDE_SDK_NODE"
var Declaration = agent.Declaration{Info: proto.SupportedAgentKind{Kind: "claude_sdk", Capabilities: proto.AgentKindCapabilities{
SubagentObservations: proto.CapabilityUnsupported,
NativeSessionRecovery: proto.CapabilityUnsupported,
MessageItems: proto.CapabilitySupported,
EnvironmentNone: proto.CapabilitySupported,
LocalEnvironment: proto.CapabilityUnsupported,
WorkspaceReadPreparation: proto.CapabilityUnsupported,
Expand Down
2 changes: 1 addition & 1 deletion apps/daemon/internal/agent/claudesdk/declaration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ func TestClaudeSDKFeatureDiscovery(t *testing.T) {

// The declaration must retain the complete baseline capability descriptor.
func TestDeclaredCapabilityBaseline(t *testing.T) {
expected := map[string]bool{"MessageItems": true, "EnvironmentNone": true, "ProgrammaticToolCallingDisable": true, "SubagentControl": true, "FunctionTools": true}
expected := map[string]bool{"EnvironmentNone": true, "ProgrammaticToolCallingDisable": true, "SubagentControl": true, "FunctionTools": true}
value := reflect.ValueOf(Declaration.Info.Capabilities)
for i := 0; i < value.NumField(); i++ {
name := value.Type().Field(i).Name
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ func runClassifiedFailureHelper(request startRequest, mode string, encode func(b
os.Exit(7)
}
if mode == "after-terminal" {
fmt.Fprintln(os.Stdout, `{"type":"delta","delta":"late"}`)
fmt.Fprintln(os.Stdout, `{"type":"delta","item_id":"message","delta":"late"}`)
}
if mode == "scanner-error" {
fmt.Fprintln(os.Stdout, strings.Repeat("x", 2*1024*1024))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ func TestStructuredOutputConfigurationReachesNativeUnchanged(t *testing.T) {
t.Setenv("OAC_RUNTIME_HOME", root)
config := Config{Entrypoint: filepath.Join(root, "worker"), StateDir: filepath.Join(root, "state")}
schema := json.RawMessage(`{"type":"object","properties":{"n":{"const":9007199254740992}}}`)
request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), ObserveMessages: true, DisableSubagents: true, Model: "model", SystemPrompt: "Original instructions.", ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium", OutputFormat: &proto.OutputFormat{Type: "json_schema", Schema: schema}}}
request := proto.PromptRequestPayload{ModelProvider: fixtureProvider(), DisableSubagents: true, Model: "model", SystemPrompt: "Original instructions.", ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium", OutputFormat: &proto.OutputFormat{Type: "json_schema", Schema: schema}}}
start, _, err := prepareConfiguration(config, request)
if err != nil {
t.Fatal(err)
Expand Down
14 changes: 10 additions & 4 deletions apps/daemon/internal/agent/claudesdk/executor_live_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
Settlement agent.TurnSettlement `json:"settlement"`
SettlementError string `json:"settlement_error,omitempty"`
Done proto.DonePayload `json:"done"`
Text string `json:"text"`
Errors []string `json:"errors,omitempty"`
}
evidence := struct {
Expand All @@ -82,7 +83,7 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
_ = os.WriteFile(filepath.Join(proof, "executor-evidence.json"), raw, 0600)
}
defer persist()
request := proto.PromptRequestPayload{DisableExecutionEnvironment: true, DisableSubagents: true, ObserveMessages: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: model, ModelProvider: provider, SystemPrompt: "Follow requested formats briefly. Remember the exact verification marker across the conversation. Use no tools."}
request := proto.PromptRequestPayload{DisableExecutionEnvironment: true, DisableSubagents: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: model, ModelProvider: provider, SystemPrompt: "Follow requested formats briefly. Remember the exact verification marker across the conversation. Use no tools."}
factory := NewExecutorFactory(config)
prepared := time.Now()
owner, err := factory(ctx, request)
Expand Down Expand Up @@ -124,6 +125,11 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
if event.ID != id {
t.Fatal("output crossed Turn identity")
}
if event.Type == proto.TypeDelta {
var delta proto.DeltaPayload
_ = event.DecodePayload(&delta)
record.Text += delta.Delta
}
if event.Type == proto.TypeDelta && record.FirstTextMS == 0 {
record.FirstTextMS = time.Since(started).Milliseconds()
record.NativePIDs = nativeChildren(record.NodePID)
Expand Down Expand Up @@ -165,11 +171,11 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
}
marker := "REUSE-" + uuid.NewString()
first := run("first", "Remember this marker: "+marker+". Reply with exactly the marker.", false)
if first.SettlementError != "" || !first.Settlement.Reusable || !strings.Contains(first.Done.Content, marker) {
if first.SettlementError != "" || !first.Settlement.Reusable || !strings.Contains(first.Text, marker) {
t.Fatal("first Turn failed or was not reusable")
}
second := run("second", "What exact marker did I give you? Reply with only that marker.", false)
if second.SettlementError != "" || !second.Settlement.Reusable || !strings.Contains(second.Done.Content, marker) || first.NodePID != second.NodePID || len(first.NativePIDs) == 0 || !slices.Equal(first.NativePIDs, second.NativePIDs) {
if second.SettlementError != "" || !second.Settlement.Reusable || !strings.Contains(second.Text, marker) || first.NodePID != second.NodePID || len(first.NativePIDs) == 0 || !slices.Equal(first.NativePIDs, second.NativePIDs) {
t.Fatal("ordinary Turns did not retain native execution and history")
}
interrupted := run("cancel", "List the numbers 1 through 10000, one number per line, without stopping early.", true)
Expand All @@ -196,7 +202,7 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
evidence.Recovered = true
}
continued := run("continued", "What exact marker did I originally give you? Reply with only the marker.", false)
if continued.SettlementError != "" || !continued.Settlement.Reusable || !strings.Contains(continued.Done.Content, marker) {
if continued.SettlementError != "" || !continued.Settlement.Reusable || !strings.Contains(continued.Text, marker) {
t.Fatal("history did not continue after cancellation")
}
if !evidence.Recovered && (continued.NodePID != second.NodePID || !slices.Equal(continued.NativePIDs, second.NativePIDs)) {
Expand Down
8 changes: 4 additions & 4 deletions apps/daemon/internal/agent/claudesdk/executor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ func runPersistentExecutorHelper() {
if cancel {
encode(bridgeEvent{Type: "error", TurnID: active, Code: "cancelled"})
} else {
encode(bridgeEvent{Type: "result", TurnID: active, SessionID: "native-persistent", Text: active})
encode(bridgeEvent{Type: "result", TurnID: active, SessionID: "native-persistent"})
}
confirmed, reusable, reason := true, true, ""
switch os.Getenv("SDK_EXECUTOR_MODE") {
Expand Down Expand Up @@ -96,12 +96,12 @@ func runPersistentExecutorHelper() {
}
active = command.TurnID
if os.Getenv("SDK_EXECUTOR_MODE") == "late" && old != "" {
encode(bridgeEvent{Type: "delta", TurnID: old, Delta: "late"})
encode(bridgeEvent{Type: "delta", TurnID: old, ItemID: "message", Delta: "late"})
continue
}
encode(bridgeEvent{Type: "turn_started", TurnID: active})
encode(bridgeEvent{Type: "input_ready", TurnID: active, SessionID: "native-persistent"})
encode(bridgeEvent{Type: "delta", TurnID: active, Delta: "partial"})
encode(bridgeEvent{Type: "delta", TurnID: active, ItemID: "message", Delta: "partial"})
if strings.HasPrefix(os.Getenv("SDK_EXECUTOR_MODE"), "pending_function") {
encode(bridgeEvent{Type: "function_call", TurnID: active, Call: &proto.FunctionCallPayload{CallID: "call", Name: "lookup", Arguments: json.RawMessage("{}")}})
}
Expand All @@ -115,7 +115,7 @@ func runPersistentExecutorHelper() {
continue
}
// A full bridge write has happened, but no native input receipt exists.
encode(bridgeEvent{Type: "delta", TurnID: active, Delta: "steer-written"})
encode(bridgeEvent{Type: "delta", TurnID: active, ItemID: "message", Delta: "steer-written"})
case "turn_cancel":
if active == command.TurnID {
settle(true)
Expand Down
18 changes: 7 additions & 11 deletions apps/daemon/internal/agent/claudesdk/executor_turn.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"fmt"
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
"strings"
"time"
)

Expand Down Expand Up @@ -46,7 +45,6 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
s.invalidate()
}
}
var content strings.Builder
var result *bridgeEvent
var classifiedFailure error
var engineCode, failedResultID string
Expand Down Expand Up @@ -98,17 +96,17 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
s.invalidate()
}
case "delta":
if start.ObserveMessages && event.ItemID == "" || !start.ObserveMessages && event.ItemID != "" {
failure = fmt.Errorf("claudesdk: invalid message delta identity")
sequence++
delta := proto.DeltaPayload{ItemID: event.ItemID, Delta: event.Delta, Sequence: sequence}
if err := delta.Validate(); err != nil {
failure = fmt.Errorf("claudesdk: %w", err)
s.invalidate()
break
}
content.WriteString(event.Delta)
sequence++
emit(proto.TypeDelta, proto.DeltaPayload{ItemID: event.ItemID, Delta: event.Delta, Sequence: sequence})
emit(proto.TypeDelta, delta)
case "output_message":
message := event.Message
if !start.ObserveMessages || message == nil || message.ID == "" ||
if message == nil || message.ID == "" ||
(message.Status != "in_progress" && message.Status != "completed") ||
(message.Status == "completed") != (message.Text != nil) {
failure = fmt.Errorf("claudesdk: invalid message observation")
Expand Down Expand Up @@ -205,11 +203,9 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
metadata[proto.DoneMetaAgentSessionID] = id
}
if failure == nil {
content.Reset()
content.WriteString(result.Text)
metadata[proto.DoneMetaAgentSessionID] = result.SessionID
}
s.outcome = proto.DonePayload{Content: content.String(), Usage: usage, Metadata: metadata}
s.outcome = proto.DonePayload{Usage: usage, Metadata: metadata}
// Publish the observed cancellation outcome before terminal delivery.
if s.settlementErr != nil || !settlementReceived || !settlementConfirmed || !reusable || outputLost || s.process.Context().Err() != nil {
s.owner.retire()
Expand Down
4 changes: 2 additions & 2 deletions apps/daemon/internal/agent/claudesdk/functions_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,7 @@ func runFunctionHelper(request startRequest, mode string, scanner *bufio.Scanner
results[value.CallID] = value
}
if mode == "functions-cancel" {
emit(bridgeEvent{Type: "delta", Delta: "waiting for confirmation"})
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "waiting for confirmation"})
if scanner.Scan() {
emit(bridgeEvent{Type: "error", Code: "cancelled"})
}
Expand All @@ -140,5 +140,5 @@ func runFunctionHelper(request startRequest, mode string, scanner *bufio.Scanner
}
emit(bridgeEvent{Type: "function_applied", CallID: id, DeliveryID: delivery})
}
emit(bridgeEvent{Type: "result", SessionID: request.Resume, Text: "done"})
emit(bridgeEvent{Type: "result", SessionID: request.Resume})
}
4 changes: 2 additions & 2 deletions apps/daemon/internal/agent/claudesdk/live_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -138,7 +138,7 @@ func TestLiveClaudeSDKTextResume(t *testing.T) {
requestStart := len(requests)
mu.Unlock()
out := make(chan proto.Envelope, 64)
request := proto.PromptRequestPayload{RunID: uuid.NewString(), Input: proto.TextInput(prompt), AgentSessionID: resume, ObserveMessages: true, DisableExecutionEnvironment: true, DisableSubagents: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: "MiniMax-M3", ModelProvider: provider, SystemPrompt: "Answer briefly and preserve the exact verification value in the conversation. Use no tools."}
request := proto.PromptRequestPayload{RunID: uuid.NewString(), Input: proto.TextInput(prompt), AgentSessionID: resume, DisableExecutionEnvironment: true, DisableSubagents: true, ExecutionControls: &proto.ExecutionControls{WebSearch: "disabled", TextVerbosity: "medium"}, Model: "MiniMax-M3", ModelProvider: provider, SystemPrompt: "Answer briefly and preserve the exact verification value in the conversation. Use no tools."}
if success != nil {
request.SystemPrompt = "Call lookup exactly once as requested, then report both result parts and any prior verification value. Never retry a failed tool."
request.FunctionTools = []proto.FunctionTool{{Name: "lookup", Description: "Return a synthetic verification value.", Parameters: json.RawMessage(`{"type":"object","properties":{"id":{"type":"string"}},"required":["id"],"additionalProperties":false}`)}}
Expand Down Expand Up @@ -265,7 +265,7 @@ func TestLiveClaudeSDKTextResume(t *testing.T) {
done = true
var payload proto.DonePayload
_ = json.Unmarshal(event.Payload, &payload)
proof.Text = payload.Content
proof.Text = messageText(proof.Events)
proof.SessionID, _ = payload.Metadata[proto.DoneMetaAgentSessionID].(string)
}
}
Expand Down
Loading
Loading