From 6e824d917dcb3674ea42ce534cba0b089db86b8b Mon Sep 17 00:00:00 2001 From: "Dr. RAGnos" Date: Sat, 8 Aug 2026 23:44:01 -0600 Subject: [PATCH] fix(photon): confirm idempotent sidecar sends --- plugins/platforms/photon/sidecar/index.mjs | 20 ++++--- .../photon/sidecar/outbound-confirmation.mjs | 47 ++++++++++++++++ .../sidecar/outbound-confirmation.test.mjs | 54 +++++++++++++++++++ .../patch-spectrum-mixed-attachments.mjs | 41 +++++++++++--- .../platforms/photon/test_spectrum_patch.py | 51 ++++++++++++++++-- 5 files changed, 196 insertions(+), 17 deletions(-) create mode 100644 plugins/platforms/photon/sidecar/outbound-confirmation.mjs create mode 100644 plugins/platforms/photon/sidecar/outbound-confirmation.test.mjs diff --git a/plugins/platforms/photon/sidecar/index.mjs b/plugins/platforms/photon/sidecar/index.mjs index 8680f336c0f1..e1a682dbdc6f 100644 --- a/plugins/platforms/photon/sidecar/index.mjs +++ b/plugins/platforms/photon/sidecar/index.mjs @@ -25,9 +25,11 @@ // - POST /attachment//(consume|release) -> finalize only // after durable consumer commit, or release for retry // - POST /healthz -> {"ok": true} -// - POST /send -> {"ok": true, "messageId": "..."} +// - POST /send -> {"ok": true, "messageId": "..."} or a confirmed +// provider receipt when clientMessageId is supplied // body: {"spaceId": "...", "text": "...", -// "format": "text" | "markdown" (default "text")} +// "format": "text" | "markdown" (default "text"), +// "clientMessageId": "..." | null} // - POST /send-richlink -> {"ok": true, "messageId": "..."} // body: {"spaceId": "...", "url": "https://..."} // - POST /send-attachment -> {"ok": true, "messageId": "..."} @@ -74,6 +76,10 @@ import crypto from "node:crypto"; import { once } from "node:events"; import { patchSpectrumTs } from "./patch-spectrum-mixed-attachments.mjs"; import { chooseSendFormat } from "./send-format.mjs"; +import { + confirmedSendResponse, + withClientMessageId, +} from "./outbound-confirmation.mjs"; import { classifyProbeRejection, shouldProbe, @@ -1052,7 +1058,7 @@ const server = http.createServer(async (req, res) => { return ok(res, { deliveryId, status }); } if (req.url === "/send") { - const { spaceId, text, format = "text" } = body || {}; + const { spaceId, text, format = "text", clientMessageId } = body || {}; if (!spaceId || typeof text !== "string") { return badRequest(res, "spaceId and text are required"); } @@ -1068,12 +1074,14 @@ const server = http.createServer(async (req, res) => { // messages that contain URLs through spectrumText while preserving // spectrumMarkdown for URL-free markdown. The decision lives in // send-format.mjs so tests can exercise it directly. - const builder = + const builder = withClientMessageId( chooseSendFormat(format, text) === "markdown" ? spectrumMarkdown(text) - : spectrumText(text); + : spectrumText(text), + clientMessageId + ); const result = await space.send(builder); - return ok(res, { messageId: result?.id || null }); + return ok(res, confirmedSendResponse(result, clientMessageId)); } if (req.url === "/send-richlink") { const { spaceId, url } = body || {}; diff --git a/plugins/platforms/photon/sidecar/outbound-confirmation.mjs b/plugins/platforms/photon/sidecar/outbound-confirmation.mjs new file mode 100644 index 000000000000..e2707cd722d9 --- /dev/null +++ b/plugins/platforms/photon/sidecar/outbound-confirmation.mjs @@ -0,0 +1,47 @@ +const CLIENT_MESSAGE_ID_RE = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,159}$/; +const PROVIDER_MESSAGE_ID_RE = /^[A-Za-z0-9][A-Za-z0-9._:-]{0,255}$/; + +function normalizedClientMessageId(value) { + if (value === undefined || value === null || value === "") return null; + if (typeof value !== "string" || !CLIENT_MESSAGE_ID_RE.test(value)) { + throw new TypeError("clientMessageId is invalid"); + } + return value; +} +export function withClientMessageId(builder, value) { + const clientMessageId = normalizedClientMessageId(value); + if (clientMessageId === null) return builder; + if (!builder || typeof builder.build !== "function") { + throw new TypeError("content builder is invalid"); + } + return { + ...builder, + async build(...args) { + const content = await builder.build(...args); + if (!content || typeof content !== "object") { + throw new TypeError("built content is invalid"); + } + return { ...content, clientMessageId }; + }, + }; +} + +export function confirmedSendResponse(result, value) { + const clientMessageId = normalizedClientMessageId(value); + const messageId = typeof result?.id === "string" ? result.id.trim() : ""; + if (!PROVIDER_MESSAGE_ID_RE.test(messageId)) { + throw new TypeError("provider message id is invalid"); + } + if (clientMessageId === null) return { messageId }; + const timestamp = result?.timestamp; + if (!(timestamp instanceof Date) || Number.isNaN(timestamp.valueOf())) { + throw new TypeError("provider timestamp is invalid"); + } + return { + messageId, + clientMessageId, + confirmed: true, + providerStatus: "accepted", + deliveredAt: timestamp.toISOString(), + }; +} diff --git a/plugins/platforms/photon/sidecar/outbound-confirmation.test.mjs b/plugins/platforms/photon/sidecar/outbound-confirmation.test.mjs new file mode 100644 index 000000000000..f83b9e39d0ae --- /dev/null +++ b/plugins/platforms/photon/sidecar/outbound-confirmation.test.mjs @@ -0,0 +1,54 @@ +import assert from "node:assert/strict"; +import test from "node:test"; + +import { + confirmedSendResponse, + withClientMessageId, +} from "./outbound-confirmation.mjs"; + +test("stable client message id reaches the built provider content", async () => { + const builder = { build: async () => ({ type: "text", text: "private" }) }; + const wrapped = withClientMessageId(builder, "kpx-lhll-123"); + + assert.deepEqual(await wrapped.build(), { + type: "text", + text: "private", + clientMessageId: "kpx-lhll-123", + }); +}); +test("provider snapshot becomes an exact confirmation response", () => { + assert.deepEqual( + confirmedSendResponse( + { id: "spc-msg-123", timestamp: new Date("2026-08-09T05:27:45.678Z") }, + "kpx-lhll-123" + ), + { + messageId: "spc-msg-123", + clientMessageId: "kpx-lhll-123", + confirmed: true, + providerStatus: "accepted", + deliveredAt: "2026-08-09T05:27:45.678Z", + } + ); +}); + +test("legacy sends stay compatible and invalid receipts fail closed", () => { + const builder = { build: async () => ({ type: "text", text: "private" }) }; + assert.equal(withClientMessageId(builder, undefined), builder); + assert.deepEqual( + confirmedSendResponse( + { id: "spc-msg-123", timestamp: new Date("2026-08-09T05:27:45.678Z") }, + undefined + ), + { messageId: "spc-msg-123" } + ); + assert.throws(() => withClientMessageId(builder, "contains spaces"), /invalid/); + assert.throws( + () => confirmedSendResponse({ id: null, timestamp: new Date() }, "kpx-123"), + /provider message id/ + ); + assert.throws( + () => confirmedSendResponse({ id: "spc-msg-123" }, "kpx-123"), + /provider timestamp/ + ); +}); diff --git a/plugins/platforms/photon/sidecar/patch-spectrum-mixed-attachments.mjs b/plugins/platforms/photon/sidecar/patch-spectrum-mixed-attachments.mjs index 151043bc9e87..33c3747788e4 100644 --- a/plugins/platforms/photon/sidecar/patch-spectrum-mixed-attachments.mjs +++ b/plugins/platforms/photon/sidecar/patch-spectrum-mixed-attachments.mjs @@ -18,7 +18,8 @@ import fs from "node:fs"; import path from "node:path"; import { fileURLToPath, pathToFileURL } from "node:url"; -const MARKER = "Hermes patch: Preserve mixed text + attachment iMessage payloads"; +const MIXED_MARKER = "Hermes patch: Preserve mixed text + attachment iMessage payloads"; +const IDEMPOTENCY_MARKER = "Hermes patch: Forward stable clientMessageId on text sends"; function scriptDir() { return path.dirname(fileURLToPath(import.meta.url)); @@ -115,6 +116,21 @@ function patchChildIndices(source) { ); } +function patchOutboundIdempotency(source) { + source = replaceOnce( + source, + `\t\tcase "text": return outboundMessage(spaceId, await remote.messages.sendText(chat, content.text, withReply(effectOption(effect), replyTo)), content);`, + `\t\tcase "text": return outboundMessage(spaceId, await remote.messages.sendText(chat, content.text, withReply({ ...effectOption(effect), clientMessageId: content.clientMessageId }, replyTo)), content);`, + "outbound text client message id" + ); + return replaceOnce( + source, + `\t\t\treturn outboundMessage(spaceId, await remote.messages.sendText(chat, rendered.text, withReply({\n\t\t\t\t...effectOption(effect),\n\t\t\t\t...formattingOption(rendered.formatting)\n\t\t\t}, replyTo)), content);`, + `\t\t\treturn outboundMessage(spaceId, await remote.messages.sendText(chat, rendered.text, withReply({\n\t\t\t\t...effectOption(effect),\n\t\t\t\t...formattingOption(rendered.formatting),\n\t\t\t\tclientMessageId: content.clientMessageId\n\t\t\t}, replyTo)), content);`, + "outbound markdown client message id" + ); +} + export function patchSpectrumTs(root = scriptDir()) { const dist = path.join( root, @@ -132,7 +148,7 @@ export function patchSpectrumTs(root = scriptDir()) { for (const file of files) { const raw = fs.readFileSync(file, "utf8"); - if (raw.includes(MARKER)) { + if (raw.includes(MIXED_MARKER) && raw.includes(IDEMPOTENCY_MARKER)) { return { patched: false, file, reason: "already patched" }; } // Normalize to LF for matching so the patch works regardless of the @@ -144,15 +160,24 @@ export function patchSpectrumTs(root = scriptDir()) { const CRLF = CR + "\n"; const usedCRLF = raw.includes(CRLF); const original = usedCRLF ? raw.split(CRLF).join("\n") : raw; - if (!original.includes("const toInboundMessages = async") || - !original.includes("const rebuildFromAppleMessage = async")) { + if (!original.includes("const sendContent = async")) { continue; } let patched = original; - patched = patchRebuild(patched); - patched = patchInbound(patched); - patched = patchChildIndices(patched); - patched = `// ${MARKER}\n${patched}`; + if (!patched.includes(MIXED_MARKER)) { + if (!patched.includes("const toInboundMessages = async") || + !patched.includes("const rebuildFromAppleMessage = async")) { + continue; + } + patched = patchRebuild(patched); + patched = patchInbound(patched); + patched = patchChildIndices(patched); + patched = `// ${MIXED_MARKER}\n${patched}`; + } + if (!patched.includes(IDEMPOTENCY_MARKER)) { + patched = patchOutboundIdempotency(patched); + patched = `// ${IDEMPOTENCY_MARKER}\n${patched}`; + } if (usedCRLF) { patched = patched.split("\n").join(CRLF); } diff --git a/tests/plugins/platforms/photon/test_spectrum_patch.py b/tests/plugins/platforms/photon/test_spectrum_patch.py index 20bc98b33c58..8b104ef153b4 100644 --- a/tests/plugins/platforms/photon/test_spectrum_patch.py +++ b/tests/plugins/platforms/photon/test_spectrum_patch.py @@ -225,7 +225,29 @@ def _tabify(src: str) -> str: cacheMessage(cache, msg); return [msg]; }; -export { rebuildFromAppleMessage, toInboundMessages }; +const outboundMessage = (spaceId, message, content) => ({ + id: message.guid, + content, + space: { id: spaceId }, + timestamp: message.dateCreated +}); +const withReply = (options, replyTo) => replyTo ? { ...options, replyTo } : options; +const effectOption = (effect) => effect ? { effect } : {}; +const formattingOption = (formatting) => formatting.length > 0 ? { formatting } : {}; +const renderMarkdown = (markdown) => ({ text: markdown, formatting: [] }); +const sendContent = async (remote, spaceId, chat, content, replyTo, effect) => { + switch (content.type) { + case "text": return outboundMessage(spaceId, await remote.messages.sendText(chat, content.text, withReply(effectOption(effect), replyTo)), content); + case "markdown": { + const rendered = renderMarkdown(content.markdown); + return outboundMessage(spaceId, await remote.messages.sendText(chat, rendered.text, withReply({ + ...effectOption(effect), + ...formattingOption(rendered.formatting) + }, replyTo)), content); + } + } +}; +export { rebuildFromAppleMessage, sendContent, toInboundMessages }; """ @@ -262,6 +284,31 @@ def test_spectrum_patch_rewrites_the_imessage_mapper(tmp_path: Path) -> None: # The text is captured in both mappers before the attachment branches run. assert "const text2 = message.content.text;" in patched assert "const text2 = event.message.content.text;" in patched + assert "clientMessageId: content.clientMessageId" in patched + + behavior = subprocess.run( + [ + "node", + "--input-type=module", + "-e", + ( + f'import {{ sendContent }} from {json.dumps(chunk.as_uri())};' + 'const seen=[]; const remote={messages:{sendText:async(...args)=>' + '(seen.push(args),{guid:"spc-msg-1",dateCreated:new Date(0)})}};' + 'await sendContent(remote,"space","any;-;+10000000000",' + '{type:"text",text:"private",clientMessageId:"stable-1"});' + 'await sendContent(remote,"space","any;-;+10000000000",' + '{type:"markdown",markdown:"private",clientMessageId:"stable-2"});' + 'console.log(JSON.stringify(seen.map(args=>args[2].clientMessageId)));' + ), + ], + cwd=Path.cwd(), + text=True, + capture_output=True, + check=False, + ) + assert behavior.returncode == 0, behavior.stderr + assert json.loads(behavior.stdout) == ["stable-1", "stable-2"] # Re-running is a no-op (idempotent self-heal on every sidecar start). again = subprocess.run( @@ -273,5 +320,3 @@ def test_spectrum_patch_rewrites_the_imessage_mapper(tmp_path: Path) -> None: ) assert again.returncode == 0, again.stderr assert chunk.read_text(encoding="utf-8") == patched - -