返回 DeepSeek-Reasonix
extensions.go
根目录 / internal / agent / extensions.go
1 package agent
2
3 import (
4 "bytes"
5 "context"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "reflect"
10 "strings"
11
12 "reasonix/internal/event"
13 "reasonix/internal/extension"
14 "reasonix/internal/extension/dispatch"
15 "reasonix/internal/extension/providerconv"
16 "reasonix/internal/i18n"
17 "reasonix/internal/provider"
18 )
19
20 // Extension Protocol v2 agent-side wiring. The agent consults the
21 // frozen dispatcher at the nine agent-loop intercept points:
22 //
23 // agent.before_start Run, before the turn is appended (block aborts the run)
24 // context.prepare stream, on the request message copy (never the session)
25 // provider.request stream, after the request is fully assembled
26 // provider.response stream, after a successful stream, before persisting
27 // tool.before executeOne, right after the call parses
28 // permission.decision executeOne, at the permission gate (allow/deny rulings)
29 // tool.after executeOne, after Execute returns (success or error)
30 // compaction.prepare compact, before the fold is archived and summarized
31 // compaction.complete compact, after the summary is produced, before persist
32 //
33 // Decision semantics are uniform: continue passes the payload through;
34 // replace substitutes it (strictly re-decoded and revalidated by the
35 // dispatcher, then checked against host invariants here); block fails the
36 // local operation — the run (before_start), the provider request (context/
37 // provider.request), the turn (provider.response), the tool call
38 // (tool.before/after, permission.decision), or the compaction pass — with the
39 // redacted reason. allow/deny is terminal at permission.decision only, where
40 // the host verdict is computed FIRST and combined: an extension allow
41 // overrides a host deny (full-trust contract, audited), an extension deny
42 // overrides a host allow, and continue leaves the host decision standing.
43 //
44 // Error policy follows the dispatcher: a required extension's failure fails
45 // the local operation; an optional extension's failure is warned about once
46 // and skipped. A nil dispatcher (no runtime packages installed) passes
47 // every point through untouched, so behavior stays byte-identical to the
48 // pre-dispatch path.
49 //
50 // Two-phase ruling at slot-mapped points: points that map to a replacement
51 // slot (context.prepare → context, provider.request → provider_request,
52 // provider.response → provider_response, compaction.prepare/complete →
53 // compaction, permission.decision → permission) first walk the intercept
54 // chain, then give the slot's OWNER the final say through RunStrategy over
55 // the (possibly interceptor-modified) payload. An owner that also declared
56 // the point under intercepts participates in both phases, in exactly this
57 // order — chain interceptor first, slot strategy last — so its strategy
58 // ruling is always the final replacement phase. A chain block short-circuits
59 // the strategy phase (the operation is already stopped). The owner is
60 // required-class by definition: its block, timeout, error, or contract
61 // violation is fatal to the local operation. Replaced values are adopted only
62 // when a replace ruling actually changed the payload, so a no-replacement
63 // walk (including an unowned slot) keeps the original values byte-identically.
64 //
65 // Ephemerality is the cache contract: context.prepare and provider.request
66 // replacements shape only the request being assembled — a.session.Messages is
67 // never mutated — while a provider.response replacement is persisted as the
68 // visible assistant turn (that IS the user's transcript), and tool.after
69 // replacements become the tool result the model reads.
70 //
71 // Observer model: after every completed intercept walk (blocked or not) the
72 // agent fires the point's fire-and-forget Event with the final payload, so
73 // observation-only extensions see exactly what the host acted on. A
74 // required-extension failure skips the event — the operation itself failed.
75
76 // extensionBlockedError reports an extension's block ruling as an operation
77 // failure. The reason is already credential-redacted by the dispatcher.
78 func extensionBlockedError(point extension.InterceptorPoint, reason string) error {
79 reason = strings.TrimSpace(reason)
80 if reason == "" {
81 reason = "no reason given"
82 }
83 return fmt.Errorf("extension blocked %s: %s — %s", point, reason, i18n.M.ExtensionBlockRecovery)
84 }
85
86 // extensionBlockReason normalizes a block reason for tool-result surfaces.
87 func extensionBlockReason(reason string) string {
88 reason = strings.TrimSpace(reason)
89 if reason == "" {
90 return "blocked by extension"
91 }
92 return reason
93 }
94
95 // strategyReplaced runs the replacement-slot owner's strategy for the point
96 // and reports whether a replace ruling actually changed the payload (the
97 // adoption signal for the caller's converted values). An unowned slot no-ops
98 // inside RunStrategy, so the fast path costs one comparison. The owner is
99 // required-class: a block, timeout, error, or contract violation is returned
100 // as a fatal error for the local operation.
101 func strategyReplaced(ctx context.Context, d *dispatch.Dispatcher, slot extension.Slot, point extension.InterceptorPoint, payloadPtr any) (bool, error) {
102 before := reflect.ValueOf(payloadPtr).Elem().Interface()
103 if err := d.RunStrategy(ctx, slot, point, payloadPtr); err != nil {
104 return false, err
105 }
106 return !reflect.DeepEqual(before, reflect.ValueOf(payloadPtr).Elem().Interface()), nil
107 }
108
109 // interceptAgentStart runs agent.before_start at the top of Run. A block (or
110 // a required extension's failure) aborts the run before the user turn is
111 // appended; the error surfaces like a normal run error.
112 func (a *Agent) interceptAgentStart(ctx context.Context) error {
113 d := a.svc.extensions
114 if d == nil {
115 return nil
116 }
117 providerCtx := a.withAgentContext(ctx)
118 payload := dispatch.AgentStartPayload{
119 Model: a.svc.prov.Name(),
120 ToolCount: len(a.svc.tools.SchemasForContext(providerCtx)),
121 SessionID: ParentSession(ctx),
122 }
123 result, err := d.Intercept(ctx, extension.PointAgentBeforeStart, &payload)
124 if err != nil {
125 return err
126 }
127 d.Event(extension.PointAgentBeforeStart, payload)
128 if result.Blocked {
129 return extensionBlockedError(extension.PointAgentBeforeStart, result.BlockReason)
130 }
131 return nil
132 }
133
134 // interceptContextPrepare runs context.prepare on the request message copy.
135 // The returned slice feeds only this provider request: the session log is
136 // never touched, so a replacement is invisible to the next turn (and to the
137 // prompt-cache prefix) — ephemerality is the cache contract.
138 func (a *Agent) interceptContextPrepare(ctx context.Context, messages []provider.Message) ([]provider.Message, error) {
139 d := a.svc.extensions
140 if d == nil {
141 return messages, nil
142 }
143 payload := dispatch.ContextPayload{Messages: providerconv.MessagesToProtocol(messages)}
144 result, err := d.Intercept(ctx, extension.PointContextPrepare, &payload)
145 if err != nil {
146 return nil, extensionRequestFailure(err)
147 }
148 if result.Blocked {
149 d.Event(extension.PointContextPrepare, payload)
150 return nil, extensionBlockedError(extension.PointContextPrepare, result.BlockReason)
151 }
152 // The context slot owner gets the final say over the chain-walked payload.
153 replaced, err := strategyReplaced(ctx, d, extension.SlotContext, extension.PointContextPrepare, &payload)
154 if err != nil {
155 return nil, extensionRequestFailure(err)
156 }
157 d.Event(extension.PointContextPrepare, payload)
158 if len(result.Applied) > 0 || replaced {
159 return providerconv.MessagesFromProtocol(payload.Messages), nil
160 }
161 return messages, nil
162 }
163
164 // interceptProviderRequest runs provider.request on the fully assembled
165 // request (post CreatedAt-strip). A replacement is revalidated by the payload
166 // registry (tool parameter schemas must be JSON objects, messages/tools must
167 // be arrays) before it may substitute the request being sent.
168 func (a *Agent) interceptProviderRequest(ctx context.Context, req provider.Request) (provider.Request, error) {
169 d := a.svc.extensions
170 if d == nil {
171 return req, nil
172 }
173 payload := dispatch.ProviderRequestPayload{Request: providerconv.RequestToProtocol(req)}
174 result, err := d.Intercept(ctx, extension.PointProviderRequest, &payload)
175 if err != nil {
176 return provider.Request{}, extensionRequestFailure(err)
177 }
178 if result.Blocked {
179 d.Event(extension.PointProviderRequest, payload)
180 return provider.Request{}, extensionBlockedError(extension.PointProviderRequest, result.BlockReason)
181 }
182 // The provider_request slot owner gets the final say over the
183 // chain-walked payload.
184 replaced, err := strategyReplaced(ctx, d, extension.SlotProviderRequest, extension.PointProviderRequest, &payload)
185 if err != nil {
186 return provider.Request{}, extensionRequestFailure(err)
187 }
188 d.Event(extension.PointProviderRequest, payload)
189 if len(result.Applied) > 0 || replaced {
190 return providerconv.RequestFromProtocol(payload.Request), nil
191 }
192 return req, nil
193 }
194
195 func extensionRequestFailure(err error) error {
196 if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
197 return err
198 }
199 return fmt.Errorf("%s: %w", i18n.M.ExtensionRequestRecovery, err)
200 }
201
202 // interceptProviderResponse runs provider.response after the stream completed
203 // successfully, before the assistant turn is persisted. A replacement is
204 // persisted as the visible turn — the user's transcript and the model's own
205 // history on the next request. The live text/reasoning deltas already
206 // streamed to the frontend are not retroactively changed; the closing Message
207 // event and the session carry the replaced values. Session-level cache
208 // counters keep the provider's real usage (they were accumulated while
209 // streaming); a replaced Usage drives only this turn's Usage event and
210 // compaction decision. A block fails the turn with the redacted reason.
211 func (a *Agent) interceptProviderResponse(ctx context.Context, text, reasoning, signature string, calls []provider.ToolCall, usage *provider.Usage) (string, string, string, []provider.ToolCall, *provider.Usage, error) {
212 d := a.svc.extensions
213 if d == nil {
214 return text, reasoning, signature, calls, usage, nil
215 }
216 payload := dispatch.ProviderResponsePayload{
217 Text: text,
218 Reasoning: reasoning,
219 Signature: signature,
220 Calls: providerconv.ToolCallsToProtocol(calls),
221 Usage: providerconv.UsageToProtocol(usage),
222 }
223 result, err := d.Intercept(ctx, extension.PointProviderResponse, &payload)
224 if err != nil {
225 return "", "", "", nil, nil, err
226 }
227 if result.Blocked {
228 d.Event(extension.PointProviderResponse, payload)
229 return "", "", "", nil, nil, extensionBlockedError(extension.PointProviderResponse, result.BlockReason)
230 }
231 // The provider_response slot owner gets the final say over the
232 // chain-walked payload.
233 replaced, err := strategyReplaced(ctx, d, extension.SlotProviderResponse, extension.PointProviderResponse, &payload)
234 if err != nil {
235 return "", "", "", nil, nil, err
236 }
237 d.Event(extension.PointProviderResponse, payload)
238 if len(result.Applied) > 0 || replaced {
239 return payload.Text, payload.Reasoning, payload.Signature,
240 providerconv.ToolCallsFromProtocol(payload.Calls), providerconv.UsageFromProtocol(payload.Usage), nil
241 }
242 return text, reasoning, signature, calls, usage, nil
243 }
244
245 // interceptToolBefore runs tool.before after the host resolved and validated
246 // the concrete target. A block
247 // fails the call with the reason as the tool-result error (mirroring a
248 // PreToolUse hook block). A replacement substitutes the provider-visible name
249 // and arguments, but only after host revalidation — the arguments must decode
250 // as a JSON object and the name must still resolve in the registry — and the
251 // substituted call is then re-parsed so policy, permission, and evidence all
252 // see the call that will actually execute. An invalid replacement fails the
253 // call with a contract-violation error result.
254 func (a *Agent) interceptToolBefore(ctx context.Context, plan *toolCallPlan) (toolOutcome, bool) {
255 d := a.svc.extensions
256 if d == nil {
257 return toolOutcome{}, false
258 }
259 payload := dispatch.ToolBeforePayload{Name: plan.call.Name, Arguments: plan.call.Arguments}
260 result, err := d.Intercept(ctx, extension.PointToolBefore, &payload)
261 if err != nil {
262 msg := fmt.Sprintf("error: %v", err)
263 return toolOutcome{output: msg, errMsg: firstLine(err.Error())}, true
264 }
265 d.Event(extension.PointToolBefore, payload)
266 if result.Blocked {
267 reason := extensionBlockReason(result.BlockReason)
268 return toolOutcome{output: "blocked: " + reason, blocked: true, errMsg: "blocked by extension"}, true
269 }
270 if len(result.Applied) == 0 {
271 return toolOutcome{}, false
272 }
273 plugin := result.Applied[len(result.Applied)-1]
274 violation := func(detail string) (toolOutcome, bool) {
275 msg := fmt.Sprintf("extension %s violated the intercept contract at %s: %s", plugin, extension.PointToolBefore, detail)
276 return toolOutcome{output: "error: " + msg, errMsg: msg}, true
277 }
278 trimmed := strings.TrimSpace(payload.Arguments)
279 if trimmed == "" || trimmed[0] != '{' {
280 return violation("arguments must decode as a JSON object")
281 }
282 t, _, ambiguous := a.svc.tools.ResolveCall(payload.Name)
283 if t == nil || len(ambiguous) > 0 {
284 return violation(fmt.Sprintf("substituted tool name %q does not resolve in the registry", payload.Name))
285 }
286 plan.call.Name = payload.Name
287 plan.call.Arguments = payload.Arguments
288 return toolOutcome{}, false
289 }
290
291 // interceptExtensionPermission runs permission.decision at the gate point.
292 // The host decision is computed first and rides the payload; the extension
293 // ruling combines with it: allow overrides a host deny (the full-trust
294 // contract — the dispatcher records the audit note, surfaced here as a
295 // warning notice), deny or block overrides a host allow, continue leaves the
296 // host decision standing. allow is updated in place; early=true carries the
297 // blocked outcome.
298 func (a *Agent) interceptExtensionPermission(ctx context.Context, plan *toolCallPlan, allow *bool) (toolOutcome, bool) {
299 d := a.svc.extensions
300 if d == nil {
301 return toolOutcome{}, false
302 }
303 hostDecision := "deny"
304 if *allow {
305 hostDecision = "allow"
306 }
307 payload := dispatch.PermissionPayload{
308 Name: plan.permName,
309 Arguments: string(plan.permArgs),
310 ReadOnly: plan.readOnly,
311 HostDecision: hostDecision,
312 }
313 result, err := d.Intercept(ctx, extension.PointPermissionDecision, &payload)
314 if err != nil {
315 return toolOutcome{
316 output: fmt.Sprintf("blocked: %v", err),
317 blocked: true,
318 errMsg: "blocked by extension permission policy",
319 }, true
320 }
321 // The permission slot owner gets the final say after the chain walk. Its
322 // effective rulings here are continue (the chain/host combination stands)
323 // and block (veto); a replace adjusts only the payload observers see —
324 // allow/deny remains the chain's terminal mechanism.
325 if !result.Blocked {
326 if serr := d.RunStrategy(ctx, extension.SlotPermission, extension.PointPermissionDecision, &payload); serr != nil {
327 reason := serr.Error()
328 var blockErr *dispatch.BlockError
329 if errors.As(serr, &blockErr) {
330 reason = extensionBlockReason(blockErr.Reason)
331 }
332 return toolOutcome{
333 output: "blocked: " + reason,
334 blocked: true,
335 errMsg: "blocked by extension permission policy",
336 }, true
337 }
338 }
339 d.Event(extension.PointPermissionDecision, payload)
340 for _, note := range result.Audit {
341 a.svc.sink.Emit(event.Event{Kind: event.Notice, Level: event.LevelWarn, Text: note})
342 }
343 switch {
344 case result.Blocked:
345 reason := extensionBlockReason(result.BlockReason)
346 return toolOutcome{
347 output: "blocked: " + reason,
348 blocked: true,
349 errMsg: "blocked by extension permission policy",
350 }, true
351 case result.Permission != nil && !*result.Permission:
352 return toolOutcome{
353 output: "blocked: denied by extension permission policy",
354 blocked: true,
355 errMsg: "blocked by extension permission policy",
356 }, true
357 case result.Permission != nil && *result.Permission:
358 *allow = true
359 }
360 return toolOutcome{}, false
361 }
362
363 // interceptToolAfter runs tool.after on the executed result. A replacement
364 // substitutes the visible result string and the error flag — clearing IsError
365 // converts a failed call into a success with the replaced text, setting it
366 // converts a success into an error result carrying the replaced text. A block
367 // (or a required extension's failure) converts the call to an error tool
368 // result with the reason; the tool itself already ran.
369 func (a *Agent) interceptToolAfter(ctx context.Context, call provider.ToolCall, result string, err error) (string, error) {
370 d := a.svc.extensions
371 if d == nil {
372 return result, err
373 }
374 payload := dispatch.ToolAfterPayload{
375 Name: call.Name,
376 Arguments: call.Arguments,
377 Result: result,
378 IsError: err != nil,
379 }
380 res, ierr := d.Intercept(ctx, extension.PointToolAfter, &payload)
381 if ierr != nil {
382 return "", ierr
383 }
384 d.Event(extension.PointToolAfter, payload)
385 if res.Blocked {
386 return "", errors.New(extensionBlockReason(res.BlockReason))
387 }
388 if len(res.Applied) > 0 {
389 result = payload.Result
390 switch {
391 case payload.IsError && err == nil:
392 err = errors.New("extension replaced this tool result with an error")
393 case !payload.IsError:
394 err = nil
395 }
396 }
397 return result, err
398 }
399
400 // interceptCompactionPrepare runs compaction.prepare before the fold is
401 // archived and summarized, colocated with the PreCompact hook so the payload's
402 // Guidance is the hook-contributed guidance (plus any /compact focus text). A
403 // replacement's messages and guidance drive only this compaction pass; a
404 // block skips the pass with the reason surfaced through the caller's notice.
405 func (a *Agent) interceptCompactionPrepare(ctx context.Context, fold []provider.Message, guidance string) ([]provider.Message, string, error) {
406 d := a.svc.extensions
407 if d == nil {
408 return fold, guidance, nil
409 }
410 payload := dispatch.CompactionPreparePayload{
411 Messages: providerconv.MessagesToProtocol(fold),
412 Guidance: guidance,
413 }
414 originalMessages, err := json.Marshal(payload.Messages)
415 if err != nil {
416 return nil, "", err
417 }
418 result, err := d.Intercept(ctx, extension.PointCompactionPrepare, &payload)
419 if err != nil {
420 return nil, "", err
421 }
422 if result.Blocked {
423 d.Event(extension.PointCompactionPrepare, payload)
424 return nil, "", extensionBlockedError(extension.PointCompactionPrepare, result.BlockReason)
425 }
426 // The compaction slot owner gets the final say over the chain-walked fold
427 // and guidance.
428 replaced, err := strategyReplaced(ctx, d, extension.SlotCompaction, extension.PointCompactionPrepare, &payload)
429 if err != nil {
430 return nil, "", err
431 }
432 d.Event(extension.PointCompactionPrepare, payload)
433 if len(result.Applied) > 0 || replaced {
434 preparedMessages, marshalErr := json.Marshal(payload.Messages)
435 if marshalErr != nil {
436 return nil, "", marshalErr
437 }
438 if bytes.Equal(preparedMessages, originalMessages) {
439 return fold, payload.Guidance, nil
440 }
441 return providerconv.MessagesFromProtocol(payload.Messages), payload.Guidance, nil
442 }
443 return fold, guidance, nil
444 }
445
446 // interceptCompactionComplete runs compaction.complete after the summary is
447 // produced (including the mechanical-fold fallback), before it is written
448 // into the session. A replacement is persisted as the summary; a block skips
449 // the pass.
450 func (a *Agent) interceptCompactionComplete(ctx context.Context, summary string) (string, error) {
451 d := a.svc.extensions
452 if d == nil {
453 return summary, nil
454 }
455 payload := dispatch.CompactionCompletePayload{Summary: summary}
456 result, err := d.Intercept(ctx, extension.PointCompactionComplete, &payload)
457 if err != nil {
458 return "", err
459 }
460 if result.Blocked {
461 d.Event(extension.PointCompactionComplete, payload)
462 return "", extensionBlockedError(extension.PointCompactionComplete, result.BlockReason)
463 }
464 // The compaction slot owner gets the final say over the chain-walked
465 // summary.
466 replaced, err := strategyReplaced(ctx, d, extension.SlotCompaction, extension.PointCompactionComplete, &payload)
467 if err != nil {
468 return "", err
469 }
470 d.Event(extension.PointCompactionComplete, payload)
471 if len(result.Applied) > 0 || replaced {
472 return payload.Summary, nil
473 }
474 return summary, nil
475 }
476
476 lines GO