Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 14 additions & 6 deletions plugins/platforms/photon/sidecar/index.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,11 @@
// - POST /attachment/<opaque-handle>/(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": "..."}
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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");
}
Expand All @@ -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 || {};
Expand Down
47 changes: 47 additions & 0 deletions plugins/platforms/photon/sidecar/outbound-confirmation.mjs
Original file line number Diff line number Diff line change
@@ -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(),
};
}
54 changes: 54 additions & 0 deletions plugins/platforms/photon/sidecar/outbound-confirmation.test.mjs
Original file line number Diff line number Diff line change
@@ -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/
);
});
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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);
}
Expand Down
51 changes: 48 additions & 3 deletions tests/plugins/platforms/photon/test_spectrum_patch.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 };
"""


Expand Down Expand Up @@ -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(
Expand All @@ -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


Loading