| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "maps" |
| 7 | "reasonix/internal/agent" |
| 8 | "reasonix/internal/sessioninbox" |
| 9 | "strings" |
| 10 | ) |
| 11 | |
| 12 | // EnqueueInbox durably queues an instruction. Only returns a receipt after |
| 13 | // blob+manifest commit. Does not auto-start a turn (call TrySubmit / dispatcher). |
| 14 | func (c *Controller) EnqueueInbox(req InboxRequest) (sessioninbox.InboxReceipt, error) { |
| 15 | return c.EnqueueInboxContext(c.attachmentContext(), req) |
| 16 | } |
| 17 | |
| 18 | func (c *Controller) EnqueueInboxContext(ctx context.Context, req InboxRequest) (sessioninbox.InboxReceipt, error) { |
| 19 | ctx, cancel := c.NewAttachmentOperationContext(ctx) |
| 20 | defer cancel() |
| 21 | c.inbox.prepareMu.Lock() |
| 22 | defer c.inbox.prepareMu.Unlock() |
| 23 | st, err := c.ensureInbox() |
| 24 | if err != nil { |
| 25 | return sessioninbox.InboxReceipt{}, err |
| 26 | } |
| 27 | if req.ExpectedSessionPath != "" && st.SessionPath() != req.ExpectedSessionPath { |
| 28 | return sessioninbox.InboxReceipt{}, ErrInboxSessionChanged |
| 29 | } |
| 30 | submit := strings.TrimSpace(firstNonEmptyStr(req.Submit, req.Raw)) |
| 31 | if submit == "" && len(req.Invocations) == 0 { |
| 32 | submit = strings.TrimSpace(req.Display) |
| 33 | } |
| 34 | if submit == "" && len(req.Invocations) == 0 { |
| 35 | return sessioninbox.InboxReceipt{}, sessioninbox.ErrEmpty |
| 36 | } |
| 37 | display := firstNonEmptyStr(req.Display, submit) |
| 38 | raw := firstNonEmptyStr(req.Raw, submit) |
| 39 | env := sessioninbox.PromptEnvelope{ |
| 40 | DisplayText: display, |
| 41 | RawText: raw, |
| 42 | SubmitText: submit, |
| 43 | Format: req.Format, |
| 44 | Source: req.Source, |
| 45 | Idempotency: req.Idempotency, |
| 46 | ExplicitRefs: append([]string(nil), req.FreezeRefs...), |
| 47 | Invocations: sessionInboxInvocations(req.Invocations), |
| 48 | Extra: maps.Clone(req.Extra), |
| 49 | } |
| 50 | for _, item := range req.Attachments { |
| 51 | env.AttachmentIdentities = append(env.AttachmentIdentities, item.ClientAttachmentID) |
| 52 | } |
| 53 | if len(req.Attachments) > 0 { |
| 54 | env.FingerprintVersion = 1 |
| 55 | env.RequestFingerprint = inboxAttachmentFingerprint(req, env) |
| 56 | } |
| 57 | if receipt, found, err := st.LookupEnvelopeReceipt(req.Idempotency, env); found || err != nil { |
| 58 | return receipt, err |
| 59 | } |
| 60 | if err := c.freezeInboxEnvelopeReferences(ctx, &env, submit, req.FreezeRefs, req.Attachments...); err != nil { |
| 61 | return sessioninbox.InboxReceipt{}, err |
| 62 | } |
| 63 | if err := ctx.Err(); err != nil { |
| 64 | return sessioninbox.InboxReceipt{}, err |
| 65 | } |
| 66 | intent := req.Intent |
| 67 | if intent != sessioninbox.IntentSteer { |
| 68 | intent = sessioninbox.IntentFollowup |
| 69 | } |
| 70 | rec, err := st.Enqueue(sessioninbox.EnqueueRequest{ |
| 71 | Intent: intent, |
| 72 | Envelope: env, |
| 73 | Source: req.Source, |
| 74 | Idempotency: req.Idempotency, |
| 75 | SessionID: agent.BranchID(st.SessionPath()), |
| 76 | }) |
| 77 | if err != nil { |
| 78 | if errors.Is(err, sessioninbox.ErrCapacityItems) || errors.Is(err, sessioninbox.ErrCapacityBytes) || errors.Is(err, sessioninbox.ErrItemTooLarge) { |
| 79 | sessioninbox.NoteCapacityReject() |
| 80 | } else { |
| 81 | sessioninbox.NoteTxFail() |
| 82 | } |
| 83 | return sessioninbox.InboxReceipt{}, err |
| 84 | } |
| 85 | if !rec.Idempotent && len(env.ReferenceErrors) > 0 { |
| 86 | reason := strings.Join(env.ReferenceErrors, "; ") |
| 87 | if stateErr := st.SetState(rec.ItemID, sessioninbox.StateBlocked, reason); stateErr != nil { |
| 88 | return sessioninbox.InboxReceipt{}, stateErr |
| 89 | } |
| 90 | if pauseErr := st.SetPaused(true); pauseErr != nil { |
| 91 | return sessioninbox.InboxReceipt{}, pauseErr |
| 92 | } |
| 93 | rec.Paused = true |
| 94 | } |
| 95 | sessioninbox.NoteEnqueue(int64(len(env.SubmitText))) |
| 96 | return rec, nil |
| 97 | } |
| 98 |