diff --git a/CLAUDE.md b/CLAUDE.md index bb650f30..d5b09758 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -162,6 +162,15 @@ return an error naming the working read tools instead. `/api/conversations//messages`), the HTTP twin of the `search_messages` MCP tool. Optional `phone`, `conversation_id`, `since`/`until` (YYYY-MM-DD, local time, `until` inclusive to end of day), `limit` (default 50, max 500). +- `POST /api/react` β€” `{conversation_id, message_id, emoji, action}`. On a + v2-primary daemon it takes the v2 IDs the read API returns, queues the + reaction on the durable outbox and waits up to 8 s: 200 delivered, 202 + queued and still retrying (don't repeat it), 502/409 not delivered, 422 + unknown target. Legacy-primary daemons call the platform senders directly. +- `POST /api/v1/outbox/reactions` β€” the durable twin (v2-primary only, + `idempotency_key` required), returns the outbox submission immediately. + Details: [docs/agent-runbook.md](docs/agent-runbook.md) ("Reactions on + v2-primary go through the outbox"). ### Schema diff --git a/cmd/e2e-server/main.go b/cmd/e2e-server/main.go index a2734dbe..9657740f 100644 --- a/cmd/e2e-server/main.go +++ b/cmd/e2e-server/main.go @@ -58,10 +58,49 @@ type e2eServer struct { type e2eAdapter struct { *scripted.Adapter + reactions *e2eReactionScript } func (e2eAdapter) DeclaredCapabilities() bridge.CapabilitySet { - return bridge.CapabilitySet{TextSend: true, MediaSend: true} + return bridge.CapabilitySet{TextSend: true, MediaSend: true, Reactions: true} +} + +// SendReaction gives a scripted account the reaction sender the scripted +// adapter lacks. It accepts each reaction unless a test queued a failure for +// the next one. +func (a e2eAdapter) SendReaction(context.Context, bridge.ReactionRequest) (bridge.SendResult, error) { + if failure := a.reactions.next(); failure != nil { + return bridge.SendResult{}, *failure + } + return bridge.SendResult{AcceptedAt: time.Now()}, nil +} + +type e2eReactionScript struct { + mu sync.Mutex + failures []bridge.OpError +} + +func (s *e2eReactionScript) enqueueFailure(failure bridge.OpError) { + s.mu.Lock() + defer s.mu.Unlock() + s.failures = append(s.failures, failure) +} + +func (s *e2eReactionScript) clear() { + s.mu.Lock() + defer s.mu.Unlock() + s.failures = nil +} + +func (s *e2eReactionScript) next() *bridge.OpError { + s.mu.Lock() + defer s.mu.Unlock() + if len(s.failures) == 0 { + return nil + } + failure := s.failures[0] + s.failures = s.failures[1:] + return &failure } func (s *e2eServer) ServeHTTP(w http.ResponseWriter, r *http.Request) { @@ -125,6 +164,7 @@ func newE2EServer(logger zerolog.Logger) (_ *e2eServer, cleanup func(), resultEr "whatsapp-primary": scripted.New("whatsapp-primary", bridge.PlatformWhatsApp), "signal-primary": scripted.New("signal-primary", bridge.PlatformSignal), } + reactionScripts := map[string]*e2eReactionScript{} for _, adapter := range adapters { for i := 0; i < 128; i++ { adapter.EnqueueMediaResult(bridge.SendResult{ @@ -132,7 +172,8 @@ func newE2EServer(logger zerolog.Logger) (_ *e2eServer, cleanup func(), resultEr AcceptedAt: time.Now(), }) } - if err := registry.Register(e2eAdapter{Adapter: adapter}); err != nil { + reactionScripts[adapter.AccountID()] = &e2eReactionScript{} + if err := registry.Register(e2eAdapter{Adapter: adapter, reactions: reactionScripts[adapter.AccountID()]}); err != nil { _ = v2Store.Close() _ = store.Close() _ = os.RemoveAll(dataDir) @@ -527,6 +568,39 @@ func newE2EServer(logger zerolog.Logger) (_ *e2eServer, cleanup func(), resultEr writeJSON(w, map[string]any{"success": true}) }) + mux.HandleFunc("POST /_e2e/bridges/{account}/next-reaction", func(w http.ResponseWriter, r *http.Request) { + script := reactionScripts[r.PathValue("account")] + if script == nil { + http.Error(w, "unknown bridge account", http.StatusNotFound) + return + } + var req struct { + Result string `json:"next_result"` + } + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + failure := bridge.OpError{Operation: "send_reaction", Cause: fmt.Errorf("scripted %s reaction", req.Result)} + switch req.Result { + case "clear": + script.clear() + writeJSON(w, map[string]any{"success": true}) + return + case "not_connected": + failure.Class, failure.Fingerprint, failure.Dispatch = bridge.FailureTransient, "e2e_not_connected", bridge.DispatchNotCalled + case "rejected": + failure.Class, failure.Fingerprint, failure.Dispatch = bridge.FailureUnsupported, "e2e_rejected", bridge.DispatchNotCalled + case "uncertain": + failure.Class, failure.Fingerprint, failure.Dispatch = bridge.FailureTransient, "e2e_uncertain", bridge.DispatchUncertain + default: + http.Error(w, "next_result must be not_connected, rejected, uncertain or clear", http.StatusBadRequest) + return + } + script.enqueueFailure(failure) + writeJSON(w, map[string]any{"success": true}) + }) + mux.Handle("/", withPrimaryFixtureProjection( withConfirmedMediaProjection( withSyntheticReadReceipts(base, &syntheticReadReceipts), diff --git a/cmd/r5_react_v2_primary_test.go b/cmd/r5_react_v2_primary_test.go new file mode 100644 index 00000000..7c82be89 --- /dev/null +++ b/cmd/r5_react_v2_primary_test.go @@ -0,0 +1,243 @@ +package cmd + +import ( + "bytes" + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "path/filepath" + "slices" + "sync" + "testing" + "time" + + "github.com/rs/zerolog" + + "github.com/maxghenis/openmessage/internal/bridge" + "github.com/maxghenis/openmessage/internal/bridgeadapters/scripted" + "github.com/maxghenis/openmessage/internal/db" + "github.com/maxghenis/openmessage/internal/v2read" + "github.com/maxghenis/openmessage/internal/web" +) + +// r5ReactRouteAdapter gives a scripted account, which has no reaction sender +// of its own, one that records each request and confirms it. +type r5ReactRouteAdapter struct { + *scripted.Adapter + + mu sync.Mutex + requests []bridge.ReactionRequest +} + +func (a *r5ReactRouteAdapter) SendReaction(_ context.Context, req bridge.ReactionRequest) (bridge.SendResult, error) { + a.mu.Lock() + a.requests = append(a.requests, req) + a.mu.Unlock() + return bridge.SendResult{}, nil +} + +func (a *r5ReactRouteAdapter) reactionRequests() []bridge.ReactionRequest { + a.mu.Lock() + defer a.mu.Unlock() + return slices.Clone(a.requests) +} + +// r5LegacyReactionCall is one call into a legacy per-platform reaction +// callback (App.SendSignalReaction / App.SendWhatsAppReaction in production). +type r5LegacyReactionCall struct { + Platform string + ConversationID string + MessageID string +} + +// TestR5ReactRoutesThroughTheOutboxOnV2Primary posts what the web UI posts to +// /api/react on a migrated v2-primary daemon β€” the v2 conversation ID and the +// v2 message ID that the read API handed it β€” for a Signal, a WhatsApp and a +// Google conversation, and requires each reaction to reach that platform's +// transport naming the right conversation and target. +func TestR5ReactRoutesThroughTheOutboxOnV2Primary(t *testing.T) { + fixture := buildR5LegacyFixture(t) + now := time.Date(2026, 7, 17, 18, 0, 0, 0, time.UTC) + runR5Migrate(t, fixture.DataDir, filepath.Join(fixture.DataDir, "v2"), false, now) + + t.Setenv("OPENMESSAGES_DATA_DIR", fixture.DataDir) + t.Setenv("OPENMESSAGES_DEMO", "0") + t.Setenv("OPENMESSAGES_APP_SANDBOX", "1") + t.Setenv("OPENMESSAGES_V2_PRIMARY", "1") + t.Setenv("OPENMESSAGES_V2_SEND", "") + t.Setenv("OPENMESSAGES_V2_INGEST", "") + + stack, err := newV2Stack(v2StackDeps{DataDir: fixture.DataDir, Logger: zerolog.Nop()}) + if err != nil { + t.Fatalf("open migrated v2 stack: %v", err) + } + t.Cleanup(func() { _ = stack.Store.Close() }) + adapters := map[string]*r5ReactRouteAdapter{ + "sms": {Adapter: scripted.New(r5GoogleAccountID, bridge.PlatformGoogle)}, + "whatsapp": {Adapter: scripted.New(r5WhatsAppAccountID, bridge.PlatformWhatsApp)}, + "signal": {Adapter: scripted.New(r5SignalAccountID, bridge.PlatformSignal)}, + } + for platform, adapter := range adapters { + if err := stack.RegisterAdapter(adapter); err != nil { + t.Fatalf("register scripted %s adapter: %v", platform, err) + } + } + inertLegacy, err := db.New(":memory:") + if err != nil { + t.Fatalf("create inert legacy store: %v", err) + } + t.Cleanup(func() { _ = inertLegacy.Close() }) + ctx, cancel := context.WithCancel(context.Background()) + stopStack := stack.Start(ctx, inertLegacy, nil, true) + t.Cleanup(func() { + stopStack() + cancel() + }) + + var ( + legacyMu sync.Mutex + legacyCalls []r5LegacyReactionCall + ) + legacyReaction := func(platform string) func(conversationID, messageID, emoji, action string) error { + return func(conversationID, messageID, _, _ string) error { + legacyMu.Lock() + legacyCalls = append(legacyCalls, r5LegacyReactionCall{platform, conversationID, messageID}) + legacyMu.Unlock() + return nil + } + } + reads := v2read.New(stack.Store) + // The daemon's wiring (cmd/serve.go), with the legacy per-platform reaction + // callbacks recorded and no live Google client. + handler := web.APIHandlerWithOptions(inertLegacy, nil, zerolog.Nop(), nil, web.APIOptions{ + Reads: reads, + V2Primary: true, + V2: &web.V2Options{ + Service: stack.Service, Media: stack.Media, V2Store: stack.Store, + Blobs: stack.Blobs, Registry: stack.Registry, + }, + SendSignalReaction: legacyReaction("signal"), + SendWhatsAppReaction: legacyReaction("whatsapp"), + }) + + for _, conversation := range fixture.Conversations { + adapter := adapters[conversation.Platform] + if adapter == nil || conversation.LegacyID == r5SignalGroup { + continue + } + t.Run(conversation.Platform, func(t *testing.T) { + // The UI reads the thread, then reacts to a message it rendered. + listing := httptest.NewRecorder() + handler.ServeHTTP(listing, httptest.NewRequest( + http.MethodGet, "http://127.0.0.1/api/conversations/"+conversation.V2ID()+"/messages?limit=100", nil, + )) + var listed []*db.Message + if err := json.Unmarshal(listing.Body.Bytes(), &listed); err != nil || listing.Code != http.StatusOK { + t.Fatalf("list messages = %d, %v: %s", listing.Code, err, listing.Body.String()) + } + incoming := conversation.Messages[0] + var target *db.Message + for _, message := range listed { + if message.Body == incoming.Body { + target = message + } + } + if target == nil || target.ConversationID != conversation.V2ID() { + t.Fatalf("listed messages %s lack the incoming fixture message under the v2 conversation ID", listing.Body.String()) + } + + body, err := json.Marshal(map[string]string{ + "conversation_id": target.ConversationID, + "message_id": target.MessageID, + "emoji": "πŸ‘", + "action": "add", + }) + if err != nil { + t.Fatalf("encode reaction request: %v", err) + } + legacyMu.Lock() + legacyBefore := len(legacyCalls) + legacyMu.Unlock() + response := httptest.NewRecorder() + handler.ServeHTTP(response, httptest.NewRequest( + http.MethodPost, "http://127.0.0.1/api/react", bytes.NewReader(body), + )) + t.Logf("POST /api/react %s -> %d %s", body, response.Code, bytes.TrimSpace(response.Body.Bytes())) + + legacyMu.Lock() + reachedLegacy := slices.Clone(legacyCalls[legacyBefore:]) + legacyMu.Unlock() + if len(reachedLegacy) != 0 { + t.Errorf("the reaction went to the legacy transport callback, which reads the frozen legacy store: %+v", reachedLegacy) + } + if response.Code < 200 || response.Code > 299 { + t.Errorf("POST /api/react = %d %s, want a 2xx", response.Code, bytes.TrimSpace(response.Body.Bytes())) + } + + deadline := time.Now().Add(3 * time.Second) + for len(adapter.reactionRequests()) == 0 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + requests := adapter.reactionRequests() + if len(requests) != 1 { + t.Fatalf("the %s transport received %d reaction requests, want 1", conversation.Platform, len(requests)) + } + request := requests[0] + if request.Conversation.RemoteID != conversation.RemoteID || request.Target.RemoteID != incoming.RemoteID || + request.Emoji != "πŸ‘" || request.Action != bridge.ReactionAdd { + t.Fatalf("reaction request = %+v, want πŸ‘ add on %q in %q", request, incoming.RemoteID, conversation.RemoteID) + } + + // The UI reloads the thread after the answer. Nothing echoes the + // reaction back here, so what it shows is the own reaction the + // outbox recorded when the transport accepted it. + reloaded := httptest.NewRecorder() + handler.ServeHTTP(reloaded, httptest.NewRequest( + http.MethodGet, "http://127.0.0.1/api/conversations/"+conversation.V2ID()+"/messages?limit=100", nil, + )) + var after []*db.Message + if err := json.Unmarshal(reloaded.Body.Bytes(), &after); err != nil { + t.Fatalf("reload messages: %v: %s", err, reloaded.Body.String()) + } + for _, message := range after { + if message.MessageID != target.MessageID { + continue + } + // The Google fixture message carries a migrated πŸ‘ with no + // reactor, so compare with what the target showed before. + before, now := r5ThumbsUp(t, target.Reactions), r5ThumbsUp(t, message.Reactions) + if now.Count != before.Count+1 || slices.Contains(before.Actors, "me") || !slices.Contains(now.Actors, "me") { + t.Fatalf("target reactions = %q before, %q after; want one more πŸ‘, by me", target.Reactions, message.Reactions) + } + return + } + t.Fatalf("reloaded thread lacks the target %q", target.MessageID) + }) + } +} + +type r5ReactionGroup struct { + Emoji string `json:"emoji"` + Count int `json:"count"` + Actors []string `json:"actors"` +} + +// r5ThumbsUp returns the πŸ‘ group of a message DTO's Reactions JSON, empty +// when there is none. +func r5ThumbsUp(t *testing.T, raw string) r5ReactionGroup { + t.Helper() + if raw == "" { + return r5ReactionGroup{} + } + var groups []r5ReactionGroup + if err := json.Unmarshal([]byte(raw), &groups); err != nil { + t.Fatalf("decode reactions %q: %v", raw, err) + } + for _, group := range groups { + if group.Emoji == "πŸ‘" { + return group + } + } + return r5ReactionGroup{} +} diff --git a/docs/agent-runbook.md b/docs/agent-runbook.md index 8da4b444..9eaa6600 100644 --- a/docs/agent-runbook.md +++ b/docs/agent-runbook.md @@ -421,9 +421,10 @@ running app: `TestNewClientPerformsNoStoreWrites` (internal/app), `TestRunServeMCPClientDoesNotRepairStore` and `TestOpenCommandReadSourceLegacyDoesNotRepairStore` (cmd). -- **Sends/reactions route through the daemon** (`/api/v1/outbox` on v2, - `/api/send`+`/api/react` on legacy), like the CLI has done since PR #140, - with the same do-not-resend idempotency contract. With the app closed, +- **Sends/reactions route through the daemon** (sends: `/api/v1/outbox` on + v2, `/api/send` on legacy; reactions: `/api/v1/outbox/reactions` on a + v2-primary app, `/api/react` otherwise), like the CLI has done since PR + #140, with the same do-not-resend idempotency contract. With the app closed, send tools return an actionable "start the OpenMessage app" error β€” they never fall back to opening their own connections. - Escape hatches: `--transports` forces the old standalone full-stack stdio @@ -1575,11 +1576,94 @@ incoming message that is a SHA-1. signal-cli 0.14.8 refuses a non-integer while parsing its arguments (`could not convert '' to integer (64 bits)`, exit 1), and the row ended `uncertain` with nothing sent. -No surface submits a reaction to the v2 outbox yet. The web UI's `/api/react` -and the `react_to_message` MCP tool both call the legacy -`signallive.Bridge.SendReaction`, which looks the target up in the legacy -`messages.db` by message ID. Only `MessageService.SendReaction` (tests, so -far) reaches the path above. +On a v2-primary install every reaction from the web UI and the +`react_to_message` MCP tool takes this path; see "Reactions on v2-primary go +through the outbox" below. + +## Reactions on v2-primary go through the outbox + +On a v2-primary daemon the read API hands out v2 IDs (32-hex conversation and +message keys). The legacy reaction senders cannot place those: `/api/react` +used to tell platforms apart by a `signal:`/`whatsapp:` prefix or a +legacy-store lookup and fall back to Google Messages, so before this every +reaction from the UI or MCP to a v2 message, on any platform, went down the +Google path with a v2 message ID: with no Google client it answered +`503 not connected to Google Messages` (reproduced in +`TestR5ReactRoutesThroughTheOutboxOnV2Primary`), and with one it asked Google +to react to an ID that is not a Google message ID. +Legacy-primary daemons, including ones with `OPENMESSAGES_V2_SEND=1`, keep the +legacy senders unchanged. + +**Routes.** + +- `POST /api/react` `{conversation_id, message_id, emoji, action}` resolves the + target by its v2 message ID and queues the reaction on the outbox + (`MessageService.SendReaction`). The target decides the account and the + conversation; `conversation_id` may be omitted, and when given must name the + target's own conversation, by v2 ID or by its remote ID (`signal:+1…`, a + WhatsApp JID, a Google thread id). It then waits up to 8 s for the reaction + to settle, and the answer says how far it got: + + | Status | Body | Meaning | + |---|---|---| + | 200 | `success: true` | The platform accepted it (`confirmed`, or `store_failed`). | + | 202 | `success: true, queued: true` | Stored, not delivered yet (`queued`, `dispatching`, `not_dispatched`). The app keeps retrying it; do not react again. | + | 502 | `error` | `uncertain` (the transport call ended without an answer: it may have been applied) or `rejected` (the platform refused it). | + | 409 | `error` | `canceled` before it was sent. | + | 422 | `reaction_target_unavailable` | No v2 message with that ID, or not in that conversation. Nothing was queued. | + | 501 | `error` | The account has no reaction sender (an import-only platform). | + + Every 2xx/409/502 answer carries `outbox_id` and `state`. The legacy handler + answered only after the transport call; this one returns while delivery may + still be in progress, which is what the 202 is for. A request without + `idempotency_key` mints a new key, so repeating the POST queues a second + reaction. +- `POST /api/v1/outbox/reactions` (same body plus a required + `idempotency_key`) is the durable twin, like `/api/v1/outbox/messages`: it + returns the submission at once and `GET /api/v1/outbox/` follows it. A + legacy-primary daemon answers it `409 legacy primary: use /api/react`. +- The MCP tool, in-process on a v2-primary daemon, queues on the outbox and + waits up to 25 s. In transportless client mode it probes `/api/status` and, + when the app reports `v2_primary`, submits to `/api/v1/outbox/reactions` + and polls the delivery; otherwise it uses `/api/react`. Its result carries + `outbox_id`, `state`, `idempotency_key` and `settled`; `settled: false` is not + an error and means do not react again. + +**What the user sees.** The adapters store nothing when they send a reaction +(the legacy senders updated the stored reactions themselves after the send). +So when the transport accepts a reaction, the outbox confirms the row and +writes the account's own reaction (`reactor_key = 'self'`, shown as "You") in +the same SQLite transaction (`OutboxRepository.ConfirmReaction`); without it +the reaction would appear only if a transport later reported it. Its time is +when the transport accepted it, and a later report for the same reactor key +replaces it when it carries a later time (the usual `ApplyReaction` +ordering). On Google Messages the phone's copy is authoritative instead: +every Google frame for a message applies that message's full reaction +snapshot (`ReplaceEmbeddedReactions`), which removes active reactors the +snapshot does not list, and the Google decoder names reactors by participant +ID rather than as self. So the next frame for the message replaces the "self" +row with the reactor the phone names, or removes it if the phone's copy does +not have the reaction yet; a reaction can disappear briefly and come back. A reaction that ends +`uncertain`, `rejected`, `not_dispatched` or `canceled` writes nothing, so the +thread shows only reactions that went out. While a reaction is `queued`, +`dispatching` or `not_dispatched` it is listed in the web UI's outbox tray +with a Cancel button. + +**Reactions stuck retrying.** A reaction the transport cannot deliver yet goes +`not_dispatched` and retries (every 5 s by default) with no cap; a Signal +reaction whose target can never be named (above) stays there until it is +canceled. Find them with: + +```bash +curl -s "http://127.0.0.1:7007/api/v1/outbox?limit=100" | jq '.[] | select(.kind=="reaction")' +``` + +and cancel one with `POST /api/v1/outbox//cancel`. + +**A count one higher than the reactors shown.** A migrated message whose legacy +reaction named nobody (`[{"emoji":"πŸ‘","count":1}]`) keeps it as an anonymous +reactor (`reactor_key = 'anon:πŸ‘'`, `migration/transform.go`). Reacting πŸ‘ to it on v2 adds "me" beside +it, so the pill reads 2 with only "You" listed. ## Deploying a new build to a live install diff --git a/e2e/ui.spec.js b/e2e/ui.spec.js index 4320f472..3ede5601 100644 --- a/e2e/ui.spec.js +++ b/e2e/ui.spec.js @@ -1070,6 +1070,98 @@ test('maps durable outbox states to honest tray labels and actions', async ({ pa } }); +test('lists a reaction in the tray only while the app is still sending it', async ({ page }) => { + await page.waitForFunction(() => window.__openMessageTestHooks?.outboxRowView); + const views = await page.evaluate(() => { + const now = 1_700_000_000_000; + const view = window.__openMessageTestHooks.outboxRowView; + const reaction = state => view({ kind: 'reaction', state, scheduled_for_ms: now }, now); + return { + queued: reaction('queued'), + dispatching: reaction('dispatching'), + retrying: reaction('not_dispatched'), + uncertain: reaction('uncertain'), + repairing: reaction('store_failed'), + confirmed: reaction('confirmed'), + rejected: reaction('rejected'), + }; + }); + + expect(views.queued).toMatchObject({ visible: true, label: 'Sending…', action: 'cancel' }); + expect(views.queued.guidance).toContain('do not react again'); + expect(views.dispatching).toMatchObject({ visible: true, label: 'Sending…', action: '' }); + expect(views.retrying).toMatchObject({ visible: true, label: 'Retrying…', action: 'cancel' }); + // A reaction cannot be sent again from the tray, so one the app has stopped + // sending is not listed there with an action that would fail. + for (const done of [views.uncertain, views.repairing, views.confirmed, views.rejected]) { + expect(done).toMatchObject({ visible: false, action: '' }); + } +}); + +async function reactToLastReceivedMessage(page) { + const target = page.locator('#messages-area .msg.received').last(); + const messageID = await target.getAttribute('data-msg-id'); + const message = page.locator(`#messages-area .msg[data-msg-id="${messageID}"]`); + const answered = page.waitForResponse(response => + response.request().method() === 'POST' && /\/api\/react(\?|$)/.test(response.url())); + await message.hover(); + await message.locator('.action-react').click(); + const emoji = (await message.locator('.emoji-picker button').first().textContent()).trim(); + await message.locator('.emoji-picker button').first().click(); + return { message, emoji, response: await answered }; +} + +test('sends a reaction through the outbox and shows it on the message', async ({ page }) => { + await openConversation(page, 'Sarah Chen'); + + const { message, emoji, response } = await reactToLastReceivedMessage(page); + expect(response.status()).toBe(200); + expect(await response.json()).toMatchObject({ success: true, state: 'confirmed' }); + + // Nothing echoes the reaction back in this fixture: the pill is the own + // reaction the app recorded when the transport accepted it. + const pill = message.locator('.reaction-pill').filter({ hasText: emoji }); + await expect(pill).toHaveCount(1); + await expect(pill).toHaveAttribute('title', /You/); +}); + +test('reports a queued reaction and lists it in the outbox while the platform is down', async ({ page, request }) => { + await openConversation(page, 'Sarah Chen'); + // Enough refusals that the app's automatic retries cannot deliver the + // reaction before the test cancels it. + for (let i = 0; i < 12; i++) { + await request.post('/_e2e/bridges/google-primary/next-reaction', { data: { next_result: 'not_connected' } }); + } + try { + const { response } = await reactToLastReceivedMessage(page); + expect(response.status()).toBe(202); + expect(await response.json()).toMatchObject({ success: true, queued: true, state: 'not_dispatched' }); + await expect(page.locator('#thread-feedback')).toContainText('Reaction not sent yet'); + + await page.locator('#outbox-toggle-btn').click(); + const row = page.locator('#outbox-tray .outbox-row').filter({ hasText: 'Reaction' }); + await expect(row).toHaveCount(1); + await expect(row).toContainText('Retrying…'); + await row.locator('.outbox-row-action').click(); + await expect(row).toHaveCount(0); + } finally { + await request.post('/_e2e/bridges/google-primary/next-reaction', { data: { next_result: 'clear' } }); + } +}); + +test('says so when the platform refuses a reaction', async ({ page, request }) => { + await openConversation(page, 'Sarah Chen'); + await request.post('/_e2e/bridges/google-primary/next-reaction', { data: { next_result: 'rejected' } }); + try { + const { response } = await reactToLastReceivedMessage(page); + expect(response.status()).toBe(502); + expect(await response.json()).toMatchObject({ success: false, state: 'rejected' }); + await expect(page.locator('#thread-feedback')).toContainText('refused the reaction'); + } finally { + await request.post('/_e2e/bridges/google-primary/next-reaction', { data: { next_result: 'clear' } }); + } +}); + test('keeps existing thread nodes mounted after sending', async ({ page }) => { const outbound = `Stable send ${Date.now()}`; diff --git a/internal/bridgeadapters/google/send.go b/internal/bridgeadapters/google/send.go index a08a15e2..73b5b4c8 100644 --- a/internal/bridgeadapters/google/send.go +++ b/internal/bridgeadapters/google/send.go @@ -327,8 +327,8 @@ func (a *Adapter) SendReaction( } } - // Reaction dispatch confirms an empty result via ConfirmWithoutResult; - // there is no reaction-message ID or echo-reconciliation consumer. + // Reaction dispatch confirms an empty result via ConfirmReaction; there + // is no reaction-message ID or echo-reconciliation consumer. return bridge.SendResult{}, nil } diff --git a/internal/ingest/worker.go b/internal/ingest/worker.go index ceca40dc..32b781ac 100644 --- a/internal/ingest/worker.go +++ b/internal/ingest/worker.go @@ -1451,9 +1451,9 @@ func (w *Worker) prepareReactionActor( ) (preparedReactionActor, error) { if reference.IsSelf { return preparedReactionActor{ - key: "self", + key: sqlite.SelfReactorKey, isSelf: true, - label: "me", + label: sqlite.SelfReactorLabel, }, nil } if identityRaw(reference) == "" { diff --git a/internal/localapi/localapi.go b/internal/localapi/localapi.go index b9f623b7..ba66b45e 100644 --- a/internal/localapi/localapi.go +++ b/internal/localapi/localapi.go @@ -93,6 +93,14 @@ func (s DaemonStatus) SendsViaOutbox() bool { return s.V2Send || s.V2Primary } +// ReactionsViaOutbox reports whether the daemon takes reactions on the durable +// /api/v1/outbox/reactions surface. Only a v2-primary daemon does: the route +// addresses messages by v2 ID, which a legacy-primary daemon does not hand +// out, even with v2 send enabled. +func (s DaemonStatus) ReactionsViaOutbox() bool { + return s.V2Primary +} + // Status probes /api/status. The second result reports reachability: false // means no daemon answered at all (connection refused/timeout), while true // with a non-nil error means something answered but the response was not a @@ -135,6 +143,18 @@ type MediaSubmission struct { Content io.Reader } +// ReactionSubmission is a durable reaction routed at +// POST /api/v1/outbox/reactions. MessageID is the target's v2 message ID. +// ConversationID is optional; when set it must name the target's conversation. +// Action is "add", "remove", or "switch", and empty means add. +type ReactionSubmission struct { + ConversationID string `json:"conversation_id,omitempty"` + MessageID string `json:"message_id"` + Emoji string `json:"emoji"` + Action string `json:"action,omitempty"` + IdempotencyKey string `json:"idempotency_key"` +} + // Submission mirrors the daemon's v1 submission response. type Submission struct { OutboxID string `json:"outbox_id"` @@ -195,6 +215,20 @@ func (c *Client) SubmitText(ctx context.Context, submission TextSubmission) (Sub return result, nil } +// SubmitReaction enqueues a reaction on a v2-primary daemon's durable outbox. +// Like SubmitText it returns once the intent is stored; Delivery and +// WaitDelivery report what became of it. +func (c *Client) SubmitReaction(ctx context.Context, submission ReactionSubmission) (Submission, error) { + var result Submission + if err := c.postJSON(ctx, "/api/v1/outbox/reactions", submission, &result); err != nil { + return Submission{}, err + } + if result.OutboxID == "" { + return Submission{}, fmt.Errorf("submission response omitted outbox_id") + } + return result, nil +} + // SubmitMedia enqueues a media send on the daemon's durable outbox, streaming // Content as a multipart upload. func (c *Client) SubmitMedia(ctx context.Context, submission MediaSubmission) (Submission, error) { @@ -372,16 +406,37 @@ func (c *Client) LegacySendMedia(ctx context.Context, submission MediaSubmission return result, nil } +// ReactResult is the daemon's answer on /api/react. +// +// A legacy-primary daemon answers only after the transport call, and sets +// Success alone. A v2-primary daemon queues the reaction on its durable outbox +// and waits a few seconds for it to settle, so its answer also carries the +// intent: Queued is true (HTTP 202) when the reaction is stored but not yet +// delivered, in which case the daemon keeps retrying it and the caller must +// not react again. OutboxID and State identify the intent for Delivery. +type ReactResult struct { + Success bool `json:"success"` + Queued bool `json:"queued"` + OutboxID string `json:"outbox_id"` + State string `json:"state"` +} + // React routes a reaction through the daemon's /api/react surface, which -// works in every daemon mode. Action is "add", "remove", or "switch". -func (c *Client) React(ctx context.Context, conversationID, messageID, emoji, action string) error { +// works in every daemon mode. Action is "add", "remove", or "switch". A +// reaction the daemon could not deliver is a ResponseError. On a v2-primary +// daemon prefer SubmitReaction, which takes an idempotency key. +func (c *Client) React(ctx context.Context, conversationID, messageID, emoji, action string) (ReactResult, error) { payload := struct { ConversationID string `json:"conversation_id"` MessageID string `json:"message_id"` Emoji string `json:"emoji"` Action string `json:"action"` }{conversationID, messageID, emoji, action} - return c.postJSON(ctx, "/api/react", payload, &map[string]any{}) + var result ReactResult + if err := c.postJSON(ctx, "/api/react", payload, &result); err != nil { + return ReactResult{}, err + } + return result, nil } func (c *Client) postJSON(ctx context.Context, path string, payload any, target any) error { diff --git a/internal/localapi/localapi_test.go b/internal/localapi/localapi_test.go index 65f6e7cb..adaf0ac6 100644 --- a/internal/localapi/localapi_test.go +++ b/internal/localapi/localapi_test.go @@ -320,3 +320,66 @@ func TestDefaultBaseURLHonorsPortEnv(t *testing.T) { t.Fatalf("DefaultBaseURL() = %q", got) } } + +func TestReactionsViaOutboxNeedsV2Primary(t *testing.T) { + for _, status := range []DaemonStatus{ + {}, {V2Send: true}, {V2Primary: true}, {V2Send: true, V2Primary: true}, + } { + if got := status.ReactionsViaOutbox(); got != status.V2Primary { + t.Fatalf("%+v.ReactionsViaOutbox() = %v, want %v", status, got, status.V2Primary) + } + } +} + +func TestSubmitReactionPostsToTheOutbox(t *testing.T) { + var path string + var body map[string]any + client := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + path = r.Method + " " + r.URL.Path + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + t.Errorf("decode: %v", err) + } + json.NewEncoder(w).Encode(map[string]any{"outbox_id": "out", "state": "queued", "deduplicated": true}) + })) + submission, err := client.SubmitReaction(context.Background(), ReactionSubmission{ + MessageID: "msg", Emoji: "πŸ‘", IdempotencyKey: "key", + }) + if err != nil || submission.OutboxID != "out" || !submission.Deduplicated { + t.Fatalf("SubmitReaction() = %+v, %v", submission, err) + } + if path != "POST /api/v1/outbox/reactions" { + t.Fatalf("request = %s", path) + } + // conversation_id and action are optional and omitted when empty. + if len(body) != 3 || body["message_id"] != "msg" || body["emoji"] != "πŸ‘" || body["idempotency_key"] != "key" { + t.Fatalf("body = %v", body) + } + + empty := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + json.NewEncoder(w).Encode(map[string]any{"state": "queued"}) + })) + if _, err := empty.SubmitReaction(context.Background(), ReactionSubmission{MessageID: "m", Emoji: "πŸ‘", IdempotencyKey: "k"}); err == nil { + t.Fatal("expected an error for a response without outbox_id") + } +} + +func TestReactDecodesTheQueuedAnswer(t *testing.T) { + client := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusAccepted) + json.NewEncoder(w).Encode(map[string]any{"success": true, "queued": true, "outbox_id": "out", "state": "not_dispatched"}) + })) + result, err := client.React(context.Background(), "conv", "msg", "πŸ‘", "add") + if err != nil || result != (ReactResult{Success: true, Queued: true, OutboxID: "out", State: "not_dispatched"}) { + t.Fatalf("React() = %+v, %v", result, err) + } + + refused := testClient(t, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusBadGateway) + json.NewEncoder(w).Encode(map[string]any{"success": false, "error": "send reaction: the transport refused the reaction"}) + })) + if _, err := refused.React(context.Background(), "conv", "msg", "πŸ‘", "add"); err == nil { + t.Fatal("expected a ResponseError for a refused reaction") + } else if responseErr, ok := AsResponseError(err); !ok || responseErr.StatusCode != http.StatusBadGateway { + t.Fatalf("React() error = %v", err) + } +} diff --git a/internal/messaging/dispatch.go b/internal/messaging/dispatch.go index 85606977..aaf6c3ed 100644 --- a/internal/messaging/dispatch.go +++ b/internal/messaging/dispatch.go @@ -557,44 +557,30 @@ func (s *MessageService) dispatchReactionLease(ctx context.Context, outboxLease if sendErr != nil { return s.recordSendError(mutationCtx, item, sendErr, reactionOperation) } - if strings.TrimSpace(result.RemoteMessageID) == "" { - // Legacy web/api.go reaction dispatch returns only a success boolean, so - // an empty remote result is a valid confirmation for this operation. - if err := s.outbox.ConfirmWithoutResult( - mutationCtx, - item.OutboxID, - *item.LeaseToken, - ); err != nil { - return fmt.Errorf( - "dispatch outbox item %q: confirm reaction without result: %w", - item.OutboxID, - err, - ) - } - s.signalChange() - return nil - } - - confirmation := sqlite.Confirmation{ + // The transport accepted the reaction. Confirming also records it in the + // read model as this account's own, in one transaction. The adapters store + // nothing when they send a reaction (the legacy senders wrote the stored + // reaction themselves after the send), so without this row a reader sees + // the reaction only if the transport later reports it back. The Google, + // WhatsApp and Signal adapters return no remote identity for a reaction, + // so an empty result is a valid confirmation. + occurredAt := result.AcceptedAt + if occurredAt.IsZero() { + occurredAt = s.clock.Now() + } + if _, err := s.outbox.ConfirmReaction(mutationCtx, sqlite.ReactionConfirmation{ OutboxID: item.OutboxID, LeaseToken: *item.LeaseToken, ResultRemoteID: result.RemoteMessageID, - } - if err := s.outbox.Confirm(mutationCtx, confirmation); err != nil { - storeErr := s.outbox.MarkStoreFailed( - mutationCtx, + OccurredAt: occurredAt, + }); err != nil { + // The row stays leased and called, so lease recovery settles it as + // uncertain: the reaction went out, and nothing here may send it again. + return fmt.Errorf( + "dispatch outbox item %q: confirm reaction: %w", item.OutboxID, - *item.LeaseToken, - result.RemoteMessageID, - err.Error(), + err, ) - if storeErr != nil { - return fmt.Errorf( - "dispatch outbox item %q: confirm: %w", - item.OutboxID, - errors.Join(err, storeErr), - ) - } } s.signalChange() return nil diff --git a/internal/messaging/reaction_projection_test.go b/internal/messaging/reaction_projection_test.go new file mode 100644 index 00000000..2be62ae6 --- /dev/null +++ b/internal/messaging/reaction_projection_test.go @@ -0,0 +1,372 @@ +package messaging + +import ( + "context" + "fmt" + "math/rand" + "reflect" + "testing" + "testing/quick" + "time" + + "github.com/maxghenis/openmessage/internal/bridge" + "github.com/maxghenis/openmessage/internal/storage/sqlite" +) + +// reactionProjectionHarness runs reactions through the real service and +// dispatcher against a scripted transport, and reads back what a reader of +// the store sees. +type reactionProjectionHarness struct { + t *testing.T + clock *manualClock + store *sqlite.Store + service *MessageService + sender *scriptedReactionSender + reactions *sqlite.ReactionRepository + targets int + keys int +} + +const reactionProjectionOther = "identity-reaction-other" + +func newReactionProjectionHarness(t *testing.T) *reactionProjectionHarness { + t.Helper() + clock := newManualClock(messagingTestTime) + store := openMessagingTestStore(t, clock.Now()) + seedDispatchIdentity(t, store, "identity-reaction-author", "author@example.test", clock.Now()) + if err := store.UpsertIdentity(sqlite.Identity{ + IdentityID: reactionProjectionOther, AccountID: "account-1", + Kind: sqlite.IdentityKind("test_address"), CanonicalValue: "other@example.test", + RawValue: "other@example.test", MetadataJSON: `{}`, + CreatedAtMS: clock.Now().UnixMilli(), UpdatedAtMS: clock.Now().UnixMilli(), + }); err != nil { + t.Fatalf("UpsertIdentity(other reactor): %v", err) + } + sender := &scriptedReactionSender{} + registry := newScriptedRegistry("reaction-projection", &scriptedTextSender{}) + registry.setReactionSender(sender) + registry.setAvailable(true) + reactions, err := sqlite.NewReactionRepository(store, clock.Now) + if err != nil { + t.Fatalf("NewReactionRepository(): %v", err) + } + return &reactionProjectionHarness{ + t: t, clock: clock, store: store, sender: sender, reactions: reactions, + service: newMessagingTestService(t, store, registry, clock), + } +} + +// target stores a fresh incoming message to react to. +func (h *reactionProjectionHarness) target() string { + h.t.Helper() + h.targets++ + author := "identity-reaction-author" + message := mustProjectDispatchMessage(h.t, h.store, h.clock, sqlite.Message{ + MessageID: fmt.Sprintf("message-reaction-projection-%d", h.targets), + ConversationID: "conversation-1", + AccountID: "account-1", + RemoteMessageID: fmt.Sprintf("remote-reaction-projection-%d", h.targets), + SenderIdentityID: &author, + Direction: sqlite.MessageDirectionIncoming, + Body: "react to this", + State: sqlite.MessageStateActive, + OccurredAtMS: h.clock.Now().Add(-time.Minute).UnixMilli(), + }) + return message.MessageID +} + +// react submits a reaction and runs the dispatcher once with the transport +// scripted to answer as given. +func (h *reactionProjectionHarness) react( + targetID, emoji string, + action bridge.ReactionAction, + answer sendStep, +) Submission { + h.t.Helper() + h.keys++ + h.sender.mu.Lock() + h.sender.steps = append(h.sender.steps, answer) + h.sender.mu.Unlock() + submission := mustSendDispatchReaction(h.t, h.service, SendReactionCommand{ + CommonCommand: testCommonCommand(fmt.Sprintf("reaction-projection-%d", h.keys)), + TargetMessageID: targetID, + Emoji: emoji, + Action: action, + }) + if processed, err := h.service.DispatchDue(context.Background(), 4); err != nil || processed != 1 { + h.t.Fatalf("DispatchDue() = %d, %v; want exactly the reaction just queued", processed, err) + } + return submission +} + +// visible returns the active reactions on a message as a reader gets them: +// reactor ("me" for this account, else the reactor's address) to emoji. +func (h *reactionProjectionHarness) visible(targetID string) map[string]string { + h.t.Helper() + rows, err := h.reactions.ReactionsForMessages(context.Background(), []string{targetID}) + if err != nil { + h.t.Fatalf("ReactionsForMessages(): %v", err) + } + visible := map[string]string{} + for _, row := range rows[targetID] { + reactor := row.ReactorCanonical + if row.ReactorIsSelf { + reactor = sqlite.SelfReactorLabel + } + if _, duplicate := visible[reactor]; duplicate { + h.t.Fatalf("reactor %q has more than one active reaction on %q: %+v", reactor, targetID, rows[targetID]) + } + visible[reactor] = row.Emoji + } + return visible +} + +func TestDeliveredReactionShowsAsThisAccountsOwn(t *testing.T) { + h := newReactionProjectionHarness(t) + target := h.target() + + // The transport says when it accepted the reaction; that time orders it. + acceptedAt := h.clock.Now().Add(-3 * time.Second) + added := h.react(target, "πŸ‘", bridge.ReactionAdd, sendStep{result: bridge.SendResult{AcceptedAt: acceptedAt}}) + if delivery := mustDelivery(t, h.service, added.OutboxID); delivery.State != OutboxConfirmed { + t.Fatalf("delivery = %+v, want confirmed", delivery) + } + if got, want := h.visible(target), (map[string]string{"me": "πŸ‘"}); !reflect.DeepEqual(got, want) { + t.Fatalf("visible reactions = %v, want %v", got, want) + } + rows, err := h.reactions.ReactionsForMessages(context.Background(), []string{target}) + if err != nil || rows[target][0].OccurredAtMS != acceptedAt.UnixMilli() { + t.Fatalf("own reaction time = %+v, %v; want the transport's accepted time %d", rows[target], err, acceptedAt.UnixMilli()) + } + + // With no accepted time from the transport, the dispatcher's clock is it. + h.clock.Advance(time.Second) + h.react(target, "❀️", bridge.ReactionSwitch, sendStep{}) + if got, want := h.visible(target), (map[string]string{"me": "❀️"}); !reflect.DeepEqual(got, want) { + t.Fatalf("visible reactions after switch = %v, want %v", got, want) + } + rows, err = h.reactions.ReactionsForMessages(context.Background(), []string{target}) + if err != nil || rows[target][0].OccurredAtMS != h.clock.Now().UnixMilli() { + t.Fatalf("switched reaction time = %+v, %v; want the dispatcher's clock", rows[target], err) + } + + h.clock.Advance(time.Second) + h.react(target, "❀️", bridge.ReactionRemove, sendStep{}) + if got := h.visible(target); len(got) != 0 { + t.Fatalf("visible reactions after remove = %v, want none", got) + } +} + +func TestUndeliveredReactionChangesNothingAReaderSees(t *testing.T) { + for _, undelivered := range []struct { + name string + answer sendStep + want OutboxState + }{ + { + name: "the platform is disconnected", + answer: sendStep{err: bridge.OpError{ + Class: bridge.FailureTransient, Operation: "send_reaction", Dispatch: bridge.DispatchNotCalled, + }}, + want: OutboxNotDispatched, + }, + { + name: "the platform refuses it", + answer: sendStep{err: bridge.OpError{ + Class: bridge.FailureUnsupported, Operation: "send_reaction", Dispatch: bridge.DispatchNotCalled, + }}, + want: OutboxRejected, + }, + { + name: "the transport call ends without an answer", + answer: sendStep{err: bridge.OpError{ + Class: bridge.FailureTransient, Operation: "send_reaction", Dispatch: bridge.DispatchUncertain, + }}, + want: OutboxUncertain, + }, + } { + t.Run(undelivered.name, func(t *testing.T) { + h := newReactionProjectionHarness(t) + target := h.target() + // An earlier delivered reaction must survive the failed one. + h.react(target, "πŸ‘", bridge.ReactionAdd, sendStep{}) + h.clock.Advance(time.Second) + + failed := h.react(target, "πŸ˜‚", bridge.ReactionSwitch, undelivered.answer) + if delivery := mustDelivery(t, h.service, failed.OutboxID); delivery.State != undelivered.want { + t.Fatalf("delivery = %+v, want %q", delivery, undelivered.want) + } + if got, want := h.visible(target), (map[string]string{"me": "πŸ‘"}); !reflect.DeepEqual(got, want) { + t.Fatalf("visible reactions = %v, want the earlier reaction only %v", got, want) + } + }) + } + + t.Run("it is canceled before it is sent", func(t *testing.T) { + h := newReactionProjectionHarness(t) + target := h.target() + submission := mustSendDispatchReaction(t, h.service, SendReactionCommand{ + CommonCommand: testCommonCommand("reaction-projection-canceled"), + TargetMessageID: target, + Emoji: "πŸ‘", + }) + if _, err := h.service.Cancel(context.Background(), submission.OutboxID); err != nil { + t.Fatalf("Cancel(): %v", err) + } + if processed, err := h.service.DispatchDue(context.Background(), 4); err != nil || processed != 0 { + t.Fatalf("DispatchDue() = %d, %v; want nothing to dispatch", processed, err) + } + if got := h.visible(target); len(got) != 0 || h.sender.requestCount() != 0 { + t.Fatalf("visible reactions = %v after %d transport calls, want none", got, h.sender.requestCount()) + } + }) +} + +// reactionHistoryStep is one event in a message's reaction history. +type reactionHistoryStep struct { + // Kind: 0 this account reacts and the transport accepts; 1 this account + // reacts and the transport does not deliver; 2 the transport reports this + // account's reaction from another device; 3 it reports another person's. + Kind int + Emoji string + Action bridge.ReactionAction + Outcome int + // Order is the reaction's place in time among the history's events. Every + // step has a different one, and it is unrelated to the step's position, + // so events are applied out of time order. + Order int +} + +type reactionHistory []reactionHistoryStep + +func (reactionHistory) Generate(r *rand.Rand, _ int) reflect.Value { + emoji := []string{"πŸ‘", "❀️", "πŸ˜‚"} + actions := []bridge.ReactionAction{bridge.ReactionAdd, bridge.ReactionRemove, bridge.ReactionSwitch} + history := make(reactionHistory, 1+r.Intn(10)) + order := r.Perm(len(history)) + for i := range history { + history[i] = reactionHistoryStep{ + Kind: r.Intn(4), + Emoji: emoji[r.Intn(len(emoji))], + Action: actions[r.Intn(len(actions))], + Outcome: r.Intn(3), + Order: order[i], + } + } + return reflect.ValueOf(history) +} + +// TestReactionsAReaderSeesAreTheLatestDeliveredOnes states what the read +// model means once reactions go through the outbox. For any history of a +// message's reactions, applied in any order: +// +// - each reactor shows at most one reaction; +// - it is the latest, by the reaction's own time, of that reactor's +// reactions that were delivered here (accepted by the transport, or +// reported by it), and a latest "remove" shows nothing; +// - a reaction the transport did not deliver changes nothing; +// - every reaction this account submitted reaches the transport exactly +// once and ends in the outbox state its answer calls for. +func TestReactionsAReaderSeesAreTheLatestDeliveredOnes(t *testing.T) { + h := newReactionProjectionHarness(t) + ctx := context.Background() + undelivered := []struct { + answer sendStep + want OutboxState + }{ + {sendStep{err: bridge.OpError{Class: bridge.FailureUnsupported, Operation: "send_reaction", Dispatch: bridge.DispatchNotCalled}}, OutboxRejected}, + {sendStep{err: bridge.OpError{Class: bridge.FailureTransient, Operation: "send_reaction", Dispatch: bridge.DispatchUncertain}}, OutboxUncertain}, + {sendStep{err: bridge.OpError{Class: bridge.FailureTransient, Operation: "send_reaction", Dispatch: bridge.DispatchNotCalled}}, OutboxNotDispatched}, + } + + property := func(history reactionHistory) bool { + target := h.target() + base := h.clock.Now() + type latest struct { + order int + emoji string + } + model := map[string]latest{} + deliver := func(reactor string, step reactionHistoryStep) { + if seen, ok := model[reactor]; ok && seen.order > step.Order { + return + } + emoji := step.Emoji + if step.Action == bridge.ReactionRemove { + emoji = "" + } + model[reactor] = latest{order: step.Order, emoji: emoji} + } + calls := h.sender.requestCount() + submitted := 0 + + for index, step := range history { + h.clock.Advance(time.Millisecond) + at := base.Add(time.Duration(step.Order+1) * time.Hour) + switch step.Kind { + case 0: + submission := h.react(target, step.Emoji, step.Action, sendStep{result: bridge.SendResult{AcceptedAt: at}}) + submitted++ + if state := mustDelivery(t, h.service, submission.OutboxID).State; state != OutboxConfirmed { + t.Errorf("step %d: delivered reaction state = %q, want confirmed", index, state) + return false + } + deliver("me", step) + case 1: + outcome := undelivered[step.Outcome] + submission := h.react(target, step.Emoji, step.Action, outcome.answer) + submitted++ + if state := mustDelivery(t, h.service, submission.OutboxID).State; state != outcome.want { + t.Errorf("step %d: undelivered reaction state = %q, want %q", index, state, outcome.want) + return false + } + if outcome.want == OutboxNotDispatched { + // The dispatcher would retry it; take it out of play so the + // transport sees each reaction of this history once. + if _, err := h.service.Cancel(ctx, submission.OutboxID); err != nil { + t.Errorf("step %d: Cancel(): %v", index, err) + return false + } + } + default: + report := sqlite.ReactionApply{ + AccountID: "account-1", ConversationID: "conversation-1", MessageID: target, + ReactorKey: sqlite.SelfReactorKey, ReactorIsSelf: true, ReactorLabel: sqlite.SelfReactorLabel, + Emoji: step.Emoji, Action: step.Action, + OccurredAtMS: at.UnixMilli(), SourceSeqMS: h.clock.Now().UnixMilli(), + } + reactor := "me" + if step.Kind == 3 { + other := reactionProjectionOther + report.ReactorKey, report.ReactorIdentityID = other, &other + report.ReactorIsSelf, report.ReactorLabel = false, "" + reactor = "other@example.test" + } + if _, err := h.reactions.ApplyReaction(ctx, report); err != nil { + t.Errorf("step %d: ApplyReaction(): %v", index, err) + return false + } + deliver(reactor, step) + } + } + + want := map[string]string{} + for reactor, seen := range model { + if seen.emoji != "" { + want[reactor] = seen.emoji + } + } + if got := h.visible(target); !reflect.DeepEqual(got, want) { + t.Errorf("history %+v\nvisible reactions = %v, want %v", history, got, want) + return false + } + if got := h.sender.requestCount() - calls; got != submitted { + t.Errorf("history %+v\ntransport calls = %d, want one per submitted reaction (%d)", history, got, submitted) + return false + } + return true + } + if err := quick.Check(property, &quick.Config{MaxCount: 150, Rand: rand.New(rand.NewSource(20261009))}); err != nil { + t.Fatal(err) + } +} diff --git a/internal/migration/transform.go b/internal/migration/transform.go index 787612ad..2fb83523 100644 --- a/internal/migration/transform.go +++ b/internal/migration/transform.go @@ -498,7 +498,7 @@ func planLegacyReactions(dataset legacyDataset, state *transformState, report *R normalizedActor := normalizeLegacyReactionActor(platform, actor) switch { case strings.EqualFold(actor, "me") || isOwnReactionActor(platform, actor, normalizedActor, ownNumbers[platform]): - row.ReactorKey, row.ReactorIsSelf, row.ReactorLabel = "self", true, "me" + row.ReactorKey, row.ReactorIsSelf, row.ReactorLabel = sqlite.SelfReactorKey, true, sqlite.SelfReactorLabel case actor == "": row.ReactorKey = "anon:" + emoji default: diff --git a/internal/storage/sqlite/outbox.go b/internal/storage/sqlite/outbox.go index b0b00460..7e322fd9 100644 --- a/internal/storage/sqlite/outbox.go +++ b/internal/storage/sqlite/outbox.go @@ -9,6 +9,8 @@ import ( "fmt" "strings" "time" + + "github.com/maxghenis/openmessage/internal/bridge" ) // OutboxKind identifies the shape of an outbound intent. @@ -210,6 +212,17 @@ type Confirmation struct { TransportResultID string } +// ReactionConfirmation records that the transport accepted a reaction for the +// active lease. ResultRemoteID is empty when the transport returns no identity +// for a reaction. OccurredAt is when the transport accepted it, and orders the +// reaction against other mutations of the same reactor. +type ReactionConfirmation struct { + OutboxID string + LeaseToken string + ResultRemoteID string + OccurredAt time.Time +} + // ReconcileRequest identifies one accepted transport request and its // authoritative remote result without requiring an active dispatch lease. type ReconcileRequest struct { @@ -1531,6 +1544,113 @@ func (r *OutboxRepository) ConfirmWithoutResult( return r.requireLeaseMutation(ctx, "confirm without result", outboxID, result) } +// ConfirmReaction terminally confirms a reaction's active called lease and, in +// the same transaction, applies the reaction to the read model as this +// account's own. Sending a reaction stores nothing else locally, so without +// this row a reader would see the reaction only if the transport later +// reported it back. The row is the one such a report addresses +// (SelfReactorKey), and it is applied with ApplyReaction's ordering rule, so a +// mutation of the same reactor that already carries a later time wins. +// +// The target's account and conversation are read from the message as it is +// stored now, not from the outbox row, because a message can be rebound to +// another conversation after its reaction was queued. applied reports whether +// the read model changed. +func (r *OutboxRepository) ConfirmReaction( + ctx context.Context, + confirmation ReactionConfirmation, +) (applied bool, err error) { + occurredAtMS := confirmation.OccurredAt.UnixMilli() + if occurredAtMS <= 0 { + return false, fmt.Errorf( + "confirm reaction outbox item %q: occurrence time is not positive", + confirmation.OutboxID, + ) + } + nowMS, err := r.nowMS("confirm reaction outbox item") + if err != nil { + return false, err + } + tx, err := r.store.db.BeginTx(ctx, nil) + if err != nil { + return false, fmt.Errorf( + "confirm reaction outbox item %q: begin transaction: %w", + confirmation.OutboxID, + err, + ) + } + defer tx.Rollback() + + result, err := tx.ExecContext(ctx, ` + UPDATE outbox + SET state = 'confirmed', + result_remote_id = ?, + error_class = NULL, + error_code = NULL, + error_detail = NULL, + next_attempt_at_ms = NULL, + lease_owner = NULL, + lease_token = NULL, + lease_expires_at_ms = NULL, + transport_called_at_ms = NULL, + updated_at_ms = ? + WHERE outbox_id = ? + AND kind = 'reaction' + AND state = 'dispatching' + AND lease_token = ? + AND transport_called_at_ms IS NOT NULL + `, + nullableOutboxText(strings.TrimSpace(confirmation.ResultRemoteID)), + nowMS, + confirmation.OutboxID, + confirmation.LeaseToken, + ) + if err != nil { + return false, fmt.Errorf("confirm reaction outbox item %q: update: %w", confirmation.OutboxID, err) + } + if err := r.requireLeaseMutationWithQueryer( + ctx, tx, "confirm reaction", confirmation.OutboxID, result, + ); err != nil { + return false, err + } + + reaction := ReactionApply{ + ReactorKey: SelfReactorKey, + ReactorIsSelf: true, + ReactorLabel: SelfReactorLabel, + OccurredAtMS: occurredAtMS, + SourceSeqMS: occurredAtMS, + } + var action string + if err := tx.QueryRowContext(ctx, ` + SELECT m.message_id, m.account_id, m.conversation_id, r.emoji, r.action + FROM outbox_reactions AS r + JOIN messages AS m ON m.message_id = r.target_message_id + WHERE r.outbox_id = ? + `, confirmation.OutboxID).Scan( + &reaction.MessageID, + &reaction.AccountID, + &reaction.ConversationID, + &reaction.Emoji, + &action, + ); err != nil { + return false, fmt.Errorf( + "confirm reaction outbox item %q: read reaction target: %w", + confirmation.OutboxID, + err, + ) + } + reaction.Action = bridge.ReactionAction(action) + applied, err = applyReaction(ctx, tx, reaction, nowMS) + if err != nil { + return false, fmt.Errorf("confirm reaction outbox item %q: %w", confirmation.OutboxID, err) + } + if err := tx.Commit(); err != nil { + return false, fmt.Errorf("confirm reaction outbox item %q: commit: %w", confirmation.OutboxID, err) + } + return applied, nil +} + // Confirm atomically records a known remote result for the active called // transport lease. func (r *OutboxRepository) Confirm( diff --git a/internal/storage/sqlite/outbox_reaction_confirm_test.go b/internal/storage/sqlite/outbox_reaction_confirm_test.go new file mode 100644 index 00000000..427622e0 --- /dev/null +++ b/internal/storage/sqlite/outbox_reaction_confirm_test.go @@ -0,0 +1,420 @@ +package sqlite + +import ( + "context" + "errors" + "fmt" + "math/rand" + "reflect" + "testing" + "testing/quick" + "time" + + "github.com/maxghenis/openmessage/internal/bridge" +) + +const confirmReactionTimeMS int64 = 1_780_000_000_000 + +// confirmReactionFixture is an outbox with one reaction target and the lease +// steps a dispatcher takes before it may confirm. +type confirmReactionFixture struct { + t *testing.T + clock *outboxTestClock + store *Store + outbox *OutboxRepository + reactions *ReactionRepository + enqueued int + targetID string +} + +// retarget points the fixture at a fresh message, so one store can serve many +// independent histories. +func (f *confirmReactionFixture) retarget(messageID string) { + f.t.Helper() + seedOutboxTestMessage(f.t, f.store, messageID, "account-a", "conversation-a") + f.targetID = messageID +} + +func newConfirmReactionFixture(t *testing.T) *confirmReactionFixture { + t.Helper() + clock := newOutboxTestClock(confirmReactionTimeMS) + store, outbox := openOutboxTestRepository(t, clock.Now) + seedMessageConversation(t, store, "conversation-a", "account-a") + seedOutboxTestMessage(t, store, "target-a", "account-a", "conversation-a") + seedMessageIdentity(t, store, "identity-a", "account-a") + seedMessageIdentity(t, store, "identity-other", "account-a") + reactions, err := NewReactionRepository(store, clock.Now) + if err != nil { + t.Fatalf("NewReactionRepository(): %v", err) + } + return &confirmReactionFixture{ + t: t, clock: clock, store: store, outbox: outbox, reactions: reactions, + targetID: "target-a", + } +} + +// called enqueues a reaction and takes it to the point where the transport +// has been called, which is where the dispatcher holds it when the transport +// answers. +func (f *confirmReactionFixture) called(emoji string, action bridge.ReactionAction) (outboxID, leaseToken string) { + f.t.Helper() + ctx := context.Background() + f.enqueued++ + item := outboxTestReactionItem(fmt.Sprintf("reaction-%d", f.enqueued)) + if _, _, err := f.outbox.EnqueueReaction(ctx, item, OutboxReaction{ + TargetMessageID: f.targetID, Emoji: emoji, Action: string(action), + }); err != nil { + f.t.Fatalf("EnqueueReaction(): %v", err) + } + leased := mustLeaseOne(f.t, f.outbox, LeaseRequest{ + Owner: "dispatcher", Now: f.clock.Now(), Duration: time.Minute, Limit: 1, + }) + if leased.OutboxID != item.OutboxID { + f.t.Fatalf("leased %q, want the reaction just enqueued %q", leased.OutboxID, item.OutboxID) + } + leaseToken = mustLeaseToken(f.t, leased) + if err := f.outbox.MarkTransportCalled(ctx, Attempt{ + OutboxID: item.OutboxID, LeaseToken: leaseToken, + AttemptToken: item.OutboxID + ":" + leaseToken, StartedAt: f.clock.Now(), + }); err != nil { + f.t.Fatalf("MarkTransportCalled(): %v", err) + } + return item.OutboxID, leaseToken +} + +func (f *confirmReactionFixture) state(outboxID string) OutboxState { + f.t.Helper() + item, err := f.outbox.FindByID(context.Background(), outboxID) + if err != nil { + f.t.Fatalf("FindByID(%q): %v", outboxID, err) + } + return item.State +} + +func TestConfirmReactionConfirmsAndRecordsTheOwnReactionTogether(t *testing.T) { + f := newConfirmReactionFixture(t) + outboxID, leaseToken := f.called("πŸ‘", bridge.ReactionAdd) + acceptedAtMS := confirmReactionTimeMS + 250 + f.clock.Set(confirmReactionTimeMS + 900) + + applied, err := f.outbox.ConfirmReaction(context.Background(), ReactionConfirmation{ + OutboxID: outboxID, LeaseToken: leaseToken, OccurredAt: time.UnixMilli(acceptedAtMS), + }) + if err != nil || !applied { + t.Fatalf("ConfirmReaction() = %v, %v; want the read model changed", applied, err) + } + item, err := f.outbox.FindByID(context.Background(), outboxID) + if err != nil { + t.Fatalf("FindByID(): %v", err) + } + if item.State != OutboxConfirmed || item.ResultRemoteID != nil || item.LeaseToken != nil { + t.Fatalf("confirmed reaction row = %+v, want confirmed with no result and no lease", item) + } + got := readStoredReaction(t, f.store, f.targetID, SelfReactorKey) + want := storedReaction{ + MessageID: f.targetID, ReactorKey: SelfReactorKey, ReactorIsSelf: true, ReactorLabel: SelfReactorLabel, + Emoji: "πŸ‘", State: "active", OccurredAtMS: acceptedAtMS, SourceSeqMS: acceptedAtMS, + CreatedAtMS: confirmReactionTimeMS + 900, UpdatedAtMS: confirmReactionTimeMS + 900, + } + if !reflect.DeepEqual(got, want) { + t.Fatalf("own reaction row = %+v, want %+v", got, want) + } + + // A remove through the same path tombstones that one row. + removeID, removeToken := f.called("πŸ‘", bridge.ReactionRemove) + if _, err := f.outbox.ConfirmReaction(context.Background(), ReactionConfirmation{ + OutboxID: removeID, LeaseToken: removeToken, ResultRemoteID: "remote-remove", + OccurredAt: time.UnixMilli(acceptedAtMS + 1), + }); err != nil { + t.Fatalf("ConfirmReaction(remove): %v", err) + } + removed, err := f.outbox.FindByID(context.Background(), removeID) + if err != nil || removed.State != OutboxConfirmed { + t.Fatalf("removal row = %+v, %v; want confirmed", removed, err) + } + assertOutboxText(t, "result remote ID", removed.ResultRemoteID, "remote-remove") + if row := readStoredReaction(t, f.store, f.targetID, SelfReactorKey); row.State != "removed" || row.Emoji != "πŸ‘" { + t.Fatalf("own reaction row after remove = %+v, want the same row removed", row) + } + assertRowCount(t, f.store.db, "reactions", 1) +} + +// A reaction that is not confirmed must leave no trace in the read model, and +// a confirmation that cannot be recorded must not confirm: the two writes are +// one transaction. +func TestConfirmReactionIsAllOrNothing(t *testing.T) { + t.Run("a lost lease writes nothing", func(t *testing.T) { + f := newConfirmReactionFixture(t) + outboxID, _ := f.called("πŸ‘", bridge.ReactionAdd) + _, err := f.outbox.ConfirmReaction(context.Background(), ReactionConfirmation{ + OutboxID: outboxID, LeaseToken: "someone-elses-lease", OccurredAt: f.clock.Now(), + }) + if !errors.Is(err, ErrLeaseLost) { + t.Fatalf("ConfirmReaction(wrong lease) = %v, want ErrLeaseLost", err) + } + if state := f.state(outboxID); state != OutboxDispatching { + t.Fatalf("state = %q, want still dispatching", state) + } + assertRowCount(t, f.store.db, "reactions", 0) + }) + + t.Run("a reaction the transport was never called for cannot be confirmed", func(t *testing.T) { + f := newConfirmReactionFixture(t) + ctx := context.Background() + item := outboxTestReactionItem("uncalled") + if _, _, err := f.outbox.EnqueueReaction(ctx, item, OutboxReaction{ + TargetMessageID: f.targetID, Emoji: "πŸ‘", Action: "add", + }); err != nil { + t.Fatalf("EnqueueReaction(): %v", err) + } + leased := mustLeaseOne(t, f.outbox, LeaseRequest{ + Owner: "dispatcher", Now: f.clock.Now(), Duration: time.Minute, Limit: 1, + }) + if _, err := f.outbox.ConfirmReaction(ctx, ReactionConfirmation{ + OutboxID: item.OutboxID, LeaseToken: mustLeaseToken(t, leased), OccurredAt: f.clock.Now(), + }); !errors.Is(err, ErrLeaseLost) { + t.Fatalf("ConfirmReaction(uncalled) = %v, want ErrLeaseLost", err) + } + assertRowCount(t, f.store.db, "reactions", 0) + }) + + t.Run("a failed read-model write rolls the confirmation back", func(t *testing.T) { + f := newConfirmReactionFixture(t) + outboxID, leaseToken := f.called("πŸ‘", bridge.ReactionAdd) + // The stored action is corrupted after enqueue, so the outbox update + // succeeds and the read-model write is the step that fails. + mustExec(t, f.store.db, `PRAGMA ignore_check_constraints = ON`) + mustExec(t, f.store.db, `UPDATE outbox_reactions SET action = 'explode' WHERE outbox_id = ?`, outboxID) + mustExec(t, f.store.db, `PRAGMA ignore_check_constraints = OFF`) + + if _, err := f.outbox.ConfirmReaction(context.Background(), ReactionConfirmation{ + OutboxID: outboxID, LeaseToken: leaseToken, OccurredAt: f.clock.Now(), + }); err == nil { + t.Fatal("ConfirmReaction(corrupt action) succeeded") + } + if state := f.state(outboxID); state != OutboxDispatching { + t.Fatalf("state = %q, want still dispatching after the rollback", state) + } + assertRowCount(t, f.store.db, "reactions", 0) + }) + + t.Run("only a reaction row can be confirmed this way", func(t *testing.T) { + f := newConfirmReactionFixture(t) + ctx := context.Background() + item := outboxTestItem("text") + if _, _, err := f.outbox.EnqueueOutgoingMessage(ctx, item, outboxTestOutgoingMessage(item, "text")); err != nil { + t.Fatalf("EnqueueOutgoingMessage(): %v", err) + } + leased := mustLeaseOne(t, f.outbox, LeaseRequest{ + Owner: "dispatcher", Now: f.clock.Now(), Duration: time.Minute, Limit: 1, + }) + leaseToken := mustLeaseToken(t, leased) + if err := f.outbox.MarkTransportCalled(ctx, Attempt{ + OutboxID: item.OutboxID, LeaseToken: leaseToken, AttemptToken: "attempt", StartedAt: f.clock.Now(), + }); err != nil { + t.Fatalf("MarkTransportCalled(): %v", err) + } + if _, err := f.outbox.ConfirmReaction(ctx, ReactionConfirmation{ + OutboxID: item.OutboxID, LeaseToken: leaseToken, OccurredAt: f.clock.Now(), + }); !errors.Is(err, ErrLeaseLost) { + t.Fatalf("ConfirmReaction(text row) = %v, want ErrLeaseLost", err) + } + if state := f.state(item.OutboxID); state != OutboxDispatching { + t.Fatalf("state = %q, want the text row untouched", state) + } + }) + + t.Run("an occurrence time is required", func(t *testing.T) { + f := newConfirmReactionFixture(t) + outboxID, leaseToken := f.called("πŸ‘", bridge.ReactionAdd) + if _, err := f.outbox.ConfirmReaction(context.Background(), ReactionConfirmation{ + OutboxID: outboxID, LeaseToken: leaseToken, + }); err == nil { + t.Fatal("ConfirmReaction(zero time) succeeded") + } + if state := f.state(outboxID); state != OutboxDispatching { + t.Fatalf("state = %q, want still dispatching", state) + } + }) +} + +// A message can move to another conversation after its reaction was queued. +// The read-model row must follow the message, or its conversation foreign key +// would name a thread the message is no longer in. +func TestConfirmReactionRecordsUnderTheTargetsCurrentConversation(t *testing.T) { + f := newConfirmReactionFixture(t) + outboxID, leaseToken := f.called("πŸ‘", bridge.ReactionAdd) + seedMessageConversation(t, f.store, "conversation-b", "account-a") + mustExec(t, f.store.db, `UPDATE messages SET conversation_id = 'conversation-b' WHERE message_id = ?`, f.targetID) + + if _, err := f.outbox.ConfirmReaction(context.Background(), ReactionConfirmation{ + OutboxID: outboxID, LeaseToken: leaseToken, OccurredAt: f.clock.Now(), + }); err != nil { + t.Fatalf("ConfirmReaction(): %v", err) + } + var conversationID string + if err := f.store.db.QueryRow( + `SELECT conversation_id FROM reactions WHERE message_id = ? AND reactor_key = ?`, + f.targetID, SelfReactorKey, + ).Scan(&conversationID); err != nil || conversationID != "conversation-b" { + t.Fatalf("own reaction conversation = %q, %v; want the message's current one", conversationID, err) + } +} + +// reactionScriptStep is one event in the life of a message's reactions. +type reactionScriptStep struct { + // Kind: 0 this account's reaction, delivered through the outbox; 1 this + // account's reaction, not delivered; 2 a transport report of this + // account's reaction (made on another device); 3 another person's. + Kind int + Emoji string + Action bridge.ReactionAction + // AtMS is the reaction's own time and SeqMS the frame time a transport + // report carries. Both come from a few values, so ties are common. + AtMS int64 + SeqMS int64 + // Outcome picks how an undelivered reaction ended. + Outcome int +} + +type reactionScript []reactionScriptStep + +func (reactionScript) Generate(r *rand.Rand, _ int) reflect.Value { + emoji := []string{"πŸ‘", "❀️", "πŸ˜‚"} + actions := []bridge.ReactionAction{bridge.ReactionAdd, bridge.ReactionRemove, bridge.ReactionSwitch} + script := make(reactionScript, 1+r.Intn(14)) + for i := range script { + script[i] = reactionScriptStep{ + Kind: r.Intn(4), + Emoji: emoji[r.Intn(len(emoji))], + Action: actions[r.Intn(len(actions))], + AtMS: confirmReactionTimeMS + int64(r.Intn(6)), + SeqMS: confirmReactionTimeMS + int64(r.Intn(6)), + Outcome: r.Intn(3), + } + } + return reflect.ValueOf(script) +} + +// storedReactionsFor returns every reaction row of a message, removed ones +// included, without the storage timestamps that record when a row was written. +func storedReactionsFor(t *testing.T, store *Store, messageID string) []storedReaction { + t.Helper() + rows, err := store.db.Query(`SELECT reactor_key FROM reactions WHERE message_id = ? ORDER BY reactor_key`, messageID) + if err != nil { + t.Fatalf("list reactions: %v", err) + } + var keys []string + for rows.Next() { + var key string + if err := rows.Scan(&key); err != nil { + t.Fatalf("scan reaction key: %v", err) + } + keys = append(keys, key) + } + if err := rows.Close(); err != nil { + t.Fatalf("close reactions: %v", err) + } + stored := make([]storedReaction, 0, len(keys)) + for _, key := range keys { + stored = append(stored, readStoredReaction(t, store, messageID, key)) + } + return stored +} + +// TestConfirmedOwnReactionIsStoredAsItsEchoWouldBe is the differential check +// between the two writers of an own reaction. For any history of a message's +// reactions, a store where this account's delivered reactions were confirmed +// by the outbox holds exactly the rows of a store where the transport +// reported each of them instead, and a reaction that was not delivered +// (rejected, left uncertain, or put back for retry) leaves no row at all. +func TestConfirmedOwnReactionIsStoredAsItsEchoWouldBe(t *testing.T) { + // Two stores serve every case; each case gets a message of its own in + // both, because opening a store costs far more than a case does. + viaOutbox := newConfirmReactionFixture(t) + viaEcho := newConfirmReactionFixture(t) + ctx := context.Background() + cases := 0 + property := func(script reactionScript) bool { + cases++ + messageID := fmt.Sprintf("target-case-%d", cases) + viaOutbox.retarget(messageID) + viaEcho.retarget(messageID) + for index, step := range script { + nowMS := confirmReactionTimeMS + 100 + int64(index) + viaOutbox.clock.Set(nowMS) + viaEcho.clock.Set(nowMS) + self := ReactionApply{ + AccountID: "account-a", ConversationID: "conversation-a", MessageID: messageID, + ReactorKey: SelfReactorKey, ReactorIsSelf: true, ReactorLabel: SelfReactorLabel, + Emoji: step.Emoji, Action: step.Action, OccurredAtMS: step.AtMS, SourceSeqMS: step.AtMS, + } + switch step.Kind { + case 0: + outboxID, leaseToken := viaOutbox.called(step.Emoji, step.Action) + if _, err := viaOutbox.outbox.ConfirmReaction(ctx, ReactionConfirmation{ + OutboxID: outboxID, LeaseToken: leaseToken, OccurredAt: time.UnixMilli(step.AtMS), + }); err != nil { + t.Errorf("step %d ConfirmReaction(): %v", index, err) + return false + } + if state := viaOutbox.state(outboxID); state != OutboxConfirmed { + t.Errorf("step %d state = %q, want confirmed", index, state) + return false + } + if _, err := viaEcho.reactions.ApplyReaction(ctx, self); err != nil { + t.Errorf("step %d ApplyReaction(echo of own): %v", index, err) + return false + } + case 1: + outboxID, leaseToken := viaOutbox.called(step.Emoji, step.Action) + var err error + want := OutboxRejected + switch step.Outcome { + case 0: + err = viaOutbox.outbox.Reject(ctx, outboxID, leaseToken, "unsupported", "send_reaction", "scripted") + case 1: + want = OutboxUncertain + err = viaOutbox.outbox.MarkUncertain(ctx, outboxID, leaseToken, "transient", "send_reaction", "scripted") + default: + want = OutboxNotDispatched + err = viaOutbox.outbox.MarkCalledNotDispatched( + ctx, outboxID, leaseToken, "transient", "send_reaction", "scripted", + time.UnixMilli(nowMS).Add(24*time.Hour), + ) + } + if err != nil { + t.Errorf("step %d undelivered outcome %d: %v", index, step.Outcome, err) + return false + } + if state := viaOutbox.state(outboxID); state != want { + t.Errorf("step %d state = %q, want %q", index, state, want) + return false + } + default: + report := self + report.SourceSeqMS = step.SeqMS + if step.Kind == 3 { + report.ReactorKey, report.ReactorIsSelf, report.ReactorLabel = "identity-other", false, "" + report.ReactorIdentityID = pointer("identity-other") + } + for _, fixture := range []*confirmReactionFixture{viaOutbox, viaEcho} { + if _, err := fixture.reactions.ApplyReaction(ctx, report); err != nil { + t.Errorf("step %d ApplyReaction(report): %v", index, err) + return false + } + } + } + } + got := storedReactionsFor(t, viaOutbox.store, messageID) + want := storedReactionsFor(t, viaEcho.store, messageID) + if !reflect.DeepEqual(got, want) { + t.Errorf("script %+v\nvia the outbox: %+v\nvia echoes: %+v", script, got, want) + return false + } + return true + } + if err := quick.Check(property, &quick.Config{MaxCount: 150, Rand: rand.New(rand.NewSource(20261009))}); err != nil { + t.Fatal(err) + } +} diff --git a/internal/storage/sqlite/reactions.go b/internal/storage/sqlite/reactions.go index 300034ec..29edfbb1 100644 --- a/internal/storage/sqlite/reactions.go +++ b/internal/storage/sqlite/reactions.go @@ -73,6 +73,21 @@ func NewReactionRepository( return &ReactionRepository{store: store, now: now}, nil } +// SelfReactorKey and SelfReactorLabel identify this account's own reaction on +// a message. Every writer of an own reaction (the ingest worker for a +// transport echo, the migration, and the outbox when it confirms a reaction +// this account sent) uses them, so they all address one row per message. +const ( + SelfReactorKey = "self" + SelfReactorLabel = "me" +) + +// reactionExecer is the write surface ApplyReaction needs; *sql.DB and +// *sql.Tx both provide it. +type reactionExecer interface { + ExecContext(ctx context.Context, query string, args ...any) (sql.Result, error) +} + // ApplyReaction applies one WhatsApp, Signal, or tapback delta. A newer // removal preserves the previous emoji for audit; a removal received before // any add inserts an empty-emoji tombstone that fences stale adds. Delta @@ -80,6 +95,21 @@ func NewReactionRepository( func (r *ReactionRepository) ApplyReaction( ctx context.Context, reaction ReactionApply, +) (bool, error) { + nowMS, err := r.nowMS("apply reaction") + if err != nil { + return false, err + } + return applyReaction(ctx, r.store.db, reaction, nowMS) +} + +// applyReaction is ApplyReaction on any execer, so the outbox can record a +// confirmed own reaction inside its confirming transaction. +func applyReaction( + ctx context.Context, + execer reactionExecer, + reaction ReactionApply, + nowMS int64, ) (bool, error) { state := "active" emoji := reaction.Emoji @@ -99,11 +129,7 @@ func (r *ReactionRepository) ApplyReaction( ) } - nowMS, err := r.nowMS("apply reaction") - if err != nil { - return false, err - } - result, err := r.store.db.ExecContext(ctx, ` + result, err := execer.ExecContext(ctx, ` INSERT INTO reactions ( message_id, reactor_key, diff --git a/internal/tools/daemon.go b/internal/tools/daemon.go index 39c0ff46..3df7bb75 100644 --- a/internal/tools/daemon.go +++ b/internal/tools/daemon.go @@ -405,12 +405,37 @@ func daemonReactToMessageHandler(options Options) server.ToolHandlerFunc { if action == "" { action = "add" } - if err := daemon.React(ctx, conversationID, messageID, emoji, action); err != nil { + // A v2-primary app takes reactions on its durable outbox, where the + // intent carries an idempotency key and its delivery can be followed. + // Any other answer to the probe keeps the reaction on /api/react, + // which every app mode serves. + status, reachable, probeErr := daemon.Status(ctx) + if probeErr != nil && !reachable { + return daemonDownResult(probeErr), nil + } + if probeErr == nil && status.ReactionsViaOutbox() { + return daemonSubmitReactionAndWait(ctx, daemon, args, conversationID, messageID, emoji, action), nil + } + result, err := daemon.React(ctx, conversationID, messageID, emoji, action) + if err != nil { if responseErr, ok := localapi.AsResponseError(err); ok { return errorResult(fmt.Sprintf("the app could not send the reaction: HTTP %d: %s", responseErr.StatusCode, responseErr.Body)), nil } return daemonDownResult(err), nil } + if !result.Success { + // The legacy Google route answers 200 with success false when the + // phone refuses the reaction. + return errorResult("the app reported that the reaction was not applied"), nil + } + if result.Queued { + // A v2-primary app answered on the compatibility route: the + // reaction is on its outbox and not delivered yet. + return v2ReactionResult( + messaging.Delivery{OutboxID: result.OutboxID, State: messaging.OutboxState(result.State)}, + false, "", nil, + ), nil + } return structuredResult(map[string]any{ "ok": true, "message_id": messageID, diff --git a/internal/tools/react_to_message.go b/internal/tools/react_to_message.go index 76a34d66..8cfc2516 100644 --- a/internal/tools/react_to_message.go +++ b/internal/tools/react_to_message.go @@ -35,19 +35,33 @@ var ( } ) -func reactToMessageTool() mcp.Tool { - return mcp.NewTool("react_to_message", - mcp.WithDescription("Add, remove, or switch a reaction on an existing message across supported platforms"), +func reactToMessageTool(v2Enabled ...bool) mcp.Tool { + description := "Add, remove, or switch a reaction on an existing message across supported platforms" + options := []mcp.ToolOption{ + mcp.WithDescription(description), mcp.WithString("conversation_id", mcp.Required(), mcp.Description("Conversation ID containing the target message")), mcp.WithString("message_id", mcp.Required(), mcp.Description("Target message ID")), mcp.WithString("emoji", mcp.Required(), mcp.Description("Emoji reaction to apply")), mcp.WithString("action", mcp.Description("Optional action: add, remove, or switch. Defaults to add.")), + } + if v2Requested(v2Enabled) { + options[0] = mcp.WithDescription(description + v2ReactionDescription) + options = append(options, mcp.WithString("idempotency_key", mcp.Description(v2ReactionIdempotencyDescription))) + } + options = append(options, mcp.WithDestructiveHintAnnotation(false), mcp.WithIdempotentHintAnnotation(false), ) + return mcp.NewTool("react_to_message", options...) } -func reactToMessageHandler(a *app.App) server.ToolHandlerFunc { +// reactToMessageHandler serves react_to_message in a process that owns its +// transports. On a v2-primary install the read tools hand out v2 IDs, which +// only the v2 outbox can place, so the reaction is queued there; otherwise it +// goes straight to the platform's legacy sender. +func reactToMessageHandler(a *app.App, configured ...Options) server.ToolHandlerFunc { + options := resolvedOptions(a, configured) + v2 := activeV2([]*V2Dependencies{options.V2}) return func(ctx context.Context, req mcp.CallToolRequest) (*mcp.CallToolResult, error) { args := req.GetArguments() conversationID := strArg(args, "conversation_id") @@ -64,6 +78,12 @@ func reactToMessageHandler(a *app.App) server.ToolHandlerFunc { if emoji == "" { return errorResult("emoji is required"), nil } + if options.V2Primary { + if v2 == nil { + return errorResult("reactions on a v2-primary install go through the v2 outbox, which this process has not configured"), nil + } + return submitV2Reaction(ctx, v2, args, conversationID, messageID, emoji, action), nil + } conv, err := a.Store.GetConversation(conversationID) if err != nil { diff --git a/internal/tools/tools.go b/internal/tools/tools.go index b2ca695a..21449637 100644 --- a/internal/tools/tools.go +++ b/internal/tools/tools.go @@ -75,9 +75,12 @@ func RegisterWithOptions(s *server.MCPServer, a *app.App, options Options) { s.AddTool(sendToConversationTool(true), sendToConversationHandler(a, configuredV2)) s.AddTool(sendMediaToConversationTool(true), sendMediaToConversationHandler(a, configuredV2)) } - if options.Daemon != nil { - s.AddTool(reactToMessageTool(), daemonReactToMessageHandler(options)) - } else { + switch { + case options.Daemon != nil: + s.AddTool(reactToMessageTool(true), daemonReactToMessageHandler(options)) + case v2Primary: + s.AddTool(reactToMessageTool(true), reactToMessageHandler(a, options)) + default: s.AddTool(reactToMessageTool(), reactToMessageHandler(a)) } s.AddTool(setMessageTranscriptTool(), setMessageTranscriptHandler(a)) diff --git a/internal/tools/v2_react.go b/internal/tools/v2_react.go new file mode 100644 index 00000000..33d33712 --- /dev/null +++ b/internal/tools/v2_react.go @@ -0,0 +1,188 @@ +package tools + +import ( + "context" + "fmt" + + "github.com/mark3labs/mcp-go/mcp" + + "github.com/maxghenis/openmessage/internal/localapi" + "github.com/maxghenis/openmessage/internal/messaging" + "github.com/maxghenis/openmessage/internal/v2wire" +) + +const v2ReactionDescription = " On a v2 install the reaction goes through the app's durable outbox, and this waits for it to settle. Use the conversation and message IDs the read tools return. The result carries outbox_id, state and idempotency_key; ok is true once the platform accepted the reaction. If settled is false the reaction is still queued and the app keeps retrying it in the background: do not react again (it is listed in the app's outbox, where it can be canceled). An uncertain result means the platform may have applied it; read the message's reactions before reacting again." + +const v2ReactionIdempotencyDescription = "Optional retry key for the exact same reaction, used on a v2 install. Every v2 result echoes the key in use; reuse it only when repeating a reaction whose response was lost. Omit it to mint a new intent." + +// submitV2Reaction queues a reaction on this process's own v2 outbox and waits +// for it to settle. It serves the in-process MCP surface of a v2-primary +// daemon. +func submitV2Reaction( + ctx context.Context, + v2 *V2Dependencies, + args map[string]any, + conversationID, messageID, emoji, action string, +) *mcp.CallToolResult { + if v2.Service == nil { + return errorResult("v2 send service is unavailable") + } + key, err := v2IdempotencyKey(args) + if err != nil { + return errorResult(err.Error()) + } + submission, err := v2wire.SubmitReactionV2(ctx, v2.nativeDeps(), v2wire.ReactionInput{ + ConversationID: conversationID, + MessageID: messageID, + Emoji: emoji, + Action: action, + IdempotencyKey: key, + }) + if err != nil { + return errorResult(fmt.Sprintf("failed to submit reaction: %v", err)) + } + // Bounded like the client-mode wait: a reaction the dispatcher keeps + // putting back (its platform has no lease to give) never settles, and the + // caller still needs an answer. + waitCtx, cancel := context.WithTimeout(ctx, daemonSettleTimeout) + defer cancel() + delivery, err := v2.Service.Wait(waitCtx, submission.OutboxID) + if delivery.OutboxID == "" { + delivery = messaging.Delivery{OutboxID: submission.OutboxID, State: submission.State} + } + return v2ReactionResult(delivery, submission.Deduplicated, key, err) +} + +// daemonSubmitReactionAndWait is submitV2Reaction for transportless client +// mode: the running app owns the outbox, and this process follows the intent +// over the app's local API. +func daemonSubmitReactionAndWait( + ctx context.Context, + daemon *localapi.Client, + args map[string]any, + conversationID, messageID, emoji, action string, +) *mcp.CallToolResult { + key, err := v2IdempotencyKey(args) + if err != nil { + return errorResult(err.Error()) + } + submission, err := daemon.SubmitReaction(ctx, localapi.ReactionSubmission{ + ConversationID: conversationID, + MessageID: messageID, + Emoji: emoji, + Action: action, + IdempotencyKey: key, + }) + if err != nil { + if localapi.IsDeterministicRejection(err) { + return errorResult(fmt.Sprintf("reaction rejected by the app: %v", err)) + } + if responseErr, ok := localapi.AsResponseError(err); ok { + return errorResult(fmt.Sprintf("the app could not queue the reaction: HTTP %d: %s", responseErr.StatusCode, responseErr.Body)) + } + // The request failed mid-flight, so the app may or may not hold the + // intent. Not an IsError result: an error invites a second reaction. + return structuredResult(map[string]any{ + "ok": false, + "settled": false, + "ambiguous": true, + "idempotency_key": key, + "error": err.Error(), + }, fmt.Sprintf( + "The reaction's outcome is unknown (%v). Do NOT react again with a new key. To replay-check this exact reaction, repeat it with the same idempotency_key: %s. If the app is not running, start it first.", + err, key, + )) + } + delivery, settled, err := daemon.WaitDelivery(ctx, submission.OutboxID, daemonSettleTimeout) + if err != nil { + // The intent is durably queued on the daemon; only our view failed. + return v2ReactionResult( + messaging.Delivery{OutboxID: submission.OutboxID, State: messaging.OutboxState(submission.State)}, + submission.Deduplicated, key, err, + ) + } + converted := deliveryFromLocalAPI(delivery) + if !settled { + return v2ReactionResult(converted, submission.Deduplicated, key, context.DeadlineExceeded) + } + return v2ReactionResult(converted, submission.Deduplicated, key, nil) +} + +// v2ReactionResult reports a queued reaction's outcome. The intent is already +// stored when this runs, so no path here returns an IsError result: a tool +// error invites the calling agent to react again, and the outbox finishes the +// first reaction regardless. waitErr is why the wait ended before the +// reaction settled, if it did. +func v2ReactionResult( + delivery messaging.Delivery, + deduplicated bool, + idempotencyKey string, + waitErr error, +) *mcp.CallToolResult { + state := delivery.State + if state == "" { + state = messaging.OutboxQueued + } + settled := v2ReactionSettled(state) + payload := map[string]any{ + "ok": v2DeliveryOK(state), + "settled": settled, + "outbox_id": delivery.OutboxID, + "state": state, + "deduplicated": deduplicated, + } + // The app's compatibility route mints its own key and does not return it. + if idempotencyKey != "" { + payload["idempotency_key"] = idempotencyKey + } + if !settled { + payload["auto_retry"] = true + if waitErr != nil { + payload["wait_error"] = waitErr.Error() + } + } + if delivery.ErrorClass != "" { + payload["error_class"] = delivery.ErrorClass + } + if delivery.ErrorCode != "" { + payload["error_code"] = delivery.ErrorCode + } + if delivery.Warning != "" { + payload["warning"] = delivery.Warning + } + return structuredResult(payload, v2ReactionText(delivery.OutboxID, state, delivery.ErrorClass, idempotencyKey)) +} + +// v2ReactionSettled reports whether the app is done with the reaction. A +// queued, in-flight or automatically retrying reaction is not. +func v2ReactionSettled(state messaging.OutboxState) bool { + switch state { + case messaging.OutboxConfirmed, messaging.OutboxStoreFailed, messaging.OutboxUncertain, + messaging.OutboxRejected, messaging.OutboxCanceled: + return true + default: + return false + } +} + +func v2ReactionText(outboxID string, state messaging.OutboxState, errorClass, idempotencyKey string) string { + switch state { + case messaging.OutboxConfirmed, messaging.OutboxStoreFailed: + return fmt.Sprintf("Reaction delivered (outbox %s).", outboxID) + case messaging.OutboxUncertain: + return fmt.Sprintf("The reaction's outcome is unknown (outbox %s): the platform may have applied it. Read the message's reactions before reacting again.", outboxID) + case messaging.OutboxRejected: + return fmt.Sprintf("The reaction was rejected (outbox %s, error class %s). The app will not retry it.", outboxID, firstNonEmpty(errorClass, "unknown")) + case messaging.OutboxCanceled: + return fmt.Sprintf("The reaction was canceled before it was sent (outbox %s).", outboxID) + default: + text := fmt.Sprintf( + "The reaction is durably queued (outbox %s, state %s) and the app keeps retrying it in the background. Do NOT react again; it is listed in the app's outbox, where it can be canceled.", + outboxID, state, + ) + if idempotencyKey != "" { + text += fmt.Sprintf(" To repeat this exact reaction deliberately, reuse idempotency_key %s.", idempotencyKey) + } + return text + } +} diff --git a/internal/tools/v2_react_test.go b/internal/tools/v2_react_test.go new file mode 100644 index 00000000..27438d88 --- /dev/null +++ b/internal/tools/v2_react_test.go @@ -0,0 +1,389 @@ +package tools + +import ( + "context" + "encoding/json" + "errors" + "net/http" + "strings" + "sync" + "testing" + "time" + + "github.com/mark3labs/mcp-go/mcp" + "github.com/mark3labs/mcp-go/server" + + "github.com/maxghenis/openmessage/internal/app" + "github.com/maxghenis/openmessage/internal/bridge" + "github.com/maxghenis/openmessage/internal/storage/sqlite" +) + +const ( + v2ReactConversationID = "v2-react-conversation" + v2ReactMessageID = "v2-react-target" +) + +// v2ToolReactionSender records each reaction and answers with the next +// scripted step, accepting once the script runs out. +type v2ToolReactionSender struct { + mu sync.Mutex + steps []v2ToolSendStep + requests []bridge.ReactionRequest +} + +func (s *v2ToolReactionSender) SendReaction(_ context.Context, request bridge.ReactionRequest) (bridge.SendResult, error) { + s.mu.Lock() + defer s.mu.Unlock() + s.requests = append(s.requests, request) + if len(s.steps) == 0 { + return bridge.SendResult{AcceptedAt: time.Now()}, nil + } + step := s.steps[0] + s.steps = s.steps[1:] + return step.result, step.err +} + +func (s *v2ToolReactionSender) snapshotRequests() []bridge.ReactionRequest { + s.mu.Lock() + defer s.mu.Unlock() + return append([]bridge.ReactionRequest(nil), s.requests...) +} + +// newV2ReactToolHarness is a v2-primary MCP surface over a v2 store holding +// one conversation and one incoming message, whose transport can react. +func newV2ReactToolHarness(t *testing.T, steps ...v2ToolSendStep) (*v2ToolHarness, *v2ToolReactionSender, *server.MCPServer) { + t.Helper() + harness := newV2ToolHarness(t) + sender := &v2ToolReactionSender{steps: steps} + harness.deps.Registry.(*v2ToolRegistry).reaction = sender + + store := harness.deps.V2Store + nowMS := time.Now().UnixMilli() + if err := store.UpsertAccount(sqlite.Account{ + AccountID: "google-primary", BridgeKey: "google_messages", DisplayName: "Google", + Mode: sqlite.AccountModeLive, Enabled: true, ConfigJSON: "{}", CreatedAtMS: nowMS, UpdatedAtMS: nowMS, + }); err != nil { + t.Fatalf("UpsertAccount(): %v", err) + } + if err := store.UpsertConversation(sqlite.Conversation{ + ConversationID: v2ReactConversationID, AccountID: "google-primary", RemoteConversationID: "remote-react-thread", + Kind: sqlite.ConversationKindDirect, Title: "React thread", NotificationMode: sqlite.NotificationModeAll, + MetadataJSON: "{}", CreatedAtMS: nowMS, UpdatedAtMS: nowMS, + }); err != nil { + t.Fatalf("UpsertConversation(): %v", err) + } + messages, err := sqlite.NewMessageRepository(store, time.Now) + if err != nil { + t.Fatalf("NewMessageRepository(): %v", err) + } + if err := messages.ImportMessage(context.Background(), sqlite.MessageProjection{Message: sqlite.Message{ + MessageID: v2ReactMessageID, ConversationID: v2ReactConversationID, AccountID: "google-primary", + RemoteMessageID: "remote-react-target", Direction: sqlite.MessageDirectionIncoming, + Body: "react to this", State: sqlite.MessageStateActive, OccurredAtMS: nowMS - 60_000, + }}); err != nil { + t.Fatalf("ImportMessage(): %v", err) + } + + // The legacy senders read the legacy store, which holds none of these + // IDs; a v2-primary reaction must never reach them. + originalWhatsApp, originalSignal := sendWhatsAppReactionMessage, sendSignalReactionMessage + failLegacy := func(*app.App, string, string, string, string) error { + t.Error("a v2-primary reaction reached a legacy sender") + return errors.New("legacy sender") + } + sendWhatsAppReactionMessage, sendSignalReactionMessage = failLegacy, failLegacy + t.Cleanup(func() { + sendWhatsAppReactionMessage, sendSignalReactionMessage = originalWhatsApp, originalSignal + }) + + mcpServer := server.NewMCPServer("v2-react-test", "test") + deps := harness.deps + RegisterWithOptions(mcpServer, harness.app, Options{Reads: harness.app.Store, V2Primary: true, V2: &deps}) + return harness, sender, mcpServer +} + +func callReactTool(t *testing.T, mcpServer *server.MCPServer, arguments map[string]any) *mcp.CallToolResult { + t.Helper() + tool := mcpServer.GetTool("react_to_message") + if tool == nil { + t.Fatal("react_to_message is not registered") + } + result, err := tool.Handler(context.Background(), v2ToolCall(arguments)) + if err != nil { + t.Fatalf("react_to_message: %v", err) + } + return result +} + +func TestV2ReactToMessageGoesThroughTheOutbox(t *testing.T) { + _, sender, mcpServer := newV2ReactToolHarness(t) + if _, ok := mcpServer.GetTool("react_to_message").Tool.InputSchema.Properties["idempotency_key"]; !ok { + t.Fatal("the v2 react_to_message descriptor does not offer idempotency_key") + } + arguments := map[string]any{ + "conversation_id": v2ReactConversationID, "message_id": v2ReactMessageID, + "emoji": "πŸ‘", "idempotency_key": "react-tool-key", + } + result := callReactTool(t, mcpServer, arguments) + payload := v2ToolPayload(t, result) + if result.IsError { + t.Fatalf("react_to_message = error %v", payload) + } + assertV2ToolBool(t, payload, "ok", true) + assertV2ToolBool(t, payload, "settled", true) + assertV2ToolString(t, payload, "state", "confirmed") + assertV2ToolString(t, payload, "idempotency_key", "react-tool-key") + if requests := sender.snapshotRequests(); len(requests) != 1 || + requests[0].Target.RemoteID != "remote-react-target" || requests[0].Conversation.RemoteID != "remote-react-thread" || + requests[0].Emoji != "πŸ‘" || requests[0].Action != bridge.ReactionAdd { + t.Fatalf("reaction requests = %+v", requests) + } + + // The same key replays the first intent and sends nothing new. + replay := v2ToolPayload(t, callReactTool(t, mcpServer, arguments)) + assertV2ToolBool(t, replay, "deduplicated", true) + assertV2ToolString(t, replay, "outbox_id", v2ToolString(payload, "outbox_id")) + if requests := sender.snapshotRequests(); len(requests) != 1 { + t.Fatalf("a replayed reaction reached the transport again: %+v", requests) + } +} + +func TestV2ReactToMessageReportsAnUndeliveredReactionWithoutAnError(t *testing.T) { + _, _, mcpServer := newV2ReactToolHarness(t, v2ToolSendStep{err: bridge.OpError{ + Class: bridge.FailureTransient, Operation: "send_reaction", Fingerprint: "disconnected", + Dispatch: bridge.DispatchNotCalled, + }}) + result := callReactTool(t, mcpServer, map[string]any{ + "conversation_id": v2ReactConversationID, "message_id": v2ReactMessageID, "emoji": "πŸ‘", + }) + payload := v2ToolPayload(t, result) + // An error result invites the agent to react again; the app retries it. + if result.IsError { + t.Fatalf("an automatically retrying reaction was reported as an error: %v", payload) + } + assertV2ToolBool(t, payload, "ok", false) + assertV2ToolBool(t, payload, "settled", false) + assertV2ToolBool(t, payload, "auto_retry", true) + assertV2ToolString(t, payload, "state", "not_dispatched") + if key := v2ToolString(payload, "idempotency_key"); key == "" { + t.Fatalf("payload = %v, want the minted idempotency key", payload) + } + if text := resultText(t, result); !strings.Contains(text, "Do NOT react again") { + t.Fatalf("text = %q", text) + } +} + +func TestV2ReactToMessageRefusesWhatItCannotPlace(t *testing.T) { + _, sender, mcpServer := newV2ReactToolHarness(t) + for name, arguments := range map[string]map[string]any{ + "a legacy message ID": {"conversation_id": v2ReactConversationID, "message_id": "signal:1700000001000", "emoji": "πŸ‘"}, + "another conversation's ID": {"conversation_id": "v2-tool-conversation", "message_id": v2ReactMessageID, "emoji": "πŸ‘"}, + "a blank idempotency key": {"conversation_id": v2ReactConversationID, "message_id": v2ReactMessageID, "emoji": "πŸ‘", "idempotency_key": " "}, + } { + t.Run(name, func(t *testing.T) { + if result := callReactTool(t, mcpServer, arguments); !result.IsError { + t.Fatalf("react_to_message = %v, want an error", v2ToolPayload(t, result)) + } + }) + } + if requests := sender.snapshotRequests(); len(requests) != 0 { + t.Fatalf("refused reactions reached the transport: %+v", requests) + } +} + +func TestV2PrimaryReactToMessageWithoutTheOutboxIsAnError(t *testing.T) { + a := testApp(t) + mcpServer := server.NewMCPServer("v2-react-unconfigured", "test") + RegisterWithOptions(mcpServer, a, Options{Reads: a.Store, V2Primary: true}) + result := callReactTool(t, mcpServer, map[string]any{ + "conversation_id": v2ReactConversationID, "message_id": v2ReactMessageID, "emoji": "πŸ‘", + }) + if !result.IsError || !strings.Contains(resultText(t, result), "v2 outbox") { + t.Fatalf("react_to_message = %v, want an error naming the v2 outbox", resultText(t, result)) + } +} + +// reactDaemon is a fake running app: its status and the reaction routes it +// answers are set per test, and it records which routes were called. +type reactDaemon struct { + t *testing.T + status map[string]any + mu sync.Mutex + calls []string + bodies []map[string]any +} + +func (d *reactDaemon) record(route string, r *http.Request) { + d.mu.Lock() + defer d.mu.Unlock() + d.calls = append(d.calls, route) + if r.Body != nil && r.Method == http.MethodPost { + var body map[string]any + _ = json.NewDecoder(r.Body).Decode(&body) + d.bodies = append(d.bodies, body) + } +} + +func (d *reactDaemon) routes() []string { + d.mu.Lock() + defer d.mu.Unlock() + return append([]string(nil), d.calls...) +} + +func (d *reactDaemon) handler(react, submit, delivery http.HandlerFunc) http.Handler { + mux := http.NewServeMux() + mux.HandleFunc("/api/status", func(w http.ResponseWriter, r *http.Request) { + json.NewEncoder(w).Encode(d.status) + }) + wrap := func(route string, next http.HandlerFunc) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + d.record(route, r) + if next == nil { + d.t.Errorf("the daemon's %s was called", route) + http.Error(w, "unexpected", http.StatusTeapot) + return + } + next(w, r) + } + } + mux.HandleFunc("POST /api/react", wrap("react", react)) + mux.HandleFunc("POST /api/v1/outbox/reactions", wrap("submit", submit)) + mux.HandleFunc("GET /api/v1/outbox/{id}", wrap("delivery", delivery)) + return mux +} + +func answerJSON(status int, body map[string]any) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + json.NewEncoder(w).Encode(body) + } +} + +func callDaemonReact(t *testing.T, mux http.Handler, arguments map[string]any) *mcp.CallToolResult { + t.Helper() + result, err := daemonReactToMessageHandler(Options{Daemon: daemonClientFor(t, mux)})(context.Background(), v2ToolCall(arguments)) + if err != nil { + t.Fatalf("handler: %v", err) + } + return result +} + +func TestDaemonReactToMessageOnAV2PrimaryApp(t *testing.T) { + arguments := map[string]any{ + "conversation_id": "v2-conv", "message_id": "v2-msg", "emoji": "❀️", "action": "Switch", + "idempotency_key": "client-react-key", + } + + t.Run("queues on the outbox and follows it to delivery", func(t *testing.T) { + daemon := &reactDaemon{t: t, status: map[string]any{"v2_primary": true, "v2_send": true}} + mux := daemon.handler(nil, + answerJSON(http.StatusOK, map[string]any{"outbox_id": "out-react", "state": "queued"}), + answerJSON(http.StatusOK, map[string]any{"outbox_id": "out-react", "state": "confirmed"})) + result := callDaemonReact(t, mux, arguments) + payload := v2ToolPayload(t, result) + if result.IsError { + t.Fatalf("result = error %v", payload) + } + assertV2ToolBool(t, payload, "ok", true) + assertV2ToolBool(t, payload, "settled", true) + assertV2ToolString(t, payload, "outbox_id", "out-react") + assertV2ToolString(t, payload, "idempotency_key", "client-react-key") + if routes := daemon.routes(); len(routes) < 2 || routes[0] != "submit" || routes[len(routes)-1] != "delivery" { + t.Fatalf("daemon routes = %v, want the outbox submit then delivery polls", routes) + } + want := map[string]any{ + "conversation_id": "v2-conv", "message_id": "v2-msg", "emoji": "❀️", + "action": "switch", "idempotency_key": "client-react-key", + } + if got := daemon.bodies[0]; len(got) != len(want) || got["action"] != want["action"] || + got["message_id"] != want["message_id"] || got["idempotency_key"] != want["idempotency_key"] { + t.Fatalf("submitted reaction = %v, want %v", got, want) + } + }) + + t.Run("a reaction the app keeps retrying is not an error", func(t *testing.T) { + daemon := &reactDaemon{t: t, status: map[string]any{"v2_primary": true}} + mux := daemon.handler(nil, + answerJSON(http.StatusOK, map[string]any{"outbox_id": "out-react", "state": "queued"}), + answerJSON(http.StatusOK, map[string]any{"outbox_id": "out-react", "state": "not_dispatched", "error_class": "transient"})) + result := callDaemonReact(t, mux, arguments) + payload := v2ToolPayload(t, result) + if result.IsError { + t.Fatalf("result = error %v", payload) + } + assertV2ToolBool(t, payload, "settled", false) + assertV2ToolBool(t, payload, "auto_retry", true) + }) + + t.Run("a refused reaction is an error", func(t *testing.T) { + daemon := &reactDaemon{t: t, status: map[string]any{"v2_primary": true}} + mux := daemon.handler(nil, answerJSON(http.StatusUnprocessableEntity, map[string]any{"error": "reaction_target_unavailable"}), nil) + result := callDaemonReact(t, mux, arguments) + if !result.IsError || !strings.Contains(resultText(t, result), "rejected by the app") { + t.Fatalf("result = %q, want a rejection", resultText(t, result)) + } + }) + + t.Run("a lost answer is reported as unknown, not as an error", func(t *testing.T) { + daemon := &reactDaemon{t: t, status: map[string]any{"v2_primary": true}} + mux := daemon.handler(nil, func(w http.ResponseWriter, r *http.Request) { + connection, _, err := http.NewResponseController(w).Hijack() + if err != nil { + t.Errorf("hijack: %v", err) + return + } + _ = connection.Close() + }, nil) + result := callDaemonReact(t, mux, arguments) + payload := v2ToolPayload(t, result) + if result.IsError { + t.Fatalf("a lost answer was reported as an error, which invites a second reaction: %v", payload) + } + assertV2ToolBool(t, payload, "ambiguous", true) + assertV2ToolString(t, payload, "idempotency_key", "client-react-key") + }) +} + +func TestDaemonReactToMessageOnALegacyPrimaryApp(t *testing.T) { + arguments := map[string]any{"conversation_id": "signal:+15551234567", "message_id": "signal:1", "emoji": "πŸ‘"} + + t.Run("v2 send without v2 primary keeps reactions on /api/react", func(t *testing.T) { + daemon := &reactDaemon{t: t, status: map[string]any{"v2_send": true, "v2_primary": false}} + mux := daemon.handler(answerJSON(http.StatusOK, map[string]any{"success": true}), nil, nil) + result := callDaemonReact(t, mux, arguments) + if result.IsError { + t.Fatalf("result = %q", resultText(t, result)) + } + if routes := daemon.routes(); len(routes) != 1 || routes[0] != "react" { + t.Fatalf("daemon routes = %v, want only /api/react", routes) + } + }) + + t.Run("an app answering success false did not apply the reaction", func(t *testing.T) { + daemon := &reactDaemon{t: t, status: map[string]any{}} + mux := daemon.handler(answerJSON(http.StatusOK, map[string]any{"success": false}), nil, nil) + if result := callDaemonReact(t, mux, arguments); !result.IsError { + t.Fatalf("result = %q, want an error", resultText(t, result)) + } + }) + + t.Run("a queued answer on /api/react is reported as queued", func(t *testing.T) { + // When the status probe fails the client falls back to /api/react, + // which a v2-primary app answers with 202 while the reaction is queued. + daemon := &reactDaemon{t: t, status: map[string]any{}} + mux := daemon.handler(answerJSON(http.StatusAccepted, map[string]any{ + "success": true, "queued": true, "outbox_id": "out-compat", "state": "not_dispatched", + }), nil, nil) + result := callDaemonReact(t, mux, arguments) + payload := v2ToolPayload(t, result) + if result.IsError { + t.Fatalf("result = error %v", payload) + } + assertV2ToolBool(t, payload, "settled", false) + assertV2ToolString(t, payload, "outbox_id", "out-compat") + if _, minted := payload["idempotency_key"]; minted { + t.Fatalf("payload = %v invents an idempotency key the app never returned", payload) + } + }) +} diff --git a/internal/tools/v2_send_test.go b/internal/tools/v2_send_test.go index 462521e5..77b77ace 100644 --- a/internal/tools/v2_send_test.go +++ b/internal/tools/v2_send_test.go @@ -626,6 +626,7 @@ type v2ToolRegistry struct { accountID string text bridge.TextSender media bridge.MediaSender + reaction bridge.ReactionSender } func (r *v2ToolRegistry) Snapshot(accountID string) (bridge.Snapshot, bool) { @@ -661,6 +662,8 @@ func (r *v2ToolRegistry) Acquire( lease.Text = r.text case bridge.CapabilityMediaSend: lease.Media = r.media + case bridge.CapabilityReactions: + lease.Reaction = r.reaction default: return nil, bridge.ErrCapabilityUnavailable } @@ -674,6 +677,7 @@ func (r *v2ToolRegistry) Capabilities(accountID string) bridge.CapabilitySet { return bridge.CapabilitySet{ TextSend: r.text != nil, MediaSend: r.media != nil, + Reactions: r.reaction != nil, } } diff --git a/internal/v2wire/submit_reaction.go b/internal/v2wire/submit_reaction.go new file mode 100644 index 00000000..99deeb6f --- /dev/null +++ b/internal/v2wire/submit_reaction.go @@ -0,0 +1,129 @@ +package v2wire + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "github.com/maxghenis/openmessage/internal/bridge" + "github.com/maxghenis/openmessage/internal/messaging" + "github.com/maxghenis/openmessage/internal/storage/sqlite" + "github.com/maxghenis/openmessage/internal/v2keys" +) + +// ErrReactionTargetUnavailable reports that a reaction names no message this +// store holds, or names it in a conversation the message is not in. +var ErrReactionTargetUnavailable = errors.New("reaction target is unavailable") + +// ReactionInput is a reaction addressed the way v2 reads hand IDs out. +type ReactionInput struct { + // ConversationID is optional: the target message decides the conversation + // and the account. A supplied value must name that conversation, by its + // v2 ID or by the remote ID that pre-cutover callers hold. + ConversationID string + // MessageID is the target's v2 message ID. + MessageID string + Emoji string + // Action is "add", "remove" or "switch"; empty means add. + Action string + IdempotencyKey string +} + +// SubmitReactionV2 resolves a reaction's account and conversation from its +// target message in the v2 store and submits it to the durable messaging +// service. It reads no legacy state. +func SubmitReactionV2( + ctx context.Context, + deps NativeDeps, + input ReactionInput, +) (messaging.Submission, error) { + if err := validateNativeSubmitDeps(ctx, deps); err != nil { + return messaging.Submission{}, err + } + messageID := strings.TrimSpace(input.MessageID) + if messageID == "" { + return messaging.Submission{}, fmt.Errorf("%w: message ID is empty", messaging.ErrInvalidCommand) + } + action, err := reactionAction(input.Action) + if err != nil { + return messaging.Submission{}, err + } + + repository, err := sqlite.NewMessageRepository(deps.V2, time.Now) + if err != nil { + return messaging.Submission{}, fmt.Errorf("resolve v2 reaction target %q: %w", messageID, err) + } + target, err := repository.GetMessage(ctx, messageID) + if errors.Is(err, sqlite.ErrNotFound) { + return messaging.Submission{}, fmt.Errorf( + "%w: no v2 message %q", + ErrReactionTargetUnavailable, + messageID, + ) + } + if err != nil { + return messaging.Submission{}, fmt.Errorf("resolve v2 reaction target %q: %w", messageID, err) + } + conversation, err := deps.V2.GetConversation(target.ConversationID) + if err != nil { + return messaging.Submission{}, fmt.Errorf( + "resolve v2 conversation %q of reaction target %q: %w", + target.ConversationID, + messageID, + err, + ) + } + if !namesConversation(input.ConversationID, conversation) { + return messaging.Submission{}, fmt.Errorf( + "%w: v2 message %q is not in conversation %q", + ErrReactionTargetUnavailable, + messageID, + strings.TrimSpace(input.ConversationID), + ) + } + if !deps.Registry.Capabilities(target.AccountID).Reactions { + return messaging.Submission{}, fmt.Errorf( + "%w: account %q does not support reactions", + ErrPlatformNotSendable, + target.AccountID, + ) + } + + return deps.Service.SendReaction(ctx, messaging.SendReactionCommand{ + CommonCommand: messaging.CommonCommand{ + AccountID: target.AccountID, + ConversationID: conversation.ConversationID, + IdempotencyKey: input.IdempotencyKey, + }, + TargetMessageID: target.MessageID, + Emoji: input.Emoji, + Action: action, + }) +} + +// namesConversation reports whether a caller-supplied conversation key names +// the conversation: empty (the caller left it to the target), its v2 ID, or +// its remote ID, which is the key the legacy store used and v2 reads still +// accept (issue #155). +func namesConversation(supplied string, conversation sqlite.Conversation) bool { + supplied = strings.TrimSpace(supplied) + if supplied == "" || supplied == conversation.ConversationID { + return true + } + // Only Signal remote IDs are normalized, and the rule (trim inside the + // "signal:" prefixes) leaves every other platform's ID as it is. + return v2keys.NormalizeRemoteConversationID("signal", supplied) == conversation.RemoteConversationID +} + +func reactionAction(raw string) (bridge.ReactionAction, error) { + switch action := bridge.ReactionAction(strings.ToLower(strings.TrimSpace(raw))); action { + case "": + return bridge.ReactionAdd, nil + case bridge.ReactionAdd, bridge.ReactionRemove, bridge.ReactionSwitch: + return action, nil + default: + return "", fmt.Errorf("%w: reaction action %q is invalid", messaging.ErrInvalidCommand, raw) + } +} diff --git a/internal/v2wire/submit_reaction_test.go b/internal/v2wire/submit_reaction_test.go new file mode 100644 index 00000000..4620a61c --- /dev/null +++ b/internal/v2wire/submit_reaction_test.go @@ -0,0 +1,230 @@ +package v2wire + +import ( + "context" + "errors" + "fmt" + "math/rand" + "reflect" + "testing" + "testing/quick" + "time" + + "github.com/maxghenis/openmessage/internal/bridge" + "github.com/maxghenis/openmessage/internal/messaging" + "github.com/maxghenis/openmessage/internal/storage/sqlite" +) + +// reactionSubmitFixture is a v2 store with three conversations on a reacting +// account and one on an account that cannot react, each holding one message. +type reactionSubmitFixture struct { + t *testing.T + store *sqlite.Store + deps NativeDeps + outbox *sqlite.OutboxRepository + keys int + // messages maps a message ID to the conversation it is in. + messages map[string]sqlite.Conversation +} + +func newReactionSubmitFixture(t *testing.T) *reactionSubmitFixture { + t.Helper() + store := openV2TestStore(t) + registry := submitTestRegistry{caps: map[string]bridge.CapabilitySet{ + "account-signal": {Reactions: true}, + "account-google": {Reactions: true}, + "account-import": {}, + }} + outbox, err := sqlite.NewOutboxRepository(store, time.Now) + if err != nil { + t.Fatalf("NewOutboxRepository(): %v", err) + } + fixture := &reactionSubmitFixture{ + t: t, store: store, outbox: outbox, messages: map[string]sqlite.Conversation{}, + deps: NativeDeps{V2: store, Service: newSubmitTestService(t, store, registry, nil), Registry: registry}, + } + for _, seed := range []struct{ account, conversation, remote, message string }{ + {"account-signal", "v2-signal-direct", "signal:+15550001111", "message-signal-direct"}, + {"account-signal", "v2-signal-group", "signal-group:Z3JvdXA=", "message-signal-group"}, + {"account-google", "v2-google-thread", "1234", "message-google"}, + {"account-import", "v2-import-thread", "imessage:chat1", "message-import"}, + } { + nowMS := time.Now().UnixMilli() + if err := store.UpsertAccount(sqlite.Account{ + AccountID: seed.account, BridgeKey: "reaction-test", DisplayName: seed.account, + Mode: sqlite.AccountModeLive, Enabled: true, ConfigJSON: "{}", + CreatedAtMS: nowMS, UpdatedAtMS: nowMS, + }); err != nil { + t.Fatalf("UpsertAccount(%q): %v", seed.account, err) + } + conversation := sqlite.Conversation{ + ConversationID: seed.conversation, AccountID: seed.account, RemoteConversationID: seed.remote, + Kind: sqlite.ConversationKindDirect, Title: seed.conversation, + NotificationMode: sqlite.NotificationModeAll, MetadataJSON: "{}", + CreatedAtMS: nowMS, UpdatedAtMS: nowMS, + } + if err := store.UpsertConversation(conversation); err != nil { + t.Fatalf("UpsertConversation(%q): %v", seed.conversation, err) + } + projectV2TestMessage(t, store, sqlite.Message{ + MessageID: seed.message, ConversationID: seed.conversation, AccountID: seed.account, + RemoteMessageID: "remote-" + seed.message, Direction: sqlite.MessageDirectionIncoming, + Body: "react to this", State: sqlite.MessageStateActive, OccurredAtMS: 1_900_000_000_000, + }) + fixture.messages[seed.message] = conversation + } + return fixture +} + +func (f *reactionSubmitFixture) submit(input ReactionInput) (messaging.Submission, error) { + if input.IdempotencyKey == "" { + f.keys++ + input.IdempotencyKey = fmt.Sprintf("reaction-submit-%d", f.keys) + } + if input.Emoji == "" { + input.Emoji = "πŸ‘" + } + return SubmitReactionV2(context.Background(), f.deps, input) +} + +func (f *reactionSubmitFixture) queuedReactions() int { + f.t.Helper() + rows, err := f.outbox.ListPending(context.Background(), sqlite.ListPendingParams{Limit: 10_000}) + if err != nil { + f.t.Fatalf("ListPending(): %v", err) + } + return len(rows) +} + +func TestSubmitReactionV2TakesAccountAndConversationFromTheTarget(t *testing.T) { + f := newReactionSubmitFixture(t) + submission, err := f.submit(ReactionInput{ + MessageID: " message-signal-group ", Emoji: " ❀️ ", Action: "Switch", IdempotencyKey: "reaction-key", + }) + if err != nil { + t.Fatalf("SubmitReactionV2(): %v", err) + } + if submission.State != messaging.OutboxQueued || submission.LocalMessageID != "" || submission.Deduplicated { + t.Fatalf("submission = %+v, want a new queued intent with no local message", submission) + } + item, err := f.outbox.FindByID(context.Background(), submission.OutboxID) + if err != nil { + t.Fatalf("FindByID(): %v", err) + } + reaction, err := f.outbox.GetOutboxReaction(context.Background(), submission.OutboxID) + if err != nil { + t.Fatalf("GetOutboxReaction(): %v", err) + } + if item.Kind != sqlite.OutboxKindReaction || item.AccountID != "account-signal" || + item.ConversationID != "v2-signal-group" || reaction.TargetMessageID != "message-signal-group" || + reaction.Emoji != "❀️" || reaction.Action != string(bridge.ReactionSwitch) { + t.Fatalf("queued reaction = %+v / %+v", item, reaction) + } + + replay, err := f.submit(ReactionInput{ + ConversationID: "v2-signal-group", MessageID: "message-signal-group", + Emoji: "❀️", Action: "switch", IdempotencyKey: "reaction-key", + }) + if err != nil || !replay.Deduplicated || replay.OutboxID != submission.OutboxID { + t.Fatalf("replay with the same key = %+v, %v; want the first intent", replay, err) + } + if _, err := f.submit(ReactionInput{ + MessageID: "message-signal-group", Emoji: "πŸ‘", IdempotencyKey: "reaction-key", + }); !errors.Is(err, messaging.ErrIdempotencyConflict) { + t.Fatalf("same key, other reaction = %v, want ErrIdempotencyConflict", err) + } + if queued := f.queuedReactions(); queued != 1 { + t.Fatalf("queued reactions = %d, want 1", queued) + } +} + +func TestSubmitReactionV2RefusesWhatItCannotPlace(t *testing.T) { + for _, refused := range []struct { + name string + input ReactionInput + want error + }{ + {"no message", ReactionInput{MessageID: " "}, messaging.ErrInvalidCommand}, + {"an action that is not a reaction action", ReactionInput{MessageID: "message-google", Action: "toggle"}, messaging.ErrInvalidCommand}, + {"a message this store does not hold", ReactionInput{MessageID: "signal:1700000001000"}, ErrReactionTargetUnavailable}, + {"a message in another conversation", ReactionInput{MessageID: "message-google", ConversationID: "v2-signal-direct"}, ErrReactionTargetUnavailable}, + {"another conversation's remote ID", ReactionInput{MessageID: "message-signal-group", ConversationID: "signal:+15550001111"}, ErrReactionTargetUnavailable}, + {"an account with no reaction sender", ReactionInput{MessageID: "message-import"}, ErrPlatformNotSendable}, + } { + t.Run(refused.name, func(t *testing.T) { + f := newReactionSubmitFixture(t) + if _, err := f.submit(refused.input); !errors.Is(err, refused.want) { + t.Fatalf("SubmitReactionV2() = %v, want %v", err, refused.want) + } + if queued := f.queuedReactions(); queued != 0 { + t.Fatalf("a refused reaction left %d outbox rows", queued) + } + }) + } +} + +// reactionAddress is a reaction as a caller might address it: one of the +// fixture's messages, and a conversation key that may or may not be its own. +type reactionAddress struct { + MessageID string + ConversationKey string +} + +func (reactionAddress) Generate(r *rand.Rand, _ int) reflect.Value { + messages := []string{"message-signal-direct", "message-signal-group", "message-google"} + keys := []string{ + "", " ", + "v2-signal-direct", " v2-signal-direct\t", "v2-signal-group", "v2-google-thread", "v2-import-thread", + "signal:+15550001111", "signal: +15550001111 ", " signal:+15550001111", "signal:+15550002222", + "signal-group:Z3JvdXA=", "signal-group: Z3JvdXA= ", "signal-group:b3RoZXI=", + "1234", " 1234 ", "12345", "imessage:chat1", "V2-SIGNAL-DIRECT", "remote-v2-signal-direct", + } + return reflect.ValueOf(reactionAddress{ + MessageID: messages[r.Intn(len(messages))], + ConversationKey: keys[r.Intn(len(keys))], + }) +} + +// TestSubmitReactionV2QueuesOnlyInTheTargetsConversation is the routing +// invariant: however a reaction is addressed, it is queued if and only if the +// conversation key is empty or names the target's own conversation (by v2 ID +// or by remote ID, whitespace aside), and a queued reaction always sits under +// the target's account and conversation. A refused one queues nothing. +func TestSubmitReactionV2QueuesOnlyInTheTargetsConversation(t *testing.T) { + f := newReactionSubmitFixture(t) + // The acceptable keys, written out per conversation rather than derived + // with the code under test. + names := map[string]map[string]bool{ + "v2-signal-direct": {"": true, " ": true, "v2-signal-direct": true, " v2-signal-direct\t": true, + "signal:+15550001111": true, "signal: +15550001111 ": true, " signal:+15550001111": true}, + "v2-signal-group": {"": true, " ": true, "v2-signal-group": true, + "signal-group:Z3JvdXA=": true, "signal-group: Z3JvdXA= ": true}, + "v2-google-thread": {"": true, " ": true, "v2-google-thread": true, "1234": true, " 1234 ": true}, + } + property := func(address reactionAddress) bool { + conversation := f.messages[address.MessageID] + before := f.queuedReactions() + submission, err := f.submit(ReactionInput{MessageID: address.MessageID, ConversationID: address.ConversationKey}) + if !names[conversation.ConversationID][address.ConversationKey] { + if !errors.Is(err, ErrReactionTargetUnavailable) || f.queuedReactions() != before { + t.Errorf("%+v: err = %v with %d new rows, want refused with none", address, err, f.queuedReactions()-before) + return false + } + return true + } + if err != nil { + t.Errorf("%+v: SubmitReactionV2() = %v, want queued", address, err) + return false + } + item, err := f.outbox.FindByID(context.Background(), submission.OutboxID) + if err != nil || item.AccountID != conversation.AccountID || item.ConversationID != conversation.ConversationID { + t.Errorf("%+v: queued under %q/%q (%v), want %q/%q", address, + item.AccountID, item.ConversationID, err, conversation.AccountID, conversation.ConversationID) + return false + } + return f.queuedReactions() == before+1 + } + if err := quick.Check(property, &quick.Config{MaxCount: 300, Rand: rand.New(rand.NewSource(20261009))}); err != nil { + t.Fatal(err) + } +} diff --git a/internal/web/api.go b/internal/web/api.go index 78a6781e..8608eafd 100644 --- a/internal/web/api.go +++ b/internal/web/api.go @@ -175,7 +175,7 @@ func APIHandlerWithOptions(store *db.Store, cli *client.Client, logger zerolog.L if reads == nil { reads = store } - registerV1Routes(mux, store, logger, opts.V2, opts.V2Primary) + v1 := registerV1Routes(mux, store, logger, opts.V2, opts.V2Primary) diagnosticsStartedAt := time.Now() getClient := func() *client.Client { if opts.Client != nil { @@ -2079,12 +2079,7 @@ func APIHandlerWithOptions(store *db.Store, cli *client.Client, logger zerolog.L httpError(w, "method not allowed", 405) return } - var req struct { - ConversationID string `json:"conversation_id"` - MessageID string `json:"message_id"` - Emoji string `json:"emoji"` - Action string `json:"action"` // "add", "remove", "switch"; default "add" - } + var req reactionRequest if err := json.NewDecoder(r.Body).Decode(&req); err != nil { httpError(w, "invalid JSON: "+err.Error(), 400) return @@ -2093,6 +2088,13 @@ func APIHandlerWithOptions(store *db.Store, cli *client.Client, logger zerolog.L httpError(w, "message_id and emoji are required", 400) return } + if opts.V2Primary { + // The read API hands out v2 IDs here, which the legacy routing + // below cannot place: it tells platforms apart by ID prefix or a + // legacy-store lookup and falls back to Google Messages. + v1.react(w, r, req) + return + } if isWhatsAppConversation(req.ConversationID) { if opts.SendWhatsAppReaction == nil { httpError(w, "whatsapp reactions are not available", 501) diff --git a/internal/web/apiv1.go b/internal/web/apiv1.go index e138a3e2..7430d073 100644 --- a/internal/web/apiv1.go +++ b/internal/web/apiv1.go @@ -2,7 +2,9 @@ package web import ( "context" + "crypto/rand" "database/sql" + "encoding/hex" "encoding/json" "errors" "fmt" @@ -33,6 +35,10 @@ type V2Options struct { V2Store *sqlite.Store Blobs *blob.BlobStore Registry bridge.Registry + + // ReactSettleTimeout bounds how long POST /api/react waits for a reaction + // to settle on a v2-primary daemon. Zero means reactSettleTimeout. + ReactSettleTimeout time.Duration } type v1API struct { @@ -40,8 +46,16 @@ type v1API struct { logger zerolog.Logger v2 *V2Options primary bool + // reactSettle bounds how long POST /api/react waits for a reaction to + // settle on a v2-primary daemon. + reactSettle time.Duration } +// reactSettleTimeout is the default reactSettle. It stays under the 10s +// request timeout of localapi clients, which report a slower answer as the +// app not running. +const reactSettleTimeout = 8 * time.Second + type v1SubmissionResponse struct { OutboxID string `json:"outbox_id"` LocalMessageID string `json:"local_message_id"` @@ -75,11 +89,15 @@ type v1PendingResponse struct { ErrorCode string `json:"error_code,omitempty"` } -func registerV1Routes(mux *http.ServeMux, legacy *db.Store, logger zerolog.Logger, v2 *V2Options, primary bool) { - api := &v1API{legacy: legacy, logger: logger, v2: v2, primary: primary} +func registerV1Routes(mux *http.ServeMux, legacy *db.Store, logger zerolog.Logger, v2 *V2Options, primary bool) *v1API { + api := &v1API{legacy: legacy, logger: logger, v2: v2, primary: primary, reactSettle: reactSettleTimeout} + if v2 != nil && v2.ReactSettleTimeout > 0 { + api.reactSettle = v2.ReactSettleTimeout + } mux.HandleFunc("POST /api/v1/outbox/messages", api.submitText) mux.HandleFunc("POST /api/v1/outbox/media", api.submitMedia) + mux.HandleFunc("POST /api/v1/outbox/reactions", api.submitReaction) mux.HandleFunc("GET /api/v1/outbox", api.listPending) mux.HandleFunc("GET /api/v1/outbox/{id}", api.getDelivery) mux.HandleFunc("POST /api/v1/outbox/{id}/cancel", api.cancel) @@ -93,6 +111,7 @@ func registerV1Routes(mux *http.ServeMux, legacy *db.Store, logger zerolog.Logge // including paths introduced by a newer client than this daemon knows. mux.HandleFunc("/api/v1", api.v1Fallback) mux.HandleFunc("/api/v1/", api.v1Fallback) + return api } func (a *v1API) v1Fallback(w http.ResponseWriter, _ *http.Request) { @@ -228,6 +247,162 @@ func (a *v1API) submitMedia(w http.ResponseWriter, r *http.Request) { writeJSON(w, submissionResponse(submission)) } +// reactionRequest is the body of POST /api/v1/outbox/reactions and of +// POST /api/react. conversation_id is optional: the target message decides +// the conversation. +type reactionRequest struct { + ConversationID string `json:"conversation_id"` + MessageID string `json:"message_id"` + Emoji string `json:"emoji"` + Action string `json:"action"` // "add", "remove", "switch"; default "add" + IdempotencyKey string `json:"idempotency_key"` +} + +func (request reactionRequest) input(idempotencyKey string) v2wire.ReactionInput { + return v2wire.ReactionInput{ + ConversationID: request.ConversationID, + MessageID: request.MessageID, + Emoji: request.Emoji, + Action: request.Action, + IdempotencyKey: idempotencyKey, + } +} + +// submitReaction queues a reaction on the durable outbox and returns as soon +// as the intent is stored, like the other outbox submissions: the response +// says nothing about delivery, which GET /api/v1/outbox/{id} reports. It +// takes v2 IDs, so it serves a v2-primary daemon only; on a legacy-primary +// daemon reactions stay on /api/react. +func (a *v1API) submitReaction(w http.ResponseWriter, r *http.Request) { + if !a.enabled(w) { + return + } + if !a.primary { + httpError(w, "legacy primary: use /api/react", http.StatusConflict) + return + } + var request reactionRequest + if err := json.NewDecoder(r.Body).Decode(&request); err != nil { + httpError(w, "invalid JSON: "+err.Error(), http.StatusBadRequest) + return + } + if strings.TrimSpace(request.MessageID) == "" || strings.TrimSpace(request.Emoji) == "" { + httpError(w, "message_id and emoji are required", http.StatusBadRequest) + return + } + idempotencyKey, err := normalizeRequiredV2IdempotencyKey(request.IdempotencyKey) + if err != nil { + httpError(w, err.Error(), http.StatusBadRequest) + return + } + if !a.submitDependenciesAvailable(w) { + return + } + submission, err := v2wire.SubmitReactionV2(r.Context(), a.nativeDeps(), request.input(idempotencyKey)) + if err != nil { + a.writeError(w, err) + return + } + writeJSON(w, submissionResponse(submission)) +} + +// react serves POST /api/react on a v2-primary daemon. The legacy handler +// answered only after the transport call; delivery is now the dispatcher's +// job, so this queues the reaction on the outbox and waits up to reactSettle +// for it to settle, and the status says how far it got (reactOutcome). +// A request without an idempotency_key is a new intent each time. +func (a *v1API) react(w http.ResponseWriter, r *http.Request, request reactionRequest) { + if !a.enabled(w) || !a.submitDependenciesAvailable(w) { + return + } + idempotencyKey, err := normalizeSendIdempotencyKey(request.IdempotencyKey) + if err != nil { + httpError(w, err.Error(), http.StatusBadRequest) + return + } + if idempotencyKey == "" { + if idempotencyKey, err = newReactionIdempotencyKey(); err != nil { + a.writeError(w, err) + return + } + } + submission, err := v2wire.SubmitReactionV2(r.Context(), a.nativeDeps(), request.input(idempotencyKey)) + if err != nil { + a.writeError(w, err) + return + } + + waitCtx, cancel := context.WithTimeout(r.Context(), a.reactSettle) + defer cancel() + delivery, err := a.v2.Service.Wait(waitCtx, submission.OutboxID) + if err != nil && delivery.OutboxID == "" { + // The intent is stored; only this read of it failed. + delivery = messaging.Delivery{OutboxID: submission.OutboxID, State: submission.State} + } + status, body := reactOutcome(delivery) + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(body) +} + +// reactOutcome maps the delivery state POST /api/react reached within its +// wait to the HTTP answer: +// +// - 200, success true: the transport accepted the reaction, and it is in +// the read model as this account's own. +// - 202, success true, queued true: the reaction is stored on the outbox +// and has not been delivered yet. The app keeps retrying it until it goes +// out or is canceled, so the caller must not repeat it. +// - 409, 502, with error: the reaction will not be delivered by this +// intent. It was canceled (409), the transport refused it (502, state +// rejected), or the transport call ended without an answer and the +// reaction may or may not have been applied (502, state uncertain). +// +// Every answer carries outbox_id and state, so a caller can follow the +// intent on GET /api/v1/outbox/{id}. +func reactOutcome(delivery messaging.Delivery) (int, map[string]any) { + body := map[string]any{ + "outbox_id": delivery.OutboxID, + "state": delivery.State, + } + if delivery.ErrorClass != "" { + body["error_class"] = delivery.ErrorClass + } + if delivery.ErrorCode != "" { + body["error_code"] = delivery.ErrorCode + } + fail := func(status int, message string) (int, map[string]any) { + body["success"] = false + body["error"] = "send reaction: " + message + return status, body + } + switch delivery.State { + case messaging.OutboxConfirmed, messaging.OutboxStoreFailed: + body["success"] = true + return http.StatusOK, body + case messaging.OutboxUncertain: + return fail(http.StatusBadGateway, "the transport gave no answer, so the reaction may or may not have been applied") + case messaging.OutboxRejected: + return fail(http.StatusBadGateway, "the transport refused the reaction") + case messaging.OutboxCanceled: + return fail(http.StatusConflict, "the reaction was canceled before it was sent") + default: + // queued, dispatching, not_dispatched, and any state added later: + // stored, not delivered, still the app's to send. + body["success"] = true + body["queued"] = true + return http.StatusAccepted, body + } +} + +func newReactionIdempotencyKey() (string, error) { + var value [16]byte + if _, err := rand.Read(value[:]); err != nil { + return "", fmt.Errorf("generate reaction idempotency key: %w", err) + } + return "react-" + hex.EncodeToString(value[:]), nil +} + func (a *v1API) listPending(w http.ResponseWriter, r *http.Request) { if !a.enabled(w) || !a.serviceAvailable(w) { return @@ -448,6 +623,8 @@ func v1ErrorResponse(err error) (int, string) { switch { case errors.Is(err, v2wire.ErrReplyTargetUnavailable): return http.StatusUnprocessableEntity, "reply_target_unavailable" + case errors.Is(err, v2wire.ErrReactionTargetUnavailable): + return http.StatusUnprocessableEntity, "reaction_target_unavailable" case errors.Is(err, messaging.ErrIdempotencyConflict): return http.StatusConflict, err.Error() case errors.Is(err, messaging.ErrInvalidState): diff --git a/internal/web/apiv1_test.go b/internal/web/apiv1_test.go index d9154fcf..7aec716f 100644 --- a/internal/web/apiv1_test.go +++ b/internal/web/apiv1_test.go @@ -33,6 +33,7 @@ func TestV1RoutesReturnServiceUnavailableWhenDisabled(t *testing.T) { }{ {http.MethodPost, "/api/v1/outbox/messages"}, {http.MethodPost, "/api/v1/outbox/media"}, + {http.MethodPost, "/api/v1/outbox/reactions"}, {http.MethodGet, "/api/v1/outbox"}, {http.MethodGet, "/api/v1/outbox/outbox-1"}, {http.MethodPost, "/api/v1/outbox/outbox-1/cancel"}, @@ -263,6 +264,7 @@ func TestV1ErrorStatusMapping(t *testing.T) { {name: "invalid state", err: messaging.ErrInvalidState, wantStatus: http.StatusConflict}, {name: "platform", err: v2wire.ErrPlatformNotSendable, wantStatus: http.StatusNotImplemented}, {name: "reply", err: v2wire.ErrReplyTargetUnavailable, wantStatus: http.StatusUnprocessableEntity, wantMessage: "reply_target_unavailable"}, + {name: "reaction", err: fmt.Errorf("%w: no v2 message", v2wire.ErrReactionTargetUnavailable), wantStatus: http.StatusUnprocessableEntity, wantMessage: "reaction_target_unavailable"}, {name: "media unavailable", err: media.ErrUnavailable, wantStatus: http.StatusServiceUnavailable}, {name: "not found", err: sqlite.ErrNotFound, wantStatus: http.StatusNotFound, wantMessage: "not found"}, {name: "too large", err: messaging.ErrTooLarge, wantStatus: http.StatusRequestEntityTooLarge}, diff --git a/internal/web/react_v2_test.go b/internal/web/react_v2_test.go new file mode 100644 index 00000000..af0f68d9 --- /dev/null +++ b/internal/web/react_v2_test.go @@ -0,0 +1,171 @@ +package web + +import ( + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/rs/zerolog" + + "github.com/maxghenis/openmessage/internal/db" + "github.com/maxghenis/openmessage/internal/messaging" +) + +// TestReactOutcomeAnswersEveryState is the contract of POST /api/react on a +// v2-primary daemon, checked over every outbox state and over a state this +// build does not know: +// +// - the answer is 200, 202, 409 or 502, and always names the intent; +// - success is true exactly on 200 and 202; +// - 200 means the transport accepted the reaction, and nothing else does; +// - 202 means stored and still the app's to send, and is the only answer +// marked queued; +// - an error is given exactly when the intent will not deliver the reaction. +func TestReactOutcomeAnswersEveryState(t *testing.T) { + want := map[messaging.OutboxState]int{ + messaging.OutboxQueued: http.StatusAccepted, + messaging.OutboxDispatching: http.StatusAccepted, + messaging.OutboxNotDispatched: http.StatusAccepted, + messaging.OutboxConfirmed: http.StatusOK, + messaging.OutboxStoreFailed: http.StatusOK, + messaging.OutboxUncertain: http.StatusBadGateway, + messaging.OutboxRejected: http.StatusBadGateway, + messaging.OutboxCanceled: http.StatusConflict, + "a-state-added-later": http.StatusAccepted, + } + for state, wantStatus := range want { + t.Run(string(state), func(t *testing.T) { + status, body := reactOutcome(messaging.Delivery{ + OutboxID: "outbox-1", State: state, ErrorClass: "transient", ErrorCode: "send_reaction", + }) + if status != wantStatus { + t.Fatalf("status = %d, want %d", status, wantStatus) + } + if body["outbox_id"] != "outbox-1" || body["state"] != state || + body["error_class"] != "transient" || body["error_code"] != "send_reaction" { + t.Fatalf("body = %v, want the intent and its error named", body) + } + if success, _ := body["success"].(bool); success != (status < 300) { + t.Fatalf("success = %v on %d", body["success"], status) + } + if _, queued := body["queued"]; queued != (status == http.StatusAccepted) { + t.Fatalf("queued present = %v on %d", queued, status) + } + message, failed := body["error"].(string) + if failed != (status >= 400) || (failed && !strings.HasPrefix(message, "send reaction: ")) { + t.Fatalf("error = %q on %d", message, status) + } + if _, err := json.Marshal(body); err != nil { + t.Fatalf("body does not encode: %v", err) + } + }) + } +} + +func postJSONTo(t *testing.T, handler http.Handler, path, body string) (int, map[string]any) { + t.Helper() + recorder := httptest.NewRecorder() + handler.ServeHTTP(recorder, httptest.NewRequest(http.MethodPost, "http://127.0.0.1"+path, strings.NewReader(body))) + raw, _ := io.ReadAll(recorder.Body) + var decoded map[string]any + if err := json.Unmarshal(raw, &decoded); err != nil { + t.Fatalf("POST %s = %d with a body that is not a JSON object: %s", path, recorder.Code, raw) + } + return recorder.Code, decoded +} + +// The outbox reaction route addresses messages by v2 ID, which only a +// v2-primary daemon hands out. +func TestV1ReactionRouteServesV2PrimaryOnly(t *testing.T) { + harness := newA3Harness(t, true) + status, body := postJSONTo(t, harness.handler, "/api/v1/outbox/reactions", + `{"message_id":"m","emoji":"πŸ‘","idempotency_key":"k"}`) + if status != http.StatusConflict || body["error"] != "legacy primary: use /api/react" { + t.Fatalf("v2 send without v2 primary = %d %v, want 409 naming /api/react", status, body) + } +} + +func TestV1ReactionRouteValidatesTheRequest(t *testing.T) { + harness := newA3Harness(t, true) + handler := APIHandlerWithOptions(harness.legacy, nil, zerolog.Nop(), nil, APIOptions{ + V2Primary: true, + V2: &V2Options{Service: harness.service, V2Store: harness.v2, Registry: harness.registry}, + }) + for name, request := range map[string]struct { + body string + wantStatus int + wantError string + }{ + "not JSON": {`{`, http.StatusBadRequest, ""}, + "no message": {`{"emoji":"πŸ‘","idempotency_key":"k"}`, http.StatusBadRequest, "message_id and emoji are required"}, + "no emoji": {`{"message_id":"m","emoji":" ","idempotency_key":"k"}`, http.StatusBadRequest, "message_id and emoji are required"}, + "no idempotency key": {`{"message_id":"m","emoji":"πŸ‘"}`, http.StatusBadRequest, "idempotency_key is required"}, + "an unusable key": {`{"message_id":"m","emoji":"πŸ‘","idempotency_key":"a b"}`, http.StatusBadRequest, "idempotency_key contains unsupported characters"}, + "an unknown message": {`{"message_id":"m","emoji":"πŸ‘","idempotency_key":"k"}`, http.StatusUnprocessableEntity, "reaction_target_unavailable"}, + "an unknown action": {`{"message_id":"m","emoji":"πŸ‘","action":"toggle","idempotency_key":"k"}`, http.StatusBadRequest, ""}, + } { + t.Run(name, func(t *testing.T) { + status, body := postJSONTo(t, handler, "/api/v1/outbox/reactions", request.body) + if status != request.wantStatus || (request.wantError != "" && body["error"] != request.wantError) { + t.Fatalf("POST = %d %v, want %d %q", status, body, request.wantStatus, request.wantError) + } + }) + } + + // /api/react applies the same rules, except that it mints a key when the + // caller gives none. + for name, request := range map[string]struct { + body string + wantStatus int + wantError string + }{ + "an unusable key": {`{"message_id":"m","emoji":"πŸ‘","idempotency_key":"a b"}`, http.StatusBadRequest, "idempotency_key contains unsupported characters"}, + "an unknown message": {`{"message_id":"m","emoji":"πŸ‘"}`, http.StatusUnprocessableEntity, "reaction_target_unavailable"}, + } { + t.Run("/api/react "+name, func(t *testing.T) { + status, body := postJSONTo(t, handler, "/api/react", request.body) + if status != request.wantStatus || body["error"] != request.wantError { + t.Fatalf("POST = %d %v, want %d %q", status, body, request.wantStatus, request.wantError) + } + }) + } + pending, err := harness.service.ListPending(t.Context(), messaging.ListPendingQuery{Limit: 10}) + if err != nil || len(pending) != 0 { + t.Fatalf("refused reactions left outbox rows: %+v, %v", pending, err) + } +} + +// With v2 send enabled but the legacy store still primary, the read API hands +// out legacy IDs, so /api/react must keep using the legacy senders and never +// touch the outbox. +func TestReactStaysOnLegacySendersUntilV2IsPrimary(t *testing.T) { + harness := newA3Harness(t, true) + var calls []string + handler := APIHandlerWithOptions(harness.legacy, nil, zerolog.Nop(), nil, APIOptions{ + V2: &V2Options{Service: harness.service, V2Store: harness.v2, Registry: harness.registry}, + SendSignalReaction: func(conversationID, messageID, emoji, action string) error { + calls = append(calls, conversationID+"|"+messageID+"|"+emoji+"|"+action) + return nil + }, + }) + if err := harness.legacy.UpsertConversation(&db.Conversation{ + ConversationID: "signal:+15551234567", Name: "Taylor Price", SourcePlatform: "signal", + }); err != nil { + t.Fatal(err) + } + status, body := postJSONTo(t, handler, "/api/react", + `{"conversation_id":"signal:+15551234567","message_id":"signal:target","emoji":"πŸ˜‚","action":"add"}`) + if status != http.StatusOK || body["success"] != true || len(body) != 1 { + t.Fatalf("POST /api/react = %d %v, want the legacy answer {success:true}", status, body) + } + if len(calls) != 1 || calls[0] != "signal:+15551234567|signal:target|πŸ˜‚|add" { + t.Fatalf("legacy Signal sender calls = %v", calls) + } + pending, err := harness.service.ListPending(t.Context(), messaging.ListPendingQuery{Limit: 10}) + if err != nil || len(pending) != 0 { + t.Fatalf("a legacy-primary reaction reached the outbox: %+v, %v", pending, err) + } +} diff --git a/internal/web/static/index.html b/internal/web/static/index.html index 83436f7c..d609e871 100644 --- a/internal/web/static/index.html +++ b/internal/web/static/index.html @@ -5990,7 +5990,10 @@

Forward message

const NOTIFICATION_MODE_MENTIONS = 'mentions'; const NOTIFICATION_MODE_MUTED = 'muted'; const OUTBOX_VISIBLE_STATES = new Set(['queued', 'dispatching', 'not_dispatched', 'uncertain', 'store_failed']); - const OUTBOX_SEND_KINDS = new Set(['text', 'media']); + const OUTBOX_TRAY_KINDS = new Set(['text', 'media', 'reaction']); + // A reaction is listed only while the app is still sending it. It has no + // "send again", so one that ends unresolved is reported when it is sent. + const OUTBOX_REACTION_STATES = new Set(['queued', 'dispatching', 'not_dispatched']); let activeConvoId = null; let conversations = []; let allConversations = []; @@ -11918,10 +11921,11 @@

Forward message

} // The v2 tray is a read-through view of the durable outbox. It deliberately - // omits terminal rows and non-message intents; confirmed sends live in the - // thread, while reactions/read receipts remain outside this surface. + // omits terminal rows and read receipts; confirmed sends and reactions live + // in the thread. function outboxRowView(row, now = Date.now()) { const stateName = String(row && row.state || '').trim().toLowerCase(); + const isReaction = String(row && row.kind || '').trim().toLowerCase() === 'reaction'; const scheduledForMS = Number(row && row.scheduled_for_ms || 0); const nextAttemptMS = Number(row && row.next_attempt_at_ms || 0); const scheduled = stateName === 'queued' && scheduledForMS > now; @@ -11942,22 +11946,28 @@

Forward message

view.label = scheduled ? 'Scheduled' : 'Sending…'; view.guidance = scheduled ? 'Waiting in the outbox until the scheduled time.' - : 'The app owns this send; do not resend it.'; + : (isReaction ? 'The app owns this reaction; do not react again.' : 'The app owns this send; do not resend it.'); view.action = 'cancel'; view.actionLabel = 'Cancel'; - view.actionTitle = scheduled ? 'Cancel scheduled send' : 'Cancel queued send'; + view.actionTitle = scheduled ? 'Cancel scheduled send' : (isReaction ? 'Cancel queued reaction' : 'Cancel queued send'); if (scheduled) view.timePrefix = 'For'; break; case 'dispatching': view.label = 'Sending…'; - view.guidance = 'The send is in progress; do not resend it.'; + view.guidance = isReaction + ? 'The reaction is in progress; do not react again.' + : 'The send is in progress; do not resend it.'; break; case 'not_dispatched': view.label = 'Retrying…'; - view.guidance = 'The app retries automatically; do not resend it.'; + view.guidance = isReaction + ? 'The app retries automatically; do not react again.' + : 'The app retries automatically; do not resend it.'; view.action = 'cancel'; view.actionLabel = 'Cancel'; - view.actionTitle = 'Cancel this automatically retrying send'; + view.actionTitle = isReaction + ? 'Cancel this automatically retrying reaction' + : 'Cancel this automatically retrying send'; break; case 'uncertain': view.label = 'Sent β€” unconfirmed'; @@ -11976,14 +11986,19 @@

Forward message

default: view.visible = false; } + if (isReaction && !OUTBOX_REACTION_STATES.has(stateName)) { + view.visible = false; + view.action = ''; + view.actionLabel = ''; + view.actionTitle = ''; + } return view; } function visibleOutboxRows() { return outboxRows.filter(row => { const kind = String(row && row.kind || '').trim().toLowerCase(); - const stateName = String(row && row.state || '').trim().toLowerCase(); - return OUTBOX_SEND_KINDS.has(kind) && OUTBOX_VISIBLE_STATES.has(stateName); + return OUTBOX_TRAY_KINDS.has(kind) && outboxRowView(row).visible; }); } @@ -12001,8 +12016,10 @@

Forward message

function outboxRowHTML(row, now) { const view = outboxRowView(row, now); - const summary = String(row && row.summary || '').trim() - || (String(row && row.kind || '').toLowerCase() === 'media' ? 'Attachment' : 'Message'); + const kind = String(row && row.kind || '').toLowerCase(); + const summary = kind === 'reaction' + ? `Reaction ${String(row && row.summary || '').trim()}`.trim() + : (String(row && row.summary || '').trim() || (kind === 'media' ? 'Attachment' : 'Message')); const time = view.timeMS > 0 ? `${view.timePrefix} ${formatScheduledWhen(view.timeMS)}` : ''; @@ -14059,18 +14076,41 @@

${escapeHtml(title)} Β· ${rows.length}

} // ─── Reactions ─── - window.sendReaction = async function(messageId, emoji) { - // Close pickers - document.querySelectorAll('.emoji-picker.show').forEach(p => p.classList.remove('show')); - document.querySelectorAll('.emoji-full-panel.show').forEach(p => p.classList.remove('show')); - closeComposeGifPanel(); + // On a v2 daemon /api/react stores the reaction on the outbox and waits a + // few seconds for it. 200 means the platform accepted it, 202 means it is + // stored and still the app's to send, and anything else means this reaction + // will not go out. + async function reactionError(response) { + if (state.v2Primary === true) { + if (response.status === 422) return new Error("Can't react to that message."); + if (response.status === 501) return new Error("This conversation can't send reactions."); + if (response.status === 503) return new Error('Reactions are starting up β€” try again.'); + } + return new Error(await responseError(response)); + } + + window.sendReaction = async function(messageId, emoji) { + // Close pickers + document.querySelectorAll('.emoji-picker.show').forEach(p => p.classList.remove('show')); + document.querySelectorAll('.emoji-full-panel.show').forEach(p => p.classList.remove('show')); + closeComposeGifPanel(); try { - await postJSON('/api/react', { - message_id: messageId, - emoji: emoji, - conversation_id: activeConvoId, - action: 'add', + const response = await fetch(API + '/api/react', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ + message_id: messageId, + emoji: emoji, + conversation_id: activeConvoId, + action: 'add', + }), }); + if (!response.ok) throw await reactionError(response); + const result = await response.json().catch(() => null); + if (response.status === 202 || (result && result.queued)) { + showThreadFeedback("Reaction not sent yet. The app keeps trying; it's in the outbox, where you can cancel it.", 'send'); + scheduleOutboxRefresh(0); + } // Reload to see updated reactions if (activeConvoId) await loadMessages(activeConvoId, {poll: true}); } catch (err) {