Skip to content

Commit e128765

Browse files
alex-w-99claude
andcommitted
Enrich per-call and per-turn context; add burst batching and channel overrides
- Voice: full per-call system instructions (own identity block, caller's contact card incl. notes, known/unknown-caller directives, outbound purpose/opening guidance, consult-tool capability list, post-call action and two-step hangup choreography, third-party privacy rule) plus a direction-aware greeting instruction; call parties resolve before the mode is chosen; voice frames and post-call prompts carry call_id. - Reactions: name the target message and frame tapbacks with reply-restraint guidance ending in the [SILENT] escape. - Burst batching: opt-in quiet-window merging of rapid SMS/iMessage fragments into one tagged turn (gateway.textBatchWindowMs / INKBOX_TEXT_BATCH_WINDOW_MS; caps on count and size; commands bypass). - Frames: group participants list, RFC Message-ID on email tags for reply threading, contact notes on the resolver (voice-only injection). - Channel overrides: gateway.channelPrompts (per-turn operator directive) and gateway.channelAgents (opencode agent override), keyed by contact id or channel. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 9c4ed6a commit e128765

15 files changed

Lines changed: 675 additions & 65 deletions

‎src/config.ts‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,15 @@ export interface GatewayOptions {
7373
agent?: string;
7474
// Optional model override for gateway sessions, "provider/model".
7575
model?: string;
76+
// Batch rapid-fire SMS/iMessage fragments arriving within this quiet
77+
// window (ms) into one merged turn. 0 (default) disables batching.
78+
textBatchWindowMs?: number;
79+
// Extra per-turn directive text, keyed by contact id or channel name
80+
// (contact id wins). Injected under the frame tag on matching turns.
81+
channelPrompts?: Record<string, string>;
82+
// opencode agent override keyed by contact id or channel name (contact id
83+
// wins); falls back to the gateway-wide agent.
84+
channelAgents?: Record<string, string>;
7685
voice?: {
7786
enabled?: boolean;
7887
realtime?: {
@@ -104,6 +113,9 @@ export interface ResolvedGatewayConfig {
104113
mediaDir?: string;
105114
agent?: string;
106115
model?: string;
116+
textBatchWindowMs: number;
117+
channelPrompts: Record<string, string>;
118+
channelAgents: Record<string, string>;
107119
voice: {
108120
enabled: boolean;
109121
realtime: {
@@ -157,6 +169,18 @@ function stringArray(value: unknown): string[] {
157169
.filter((entry): entry is string => Boolean(entry));
158170
}
159171

172+
// Keep only entries whose key and value are both non-empty strings.
173+
function stringRecord(value: unknown): Record<string, string> {
174+
if (!isRecord(value)) return {};
175+
const out: Record<string, string> = {};
176+
for (const [k, v] of Object.entries(value)) {
177+
const key = nonEmptyString(k);
178+
const val = nonEmptyString(v);
179+
if (key && val) out[key] = val;
180+
}
181+
return out;
182+
}
183+
160184
// ~/.inkbox/config — `key = value` lines, the same file the Inkbox SDK and CLI
161185
// read. Lowest-precedence credential source, after options and env vars.
162186
function readInkboxConfigFile(): Record<string, string> {
@@ -313,6 +337,10 @@ function resolveGatewayConfig(
313337
mediaDir: nonEmptyString(opts.mediaDir) ?? nonEmptyString(env.INKBOX_OPENCODE_MEDIA_DIR),
314338
agent: nonEmptyString(opts.agent),
315339
model: nonEmptyString(opts.model),
340+
textBatchWindowMs:
341+
numeric(opts.textBatchWindowMs) ?? numeric(env.INKBOX_TEXT_BATCH_WINDOW_MS) ?? 0,
342+
channelPrompts: stringRecord(opts.channelPrompts),
343+
channelAgents: stringRecord(opts.channelAgents),
316344
voice: {
317345
enabled: voice.enabled ?? boolEnv(env.INKBOX_VOICE_ENABLED) ?? false,
318346
realtime: {

‎src/gateway/burst.ts‎

Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
1+
import type { InboundMessage } from "./types.js";
2+
3+
// Humans text in fragments and corrections. Batching fragments that arrive
4+
// within a quiet window into one merged turn keeps a single thought from
5+
// becoming several self-interrupting agent turns.
6+
export const DEFAULT_BURST_MAX_MESSAGES = 8;
7+
export const DEFAULT_BURST_MAX_CHARS = 4000;
8+
9+
export interface BurstBuffer {
10+
// Queue a message for its chat's pending batch (starting one if needed).
11+
add(msg: InboundMessage): void;
12+
// Flush every pending batch immediately (shutdown).
13+
flushAll(): void;
14+
}
15+
16+
interface Pending {
17+
msgs: InboundMessage[];
18+
chars: number;
19+
timer: ReturnType<typeof setTimeout>;
20+
}
21+
22+
export function createBurstBuffer(opts: {
23+
windowMs: number;
24+
maxMessages?: number;
25+
maxChars?: number;
26+
deliver(msg: InboundMessage): void;
27+
}): BurstBuffer {
28+
const maxMessages = opts.maxMessages ?? DEFAULT_BURST_MAX_MESSAGES;
29+
const maxChars = opts.maxChars ?? DEFAULT_BURST_MAX_CHARS;
30+
const pending = new Map<string, Pending>();
31+
32+
function flush(chatKey: string): void {
33+
const batch = pending.get(chatKey);
34+
if (!batch) return;
35+
pending.delete(chatKey);
36+
clearTimeout(batch.timer);
37+
opts.deliver(mergeBurst(batch.msgs));
38+
}
39+
40+
return {
41+
add(msg) {
42+
const existing = pending.get(msg.chatKey);
43+
if (!existing) {
44+
const timer = setTimeout(() => flush(msg.chatKey), opts.windowMs);
45+
timer.unref?.();
46+
pending.set(msg.chatKey, { msgs: [msg], chars: msg.text.length, timer });
47+
return;
48+
}
49+
existing.msgs.push(msg);
50+
existing.chars += msg.text.length;
51+
// Sliding quiet window: each fragment restarts the countdown.
52+
clearTimeout(existing.timer);
53+
if (existing.msgs.length >= maxMessages || existing.chars >= maxChars) {
54+
flush(msg.chatKey);
55+
return;
56+
}
57+
existing.timer = setTimeout(() => flush(msg.chatKey), opts.windowMs);
58+
existing.timer.unref?.();
59+
},
60+
61+
flushAll() {
62+
for (const chatKey of [...pending.keys()]) flush(chatKey);
63+
},
64+
};
65+
}
66+
67+
// Collapse a batch into one message: newest metadata (ids, threading), all
68+
// fragment texts in arrival order, and every attachment.
69+
export function mergeBurst(msgs: InboundMessage[]): InboundMessage {
70+
if (msgs.length === 1) return msgs[0];
71+
const last = msgs[msgs.length - 1];
72+
return {
73+
...last,
74+
text: msgs
75+
.map((m) => m.text)
76+
.filter((t) => t.trim() !== "")
77+
.join("\n"),
78+
mediaPaths: msgs.flatMap((m) => m.mediaPaths),
79+
burst: msgs.length,
80+
};
81+
}

‎src/gateway/contacts.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,9 @@ export interface ResolvedContact {
1010
contactCompany?: string;
1111
contactEmails?: string[];
1212
contactPhones?: string[];
13+
// Free-form notes from the contact record. Injected into voice-call
14+
// instructions only — never into per-message frame tags (can be long).
15+
contactNotes?: string;
1316
}
1417

1518
// One-line contact card for [inkbox:...] frame tags: the addresses the agent
@@ -80,12 +83,14 @@ export function createContactResolver(
8083
const company = contact.companyName?.trim();
8184
const emails = (contact.emails ?? []).map((e) => e.value).filter(Boolean);
8285
const phones = (contact.phones ?? []).map((p) => p.value).filter(Boolean);
86+
const notes = contact.notes?.trim();
8387
return {
8488
contactId: contact.id,
8589
...(name ? { contactName: name } : {}),
8690
...(company ? { contactCompany: company } : {}),
8791
...(emails.length ? { contactEmails: emails } : {}),
8892
...(phones.length ? { contactPhones: phones } : {}),
93+
...(notes ? { contactNotes: notes } : {}),
8994
};
9095
} catch (error) {
9196
// Resolution must never drop a message: warn and fall through to the

‎src/gateway/dispatch.ts‎

Lines changed: 44 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,10 @@
11
import type { InkboxRuntime } from "../client.js";
22
import type { ResolvedConfig, ResolvedGatewayConfig } from "../config.js";
3+
import type { BurstBuffer } from "./burst.js";
34
import type { ContactResolver } from "./contacts.js";
45
import type { NotifyOnce } from "./dedup.js";
56
import { downloadMedia, mediaDir } from "./media.js";
7+
import { SILENT } from "./prompts.js";
68
import type {
79
Channel,
810
GatewayLogger,
@@ -30,6 +32,9 @@ export interface DispatchDeps {
3032
sessions: SessionManager;
3133
notify: NotifyOnce;
3234
logger: GatewayLogger;
35+
// When set, rapid-fire SMS/iMessage fragments are batched per chat into
36+
// one merged turn instead of dispatched individually.
37+
bursts?: BurstBuffer;
3338
// Handle a verified non-Inkbox (external) webhook.
3439
onExternal?(event: VerifiedEvent): Promise<void>;
3540
}
@@ -227,7 +232,9 @@ async function handleInbound(
227232
...resolved,
228233
text: info.text,
229234
mediaPaths,
230-
...(participants > 1 ? { group: { participantCount: participants } } : {}),
235+
...(participants > 1
236+
? { group: { participantCount: participants, participants: participantNames(event.body) } }
237+
: {}),
231238
};
232239

233240
// A media-only message still wakes the agent.
@@ -236,6 +243,13 @@ async function handleInbound(
236243
return true;
237244
}
238245

246+
// Batch phone-channel fragments when enabled; slash commands bypass the
247+
// window so control replies stay immediate.
248+
if (deps.bursts && channel !== "email" && !msg.text.trim().startsWith("/")) {
249+
deps.bursts.add(msg);
250+
return true;
251+
}
252+
239253
// Kick off the turn without blocking the webhook response — a model run can
240254
// take minutes, and holding the HTTP request open makes the provider retry.
241255
void deps.sessions
@@ -253,11 +267,28 @@ function countParticipants(body: Record<string, unknown>): number {
253267
return Math.max(contacts, identities);
254268
}
255269

270+
// Names of the resolved remote parties on the event, for the group frame tag.
271+
function participantNames(body: Record<string, unknown>): string[] {
272+
const data = record(body.data);
273+
const names: string[] = [];
274+
for (const c of Array.isArray(data?.contacts) ? data.contacts : []) {
275+
const n = str(record(c)?.name);
276+
if (n) names.push(n);
277+
}
278+
for (const a of Array.isArray(data?.agent_identities) ? data.agent_identities : []) {
279+
const rec = record(a);
280+
const n = str(rec?.display_name) ?? str(rec?.agent_handle);
281+
if (n) names.push(n);
282+
}
283+
return names;
284+
}
285+
256286
async function handleReaction(deps: DispatchDeps, event: VerifiedEvent): Promise<boolean> {
257287
const r = resourceOf(event.body, "reaction");
258288
const from = str(r?.remote_number);
259289
const conversationId = str(r?.conversation_id);
260-
const reaction = str(r?.reaction) ?? "reaction";
290+
const reaction = str(r?.custom_emoji) ?? str(r?.reaction) ?? "reaction";
291+
const targetMessageId = str(r?.target_message_id);
261292
if (!from) return true;
262293
const resolved = await deps.contacts.resolve(from);
263294
if (!senderAllowed(from, resolved.contactId, deps.config.gateway)) return true;
@@ -267,14 +298,24 @@ async function handleReaction(deps: DispatchDeps, event: VerifiedEvent): Promise
267298
conversationId,
268299
from,
269300
});
301+
// Reactions carry reply-restraint guidance: a tapback is a lightweight
302+
// signal, and most warrant no visible reply at all.
303+
const who = resolved.contactName ?? from;
304+
const text = [
305+
`[reaction: ${reaction}${targetMessageId ? ` target_message_id=${targetMessageId}` : ""}]`,
306+
`${who} reacted with a '${reaction}' tapback to your message.`,
307+
"A reaction is a lightweight signal, not always a request for a reply — a 'question' " +
308+
"tapback usually wants a follow-up; 'love', 'like', 'laugh', or 'dislike' are usually " +
309+
`just acknowledgement. If no reply is warranted, reply with exactly ${SILENT}.`,
310+
].join("\n");
270311
void deps.sessions
271312
.handleInbound({
272313
channel: "imessage",
273314
chatKey,
274315
from,
275316
conversationId,
276317
...resolved,
277-
text: `[reaction: ${reaction}]`,
318+
text,
278319
mediaPaths: [],
279320
messageId: str(r?.id),
280321
})

‎src/gateway/index.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import type { OpencodeClient } from "@opencode-ai/sdk";
22
import type { InkboxRuntime } from "../client.js";
33
import type { ResolvedConfig } from "../config.js";
4+
import { createBurstBuffer } from "./burst.js";
45
import { handleCommand } from "./commands.js";
56
import { createContactResolver } from "./contacts.js";
67
import { createNotifyOnce, createRequestDedup } from "./dedup.js";
@@ -147,6 +148,19 @@ export async function startGateway(opts: StartGatewayOptions): Promise<GatewayHa
147148

148149
const events = subscribeEvents(opts.opencode, escalation, logger, opts.directory);
149150

151+
// Fragment batching for phone channels, when a quiet window is configured.
152+
const bursts =
153+
g.textBatchWindowMs > 0
154+
? createBurstBuffer({
155+
windowMs: g.textBatchWindowMs,
156+
deliver: (msg) => {
157+
void wrapSessions()
158+
.handleInbound(msg)
159+
.catch((err) => logger.error("turn.dispatch_failed", { error: String(err) }));
160+
},
161+
})
162+
: undefined;
163+
150164
async function onEvent(event: VerifiedEvent): Promise<boolean | undefined> {
151165
return dispatchEvent(
152166
{
@@ -156,6 +170,7 @@ export async function startGateway(opts: StartGatewayOptions): Promise<GatewayHa
156170
sessions: wrapSessions(),
157171
notify,
158172
logger,
173+
bursts,
159174
onExternal: g.externalEvents ? handleExternal : undefined,
160175
},
161176
event,
@@ -245,6 +260,7 @@ export async function startGateway(opts: StartGatewayOptions): Promise<GatewayHa
245260
publicUrl: transport.publicUrl,
246261
async close() {
247262
events.close();
263+
bursts?.flushAll();
248264
await sessions.close();
249265
await server.close();
250266
await transport.close();

‎src/gateway/prompts.ts‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -126,18 +126,30 @@ function groupReminder(participantCount: number): string {
126126
// Prefix an inbound message with a one-line routing tag — the per-turn
127127
// context the static system prompt can't carry: which channel this message
128128
// arrived on and who sent it. Quoted fields may contain spaces.
129-
export function frameInbound(msg: InboundMessage): string {
130-
const fields = [`inkbox:${msg.channel}`, `from=${msg.from}`];
129+
export function frameInbound(msg: InboundMessage, directive?: string): string {
130+
// Merged fragment bursts are tagged as such, with the fragment count.
131+
const burst = (msg.burst ?? 1) > 1;
132+
const fields = [
133+
burst ? `inkbox:${msg.channel}_burst messages=${msg.burst}` : `inkbox:${msg.channel}`,
134+
`from=${msg.from}`,
135+
];
131136
if (msg.channel === "email") {
132137
if (msg.subject) fields.push(`subject=${JSON.stringify(msg.subject)}`);
138+
// RFC 5322 Message-ID, so the agent can thread its own follow-up sends
139+
// via inkbox_send_email's inReplyToMessageId.
140+
if (msg.rfcMessageId) fields.push(`message_id=${JSON.stringify(msg.rfcMessageId)}`);
133141
} else if (msg.conversationId) {
134142
fields.push(`conversation_id=${msg.conversationId}`);
135143
}
144+
if (msg.group?.participants?.length) {
145+
fields.push(`participants=${JSON.stringify(msg.group.participants.join(", "))}`);
146+
}
136147
// The contact card carries the addresses the agent may reach this person
137148
// at, so cross-channel follow-ups never have to guess.
138149
fields.push("|", contactCard(msg));
139150

140151
const lines = [`[${fields.join(" ")}]`];
152+
if (directive) lines.push(`Operator directive for this channel: ${directive}`);
141153
if (msg.group) lines.push(groupReminder(msg.group.participantCount));
142154
lines.push(msg.text);
143155
if (msg.mediaPaths.length > 0) lines.push(`[attached files: ${msg.mediaPaths.join(", ")}]`);

‎src/gateway/sessions.ts‎

Lines changed: 16 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@ interface QueuedTurn {
1616
kind: TurnKind;
1717
text: string;
1818
deliver: boolean;
19+
// Per-contact/per-channel opencode agent override for this turn.
20+
agent?: string;
1921
replyTarget?: ReplyTarget;
2022
// True for a follow-up turn enqueued after a delivery failure, so a second
2123
// failure doesn't spawn another recovery (bounded to one attempt).
@@ -78,13 +80,18 @@ export function createSessionManager(deps: SessionManagerDeps): SessionManager {
7880
return id;
7981
}
8082

81-
async function runPrompt(sessionID: string, text: string): Promise<string | undefined> {
83+
async function runPrompt(
84+
sessionID: string,
85+
text: string,
86+
agentOverride?: string,
87+
): Promise<string | undefined> {
8288
const g = deps.config.gateway;
89+
const agent = agentOverride ?? g.agent;
8390
const res = await deps.opencode.session.prompt({
8491
path: { id: sessionID },
8592
query: { directory: deps.directory },
8693
body: {
87-
...(g.agent ? { agent: g.agent } : {}),
94+
...(agent ? { agent } : {}),
8895
...(g.model?.includes("/")
8996
? {
9097
model: {
@@ -113,7 +120,7 @@ export function createSessionManager(deps: SessionManagerDeps): SessionManager {
113120
if (turn.kind === "normal") entry.interruptNormal = false;
114121
try {
115122
const sessionID = await ensureSession(chatKey);
116-
const out = await runPrompt(sessionID, turn.text);
123+
const out = await runPrompt(sessionID, turn.text, turn.agent);
117124
// If a newer message interrupted this normal turn, drop its output.
118125
if (turn.kind === "normal" && entry.interruptNormal) {
119126
deps.logger.info("turn.interrupted", { chatKey });
@@ -187,11 +194,16 @@ export function createSessionManager(deps: SessionManagerDeps): SessionManager {
187194
const sessionID = deps.state.getSession(msg.chatKey);
188195
if (sessionID) await interruptInFlightNormal(msg.chatKey, sessionID);
189196
}
197+
// Operator overrides, keyed by contact id first, then channel.
198+
const g = deps.config.gateway;
199+
const overrideFor = (map: Record<string, string>): string | undefined =>
200+
(msg.contactId ? map[msg.contactId] : undefined) ?? map[msg.channel];
190201
await new Promise<string | undefined>((resolve, reject) => {
191202
entry.queue.push({
192203
kind: "normal",
193-
text: frameInbound(msg),
204+
text: frameInbound(msg, overrideFor(g.channelPrompts)),
194205
deliver: true,
206+
agent: overrideFor(g.channelAgents),
195207
replyTarget,
196208
resolve,
197209
reject,

0 commit comments

Comments
 (0)