| 1 | package main |
| 2 | |
| 3 | import ( |
| 4 | "encoding/json" |
| 5 | "fmt" |
| 6 | "io" |
| 7 | "net/http" |
| 8 | "net/url" |
| 9 | |
| 10 | "reasonix/internal/agent" |
| 11 | "reasonix/internal/servecontract" |
| 12 | "reasonix/internal/sessioninbox" |
| 13 | ) |
| 14 | |
| 15 | // remoteInboxSnapshot reads the selected durable queue using the same route |
| 16 | // identity as remote runtime synchronization. A reply from a replaced tunnel |
| 17 | // or selection must not clear or populate the current composer's queue. |
| 18 | func (a *App) remoteInboxSnapshot(tabID string) (InboxSnapshotView, error) { |
| 19 | a.remoteTabMu.Lock() |
| 20 | tab := a.remoteTabs[tabID] |
| 21 | if tab == nil || tab.client == nil || tab.state != "ready" || tab.routing.currentPath == "" || tab.routing.rehydratingPath != "" { |
| 22 | a.remoteTabMu.Unlock() |
| 23 | return InboxSnapshotView{}, fmt.Errorf("remote tab %q is not ready for inbox reads", tabID) |
| 24 | } |
| 25 | client, base, path := tab.client, tab.base, tab.routing.currentPath |
| 26 | gen, selection := tab.gen, tab.selectionRevision |
| 27 | a.remoteTabMu.Unlock() |
| 28 | |
| 29 | ctx, cancel := commandContext(a) |
| 30 | defer cancel() |
| 31 | data, err := serveGet(ctx, client, serveURL(base, "/inbox?session="+url.QueryEscape(path))) |
| 32 | if err != nil { |
| 33 | return InboxSnapshotView{}, err |
| 34 | } |
| 35 | var snap sessioninbox.InboxSnapshot |
| 36 | if err := json.Unmarshal(data, &snap); err != nil { |
| 37 | return InboxSnapshotView{}, err |
| 38 | } |
| 39 | if snap.SessionPath == "" || agent.CanonicalSessionPath(snap.SessionPath) != agent.CanonicalSessionPath(path) { |
| 40 | return InboxSnapshotView{}, fmt.Errorf("remote inbox session changed") |
| 41 | } |
| 42 | a.remoteTabMu.Lock() |
| 43 | defer a.remoteTabMu.Unlock() |
| 44 | current := a.remoteTabs[tabID] |
| 45 | if current != tab || current.gen != gen || current.selectionRevision != selection || current.client != client || current.base != base || current.state != "ready" || current.routing.currentPath != path || current.routing.rehydratingPath != "" { |
| 46 | return InboxSnapshotView{}, fmt.Errorf("remote inbox route changed during read") |
| 47 | } |
| 48 | view := inboxSnapshotView(snap) |
| 49 | view.MutationsSupported = current.capabilities[servecontract.InboxMutationsV1] |
| 50 | return view, nil |
| 51 | } |
| 52 | |
| 53 | // enqueueRemoteFollowup preserves the route, rich input and caller's stable |
| 54 | // key. An uncertain POST is reconciled by reading its receipt, never replayed. |
| 55 | func (a *App) enqueueRemoteFollowup(tabID, display, submit string, invocations []InvocationRequest, idempotency string) (InboxReceiptView, error) { |
| 56 | if err := a.requireRemoteExecutionProtocol(tabID); err != nil { |
| 57 | return InboxReceiptView{}, err |
| 58 | } |
| 59 | client, base, path, err := a.remoteTabCommandTarget(tabID) |
| 60 | if err != nil { |
| 61 | return InboxReceiptView{}, err |
| 62 | } |
| 63 | return a.enqueueRemoteFollowupAt(client, base, path, display, submit, invocations, idempotency) |
| 64 | } |
| 65 | |
| 66 | func (a *App) enqueueRemoteFollowupAt(client *http.Client, base, path, display, submit string, invocations []InvocationRequest, idempotency string) (InboxReceiptView, error) { |
| 67 | ctx, cancel := commandContext(a) |
| 68 | defer cancel() |
| 69 | body, err := json.Marshal(map[string]any{"intent": "followup", "input": submit, "display": display, "invocations": invocations, "idempotencyKey": idempotency}) |
| 70 | if err != nil { |
| 71 | return InboxReceiptView{}, err |
| 72 | } |
| 73 | resp, err := serveDoForSession(ctx, client, http.MethodPost, serveURL(base, "/inbox/items"), body, path) |
| 74 | if err == nil { |
| 75 | defer resp.Body.Close() |
| 76 | if resp.StatusCode >= 200 && resp.StatusCode < 300 { |
| 77 | var receipt InboxReceiptView |
| 78 | err = json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&receipt) |
| 79 | if err == nil && receipt.ItemID != "" { |
| 80 | return receipt, nil |
| 81 | } |
| 82 | } else { |
| 83 | err = fmt.Errorf("follow-up enqueue failed (%d)", resp.StatusCode) |
| 84 | if resp.StatusCode >= 400 && resp.StatusCode < 500 { |
| 85 | if resp.StatusCode != http.StatusRequestTimeout { |
| 86 | return InboxReceiptView{}, inboxNotSubmitted(err) |
| 87 | } |
| 88 | return InboxReceiptView{}, err |
| 89 | } |
| 90 | } |
| 91 | } |
| 92 | if idempotency != "" { |
| 93 | lookupCtx, lookupCancel := commandContext(a) |
| 94 | defer lookupCancel() |
| 95 | data, lookupErr := serveGet(lookupCtx, client, serveURL(base, "/inbox/receipt?key="+url.QueryEscape(idempotency)+"&session="+url.QueryEscape(path))) |
| 96 | var receipt InboxReceiptView |
| 97 | if lookupErr == nil && json.Unmarshal(data, &receipt) == nil && receipt.ItemID != "" { |
| 98 | return receipt, nil |
| 99 | } |
| 100 | } |
| 101 | if err == nil { |
| 102 | err = fmt.Errorf("follow-up receipt unavailable") |
| 103 | } |
| 104 | return InboxReceiptView{}, err |
| 105 | } |
| 106 |