Skip to content

Commit 1df24db

Browse files
authored
Require message identity on every assistant message (#524)
* Require message identity on every assistant message Every adapter now delivers each assistant message as an identified output_message start, deltas that name it, and a completed snapshot. The shared validator rejects a delta without item_id, so Core keeps one Item path: the legacy message projection, the native-message query, the observe_messages request field, the message_items capability and DonePayload.Content are deleted. MiniMax Code completes each ACP message by its native messageId when the next one starts or the Turn ends. The agent-host qualification rig binds the Session's assignment before it prepares, which execution admission requires, and reads the answer from completed message Items. * Reject a resumed mcode message and drop a dead merge branch
1 parent df44168 commit 1df24db

132 files changed

Lines changed: 605 additions & 712 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎apps/daemon/internal/agent/claudesdk/cancellation_live_linux_test.go‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,7 @@ func TestLiveClaudeSDKCancelResume(t *testing.T) {
6666
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
6767
defer cancel()
6868
out := make(chan proto.Envelope, 64)
69-
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."}
69+
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."}
7070
running, err := startSingleTurn(ctx, config, request, out)
7171
if err != nil {
7272
t.Fatal(err)
@@ -148,11 +148,11 @@ func TestLiveClaudeSDKCancelResume(t *testing.T) {
148148
nonce := "cancel-history-" + uuid.NewString()
149149
first := run("Remember this exact verification value: "+nonce+". First repeat it, then write two hundred numbered sentences about trees. Do not use tools.", "", true)
150150
id, _ := first.Outcome.Metadata[proto.DoneMetaAgentSessionID].(string)
151-
if id == "" || first.Outcome.Content == "" || first.Failure == "" {
151+
if id == "" || messageText(first.Events) == "" || first.Failure == "" {
152152
t.Fatal("live cancellation lost identity, partial output or interruption evidence")
153153
}
154154
second := run("Return only the exact cancel-history verification value in the earlier user request. Ignore the earlier request for numbered sentences.", id, false)
155-
if second.Failure != "" || second.Outcome.Metadata[proto.DoneMetaAgentSessionID] != id || !strings.Contains(second.Outcome.Content, nonce) || first.NodePID == second.NodePID {
155+
if second.Failure != "" || second.Outcome.Metadata[proto.DoneMetaAgentSessionID] != id || !strings.Contains(messageText(second.Events), nonce) || first.NodePID == second.NodePID {
156156
t.Fatalf("cold continuation did not preserve identity/history; evidence %s", root)
157157
}
158158
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}}, "", " ")

‎apps/daemon/internal/agent/claudesdk/cancellation_test.go‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ func TestCancellationWaitsForDrainAndPublishesOutcome(t *testing.T) {
5959
t.Fatal("successful cancellation preceded owned process release")
6060
}
6161
got := running.CancellationOutcome()
62-
if got.Content != "partialtaildrained" || got.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || got.Usage.Raw["claude_sdk_result"] == nil || got.Usage.Tokens != nil {
62+
if got.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || got.Usage.Raw["claude_sdk_result"] == nil || got.Usage.Tokens != nil {
6363
t.Fatalf("lost drained cancellation outcome: %+v", got)
6464
}
6565
var done proto.DonePayload
@@ -108,7 +108,7 @@ func TestFailureKeepsOnlyVerifiedNativeIdentity(t *testing.T) {
108108
t.Fatal("native failure was not reported")
109109
}
110110
if mode == "failure" {
111-
if done.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || done.Content != "partial" || done.Usage.Raw["claude_sdk_result"] == nil {
111+
if done.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || done.Usage.Raw["claude_sdk_result"] == nil {
112112
t.Fatal("verified failure outcome was lost", done)
113113
}
114114
} else if done.Metadata[proto.DoneMetaAgentSessionID] != nil {
@@ -168,7 +168,7 @@ func cancellationRequest() proto.PromptRequestPayload {
168168

169169
func runCancellationHelper(request startRequest, mode string, scanner *bufio.Scanner, emit func(bridgeEvent)) {
170170
if mode == "cancellation-wait" {
171-
emit(bridgeEvent{Type: "delta", Delta: "partial"})
171+
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "partial"})
172172
if !scanner.Scan() {
173173
return
174174
}
@@ -179,15 +179,15 @@ func runCancellationHelper(request startRequest, mode string, scanner *bufio.Sca
179179
// These valid observations were in flight when cancellation started.
180180
emit(bridgeEvent{Type: "input_ready", SessionID: request.Resume})
181181
emit(bridgeEvent{Type: "usage", ResultID: "native-result", SessionID: request.Resume, Usage: json.RawMessage(usageFixture)})
182-
emit(bridgeEvent{Type: "delta", Delta: "tail"})
182+
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "tail"})
183183
for {
184184
if _, err := os.Stat(filepath.Join(os.Getenv("CLAUDE_CONFIG_DIR"), "release")); err == nil {
185185
break
186186
}
187187
time.Sleep(time.Millisecond)
188188
}
189189
_, _ = os.Stderr.WriteString(strings.Repeat("x", 2*1024*1024))
190-
emit(bridgeEvent{Type: "delta", Delta: "drained"})
190+
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "drained"})
191191
emit(bridgeEvent{Type: "input_closed", SessionID: request.Resume})
192192
emit(bridgeEvent{Type: "error", Code: "cancelled"})
193193
return
@@ -199,7 +199,7 @@ func runCancellationHelper(request startRequest, mode string, scanner *bufio.Sca
199199
}
200200
emit(bridgeEvent{Type: "input_ready", SessionID: id})
201201
emit(bridgeEvent{Type: "usage", ResultID: "native-result", SessionID: id, Usage: json.RawMessage(usageFixture)})
202-
emit(bridgeEvent{Type: "delta", Delta: "partial"})
202+
emit(bridgeEvent{Type: "delta", ItemID: "message", Delta: "partial"})
203203
}
204204
emit(bridgeEvent{Type: "error", Code: "execution_failed"})
205205
}

‎apps/daemon/internal/agent/claudesdk/commands_session_test.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -175,5 +175,5 @@ func runCommandsHelper(request startRequest, mode string, scanner *bufio.Scanner
175175
return
176176
}
177177
emit(bridgeEvent{Type: "input_closed", SessionID: request.Resume})
178-
emit(bridgeEvent{Type: "result", SessionID: request.Resume, Text: "completed"})
178+
emit(bridgeEvent{Type: "result", SessionID: request.Resume})
179179
}

‎apps/daemon/internal/agent/claudesdk/declaration.go‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,6 @@ const claudeSDKNodeEnv = "OAC_RUNTIME_CLAUDE_SDK_NODE"
2222
var Declaration = agent.Declaration{Info: proto.SupportedAgentKind{Kind: "claude_sdk", Capabilities: proto.AgentKindCapabilities{
2323
SubagentObservations: proto.CapabilityUnsupported,
2424
NativeSessionRecovery: proto.CapabilityUnsupported,
25-
MessageItems: proto.CapabilitySupported,
2625
EnvironmentNone: proto.CapabilitySupported,
2726
LocalEnvironment: proto.CapabilityUnsupported,
2827
WorkspaceReadPreparation: proto.CapabilityUnsupported,

‎apps/daemon/internal/agent/claudesdk/declaration_test.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,7 @@ func TestClaudeSDKFeatureDiscovery(t *testing.T) {
6666

6767
// The declaration must retain the complete baseline capability descriptor.
6868
func TestDeclaredCapabilityBaseline(t *testing.T) {
69-
expected := map[string]bool{"MessageItems": true, "EnvironmentNone": true, "ProgrammaticToolCallingDisable": true, "SubagentControl": true, "FunctionTools": true}
69+
expected := map[string]bool{"EnvironmentNone": true, "ProgrammaticToolCallingDisable": true, "SubagentControl": true, "FunctionTools": true}
7070
value := reflect.ValueOf(Declaration.Info.Capabilities)
7171
for i := 0; i < value.NumField(); i++ {
7272
name := value.Type().Field(i).Name

‎apps/daemon/internal/agent/claudesdk/error_classification_test.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -92,7 +92,7 @@ func runClassifiedFailureHelper(request startRequest, mode string, encode func(b
9292
os.Exit(7)
9393
}
9494
if mode == "after-terminal" {
95-
fmt.Fprintln(os.Stdout, `{"type":"delta","delta":"late"}`)
95+
fmt.Fprintln(os.Stdout, `{"type":"delta","item_id":"message","delta":"late"}`)
9696
}
9797
if mode == "scanner-error" {
9898
fmt.Fprintln(os.Stdout, strings.Repeat("x", 2*1024*1024))

‎apps/daemon/internal/agent/claudesdk/execution_controls_test.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -81,7 +81,7 @@ func TestStructuredOutputConfigurationReachesNativeUnchanged(t *testing.T) {
8181
t.Setenv("OAC_RUNTIME_HOME", root)
8282
config := Config{Entrypoint: filepath.Join(root, "worker"), StateDir: filepath.Join(root, "state")}
8383
schema := json.RawMessage(`{"type":"object","properties":{"n":{"const":9007199254740992}}}`)
84-
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}}}
84+
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}}}
8585
start, _, err := prepareConfiguration(config, request)
8686
if err != nil {
8787
t.Fatal(err)

‎apps/daemon/internal/agent/claudesdk/executor_live_linux_test.go‎

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
6868
Settlement agent.TurnSettlement `json:"settlement"`
6969
SettlementError string `json:"settlement_error,omitempty"`
7070
Done proto.DonePayload `json:"done"`
71+
Text string `json:"text"`
7172
Errors []string `json:"errors,omitempty"`
7273
}
7374
evidence := struct {
@@ -82,7 +83,7 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
8283
_ = os.WriteFile(filepath.Join(proof, "executor-evidence.json"), raw, 0600)
8384
}
8485
defer persist()
85-
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."}
86+
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."}
8687
factory := NewExecutorFactory(config)
8788
prepared := time.Now()
8889
owner, err := factory(ctx, request)
@@ -124,6 +125,11 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
124125
if event.ID != id {
125126
t.Fatal("output crossed Turn identity")
126127
}
128+
if event.Type == proto.TypeDelta {
129+
var delta proto.DeltaPayload
130+
_ = event.DecodePayload(&delta)
131+
record.Text += delta.Delta
132+
}
127133
if event.Type == proto.TypeDelta && record.FirstTextMS == 0 {
128134
record.FirstTextMS = time.Since(started).Milliseconds()
129135
record.NativePIDs = nativeChildren(record.NodePID)
@@ -165,11 +171,11 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
165171
}
166172
marker := "REUSE-" + uuid.NewString()
167173
first := run("first", "Remember this marker: "+marker+". Reply with exactly the marker.", false)
168-
if first.SettlementError != "" || !first.Settlement.Reusable || !strings.Contains(first.Done.Content, marker) {
174+
if first.SettlementError != "" || !first.Settlement.Reusable || !strings.Contains(first.Text, marker) {
169175
t.Fatal("first Turn failed or was not reusable")
170176
}
171177
second := run("second", "What exact marker did I give you? Reply with only that marker.", false)
172-
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) {
178+
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) {
173179
t.Fatal("ordinary Turns did not retain native execution and history")
174180
}
175181
interrupted := run("cancel", "List the numbers 1 through 10000, one number per line, without stopping early.", true)
@@ -196,7 +202,7 @@ func TestLiveClaudeExecutorReuseAndCancel(t *testing.T) {
196202
evidence.Recovered = true
197203
}
198204
continued := run("continued", "What exact marker did I originally give you? Reply with only the marker.", false)
199-
if continued.SettlementError != "" || !continued.Settlement.Reusable || !strings.Contains(continued.Done.Content, marker) {
205+
if continued.SettlementError != "" || !continued.Settlement.Reusable || !strings.Contains(continued.Text, marker) {
200206
t.Fatal("history did not continue after cancellation")
201207
}
202208
if !evidence.Recovered && (continued.NodePID != second.NodePID || !slices.Equal(continued.NativePIDs, second.NativePIDs)) {

‎apps/daemon/internal/agent/claudesdk/executor_test.go‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,7 @@ func runPersistentExecutorHelper() {
6363
if cancel {
6464
encode(bridgeEvent{Type: "error", TurnID: active, Code: "cancelled"})
6565
} else {
66-
encode(bridgeEvent{Type: "result", TurnID: active, SessionID: "native-persistent", Text: active})
66+
encode(bridgeEvent{Type: "result", TurnID: active, SessionID: "native-persistent"})
6767
}
6868
confirmed, reusable, reason := true, true, ""
6969
switch os.Getenv("SDK_EXECUTOR_MODE") {
@@ -96,12 +96,12 @@ func runPersistentExecutorHelper() {
9696
}
9797
active = command.TurnID
9898
if os.Getenv("SDK_EXECUTOR_MODE") == "late" && old != "" {
99-
encode(bridgeEvent{Type: "delta", TurnID: old, Delta: "late"})
99+
encode(bridgeEvent{Type: "delta", TurnID: old, ItemID: "message", Delta: "late"})
100100
continue
101101
}
102102
encode(bridgeEvent{Type: "turn_started", TurnID: active})
103103
encode(bridgeEvent{Type: "input_ready", TurnID: active, SessionID: "native-persistent"})
104-
encode(bridgeEvent{Type: "delta", TurnID: active, Delta: "partial"})
104+
encode(bridgeEvent{Type: "delta", TurnID: active, ItemID: "message", Delta: "partial"})
105105
if strings.HasPrefix(os.Getenv("SDK_EXECUTOR_MODE"), "pending_function") {
106106
encode(bridgeEvent{Type: "function_call", TurnID: active, Call: &proto.FunctionCallPayload{CallID: "call", Name: "lookup", Arguments: json.RawMessage("{}")}})
107107
}
@@ -115,7 +115,7 @@ func runPersistentExecutorHelper() {
115115
continue
116116
}
117117
// A full bridge write has happened, but no native input receipt exists.
118-
encode(bridgeEvent{Type: "delta", TurnID: active, Delta: "steer-written"})
118+
encode(bridgeEvent{Type: "delta", TurnID: active, ItemID: "message", Delta: "steer-written"})
119119
case "turn_cancel":
120120
if active == command.TurnID {
121121
settle(true)

‎apps/daemon/internal/agent/claudesdk/executor_turn.go‎

Lines changed: 7 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@ import (
55
"fmt"
66
"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
77
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
8-
"strings"
98
"time"
109
)
1110

@@ -46,7 +45,6 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
4645
s.invalidate()
4746
}
4847
}
49-
var content strings.Builder
5048
var result *bridgeEvent
5149
var classifiedFailure error
5250
var engineCode, failedResultID string
@@ -98,17 +96,17 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
9896
s.invalidate()
9997
}
10098
case "delta":
101-
if start.ObserveMessages && event.ItemID == "" || !start.ObserveMessages && event.ItemID != "" {
102-
failure = fmt.Errorf("claudesdk: invalid message delta identity")
99+
sequence++
100+
delta := proto.DeltaPayload{ItemID: event.ItemID, Delta: event.Delta, Sequence: sequence}
101+
if err := delta.Validate(); err != nil {
102+
failure = fmt.Errorf("claudesdk: %w", err)
103103
s.invalidate()
104104
break
105105
}
106-
content.WriteString(event.Delta)
107-
sequence++
108-
emit(proto.TypeDelta, proto.DeltaPayload{ItemID: event.ItemID, Delta: event.Delta, Sequence: sequence})
106+
emit(proto.TypeDelta, delta)
109107
case "output_message":
110108
message := event.Message
111-
if !start.ObserveMessages || message == nil || message.ID == "" ||
109+
if message == nil || message.ID == "" ||
112110
(message.Status != "in_progress" && message.Status != "completed") ||
113111
(message.Status == "completed") != (message.Text != nil) {
114112
failure = fmt.Errorf("claudesdk: invalid message observation")
@@ -205,11 +203,9 @@ func (s *session) runTurn(start startRequest, out chan<- proto.Envelope) {
205203
metadata[proto.DoneMetaAgentSessionID] = id
206204
}
207205
if failure == nil {
208-
content.Reset()
209-
content.WriteString(result.Text)
210206
metadata[proto.DoneMetaAgentSessionID] = result.SessionID
211207
}
212-
s.outcome = proto.DonePayload{Content: content.String(), Usage: usage, Metadata: metadata}
208+
s.outcome = proto.DonePayload{Usage: usage, Metadata: metadata}
213209
// Publish the observed cancellation outcome before terminal delivery.
214210
if s.settlementErr != nil || !settlementReceived || !settlementConfirmed || !reusable || outputLost || s.process.Context().Err() != nil {
215211
s.owner.retire()

0 commit comments

Comments
 (0)