merge: bring gallery frontend branch up to date with goutamnextflow
goutamnextflow adopts @insignia/iios-messaging-ui (drops the bespoke messenger/inbox/mail) and adds Org Settings → Integrations. Resolve against the Smart Gallery work: - dashboard.tsx: route messenger→MessengerSdk, inbox→InboxSdk, add settings→ Settings (goutamnextflow) AND keep gallery→SmartGallery (ours). Drop the now-deleted Messenger/Inbox imports. - next.config.ts / package.json: union transpilePackages + deps (keep @photo-gallery/sdk + ML deps AND @insignia/iios-messaging-ui). - dashboard.css: keep both appended blocks (.gal* gallery + .settings-* org settings). Gallery integrates with zero errors. Vendored @photo-gallery/sdk re-synced (0-diff). Verified: tsc clean + `next build` compiles the whole merged app (gallery AI routes + dashboard). The new @insignia/iios-messaging-ui private dep is stubbed locally (see .claude/dev-scripts); it resolves on CI/Vercel.
This commit is contained in:
@@ -0,0 +1,121 @@
|
||||
// The CRM's InboxAdapter — the SDK <Inbox> rendered over the be-crm data door
|
||||
// (crm.inbox.* + crm.mail.*). Folds mail threads into the unified inbox exactly as the old
|
||||
// inbox-api did; the CRM keeps auth/tenancy server-side.
|
||||
|
||||
import type {
|
||||
InboxAdapter,
|
||||
InboxItem,
|
||||
InboxState,
|
||||
MailAttachment,
|
||||
MailMessage,
|
||||
MailPerson,
|
||||
} from "@insignia/iios-messaging-ui";
|
||||
import type { DataDoor } from "./crm-messaging-adapter";
|
||||
|
||||
const MAX_ATTACHMENT_BYTES = 26 * 1024 * 1024; // matches IIOS's cap
|
||||
|
||||
// Some types (notably .md) have no OS-registered MIME, so the browser reports an empty file.type.
|
||||
const EXT_MIME: Record<string, string> = {
|
||||
md: "text/markdown", markdown: "text/markdown", html: "text/html", htm: "text/html", txt: "text/plain", csv: "text/csv",
|
||||
};
|
||||
function mimeForFile(file: File): string {
|
||||
if (file.type) return file.type;
|
||||
const ext = file.name.toLowerCase().split(".").pop() ?? "";
|
||||
return EXT_MIME[ext] ?? "application/octet-stream";
|
||||
}
|
||||
|
||||
interface InboxItemDTO {
|
||||
id: string; kind: string; state: InboxState; title: string; summary?: string; priority: string; threadId?: string; createdAt: string;
|
||||
}
|
||||
interface MailThreadDTO { threadId: string; subject: string | null; participants: string[]; unread: number; lastMessage?: string; lastAt?: string }
|
||||
interface MailMessageDTO {
|
||||
interactionId: string; actorId: string | null; kind: string; occurredAt: string;
|
||||
html: string | null; text: string | null;
|
||||
attachment: { contentRef: string; mimeType: string; sizeBytes: number; filename: string | null } | null;
|
||||
}
|
||||
interface DirectoryDTO { id: string; displayName: string; kind: "staff" | "customer" }
|
||||
|
||||
const escapeHtml = (s: string): string => s.replace(/&/g, "&").replace(/</g, "<").replace(/>/g, ">");
|
||||
|
||||
export class CrmInboxAdapter implements InboxAdapter {
|
||||
constructor(private readonly sdk: DataDoor) {}
|
||||
|
||||
async listInbox(state?: InboxState): Promise<InboxItem[]> {
|
||||
const showMail = !state || state === "OPEN";
|
||||
const [items, mail] = await Promise.all([
|
||||
this.sdk.query<InboxItemDTO[]>("crm.inbox.list", state ? { state } : {}),
|
||||
showMail ? this.sdk.query<MailThreadDTO[]>("crm.mail.list", {}) : Promise.resolve([] as MailThreadDTO[]),
|
||||
]);
|
||||
const mailItems: InboxItem[] = mail.map((t) => ({
|
||||
id: `mail:${t.threadId}`,
|
||||
kind: "MAIL",
|
||||
state: "OPEN",
|
||||
title: t.subject || "(no subject)",
|
||||
...(t.lastMessage ? { summary: t.lastMessage } : {}),
|
||||
priority: t.unread > 0 ? "HIGH" : "LOW",
|
||||
threadId: t.threadId,
|
||||
createdAt: t.lastAt ?? "",
|
||||
}));
|
||||
return [...mailItems, ...items].sort((a, b) => (b.createdAt ?? "").localeCompare(a.createdAt ?? ""));
|
||||
}
|
||||
|
||||
async transition(id: string, state: InboxState): Promise<void> {
|
||||
await this.sdk.command("crm.inbox.transition", { id, state });
|
||||
}
|
||||
|
||||
async mailHistory(threadId: string): Promise<MailMessage[]> {
|
||||
const rows = await this.sdk.query<MailMessageDTO[]>("crm.mail.history", { threadId });
|
||||
return rows.map((m) => ({
|
||||
id: m.interactionId,
|
||||
actorId: m.actorId,
|
||||
kind: m.kind,
|
||||
at: m.occurredAt,
|
||||
html: m.html,
|
||||
text: m.text,
|
||||
attachment: m.attachment,
|
||||
}));
|
||||
}
|
||||
|
||||
async mailReply(threadId: string, content: string, attachment?: MailAttachment): Promise<void> {
|
||||
await this.sdk.command("crm.mail.reply", {
|
||||
threadId,
|
||||
content,
|
||||
...(attachment ? { attachment: { filename: attachment.filename ?? "attachment", contentRef: attachment.contentRef, mimeType: attachment.mimeType, sizeBytes: attachment.sizeBytes } } : {}),
|
||||
});
|
||||
}
|
||||
|
||||
async uploadAttachment(file: File): Promise<MailAttachment> {
|
||||
if (file.size > MAX_ATTACHMENT_BYTES) throw new Error("File is too large (max 25 MB).");
|
||||
const mime = mimeForFile(file);
|
||||
const { objectKey, uploadUrl } = await this.sdk.command<{ objectKey: string; uploadUrl: string }>("crm.media.presignUpload", { mime, sizeBytes: file.size });
|
||||
const res = await fetch(uploadUrl, { method: "PUT", body: file });
|
||||
if (!res.ok) throw new Error(`Upload failed (${res.status}).`);
|
||||
return { contentRef: objectKey, mimeType: mime, sizeBytes: file.size, filename: file.name };
|
||||
}
|
||||
|
||||
async downloadAttachment(attachment: MailAttachment): Promise<string> {
|
||||
const { url } = await this.sdk.command<{ url: string }>("crm.media.presignDownload", {
|
||||
contentRef: attachment.contentRef,
|
||||
...(attachment.mimeType ? { mime: attachment.mimeType } : {}),
|
||||
});
|
||||
return url;
|
||||
}
|
||||
|
||||
async directory(): Promise<MailPerson[]> {
|
||||
const rows = await this.sdk.query<DirectoryDTO[]>("crm.messenger.directory", { kind: "all", limit: 100 });
|
||||
return rows.map((d) => ({ id: d.id, name: d.displayName, kind: d.kind }));
|
||||
}
|
||||
|
||||
async composeInternal(recipientUserId: string, subject: string, text: string, attachments?: MailAttachment[]): Promise<void> {
|
||||
await this.sdk.command("crm.mail.internal", { recipientUserId, subject, text, html: `<p>${escapeHtml(text)}</p>`, ...attachmentsVar(attachments) });
|
||||
}
|
||||
|
||||
async composeExternal(target: string, subject: string, text: string, attachments?: MailAttachment[]): Promise<void> {
|
||||
await this.sdk.command("crm.mail.send", { target, subject, text, html: `<p>${escapeHtml(text)}</p>`, ...attachmentsVar(attachments) });
|
||||
}
|
||||
}
|
||||
|
||||
function attachmentsVar(attachments?: MailAttachment[]): { attachments?: Array<{ filename: string; contentRef: string; mimeType: string; sizeBytes: number }> } {
|
||||
if (!attachments || attachments.length === 0) return {};
|
||||
return { attachments: attachments.map((a) => ({ filename: a.filename ?? "attachment", contentRef: a.contentRef, mimeType: a.mimeType, sizeBytes: a.sizeBytes })) };
|
||||
}
|
||||
@@ -0,0 +1,344 @@
|
||||
// The CRM's implementation of the SDK's MessagingAdapter. HYBRID transport:
|
||||
// • BFF (appshell crm.messenger.*) for the conversation list, thread creation, and directory
|
||||
// — these need server-side tenancy/auth.
|
||||
// • IIOS MessageSocket (delegated token from crm.messenger.realtime) for everything live:
|
||||
// history+join, send, typing, read receipts, reactions.
|
||||
// When no socket is available (token failed / demo), it degrades to a 4s history poll.
|
||||
|
||||
import type {
|
||||
Attachment,
|
||||
ChannelSummary,
|
||||
ChannelVisibility,
|
||||
Conversation,
|
||||
CreateChannelInput,
|
||||
Membership,
|
||||
Message,
|
||||
MessageEvent,
|
||||
MessagingAdapter,
|
||||
Person,
|
||||
Reaction,
|
||||
SendOpts,
|
||||
Unsubscribe,
|
||||
} from "@insignia/iios-messaging-ui";
|
||||
import type { MessageSocket, Message as KernelMessage } from "@insignia/iios-kernel-client";
|
||||
|
||||
/** The imperative appshell data door (useAppShell().sdk). Typed structurally, not to its class. */
|
||||
export interface DataDoor {
|
||||
query<T>(action: string, variables?: Record<string, unknown>): Promise<T>;
|
||||
command<T>(action: string, variables?: Record<string, unknown>): Promise<T>;
|
||||
}
|
||||
|
||||
interface DirectoryDTO { id: string; displayName: string; kind: "staff" | "customer" }
|
||||
interface ConversationDTO {
|
||||
threadId: string; subject: string | null; membership: Membership | null;
|
||||
participants: string[]; unread: number; lastMessage?: string; lastAt?: string;
|
||||
}
|
||||
interface MessageDTO { interactionId: string; actorId: string | null; kind: string; occurredAt: string; text: string | null; attachment?: { contentRef: string; mimeType: string; sizeBytes: number; filename: string | null } | null }
|
||||
|
||||
const POLL_MS = 4000;
|
||||
const REACTION = "reaction";
|
||||
|
||||
interface Poll { seen: Set<string>; primed: boolean; timer: ReturnType<typeof setInterval> | null }
|
||||
|
||||
export class CrmMessagingAdapter implements MessagingAdapter {
|
||||
private names: Map<string, string> | null = null;
|
||||
private readonly listeners = new Map<string, Set<(e: MessageEvent) => void>>();
|
||||
private readonly polls = new Map<string, Poll>();
|
||||
private readonly joined = new Set<string>();
|
||||
/** messageId → emoji → userSet, so a single annotation delta can be re-emitted as a full set. */
|
||||
private readonly reactions = new Map<string, Map<string, Set<string>>>();
|
||||
|
||||
/** Only present with a socket — the UI hides the reaction affordance without it. */
|
||||
react?: (threadId: string, messageId: string, emoji: string) => Promise<void>;
|
||||
|
||||
constructor(
|
||||
private readonly sdk: DataDoor,
|
||||
private readonly me: string,
|
||||
private readonly socket?: MessageSocket,
|
||||
) {
|
||||
if (socket) {
|
||||
socket.on("message", (m) => {
|
||||
this.ingestReactions(m);
|
||||
void this.toKernelMessage(m).then((message) => this.emit(m.threadId, { kind: "message", message }));
|
||||
});
|
||||
socket.on("typing", (e) => this.emit(e.threadId, { kind: "typing", userId: e.userId }));
|
||||
// Receipts carry no threadId → fan to all open threads; the UI filters by messageId.
|
||||
socket.on("receipt", (e) => this.broadcast({ kind: "receipt", messageId: e.interactionId, actorId: e.actorId }));
|
||||
socket.on("annotation", (e) => {
|
||||
if (e.type !== REACTION) return;
|
||||
this.setReactionUsers(e.interactionId, e.value, e.users);
|
||||
this.emit(e.threadId, { kind: "reaction", messageId: e.interactionId, reactions: this.reactionsOf(e.interactionId) });
|
||||
});
|
||||
this.react = async (threadId, messageId, emoji) => {
|
||||
await socket.react(threadId, messageId, emoji);
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
currentActorId(): string {
|
||||
return this.me;
|
||||
}
|
||||
|
||||
async listConversations(): Promise<Conversation[]> {
|
||||
const [convs, names] = await Promise.all([
|
||||
this.sdk.query<ConversationDTO[]>("crm.messenger.conversation.list", {}),
|
||||
this.directoryMap(),
|
||||
]);
|
||||
return convs.map((c) => this.toConversation(c, names));
|
||||
}
|
||||
|
||||
async openThread(p: { participantIds: string[]; membership?: Membership; subject?: string }): Promise<{ threadId: string }> {
|
||||
const res = await this.sdk.command<{ threadId: string }>("crm.messenger.conversation.open", {
|
||||
participantIds: p.participantIds,
|
||||
...(p.membership ? { membership: p.membership } : {}),
|
||||
...(p.subject ? { subject: p.subject } : {}),
|
||||
});
|
||||
return { threadId: res.threadId };
|
||||
}
|
||||
|
||||
async history(threadId: string): Promise<Message[]> {
|
||||
if (this.socket) {
|
||||
const res = await this.socket.openThread(threadId); // joins so live events flow
|
||||
this.joined.add(threadId);
|
||||
return Promise.all(
|
||||
res.history.map((m) => {
|
||||
this.ingestReactions(m);
|
||||
return this.toKernelMessage(m);
|
||||
}),
|
||||
);
|
||||
}
|
||||
const msgs = await this.sdk.query<MessageDTO[]>("crm.messenger.history", { threadId });
|
||||
return Promise.all(msgs.map((m) => this.toDtoMessage(m)));
|
||||
}
|
||||
|
||||
async send(threadId: string, content: string, opts?: SendOpts): Promise<Message> {
|
||||
const att = opts?.attachment;
|
||||
if (this.socket) {
|
||||
const sendOpts = {
|
||||
...(opts?.parentInteractionId ? { parentInteractionId: opts.parentInteractionId } : {}),
|
||||
...(opts?.mentions && opts.mentions.length ? { mentions: opts.mentions } : {}),
|
||||
...(att?.contentRef ? { attachment: { contentRef: att.contentRef, mimeType: att.mime, sizeBytes: att.sizeBytes ?? 0 } } : {}),
|
||||
};
|
||||
const m = await this.socket.sendMessage(threadId, content, Object.keys(sendOpts).length ? sendOpts : undefined);
|
||||
const msg = this.fromKernel(m);
|
||||
// Reuse the staged attachment (already carries a display URL from upload) for instant render.
|
||||
return att ? { ...msg, attachment: att } : msg;
|
||||
}
|
||||
const m = await this.sdk.command<MessageDTO>("crm.messenger.send", { threadId, content });
|
||||
const msg = this.fromDto(m);
|
||||
this.polls.get(threadId)?.seen.add(msg.id);
|
||||
return att ? { ...msg, attachment: att } : msg;
|
||||
}
|
||||
|
||||
async upload(file: File): Promise<Attachment> {
|
||||
const mime = file.type || "application/octet-stream";
|
||||
const { objectKey, uploadUrl } = await this.sdk.command<{ objectKey: string; uploadUrl: string }>("crm.media.presignUpload", { mime, sizeBytes: file.size });
|
||||
const res = await fetch(uploadUrl, { method: "PUT", body: file });
|
||||
if (!res.ok) throw new Error(`Upload failed (${res.status}).`);
|
||||
const url = await this.downloadUrl(objectKey, mime);
|
||||
return { url, mime, name: file.name, contentRef: objectKey, sizeBytes: file.size };
|
||||
}
|
||||
|
||||
subscribe(threadId: string, cb: (e: MessageEvent) => void): Unsubscribe {
|
||||
if (!this.listeners.has(threadId)) this.listeners.set(threadId, new Set());
|
||||
this.listeners.get(threadId)!.add(cb);
|
||||
|
||||
if (this.socket) {
|
||||
if (!this.joined.has(threadId)) {
|
||||
this.joined.add(threadId);
|
||||
void this.socket.openThread(threadId).catch(() => this.joined.delete(threadId));
|
||||
}
|
||||
} else {
|
||||
this.startPoll(threadId);
|
||||
}
|
||||
|
||||
return () => {
|
||||
const set = this.listeners.get(threadId);
|
||||
set?.delete(cb);
|
||||
if (set && set.size === 0) {
|
||||
this.listeners.delete(threadId);
|
||||
const poll = this.polls.get(threadId);
|
||||
if (poll?.timer) clearInterval(poll.timer);
|
||||
this.polls.delete(threadId);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
sendTyping(threadId: string): void {
|
||||
this.socket?.typing(threadId);
|
||||
}
|
||||
|
||||
async markRead(threadId: string, messageId: string): Promise<void> {
|
||||
if (this.socket) await this.socket.markRead(threadId, messageId);
|
||||
}
|
||||
|
||||
// ── channels + members (BFF, except join which is a governed socket self-join) ──
|
||||
async browseChannels(): Promise<ChannelSummary[]> {
|
||||
const rows = await this.sdk.query<Array<{ threadId: string; name: string; topic: string | null; visibility: string; memberCount: number; joined: boolean }>>(
|
||||
"crm.messenger.channel.browse",
|
||||
{},
|
||||
);
|
||||
return rows.map((c) => ({
|
||||
threadId: c.threadId,
|
||||
name: c.name,
|
||||
topic: c.topic,
|
||||
visibility: (c.visibility === "private" ? "private" : "public") as ChannelVisibility,
|
||||
memberCount: c.memberCount,
|
||||
joined: c.joined,
|
||||
}));
|
||||
}
|
||||
|
||||
async createChannel(input: CreateChannelInput): Promise<{ threadId: string }> {
|
||||
return this.sdk.command<{ threadId: string }>("crm.messenger.channel.create", {
|
||||
name: input.name,
|
||||
...(input.topic ? { topic: input.topic } : {}),
|
||||
visibility: input.visibility,
|
||||
});
|
||||
}
|
||||
|
||||
async joinChannel(threadId: string): Promise<void> {
|
||||
// Governed public self-join over the socket (the BFF has no join verb; OPA enforces it).
|
||||
if (!this.socket) throw new Error("joining a channel needs a live connection");
|
||||
await this.socket.openThread(threadId);
|
||||
this.joined.add(threadId);
|
||||
}
|
||||
|
||||
async leaveChannel(threadId: string): Promise<void> {
|
||||
await this.sdk.command("crm.messenger.channel.leave", { threadId });
|
||||
}
|
||||
|
||||
async listMembers(threadId: string): Promise<Person[]> {
|
||||
const rows = await this.sdk.query<Array<{ userId: string; displayName: string; role: string }>>("crm.messenger.members", { threadId });
|
||||
return rows.map((r) => ({ id: r.userId, name: r.displayName, kind: "staff" as const }));
|
||||
}
|
||||
|
||||
// ── polling fallback (no socket) ───────────────────────────────
|
||||
private startPoll(threadId: string): void {
|
||||
if (this.polls.has(threadId)) return;
|
||||
const poll: Poll = { seen: new Set(), primed: false, timer: null };
|
||||
this.polls.set(threadId, poll);
|
||||
const tick = async (): Promise<void> => {
|
||||
if (!this.polls.has(threadId)) return;
|
||||
try {
|
||||
const msgs = await this.sdk.query<MessageDTO[]>("crm.messenger.history", { threadId });
|
||||
for (const m of msgs) {
|
||||
if (poll.seen.has(m.interactionId)) continue;
|
||||
poll.seen.add(m.interactionId);
|
||||
if (poll.primed) this.emit(threadId, { kind: "message", message: this.fromDto(m) });
|
||||
}
|
||||
poll.primed = true;
|
||||
} catch {
|
||||
/* transient — retry next tick */
|
||||
}
|
||||
};
|
||||
void tick();
|
||||
poll.timer = setInterval(tick, POLL_MS);
|
||||
}
|
||||
|
||||
/** The org directory — people you can start a DM/group with. Drives the "New message" picker. */
|
||||
async directory(): Promise<Person[]> {
|
||||
const dir = await this.sdk.query<DirectoryDTO[]>("crm.messenger.directory", { kind: "all", limit: 100 });
|
||||
return dir.map((d) => ({ id: d.id, name: d.displayName, kind: d.kind }));
|
||||
}
|
||||
|
||||
// ── mapping ────────────────────────────────────────────────────
|
||||
private async directoryMap(): Promise<Map<string, string>> {
|
||||
if (!this.names) {
|
||||
this.names = new Map((await this.directory()).map((p) => [p.id, p.name]));
|
||||
}
|
||||
return this.names;
|
||||
}
|
||||
|
||||
private toConversation(c: ConversationDTO, names: Map<string, string>): Conversation {
|
||||
const others = c.participants.filter((p) => p !== this.me);
|
||||
const title = c.subject?.trim() || others.map((id) => names.get(id) ?? id).join(", ") || "Conversation";
|
||||
return {
|
||||
threadId: c.threadId,
|
||||
title,
|
||||
subject: c.subject,
|
||||
membership: c.membership,
|
||||
participants: c.participants,
|
||||
unread: c.unread,
|
||||
...(c.lastMessage ? { lastMessage: c.lastMessage } : {}),
|
||||
...(c.lastAt ? { lastAt: c.lastAt } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
/** Kernel Message (socket) → SDK Message. actorId = senderId (userId space), matching currentActorId. */
|
||||
private fromKernel(m: KernelMessage): Message {
|
||||
return {
|
||||
id: m.id,
|
||||
actorId: m.senderId ?? null,
|
||||
text: m.content ?? "",
|
||||
at: m.createdAt,
|
||||
parentInteractionId: m.parentInteractionId ?? null,
|
||||
reactions: this.reactionsOf(m.id),
|
||||
};
|
||||
}
|
||||
|
||||
/** BFF DTO (poll fallback) → SDK Message. Note: actorId is IIOS actor-id space here. */
|
||||
private fromDto(m: MessageDTO): Message {
|
||||
return { id: m.interactionId, actorId: m.actorId, text: m.text ?? "", at: m.occurredAt };
|
||||
}
|
||||
|
||||
// ── attachments ────────────────────────────────────────────────
|
||||
/** A short-lived signed URL to display/download a stored object. */
|
||||
private async downloadUrl(contentRef: string, mime?: string): Promise<string> {
|
||||
const { url } = await this.sdk.command<{ url: string }>("crm.media.presignDownload", { contentRef, ...(mime ? { mime } : {}) });
|
||||
return url;
|
||||
}
|
||||
|
||||
private async resolveAttachment(a: { contentRef: string; mimeType: string; sizeBytes: number; filename?: string | null } | null | undefined): Promise<Attachment | undefined> {
|
||||
if (!a?.contentRef) return undefined;
|
||||
const url = await this.downloadUrl(a.contentRef, a.mimeType);
|
||||
return { url, mime: a.mimeType, name: a.filename ?? "attachment", contentRef: a.contentRef, sizeBytes: a.sizeBytes };
|
||||
}
|
||||
|
||||
private async toKernelMessage(m: KernelMessage): Promise<Message> {
|
||||
const base = this.fromKernel(m);
|
||||
const att = await this.resolveAttachment(m.attachment ?? null);
|
||||
return att ? { ...base, attachment: att } : base;
|
||||
}
|
||||
|
||||
private async toDtoMessage(m: MessageDTO): Promise<Message> {
|
||||
const base = this.fromDto(m);
|
||||
const att = await this.resolveAttachment(m.attachment ?? null);
|
||||
return att ? { ...base, attachment: att } : base;
|
||||
}
|
||||
|
||||
// ── reaction state ─────────────────────────────────────────────
|
||||
private ingestReactions(m: KernelMessage): void {
|
||||
for (const a of m.annotations ?? []) {
|
||||
if (a.type === REACTION) this.setReactionUsers(m.id, a.value, a.users);
|
||||
}
|
||||
}
|
||||
|
||||
private setReactionUsers(messageId: string, emoji: string, users: string[]): void {
|
||||
let byEmoji = this.reactions.get(messageId);
|
||||
if (!byEmoji) {
|
||||
byEmoji = new Map();
|
||||
this.reactions.set(messageId, byEmoji);
|
||||
}
|
||||
if (users.length === 0) byEmoji.delete(emoji);
|
||||
else byEmoji.set(emoji, new Set(users));
|
||||
}
|
||||
|
||||
private reactionsOf(messageId: string): Reaction[] {
|
||||
const byEmoji = this.reactions.get(messageId);
|
||||
if (!byEmoji) return [];
|
||||
const out: Reaction[] = [];
|
||||
for (const [emoji, users] of byEmoji) {
|
||||
if (users.size > 0) out.push({ emoji, count: users.size, mine: users.has(this.me) });
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
// ── event fan-out ──────────────────────────────────────────────
|
||||
private emit(threadId: string, e: MessageEvent): void {
|
||||
this.listeners.get(threadId)?.forEach((cb) => cb(e));
|
||||
}
|
||||
|
||||
private broadcast(e: MessageEvent): void {
|
||||
for (const set of this.listeners.values()) set.forEach((cb) => cb(e));
|
||||
}
|
||||
}
|
||||
@@ -1,95 +0,0 @@
|
||||
"use client";
|
||||
|
||||
// Inbox data layer. The inbox is a personalized work/awareness feed IIOS projects from events
|
||||
// (NEEDS_REPLY, MENTION, …). The CRM lists it and transitions item state; items are never created
|
||||
// here. Mock when the Shell isn't configured; live via the be-crm data door (crm.inbox.*) otherwise.
|
||||
//
|
||||
// Live contract (be-crm):
|
||||
// query crm.inbox.list { state? } -> InboxItem[]
|
||||
// cmd crm.inbox.transition { id, state, reason? } -> InboxItem
|
||||
|
||||
import { useCallback, useEffect, useMemo, useState } from "react";
|
||||
import { useAppShell, useQuery } from "@abe-kap/appshell-sdk/react";
|
||||
import { isShellConfigured } from "./appshell";
|
||||
import type { MailThread } from "./mail-api";
|
||||
|
||||
export type InboxState = "OPEN" | "SNOOZED" | "DONE" | "ARCHIVED" | "CANCELLED" | "STALE";
|
||||
export interface UiInboxItem {
|
||||
id: string; kind: string; state: InboxState; title: string; summary?: string;
|
||||
priority: string; threadId?: string; createdAt: string;
|
||||
}
|
||||
|
||||
const SHELL = isShellConfigured();
|
||||
|
||||
export interface InboxData {
|
||||
live: boolean; loading: boolean; error: string | null;
|
||||
items: UiInboxItem[];
|
||||
transition: (id: string, state: InboxState) => Promise<void>;
|
||||
refetch: () => void;
|
||||
}
|
||||
|
||||
export function useInboxData(state?: InboxState): InboxData {
|
||||
return SHELL ? useLiveInbox(state) : useMockInbox(state);
|
||||
}
|
||||
|
||||
function useLiveInbox(state?: InboxState): InboxData {
|
||||
const { sdk } = useAppShell();
|
||||
const q = useQuery<UiInboxItem[]>("crm.inbox.list", state ? { state } : {});
|
||||
// Mail lives in crm-mail threads, NOT the inbox projection — fold it into the one unified
|
||||
// surface. Mail has no inbox work-item state, so it only shows in the Open (or unfiltered) view.
|
||||
const showMail = !state || state === "OPEN";
|
||||
const mq = useQuery<MailThread[]>("crm.mail.list", {});
|
||||
|
||||
// The SDK's useQuery only refetches when the ACTION changes, not the variables — so a filter
|
||||
// change (same action, new { state }) wouldn't reload. Force a refetch when the filter changes.
|
||||
const refetchInbox = q.refetch;
|
||||
useEffect(() => { refetchInbox(); }, [state, refetchInbox]);
|
||||
|
||||
const items = useMemo<UiInboxItem[]>(() => {
|
||||
const inboxItems = q.data ?? [];
|
||||
const mailItems: UiInboxItem[] = showMail
|
||||
? (mq.data ?? []).map((t) => ({
|
||||
id: `mail:${t.threadId}`,
|
||||
kind: "MAIL",
|
||||
state: "OPEN" as InboxState,
|
||||
title: t.subject || "(no subject)",
|
||||
...(t.lastMessage ? { summary: t.lastMessage } : {}),
|
||||
priority: t.unread > 0 ? "HIGH" : "LOW",
|
||||
threadId: t.threadId,
|
||||
createdAt: t.lastAt ?? "",
|
||||
}))
|
||||
: [];
|
||||
// Newest first; mail and inbox items interleave by time.
|
||||
return [...mailItems, ...inboxItems].sort((a, b) => (b.createdAt ?? "").localeCompare(a.createdAt ?? ""));
|
||||
}, [q.data, mq.data, showMail]);
|
||||
|
||||
const transition = useCallback(async (id: string, next: InboxState) => {
|
||||
await sdk.command("crm.inbox.transition", { id, state: next });
|
||||
q.refetch();
|
||||
}, [sdk, q]);
|
||||
|
||||
return {
|
||||
live: true,
|
||||
loading: q.loading || (showMail && mq.loading),
|
||||
// Don't let a mail-list hiccup blank the whole inbox — surface only the inbox error.
|
||||
error: q.error?.message ?? null,
|
||||
items,
|
||||
transition,
|
||||
refetch: () => { q.refetch(); mq.refetch(); },
|
||||
};
|
||||
}
|
||||
|
||||
const MOCK_ITEMS: UiInboxItem[] = [
|
||||
{ id: "in_1", kind: "MENTION", state: "OPEN", title: "Sofia mentioned you", summary: "@you — can you confirm the Henderson scope?", priority: "HIGH", threadId: "th_mock_1", createdAt: new Date().toISOString() },
|
||||
{ id: "in_2", kind: "NEEDS_REPLY", state: "OPEN", title: "Reply needed — Storm response", summary: "Dan: Crew is rolling out at 7.", priority: "MEDIUM", threadId: "th_mock_2", createdAt: new Date().toISOString() },
|
||||
{ id: "in_3", kind: "SUPPORT_UPDATE", state: "OPEN", title: "Ticket TK-204 updated", summary: "Customer replied on the roof-leak case.", priority: "LOW", createdAt: new Date().toISOString() },
|
||||
];
|
||||
|
||||
function useMockInbox(state?: InboxState): InboxData {
|
||||
const [items, setItems] = useState<UiInboxItem[]>(MOCK_ITEMS);
|
||||
const filtered = useMemo(() => (state ? items.filter((i) => i.state === state) : items), [items, state]);
|
||||
const transition = useCallback(async (id: string, next: InboxState) => {
|
||||
setItems((l) => l.map((i) => (i.id === id ? { ...i, state: next } : i)));
|
||||
}, []);
|
||||
return { live: false, loading: false, error: null, items: filtered, transition, refetch: () => {} };
|
||||
}
|
||||
@@ -1,121 +0,0 @@
|
||||
"use client";
|
||||
|
||||
// Mail data layer. A dedicated Mail reader over the be-crm data door (crm.mail.*), distinct from
|
||||
// the Messenger chat and from the work-item Inbox. Live via the AppShell SDK; a small mock keeps the
|
||||
// demo working before the Shell + be-crm are connected.
|
||||
//
|
||||
// Live contract (be-crm):
|
||||
// query crm.mail.list {} -> MailThread[]
|
||||
// query crm.mail.history { threadId } -> MailMessage[]
|
||||
// cmd crm.mail.reply { threadId, content } -> { interactionId, threadId }
|
||||
// cmd crm.mail.internal { recipientUserId, subject?, text?, html? } -> { threadId }
|
||||
// cmd crm.mail.send { target, subject?, text?, html?, mirrorToUserId? } -> { commandId }
|
||||
// query crm.messenger.directory { kind, limit } -> people to compose to (reused)
|
||||
|
||||
import { useCallback, useMemo, useState } from "react";
|
||||
import { useAppShell, useQuery } from "@abe-kap/appshell-sdk/react";
|
||||
import { isShellConfigured } from "./appshell";
|
||||
|
||||
export interface MailThread {
|
||||
threadId: string; subject: string | null; participants: string[]; unread: number; lastMessage?: string; lastAt?: string;
|
||||
}
|
||||
export interface MailAttachment { contentRef: string; mimeType: string; sizeBytes: number; filename: string | null }
|
||||
export interface MailMessage {
|
||||
interactionId: string; actorId: string | null; kind: string; occurredAt: string; html: string | null; text: string | null; attachment: MailAttachment | null;
|
||||
}
|
||||
export interface MailPerson { id: string; name: string; kind: "staff" | "customer" }
|
||||
|
||||
/** Shape produced by media-api's useUploadAttachment, passed into a reply/compose. */
|
||||
export interface OutgoingAttachment { contentRef: string; mimeType: string; sizeBytes: number; filename: string }
|
||||
|
||||
const SHELL = isShellConfigured();
|
||||
|
||||
/* ============================ Thread list ============================ */
|
||||
|
||||
export interface MailListData {
|
||||
live: boolean; loading: boolean; error: string | null; threads: MailThread[]; refetch: () => void;
|
||||
}
|
||||
|
||||
export function useMailThreads(): MailListData {
|
||||
if (SHELL) {
|
||||
const q = useQuery<MailThread[]>("crm.mail.list", {});
|
||||
return { live: true, loading: q.loading, error: q.error?.message ?? null, threads: q.data ?? [], refetch: q.refetch };
|
||||
}
|
||||
return { live: false, loading: false, error: null, threads: MOCK_THREADS, refetch: () => {} };
|
||||
}
|
||||
|
||||
/* ============================ One thread ============================ */
|
||||
|
||||
export interface MailThreadData {
|
||||
loading: boolean; error: string | null; messages: MailMessage[]; reply: (content: string, attachment?: OutgoingAttachment) => Promise<void>; refetch: () => void;
|
||||
}
|
||||
|
||||
export function useMailThread(threadId: string | null): MailThreadData {
|
||||
if (SHELL) return useLiveThread(threadId);
|
||||
return useMockThread(threadId);
|
||||
}
|
||||
|
||||
function useLiveThread(threadId: string | null): MailThreadData {
|
||||
const { sdk } = useAppShell();
|
||||
const q = useQuery<MailMessage[]>("crm.mail.history", threadId ? { threadId } : { threadId: "" });
|
||||
const reply = useCallback(async (content: string, attachment?: OutgoingAttachment) => {
|
||||
if (!threadId) return;
|
||||
await sdk.command("crm.mail.reply", { threadId, content, ...(attachment ? { attachment } : {}) });
|
||||
q.refetch();
|
||||
}, [sdk, threadId, q]);
|
||||
return { loading: q.loading, error: q.error?.message ?? null, messages: threadId ? (q.data ?? []) : [], reply, refetch: q.refetch };
|
||||
}
|
||||
|
||||
/* ============================ Compose ============================ */
|
||||
|
||||
export interface ComposeData {
|
||||
directory: MailPerson[];
|
||||
sendInternal: (recipientUserId: string, subject: string, text: string, attachments?: OutgoingAttachment[]) => Promise<void>;
|
||||
sendExternal: (target: string, subject: string, text: string, opts?: { mirrorToUserId?: string; attachments?: OutgoingAttachment[] }) => Promise<void>;
|
||||
}
|
||||
|
||||
export function useMailCompose(onSent: () => void): ComposeData {
|
||||
if (SHELL) {
|
||||
const { sdk } = useAppShell();
|
||||
const dirQ = useQuery<MailPerson[]>("crm.messenger.directory", { kind: "all", limit: 100 });
|
||||
const directory = useMemo(() => (dirQ.data ?? []).map((d) => ({ id: (d as unknown as { id: string }).id, name: (d as unknown as { displayName?: string; name?: string }).displayName ?? (d as unknown as { name?: string }).name ?? "", kind: (d as MailPerson).kind })), [dirQ.data]);
|
||||
const sendInternal = useCallback(async (recipientUserId: string, subject: string, text: string, attachments?: OutgoingAttachment[]) => {
|
||||
await sdk.command("crm.mail.internal", { recipientUserId, subject, text, html: `<p>${escapeHtml(text)}</p>`, ...(attachments && attachments.length ? { attachments } : {}) });
|
||||
onSent();
|
||||
}, [sdk, onSent]);
|
||||
const sendExternal = useCallback(async (target: string, subject: string, text: string, opts?: { mirrorToUserId?: string; attachments?: OutgoingAttachment[] }) => {
|
||||
await sdk.command("crm.mail.send", { target, subject, text, html: `<p>${escapeHtml(text)}</p>`, ...(opts?.mirrorToUserId ? { mirrorToUserId: opts.mirrorToUserId } : {}), ...(opts?.attachments && opts.attachments.length ? { attachments: opts.attachments } : {}) });
|
||||
onSent();
|
||||
}, [sdk, onSent]);
|
||||
return { directory, sendInternal, sendExternal };
|
||||
}
|
||||
return { directory: MOCK_PEOPLE, sendInternal: async () => onSent(), sendExternal: async () => onSent() };
|
||||
}
|
||||
|
||||
function escapeHtml(s: string): string {
|
||||
return s.replace(/&/g, "&").replace(/</g, "<").replace(/>/g, ">");
|
||||
}
|
||||
|
||||
/* ============================ Mock (demo mode) ============================ */
|
||||
|
||||
const now = () => new Date().toISOString();
|
||||
const MOCK_PEOPLE: MailPerson[] = [
|
||||
{ id: "pp_sofia", name: "Sofia Ramirez", kind: "staff" },
|
||||
{ id: "cust_acme", name: "Acme Roofing (Client)", kind: "customer" },
|
||||
];
|
||||
const MOCK_THREADS: MailThread[] = [
|
||||
{ threadId: "mt_1", subject: "Welcome to the Founders Club", participants: ["you", "system"], unread: 1, lastMessage: "Thanks for joining…", lastAt: now() },
|
||||
{ threadId: "mt_2", subject: "Storm response — East side", participants: ["you", "pp_sofia"], unread: 0, lastMessage: "Crew rolling out at 7", lastAt: now() },
|
||||
];
|
||||
function useMockThread(threadId: string | null): MailThreadData {
|
||||
const [extra, setExtra] = useState<MailMessage[]>([]);
|
||||
const base: MailMessage[] = threadId === "mt_1"
|
||||
? [{ interactionId: "m1", actorId: "system", kind: "EMAIL", occurredAt: now(), html: "<p>Thanks for joining the <b>Founders Club</b>. Set up your account to get started.</p>", text: "Thanks for joining the Founders Club.", attachment: null }]
|
||||
: threadId === "mt_2"
|
||||
? [{ interactionId: "m2", actorId: "pp_sofia", kind: "EMAIL", occurredAt: now(), html: "<p>Crew is rolling out at 7. Confirm the Henderson scope?</p>", text: "Crew rolling out at 7.", attachment: null }]
|
||||
: [];
|
||||
const reply = useCallback(async (content: string, attachment?: OutgoingAttachment) => {
|
||||
setExtra((l) => [...l, { interactionId: `r_${l.length}`, actorId: "you", kind: "MESSAGE", occurredAt: now(), html: null, text: content, attachment: attachment ? { contentRef: attachment.contentRef, mimeType: attachment.mimeType, sizeBytes: attachment.sizeBytes, filename: attachment.filename } : null }]);
|
||||
}, []);
|
||||
return { loading: false, error: null, messages: threadId ? [...base, ...extra] : [], reply, refetch: () => {} };
|
||||
}
|
||||
+14
-1
@@ -14,12 +14,25 @@ export function isImage(mime?: string | null): boolean {
|
||||
return !!mime && mime.startsWith("image/");
|
||||
}
|
||||
|
||||
// Some types (notably .md) have no OS-registered MIME, so the browser reports an empty file.type.
|
||||
// Fall back to the extension for the text types IIOS allows, else a generic binary.
|
||||
const EXT_MIME: Record<string, string> = {
|
||||
md: "text/markdown", markdown: "text/markdown",
|
||||
html: "text/html", htm: "text/html",
|
||||
txt: "text/plain", csv: "text/csv",
|
||||
};
|
||||
function mimeForFile(file: File): string {
|
||||
if (file.type) return file.type;
|
||||
const ext = file.name.toLowerCase().split(".").pop() ?? "";
|
||||
return EXT_MIME[ext] ?? "application/octet-stream";
|
||||
}
|
||||
|
||||
/** Upload a File → { contentRef, mimeType, sizeBytes, filename }. Throws on oversize / failure. */
|
||||
export function useUploadAttachment() {
|
||||
const { sdk } = useAppShell();
|
||||
return useCallback(async (file: File): Promise<UploadedAttachment> => {
|
||||
if (file.size > MAX_ATTACHMENT_BYTES) throw new Error("File is too large (max 25 MB).");
|
||||
const mime = file.type || "application/octet-stream";
|
||||
const mime = mimeForFile(file);
|
||||
const { objectKey, uploadUrl } = (await sdk.command("crm.media.presignUpload", { mime, sizeBytes: file.size })) as { objectKey: string; uploadUrl: string };
|
||||
const res = await fetch(uploadUrl, { method: "PUT", body: file });
|
||||
if (!res.ok) throw new Error(`Upload failed (${res.status}).`);
|
||||
|
||||
@@ -1,435 +0,0 @@
|
||||
"use client";
|
||||
|
||||
// Messenger data layer. Serves EITHER a local mock (when the Shell isn't configured — the demo
|
||||
// keeps working) OR the live be-crm data door (crm.messenger.*), behind one interface so the UI is
|
||||
// mode-agnostic. DM-vs-group + who-can-chat are enforced server-side by IIOS/OPA; this is just glue.
|
||||
//
|
||||
// Live contract (be-crm):
|
||||
// query crm.messenger.directory { kind, query?, limit } -> DirectoryEntry[]
|
||||
// query crm.messenger.conversation.list {} -> ConversationSummary[]
|
||||
// cmd crm.messenger.conversation.open { participantIds[], membership?, subject? } -> { threadId, ... }
|
||||
// query crm.messenger.history { threadId } -> MessengerMessage[]
|
||||
// cmd crm.messenger.send { threadId, content } -> MessengerMessage
|
||||
// cmd crm.messenger.participant.add { threadId, userId }
|
||||
//
|
||||
// v1 uses REST + polling for the live stream; v2 layers the IIOS MessageSocket (messenger-socket.tsx)
|
||||
// on top for live messages, typing, read receipts, and reactions.
|
||||
|
||||
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
|
||||
import { useAppShell, useAuth, useQuery } from "@abe-kap/appshell-sdk/react";
|
||||
import type { AnnotationEvent, AnnotationGroup } from "@insignia/iios-kernel-client";
|
||||
import { isShellConfigured } from "./appshell";
|
||||
import { useMessengerSocket } from "./messenger-socket";
|
||||
|
||||
export type Membership = "dm" | "group";
|
||||
export interface UiPerson { id: string; name: string; kind: "staff" | "customer" }
|
||||
export interface UiConversation {
|
||||
threadId: string; title: string; subject: string | null; membership: Membership | null;
|
||||
participants: string[]; unread: number; lastMessage?: string; lastAt?: string;
|
||||
}
|
||||
export interface UiReaction { emoji: string; count: number; mine: boolean }
|
||||
export interface UiAttachment { contentRef: string; mimeType: string; sizeBytes: number }
|
||||
export interface UiMessage {
|
||||
id: string; actorId: string | null; senderId?: string | null; text: string; at: string; mine: boolean;
|
||||
parentInteractionId?: string | null;
|
||||
attachment?: UiAttachment;
|
||||
reactions?: UiReaction[];
|
||||
}
|
||||
|
||||
interface DirectoryDTO { id: string; displayName: string; kind: "staff" | "customer" }
|
||||
interface ConversationDTO {
|
||||
threadId: string; subject: string | null; membership: Membership | null;
|
||||
participants: string[]; unread: number; lastMessage?: string; lastAt?: string;
|
||||
}
|
||||
interface MessageDTO { interactionId: string; actorId: string | null; kind: string; occurredAt: string; text: string | null }
|
||||
|
||||
const SHELL = isShellConfigured();
|
||||
const POLL_MS = 4000;
|
||||
const TYPING_TTL_MS = 3500;
|
||||
|
||||
const shortId = (id: string) => id.replace(/^(pp_|cust_)/, "").slice(0, 6);
|
||||
|
||||
/** Turn the kernel's generic annotation aggregates into reaction chips. `users` may hold user or
|
||||
* actor ids depending on the source, so `mine` is best-effort; a fresh annotation event corrects it. */
|
||||
export function toReactions(annotations: AnnotationGroup[] | undefined, myId?: string): UiReaction[] {
|
||||
if (!annotations) return [];
|
||||
return annotations
|
||||
.filter((a) => a.type === "reaction" && a.users.length > 0)
|
||||
.map((a) => ({ emoji: a.value, count: a.users.length, mine: !!myId && a.users.includes(myId) }));
|
||||
}
|
||||
|
||||
function applyAnnotation(prev: UiReaction[] | undefined, e: AnnotationEvent, myId?: string): UiReaction[] {
|
||||
const base = (prev ?? []).filter((r) => r.emoji !== e.value);
|
||||
if (e.type !== "reaction" || e.users.length === 0) return base;
|
||||
return [...base, { emoji: e.value, count: e.users.length, mine: !!myId && e.users.includes(myId) }];
|
||||
}
|
||||
|
||||
/* ======================================================================== */
|
||||
/* Public hooks */
|
||||
/* ======================================================================== */
|
||||
|
||||
export interface MessengerData {
|
||||
live: boolean; loading: boolean; error: string | null;
|
||||
directory: UiPerson[];
|
||||
conversations: UiConversation[];
|
||||
nameOf: (id: string) => string;
|
||||
openConversation: (participantIds: string[], opts?: { membership?: Membership; subject?: string }) => Promise<string>;
|
||||
refetch: () => void;
|
||||
}
|
||||
|
||||
export interface ThreadData {
|
||||
loading: boolean; error: string | null;
|
||||
messages: UiMessage[];
|
||||
send: (content: string, opts?: { parentInteractionId?: string; attachment?: UiAttachment }) => Promise<void>;
|
||||
react: (interactionId: string, emoji: string) => void;
|
||||
typingUserIds: string[];
|
||||
seenIds: Set<string>;
|
||||
refetch: () => void;
|
||||
}
|
||||
|
||||
export interface UiMember { userId: string; displayName: string; role: string }
|
||||
export interface GroupSettingsData {
|
||||
loading: boolean; error: string | null;
|
||||
members: UiMember[];
|
||||
isAdmin: boolean;
|
||||
rename: (subject: string) => Promise<void>;
|
||||
addMember: (userId: string) => Promise<void>;
|
||||
removeMember: (userId: string) => Promise<void>;
|
||||
refetch: () => void;
|
||||
}
|
||||
|
||||
export function useMessengerData(): MessengerData {
|
||||
return SHELL ? useLiveMessenger() : useMockMessenger();
|
||||
}
|
||||
export function useThread(threadId: string): ThreadData {
|
||||
return SHELL ? useLiveThread(threadId) : useMockThread(threadId);
|
||||
}
|
||||
export function useGroupSettings(threadId: string): GroupSettingsData {
|
||||
return SHELL ? useLiveGroupSettings(threadId) : useMockGroupSettings(threadId);
|
||||
}
|
||||
|
||||
/* ======================================================================== */
|
||||
/* Live implementation (be-crm data door + IIOS socket) */
|
||||
/* ======================================================================== */
|
||||
|
||||
function useLiveMessenger(): MessengerData {
|
||||
const { sdk } = useAppShell();
|
||||
const { user } = useAuth();
|
||||
const socket = useMessengerSocket();
|
||||
const myId = user?.id;
|
||||
const dirQ = useQuery<DirectoryDTO[]>("crm.messenger.directory", { kind: "all", limit: 100 });
|
||||
const convQ = useQuery<ConversationDTO[]>("crm.messenger.conversation.list", {});
|
||||
|
||||
const directory: UiPerson[] = useMemo(
|
||||
() => (dirQ.data ?? []).map((d) => ({ id: d.id, name: d.displayName, kind: d.kind })),
|
||||
[dirQ.data],
|
||||
);
|
||||
const nameById = useMemo(() => Object.fromEntries(directory.map((p) => [p.id, p.name])), [directory]);
|
||||
const nameOf = useCallback((id: string) => nameById[id] ?? `User ${shortId(id)}`, [nameById]);
|
||||
|
||||
// Live sidebar previews: patch lastMessage/lastAt the instant a message arrives on any thread,
|
||||
// then reconcile authoritative unread/order with a debounced refetch.
|
||||
const [previews, setPreviews] = useState<Record<string, { lastMessage: string; lastAt: string }>>({});
|
||||
const refetchRef = useRef(convQ.refetch);
|
||||
refetchRef.current = convQ.refetch;
|
||||
useEffect(() => {
|
||||
if (!socket) return;
|
||||
let timer: ReturnType<typeof setTimeout> | null = null;
|
||||
const off = socket.onAnyMessage((threadId, m) => {
|
||||
setPreviews((p) => ({ ...p, [threadId]: { lastMessage: m.text, lastAt: m.at } }));
|
||||
if (timer) clearTimeout(timer);
|
||||
timer = setTimeout(() => refetchRef.current(), 600);
|
||||
});
|
||||
return () => { off(); if (timer) clearTimeout(timer); };
|
||||
}, [socket]);
|
||||
|
||||
const conversations: UiConversation[] = useMemo(
|
||||
() => (convQ.data ?? []).map((c) => shape(c, nameOf, myId, previews[c.threadId])),
|
||||
[convQ.data, nameOf, myId, previews],
|
||||
);
|
||||
|
||||
const refetch = useCallback(() => { dirQ.refetch(); convQ.refetch(); }, [dirQ, convQ]);
|
||||
const openConversation = useCallback(async (participantIds: string[], opts?: { membership?: Membership; subject?: string }) => {
|
||||
const res = (await sdk.command("crm.messenger.conversation.open", {
|
||||
participantIds, ...(opts?.membership ? { membership: opts.membership } : {}), ...(opts?.subject ? { subject: opts.subject } : {}),
|
||||
})) as { threadId: string };
|
||||
convQ.refetch();
|
||||
return res.threadId;
|
||||
}, [sdk, convQ]);
|
||||
|
||||
return {
|
||||
live: true,
|
||||
loading: dirQ.loading || convQ.loading,
|
||||
error: (dirQ.error ?? convQ.error)?.message ?? null,
|
||||
directory, conversations, nameOf, openConversation, refetch,
|
||||
};
|
||||
}
|
||||
|
||||
function useLiveThread(threadId: string): ThreadData {
|
||||
const { sdk } = useAppShell();
|
||||
const socket = useMessengerSocket();
|
||||
const socketReady = socket?.ready ?? false;
|
||||
const q = useQuery<MessageDTO[]>("crm.messenger.history", { threadId });
|
||||
const [socketMsgs, setSocketMsgs] = useState<UiMessage[]>([]);
|
||||
const [myActorId, setMyActorId] = useState<string | null>(null);
|
||||
const myActorIdRef = useRef<string | null>(null);
|
||||
myActorIdRef.current = myActorId;
|
||||
const [typing, setTyping] = useState<Record<string, number>>({}); // userId -> expiry ts
|
||||
const [seenIds, setSeenIds] = useState<Set<string>>(new Set());
|
||||
const myId = socket?.myUserId;
|
||||
|
||||
// REST poll — the fallback whenever the live socket isn't connected.
|
||||
const refetchRef = useRef(q.refetch);
|
||||
refetchRef.current = q.refetch;
|
||||
useEffect(() => {
|
||||
if (socketReady) return;
|
||||
const t = setInterval(() => refetchRef.current(), POLL_MS);
|
||||
return () => clearInterval(t);
|
||||
}, [socketReady, threadId]);
|
||||
|
||||
// Socket (primary): load history + subscribe to live messages, typing, receipts, reactions.
|
||||
useEffect(() => {
|
||||
if (!socket || !socketReady) return;
|
||||
let alive = true;
|
||||
setSocketMsgs([]); setSeenIds(new Set()); setTyping({});
|
||||
void socket.openThread(threadId).then((hist) => { if (alive) setSocketMsgs(hist); }).catch(() => {});
|
||||
|
||||
const offMsg = socket.subscribe(threadId, (m) =>
|
||||
setSocketMsgs((l) => (l.some((x) => x.id === m.id) ? l : [...l, m])),
|
||||
);
|
||||
const offTyping = socket.onTyping(threadId, (userId) =>
|
||||
setTyping((t) => ({ ...t, [userId]: Date.now() + TYPING_TTL_MS })),
|
||||
);
|
||||
// Receipts are a global stream (no threadId). Count only reads by the OTHER side; seenMine then
|
||||
// narrows to my messages in this thread.
|
||||
const offReceipt = socket.onReceipt((e) => {
|
||||
if (e.actorId === myActorIdRef.current) return;
|
||||
setSeenIds((s) => (s.has(e.interactionId) ? s : new Set(s).add(e.interactionId)));
|
||||
});
|
||||
const offAnn = socket.onAnnotation(threadId, (e) =>
|
||||
setSocketMsgs((l) => l.map((m) => (m.id === e.interactionId ? { ...m, reactions: applyAnnotation(m.reactions, e, myId) } : m))),
|
||||
);
|
||||
return () => { alive = false; offMsg(); offTyping(); offReceipt(); offAnn(); };
|
||||
}, [socket, socketReady, threadId, myId]);
|
||||
|
||||
// Learn my own actor id from a message I sent, so receipts from OTHER actors read as "seen".
|
||||
useEffect(() => {
|
||||
const mine = socketMsgs.find((m) => m.mine && m.actorId);
|
||||
if (mine?.actorId && mine.actorId !== myActorId) setMyActorId(mine.actorId);
|
||||
}, [socketMsgs, myActorId]);
|
||||
|
||||
// Tell the server I've read the latest message (drives the other side's "seen" tick).
|
||||
useEffect(() => {
|
||||
if (!socket || !socketReady || socketMsgs.length === 0) return;
|
||||
socket.markRead(threadId, socketMsgs[socketMsgs.length - 1].id);
|
||||
}, [socket, socketReady, threadId, socketMsgs]);
|
||||
|
||||
// Expire stale typing entries.
|
||||
const typingUserIds = useMemo(() => {
|
||||
const now = Date.now();
|
||||
return Object.entries(typing).filter(([, exp]) => exp > now).map(([u]) => u);
|
||||
}, [typing]);
|
||||
useEffect(() => {
|
||||
if (typingUserIds.length === 0) return;
|
||||
const t = setTimeout(() => setTyping((p) => ({ ...p })), TYPING_TTL_MS);
|
||||
return () => clearTimeout(t);
|
||||
}, [typingUserIds.length, typing]);
|
||||
|
||||
const restMsgs: UiMessage[] = useMemo(
|
||||
() => (q.data ?? []).map((m) => ({
|
||||
id: m.interactionId, actorId: m.actorId, senderId: null, text: m.text ?? "", at: m.occurredAt,
|
||||
mine: !!myActorId && m.actorId === myActorId, reactions: [],
|
||||
})),
|
||||
[q.data, myActorId],
|
||||
);
|
||||
|
||||
const messages = socketReady ? socketMsgs : restMsgs;
|
||||
|
||||
// My messages the other side has read (receipts carry the other actor's id).
|
||||
const seenMine = useMemo(() => {
|
||||
const out = new Set<string>();
|
||||
for (const id of seenIds) if (messages.some((m) => m.id === id && m.mine)) out.add(id);
|
||||
return out;
|
||||
}, [seenIds, messages]);
|
||||
|
||||
const send = useCallback(async (content: string, opts?: { parentInteractionId?: string; attachment?: UiAttachment }) => {
|
||||
if (socket && socketReady) {
|
||||
await socket.send(threadId, content, opts); // echoes back over the socket as a 'message' event
|
||||
} else {
|
||||
// REST fallback carries the attachment ref too; a socket reconnect will replace with the live copy.
|
||||
const m = (await sdk.command("crm.messenger.send", { threadId, content, ...(opts?.attachment ? { attachment: opts.attachment } : {}) })) as MessageDTO;
|
||||
if (m.actorId) setMyActorId(m.actorId);
|
||||
q.refetch();
|
||||
}
|
||||
}, [socket, socketReady, threadId, sdk, q]);
|
||||
|
||||
const react = useCallback((interactionId: string, emoji: string) => {
|
||||
if (socket && socketReady) socket.react(threadId, interactionId, emoji);
|
||||
}, [socket, socketReady, threadId]);
|
||||
|
||||
return {
|
||||
loading: q.loading && !socketReady, error: q.error?.message ?? null,
|
||||
messages, send, react, typingUserIds, seenIds: seenMine, refetch: q.refetch,
|
||||
};
|
||||
}
|
||||
|
||||
function useLiveGroupSettings(threadId: string): GroupSettingsData {
|
||||
const { sdk } = useAppShell();
|
||||
const { user } = useAuth();
|
||||
const q = useQuery<UiMember[]>("crm.messenger.members", { threadId });
|
||||
const members = useMemo(() => q.data ?? [], [q.data]);
|
||||
const isAdmin = useMemo(() => members.some((m) => m.userId === user?.id && m.role === "ADMIN"), [members, user?.id]);
|
||||
|
||||
const rename = useCallback(async (subject: string) => {
|
||||
await sdk.command("crm.messenger.group.rename", { threadId, subject });
|
||||
q.refetch();
|
||||
}, [sdk, threadId, q]);
|
||||
const addMember = useCallback(async (userId: string) => {
|
||||
await sdk.command("crm.messenger.participant.add", { threadId, userId });
|
||||
q.refetch();
|
||||
}, [sdk, threadId, q]);
|
||||
const removeMember = useCallback(async (userId: string) => {
|
||||
await sdk.command("crm.messenger.participant.remove", { threadId, userId });
|
||||
q.refetch();
|
||||
}, [sdk, threadId, q]);
|
||||
|
||||
return { loading: q.loading, error: q.error?.message ?? null, members, isAdmin, rename, addMember, removeMember, refetch: q.refetch };
|
||||
}
|
||||
|
||||
function shape(
|
||||
c: ConversationDTO,
|
||||
nameOf: (id: string) => string,
|
||||
myId: string | undefined,
|
||||
overlay?: { lastMessage: string; lastAt: string },
|
||||
): UiConversation {
|
||||
// A DM's title is the OTHER person — never yourself, and never the raw unknown-id fallback for both.
|
||||
const others = myId ? c.participants.filter((p) => p !== myId) : c.participants;
|
||||
const title = c.subject?.trim()
|
||||
|| (c.membership === "group"
|
||||
? `Group · ${c.participants.length}`
|
||||
: (others.map(nameOf).join(", ") || nameOf(c.participants[0] ?? "") || "Conversation"));
|
||||
const lastMessage = overlay?.lastMessage ?? c.lastMessage;
|
||||
const lastAt = overlay?.lastAt ?? c.lastAt;
|
||||
return {
|
||||
threadId: c.threadId, title, subject: c.subject, membership: c.membership,
|
||||
participants: c.participants, unread: c.unread,
|
||||
...(lastMessage ? { lastMessage } : {}), ...(lastAt ? { lastAt } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
/* ======================================================================== */
|
||||
/* Mock implementation (no Shell configured — the demo keeps working) */
|
||||
/* ======================================================================== */
|
||||
|
||||
const MOCK_PEOPLE: UiPerson[] = [
|
||||
{ id: "pp_sofia", name: "Sofia Ramirez", kind: "staff" },
|
||||
{ id: "pp_dan", name: "Dan Whitaker", kind: "staff" },
|
||||
{ id: "pp_priya", name: "Priya Nair", kind: "staff" },
|
||||
{ id: "cust_acme", name: "Acme Roofing (Client)", kind: "customer" },
|
||||
{ id: "cust_globex", name: "Globex Homes (Client)", kind: "customer" },
|
||||
];
|
||||
|
||||
interface MockThread { threadId: string; membership: Membership; subject: string | null; participants: string[]; messages: UiMessage[] }
|
||||
const now = () => new Date().toISOString();
|
||||
let MOCK_SEQ = 100;
|
||||
|
||||
// A tiny module-level store both mock hooks share, with a subscribe-on-change so the
|
||||
// conversation list and the open thread stay in sync (no globalThis, no render writes).
|
||||
const MOCK_STORE = new Map<string, MockThread>([
|
||||
["th_mock_1", { threadId: "th_mock_1", membership: "dm", subject: null, participants: ["me", "pp_sofia"],
|
||||
messages: [{ id: "m1", actorId: "pp_sofia", text: "Can you review the Henderson estimate?", at: now(), mine: false, reactions: [] }] }],
|
||||
["th_mock_2", { threadId: "th_mock_2", membership: "group", subject: "Storm response — East side", participants: ["me", "pp_dan", "pp_priya"],
|
||||
messages: [{ id: "m2", actorId: "pp_dan", text: "Crew is rolling out at 7.", at: now(), mine: false, reactions: [] }] }],
|
||||
]);
|
||||
const mockListeners = new Set<() => void>();
|
||||
const notifyMock = () => mockListeners.forEach((l) => l());
|
||||
function useMockSubscription(): void {
|
||||
const [, setV] = useState(0);
|
||||
useEffect(() => {
|
||||
const l = () => setV((n) => n + 1);
|
||||
mockListeners.add(l);
|
||||
return () => { mockListeners.delete(l); };
|
||||
}, []);
|
||||
}
|
||||
|
||||
function useMockMessenger(): MessengerData {
|
||||
useMockSubscription();
|
||||
const nameById = useMemo(() => Object.fromEntries(MOCK_PEOPLE.map((p) => [p.id, p.name])), []);
|
||||
const nameOf = useCallback((id: string) => nameById[id] ?? `User ${shortId(id)}`, [nameById]);
|
||||
|
||||
const conversations: UiConversation[] = [...MOCK_STORE.values()].map((t) => {
|
||||
const last = t.messages[t.messages.length - 1];
|
||||
return {
|
||||
threadId: t.threadId,
|
||||
title: t.subject || t.participants.filter((p) => p !== "me").map(nameOf).join(", ") || "Conversation",
|
||||
subject: t.subject, membership: t.membership, participants: t.participants, unread: 0,
|
||||
...(last ? { lastMessage: last.text, lastAt: last.at } : {}),
|
||||
};
|
||||
});
|
||||
|
||||
const openConversation = useCallback(async (participantIds: string[], opts?: { membership?: Membership; subject?: string }) => {
|
||||
const membership = opts?.membership ?? (participantIds.length === 1 ? "dm" : "group");
|
||||
const threadId = `th_mock_${MOCK_SEQ++}`;
|
||||
MOCK_STORE.set(threadId, { threadId, membership, subject: opts?.subject ?? null, participants: ["me", ...participantIds], messages: [] });
|
||||
notifyMock();
|
||||
return threadId;
|
||||
}, []);
|
||||
|
||||
return { live: false, loading: false, error: null, directory: MOCK_PEOPLE, conversations, nameOf, openConversation, refetch: () => {} };
|
||||
}
|
||||
|
||||
function useMockThread(threadId: string): ThreadData {
|
||||
useMockSubscription();
|
||||
const thread = MOCK_STORE.get(threadId);
|
||||
const send = useCallback(async (content: string, opts?: { parentInteractionId?: string }) => {
|
||||
const t = MOCK_STORE.get(threadId);
|
||||
if (t) {
|
||||
t.messages = [...t.messages, {
|
||||
id: `m_${MOCK_SEQ++}`, actorId: "me", text: content, at: now(), mine: true, reactions: [],
|
||||
...(opts?.parentInteractionId ? { parentInteractionId: opts.parentInteractionId } : {}),
|
||||
}];
|
||||
notifyMock();
|
||||
}
|
||||
}, [threadId]);
|
||||
const react = useCallback((interactionId: string, emoji: string) => {
|
||||
const t = MOCK_STORE.get(threadId);
|
||||
if (!t) return;
|
||||
t.messages = t.messages.map((m) => {
|
||||
if (m.id !== interactionId) return m;
|
||||
const has = (m.reactions ?? []).find((r) => r.emoji === emoji);
|
||||
const reactions = has
|
||||
? (m.reactions ?? []).filter((r) => r.emoji !== emoji)
|
||||
: [...(m.reactions ?? []), { emoji, count: 1, mine: true }];
|
||||
return { ...m, reactions };
|
||||
});
|
||||
notifyMock();
|
||||
}, [threadId]);
|
||||
return {
|
||||
loading: false, error: null, messages: thread?.messages ?? [], send, react,
|
||||
typingUserIds: [], seenIds: new Set(), refetch: notifyMock,
|
||||
};
|
||||
}
|
||||
|
||||
function useMockGroupSettings(threadId: string): GroupSettingsData {
|
||||
useMockSubscription();
|
||||
const nameById = useMemo(() => Object.fromEntries(MOCK_PEOPLE.map((p) => [p.id, p.name])), []);
|
||||
const t = MOCK_STORE.get(threadId);
|
||||
const members: UiMember[] = (t?.participants ?? []).map((id) => ({
|
||||
userId: id,
|
||||
displayName: id === "me" ? "You" : (nameById[id] ?? `User ${shortId(id)}`),
|
||||
role: id === "me" ? "ADMIN" : "MEMBER",
|
||||
}));
|
||||
const rename = useCallback(async (subject: string) => {
|
||||
const th = MOCK_STORE.get(threadId);
|
||||
if (th) { th.subject = subject; notifyMock(); }
|
||||
}, [threadId]);
|
||||
const addMember = useCallback(async (userId: string) => {
|
||||
const th = MOCK_STORE.get(threadId);
|
||||
if (th && !th.participants.includes(userId)) { th.participants = [...th.participants, userId]; notifyMock(); }
|
||||
}, [threadId]);
|
||||
const removeMember = useCallback(async (userId: string) => {
|
||||
const th = MOCK_STORE.get(threadId);
|
||||
if (th) { th.participants = th.participants.filter((p) => p !== userId); notifyMock(); }
|
||||
}, [threadId]);
|
||||
return { loading: false, error: null, members, isAdmin: true, rename, addMember, removeMember, refetch: notifyMock };
|
||||
}
|
||||
@@ -1,149 +0,0 @@
|
||||
"use client";
|
||||
|
||||
// v2 live stream: one IIOS MessageSocket for the whole Messenger panel, using the SDK
|
||||
// (@insignia/iios-kernel-client) — not raw socket.io. The delegated realtime token comes from
|
||||
// the be-crm data door (crm.messenger.realtime). Threads subscribe through a context; the socket
|
||||
// re-opens every joined thread on reconnect (handled inside the SDK). In mock mode this is a no-op
|
||||
// passthrough and the thread hook falls back to the REST poll.
|
||||
//
|
||||
// Beyond plain messages, the kernel exposes typing, read receipts, and reactions (generic
|
||||
// annotations). This provider fans each server event out to per-thread listeners so the UI can
|
||||
// render typing indicators, "seen" ticks, and emoji reactions live.
|
||||
|
||||
import { createContext, useCallback, useContext, useEffect, useRef, useState, type ReactNode } from "react";
|
||||
import { MessageSocket, type Message, type AnnotationEvent } from "@insignia/iios-kernel-client";
|
||||
import { useAuth, useQuery } from "@abe-kap/appshell-sdk/react";
|
||||
import { isShellConfigured } from "./appshell";
|
||||
import { toReactions, type UiMessage } from "./messenger-api";
|
||||
|
||||
interface RealtimeDTO { url: string; audience: string; token?: string }
|
||||
|
||||
export interface ReceiptHit { interactionId: string; actorId: string }
|
||||
|
||||
export interface MessengerSocket {
|
||||
ready: boolean;
|
||||
myUserId?: string;
|
||||
openThread: (threadId: string) => Promise<UiMessage[]>;
|
||||
send: (threadId: string, content: string, opts?: { parentInteractionId?: string; attachment?: { contentRef: string; mimeType: string; sizeBytes: number } }) => Promise<void>;
|
||||
subscribe: (threadId: string, cb: (m: UiMessage) => void) => () => void;
|
||||
/** Fires for EVERY inbound message regardless of thread — drives live sidebar previews. */
|
||||
onAnyMessage: (cb: (threadId: string, m: UiMessage) => void) => () => void;
|
||||
sendTyping: (threadId: string) => void;
|
||||
onTyping: (threadId: string, cb: (userId: string) => void) => () => void;
|
||||
markRead: (threadId: string, interactionId: string) => void;
|
||||
/** The kernel's receipt event carries no threadId, so this is a global stream; the thread hook
|
||||
* filters to receipts for its own (mine) messages. */
|
||||
onReceipt: (cb: (e: ReceiptHit) => void) => () => void;
|
||||
react: (threadId: string, interactionId: string, emoji: string) => void;
|
||||
onAnnotation: (threadId: string, cb: (e: AnnotationEvent) => void) => () => void;
|
||||
}
|
||||
|
||||
const Ctx = createContext<MessengerSocket | null>(null);
|
||||
export function useMessengerSocket(): MessengerSocket | null { return useContext(Ctx); }
|
||||
|
||||
const SHELL = isShellConfigured();
|
||||
|
||||
const toUi = (m: Message, myUserId?: string): UiMessage => ({
|
||||
id: m.id, actorId: m.senderActorId ?? null, senderId: m.senderId ?? null, text: m.content ?? "", at: m.createdAt,
|
||||
mine: !!myUserId && m.senderId === myUserId,
|
||||
...(m.parentInteractionId ? { parentInteractionId: m.parentInteractionId } : {}),
|
||||
...(m.attachment ? { attachment: { contentRef: m.attachment.contentRef, mimeType: m.attachment.mimeType, sizeBytes: m.attachment.sizeBytes } } : {}),
|
||||
reactions: toReactions(m.annotations, myUserId),
|
||||
});
|
||||
|
||||
export function MessengerSocketProvider({ children }: { children: ReactNode }) {
|
||||
// SHELL is a build-time constant, so the branch is stable across renders (Rules-of-Hooks safe).
|
||||
if (!SHELL) return <>{children}</>;
|
||||
return <LiveSocketProvider>{children}</LiveSocketProvider>;
|
||||
}
|
||||
|
||||
// A tiny per-thread listener registry, reused for messages / typing / receipts / annotations.
|
||||
function makeRegistry<T>() {
|
||||
const map = new Map<string, Set<(v: T) => void>>();
|
||||
const add = (key: string, cb: (v: T) => void) => {
|
||||
if (!map.has(key)) map.set(key, new Set());
|
||||
map.get(key)!.add(cb);
|
||||
return () => { map.get(key)?.delete(cb); };
|
||||
};
|
||||
const emit = (key: string, v: T) => map.get(key)?.forEach((cb) => cb(v));
|
||||
return { add, emit };
|
||||
}
|
||||
|
||||
function LiveSocketProvider({ children }: { children: ReactNode }) {
|
||||
const { user } = useAuth();
|
||||
const rt = useQuery<RealtimeDTO>("crm.messenger.realtime", {});
|
||||
const [ready, setReady] = useState(false);
|
||||
const socketRef = useRef<MessageSocket | null>(null);
|
||||
const myRef = useRef<string | undefined>(user?.id);
|
||||
myRef.current = user?.id;
|
||||
|
||||
// One registry per event kind, keyed by threadId (plus a global message fan-out).
|
||||
const msgReg = useRef(makeRegistry<UiMessage>()).current;
|
||||
const anyMsg = useRef(new Set<(threadId: string, m: UiMessage) => void>()).current;
|
||||
const typingReg = useRef(makeRegistry<string>()).current;
|
||||
const receiptSet = useRef(new Set<(e: ReceiptHit) => void>()).current;
|
||||
const annReg = useRef(makeRegistry<AnnotationEvent>()).current;
|
||||
|
||||
const url = rt.data?.url;
|
||||
const token = rt.data?.token;
|
||||
|
||||
useEffect(() => {
|
||||
if (!url || !token) return;
|
||||
const socket = new MessageSocket({ serviceUrl: url, token, autoConnect: false });
|
||||
socketRef.current = socket;
|
||||
const offConnected = socket.onConnected(() => setReady(true));
|
||||
const offMessage = socket.on("message", (m) => {
|
||||
const ui = toUi(m, myRef.current);
|
||||
msgReg.emit(m.threadId, ui);
|
||||
anyMsg.forEach((cb) => cb(m.threadId, ui));
|
||||
});
|
||||
const offTyping = socket.on("typing", (e) => { if (e.userId !== myRef.current) typingReg.emit(e.threadId, e.userId); });
|
||||
const offReceipt = socket.on("receipt", (e) => receiptSet.forEach((cb) => cb({ interactionId: e.interactionId, actorId: e.actorId })));
|
||||
const offAnn = socket.on("annotation", (e) => annReg.emit(e.threadId, e));
|
||||
socket.connect();
|
||||
return () => {
|
||||
offConnected(); offMessage(); offTyping(); offReceipt(); offAnn();
|
||||
socket.disconnect(); socketRef.current = null; setReady(false);
|
||||
};
|
||||
}, [url, token, msgReg, anyMsg, typingReg, receiptSet, annReg]);
|
||||
|
||||
const openThread = useCallback(async (threadId: string): Promise<UiMessage[]> => {
|
||||
const s = socketRef.current;
|
||||
if (!s) return [];
|
||||
const res = await s.openThread(threadId);
|
||||
return res.history.map((m) => toUi(m, myRef.current));
|
||||
}, []);
|
||||
|
||||
const send = useCallback(async (threadId: string, content: string, opts?: { parentInteractionId?: string; attachment?: { contentRef: string; mimeType: string; sizeBytes: number } }) => {
|
||||
const s = socketRef.current;
|
||||
if (!s) throw new Error("Not connected");
|
||||
const sendOpts = {
|
||||
...(opts?.parentInteractionId ? { parentInteractionId: opts.parentInteractionId } : {}),
|
||||
...(opts?.attachment ? { attachment: opts.attachment } : {}),
|
||||
};
|
||||
await s.sendMessage(threadId, content, Object.keys(sendOpts).length ? sendOpts : undefined);
|
||||
}, []);
|
||||
|
||||
const subscribe = useCallback((threadId: string, cb: (m: UiMessage) => void) => msgReg.add(threadId, cb), [msgReg]);
|
||||
const onAnyMessage = useCallback((cb: (threadId: string, m: UiMessage) => void) => {
|
||||
anyMsg.add(cb); return () => { anyMsg.delete(cb); };
|
||||
}, [anyMsg]);
|
||||
const onTyping = useCallback((threadId: string, cb: (userId: string) => void) => typingReg.add(threadId, cb), [typingReg]);
|
||||
const onReceipt = useCallback((cb: (e: ReceiptHit) => void) => {
|
||||
receiptSet.add(cb); return () => { receiptSet.delete(cb); };
|
||||
}, [receiptSet]);
|
||||
const onAnnotation = useCallback((threadId: string, cb: (e: AnnotationEvent) => void) => annReg.add(threadId, cb), [annReg]);
|
||||
|
||||
const sendTyping = useCallback((threadId: string) => socketRef.current?.typing(threadId), []);
|
||||
const markRead = useCallback((threadId: string, interactionId: string) => { void socketRef.current?.markRead(threadId, interactionId); }, []);
|
||||
const react = useCallback((threadId: string, interactionId: string, emoji: string) => { void socketRef.current?.react(threadId, interactionId, emoji); }, []);
|
||||
|
||||
return (
|
||||
<Ctx.Provider value={{
|
||||
ready, myUserId: user?.id, openThread, send, subscribe, onAnyMessage,
|
||||
sendTyping, onTyping, markRead, onReceipt, react, onAnnotation,
|
||||
}}>
|
||||
{children}
|
||||
</Ctx.Provider>
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,78 @@
|
||||
"use client";
|
||||
|
||||
// SMS settings data layer. Serves EITHER the local mock (Shell not configured — the demo keeps
|
||||
// working) OR the live be-crm data door (crm.settings.sms.*), behind one interface.
|
||||
//
|
||||
// Live contract (be-crm → IIOS BYO credential store):
|
||||
// query crm.settings.sms.status {} -> { configured, enabled?, hints? }
|
||||
// cmd crm.settings.sms.configure { accountSid, authToken, fromNumber } -> masked status
|
||||
// The auth token is write-only: it is sealed in IIOS and NEVER returned — status carries only
|
||||
// non-secret hints (from-number + SID last-4).
|
||||
|
||||
import { useCallback, useState } from "react";
|
||||
import { useAppShell, useQuery } from "@abe-kap/appshell-sdk/react";
|
||||
import { isShellConfigured } from "./appshell";
|
||||
|
||||
export interface SmsCredentials { accountSid: string; authToken: string; fromNumber: string }
|
||||
|
||||
export interface SmsStatus {
|
||||
configured: boolean;
|
||||
enabled: boolean;
|
||||
fromNumber?: string;
|
||||
sidLast4?: string;
|
||||
}
|
||||
|
||||
export interface SmsSettingsData {
|
||||
live: boolean;
|
||||
loading: boolean;
|
||||
error: string | null;
|
||||
status: SmsStatus;
|
||||
configure: (input: SmsCredentials) => Promise<void>;
|
||||
refetch: () => void;
|
||||
}
|
||||
|
||||
interface StatusDTO { configured: boolean; enabled?: boolean; hints?: { fromNumber?: string; sidLast4?: string } }
|
||||
|
||||
function toStatus(dto?: StatusDTO | null): SmsStatus {
|
||||
return {
|
||||
configured: !!dto?.configured,
|
||||
enabled: dto?.enabled ?? false,
|
||||
fromNumber: dto?.hints?.fromNumber,
|
||||
sidLast4: dto?.hints?.sidLast4,
|
||||
};
|
||||
}
|
||||
|
||||
/* ---- Mock (demo mode) — stores only the non-secret hints, mirroring the masked live status ---- */
|
||||
function useMockSms(): SmsSettingsData {
|
||||
const [status, setStatus] = useState<SmsStatus>({ configured: false, enabled: false });
|
||||
const configure = useCallback(async ({ accountSid, fromNumber }: SmsCredentials) => {
|
||||
setStatus({ configured: true, enabled: true, fromNumber, sidLast4: accountSid.slice(-4) });
|
||||
}, []);
|
||||
return { live: false, loading: false, error: null, status, configure, refetch: () => {} };
|
||||
}
|
||||
|
||||
/* ---- Live (be-crm data door) ---- */
|
||||
function useLiveSms(): SmsSettingsData {
|
||||
const { sdk } = useAppShell();
|
||||
const q = useQuery<StatusDTO>("crm.settings.sms.status", {});
|
||||
const configure = useCallback(async (input: SmsCredentials) => {
|
||||
await sdk.command("crm.settings.sms.configure", { ...input });
|
||||
q.refetch();
|
||||
}, [sdk, q]);
|
||||
return {
|
||||
live: true,
|
||||
loading: q.loading,
|
||||
error: q.error ? String(q.error) : null,
|
||||
status: toStatus(q.data),
|
||||
configure,
|
||||
refetch: q.refetch,
|
||||
};
|
||||
}
|
||||
|
||||
const SHELL = isShellConfigured();
|
||||
|
||||
export function useSmsSettings(): SmsSettingsData {
|
||||
// SHELL is constant for the bundle's life (NEXT_PUBLIC_* is build-time), so the same hook path
|
||||
// runs every render — Rules-of-Hooks safe.
|
||||
return SHELL ? useLiveSms() : useMockSms();
|
||||
}
|
||||
Reference in New Issue
Block a user