diff --git a/.env.example b/.env.example
index 0cfe07a..30770bb 100644
--- a/.env.example
+++ b/.env.example
@@ -26,3 +26,14 @@ JUDGE_BASE_URL="http://localhost:2000/api/v2"
JUDGE_API_KEY="replace-with-proxy-api-key"
PISTON_PYTHON_VERSION="3.10.0"
PISTON_CPP_VERSION="10.2.0"
+
+# Ephemeral chat. Run `npm run chat:gateway` on a persistent host.
+CHAT_REALTIME_TOKEN_SECRET="replace-with-a-random-secret"
+CHAT_REALTIME_PUBLISH_SECRET="replace-with-a-different-random-secret"
+CHAT_REALTIME_INTERNAL_URL="http://localhost:8787"
+NEXT_PUBLIC_CHAT_WS_URL="ws://localhost:8787/socket"
+CHAT_ALLOWED_ORIGINS="http://localhost:3000"
+CHAT_GATEWAY_PORT="8787"
+
+# Vercel sends this bearer secret to authenticated cron routes.
+CRON_SECRET="replace-with-a-random-secret"
diff --git a/app/(auth)/join/page.tsx b/app/(auth)/join/page.tsx
index 43f9694..78f22ed 100644
--- a/app/(auth)/join/page.tsx
+++ b/app/(auth)/join/page.tsx
@@ -53,6 +53,9 @@ export default async function JoinPage() {
>
Continue as applicant (dev)
+
+ Continue as member (dev)
+
Continue as admin (dev)
diff --git a/app/account-bar.tsx b/app/account-bar.tsx
index 57dee86..3a7bf74 100644
--- a/app/account-bar.tsx
+++ b/app/account-bar.tsx
@@ -5,6 +5,7 @@ import {
relativeTimeFromNow,
unreadNotificationCount,
} from "../lib/notifications";
+import MessageIndicator from "./message-indicator";
import NotificationBell, { type NotificationItem } from "./notification-bell";
import SiteHeader from "./site-header";
@@ -46,6 +47,7 @@ export default async function AccountBar({
>
+
);
diff --git a/app/api/chat/ack/route.ts b/app/api/chat/ack/route.ts
new file mode 100644
index 0000000..34b0461
--- /dev/null
+++ b/app/api/chat/ack/route.ts
@@ -0,0 +1,19 @@
+import { NextResponse } from "next/server";
+import { chatAckSchema } from "../../../../lib/chat";
+import { activeChatUser } from "../../../../lib/chat-auth";
+import { chatErrorResponse, unauthorizedChatResponse } from "../../../../lib/chat-http";
+import { acknowledgeChatMessages } from "../../../../lib/chat-service";
+
+export async function POST(request: Request) {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+
+ try {
+ const parsed = chatAckSchema.safeParse(await request.json());
+ if (!parsed.success) return NextResponse.json({ error: "INVALID_ACK" }, { status: 400 });
+
+ return NextResponse.json(await acknowledgeChatMessages(user.id, parsed.data.messageIds));
+ } catch (error) {
+ return chatErrorResponse(error);
+ }
+}
diff --git a/app/api/chat/conversations/route.ts b/app/api/chat/conversations/route.ts
new file mode 100644
index 0000000..d7a8a6a
--- /dev/null
+++ b/app/api/chat/conversations/route.ts
@@ -0,0 +1,51 @@
+import { ChatConversationType } from "@/prisma-client";
+import { NextResponse } from "next/server";
+import {
+ createDirectConversationSchema,
+ createGroupConversationSchema,
+} from "../../../../lib/chat";
+import { activeChatUser } from "../../../../lib/chat-auth";
+import { chatErrorResponse, unauthorizedChatResponse } from "../../../../lib/chat-http";
+import {
+ createDirectConversation,
+ createGroupConversation,
+ listChatConversations,
+} from "../../../../lib/chat-service";
+
+export const dynamic = "force-dynamic";
+
+export async function GET() {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+ return NextResponse.json({ conversations: await listChatConversations(user.id) });
+}
+
+export async function POST(request: Request) {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+
+ try {
+ const body = (await request.json()) as { type?: unknown };
+ if (body.type === ChatConversationType.DIRECT) {
+ const parsed = createDirectConversationSchema.safeParse(body);
+ if (!parsed.success) return NextResponse.json({ error: "INVALID_DIRECT" }, { status: 400 });
+ const conversation = await createDirectConversation(user.id, parsed.data.recipientId);
+ return NextResponse.json({ conversationId: conversation.id });
+ }
+
+ if (body.type === ChatConversationType.GROUP) {
+ const parsed = createGroupConversationSchema.safeParse(body);
+ if (!parsed.success) return NextResponse.json({ error: "INVALID_GROUP" }, { status: 400 });
+ const conversation = await createGroupConversation(
+ user.id,
+ parsed.data.title,
+ parsed.data.memberIds,
+ );
+ return NextResponse.json({ conversationId: conversation.id }, { status: 201 });
+ }
+
+ return NextResponse.json({ error: "INVALID_CONVERSATION_TYPE" }, { status: 400 });
+ } catch (error) {
+ return chatErrorResponse(error);
+ }
+}
diff --git a/app/api/chat/inbox/route.ts b/app/api/chat/inbox/route.ts
new file mode 100644
index 0000000..8b8befb
--- /dev/null
+++ b/app/api/chat/inbox/route.ts
@@ -0,0 +1,18 @@
+import { NextResponse } from "next/server";
+import { activeChatUser } from "../../../../lib/chat-auth";
+import { chatErrorResponse, unauthorizedChatResponse } from "../../../../lib/chat-http";
+import { listPendingChatMessages } from "../../../../lib/chat-service";
+
+export const dynamic = "force-dynamic";
+
+export async function GET(request: Request) {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+
+ try {
+ const cursor = new URL(request.url).searchParams.get("cursor")?.trim() || undefined;
+ return NextResponse.json(await listPendingChatMessages(user.id, cursor));
+ } catch (error) {
+ return chatErrorResponse(error);
+ }
+}
diff --git a/app/api/chat/messages/route.ts b/app/api/chat/messages/route.ts
new file mode 100644
index 0000000..862737b
--- /dev/null
+++ b/app/api/chat/messages/route.ts
@@ -0,0 +1,37 @@
+import { NextResponse } from "next/server";
+import { sendChatMessageSchema } from "../../../../lib/chat";
+import { activeChatUser } from "../../../../lib/chat-auth";
+import { chatErrorResponse, unauthorizedChatResponse } from "../../../../lib/chat-http";
+import { publishChatEnvelope } from "../../../../lib/chat-realtime";
+import { enqueueChatMessage } from "../../../../lib/chat-service";
+
+export async function POST(request: Request) {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+
+ try {
+ const parsed = sendChatMessageSchema.safeParse(await request.json());
+ if (!parsed.success) {
+ return NextResponse.json({ error: "INVALID_MESSAGE" }, { status: 400 });
+ }
+
+ const result = await enqueueChatMessage({ ...parsed.data, senderId: user.id });
+ const realtime =
+ result.envelope && result.recipientIds.length > 0
+ ? await publishChatEnvelope(result.recipientIds, result.envelope)
+ : { attempted: 0, delivered: 0, unavailable: false };
+
+ return NextResponse.json(
+ {
+ messageId: result.messageId,
+ conversationId: result.conversationId,
+ createdAt: result.createdAt,
+ duplicate: result.duplicate,
+ realtime,
+ },
+ { status: result.duplicate ? 200 : 202 },
+ );
+ } catch (error) {
+ return chatErrorResponse(error);
+ }
+}
diff --git a/app/api/chat/read/route.ts b/app/api/chat/read/route.ts
new file mode 100644
index 0000000..95825c6
--- /dev/null
+++ b/app/api/chat/read/route.ts
@@ -0,0 +1,19 @@
+import { NextResponse } from "next/server";
+import { markChatReadSchema } from "../../../../lib/chat";
+import { activeChatUser } from "../../../../lib/chat-auth";
+import { chatErrorResponse, unauthorizedChatResponse } from "../../../../lib/chat-http";
+import { markChatRead } from "../../../../lib/chat-service";
+
+export async function POST(request: Request) {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+
+ try {
+ const parsed = markChatReadSchema.safeParse(await request.json());
+ if (!parsed.success) return NextResponse.json({ error: "INVALID_READ" }, { status: 400 });
+ await markChatRead(user.id, parsed.data.conversationId);
+ return NextResponse.json({ ok: true });
+ } catch (error) {
+ return chatErrorResponse(error);
+ }
+}
diff --git a/app/api/chat/report/route.ts b/app/api/chat/report/route.ts
new file mode 100644
index 0000000..4574e34
--- /dev/null
+++ b/app/api/chat/report/route.ts
@@ -0,0 +1,19 @@
+import { NextResponse } from "next/server";
+import { chatReportSchema } from "../../../../lib/chat";
+import { activeChatUser } from "../../../../lib/chat-auth";
+import { chatErrorResponse, unauthorizedChatResponse } from "../../../../lib/chat-http";
+import { reportChatMessage } from "../../../../lib/chat-service";
+
+export async function POST(request: Request) {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+
+ try {
+ const parsed = chatReportSchema.safeParse(await request.json());
+ if (!parsed.success) return NextResponse.json({ error: "INVALID_REPORT" }, { status: 400 });
+ const report = await reportChatMessage({ ...parsed.data, reporterId: user.id });
+ return NextResponse.json({ reportId: report.id }, { status: 201 });
+ } catch (error) {
+ return chatErrorResponse(error);
+ }
+}
diff --git a/app/api/chat/token/route.ts b/app/api/chat/token/route.ts
new file mode 100644
index 0000000..1540803
--- /dev/null
+++ b/app/api/chat/token/route.ts
@@ -0,0 +1,14 @@
+import { NextResponse } from "next/server";
+import { activeChatUser } from "../../../../lib/chat-auth";
+import { unauthorizedChatResponse } from "../../../../lib/chat-http";
+import { createChatSocketToken } from "../../../../lib/chat-realtime";
+
+export const dynamic = "force-dynamic";
+
+export async function POST() {
+ const user = await activeChatUser();
+ if (!user) return unauthorizedChatResponse();
+
+ const token = await createChatSocketToken(user.id);
+ return NextResponse.json({ token, socketUrl: process.env.NEXT_PUBLIC_CHAT_WS_URL ?? null });
+}
diff --git a/app/api/cron/chat-cleanup/route.ts b/app/api/cron/chat-cleanup/route.ts
new file mode 100644
index 0000000..1117df9
--- /dev/null
+++ b/app/api/cron/chat-cleanup/route.ts
@@ -0,0 +1,27 @@
+import { timingSafeEqual } from "node:crypto";
+import { NextResponse } from "next/server";
+import { cleanupExpiredChatData } from "../../../../lib/chat-service";
+
+function authorized(request: Request) {
+ const expected = process.env.CRON_SECRET;
+ const authorization = request.headers.get("authorization");
+ const provided = authorization?.startsWith("Bearer ") ? authorization.slice(7) : "";
+
+ if (!expected) return false;
+ const expectedBuffer = Buffer.from(expected);
+ const providedBuffer = Buffer.from(provided);
+ return (
+ expectedBuffer.length === providedBuffer.length &&
+ timingSafeEqual(expectedBuffer, providedBuffer)
+ );
+}
+
+export async function GET(request: Request) {
+ if (!authorized(request)) {
+ return NextResponse.json({ error: "Unauthorized" }, { status: 401 });
+ }
+
+ const result = await cleanupExpiredChatData();
+ console.info("Chat cleanup completed", result);
+ return NextResponse.json({ ok: true, ...result });
+}
diff --git a/app/api/health/route.ts b/app/api/health/route.ts
index 0679c1a..f537798 100644
--- a/app/api/health/route.ts
+++ b/app/api/health/route.ts
@@ -17,6 +17,12 @@ export async function GET() {
googleClientSecret: isConfigured(process.env.AUTH_GOOGLE_SECRET),
judgeBaseUrl: process.env.JUDGE_PROVIDER === "fake" || isConfigured(process.env.JUDGE_BASE_URL),
judgeApiKey: process.env.JUDGE_PROVIDER === "fake" || isConfigured(process.env.JUDGE_API_KEY),
+ chatRealtime:
+ process.env.NODE_ENV !== "production" ||
+ (isConfigured(process.env.CHAT_REALTIME_TOKEN_SECRET) &&
+ isConfigured(process.env.CHAT_REALTIME_PUBLISH_SECRET) &&
+ isConfigured(process.env.CHAT_REALTIME_INTERNAL_URL) &&
+ isConfigured(process.env.NEXT_PUBLIC_CHAT_WS_URL)),
};
let database = false;
@@ -44,6 +50,8 @@ export async function GET() {
prisma.problem.count(),
prisma.testCase.count(),
prisma.submission.count(),
+ prisma.chatConversation.count(),
+ prisma.chatDelivery.count(),
]);
schema = true;
} catch (error) {
diff --git a/app/dev-login/route.ts b/app/dev-login/route.ts
index ca7f247..477f66d 100644
--- a/app/dev-login/route.ts
+++ b/app/dev-login/route.ts
@@ -36,10 +36,11 @@ export async function GET(request: Request) {
return NextResponse.json({ error: "Not found" }, { status: 404 });
}
+ const requestedRole = new URL(request.url).searchParams.get("role");
const role: LocalDevRole =
- new URL(request.url).searchParams.get("role") === "admin" ? "admin" : "member";
+ requestedRole === "admin" ? "admin" : requestedRole === "active" ? "active" : "member";
await signIn(localDevProviderId(role), {
- redirectTo: role === "admin" ? "/admin/cohort" : "/apply",
+ redirectTo: role === "admin" ? "/admin/cohort" : role === "active" ? "/dashboard" : "/apply",
});
}
diff --git a/app/globals.css b/app/globals.css
index 21a2e3f..41d4f2b 100644
--- a/app/globals.css
+++ b/app/globals.css
@@ -319,7 +319,13 @@ img {
margin-left: 8px;
}
-.notification-bell {
+.message-indicator-wrap {
+ display: inline-flex;
+ margin-left: 8px;
+}
+
+.notification-bell,
+.message-indicator {
position: relative;
display: inline-grid;
place-items: center;
@@ -333,10 +339,16 @@ img {
cursor: pointer;
}
-.notification-bell:hover {
+.notification-bell:hover,
+.message-indicator:hover {
background: var(--surface-hover);
}
+.message-indicator.is-unread {
+ border-color: var(--accent);
+ background: color-mix(in srgb, var(--accent) 20%, var(--surface));
+}
+
.notification-bell-badge {
position: absolute;
top: -6px;
@@ -3117,6 +3129,418 @@ body.is-resizing-panes {
}
}
+/* Ephemeral messaging workspace. */
+.chat-page {
+ width: min(1180px, calc(100% - 32px));
+ min-height: calc(100vh - 112px);
+ margin: 0 auto;
+ padding: 24px 0;
+}
+
+.chat-page + .feedback-button {
+ display: none;
+}
+
+.chat-workspace {
+ display: grid;
+ grid-template-columns: minmax(260px, 340px) minmax(0, 1fr);
+ height: min(760px, calc(100vh - 136px));
+ min-height: 560px;
+ border: 2px solid var(--ink);
+ background: var(--surface);
+ box-shadow: 7px 7px 0 var(--shadow);
+ overflow: hidden;
+}
+
+.chat-sidebar {
+ position: relative;
+ display: flex;
+ min-width: 0;
+ flex-direction: column;
+ border-right: 1px solid var(--line);
+ background: var(--paper);
+}
+
+.chat-sidebar-header,
+.chat-thread-header {
+ display: flex;
+ align-items: center;
+ justify-content: space-between;
+ gap: 16px;
+ padding: 18px;
+ border-bottom: 1px solid var(--line);
+}
+
+.chat-sidebar-header h1,
+.chat-thread-header h2 {
+ margin: 0;
+ font-size: 1.45rem;
+ line-height: 1.1;
+ letter-spacing: 0;
+}
+
+.chat-sidebar-header .section-label,
+.chat-thread-header p {
+ margin: 0 0 4px;
+ color: var(--muted);
+ font-size: 0.78rem;
+}
+
+.chat-connection {
+ padding: 4px 8px;
+ border: 1px solid var(--line);
+ color: var(--muted);
+ font-family: var(--font-main);
+ font-size: 0.72rem;
+ font-weight: 800;
+}
+
+.chat-connection.is-live {
+ border-color: var(--accent);
+ color: var(--ink);
+}
+
+.chat-new-actions {
+ display: grid;
+ grid-template-columns: 1fr 1fr;
+ gap: 8px;
+ padding: 12px;
+ border-bottom: 1px solid var(--line);
+}
+
+.chat-new-actions button,
+.chat-compose-actions button,
+.chat-back {
+ min-height: 36px;
+ border: 1px solid var(--line);
+ background: var(--surface);
+ color: var(--ink);
+ font-family: var(--font-main);
+ font-weight: 800;
+ cursor: pointer;
+}
+
+.chat-compose-panel {
+ position: absolute;
+ z-index: 4;
+ top: 143px;
+ right: 12px;
+ left: 12px;
+ display: grid;
+ gap: 10px;
+ max-height: calc(100% - 156px);
+ padding: 16px;
+ overflow-y: auto;
+ border: 1px solid var(--ink);
+ background: var(--paper);
+ box-shadow: 5px 5px 0 var(--shadow);
+}
+
+.chat-compose-panel input,
+.chat-compose-panel select,
+.chat-send-form textarea {
+ width: 100%;
+ border: 1px solid var(--line);
+ background: var(--surface);
+ color: var(--ink);
+ font: inherit;
+}
+
+.chat-compose-panel input,
+.chat-compose-panel select {
+ min-height: 42px;
+ padding: 8px 10px;
+}
+
+.chat-member-picker {
+ max-height: min(220px, calc(100vh - 390px));
+ overflow: auto;
+ border: 1px solid var(--line);
+}
+
+.chat-member-picker label {
+ display: flex;
+ gap: 8px;
+ align-items: center;
+ min-height: 42px;
+ padding: 8px 10px;
+ border-bottom: 1px solid var(--line);
+ cursor: pointer;
+}
+
+.chat-member-picker label:last-child {
+ border-bottom: 0;
+}
+
+.chat-member-picker label:hover {
+ background: var(--surface-hover);
+}
+
+.chat-member-picker input {
+ width: 18px;
+ height: 18px;
+ min-height: 18px;
+ flex: 0 0 18px;
+ margin: 0;
+ accent-color: var(--accent);
+}
+
+.chat-compose-actions {
+ display: flex;
+ justify-content: flex-end;
+ gap: 8px;
+}
+
+.chat-compose-actions .button {
+ min-height: 38px;
+ padding: 0 14px;
+ box-shadow: 3px 3px 0 var(--shadow);
+}
+
+.chat-conversation-list,
+.chat-message-list {
+ min-height: 0;
+ overflow-y: auto;
+ overscroll-behavior: contain;
+}
+
+.chat-conversation-list > button {
+ width: 100%;
+ display: flex;
+ align-items: center;
+ justify-content: space-between;
+ gap: 12px;
+ padding: 14px 16px;
+ border: 0;
+ border-bottom: 1px solid var(--line);
+ background: transparent;
+ color: var(--ink);
+ text-align: left;
+ cursor: pointer;
+}
+
+.chat-conversation-list > button:hover,
+.chat-conversation-list > button.is-active {
+ background: var(--surface-hover);
+}
+
+.chat-conversation-list > button.is-active {
+ box-shadow: inset 4px 0 0 var(--accent);
+}
+
+.chat-conversation-copy {
+ min-width: 0;
+}
+
+.chat-conversation-copy strong,
+.chat-conversation-copy small {
+ display: block;
+ overflow: hidden;
+ text-overflow: ellipsis;
+ white-space: nowrap;
+}
+
+.chat-conversation-copy small {
+ margin-top: 4px;
+ color: var(--muted);
+}
+
+.chat-unread-badge {
+ min-width: 22px;
+ height: 22px;
+ display: inline-grid;
+ place-items: center;
+ padding: 0 6px;
+ background: var(--accent);
+ border: 1px solid var(--ink);
+ font-size: 0.72rem;
+ font-weight: 900;
+}
+
+.chat-thread {
+ min-width: 0;
+ display: grid;
+ grid-template-rows: auto minmax(0, 1fr) auto auto;
+ background: var(--surface);
+}
+
+.chat-back {
+ display: none;
+}
+
+.chat-message-list {
+ display: flex;
+ flex-direction: column;
+ gap: 10px;
+ padding: 20px;
+}
+
+.chat-message {
+ width: min(76%, 620px);
+ align-self: flex-start;
+ padding: 11px 13px;
+ border: 1px solid var(--line);
+ background: var(--paper);
+}
+
+.chat-message.is-own {
+ align-self: flex-end;
+ border-color: var(--accent);
+ background: var(--surface-hover);
+}
+
+.chat-message-meta {
+ display: flex;
+ justify-content: space-between;
+ gap: 14px;
+ color: var(--muted);
+ font-size: 0.75rem;
+}
+
+.chat-message p {
+ margin: 7px 0 0;
+ color: var(--ink);
+ line-height: 1.45;
+ overflow-wrap: anywhere;
+ white-space: pre-wrap;
+}
+
+.chat-message small {
+ display: block;
+ margin-top: 5px;
+ color: var(--muted);
+}
+
+.chat-message-error,
+.chat-form-error {
+ color: var(--danger, #b42318) !important;
+}
+
+.chat-send-form {
+ display: grid;
+ grid-template-columns: minmax(0, 1fr) auto;
+ align-items: end;
+ gap: 10px;
+ padding: 14px;
+ border-top: 1px solid var(--line);
+}
+
+.sr-only {
+ position: absolute;
+ width: 1px;
+ height: 1px;
+ padding: 0;
+ margin: -1px;
+ overflow: hidden;
+ clip: rect(0, 0, 0, 0);
+ white-space: nowrap;
+ border: 0;
+}
+
+.chat-send-form textarea {
+ min-height: 48px;
+ max-height: 150px;
+ padding: 10px 12px;
+ resize: vertical;
+}
+
+.chat-send-form .button {
+ min-height: 48px;
+ padding: 0 18px;
+ box-shadow: 3px 3px 0 var(--shadow);
+}
+
+.chat-form-error,
+.chat-empty-copy {
+ margin: 0;
+ padding: 0 14px 12px;
+ font-size: 0.85rem;
+}
+
+.chat-thread-empty {
+ align-self: center;
+ justify-self: center;
+ padding: 24px;
+ color: var(--muted);
+ text-align: center;
+}
+
+.chat-thread-empty strong {
+ color: var(--ink);
+}
+
+.chat-thread-empty p {
+ margin: 8px 0 0;
+}
+
+@media (max-width: 720px) {
+ .chat-page {
+ width: 100%;
+ height: calc(100dvh - 184px);
+ min-height: 0;
+ padding: 0;
+ }
+
+ .chat-workspace {
+ display: block;
+ height: 100%;
+ min-height: 0;
+ border-right: 0;
+ border-left: 0;
+ box-shadow: none;
+ }
+
+ .chat-sidebar,
+ .chat-thread {
+ height: 100%;
+ }
+
+ .chat-sidebar.has-selection {
+ display: none;
+ }
+
+ .chat-thread:not(.is-open) {
+ display: none;
+ }
+
+ .chat-thread.is-open {
+ display: grid;
+ }
+
+ .chat-back {
+ display: inline-flex;
+ align-items: center;
+ padding: 0 10px;
+ }
+
+ .chat-thread-header {
+ justify-content: flex-start;
+ }
+
+ .chat-message {
+ width: 88%;
+ }
+
+ .chat-compose-panel {
+ top: 143px;
+ max-height: calc(100% - 156px);
+ }
+
+ .chat-member-picker {
+ max-height: min(280px, calc(100dvh - 440px));
+ }
+
+ .chat-send-form {
+ grid-template-columns: minmax(0, 1fr) auto;
+ padding: 10px;
+ }
+
+ .chat-send-form .button {
+ min-width: 72px;
+ padding: 0 12px;
+ }
+}
+
/* Application portal: sections, status tabs, and the cohort question editor. */
.portal-section {
margin-top: 40px;
diff --git a/app/members/[id]/page.tsx b/app/members/[id]/page.tsx
index 7b5ecd8..3489d4b 100644
--- a/app/members/[id]/page.tsx
+++ b/app/members/[id]/page.tsx
@@ -52,6 +52,7 @@ export default async function MemberProfilePage({
session?.user?.id === member.id ||
(session?.user?.role === Role.ADMIN && session.user.status === UserStatus.ACTIVE);
const canNudge = session?.user?.status === UserStatus.ACTIVE && session.user.id !== member.id;
+ const canMessage = session?.user?.status === UserStatus.ACTIVE && session.user.id !== member.id;
const isOwnProfile = session?.user?.status === UserStatus.ACTIVE && session.user.id === member.id;
const isAdmin = session?.user?.role === Role.ADMIN && session.user.status === UserStatus.ACTIVE;
const badges = isAdmin ? await prisma.badge.findMany({ orderBy: { name: "asc" } }) : [];
@@ -116,6 +117,11 @@ export default async function MemberProfilePage({
/>
>
) : null}
+ {canMessage ? (
+
+ Message
+
+ ) : null}
diff --git a/app/message-indicator.tsx b/app/message-indicator.tsx
new file mode 100644
index 0000000..d0e1f53
--- /dev/null
+++ b/app/message-indicator.tsx
@@ -0,0 +1,105 @@
+"use client";
+
+import { useCallback, useEffect, useState } from "react";
+import { listLocalChatMessages } from "../lib/chat-client-db";
+
+const REFRESH_INTERVAL_MS = 15_000;
+
+async function pendingMessageIds() {
+ const messageIds = new Set
();
+ let cursor: string | null = null;
+
+ do {
+ const url = cursor ? `/api/chat/inbox?cursor=${encodeURIComponent(cursor)}` : "/api/chat/inbox";
+ const response = await fetch(url, { cache: "no-store" });
+ if (!response.ok) return messageIds;
+
+ const page = (await response.json()) as {
+ messages?: Array<{ id?: unknown }>;
+ nextCursor?: string | null;
+ };
+
+ for (const message of page.messages ?? []) {
+ if (typeof message.id === "string") messageIds.add(message.id);
+ }
+ cursor = page.nextCursor ?? null;
+ } while (cursor);
+
+ return messageIds;
+}
+
+export default function MessageIndicator({ userId }: Readonly<{ userId: string }>) {
+ const [unreadCount, setUnreadCount] = useState(0);
+
+ const refresh = useCallback(async () => {
+ if (typeof indexedDB === "undefined") return;
+
+ try {
+ const [localMessages, pendingIds] = await Promise.all([
+ listLocalChatMessages(userId),
+ pendingMessageIds(),
+ ]);
+ const unreadIds = new Set(
+ localMessages.filter((message) => message.unread).map((message) => message.id),
+ );
+ for (const messageId of Array.from(pendingIds)) unreadIds.add(messageId);
+ setUnreadCount(unreadIds.size);
+ } catch {
+ // Realtime and the messages inbox remain usable if a badge refresh fails.
+ }
+ }, [userId]);
+
+ useEffect(() => {
+ void refresh();
+ const channel =
+ typeof BroadcastChannel === "undefined"
+ ? null
+ : new BroadcastChannel(`shardup-chat-${userId}`);
+ if (channel) channel.onmessage = () => void refresh();
+ const interval = window.setInterval(() => void refresh(), REFRESH_INTERVAL_MS);
+ const onFocus = () => void refresh();
+ window.addEventListener("focus", onFocus);
+ window.addEventListener("online", onFocus);
+
+ return () => {
+ channel?.close();
+ window.clearInterval(interval);
+ window.removeEventListener("focus", onFocus);
+ window.removeEventListener("online", onFocus);
+ };
+ }, [refresh, userId]);
+
+ const label = unreadCount > 0 ? `Messages, ${unreadCount} unread` : "Messages";
+
+ return (
+
+ );
+}
diff --git a/app/messages/layout.tsx b/app/messages/layout.tsx
new file mode 100644
index 0000000..1560b38
--- /dev/null
+++ b/app/messages/layout.tsx
@@ -0,0 +1 @@
+export { default } from "../public-section-layout";
diff --git a/app/messages/messages-client.tsx b/app/messages/messages-client.tsx
new file mode 100644
index 0000000..2953fdc
--- /dev/null
+++ b/app/messages/messages-client.tsx
@@ -0,0 +1,500 @@
+"use client";
+
+import { useCallback, useEffect, useMemo, useRef, useState } from "react";
+import { useRouter } from "next/navigation";
+import {
+ createOptimisticChatMessage,
+ failOptimisticChatMessage,
+ listLocalChatMessages,
+ markChatAcknowledged,
+ markLocalConversationRead,
+ pendingChatAcknowledgements,
+ reconcileOptimisticChatMessage,
+ storeIncomingMessage,
+ type IncomingChatEnvelope,
+ type LocalChatMessage,
+ type LocalChatSender,
+} from "../../lib/chat-client-db";
+
+type MemberOption = LocalChatSender;
+type ClientConversationType = "DIRECT" | "GROUP";
+type ClientMemberRole = "OWNER" | "MEMBER";
+
+type ConversationSummary = {
+ id: string;
+ type: ClientConversationType;
+ title: string | null;
+ lastActivityAt: string;
+ members: Array;
+};
+
+type Props = {
+ currentUser: LocalChatSender;
+ conversations: ConversationSummary[];
+ members: MemberOption[];
+ initialConversationId: string | null;
+ initialRecipientId: string | null;
+};
+
+function conversationName(conversation: ConversationSummary, currentUserId: string) {
+ if (conversation.type === "GROUP") return conversation.title || "Group";
+ return (
+ conversation.members.find((member) => member.id !== currentUserId)?.name || "Direct message"
+ );
+}
+
+function isIncomingEnvelope(value: unknown): value is IncomingChatEnvelope {
+ if (!value || typeof value !== "object") return false;
+ const envelope = value as Partial;
+ return (
+ typeof envelope.id === "string" &&
+ typeof envelope.conversationId === "string" &&
+ typeof envelope.body === "string" &&
+ typeof envelope.createdAt === "string" &&
+ typeof envelope.signature === "string" &&
+ Boolean(envelope.sender && typeof envelope.sender.id === "string")
+ );
+}
+
+export default function MessagesClient({
+ currentUser,
+ conversations,
+ members,
+ initialConversationId,
+ initialRecipientId,
+}: Props) {
+ const router = useRouter();
+ const [selectedId, setSelectedId] = useState(
+ initialConversationId &&
+ conversations.some((conversation) => conversation.id === initialConversationId)
+ ? initialConversationId
+ : (conversations[0]?.id ?? null),
+ );
+ const [localMessages, setLocalMessages] = useState([]);
+ const [connection, setConnection] = useState<"connecting" | "live" | "polling">("connecting");
+ const [sendError, setSendError] = useState(null);
+ const [composeMode, setComposeMode] = useState<"closed" | "direct" | "group">("closed");
+ const [creating, setCreating] = useState(false);
+ const startedRecipient = useRef(false);
+ const broadcastRef = useRef(null);
+ const selectedIdRef = useRef(selectedId);
+
+ const refreshLocal = useCallback(async () => {
+ setLocalMessages(await listLocalChatMessages(currentUser.id));
+ }, [currentUser.id]);
+
+ const notifyLocalChange = useCallback(() => {
+ broadcastRef.current?.postMessage("changed");
+ void refreshLocal();
+ }, [refreshLocal]);
+
+ const flushAcknowledgements = useCallback(async () => {
+ const messageIds = await pendingChatAcknowledgements(currentUser.id);
+ if (messageIds.length === 0) return;
+
+ for (let offset = 0; offset < messageIds.length; offset += 100) {
+ const batch = messageIds.slice(offset, offset + 100);
+ const response = await fetch("/api/chat/ack", {
+ method: "POST",
+ headers: { "content-type": "application/json" },
+ body: JSON.stringify({ messageIds: batch }),
+ });
+ if (response.ok) await markChatAcknowledged(currentUser.id, batch);
+ }
+ }, [currentUser.id]);
+
+ const storeAndAcknowledge = useCallback(
+ async (envelopes: IncomingChatEnvelope[]) => {
+ for (const envelope of envelopes) await storeIncomingMessage(currentUser.id, envelope);
+ const openConversationId = selectedIdRef.current;
+ if (
+ openConversationId &&
+ envelopes.some((envelope) => envelope.conversationId === openConversationId)
+ ) {
+ await markLocalConversationRead(currentUser.id, openConversationId);
+ }
+ await flushAcknowledgements();
+ notifyLocalChange();
+ },
+ [currentUser.id, flushAcknowledgements, notifyLocalChange],
+ );
+
+ const recoverInbox = useCallback(async () => {
+ let cursor: string | null = null;
+ do {
+ const url = cursor
+ ? `/api/chat/inbox?cursor=${encodeURIComponent(cursor)}`
+ : "/api/chat/inbox";
+ const response = await fetch(url, { cache: "no-store" });
+ if (!response.ok) return;
+ const page = (await response.json()) as {
+ messages?: unknown[];
+ nextCursor?: string | null;
+ };
+ const envelopes = (page.messages ?? []).filter(isIncomingEnvelope);
+ if (envelopes.length > 0) await storeAndAcknowledge(envelopes);
+ cursor = page.nextCursor ?? null;
+ } while (cursor);
+ await flushAcknowledgements();
+ }, [flushAcknowledgements, storeAndAcknowledge]);
+
+ useEffect(() => {
+ void refreshLocal();
+ const channel = new BroadcastChannel(`shardup-chat-${currentUser.id}`);
+ channel.onmessage = () => void refreshLocal();
+ broadcastRef.current = channel;
+ return () => channel.close();
+ }, [currentUser.id, refreshLocal]);
+
+ useEffect(() => {
+ let socket: WebSocket | null = null;
+ let stopped = false;
+ let retryTimer: ReturnType | null = null;
+ let retryDelay = 1000;
+
+ const connect = async () => {
+ setConnection("connecting");
+ try {
+ const response = await fetch("/api/chat/token", { method: "POST" });
+ if (!response.ok) throw new Error("Token unavailable");
+ const credentials = (await response.json()) as { token: string; socketUrl: string | null };
+ if (!credentials.socketUrl) {
+ setConnection("polling");
+ await recoverInbox();
+ return;
+ }
+
+ socket = new WebSocket(credentials.socketUrl, ["shardup-chat", credentials.token]);
+ socket.onopen = () => {
+ retryDelay = 1000;
+ setConnection("live");
+ void recoverInbox();
+ };
+ socket.onmessage = (event) => {
+ try {
+ const payload = JSON.parse(String(event.data)) as { type?: string; envelope?: unknown };
+ if (payload.type === "message" && isIncomingEnvelope(payload.envelope)) {
+ void storeAndAcknowledge([payload.envelope]);
+ }
+ } catch {
+ // Ignore malformed transport frames; the inbox remains authoritative.
+ }
+ };
+ socket.onclose = () => {
+ if (stopped) return;
+ setConnection("polling");
+ retryTimer = setTimeout(() => {
+ retryDelay = Math.min(retryDelay * 2, 30_000);
+ void connect();
+ }, retryDelay);
+ };
+ } catch {
+ setConnection("polling");
+ await recoverInbox();
+ }
+ };
+
+ void connect();
+ const recovery = setInterval(() => void recoverInbox(), 15_000);
+ const onOnline = () => {
+ void recoverInbox();
+ if (!socket || socket.readyState === WebSocket.CLOSED) void connect();
+ };
+ window.addEventListener("online", onOnline);
+
+ return () => {
+ stopped = true;
+ clearInterval(recovery);
+ if (retryTimer) clearTimeout(retryTimer);
+ window.removeEventListener("online", onOnline);
+ socket?.close();
+ };
+ }, [currentUser.id, recoverInbox, storeAndAcknowledge]);
+
+ useEffect(() => {
+ selectedIdRef.current = selectedId;
+ if (!selectedId) return;
+ void markLocalConversationRead(currentUser.id, selectedId).then(notifyLocalChange);
+ void fetch("/api/chat/read", {
+ method: "POST",
+ headers: { "content-type": "application/json" },
+ body: JSON.stringify({ conversationId: selectedId }),
+ });
+ }, [currentUser.id, notifyLocalChange, selectedId]);
+
+ useEffect(() => {
+ if (!initialRecipientId || startedRecipient.current) return;
+ startedRecipient.current = true;
+ void (async () => {
+ const response = await fetch("/api/chat/conversations", {
+ method: "POST",
+ headers: { "content-type": "application/json" },
+ body: JSON.stringify({
+ type: "DIRECT",
+ recipientId: initialRecipientId,
+ }),
+ });
+ if (!response.ok) return;
+ const result = (await response.json()) as { conversationId: string };
+ setSelectedId(result.conversationId);
+ router.replace(`/messages?conversation=${encodeURIComponent(result.conversationId)}`);
+ router.refresh();
+ })();
+ }, [initialRecipientId, router]);
+
+ const selected = conversations.find((conversation) => conversation.id === selectedId) ?? null;
+ const selectedMessages = localMessages.filter((message) => message.conversationId === selectedId);
+ const latestByConversation = useMemo(() => {
+ const latest = new Map();
+ for (const message of localMessages) latest.set(message.conversationId, message);
+ return latest;
+ }, [localMessages]);
+
+ const selectConversation = (conversationId: string) => {
+ setSelectedId(conversationId);
+ router.replace(`/messages?conversation=${encodeURIComponent(conversationId)}`, {
+ scroll: false,
+ });
+ };
+
+ const sendMessage = async (event: React.FormEvent) => {
+ event.preventDefault();
+ if (!selectedId) return;
+ const form = event.currentTarget;
+ const formData = new FormData(form);
+ const body = String(formData.get("body") ?? "").trim();
+ if (!body) return;
+
+ setSendError(null);
+ form.reset();
+ const clientMessageId = crypto.randomUUID();
+ const createdAt = new Date().toISOString();
+ await createOptimisticChatMessage(currentUser.id, {
+ conversationId: selectedId,
+ clientMessageId,
+ sender: currentUser,
+ body,
+ createdAt,
+ });
+ notifyLocalChange();
+
+ try {
+ const response = await fetch("/api/chat/messages", {
+ method: "POST",
+ headers: { "content-type": "application/json" },
+ body: JSON.stringify({ conversationId: selectedId, clientMessageId, body }),
+ });
+ if (!response.ok) throw new Error("Message was not accepted");
+ const result = (await response.json()) as { messageId: string; createdAt: string };
+ await reconcileOptimisticChatMessage(
+ currentUser.id,
+ clientMessageId,
+ result.messageId,
+ result.createdAt,
+ );
+ notifyLocalChange();
+ router.refresh();
+ } catch {
+ await failOptimisticChatMessage(currentUser.id, clientMessageId);
+ setSendError("Message was not sent. Check your connection and try again.");
+ notifyLocalChange();
+ }
+ };
+
+ const createConversation = async (event: React.FormEvent) => {
+ event.preventDefault();
+ setCreating(true);
+ const form = event.currentTarget;
+ const formData = new FormData(form);
+ const payload =
+ composeMode === "direct"
+ ? { type: "DIRECT", recipientId: String(formData.get("recipientId")) }
+ : {
+ type: "GROUP",
+ title: String(formData.get("title")),
+ memberIds: formData.getAll("memberIds").map(String),
+ };
+
+ const response = await fetch("/api/chat/conversations", {
+ method: "POST",
+ headers: { "content-type": "application/json" },
+ body: JSON.stringify(payload),
+ });
+ setCreating(false);
+ if (!response.ok) {
+ setSendError("Conversation could not be created.");
+ return;
+ }
+ const result = (await response.json()) as { conversationId: string };
+ setComposeMode("closed");
+ setSelectedId(result.conversationId);
+ router.push(`/messages?conversation=${encodeURIComponent(result.conversationId)}`);
+ router.refresh();
+ };
+
+ return (
+
+
+
+
+
+ {selected ? (
+ <>
+
+
+ {selectedMessages.length ? (
+ selectedMessages.map((message) => (
+
+
+ {message.outgoing ? "You" : message.sender.name}
+
+
+ {message.body}
+ {message.status === "sending" ? Sending... : null}
+ {message.status === "failed" ? (
+ Not sent
+ ) : null}
+
+ ))
+ ) : (
+
+
No messages on this browser
+
Delivered history stays on the device that received it.
+
+ )}
+
+
+ {sendError ? {sendError}
: null}
+ >
+ ) : (
+
+
Select a conversation
+
Your delivered messages are kept only in this browser.
+
+ )}
+
+
+
+ );
+}
diff --git a/app/messages/page.tsx b/app/messages/page.tsx
new file mode 100644
index 0000000..12bccd8
--- /dev/null
+++ b/app/messages/page.tsx
@@ -0,0 +1,57 @@
+import { UserStatus } from "@/prisma-client";
+import { requireActiveUser } from "../../lib/guards";
+import { listChatConversations } from "../../lib/chat-service";
+import { memberDisplayName } from "../../lib/members";
+import { prisma } from "../../lib/prisma";
+import MessagesClient from "./messages-client";
+
+export const dynamic = "force-dynamic";
+
+export default async function MessagesPage({
+ searchParams,
+}: Readonly<{ searchParams?: { conversation?: string; recipient?: string } }>) {
+ const user = await requireActiveUser();
+ const [conversations, members] = await Promise.all([
+ listChatConversations(user.id),
+ prisma.user.findMany({
+ where: { status: UserStatus.ACTIVE, id: { not: user.id } },
+ orderBy: { name: "asc" },
+ select: {
+ id: true,
+ name: true,
+ email: true,
+ image: true,
+ profile: { select: { displayName: true, photoUrl: true } },
+ },
+ }),
+ ]);
+
+ return (
+ ({
+ id: member.id,
+ name: memberDisplayName(member),
+ image: member.profile?.photoUrl ?? member.image,
+ }))}
+ conversations={conversations.map((conversation) => ({
+ id: conversation.id,
+ type: conversation.type,
+ title: conversation.title,
+ lastActivityAt: conversation.lastActivityAt.toISOString(),
+ members: conversation.members.map((membership) => ({
+ id: membership.user.id,
+ name: memberDisplayName(membership.user),
+ image: membership.user.profile?.photoUrl ?? membership.user.image,
+ role: membership.role,
+ })),
+ }))}
+ />
+ );
+}
diff --git a/auth.ts b/auth.ts
index 5df2407..c54d33b 100644
--- a/auth.ts
+++ b/auth.ts
@@ -15,11 +15,12 @@ export const isGoogleOAuthConfigured = Boolean(
process.env.AUTH_GOOGLE_ID && process.env.AUTH_GOOGLE_SECRET,
);
-export type LocalDevRole = "admin" | "member";
+export type LocalDevRole = "admin" | "member" | "active";
const LOCAL_DEV_IDENTITIES: Record = {
admin: { id: "local-dev-admin", email: "admin@shardup.local", name: "Local Admin" },
member: { id: "local-dev-member", email: "applicant@shardup.local", name: "Local Applicant" },
+ active: { id: "local-dev-active", email: "member@shardup.local", name: "Local Member" },
};
export function localDevProviderId(role: LocalDevRole) {
@@ -29,6 +30,7 @@ export function localDevProviderId(role: LocalDevRole) {
// Dev-only logins (admin + applicant) for local testing without Google. Off in production.
function makeLocalDevProvider(role: LocalDevRole) {
const isAdmin = role === "admin";
+ const isApplicant = role === "member";
const { id, email, name } = LOCAL_DEV_IDENTITIES[role];
return Credentials({
@@ -42,18 +44,18 @@ function makeLocalDevProvider(role: LocalDevRole) {
name,
role: isAdmin ? Role.ADMIN : Role.MEMBER,
// The dev applicant is a test account — reset to PENDING on every login.
- status: isAdmin ? UserStatus.ACTIVE : UserStatus.PENDING,
+ status: isApplicant ? UserStatus.PENDING : UserStatus.ACTIVE,
},
create: {
email,
name,
role: isAdmin ? Role.ADMIN : Role.MEMBER,
- status: isAdmin ? UserStatus.ACTIVE : UserStatus.PENDING,
+ status: isApplicant ? UserStatus.PENDING : UserStatus.ACTIVE,
},
});
// Clear the old application so each login starts a fresh draft. Admins don't get one.
- if (!isAdmin) {
+ if (isApplicant) {
await prisma.application.deleteMany({ where: { userId: user.id } });
}
@@ -82,7 +84,13 @@ const providers = [
}),
]
: []),
- ...(isLocalDevAuthEnabled ? [makeLocalDevProvider("member"), makeLocalDevProvider("admin")] : []),
+ ...(isLocalDevAuthEnabled
+ ? [
+ makeLocalDevProvider("member"),
+ makeLocalDevProvider("active"),
+ makeLocalDevProvider("admin"),
+ ]
+ : []),
];
export const { handlers, signIn, signOut, auth } = NextAuth({
diff --git a/docs/chat-gateway.md b/docs/chat-gateway.md
new file mode 100644
index 0000000..c8d5c67
--- /dev/null
+++ b/docs/chat-gateway.md
@@ -0,0 +1,79 @@
+# Self-Hosted Chat Gateway
+
+ShardUp chat uses PostgreSQL as a transient store-and-forward inbox. The gateway only forwards
+committed messages to currently connected browsers; it never stores message content. If the
+gateway is stopped or restarted, clients reconnect and recover their pending messages from the
+database.
+
+Redis is intentionally not required for the first deployment. A single gateway process keeps a
+`userId -> WebSocket connections` map in memory. Add Redis Pub/Sub only when running more than one
+gateway replica, so a publish received by one replica can reach sockets connected to another.
+
+## Local Development
+
+Add the chat variables from `.env.example` to `.env.local`, then run the app and gateway in separate
+terminals:
+
+```bash
+npm run dev
+npm run chat:gateway
+```
+
+Use two browser profiles and the development member/admin accounts to test delivery. When the
+gateway is not running, chat still works through the 15-second pending-inbox recovery loop.
+
+## Oracle VM Deployment
+
+The gateway can run beside the existing Piston deployment as an isolated Node 20 container. Copy
+`infra/chat-gateway.mjs`, `infra/chat-gateway.Dockerfile`, and
+`infra/chat-gateway-compose.yaml` into `/opt/shardup-chat/infra`, then build it:
+
+```bash
+cd /opt/shardup-chat
+docker compose -f infra/chat-gateway-compose.yaml up -d --build
+```
+
+Create `/etc/shardup-chat.env` with permissions `0600`:
+
+```text
+CHAT_GATEWAY_PORT=8787
+CHAT_REALTIME_TOKEN_SECRET=
+CHAT_REALTIME_PUBLISH_SECRET=
+CHAT_ALLOWED_ORIGINS=https://YOUR_SHARDUP_DOMAIN
+```
+
+Add a path handler before the existing judge authorization handler:
+
+```caddy
+YOUR_VM_DOMAIN {
+ handle_path /chat/* {
+ reverse_proxy 127.0.0.1:8787
+ }
+
+ # Existing judge handlers continue below this route.
+}
+```
+
+Only ports 80/443 should be public; do not expose port 8787 through the Oracle security list.
+
+Set these variables in Vercel and redeploy:
+
+```text
+CHAT_REALTIME_TOKEN_SECRET=
+CHAT_REALTIME_PUBLISH_SECRET=
+CHAT_REALTIME_INTERNAL_URL=https://YOUR_VM_DOMAIN/chat
+NEXT_PUBLIC_CHAT_WS_URL=wss://YOUR_VM_DOMAIN/chat/socket
+CRON_SECRET=
+```
+
+`CHAT_REALTIME_PUBLISH_SECRET` protects the server-only `/publish` endpoint. Browsers receive only a
+short-lived, user-bound connection token. The gateway rejects all client-originated message frames.
+
+## Operations
+
+- Health: `GET https://YOUR_VM_DOMAIN/chat/health`
+- Logs: `docker logs -f shardup_chat_gateway`
+- Restart: `docker restart shardup_chat_gateway`
+- Rotate both chat secrets in the gateway and Vercel together.
+- Pending messages expire after 30 days. Delivered message payloads are deleted after browser
+ IndexedDB persistence and acknowledgement.
\ No newline at end of file
diff --git a/infra/chat-gateway-compose.yaml b/infra/chat-gateway-compose.yaml
new file mode 100644
index 0000000..eb13046
--- /dev/null
+++ b/infra/chat-gateway-compose.yaml
@@ -0,0 +1,18 @@
+services:
+ gateway:
+ build:
+ context: ..
+ dockerfile: infra/chat-gateway.Dockerfile
+ container_name: shardup_chat_gateway
+ restart: unless-stopped
+ env_file:
+ - /etc/shardup-chat.env
+ ports:
+ - "127.0.0.1:8787:8787"
+ read_only: true
+ tmpfs:
+ - /tmp
+ cap_drop:
+ - ALL
+ security_opt:
+ - no-new-privileges:true
diff --git a/infra/chat-gateway.Dockerfile b/infra/chat-gateway.Dockerfile
new file mode 100644
index 0000000..8bbb1aa
--- /dev/null
+++ b/infra/chat-gateway.Dockerfile
@@ -0,0 +1,18 @@
+FROM node:20-alpine
+
+WORKDIR /app
+
+RUN npm init -y \
+ && npm install --omit=dev jose@6.2.8 ws@8.21.0 \
+ && npm cache clean --force
+
+COPY infra/chat-gateway.mjs /app/server.mjs
+
+USER node
+
+EXPOSE 8787
+
+HEALTHCHECK --interval=30s --timeout=3s --start-period=5s --retries=3 \
+ CMD wget -qO- http://127.0.0.1:8787/health >/dev/null || exit 1
+
+CMD ["node", "/app/server.mjs"]
\ No newline at end of file
diff --git a/infra/chat-gateway.mjs b/infra/chat-gateway.mjs
new file mode 100644
index 0000000..d6a3077
--- /dev/null
+++ b/infra/chat-gateway.mjs
@@ -0,0 +1,205 @@
+import { createServer } from "node:http";
+import { timingSafeEqual } from "node:crypto";
+import { jwtVerify } from "jose";
+import { WebSocket, WebSocketServer } from "ws";
+
+const port = Number(process.env.CHAT_GATEWAY_PORT ?? 8787);
+const publishSecret = process.env.CHAT_REALTIME_PUBLISH_SECRET;
+const tokenSecret = process.env.CHAT_REALTIME_TOKEN_SECRET ?? process.env.AUTH_SECRET;
+const allowedOrigins = new Set(
+ (process.env.CHAT_ALLOWED_ORIGINS ?? "")
+ .split(",")
+ .map((origin) => origin.trim())
+ .filter(Boolean),
+);
+const socketsByUser = new Map();
+const encoder = new TextEncoder();
+
+if (!publishSecret || !tokenSecret) {
+ throw new Error("CHAT_REALTIME_PUBLISH_SECRET and CHAT_REALTIME_TOKEN_SECRET are required");
+}
+
+function json(response, status, body) {
+ response.writeHead(status, { "content-type": "application/json", "cache-control": "no-store" });
+ response.end(JSON.stringify(body));
+}
+
+function safeSecretMatch(value, expected) {
+ const actualBuffer = Buffer.from(value ?? "");
+ const expectedBuffer = Buffer.from(expected);
+ return (
+ actualBuffer.length === expectedBuffer.length && timingSafeEqual(actualBuffer, expectedBuffer)
+ );
+}
+
+function readJson(request) {
+ return new Promise((resolve, reject) => {
+ const chunks = [];
+ let size = 0;
+
+ request.on("data", (chunk) => {
+ size += chunk.length;
+ if (size > 128 * 1024) {
+ reject(new Error("Payload too large"));
+ request.destroy();
+ return;
+ }
+ chunks.push(chunk);
+ });
+ request.on("end", () => {
+ try {
+ resolve(JSON.parse(Buffer.concat(chunks).toString("utf8")));
+ } catch (error) {
+ reject(error);
+ }
+ });
+ request.on("error", reject);
+ });
+}
+
+function validPublish(body) {
+ return (
+ body &&
+ Array.isArray(body.recipientIds) &&
+ body.recipientIds.length > 0 &&
+ body.recipientIds.length <= 20 &&
+ body.recipientIds.every((id) => typeof id === "string" && id.length > 0 && id.length <= 128) &&
+ body.envelope &&
+ typeof body.envelope.id === "string" &&
+ typeof body.envelope.body === "string" &&
+ body.envelope.body.length <= 4000
+ );
+}
+
+const server = createServer(async (request, response) => {
+ const url = new URL(request.url ?? "/", `http://${request.headers.host ?? "localhost"}`);
+
+ if (request.method === "GET" && url.pathname === "/health") {
+ json(response, 200, {
+ ok: true,
+ users: socketsByUser.size,
+ connections: Array.from(socketsByUser.values()).reduce(
+ (total, sockets) => total + sockets.size,
+ 0,
+ ),
+ });
+ return;
+ }
+
+ if (request.method !== "POST" || url.pathname !== "/publish") {
+ json(response, 404, { error: "Not found" });
+ return;
+ }
+
+ const authorization = request.headers.authorization;
+ const bearer = authorization?.startsWith("Bearer ") ? authorization.slice(7) : "";
+ if (!safeSecretMatch(bearer, publishSecret)) {
+ json(response, 401, { error: "Unauthorized" });
+ return;
+ }
+
+ try {
+ const body = await readJson(request);
+ if (!validPublish(body)) {
+ json(response, 400, { error: "Invalid publish payload" });
+ return;
+ }
+
+ const payload = JSON.stringify({ type: "message", envelope: body.envelope });
+ let delivered = 0;
+ for (const recipientId of new Set(body.recipientIds)) {
+ for (const socket of socketsByUser.get(recipientId) ?? []) {
+ if (socket.readyState === WebSocket.OPEN) {
+ socket.send(payload);
+ delivered += 1;
+ }
+ }
+ }
+
+ json(response, 200, { delivered });
+ } catch (error) {
+ json(response, 400, { error: error instanceof Error ? error.message : "Invalid request" });
+ }
+});
+
+const webSocketServer = new WebSocketServer({
+ noServer: true,
+ maxPayload: 1024,
+ handleProtocols(protocols) {
+ return protocols.has("shardup-chat") ? "shardup-chat" : false;
+ },
+});
+
+server.on("upgrade", async (request, socket, head) => {
+ try {
+ const url = new URL(request.url ?? "/", `http://${request.headers.host ?? "localhost"}`);
+ const origin = request.headers.origin;
+ const protocols = (request.headers["sec-websocket-protocol"] ?? "")
+ .split(",")
+ .map((protocol) => protocol.trim());
+ const token = protocols.find((protocol) => protocol !== "shardup-chat");
+
+ if (
+ url.pathname !== "/socket" ||
+ !protocols.includes("shardup-chat") ||
+ !token ||
+ (allowedOrigins.size > 0 && (!origin || !allowedOrigins.has(origin)))
+ ) {
+ socket.destroy();
+ return;
+ }
+
+ const { payload } = await jwtVerify(token, encoder.encode(tokenSecret), {
+ algorithms: ["HS256"],
+ issuer: "shardup-web",
+ audience: "shardup-chat-gateway",
+ });
+
+ if (payload.scope !== "chat:connect" || typeof payload.sub !== "string") {
+ socket.destroy();
+ return;
+ }
+
+ request.chatUserId = payload.sub;
+ webSocketServer.handleUpgrade(request, socket, head, (webSocket) => {
+ webSocketServer.emit("connection", webSocket, request);
+ });
+ } catch {
+ socket.destroy();
+ }
+});
+
+webSocketServer.on("connection", (socket, request) => {
+ const userId = request.chatUserId;
+ const userSockets = socketsByUser.get(userId) ?? new Set();
+ socket.isAlive = true;
+ userSockets.add(socket);
+ socketsByUser.set(userId, userSockets);
+
+ socket.on("pong", () => {
+ socket.isAlive = true;
+ });
+ socket.on("message", () => {
+ socket.close(1008, "Client publishing is not allowed");
+ });
+ socket.on("close", () => {
+ userSockets.delete(socket);
+ if (userSockets.size === 0) socketsByUser.delete(userId);
+ });
+});
+
+const heartbeat = setInterval(() => {
+ for (const socket of webSocketServer.clients) {
+ if (!socket.isAlive) {
+ socket.terminate();
+ continue;
+ }
+ socket.isAlive = false;
+ socket.ping();
+ }
+}, 30_000);
+
+server.on("close", () => clearInterval(heartbeat));
+server.listen(port, "0.0.0.0", () => {
+ console.log(`ShardUp chat gateway listening on :${port}`);
+});
diff --git a/lib/chat-auth.ts b/lib/chat-auth.ts
new file mode 100644
index 0000000..6a3e4a1
--- /dev/null
+++ b/lib/chat-auth.ts
@@ -0,0 +1,12 @@
+import { UserStatus } from "@/prisma-client";
+import { auth } from "../auth";
+
+export async function activeChatUser() {
+ const session = await auth();
+
+ if (!session?.user || session.user.status !== UserStatus.ACTIVE) {
+ return null;
+ }
+
+ return session.user;
+}
diff --git a/lib/chat-client-db.ts b/lib/chat-client-db.ts
new file mode 100644
index 0000000..500ec43
--- /dev/null
+++ b/lib/chat-client-db.ts
@@ -0,0 +1,178 @@
+import { openDB, type DBSchema, type IDBPDatabase } from "idb";
+
+export type LocalChatSender = {
+ id: string;
+ name: string;
+ image: string | null;
+};
+
+export type LocalChatMessage = {
+ id: string;
+ conversationId: string;
+ sender: LocalChatSender;
+ body: string;
+ createdAt: string;
+ signature: string | null;
+ outgoing: boolean;
+ status: "sending" | "sent" | "failed" | "received";
+ unread: boolean;
+ ackPending: 0 | 1;
+ clientMessageId: string | null;
+};
+
+export type IncomingChatEnvelope = {
+ id: string;
+ conversationId: string;
+ sender: LocalChatSender;
+ body: string;
+ createdAt: string;
+ signature: string;
+};
+
+interface ChatDatabase extends DBSchema {
+ messages: {
+ key: string;
+ value: LocalChatMessage;
+ indexes: {
+ "by-conversation": string;
+ "by-ack-pending": number;
+ };
+ };
+}
+
+const databases = new Map>>();
+
+function databaseFor(userId: string) {
+ let database = databases.get(userId);
+ if (!database) {
+ database = openDB(`shardup-chat-${userId}`, 1, {
+ upgrade(db) {
+ const messages = db.createObjectStore("messages", { keyPath: "id" });
+ messages.createIndex("by-conversation", "conversationId");
+ messages.createIndex("by-ack-pending", "ackPending");
+ },
+ });
+ databases.set(userId, database);
+ }
+ return database;
+}
+
+export async function storeIncomingMessage(userId: string, envelope: IncomingChatEnvelope) {
+ const db = await databaseFor(userId);
+ const transaction = db.transaction("messages", "readwrite");
+ const existing = await transaction.store.get(envelope.id);
+
+ if (!existing) {
+ await transaction.store.put({
+ ...envelope,
+ outgoing: false,
+ status: "received",
+ unread: true,
+ ackPending: 1,
+ clientMessageId: null,
+ });
+ } else if (!existing.outgoing && !existing.ackPending) {
+ await transaction.store.put({ ...existing, ackPending: 1 });
+ }
+
+ await transaction.done;
+}
+
+export async function listLocalChatMessages(userId: string) {
+ const db = await databaseFor(userId);
+ const messages = await db.getAll("messages");
+ return messages.sort(
+ (first, second) =>
+ first.createdAt.localeCompare(second.createdAt) || first.id.localeCompare(second.id),
+ );
+}
+
+export async function createOptimisticChatMessage(
+ userId: string,
+ input: {
+ conversationId: string;
+ clientMessageId: string;
+ sender: LocalChatSender;
+ body: string;
+ createdAt: string;
+ },
+) {
+ const db = await databaseFor(userId);
+ const message: LocalChatMessage = {
+ id: `local:${input.clientMessageId}`,
+ conversationId: input.conversationId,
+ sender: input.sender,
+ body: input.body,
+ createdAt: input.createdAt,
+ signature: null,
+ outgoing: true,
+ status: "sending",
+ unread: false,
+ ackPending: 0,
+ clientMessageId: input.clientMessageId,
+ };
+ await db.put("messages", message);
+ return message;
+}
+
+export async function reconcileOptimisticChatMessage(
+ userId: string,
+ clientMessageId: string,
+ messageId: string,
+ createdAt: string,
+) {
+ const db = await databaseFor(userId);
+ const transaction = db.transaction("messages", "readwrite");
+ const temporaryId = `local:${clientMessageId}`;
+ const existing = await transaction.store.get(temporaryId);
+
+ if (existing) {
+ await transaction.store.delete(temporaryId);
+ await transaction.store.put({
+ ...existing,
+ id: messageId,
+ createdAt,
+ status: "sent",
+ });
+ }
+
+ await transaction.done;
+}
+
+export async function failOptimisticChatMessage(userId: string, clientMessageId: string) {
+ const db = await databaseFor(userId);
+ const id = `local:${clientMessageId}`;
+ const existing = await db.get("messages", id);
+ if (existing) await db.put("messages", { ...existing, status: "failed" });
+}
+
+export async function pendingChatAcknowledgements(userId: string) {
+ const db = await databaseFor(userId);
+ const messages = await db.getAllFromIndex("messages", "by-ack-pending", 1);
+ return messages.filter((message) => !message.outgoing).map((message) => message.id);
+}
+
+export async function markChatAcknowledged(userId: string, messageIds: string[]) {
+ const db = await databaseFor(userId);
+ const transaction = db.transaction("messages", "readwrite");
+
+ for (const messageId of messageIds) {
+ const message = await transaction.store.get(messageId);
+ if (message) await transaction.store.put({ ...message, ackPending: 0 });
+ }
+
+ await transaction.done;
+}
+
+export async function markLocalConversationRead(userId: string, conversationId: string) {
+ const db = await databaseFor(userId);
+ const transaction = db.transaction("messages", "readwrite");
+ let cursor = await transaction.store.index("by-conversation").openCursor(conversationId);
+
+ while (cursor) {
+ if (cursor.value.unread) await cursor.update({ ...cursor.value, unread: false });
+ cursor = await cursor.continue();
+ }
+
+ await transaction.done;
+}
diff --git a/lib/chat-http.ts b/lib/chat-http.ts
new file mode 100644
index 0000000..b9e5258
--- /dev/null
+++ b/lib/chat-http.ts
@@ -0,0 +1,29 @@
+import { NextResponse } from "next/server";
+import { ChatServiceError } from "./chat-service";
+
+export function chatErrorResponse(error: unknown) {
+ if (error instanceof ChatServiceError) {
+ const status =
+ error.code === "RATE_LIMITED"
+ ? 429
+ : error.code === "FORBIDDEN"
+ ? 403
+ : error.code === "NOT_FOUND"
+ ? 404
+ : 400;
+ const response = NextResponse.json({ error: error.code, message: error.message }, { status });
+
+ if (error.retryAfterSeconds) {
+ response.headers.set("Retry-After", String(error.retryAfterSeconds));
+ }
+
+ return response;
+ }
+
+ console.error("Chat request failed", error);
+ return NextResponse.json({ error: "INTERNAL_ERROR" }, { status: 500 });
+}
+
+export function unauthorizedChatResponse() {
+ return NextResponse.json({ error: "UNAUTHORIZED" }, { status: 401 });
+}
diff --git a/lib/chat-rate-limit.ts b/lib/chat-rate-limit.ts
new file mode 100644
index 0000000..414d32d
--- /dev/null
+++ b/lib/chat-rate-limit.ts
@@ -0,0 +1,54 @@
+import { ChatRateLimitScope, type Prisma } from "@/prisma-client";
+
+type RateLimitConfig = {
+ limit: number;
+ windowMs: number;
+};
+
+const CHAT_RATE_LIMITS: Record = {
+ [ChatRateLimitScope.MESSAGE_SEND]: { limit: 20, windowMs: 60_000 },
+ [ChatRateLimitScope.CONVERSATION_CREATE]: { limit: 5, windowMs: 60 * 60_000 },
+};
+
+export type ChatRateLimitResult =
+ | { allowed: true; remaining: number }
+ | { allowed: false; retryAfterSeconds: number };
+
+export async function consumeChatRateLimit(
+ transaction: Prisma.TransactionClient,
+ userId: string,
+ scope: ChatRateLimitScope,
+ now: Date = new Date(),
+): Promise {
+ const config = CHAT_RATE_LIMITS[scope];
+ const existing = await transaction.chatRateLimit.findUnique({
+ where: { userId_scope: { userId, scope } },
+ });
+
+ if (!existing || existing.expiresAt <= now) {
+ const expiresAt = new Date(now.getTime() + config.windowMs);
+ await transaction.chatRateLimit.upsert({
+ where: { userId_scope: { userId, scope } },
+ create: { userId, scope, windowStart: now, count: 1, expiresAt },
+ update: { windowStart: now, count: 1, expiresAt },
+ });
+ return { allowed: true, remaining: config.limit - 1 };
+ }
+
+ if (existing.count >= config.limit) {
+ return {
+ allowed: false,
+ retryAfterSeconds: Math.max(
+ 1,
+ Math.ceil((existing.expiresAt.getTime() - now.getTime()) / 1000),
+ ),
+ };
+ }
+
+ await transaction.chatRateLimit.update({
+ where: { id: existing.id },
+ data: { count: { increment: 1 } },
+ });
+
+ return { allowed: true, remaining: config.limit - existing.count - 1 };
+}
diff --git a/lib/chat-realtime.ts b/lib/chat-realtime.ts
new file mode 100644
index 0000000..2b29f06
--- /dev/null
+++ b/lib/chat-realtime.ts
@@ -0,0 +1,78 @@
+import { SignJWT } from "jose";
+import type { ChatEnvelope } from "./chat-service";
+
+const CHAT_TOKEN_ISSUER = "shardup-web";
+const CHAT_TOKEN_AUDIENCE = "shardup-chat-gateway";
+
+function tokenSecret() {
+ const value = process.env.CHAT_REALTIME_TOKEN_SECRET ?? process.env.AUTH_SECRET;
+
+ if (!value) {
+ throw new Error("CHAT_REALTIME_TOKEN_SECRET or AUTH_SECRET is required");
+ }
+
+ return new TextEncoder().encode(value);
+}
+
+export async function createChatSocketToken(userId: string) {
+ return new SignJWT({ scope: "chat:connect" })
+ .setProtectedHeader({ alg: "HS256", typ: "JWT" })
+ .setSubject(userId)
+ .setIssuer(CHAT_TOKEN_ISSUER)
+ .setAudience(CHAT_TOKEN_AUDIENCE)
+ .setIssuedAt()
+ .setExpirationTime("10m")
+ .sign(tokenSecret());
+}
+
+export type ChatPublishResult = {
+ attempted: number;
+ delivered: number;
+ unavailable: boolean;
+};
+
+export async function publishChatEnvelope(
+ recipientIds: string[],
+ envelope: ChatEnvelope,
+): Promise {
+ const baseUrl = process.env.CHAT_REALTIME_INTERNAL_URL;
+ const secret = process.env.CHAT_REALTIME_PUBLISH_SECRET;
+
+ if (!baseUrl || !secret) {
+ return { attempted: recipientIds.length, delivered: 0, unavailable: true };
+ }
+
+ try {
+ const response = await fetch(`${baseUrl.replace(/\/$/, "")}/publish`, {
+ method: "POST",
+ headers: {
+ authorization: `Bearer ${secret}`,
+ "content-type": "application/json",
+ },
+ body: JSON.stringify({ recipientIds, envelope }),
+ cache: "no-store",
+ signal: AbortSignal.timeout(3000),
+ });
+
+ if (!response.ok) {
+ console.error("Chat realtime publish failed", {
+ messageId: envelope.id,
+ status: response.status,
+ });
+ return { attempted: recipientIds.length, delivered: 0, unavailable: true };
+ }
+
+ const result = (await response.json()) as { delivered?: number };
+ return {
+ attempted: recipientIds.length,
+ delivered: typeof result.delivered === "number" ? result.delivered : 0,
+ unavailable: false,
+ };
+ } catch (error) {
+ console.error("Chat realtime gateway unavailable", {
+ messageId: envelope.id,
+ error: error instanceof Error ? error.message : "Unknown error",
+ });
+ return { attempted: recipientIds.length, delivered: 0, unavailable: true };
+ }
+}
diff --git a/lib/chat-service.ts b/lib/chat-service.ts
new file mode 100644
index 0000000..0d9486b
--- /dev/null
+++ b/lib/chat-service.ts
@@ -0,0 +1,573 @@
+import {
+ ChatConversationType,
+ ChatMemberRole,
+ ChatRateLimitScope,
+ Prisma,
+ UserStatus,
+} from "@/prisma-client";
+import {
+ CHAT_SEND_RECEIPT_TTL_HOURS,
+ CHAT_UNDELIVERED_TTL_DAYS,
+ allChatParticipantsActive,
+ chatRecipientIds,
+ directConversationKey,
+} from "./chat";
+import { consumeChatRateLimit } from "./chat-rate-limit";
+import { signChatEnvelope, verifyChatEnvelope } from "./chat-signing";
+import { memberDisplayName } from "./members";
+import { prisma } from "./prisma";
+
+const CHAT_INBOX_PAGE_SIZE = 50;
+
+export type ChatServiceErrorCode =
+ | "FORBIDDEN"
+ | "INVALID_PARTICIPANTS"
+ | "NOT_FOUND"
+ | "RATE_LIMITED";
+
+export class ChatServiceError extends Error {
+ constructor(
+ readonly code: ChatServiceErrorCode,
+ message: string,
+ readonly retryAfterSeconds?: number,
+ ) {
+ super(message);
+ this.name = "ChatServiceError";
+ }
+}
+
+export type ChatEnvelope = {
+ id: string;
+ conversationId: string;
+ sender: {
+ id: string;
+ name: string;
+ image: string | null;
+ };
+ body: string;
+ createdAt: string;
+ signature: string;
+};
+
+type EnqueueChatMessageInput = {
+ conversationId: string;
+ senderId: string;
+ clientMessageId: string;
+ body: string;
+};
+
+function addMilliseconds(date: Date, milliseconds: number) {
+ return new Date(date.getTime() + milliseconds);
+}
+
+function isPrismaError(error: unknown, code: string) {
+ return error instanceof Prisma.PrismaClientKnownRequestError && error.code === code;
+}
+
+async function serializableTransaction(
+ operation: (transaction: Prisma.TransactionClient) => Promise,
+) {
+ for (let attempt = 0; attempt < 3; attempt += 1) {
+ try {
+ return await prisma.$transaction(operation, {
+ isolationLevel: Prisma.TransactionIsolationLevel.Serializable,
+ });
+ } catch (error) {
+ if (!isPrismaError(error, "P2034") || attempt === 2) {
+ throw error;
+ }
+ }
+ }
+
+ throw new Error("Chat transaction retry limit reached");
+}
+
+async function enforceRateLimit(
+ transaction: Prisma.TransactionClient,
+ userId: string,
+ scope: ChatRateLimitScope,
+ now: Date,
+) {
+ const rateLimit = await consumeChatRateLimit(transaction, userId, scope, now);
+
+ if (!rateLimit.allowed) {
+ throw new ChatServiceError(
+ "RATE_LIMITED",
+ "Too many chat requests",
+ rateLimit.retryAfterSeconds,
+ );
+ }
+}
+
+export async function createDirectConversation(creatorId: string, recipientId: string) {
+ if (creatorId === recipientId) {
+ throw new ChatServiceError("INVALID_PARTICIPANTS", "Choose another member");
+ }
+
+ const directKey = directConversationKey(creatorId, recipientId);
+
+ try {
+ return await serializableTransaction(async (transaction) => {
+ const existing = await transaction.chatConversation.findUnique({ where: { directKey } });
+
+ if (existing) {
+ return existing;
+ }
+
+ const recipient = await transaction.user.findFirst({
+ where: { id: recipientId, status: UserStatus.ACTIVE },
+ select: { id: true },
+ });
+
+ if (!recipient) {
+ throw new ChatServiceError("INVALID_PARTICIPANTS", "Member is unavailable");
+ }
+
+ const now = new Date();
+ await enforceRateLimit(transaction, creatorId, ChatRateLimitScope.CONVERSATION_CREATE, now);
+
+ return await transaction.chatConversation.create({
+ data: {
+ type: ChatConversationType.DIRECT,
+ directKey,
+ createdById: creatorId,
+ members: {
+ create: [
+ { userId: creatorId, role: ChatMemberRole.MEMBER },
+ { userId: recipient.id, role: ChatMemberRole.MEMBER },
+ ],
+ },
+ },
+ });
+ });
+ } catch (error) {
+ if (isPrismaError(error, "P2002")) {
+ const concurrent = await prisma.chatConversation.findUnique({ where: { directKey } });
+ if (concurrent) return concurrent;
+ }
+ throw error;
+ }
+}
+
+export async function createGroupConversation(
+ creatorId: string,
+ title: string,
+ requestedMemberIds: string[],
+) {
+ const memberIds = Array.from(new Set(requestedMemberIds));
+
+ if (memberIds.includes(creatorId)) {
+ throw new ChatServiceError("INVALID_PARTICIPANTS", "The creator is already included");
+ }
+
+ return serializableTransaction(async (transaction) => {
+ const users = await transaction.user.findMany({
+ where: { id: { in: memberIds } },
+ select: { id: true, status: true },
+ });
+
+ if (!allChatParticipantsActive(users, memberIds.length)) {
+ throw new ChatServiceError("INVALID_PARTICIPANTS", "Every member must be active");
+ }
+
+ const now = new Date();
+ await enforceRateLimit(transaction, creatorId, ChatRateLimitScope.CONVERSATION_CREATE, now);
+
+ return transaction.chatConversation.create({
+ data: {
+ type: ChatConversationType.GROUP,
+ title,
+ createdById: creatorId,
+ members: {
+ create: [
+ { userId: creatorId, role: ChatMemberRole.OWNER },
+ ...memberIds.map((userId) => ({ userId, role: ChatMemberRole.MEMBER })),
+ ],
+ },
+ },
+ });
+ });
+}
+
+export async function listChatConversations(userId: string) {
+ return prisma.chatConversation.findMany({
+ where: { members: { some: { userId, leftAt: null } } },
+ orderBy: { lastActivityAt: "desc" },
+ include: {
+ members: {
+ where: { leftAt: null },
+ orderBy: { joinedAt: "asc" },
+ include: {
+ user: {
+ select: {
+ id: true,
+ name: true,
+ email: true,
+ image: true,
+ status: true,
+ profile: { select: { displayName: true, photoUrl: true } },
+ },
+ },
+ },
+ },
+ },
+ });
+}
+
+function receiptResult(receipt: { messageId: string; conversationId: string; createdAt: Date }) {
+ return {
+ duplicate: true as const,
+ messageId: receipt.messageId,
+ conversationId: receipt.conversationId,
+ createdAt: receipt.createdAt.toISOString(),
+ recipientIds: [] as string[],
+ envelope: null,
+ };
+}
+
+function pendingMessageResult(message: { id: string; conversationId: string; createdAt: Date }) {
+ return receiptResult({
+ messageId: message.id,
+ conversationId: message.conversationId,
+ createdAt: message.createdAt,
+ });
+}
+
+export async function enqueueChatMessage(input: EnqueueChatMessageInput) {
+ const priorReceipt = await prisma.chatSendReceipt.findUnique({
+ where: {
+ senderId_clientMessageId: {
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ },
+ },
+ });
+
+ if (priorReceipt) {
+ return receiptResult(priorReceipt);
+ }
+
+ const priorPendingMessage = await prisma.pendingChatMessage.findUnique({
+ where: {
+ senderId_clientMessageId: {
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ },
+ },
+ });
+
+ if (priorPendingMessage) {
+ return pendingMessageResult(priorPendingMessage);
+ }
+
+ try {
+ return await serializableTransaction(async (transaction) => {
+ const receipt = await transaction.chatSendReceipt.findUnique({
+ where: {
+ senderId_clientMessageId: {
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ },
+ },
+ });
+
+ if (receipt) {
+ return receiptResult(receipt);
+ }
+
+ const pendingMessage = await transaction.pendingChatMessage.findUnique({
+ where: {
+ senderId_clientMessageId: {
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ },
+ },
+ });
+
+ if (pendingMessage) {
+ return pendingMessageResult(pendingMessage);
+ }
+
+ const conversation = await transaction.chatConversation.findFirst({
+ where: {
+ id: input.conversationId,
+ members: { some: { userId: input.senderId, leftAt: null } },
+ },
+ include: { members: { where: { leftAt: null } } },
+ });
+
+ if (!conversation) {
+ throw new ChatServiceError("FORBIDDEN", "Conversation is unavailable");
+ }
+
+ const memberRecipientIds = chatRecipientIds(conversation, input.senderId);
+ const activeRecipients = await transaction.user.findMany({
+ where: { id: { in: memberRecipientIds }, status: UserStatus.ACTIVE },
+ select: { id: true },
+ });
+ const recipientIds = activeRecipients.map((recipient) => recipient.id);
+
+ if (recipientIds.length === 0) {
+ throw new ChatServiceError("INVALID_PARTICIPANTS", "No active recipient is available");
+ }
+
+ const now = new Date();
+ await enforceRateLimit(transaction, input.senderId, ChatRateLimitScope.MESSAGE_SEND, now);
+
+ const messageId = crypto.randomUUID();
+ const expiresAt = addMilliseconds(now, CHAT_UNDELIVERED_TTL_DAYS * 24 * 60 * 60_000);
+ const receiptExpiresAt = addMilliseconds(now, CHAT_SEND_RECEIPT_TTL_HOURS * 60 * 60_000);
+
+ const sender = await transaction.user.findUniqueOrThrow({
+ where: { id: input.senderId },
+ select: {
+ id: true,
+ name: true,
+ email: true,
+ image: true,
+ profile: { select: { displayName: true, photoUrl: true } },
+ },
+ });
+
+ await transaction.pendingChatMessage.create({
+ data: {
+ id: messageId,
+ conversationId: conversation.id,
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ body: input.body,
+ createdAt: now,
+ expiresAt,
+ deliveries: { create: recipientIds.map((recipientId) => ({ recipientId })) },
+ },
+ });
+ await transaction.chatSendReceipt.create({
+ data: {
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ messageId,
+ conversationId: conversation.id,
+ createdAt: now,
+ expiresAt: receiptExpiresAt,
+ },
+ });
+ await transaction.chatConversation.update({
+ where: { id: conversation.id },
+ data: { lastActivityAt: now },
+ });
+
+ const createdAt = now.toISOString();
+ const signable = {
+ id: messageId,
+ conversationId: conversation.id,
+ senderId: sender.id,
+ body: input.body,
+ createdAt,
+ };
+ const envelope: ChatEnvelope = {
+ id: messageId,
+ conversationId: conversation.id,
+ sender: {
+ id: sender.id,
+ name: memberDisplayName(sender),
+ image: sender.profile?.photoUrl ?? sender.image,
+ },
+ body: input.body,
+ createdAt,
+ signature: signChatEnvelope(signable),
+ };
+
+ return {
+ duplicate: false as const,
+ messageId,
+ conversationId: conversation.id,
+ createdAt,
+ recipientIds,
+ envelope,
+ };
+ });
+ } catch (error) {
+ if (isPrismaError(error, "P2002")) {
+ const receipt = await prisma.chatSendReceipt.findUnique({
+ where: {
+ senderId_clientMessageId: {
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ },
+ },
+ });
+ if (receipt) return receiptResult(receipt);
+ const pendingMessage = await prisma.pendingChatMessage.findUnique({
+ where: {
+ senderId_clientMessageId: {
+ senderId: input.senderId,
+ clientMessageId: input.clientMessageId,
+ },
+ },
+ });
+ if (pendingMessage) return pendingMessageResult(pendingMessage);
+ }
+ throw error;
+ }
+}
+
+function envelopeFromDelivery(delivery: {
+ message: {
+ id: string;
+ conversationId: string;
+ body: string;
+ createdAt: Date;
+ sender: {
+ id: string;
+ name: string | null;
+ email: string;
+ image: string | null;
+ profile: { displayName: string | null; photoUrl: string | null } | null;
+ };
+ };
+}): ChatEnvelope {
+ const message = delivery.message;
+ const createdAt = message.createdAt.toISOString();
+ const signable = {
+ id: message.id,
+ conversationId: message.conversationId,
+ senderId: message.sender.id,
+ body: message.body,
+ createdAt,
+ };
+
+ return {
+ id: message.id,
+ conversationId: message.conversationId,
+ sender: {
+ id: message.sender.id,
+ name: memberDisplayName(message.sender),
+ image: message.sender.profile?.photoUrl ?? message.sender.image,
+ },
+ body: message.body,
+ createdAt,
+ signature: signChatEnvelope(signable),
+ };
+}
+
+export async function listPendingChatMessages(userId: string, cursor?: string) {
+ const deliveries = await prisma.chatDelivery.findMany({
+ where: { recipientId: userId, message: { expiresAt: { gt: new Date() } } },
+ orderBy: [{ createdAt: "asc" }, { id: "asc" }],
+ take: CHAT_INBOX_PAGE_SIZE + 1,
+ ...(cursor ? { cursor: { id: cursor }, skip: 1 } : {}),
+ include: {
+ message: {
+ include: {
+ sender: {
+ select: {
+ id: true,
+ name: true,
+ email: true,
+ image: true,
+ profile: { select: { displayName: true, photoUrl: true } },
+ },
+ },
+ },
+ },
+ },
+ });
+ const hasMore = deliveries.length > CHAT_INBOX_PAGE_SIZE;
+ const page = deliveries.slice(0, CHAT_INBOX_PAGE_SIZE);
+
+ return {
+ messages: page.map(envelopeFromDelivery),
+ nextCursor: hasMore ? (page.at(-1)?.id ?? null) : null,
+ };
+}
+
+export async function acknowledgeChatMessages(userId: string, messageIds: string[]) {
+ return serializableTransaction(async (transaction) => {
+ const deleted = await transaction.chatDelivery.deleteMany({
+ where: { recipientId: userId, messageId: { in: messageIds } },
+ });
+ const payloads = await transaction.pendingChatMessage.deleteMany({
+ where: { id: { in: messageIds }, deliveries: { none: {} } },
+ });
+
+ return { acknowledged: deleted.count, payloadsDeleted: payloads.count };
+ });
+}
+
+export async function markChatRead(userId: string, conversationId: string) {
+ const updated = await prisma.chatMember.updateMany({
+ where: { userId, conversationId, leftAt: null },
+ data: { lastReadAt: new Date() },
+ });
+
+ if (updated.count === 0) {
+ throw new ChatServiceError("FORBIDDEN", "Conversation is unavailable");
+ }
+}
+
+type ReportChatMessageInput = {
+ conversationId: string;
+ reporterId: string;
+ messageId: string;
+ reportedSenderId: string;
+ body: string;
+ messageCreatedAt: Date;
+ note: string | null;
+ signature: string;
+};
+
+export async function reportChatMessage(input: ReportChatMessageInput) {
+ const membership = await prisma.chatMember.findFirst({
+ where: { conversationId: input.conversationId, userId: input.reporterId, leftAt: null },
+ select: { id: true },
+ });
+
+ if (!membership) {
+ throw new ChatServiceError("FORBIDDEN", "Conversation is unavailable");
+ }
+
+ const valid = verifyChatEnvelope(
+ {
+ id: input.messageId,
+ conversationId: input.conversationId,
+ senderId: input.reportedSenderId,
+ body: input.body,
+ createdAt: input.messageCreatedAt.toISOString(),
+ },
+ input.signature,
+ );
+
+ if (!valid) {
+ throw new ChatServiceError("FORBIDDEN", "Message authenticity check failed");
+ }
+
+ return prisma.chatReport.upsert({
+ where: { reporterId_messageId: { reporterId: input.reporterId, messageId: input.messageId } },
+ create: {
+ conversationId: input.conversationId,
+ reporterId: input.reporterId,
+ reportedSenderId: input.reportedSenderId,
+ messageId: input.messageId,
+ body: input.body,
+ messageCreatedAt: input.messageCreatedAt,
+ note: input.note,
+ },
+ update: { note: input.note },
+ });
+}
+
+export async function cleanupExpiredChatData(now: Date = new Date()) {
+ return serializableTransaction(async (transaction) => {
+ const [messages, receipts, rateLimits] = await Promise.all([
+ transaction.pendingChatMessage.deleteMany({ where: { expiresAt: { lte: now } } }),
+ transaction.chatSendReceipt.deleteMany({ where: { expiresAt: { lte: now } } }),
+ transaction.chatRateLimit.deleteMany({ where: { expiresAt: { lte: now } } }),
+ ]);
+
+ return {
+ messagesDeleted: messages.count,
+ receiptsDeleted: receipts.count,
+ rateLimitsDeleted: rateLimits.count,
+ };
+ });
+}
diff --git a/lib/chat-signing.ts b/lib/chat-signing.ts
new file mode 100644
index 0000000..ecc1103
--- /dev/null
+++ b/lib/chat-signing.ts
@@ -0,0 +1,40 @@
+import { createHmac, timingSafeEqual } from "node:crypto";
+
+export type SignableChatEnvelope = {
+ id: string;
+ conversationId: string;
+ senderId: string;
+ body: string;
+ createdAt: string;
+};
+
+function signingSecret() {
+ const secret = process.env.CHAT_MESSAGE_SIGNING_SECRET ?? process.env.AUTH_SECRET;
+
+ if (!secret) {
+ throw new Error("CHAT_MESSAGE_SIGNING_SECRET or AUTH_SECRET is required");
+ }
+
+ return secret;
+}
+
+function signableValue(envelope: SignableChatEnvelope) {
+ return JSON.stringify([
+ "shardup-chat-v1",
+ envelope.id,
+ envelope.conversationId,
+ envelope.senderId,
+ envelope.body,
+ envelope.createdAt,
+ ]);
+}
+
+export function signChatEnvelope(envelope: SignableChatEnvelope) {
+ return createHmac("sha256", signingSecret()).update(signableValue(envelope)).digest("base64url");
+}
+
+export function verifyChatEnvelope(envelope: SignableChatEnvelope, signature: string) {
+ const expected = Buffer.from(signChatEnvelope(envelope));
+ const provided = Buffer.from(signature);
+ return expected.length === provided.length && timingSafeEqual(expected, provided);
+}
diff --git a/lib/chat.ts b/lib/chat.ts
new file mode 100644
index 0000000..5e55b9d
--- /dev/null
+++ b/lib/chat.ts
@@ -0,0 +1,114 @@
+import { ChatConversationType, ChatMemberRole, UserStatus } from "@/prisma-client";
+import { z } from "zod";
+
+export const CHAT_MESSAGE_MAX_LENGTH = 4000;
+export const CHAT_GROUP_MAX_MEMBERS = 20;
+export const CHAT_UNDELIVERED_TTL_DAYS = 30;
+export const CHAT_SEND_RECEIPT_TTL_HOURS = 24;
+
+const chatId = z.string().trim().min(1).max(128);
+
+export const sendChatMessageSchema = z.object({
+ conversationId: chatId,
+ clientMessageId: chatId,
+ body: z.string().trim().min(1, "Write a message").max(CHAT_MESSAGE_MAX_LENGTH),
+});
+
+export const createDirectConversationSchema = z.object({
+ recipientId: chatId,
+});
+
+export const createGroupConversationSchema = z.object({
+ title: z.string().trim().min(2, "Add a group name").max(80),
+ memberIds: z
+ .array(chatId)
+ .min(1, "Choose at least one member")
+ .max(CHAT_GROUP_MAX_MEMBERS - 1)
+ .refine((memberIds) => new Set(memberIds).size === memberIds.length, {
+ message: "Choose each member once",
+ }),
+});
+
+export const chatAckSchema = z.object({
+ messageIds: z
+ .array(chatId)
+ .min(1)
+ .max(100)
+ .transform((messageIds) => Array.from(new Set(messageIds))),
+});
+
+export const markChatReadSchema = z.object({
+ conversationId: chatId,
+});
+
+export const chatReportSchema = z.object({
+ conversationId: chatId,
+ messageId: chatId,
+ reportedSenderId: chatId,
+ body: z.string().min(1).max(CHAT_MESSAGE_MAX_LENGTH),
+ messageCreatedAt: z.coerce.date(),
+ signature: z.string().trim().min(1).max(256),
+ note: z
+ .string()
+ .trim()
+ .max(1000)
+ .optional()
+ .transform((value) => value || null),
+});
+
+export function directConversationKey(firstUserId: string, secondUserId: string) {
+ const userIds = [firstUserId.trim(), secondUserId.trim()];
+
+ if (userIds.some((userId) => !userId) || userIds[0] === userIds[1]) {
+ throw new Error("A direct conversation requires two different users");
+ }
+
+ return `direct:${userIds.sort().map(encodeURIComponent).join(":")}`;
+}
+
+type ChatMembership = {
+ userId: string;
+ role: ChatMemberRole;
+ leftAt: Date | null;
+};
+
+type ChatConversationAccess = {
+ type: ChatConversationType;
+ members: ChatMembership[];
+};
+
+export function activeChatMember(
+ conversation: ChatConversationAccess,
+ userId: string,
+): ChatMembership | null {
+ return (
+ conversation.members.find((member) => member.userId === userId && member.leftAt === null) ??
+ null
+ );
+}
+
+export function canAccessChat(conversation: ChatConversationAccess, userId: string) {
+ return activeChatMember(conversation, userId) !== null;
+}
+
+export function canManageChat(conversation: ChatConversationAccess, userId: string) {
+ const member = activeChatMember(conversation, userId);
+ return conversation.type === ChatConversationType.GROUP && member?.role === ChatMemberRole.OWNER;
+}
+
+export function chatRecipientIds(conversation: ChatConversationAccess, senderId: string) {
+ return Array.from(
+ new Set(
+ conversation.members
+ .filter((member) => member.leftAt === null && member.userId !== senderId)
+ .map((member) => member.userId),
+ ),
+ );
+}
+
+export function allChatParticipantsActive(
+ users: Array<{ status: UserStatus }>,
+ expectedCount: number,
+) {
+ return users.length === expectedCount && users.every((user) => user.status === UserStatus.ACTIVE);
+}
diff --git a/package-lock.json b/package-lock.json
index 2c01718..02da6e6 100644
--- a/package-lock.json
+++ b/package-lock.json
@@ -20,6 +20,8 @@
"@prisma/client": "^7.8.0",
"@uiw/react-codemirror": "^4.25.10",
"@vercel/blob": "^2.4.1",
+ "idb": "^8.0.3",
+ "jose": "^6.2.8",
"lenis": "^1.3.25",
"next": "^14.2.4",
"next-auth": "^5.0.0-beta.31",
@@ -7915,6 +7917,12 @@
"node": ">=0.10.0"
}
},
+ "node_modules/idb": {
+ "version": "8.0.3",
+ "resolved": "https://registry.npmjs.org/idb/-/idb-8.0.3.tgz",
+ "integrity": "sha512-LtwtVyVYO5BqRvcsKuB2iUMnHwPVByPCXFXOpuU96IZPPoPN6xjOGxZQ74pgSVVLQWtUOYgyeL4GE98BY5D3wg==",
+ "license": "ISC"
+ },
"node_modules/ignore": {
"version": "5.3.2",
"resolved": "https://registry.npmjs.org/ignore/-/ignore-5.3.2.tgz",
@@ -8543,7 +8551,9 @@
}
},
"node_modules/jose": {
- "version": "6.2.3",
+ "version": "6.2.8",
+ "resolved": "https://registry.npmjs.org/jose/-/jose-6.2.8.tgz",
+ "integrity": "sha512-Bsdjwm3Qsd/P0jR+BHDe3LytDfY7WBq2HmCCLIwuVRHMuEC9ae7/R474GIUdF1NgCyZjzVo/A9DOiOBtXq8ZoQ==",
"license": "MIT",
"funding": {
"url": "https://github.com/sponsors/panva"
diff --git a/package.json b/package.json
index 8e426aa..ea8f27c 100644
--- a/package.json
+++ b/package.json
@@ -4,6 +4,7 @@
"private": true,
"scripts": {
"dev": "next dev",
+ "chat:gateway": "node infra/chat-gateway.mjs",
"postinstall": "node scripts/copy-excalidraw-assets.mjs",
"build": "prisma generate && next build",
"start": "next start",
@@ -37,6 +38,8 @@
"@prisma/client": "^7.8.0",
"@uiw/react-codemirror": "^4.25.10",
"@vercel/blob": "^2.4.1",
+ "idb": "^8.0.3",
+ "jose": "^6.2.8",
"lenis": "^1.3.25",
"next": "^14.2.4",
"next-auth": "^5.0.0-beta.31",
diff --git a/prisma/migrations/20260808120000_add_ephemeral_chat/migration.sql b/prisma/migrations/20260808120000_add_ephemeral_chat/migration.sql
new file mode 100644
index 0000000..29e6f19
--- /dev/null
+++ b/prisma/migrations/20260808120000_add_ephemeral_chat/migration.sql
@@ -0,0 +1,126 @@
+-- Durable conversation metadata with transient store-and-forward payloads.
+CREATE TYPE "ChatConversationType" AS ENUM ('DIRECT', 'GROUP');
+CREATE TYPE "ChatMemberRole" AS ENUM ('OWNER', 'MEMBER');
+CREATE TYPE "ChatReportStatus" AS ENUM ('OPEN', 'RESOLVED', 'DISMISSED');
+CREATE TYPE "ChatRateLimitScope" AS ENUM ('MESSAGE_SEND', 'CONVERSATION_CREATE');
+
+CREATE TABLE "ChatConversation" (
+ "id" TEXT NOT NULL,
+ "type" "ChatConversationType" NOT NULL,
+ "title" TEXT,
+ "directKey" TEXT,
+ "createdById" TEXT NOT NULL,
+ "lastActivityAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ "updatedAt" TIMESTAMP(3) NOT NULL,
+
+ CONSTRAINT "ChatConversation_pkey" PRIMARY KEY ("id")
+);
+
+CREATE TABLE "ChatMember" (
+ "id" TEXT NOT NULL,
+ "conversationId" TEXT NOT NULL,
+ "userId" TEXT NOT NULL,
+ "role" "ChatMemberRole" NOT NULL DEFAULT 'MEMBER',
+ "joinedAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ "leftAt" TIMESTAMP(3),
+ "lastReadAt" TIMESTAMP(3),
+
+ CONSTRAINT "ChatMember_pkey" PRIMARY KEY ("id")
+);
+
+CREATE TABLE "PendingChatMessage" (
+ "id" TEXT NOT NULL,
+ "conversationId" TEXT NOT NULL,
+ "senderId" TEXT NOT NULL,
+ "clientMessageId" TEXT NOT NULL,
+ "body" TEXT NOT NULL,
+ "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ "expiresAt" TIMESTAMP(3) NOT NULL,
+
+ CONSTRAINT "PendingChatMessage_pkey" PRIMARY KEY ("id")
+);
+
+CREATE TABLE "ChatDelivery" (
+ "id" TEXT NOT NULL,
+ "messageId" TEXT NOT NULL,
+ "recipientId" TEXT NOT NULL,
+ "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+
+ CONSTRAINT "ChatDelivery_pkey" PRIMARY KEY ("id")
+);
+
+CREATE TABLE "ChatSendReceipt" (
+ "id" TEXT NOT NULL,
+ "senderId" TEXT NOT NULL,
+ "clientMessageId" TEXT NOT NULL,
+ "messageId" TEXT NOT NULL,
+ "conversationId" TEXT NOT NULL,
+ "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ "expiresAt" TIMESTAMP(3) NOT NULL,
+
+ CONSTRAINT "ChatSendReceipt_pkey" PRIMARY KEY ("id")
+);
+
+CREATE TABLE "ChatReport" (
+ "id" TEXT NOT NULL,
+ "conversationId" TEXT NOT NULL,
+ "reporterId" TEXT NOT NULL,
+ "reportedSenderId" TEXT NOT NULL,
+ "resolvedById" TEXT,
+ "messageId" TEXT NOT NULL,
+ "body" TEXT NOT NULL,
+ "note" TEXT,
+ "messageCreatedAt" TIMESTAMP(3) NOT NULL,
+ "status" "ChatReportStatus" NOT NULL DEFAULT 'OPEN',
+ "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ "resolvedAt" TIMESTAMP(3),
+
+ CONSTRAINT "ChatReport_pkey" PRIMARY KEY ("id")
+);
+
+CREATE TABLE "ChatRateLimit" (
+ "id" TEXT NOT NULL,
+ "userId" TEXT NOT NULL,
+ "scope" "ChatRateLimitScope" NOT NULL,
+ "windowStart" TIMESTAMP(3) NOT NULL,
+ "count" INTEGER NOT NULL DEFAULT 0,
+ "expiresAt" TIMESTAMP(3) NOT NULL,
+ "updatedAt" TIMESTAMP(3) NOT NULL,
+
+ CONSTRAINT "ChatRateLimit_pkey" PRIMARY KEY ("id")
+);
+
+CREATE UNIQUE INDEX "ChatConversation_directKey_key" ON "ChatConversation"("directKey");
+CREATE INDEX "ChatConversation_lastActivityAt_idx" ON "ChatConversation"("lastActivityAt");
+CREATE INDEX "ChatConversation_createdById_idx" ON "ChatConversation"("createdById");
+CREATE UNIQUE INDEX "ChatMember_conversationId_userId_key" ON "ChatMember"("conversationId", "userId");
+CREATE INDEX "ChatMember_userId_leftAt_conversationId_idx" ON "ChatMember"("userId", "leftAt", "conversationId");
+CREATE UNIQUE INDEX "PendingChatMessage_senderId_clientMessageId_key" ON "PendingChatMessage"("senderId", "clientMessageId");
+CREATE INDEX "PendingChatMessage_conversationId_createdAt_id_idx" ON "PendingChatMessage"("conversationId", "createdAt", "id");
+CREATE INDEX "PendingChatMessage_expiresAt_idx" ON "PendingChatMessage"("expiresAt");
+CREATE UNIQUE INDEX "ChatDelivery_messageId_recipientId_key" ON "ChatDelivery"("messageId", "recipientId");
+CREATE INDEX "ChatDelivery_recipientId_createdAt_messageId_idx" ON "ChatDelivery"("recipientId", "createdAt", "messageId");
+CREATE UNIQUE INDEX "ChatSendReceipt_senderId_clientMessageId_key" ON "ChatSendReceipt"("senderId", "clientMessageId");
+CREATE INDEX "ChatSendReceipt_expiresAt_idx" ON "ChatSendReceipt"("expiresAt");
+CREATE INDEX "ChatSendReceipt_conversationId_createdAt_idx" ON "ChatSendReceipt"("conversationId", "createdAt");
+CREATE UNIQUE INDEX "ChatReport_reporterId_messageId_key" ON "ChatReport"("reporterId", "messageId");
+CREATE INDEX "ChatReport_status_createdAt_idx" ON "ChatReport"("status", "createdAt");
+CREATE INDEX "ChatReport_conversationId_idx" ON "ChatReport"("conversationId");
+CREATE UNIQUE INDEX "ChatRateLimit_userId_scope_key" ON "ChatRateLimit"("userId", "scope");
+CREATE INDEX "ChatRateLimit_expiresAt_idx" ON "ChatRateLimit"("expiresAt");
+
+ALTER TABLE "ChatConversation" ADD CONSTRAINT "ChatConversation_createdById_fkey" FOREIGN KEY ("createdById") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatMember" ADD CONSTRAINT "ChatMember_conversationId_fkey" FOREIGN KEY ("conversationId") REFERENCES "ChatConversation"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatMember" ADD CONSTRAINT "ChatMember_userId_fkey" FOREIGN KEY ("userId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "PendingChatMessage" ADD CONSTRAINT "PendingChatMessage_conversationId_fkey" FOREIGN KEY ("conversationId") REFERENCES "ChatConversation"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "PendingChatMessage" ADD CONSTRAINT "PendingChatMessage_senderId_fkey" FOREIGN KEY ("senderId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatDelivery" ADD CONSTRAINT "ChatDelivery_messageId_fkey" FOREIGN KEY ("messageId") REFERENCES "PendingChatMessage"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatDelivery" ADD CONSTRAINT "ChatDelivery_recipientId_fkey" FOREIGN KEY ("recipientId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatSendReceipt" ADD CONSTRAINT "ChatSendReceipt_senderId_fkey" FOREIGN KEY ("senderId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatSendReceipt" ADD CONSTRAINT "ChatSendReceipt_conversationId_fkey" FOREIGN KEY ("conversationId") REFERENCES "ChatConversation"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatReport" ADD CONSTRAINT "ChatReport_conversationId_fkey" FOREIGN KEY ("conversationId") REFERENCES "ChatConversation"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatReport" ADD CONSTRAINT "ChatReport_reporterId_fkey" FOREIGN KEY ("reporterId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatReport" ADD CONSTRAINT "ChatReport_reportedSenderId_fkey" FOREIGN KEY ("reportedSenderId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
+ALTER TABLE "ChatReport" ADD CONSTRAINT "ChatReport_resolvedById_fkey" FOREIGN KEY ("resolvedById") REFERENCES "User"("id") ON DELETE SET NULL ON UPDATE CASCADE;
+ALTER TABLE "ChatRateLimit" ADD CONSTRAINT "ChatRateLimit_userId_fkey" FOREIGN KEY ("userId") REFERENCES "User"("id") ON DELETE CASCADE ON UPDATE CASCADE;
\ No newline at end of file
diff --git a/prisma/schema.prisma b/prisma/schema.prisma
index 8b159a3..4e2934f 100644
--- a/prisma/schema.prisma
+++ b/prisma/schema.prisma
@@ -76,6 +76,27 @@ enum NotificationType {
NUDGE_RECEIVED
}
+enum ChatConversationType {
+ DIRECT
+ GROUP
+}
+
+enum ChatMemberRole {
+ OWNER
+ MEMBER
+}
+
+enum ChatReportStatus {
+ OPEN
+ RESOLVED
+ DISMISSED
+}
+
+enum ChatRateLimitScope {
+ MESSAGE_SEND
+ CONVERSATION_CREATE
+}
+
model User {
id String @id @default(cuid())
name String?
@@ -112,6 +133,15 @@ model User {
recommendedResources Resource[] @relation("ResourceRecommender")
notifications Notification[] @relation("NotificationActor")
receivedNotifications Notification[] @relation("NotificationRecipient")
+ createdChatConversations ChatConversation[] @relation("ChatConversationCreator")
+ chatMemberships ChatMember[]
+ sentPendingChatMessages PendingChatMessage[] @relation("PendingChatMessageSender")
+ chatDeliveries ChatDelivery[] @relation("ChatDeliveryRecipient")
+ chatSendReceipts ChatSendReceipt[] @relation("ChatSendReceiptSender")
+ chatReports ChatReport[] @relation("ChatReportReporter")
+ reportedChatMessages ChatReport[] @relation("ChatReportSender")
+ resolvedChatReports ChatReport[] @relation("ChatReportResolver")
+ chatRateLimits ChatRateLimit[]
}
// The yearly application round. One cohort is active at a time; admins edit its
@@ -598,3 +628,132 @@ model Notification {
@@index([actorId])
@@index([recipientId])
}
+
+// Conversation and membership metadata is durable. Message content is not:
+// pending payloads are removed after every recipient durably acknowledges them.
+model ChatConversation {
+ id String @id @default(cuid())
+ type ChatConversationType
+ title String?
+ directKey String? @unique
+ createdById String
+ lastActivityAt DateTime @default(now())
+ createdAt DateTime @default(now())
+ updatedAt DateTime @updatedAt
+
+ creator User @relation("ChatConversationCreator", fields: [createdById], references: [id], onDelete: Cascade)
+ members ChatMember[]
+ pendingMessages PendingChatMessage[]
+ sendReceipts ChatSendReceipt[]
+ reports ChatReport[]
+
+ @@index([lastActivityAt])
+ @@index([createdById])
+}
+
+model ChatMember {
+ id String @id @default(cuid())
+ conversationId String
+ userId String
+ role ChatMemberRole @default(MEMBER)
+ joinedAt DateTime @default(now())
+ leftAt DateTime?
+ lastReadAt DateTime?
+
+ conversation ChatConversation @relation(fields: [conversationId], references: [id], onDelete: Cascade)
+ user User @relation(fields: [userId], references: [id], onDelete: Cascade)
+
+ @@unique([conversationId, userId])
+ @@index([userId, leftAt, conversationId])
+}
+
+model PendingChatMessage {
+ id String @id @default(cuid())
+ conversationId String
+ senderId String
+ clientMessageId String
+ body String
+ createdAt DateTime @default(now())
+ expiresAt DateTime
+
+ conversation ChatConversation @relation(fields: [conversationId], references: [id], onDelete: Cascade)
+ sender User @relation("PendingChatMessageSender", fields: [senderId], references: [id], onDelete: Cascade)
+ deliveries ChatDelivery[]
+
+ @@unique([senderId, clientMessageId])
+ @@index([conversationId, createdAt, id])
+ @@index([expiresAt])
+}
+
+model ChatDelivery {
+ id String @id @default(cuid())
+ messageId String
+ recipientId String
+ createdAt DateTime @default(now())
+
+ message PendingChatMessage @relation(fields: [messageId], references: [id], onDelete: Cascade)
+ recipient User @relation("ChatDeliveryRecipient", fields: [recipientId], references: [id], onDelete: Cascade)
+
+ @@unique([messageId, recipientId])
+ @@index([recipientId, createdAt, messageId])
+}
+
+// Receipts contain no message content. They keep send retries idempotent after
+// the corresponding pending payload has already been acknowledged and deleted.
+model ChatSendReceipt {
+ id String @id @default(cuid())
+ senderId String
+ clientMessageId String
+ messageId String
+ conversationId String
+ createdAt DateTime @default(now())
+ expiresAt DateTime
+
+ sender User @relation("ChatSendReceiptSender", fields: [senderId], references: [id], onDelete: Cascade)
+ conversation ChatConversation @relation(fields: [conversationId], references: [id], onDelete: Cascade)
+
+ @@unique([senderId, clientMessageId])
+ @@index([expiresAt])
+ @@index([conversationId, createdAt])
+}
+
+// A report is an explicit user-provided moderation snapshot and is the only
+// server-side record allowed to retain content after normal delivery.
+model ChatReport {
+ id String @id @default(cuid())
+ conversationId String
+ reporterId String
+ reportedSenderId String
+ resolvedById String?
+ messageId String
+ body String
+ note String?
+ messageCreatedAt DateTime
+ status ChatReportStatus @default(OPEN)
+ createdAt DateTime @default(now())
+ resolvedAt DateTime?
+
+ conversation ChatConversation @relation(fields: [conversationId], references: [id], onDelete: Cascade)
+ reporter User @relation("ChatReportReporter", fields: [reporterId], references: [id], onDelete: Cascade)
+ reportedSender User @relation("ChatReportSender", fields: [reportedSenderId], references: [id], onDelete: Cascade)
+ resolvedBy User? @relation("ChatReportResolver", fields: [resolvedById], references: [id], onDelete: SetNull)
+
+ @@unique([reporterId, messageId])
+ @@index([status, createdAt])
+ @@index([conversationId])
+}
+
+model ChatRateLimit {
+ id String @id @default(cuid())
+ userId String
+ scope ChatRateLimitScope
+ windowStart DateTime
+ count Int @default(0)
+ expiresAt DateTime
+ updatedAt DateTime @updatedAt
+
+ user User @relation(fields: [userId], references: [id], onDelete: Cascade)
+
+ @@unique([userId, scope])
+ @@index([expiresAt])
+}
diff --git a/readme.md b/readme.md
index abf0adf..c3c71d7 100644
--- a/readme.md
+++ b/readme.md
@@ -165,6 +165,7 @@ Google OAuth is the real sign-in method. Callback URLs:
For local development, Google credentials are optional. Set `LOCAL_DEV_AUTH_ENABLED=true` in `.env.local` and the `/join` page shows two development-only sign-ins:
- **Continue as applicant (dev)** — signs in as `applicant@shardup.local` to test the application flow. This is a throwaway test account: it is reset to `PENDING` with a fresh blank application on every login, so you can re-run the flow repeatedly.
+- **Continue as member (dev)** — signs in as `member@shardup.local`, an ACTIVE member used for chat and member-only feature testing.
- **Continue as admin (dev)** — signs in as `admin@shardup.local` to test application review. Make sure `admin@shardup.local` is in `ADMIN_EMAILS`.
### Useful commands
@@ -174,8 +175,17 @@ For local development, Google credentials are optional. Set `LOCAL_DEV_AUTH_ENAB
- `npm run prisma:seed`
- `npm run prisma:studio`
- `npm run dev`
+- `npm run chat:gateway`
- `npm run build`
+### Ephemeral messaging
+
+Messages are queued in PostgreSQL before realtime fan-out. Browsers persist delivered messages in
+IndexedDB and then acknowledge them; the server removes each recipient delivery and deletes the
+payload after the final acknowledgement. Offline messages expire after 30 days. Run the open-source
+WebSocket gateway with `npm run chat:gateway`; Redis is not needed for a single gateway process.
+See `docs/chat-gateway.md` for local and Oracle VM deployment.
+
Events are published manually for now. Seed sample events with `npm run prisma:seed`, or manage rows directly in Prisma Studio. Add an optional `imageUrl` to show an event image on the list and detail pages. RSVP is available only to active members; signed-out users can view events but must sign in before RSVPing.
Practice problems are also seed-managed for now. `npm run prisma:seed` publishes the sample `Sum Two Numbers` problem with sample and hidden test cases. Anyone can view problems; only active members can submit solutions.
diff --git a/tests/e2e/messages.spec.ts b/tests/e2e/messages.spec.ts
new file mode 100644
index 0000000..723986f
--- /dev/null
+++ b/tests/e2e/messages.spec.ts
@@ -0,0 +1,65 @@
+import { expect, test } from "@playwright/test";
+import { devLogin } from "./utils";
+
+test.describe("messages", () => {
+ test("stores, recovers, and acknowledges an offline direct message", async ({ browser }) => {
+ const recipientContext = await browser.newContext();
+ const senderContext = await browser.newContext();
+ const recipient = await recipientContext.newPage();
+ const sender = await senderContext.newPage();
+ const message = `Queued message ${Date.now()}`;
+
+ try {
+ await devLogin(recipient, "admin");
+ await devLogin(sender, "active");
+
+ await sender.goto("/messages");
+ await sender.getByRole("button", { name: "New message" }).click();
+ await sender.locator('select[name="recipientId"]').selectOption({ label: "Local Admin" });
+ await sender.getByRole("button", { name: "Create", exact: true }).click();
+
+ await expect(sender.getByRole("heading", { name: "Local Admin" })).toBeVisible();
+ await sender.getByRole("textbox", { name: "Message", exact: true }).fill(message);
+ const sendResponsePromise = sender.waitForResponse(
+ (response) =>
+ response.request().method() === "POST" && response.url().endsWith("/api/chat/messages"),
+ );
+ await sender.getByRole("button", { name: "Send", exact: true }).click();
+ const sendResponse = await sendResponsePromise;
+ expect(sendResponse.ok()).toBe(true);
+ const { messageId } = (await sendResponse.json()) as { messageId: string };
+ const sentMessage = sender.locator(".chat-message", { hasText: message });
+ await expect(sentMessage).toContainText(message);
+ await expect(sentMessage).not.toContainText("Sending...");
+ await expect(sentMessage).not.toContainText("Not sent");
+
+ await recipient.bringToFront();
+ await recipient.evaluate(() => window.dispatchEvent(new Event("focus")));
+ const messageIndicator = recipient.getByRole("link", {
+ name: /Messages, \d+ unread/,
+ });
+ await expect(messageIndicator).toBeVisible();
+ await messageIndicator.click();
+ await expect(recipient.locator(".chat-message", { hasText: message })).toContainText(
+ message,
+ {
+ timeout: 10_000,
+ },
+ );
+
+ await recipient.reload();
+ await expect(recipient.locator(".chat-message", { hasText: message })).toContainText(message);
+
+ await expect
+ .poll(async () => {
+ const response = await recipient.request.get("/api/chat/inbox");
+ const inbox = (await response.json()) as { messages: Array<{ id: string }> };
+ return inbox.messages.some((entry) => entry.id === messageId);
+ })
+ .toBe(false);
+ } finally {
+ await senderContext.close();
+ await recipientContext.close();
+ }
+ });
+});
diff --git a/tests/e2e/utils.ts b/tests/e2e/utils.ts
index 519de37..35d5641 100644
--- a/tests/e2e/utils.ts
+++ b/tests/e2e/utils.ts
@@ -1,10 +1,12 @@
import { type Page } from "@playwright/test";
-type DevRole = "admin" | "member";
+type DevRole = "admin" | "member" | "active";
// Uses the app's own development-only sign-in route. The route signs the user
// in server-side and redirects: admin -> /admin/cohort, member -> /apply.
export async function devLogin(page: Page, role: DevRole) {
await page.goto(`/dev-login?role=${role}`);
- await page.waitForURL(role === "admin" ? "**/admin/cohort" : "**/apply");
+ await page.waitForURL(
+ role === "admin" ? "**/admin/cohort" : role === "active" ? "**/dashboard" : "**/apply",
+ );
}
diff --git a/tests/unit/components/account-bar.test.tsx b/tests/unit/components/account-bar.test.tsx
index 6f6b565..53f05ea 100644
--- a/tests/unit/components/account-bar.test.tsx
+++ b/tests/unit/components/account-bar.test.tsx
@@ -44,6 +44,7 @@ describe("AccountBar", () => {
"href",
"/masterclass",
);
+ expect(screen.getByRole("link", { name: "Messages" })).toHaveAttribute("href", "/messages");
expect(screen.getByRole("button", { name: "Sign out" })).toBeInTheDocument();
expect(screen.queryByRole("link", { name: "Cohort" })).not.toBeInTheDocument();
expect(screen.queryByRole("link", { name: "Applications" })).not.toBeInTheDocument();
diff --git a/tests/unit/lib/chat.test.ts b/tests/unit/lib/chat.test.ts
new file mode 100644
index 0000000..376ef43
--- /dev/null
+++ b/tests/unit/lib/chat.test.ts
@@ -0,0 +1,98 @@
+import { ChatConversationType, ChatMemberRole, UserStatus } from "@/prisma-client";
+import { describe, expect, it } from "vitest";
+import {
+ allChatParticipantsActive,
+ canAccessChat,
+ canManageChat,
+ chatAckSchema,
+ chatRecipientIds,
+ createGroupConversationSchema,
+ directConversationKey,
+ sendChatMessageSchema,
+} from "../../../lib/chat";
+
+const group = {
+ type: ChatConversationType.GROUP,
+ members: [
+ { userId: "owner", role: ChatMemberRole.OWNER, leftAt: null },
+ { userId: "member", role: ChatMemberRole.MEMBER, leftAt: null },
+ { userId: "former", role: ChatMemberRole.MEMBER, leftAt: new Date() },
+ ],
+};
+
+describe("chat schemas", () => {
+ it("trims valid messages and rejects empty or oversized bodies", () => {
+ expect(
+ sendChatMessageSchema.parse({
+ conversationId: "conversation",
+ clientMessageId: "client-message",
+ body: " hello ",
+ }).body,
+ ).toBe("hello");
+ expect(
+ sendChatMessageSchema.safeParse({
+ conversationId: "conversation",
+ clientMessageId: "client-message",
+ body: " ",
+ }).success,
+ ).toBe(false);
+ expect(
+ sendChatMessageSchema.safeParse({
+ conversationId: "conversation",
+ clientMessageId: "client-message",
+ body: "x".repeat(4001),
+ }).success,
+ ).toBe(false);
+ });
+
+ it("requires unique group members and caps the group at twenty people", () => {
+ expect(
+ createGroupConversationSchema.safeParse({ title: "Study group", memberIds: ["one"] }).success,
+ ).toBe(true);
+ expect(
+ createGroupConversationSchema.safeParse({
+ title: "Study group",
+ memberIds: ["one", "one"],
+ }).success,
+ ).toBe(false);
+ expect(
+ createGroupConversationSchema.safeParse({
+ title: "Study group",
+ memberIds: Array.from({ length: 20 }, (_, index) => `member-${index}`),
+ }).success,
+ ).toBe(false);
+ });
+
+ it("deduplicates acknowledgement IDs", () => {
+ expect(chatAckSchema.parse({ messageIds: ["one", "one", "two"] }).messageIds).toEqual([
+ "one",
+ "two",
+ ]);
+ });
+});
+
+describe("chat identity and permissions", () => {
+ it("builds the same direct key in either user order", () => {
+ expect(directConversationKey("alice", "bob")).toBe(directConversationKey("bob", "alice"));
+ expect(() => directConversationKey("alice", "alice")).toThrow();
+ });
+
+ it("allows active members to access and only the group owner to manage", () => {
+ expect(canAccessChat(group, "member")).toBe(true);
+ expect(canAccessChat(group, "former")).toBe(false);
+ expect(canManageChat(group, "owner")).toBe(true);
+ expect(canManageChat(group, "member")).toBe(false);
+ });
+
+ it("fans out to current members other than the sender", () => {
+ expect(chatRecipientIds(group, "owner")).toEqual(["member"]);
+ });
+
+ it("requires every requested participant to exist and be active", () => {
+ expect(
+ allChatParticipantsActive([{ status: UserStatus.ACTIVE }, { status: UserStatus.ACTIVE }], 2),
+ ).toBe(true);
+ expect(allChatParticipantsActive([{ status: UserStatus.ACTIVE }], 2)).toBe(false);
+ expect(allChatParticipantsActive([{ status: UserStatus.SUSPENDED }], 1)).toBe(false);
+ });
+});
diff --git a/vercel.json b/vercel.json
index cbf6517..e828892 100644
--- a/vercel.json
+++ b/vercel.json
@@ -1,5 +1,11 @@
{
"buildCommand": "prisma generate && next build",
"installCommand": "npm install",
- "framework": "nextjs"
+ "framework": "nextjs",
+ "crons": [
+ {
+ "path": "/api/cron/chat-cleanup",
+ "schedule": "17 2 * * *"
+ }
+ ]
}