| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "fmt" |
| 6 | "strings" |
| 7 | |
| 8 | "reasonix/internal/attachment" |
| 9 | "reasonix/internal/sessioninbox" |
| 10 | ) |
| 11 | |
| 12 | // RunInboxTurn synchronously claims and executes one durable item. Bot and ACP |
| 13 | // use this path so their blocking response sink remains attached through every |
| 14 | // queued follow-up while Controller still owns durable state and ack semantics. |
| 15 | func (c *Controller) RunInboxTurn(ctx context.Context, id string) error { |
| 16 | st, err := c.ensureInbox() |
| 17 | if err != nil { |
| 18 | return err |
| 19 | } |
| 20 | meta, env, err := st.ReadItem(id) |
| 21 | if err != nil { |
| 22 | return err |
| 23 | } |
| 24 | if meta.State != sessioninbox.StateQueued { |
| 25 | return sessioninbox.ErrInvalidState |
| 26 | } |
| 27 | run, block, err := c.prepareInboxRunContext(ctx, env) |
| 28 | if err != nil { |
| 29 | return err |
| 30 | } |
| 31 | if block != "" { |
| 32 | _ = st.SetState(id, sessioninbox.StateBlocked, block) |
| 33 | _ = st.SetPaused(true) |
| 34 | return fmt.Errorf("%w: %s", sessioninbox.ErrInvalidState, block) |
| 35 | } |
| 36 | err = c.runSynchronousTurn(ctx, func() error { |
| 37 | c.inbox.admissionMu.Lock() |
| 38 | defer c.inbox.admissionMu.Unlock() |
| 39 | c.inbox.trackAdmission(id) |
| 40 | defer c.inbox.untrackAdmission(id) |
| 41 | if err := st.ClaimItem(id); err != nil { |
| 42 | return err |
| 43 | } |
| 44 | c.inbox.mu.Lock() |
| 45 | c.inbox.trackActive(id) |
| 46 | c.inbox.mu.Unlock() |
| 47 | return nil |
| 48 | }, run) |
| 49 | if err != nil { |
| 50 | return err |
| 51 | } |
| 52 | return c.waitForGoalTerminal(ctx) |
| 53 | } |
| 54 | |
| 55 | func (c *Controller) prepareInboxRun(env sessioninbox.PromptEnvelope) (func(context.Context) error, string, error) { |
| 56 | return c.prepareInboxRunContext(c.attachmentContext(), env) |
| 57 | } |
| 58 | |
| 59 | func (c *Controller) prepareInboxRunContext(ctx context.Context, env sessioninbox.PromptEnvelope) (func(context.Context) error, string, error) { |
| 60 | var sources []attachment.Source |
| 61 | for _, input := range env.ImageInputs { |
| 62 | if input.Attachment != nil { |
| 63 | sources = append(sources, attachment.Source{Existing: input.Attachment, DisplayName: input.Attachment.DisplayName}) |
| 64 | } |
| 65 | } |
| 66 | if _, err := c.attachmentService().PrepareBatch(ctx, sources); err != nil { |
| 67 | return nil, ImageReferenceFailures(imageFailuresFromAttachment(err)).Error(), nil |
| 68 | } |
| 69 | submit, frozenImages, block, err := applyInboxReferences(env) |
| 70 | if err != nil || block != "" { |
| 71 | return nil, block, err |
| 72 | } |
| 73 | display := firstNonEmptyStr(env.DisplayText, submit) |
| 74 | raw := firstNonEmptyStr(env.RawText, submit) |
| 75 | requests := controlInvocationsFromInbox(env) |
| 76 | if len(requests) == 0 { |
| 77 | return func(ctx context.Context) error { |
| 78 | ctx = contextWithPreparedImageReferences(ctx, preparedImageReferences{inputs: env.ImageInputs}) |
| 79 | return c.runGoalLoopWithFrozenImagesRawDisplay(c.withTurnFormat(ctx, strings.TrimSpace(env.Format)), submit, raw, display, frozenImages) |
| 80 | }, "", nil |
| 81 | } |
| 82 | prepared, err := c.prepareInvocationTurn(submit, requests) |
| 83 | if err != nil { |
| 84 | return nil, err.Error(), nil |
| 85 | } |
| 86 | return func(ctx context.Context) error { |
| 87 | ctx = contextWithPreparedImageReferences(ctx, preparedImageReferences{inputs: env.ImageInputs}) |
| 88 | return c.runPreparedInvocationTurn(c.withTurnFormat(ctx, strings.TrimSpace(env.Format)), prepared, submit, raw, display, frozenImages) |
| 89 | }, "", nil |
| 90 | } |
| 91 | |
| 92 | // submitPreparedInboxTurn starts an already-classified inbox envelope without |
| 93 | // interpreting slash commands, shell shortcuts, or @references a second time. |
| 94 | func (c *Controller) submitPreparedInboxTurn(itemID string, run func(context.Context) error) admissionResult { |
| 95 | return c.runGuardedInbox(run, func() { |
| 96 | c.inbox.mu.Lock() |
| 97 | c.inbox.trackActive(itemID) |
| 98 | c.inbox.mu.Unlock() |
| 99 | }) |
| 100 | } |
| 101 | |
| 102 | func sessionInboxInvocations(requests []InvocationRequest) []sessioninbox.StructuredInvocation { |
| 103 | if len(requests) == 0 { |
| 104 | return nil |
| 105 | } |
| 106 | out := make([]sessioninbox.StructuredInvocation, 0, len(requests)) |
| 107 | for _, request := range requests { |
| 108 | out = append(out, sessioninbox.StructuredInvocation{Name: request.Name, Kind: request.Kind, Offset: request.Offset}) |
| 109 | } |
| 110 | return out |
| 111 | } |
| 112 | |
| 113 | func controlInvocationsFromInbox(env sessioninbox.PromptEnvelope) []InvocationRequest { |
| 114 | stored := env.Invocations |
| 115 | if len(stored) == 0 && env.Invocation != nil { |
| 116 | stored = []sessioninbox.StructuredInvocation{*env.Invocation} |
| 117 | } |
| 118 | out := make([]InvocationRequest, 0, len(stored)) |
| 119 | for _, invocation := range stored { |
| 120 | out = append(out, InvocationRequest{Name: invocation.Name, Kind: invocation.Kind, Offset: invocation.Offset}) |
| 121 | } |
| 122 | return out |
| 123 | } |
| 124 |