返回 DeepSeek-Reasonix
submission_identity.go
根目录 / internal / control / submission_identity.go
1 package control
2
3 import (
4 "context"
5 "crypto/sha256"
6 "encoding/hex"
7 "encoding/json"
8 "errors"
9 "log/slog"
10 "regexp"
11 "strings"
12 "sync"
13 "sync/atomic"
14
15 "reasonix/internal/agent"
16 "reasonix/internal/attachment"
17 "reasonix/internal/event"
18 "reasonix/internal/session"
19 )
20
21 // SubmissionRequest fingerprints the actual operation, never only its label.
22 type SubmissionRequest struct {
23 ID string `json:"-"`
24 HTTP bool `json:"http,omitempty"`
25 Input string `json:"input"`
26 Display string `json:"display,omitempty"`
27 Format string `json:"format,omitempty"`
28 Action string `json:"action,omitempty"`
29 RecoveryID string `json:"recoveryId,omitempty"`
30 Original string `json:"original,omitempty"`
31 Goal string `json:"goal,omitempty"`
32 ToolApprovalMode string `json:"toolApprovalMode,omitempty"`
33 Invocations []InvocationRequest `json:"invocations,omitempty"`
34 DraftIDs []string `json:"draftIds,omitempty"`
35 AttachmentDigests []string `json:"attachmentDigests,omitempty"`
36 Attachments []SubmissionAttachment `json:"attachments,omitempty"`
37 frozenSources map[string]*attachment.AttachmentRef
38 inheritedSources []attachment.Source
39 }
40
41 // SubmissionAttachment describes an ordered logical item, independently of
42 // its temporary credential. Reference reads still require session authority.
43 type SubmissionAttachment struct {
44 ClientAttachmentID string `json:"clientAttachmentId"`
45 DraftID string `json:"draftId,omitempty"`
46 Path string `json:"path,omitempty"`
47 Reference *attachment.AttachmentRef `json:"reference,omitempty"`
48 }
49
50 // ErrSubmissionNotAccepted is only attached to failures that precede durable
51 // admission. Transport cancellation and flush failures intentionally omit it:
52 // the caller must query the receipt before deciding whether a retry may run.
53 var ErrSubmissionNotAccepted = errors.New("submission not accepted")
54
55 // ErrSubmissionIdentityUnavailable marks a runtime with no durable store to
56 // record an identified submission in, such as a native legacy transcript.
57 var ErrSubmissionIdentityUnavailable = errors.New("durable submission identity unavailable")
58
59 type submissionIdentityState struct {
60 mu sync.Mutex
61 pending atomic.Pointer[pendingSubmissionAdmission]
62 reused atomic.Uint64
63 conflicts atomic.Uint64
64 unknown atomic.Uint64
65 }
66
67 // pendingSubmissionAdmission is the immutable fact persisted with turn/start.
68 // The execution path receives the same images through turnAdmission, so it does
69 // not need to consult controller-scoped temporary state.
70 type pendingSubmissionAdmission struct {
71 receipt session.SubmissionReceipt
72 imageInputs []attachment.ImageInput
73 durableCtx context.Context
74 }
75
76 type turnAdmission struct {
77 images preparedImageReferences
78 durableCtx context.Context
79 result *admissionResult
80 }
81
82 func newTurnAdmission(ctx context.Context, images preparedImageReferences) turnAdmission {
83 if ctx == nil {
84 ctx = context.Background()
85 }
86 return turnAdmission{images: images, durableCtx: ctx}
87 }
88
89 func (c *Controller) appendSubmissionEvent(events []session.Event, turnID string) []session.Event {
90 if pending := c.submissions.pending.Load(); pending != nil {
91 receipt := pending.receipt
92 receipt.TurnID = turnID
93 payload := struct {
94 session.SubmissionReceipt
95 ImageInputs []attachment.ImageInput `json:"imageInputs,omitempty"`
96 }{SubmissionReceipt: receipt, ImageInputs: pending.imageInputs}
97 data, _ := json.Marshal(payload)
98 return append(events, session.Event{Kind: "submission/accepted", Optional: true, Payload: data})
99 }
100 return events
101 }
102
103 func (c *Controller) trySubmissionAdmissionLock() func() {
104 if !c.submissions.mu.TryLock() {
105 return nil
106 }
107 return sync.OnceFunc(c.submissions.mu.Unlock)
108 }
109
110 func (c *Controller) releaseSubmissionAdmission() {
111 c.submissions.mu.Unlock()
112 // A synchronous queue dispatch may have deferred while submit owned the gate.
113 c.maybeDispatchInbox()
114 }
115
116 func submissionFingerprint(req SubmissionRequest) string {
117 data, _ := json.Marshal(req)
118 digest := sha256.Sum256(data)
119 return hex.EncodeToString(digest[:])
120 }
121
122 const submissionFingerprintVersion = 1
123
124 func canonicalSubmissionFingerprint(req SubmissionRequest) string {
125 req.AttachmentDigests = nil
126 for i, item := range req.Attachments {
127 transport := item.Path
128 if item.DraftID != "" {
129 transport = "draft:" + item.DraftID
130 }
131 if transport != "" {
132 req.Display = strings.ReplaceAll(req.Display, "("+transport+")", "(attachment:"+item.ClientAttachmentID+")")
133 req.Input = canonicalAttachmentToken(req.Input, transport, item.ClientAttachmentID)
134 req.Original = canonicalAttachmentToken(req.Original, transport, item.ClientAttachmentID)
135 }
136 if item.DraftID != "" {
137 req.DraftIDs = withoutDraftID(req.DraftIDs, item.DraftID)
138 }
139 if i == 0 {
140 req.Attachments = append([]SubmissionAttachment(nil), req.Attachments...)
141 }
142 req.Attachments[i] = SubmissionAttachment{ClientAttachmentID: item.ClientAttachmentID}
143 }
144 return submissionFingerprint(req)
145 }
146
147 // MatchesSubmissionReceipt validates a cold durable receipt without creating
148 // a Controller. Version zero keeps the exact pre-attachment fingerprint.
149 func MatchesSubmissionReceipt(req SubmissionRequest, receipt session.SubmissionReceipt) bool {
150 if req.ID != receipt.SubmissionID {
151 return false
152 }
153 switch receipt.FingerprintVersion {
154 case 0:
155 return receipt.Fingerprint == submissionFingerprint(req)
156 case submissionFingerprintVersion:
157 return receipt.Fingerprint == canonicalSubmissionFingerprint(req)
158 default:
159 return false
160 }
161 }
162
163 // LookupSubmission also detects accidental reuse of a key for different input.
164 func (c *Controller) LookupSubmission(req SubmissionRequest) (session.SubmissionReceipt, bool, error) {
165 return c.LookupSubmissionContext(c.attachmentContext(), req)
166 }
167
168 func (c *Controller) LookupSubmissionContext(ctx context.Context, req SubmissionRequest) (session.SubmissionReceipt, bool, error) {
169 store := c.sessionEventStore()
170 if store == nil || req.ID == "" {
171 return session.SubmissionReceipt{}, false, nil
172 }
173 receipt, ok := store.Submission(req.ID)
174 fingerprint := canonicalSubmissionFingerprint(req)
175 if receipt.FingerprintVersion == 0 {
176 fingerprint = submissionFingerprint(req)
177 }
178 if ok && receipt.FingerprintVersion > submissionFingerprintVersion {
179 return receipt, true, errors.New("submission receipt version is unsupported; automatic replay is disabled")
180 }
181 if ok && receipt.Fingerprint != fingerprint {
182 if receipt.FingerprintVersion == 0 {
183 return receipt, true, errors.New("legacy submission receipt cannot be verified; automatic replay is disabled")
184 }
185 c.submissions.conflicts.Add(1)
186 slog.Warn("submission conflict", "submissionId", req.ID)
187 return receipt, true, errors.New("submission identity conflicts with different input")
188 }
189 if ok {
190 c.submissions.reused.Add(1)
191 // The projection may already contain an asynchronously accepted batch.
192 // A retry is acknowledged only after its canonical journal is durable.
193 if _, err := store.Flush(ctx); err != nil {
194 c.submissions.unknown.Add(1)
195 return receipt, true, err
196 }
197 slog.Debug("submission already accepted", "submissionId", req.ID, "turnId", receipt.TurnID)
198 }
199 return receipt, ok, nil
200 }
201
202 // SubmitIdentified serializes identity checking with synchronous turn admission.
203 func (c *Controller) SubmitIdentified(req SubmissionRequest) (session.SubmissionReceipt, error) {
204 return c.SubmitIdentifiedContext(c.attachmentContext(), req)
205 }
206
207 func (c *Controller) SubmitIdentifiedContext(ctx context.Context, req SubmissionRequest) (session.SubmissionReceipt, error) {
208 if len(req.ID) > 256 || strings.ContainsAny(req.ID, "\x00\r\n") {
209 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, errors.New("invalid submission identity"))
210 }
211 return c.submitIdentifiedWithSetupContext(ctx, req, nil, func(admission turnAdmission) {
212 c.submitIdentifiedRequestLocked(req, admission)
213 })
214 }
215
216 func (c *Controller) submitIdentified(req SubmissionRequest, submit func()) (session.SubmissionReceipt, error) {
217 return c.submitIdentifiedWithSetup(req, nil, func(turnAdmission) { submit() })
218 }
219
220 // SubmitIdentifiedWithSetup validates and freezes explicit image attachments
221 // before setup mutates host-visible session state. Desktop uses this for the
222 // initial Goal transaction, whose Goal/profile changes must not survive a
223 // rejected image submission.
224 func (c *Controller) SubmitIdentifiedWithSetup(req SubmissionRequest, setup func() error) (session.SubmissionReceipt, error) {
225 return c.SubmitIdentifiedWithSetupContext(c.attachmentContext(), req, setup)
226 }
227
228 func (c *Controller) SubmitIdentifiedWithSetupContext(ctx context.Context, req SubmissionRequest, setup func() error) (session.SubmissionReceipt, error) {
229 if len(req.ID) > 256 || strings.ContainsAny(req.ID, "\x00\r\n") {
230 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, errors.New("invalid submission identity"))
231 }
232 return c.submitIdentifiedWithSetupContext(ctx, req, setup, func(admission turnAdmission) {
233 c.submitIdentifiedRequestLocked(req, admission)
234 })
235 }
236
237 func (c *Controller) submitIdentifiedRequestLocked(req SubmissionRequest, admission turnAdmission) {
238 switch {
239 case req.Action == ProtocolRecoveryAction:
240 c.submitProtocolRecoveryLocked(req.RecoveryID, req.Input, admission)
241 case req.Action == "delivery-recovery" || req.Action == FinalReadinessRecoveryAction:
242 c.submitFinalReadinessRecoveryLocked(req.Display, req.Input, admission)
243 case req.Action == "shell":
244 c.runShell(req.Input, admission)
245 case req.HTTP:
246 c.submitHTTPWithFormatLocked(req.Input, req.Display, req.Format, admission)
247 case len(req.Invocations) > 0:
248 c.submitInvocationsLocked(req.Input, req.Display, req.Invocations, admission)
249 default:
250 c.submitLocked(req.Input, req.Display, req.Original, admission)
251 }
252 }
253
254 func (c *Controller) submitIdentifiedWithSetup(req SubmissionRequest, setup func() error, submit func(turnAdmission)) (session.SubmissionReceipt, error) {
255 return c.submitIdentifiedWithSetupContext(c.attachmentContext(), req, setup, submit)
256 }
257
258 func (c *Controller) submitIdentifiedWithSetupContext(ctx context.Context, req SubmissionRequest, setup func() error, submit func(turnAdmission)) (session.SubmissionReceipt, error) {
259 prepared, err := c.PrepareSubmission(ctx, req)
260 if err != nil {
261 c.notice(err.Error())
262 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, err)
263 }
264 return c.acceptPreparedSubmission(ctx, prepared, setup, submit)
265 }
266
267 func (c *Controller) acceptPreparedSubmission(ctx context.Context, candidate *PreparedSubmission, setup func() error, submit func(turnAdmission)) (session.SubmissionReceipt, error) {
268 ctx, cancel := c.NewAttachmentOperationContext(ctx)
269 defer cancel()
270 c.submissions.mu.Lock()
271 defer c.releaseSubmissionAdmission()
272 req := candidate.request
273 if candidate.scope != c.attachmentScope() {
274 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, errors.New("attachment target changed; please retry"))
275 }
276 if receipt, ok, err := c.LookupSubmissionContext(ctx, req); ok || err != nil {
277 return receipt, err
278 }
279 prepared := candidate.images
280 if err := ctx.Err(); err != nil {
281 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, err)
282 }
283 store := c.sessionEventStore()
284 if req.ID != "" {
285 if store == nil {
286 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, ErrSubmissionIdentityUnavailable)
287 }
288 if c.Running() {
289 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, ErrTurnRunning)
290 }
291 }
292 if setup != nil {
293 if err := setup(); err != nil {
294 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, err)
295 }
296 }
297 admission := newTurnAdmission(ctx, prepared)
298 if req.ID == "" {
299 result := admissionResult(-1)
300 admission.result = &result
301 submit(admission)
302 if result != turnStarted {
303 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, errors.New("submission was not admitted"))
304 }
305 return session.SubmissionReceipt{}, nil
306 }
307 receipt := &session.SubmissionReceipt{SessionID: store.ID(), SubmissionID: req.ID,
308 FingerprintVersion: submissionFingerprintVersion, Fingerprint: canonicalSubmissionFingerprint(req), MessageID: agent.NewMessageID(), AcceptedAttachmentDigests: strings.Join(attachmentDigests(prepared), ",")}
309 pending := &pendingSubmissionAdmission{
310 receipt: *receipt,
311 imageInputs: append([]attachment.ImageInput(nil), prepared.inputs...),
312 durableCtx: ctx,
313 }
314 c.submissions.pending.Store(pending)
315 defer c.submissions.pending.Store(nil)
316 c.SetTurnSubmissionID(req.ID)
317 submit(admission)
318 if err := c.flushSubmissionAdmission(ctx); err != nil {
319 c.submissions.unknown.Add(1)
320 return session.SubmissionReceipt{}, err
321 }
322 accepted, ok := store.Submission(req.ID)
323 if !ok {
324 return session.SubmissionReceipt{}, errors.Join(ErrSubmissionNotAccepted, errors.New("submission was not durably admitted"))
325 }
326 return accepted, nil
327 }
328
329 func (c *Controller) submissionForTurn(turnID string) (session.SubmissionReceipt, bool) {
330 store := c.sessionEventStore()
331 if store == nil {
332 return session.SubmissionReceipt{}, false
333 }
334 return store.SubmissionForTurn(turnID)
335 }
336
337 func (c *Controller) flushSubmissionAdmission(ctx context.Context) error {
338 if c.submissions.pending.Load() == nil {
339 return nil
340 }
341 _, err := c.sessionEventStore().Flush(ctx)
342 return err
343 }
344
345 func (c *Controller) submissionAdmissionContext() context.Context {
346 if pending := c.submissions.pending.Load(); pending != nil && pending.durableCtx != nil {
347 return pending.durableCtx
348 }
349 return context.Background()
350 }
351
352 func (c *Controller) prepareSubmissionImages(req SubmissionRequest) (preparedImageReferences, []ImageReferenceFailure) {
353 return c.prepareSubmissionImagesContext(c.attachmentContext(), req)
354 }
355
356 func (c *Controller) prepareSubmissionImagesContext(ctx context.Context, req SubmissionRequest) (preparedImageReferences, []ImageReferenceFailure) {
357 ids := append([]string(nil), req.DraftIDs...)
358 ids = append(ids, draftIDsFromInput(req.Input)...)
359 seen := make(map[string]bool)
360 unique := ids[:0]
361 for _, id := range ids {
362 if !seen[id] {
363 seen[id] = true
364 unique = append(unique, id)
365 }
366 }
367 ids = unique
368 prepared := preparedImageReferences{byPath: map[string]string{}}
369 svc := c.attachmentService()
370 sources := append([]attachment.Source(nil), req.inheritedSources...)
371 structured, failures := c.structuredImageSources(ctx, req.Attachments)
372 if len(failures) > 0 {
373 return preparedImageReferences{}, failures
374 }
375 sources = append(sources, structured...)
376 if len(req.Attachments) > 0 {
377 filtered := ids[:0]
378 for _, id := range ids {
379 found := false
380 for _, item := range req.Attachments {
381 if item.DraftID == id {
382 found = true
383 break
384 }
385 }
386 if !found {
387 filtered = append(filtered, id)
388 }
389 }
390 ids = filtered
391 }
392 for _, id := range ids {
393 key := "draft:" + id
394 ref := req.frozenSources[key]
395 if ref == nil {
396 draft, ok := svc.Drafts().Lookup(c.attachmentScope(), id)
397 if !ok {
398 return preparedImageReferences{}, imageFailuresFromAttachment(attachment.Error{Code: attachment.CodeMissing, Message: "draft credential is not valid"})
399 }
400 ref = &draft.Ref
401 }
402 sources = append(sources, attachment.Source{Existing: ref, DisplayName: ref.DisplayName, Path: key})
403 }
404 // Structured attachments promise image understanding for this turn. Legacy
405 // @.reasonix/attachments paths are also frozen below, but a text-only model may
406 // retain them as tool-readable references without a vision fallback.
407 prepared.requiresImageUnderstanding = len(sources) > 0
408 for _, source := range c.explicitImageSources(req.Input) {
409 if frozen := req.frozenSources[normalizedImageReferencePath(source.Path)]; frozen != nil {
410 source.Existing = frozen
411 }
412 duplicate := false
413 for _, existing := range sources {
414 if existing.Path != "" && normalizedImageReferencePath(existing.Path) == normalizedImageReferencePath(source.Path) {
415 duplicate = true
416 break
417 }
418 }
419 if !duplicate {
420 sources = append(sources, source)
421 }
422 }
423 if len(sources) == 0 {
424 return prepared, nil
425 }
426 sources = c.appendOrdinaryImageSources(req.Input, sources)
427 batch, err := svc.PrepareBatch(ctx, sources)
428 if err != nil {
429 return preparedImageReferences{}, imageFailuresFromAttachment(err)
430 }
431 refs, err := svc.CommitBatch(ctx, batch)
432 if err != nil {
433 return preparedImageReferences{}, imageFailuresFromAttachment(err)
434 }
435 if err := c.rebindPreparedDrafts(req, ids); err != nil {
436 return preparedImageReferences{}, imageFailuresFromAttachment(err)
437 }
438 prepared.inputs = svc.InputsFromRefs(refs)
439 for _, ref := range refs {
440 prepared.ordered = append(prepared.ordered, ref.Content.Digest)
441 }
442 for i, source := range sources {
443 if source.Path != "" {
444 prepared.byPath[normalizedImageReferencePath(source.Path)] = refs[i].Content.Digest
445 }
446 }
447 for i, item := range req.Attachments {
448 prepared.byPath["attachment:"+item.ClientAttachmentID] = refs[len(req.inheritedSources)+i].Content.Digest
449 }
450 prepared.inputs = append(prepared.inputs, legacyRemoteImageInputs(req.Input)...)
451 return prepared, nil
452 }
453
454 var draftIDPattern = regexp.MustCompile(`^draft:([0-9a-fA-F]{32})$`)
455
456 func draftIDsFromInput(input string) []string {
457 var ids []string
458 seen := map[string]bool{}
459 for _, token := range parseRefTokens(input) {
460 match := draftIDPattern.FindStringSubmatch(token)
461 if len(match) != 2 {
462 continue
463 }
464 id := match[1]
465 if seen[id] {
466 continue
467 }
468 seen[id] = true
469 ids = append(ids, id)
470 }
471 return ids
472 }
473
474 func attachmentDigests(prepared preparedImageReferences) []string {
475 if len(prepared.inputs) == 0 {
476 return nil
477 }
478 out := make([]string, 0, len(prepared.inputs))
479 for _, in := range prepared.inputs {
480 if in.Kind == attachment.KindAttachment && in.Attachment != nil {
481 out = append(out, in.Attachment.Content.Digest)
482 }
483 }
484 return out
485 }
486
487 func (c *Controller) flushSubmissionStart(ctx context.Context, kind event.Kind) error {
488 if kind != event.TurnStarted {
489 return nil
490 }
491 return c.flushSubmissionAdmission(ctx)
492 }
493
494 // TurnIDForSubmission exposes the synchronous admission receipt without
495 // depending on whether the provider is still running when the desktop call returns.
496 func (c *Controller) TurnIDForSubmission(submissionID string) string {
497 if store := c.sessionEventStore(); store != nil {
498 if receipt, ok := store.Submission(submissionID); ok {
499 return receipt.TurnID
500 }
501 }
502 ledger := c.turnEventLedger()
503 if ledger == nil {
504 return ""
505 }
506 return ledger.TurnIDForSubmission(submissionID)
507 }
508
508 lines GO