From 45575bd37025ab6ab0340e75750a015a7f82bd82 Mon Sep 17 00:00:00 2001 From: John Chantzigoulas Date: Fri, 11 Sep 2026 14:08:35 +0300 Subject: [PATCH 1/2] feat(router): add event relays administration page Give operators CRUD and an in-page guide for mapping exact bus events to ReBAC-protected socket rooms, with a local template preview before save. --- package.json | 1 + .../(modules)/router/event-relays/page.tsx | 32 ++ src/app/(dashboard)/(modules)/router/page.tsx | 15 +- src/components/navigation/navList.config.ts | 1 + .../router/event-relays/event-relay-docs.tsx | 230 +++++++++++++ .../router/event-relays/event-relay-form.tsx | 213 ++++++++++++ .../router/event-relays/event-relay-list.tsx | 304 ++++++++++++++++++ src/components/router/event-relays/zod.ts | 72 +++++ src/lib/api/router/index.ts | 39 ++- src/lib/event-relays/path.ts | 56 ++++ src/lib/event-relays/preview.test.ts | 31 ++ src/lib/event-relays/preview.ts | 33 ++ src/lib/event-relays/template.ts | 59 ++++ src/lib/models/Router.ts | 32 ++ 14 files changed, 1116 insertions(+), 2 deletions(-) create mode 100644 src/app/(dashboard)/(modules)/router/event-relays/page.tsx create mode 100644 src/components/router/event-relays/event-relay-docs.tsx create mode 100644 src/components/router/event-relays/event-relay-form.tsx create mode 100644 src/components/router/event-relays/event-relay-list.tsx create mode 100644 src/components/router/event-relays/zod.ts create mode 100644 src/lib/event-relays/path.ts create mode 100644 src/lib/event-relays/preview.test.ts create mode 100644 src/lib/event-relays/preview.ts create mode 100644 src/lib/event-relays/template.ts diff --git a/package.json b/package.json index e5f90b63c..62ec88456 100644 --- a/package.json +++ b/package.json @@ -9,6 +9,7 @@ "start": "next start -p 8080", "lint": "eslint .", "test:semantic-colors": "node --experimental-strip-types --test src/lib/semantic-colors.test.ts src/lib/reactflow-edge-colors.test.ts", + "test:event-relays": "node --experimental-strip-types --test src/lib/event-relays/preview.test.ts", "test:schema-fields": "node --experimental-strip-types --test src/lib/database/schema-field-definition.test.ts src/lib/database/system-schema-fields.test.ts", "build:docker": "docker build --platform linux/amd64 -t ghcr.io/conduitplatform/conduit-ui:latest .", "prepare": "husky", diff --git a/src/app/(dashboard)/(modules)/router/event-relays/page.tsx b/src/app/(dashboard)/(modules)/router/event-relays/page.tsx new file mode 100644 index 000000000..509f7094a --- /dev/null +++ b/src/app/(dashboard)/(modules)/router/event-relays/page.tsx @@ -0,0 +1,32 @@ +import { getEventRelays, getRouterSettings } from '@/lib/api/router'; +import { EventRelayList } from '@/components/router/event-relays/event-relay-list'; + +export default async function EventRelaysPage(props: { + searchParams: Promise<{ + skip?: string; + limit?: string; + search?: string; + }>; +}) { + const searchParams = await props.searchParams; + const skip = Number(searchParams.skip ?? 0); + const limit = Number(searchParams.limit ?? 10); + const [{ relays, count }, { config }] = await Promise.all([ + getEventRelays({ + skip, + limit, + search: searchParams.search, + }), + getRouterSettings(), + ]); + + return ( +
+ +
+ ); +} diff --git a/src/app/(dashboard)/(modules)/router/page.tsx b/src/app/(dashboard)/(modules)/router/page.tsx index 64eaa8897..11e104ea1 100644 --- a/src/app/(dashboard)/(modules)/router/page.tsx +++ b/src/app/(dashboard)/(modules)/router/page.tsx @@ -1,5 +1,12 @@ import React from 'react'; -import { BarChart3, Network, Route, Settings, Shield } from 'lucide-react'; +import { + BarChart3, + Network, + Radio, + Route, + Settings, + Shield, +} from 'lucide-react'; import { ModuleDashboard } from '@/components/dashboard/ModuleDashboard'; import { getModuleStatus, @@ -66,6 +73,12 @@ export default async function RouterDashboard() { icon: , href: '/router/vizualize', }, + { + title: 'Event Relays', + description: 'Forward bus events to socket subscribers', + icon: , + href: '/router/event-relays', + }, { title: 'Settings', description: 'Router module configuration', diff --git a/src/components/navigation/navList.config.ts b/src/components/navigation/navList.config.ts index ff9d42677..017235613 100644 --- a/src/components/navigation/navList.config.ts +++ b/src/components/navigation/navList.config.ts @@ -133,6 +133,7 @@ export const navGroups: NavGroup[] = [ items: [ { title: 'Visualize', url: '/router/vizualize' }, { title: 'Security', url: '/router/security' }, + { title: 'Event Relays', url: '/router/event-relays' }, { title: 'Settings', url: '/router/settings' }, ], }, diff --git a/src/components/router/event-relays/event-relay-docs.tsx b/src/components/router/event-relays/event-relay-docs.tsx new file mode 100644 index 000000000..a1c3ed895 --- /dev/null +++ b/src/components/router/event-relays/event-relay-docs.tsx @@ -0,0 +1,230 @@ +'use client'; + +import { ChevronDown } from 'lucide-react'; +import { Card } from '@/components/ui/card'; +import { + Collapsible, + CollapsibleContent, + CollapsibleTrigger, +} from '@/components/ui/collapsible'; +import { cn } from '@/lib/utils'; + +const CLIENT_SNIPPET = `const socket = io(\`\${SOCKET_URL}/events/\`, { + path: '/realtime', + extraHeaders: { authorization: \`Bearer \${accessToken}\` }, +}); +socket.emit('subscribe', relayId, resourceId); +socket.on('order-updated', payload => {}); +socket.emit('unsubscribe', relayId, resourceId);`; + +const STEPS = [ + { + title: 'Bus event', + body: 'A module publishes JSON on an exact Redis channel — for example database realtime or a custom module.', + }, + { + title: 'Active relay', + body: 'The matching relay reads the resource id from the payload and renders the message template.', + }, + { + title: 'ReBAC subscribe', + body: 'Clients never pick a room name. Subscribe succeeds only if the user has the relay permission on that resource.', + }, + { + title: 'Socket emit', + body: 'Router emits socketEvent to the hashed /events/ room. Delivery is ephemeral — missed events are gone.', + }, +] as const; + +const SCOPE = [ + { + title: 'Use when', + body: 'You already publish JSON on an exact Redis bus channel and need live, per-resource UI updates with a ReBAC check.', + }, + { + title: 'Skip when', + body: 'You need replay, history, guaranteed delivery, wildcard channels, or a broadcast with no permission check. Relays are not a queue.', + }, + { + title: 'Requires', + body: 'Router sockets enabled, the Authorization module available, and something publishing on that exact channel.', + }, +] as const; + +interface EventRelayDocsProps { + open: boolean; + onOpenChange: (open: boolean) => void; +} + +export function EventRelayDocs({ open, onOpenChange }: EventRelayDocsProps) { + return ( + + + svg]:rotate-180' + )} + > + + + How Event Relays work + + + Forward an exact bus event to permission-scoped socket + subscribers. Not a queue, and not a generic websocket broadcast. + + + + + +
+
+

Scope

+
+ {SCOPE.map(item => ( +
+

+ {item.title} +

+

+ {item.body} +

+
+ ))} +
+
+ +
+

+ How a message moves +

+
    + {STEPS.map((step, index) => ( +
  1. + + {index + 1} + +
    +

    + {step.title} +

    +

    + {step.body} +

    +
    +
  2. + ))} +
+
+ +
+

+ Configure a relay +

+

+ Example: notify clients when an Order document changes. +

+
+ + + + + +
+
+ +
+

+ Subscribe from a client +

+

+ Connect to {'/events/'} with{' '} + {'path: /realtime'} and a user bearer token. Then + subscribe with the relay id and resource id. +

+
+                {CLIENT_SNIPPET}
+              
+
+ +
+

Limits

+
    +
  • + Bus channels must match exactly. Patterns like{' '} + {'database:change:*'} are not supported. +
  • +
  • + Subscribe fails closed if Authorization is unavailable or the + user lacks permission. +
  • +
  • + Turn a relay off with Active to stop forwarding without + deleting it. Deleting drops current subscribers immediately. +
  • +
+
+
+
+
+
+ ); +} + +function Field({ + name, + value, + hint, +}: { + name: string; + value: string; + hint: string; +}) { + return ( +
+
+ {name} +
+
+ {value} +

{hint}

+
+
+ ); +} + +function Code({ children }: { children: string }) { + return ( + + {children} + + ); +} diff --git a/src/components/router/event-relays/event-relay-form.tsx b/src/components/router/event-relays/event-relay-form.tsx new file mode 100644 index 000000000..92694b7e3 --- /dev/null +++ b/src/components/router/event-relays/event-relay-form.tsx @@ -0,0 +1,213 @@ +'use client'; + +import { useMemo } from 'react'; +import { useForm, useWatch } from 'react-hook-form'; +import { rhfZodResolver } from '@/lib/zod-form'; +import { Form } from '@/components/ui/form'; +import { InputField } from '@/components/ui/form-inputs/InputField'; +import SwitchField from '@/components/ui/form-inputs/SwitchField'; +import { CodeField } from '@/components/ui/form-inputs/CodeField'; +import { Button } from '@/components/ui/button'; +import { Alert, AlertDescription, AlertTitle } from '@/components/ui/alert'; +import { + EventRelayFormSchema, + EventRelayFormValues, + parseJsonField, +} from '@/components/router/event-relays/zod'; +import { previewEventRelay } from '@/lib/event-relays/preview'; +import { EventRelay, EventRelayWriteRequest } from '@/lib/models/Router'; + +const DEFAULT_TEMPLATE = '{\n "id": "{{payload._id}}"\n}'; +const DEFAULT_SAMPLE = + '{\n "_id": "64f1c0a2b4d0e1f2a3b4c5d6",\n "status": "paid"\n}'; + +interface EventRelayFormProps { + relay?: EventRelay | null; + onSubmit: (data: EventRelayWriteRequest) => Promise; + onCancel: () => void; + isSaving?: boolean; +} + +export function EventRelayForm({ + relay, + onSubmit, + onCancel, + isSaving, +}: EventRelayFormProps) { + const form = useForm({ + resolver: rhfZodResolver(EventRelayFormSchema), + defaultValues: { + name: relay?.name ?? '', + notes: relay?.notes ?? '', + active: relay?.active ?? true, + busEvent: relay?.busEvent ?? '', + socketEvent: relay?.socketEvent ?? '', + resourceType: relay?.resourceType ?? '', + resourceIdPath: relay?.resourceIdPath ?? '_id', + permission: relay?.permission ?? 'read', + messageTemplate: relay + ? JSON.stringify(relay.messageTemplate, null, 2) + : DEFAULT_TEMPLATE, + samplePayload: DEFAULT_SAMPLE, + }, + }); + + const watched = useWatch({ control: form.control }); + const preview = useMemo(() => { + try { + const template = parseJsonField( + watched.messageTemplate ?? '', + 'Message template' + ); + const sample = parseJsonField( + watched.samplePayload?.trim() || '{}', + 'Sample payload' + ); + return previewEventRelay({ + resourceIdPath: watched.resourceIdPath || '_id', + messageTemplate: template, + samplePayload: sample, + }); + } catch (err) { + return { + error: err instanceof Error ? err.message : String(err), + }; + } + }, [watched.messageTemplate, watched.samplePayload, watched.resourceIdPath]); + + const handleSubmit = form.handleSubmit(async values => { + const messageTemplate = parseJsonField( + values.messageTemplate, + 'Message template' + ); + const notes = values.notes?.trim(); + await onSubmit({ + name: values.name, + notes: notes || undefined, + active: values.active, + busEvent: values.busEvent, + socketEvent: values.socketEvent, + resourceType: values.resourceType, + resourceIdPath: values.resourceIdPath, + permission: values.permission, + messageTemplate, + }); + }); + + return ( +
+ +
+ + +
+ +
+ + +
+
+ + + +
+ + +
+

Preview

+

+ Local only. Nothing is published to the bus. +

+ {preview.error ? ( +

{preview.error}

+ ) : ( +
+

+ Resource{' '} + + {watched.resourceType || 'Type'}:{preview.resourceId} + +

+
+                {JSON.stringify(preview.payload, null, 2)}
+              
+
+ )} +
+ + Client contract + +

+ Connect to /events/ with{' '} + path: /realtime and a bearer + token. Then emit{' '} + subscribe(relayId, resourceId). +

+
+              {`const socket = io(\`\${SOCKET_URL}/events/\`, {
+  path: '/realtime',
+  extraHeaders: { authorization: \`Bearer \${accessToken}\` },
+});
+socket.emit('subscribe', relayId, resourceId);
+socket.on('${watched.socketEvent || 'your-event'}', payload => {});
+socket.emit('unsubscribe', relayId, resourceId);`}
+            
+
+
+
+ + +
+ + + ); +} diff --git a/src/components/router/event-relays/event-relay-list.tsx b/src/components/router/event-relays/event-relay-list.tsx new file mode 100644 index 000000000..c64e5a456 --- /dev/null +++ b/src/components/router/event-relays/event-relay-list.tsx @@ -0,0 +1,304 @@ +'use client'; + +import { useCallback, useMemo, useState } from 'react'; +import { ColumnDef } from '@tanstack/react-table'; +import { BookOpen, Plus, Radio, Zap } from 'lucide-react'; +import { Button } from '@/components/ui/button'; +import { Badge } from '@/components/ui/badge'; +import { DataTable } from '@/components/ui/data-table'; +import { EmptyState } from '@/components/ui/empty-state'; +import { Alert, AlertDescription, AlertTitle } from '@/components/ui/alert'; +import { + Dialog, + DialogContent, + DialogDescription, + DialogHeader, + DialogTitle, +} from '@/components/ui/dialog'; +import { + Tooltip, + TooltipContent, + TooltipProvider, + TooltipTrigger, +} from '@/components/ui/tooltip'; +import { + PageActions, + PageDescription, + PageHeader, + PageTitle, +} from '@/components/ui/page-header'; +import { SearchInput } from '@/components/ui/form-inputs/SearchInput'; +import { DeleteAlert } from '@/components/helpers/delete'; +import { EventRelayDocs } from '@/components/router/event-relays/event-relay-docs'; +import { EventRelayForm } from '@/components/router/event-relays/event-relay-form'; +import { EventRelay, EventRelayWriteRequest } from '@/lib/models/Router'; +import { + createEventRelay, + deleteEventRelay, + patchEventRelay, +} from '@/lib/api/router'; +import { useSettingsSave } from '@/lib/hooks/use-settings-save'; +import { useRouter } from 'next/navigation'; + +interface EventRelayListProps { + relays: EventRelay[]; + count: number; + socketsEnabled: boolean; +} + +export function EventRelayList({ + relays, + count, + socketsEnabled, +}: EventRelayListProps) { + const router = useRouter(); + const [isCreateOpen, setIsCreateOpen] = useState(false); + const [editing, setEditing] = useState(null); + const [docsOpen, setDocsOpen] = useState(count === 0); + const { save, isSaving } = useSettingsSave('Event Relay'); + + const toggleDocs = useCallback(() => { + setDocsOpen(open => { + const next = !open; + if (next) { + requestAnimationFrame(() => { + document + .getElementById('event-relay-docs') + ?.scrollIntoView({ behavior: 'smooth', block: 'start' }); + }); + } + return next; + }); + }, []); + + const refresh = useCallback(() => { + router.refresh(); + }, [router]); + + const handleCreate = async (data: EventRelayWriteRequest) => { + const result = await save({ + action: async () => { + await createEventRelay(data); + await refresh(); + }, + successMessage: 'Event relay created', + }); + if (result.ok) { + setIsCreateOpen(false); + } + }; + + const handleUpdate = async (data: EventRelayWriteRequest) => { + if (!editing) return; + const result = await save({ + action: async () => { + await patchEventRelay(editing._id, data); + await refresh(); + }, + successMessage: 'Event relay updated', + }); + if (result.ok) { + setEditing(null); + } + }; + + const columns = useMemo[]>( + () => [ + { + accessorKey: 'name', + header: 'Name', + cell: ({ row }) => ( +
+

{row.original.name}

+ {row.original.notes ? ( +

+ {row.original.notes} +

+ ) : null} +
+ ), + }, + { + accessorKey: 'busEvent', + header: 'Bus event', + cell: ({ row }) => ( + + {row.original.busEvent} + + ), + }, + { + accessorKey: 'socketEvent', + header: 'Socket event', + cell: ({ row }) => ( + + {row.original.socketEvent} + + ), + }, + { + accessorKey: 'resourceType', + header: 'Resource', + cell: ({ row }) => ( + + {row.original.resourceType}:{row.original.permission} + + ), + }, + { + accessorKey: 'active', + header: 'Status', + cell: ({ row }) => ( + + {row.original.active ? 'Active' : 'Disabled'} + + ), + }, + { + id: 'actions', + cell: ({ row }) => ( +
+ + + save({ + action: async () => { + await deleteEventRelay(row.original._id); + await refresh(); + }, + successMessage: 'Event relay deleted', + }) + } + /> +
+ ), + }, + ], + [save, refresh] + ); + + return ( +
+ +
+ Event Relays + + Send a templated socket message when an exact bus event arrives. + Subscribers join resource rooms on /events/ after a ReBAC check. + +
+ + + + + + + + {docsOpen + ? 'Hide how Event Relays work' + : 'Show how Event Relays work'} + + + + + +
+ + {!socketsEnabled ? ( + + + WebSockets are disabled + + Enable sockets in Router Settings before clients can subscribe to + relays. + + + ) : null} + + + +
+ +
+ + {relays.length === 0 ? ( + setIsCreateOpen(true)}> + + New relay + + } + /> + ) : ( + + )} + + + + + Create event relay + + Match one bus channel and emit a JSON template to the related + resource room. + + + setIsCreateOpen(false)} + isSaving={isSaving} + /> + + + + { + if (!open) setEditing(null); + }} + > + + + Edit event relay + + Changes apply immediately to new bus events. Delivery is ephemeral + and is not replayed. + + + {editing ? ( + setEditing(null)} + isSaving={isSaving} + /> + ) : null} + + +
+ ); +} diff --git a/src/components/router/event-relays/zod.ts b/src/components/router/event-relays/zod.ts new file mode 100644 index 000000000..bbffddbc4 --- /dev/null +++ b/src/components/router/event-relays/zod.ts @@ -0,0 +1,72 @@ +import { z } from 'zod'; +import { parseDotPath, RESERVED_SOCKET_EVENTS } from '@/lib/event-relays/path'; + +export const EventRelayFormSchema = z.object({ + name: z + .string() + .trim() + .min(1, 'Name is required') + .max(64, 'Name must be at most 64 characters') + .regex( + /^[A-Za-z0-9][A-Za-z0-9 _.-]{0,63}$/, + 'Name must start with a letter or number' + ), + notes: z.string().max(256).optional(), + active: z.boolean(), + busEvent: z + .string() + .trim() + .min(1, 'Bus event is required') + .max(128) + .regex( + /^[A-Za-z0-9][A-Za-z0-9_.:-]{0,127}$/, + 'Use an exact channel name with no wildcards' + ) + .refine(value => !value.includes('*'), 'Wildcards are not supported'), + socketEvent: z + .string() + .trim() + .min(1, 'Socket event is required') + .max(64) + .regex(/^[A-Za-z][A-Za-z0-9_:-]{0,63}$/, 'Socket event name is invalid') + .refine( + value => !RESERVED_SOCKET_EVENTS.has(value), + 'This socket event name is reserved' + ), + resourceType: z + .string() + .trim() + .min(1, 'Resource type is required') + .regex(/^[A-Za-z][A-Za-z0-9_]{0,63}$/, 'Resource type is invalid'), + resourceIdPath: z + .string() + .trim() + .min(1, 'Resource ID path is required') + .superRefine((value, ctx) => { + try { + parseDotPath(value); + } catch (err) { + ctx.addIssue({ + code: 'custom', + message: err instanceof Error ? err.message : 'Invalid path', + }); + } + }), + permission: z + .string() + .trim() + .min(1, 'Permission is required') + .regex(/^[A-Za-z][A-Za-z0-9_]{0,63}$/, 'Permission is invalid'), + messageTemplate: z.string().trim().min(1, 'Message template is required'), + samplePayload: z.string().optional(), +}); + +export type EventRelayFormValues = z.infer; + +export function parseJsonField(value: string, label: string): unknown { + try { + return JSON.parse(value); + } catch { + throw new Error(`${label} must be valid JSON`); + } +} diff --git a/src/lib/api/router/index.ts b/src/lib/api/router/index.ts index 006227974..5077e9b75 100644 --- a/src/lib/api/router/index.ts +++ b/src/lib/api/router/index.ts @@ -1,6 +1,11 @@ 'use server'; import { getApiClient } from '@/lib/api'; -import { RouterSettings } from '@/lib/models/Router'; +import { + EventRelay, + EventRelaysResponse, + EventRelayWriteRequest, + RouterSettings, +} from '@/lib/models/Router'; import { afterPatchServing } from '@/lib/api/modules/afterPatchServing'; import { PatchSettingsOptions } from '@/lib/api/modules/patch-settings-options'; @@ -94,3 +99,35 @@ export const patchAppRouteMiddlewares = async ( ); return res.data; }; + +export const getEventRelays = async (params?: { + skip?: number; + limit?: number; + search?: string; +}) => { + const res = await ( + await getApiClient() + ).get('/router/event-relays', { params }); + return res.data; +}; + +export const createEventRelay = async (data: EventRelayWriteRequest) => { + const res = await ( + await getApiClient() + ).post('/router/event-relays', data); + return res.data; +}; + +export const patchEventRelay = async ( + id: string, + data: Partial +) => { + const res = await ( + await getApiClient() + ).patch(`/router/event-relays/${id}`, data); + return res.data; +}; + +export const deleteEventRelay = async (id: string) => { + await (await getApiClient()).delete(`/router/event-relays/${id}`); +}; diff --git a/src/lib/event-relays/path.ts b/src/lib/event-relays/path.ts new file mode 100644 index 000000000..e364ada57 --- /dev/null +++ b/src/lib/event-relays/path.ts @@ -0,0 +1,56 @@ +export const RESERVED_SOCKET_EVENTS = new Set([ + 'connect', + 'disconnect', + 'connect_error', + 'error', + 'join-room', + 'leave-room', + 'conduit_error', + 'subscribe', + 'unsubscribe', + 'ping', + 'pong', +]); + +const FORBIDDEN_PATH_SEGMENTS = new Set([ + '__proto__', + 'constructor', + 'prototype', +]); + +const MAX_PATH_SEGMENTS = 8; +export const MAX_TEMPLATE_BYTES = 16 * 1024; + +const PATH_SEGMENT = /^[A-Za-z_][A-Za-z0-9_]*$/; + +export function parseDotPath(path: string): string[] { + const trimmed = path.trim(); + if (!trimmed) { + throw new Error('Path is required'); + } + const segments = trimmed.split('.'); + if (segments.length > MAX_PATH_SEGMENTS) { + throw new Error(`Path exceeds ${MAX_PATH_SEGMENTS} segments`); + } + for (const segment of segments) { + if (FORBIDDEN_PATH_SEGMENTS.has(segment) || !PATH_SEGMENT.test(segment)) { + throw new Error('Path contains an invalid segment'); + } + } + return segments; +} + +export function lookupOwnPath(source: unknown, path: string): unknown { + const segments = parseDotPath(path); + let current: unknown = source; + for (const segment of segments) { + if (current === null || typeof current !== 'object') { + return undefined; + } + if (!Object.prototype.hasOwnProperty.call(current, segment)) { + return undefined; + } + current = (current as Record)[segment]; + } + return current; +} diff --git a/src/lib/event-relays/preview.test.ts b/src/lib/event-relays/preview.test.ts new file mode 100644 index 000000000..50ecbd9df --- /dev/null +++ b/src/lib/event-relays/preview.test.ts @@ -0,0 +1,31 @@ +import assert from 'node:assert/strict'; +import { describe, it } from 'node:test'; +import { lookupOwnPath } from './path.ts'; +import { renderMessageTemplate } from './template.ts'; +import { previewEventRelay } from './preview.ts'; + +describe('event relay preview helpers', () => { + it('resolves the resource id and rendered payload', () => { + const result = previewEventRelay({ + resourceIdPath: '_id', + messageTemplate: { id: '{{payload._id}}', status: '{{payload.status}}' }, + samplePayload: { _id: 'order-1', status: 'paid' }, + }); + assert.equal(result.error, undefined); + assert.equal(result.resourceId, 'order-1'); + assert.deepEqual(result.payload, { id: 'order-1', status: 'paid' }); + }); + + it('fails closed on missing fields and prototype paths', () => { + assert.equal(lookupOwnPath({ a: 1 }, 'b'), undefined); + assert.throws(() => + renderMessageTemplate({ x: '{{payload.missing}}' }, {}) + ); + const preview = previewEventRelay({ + resourceIdPath: '__proto__', + messageTemplate: { id: '{{payload._id}}' }, + samplePayload: { _id: '1' }, + }); + assert.equal(typeof preview.error, 'string'); + }); +}); diff --git a/src/lib/event-relays/preview.ts b/src/lib/event-relays/preview.ts new file mode 100644 index 000000000..c085cbced --- /dev/null +++ b/src/lib/event-relays/preview.ts @@ -0,0 +1,33 @@ +import { lookupOwnPath } from './path.ts'; +import { renderMessageTemplate } from './template.ts'; + +export type RelayPreview = { + resourceId?: string; + payload?: unknown; + error?: string; +}; + +export function previewEventRelay(options: { + resourceIdPath: string; + messageTemplate: unknown; + samplePayload: unknown; +}): RelayPreview { + try { + const resourceId = lookupOwnPath( + options.samplePayload, + options.resourceIdPath + ); + if (resourceId === undefined) { + return { + error: `Resource ID path '${options.resourceIdPath}' was not found`, + }; + } + const payload = renderMessageTemplate( + options.messageTemplate, + options.samplePayload + ); + return { resourceId: String(resourceId), payload }; + } catch (err) { + return { error: err instanceof Error ? err.message : String(err) }; + } +} diff --git a/src/lib/event-relays/template.ts b/src/lib/event-relays/template.ts new file mode 100644 index 000000000..06b0808c7 --- /dev/null +++ b/src/lib/event-relays/template.ts @@ -0,0 +1,59 @@ +import { lookupOwnPath, MAX_TEMPLATE_BYTES } from './path.ts'; + +const PLACEHOLDER = /\{\{\s*payload\.([A-Za-z_][A-Za-z0-9_.]*)\s*\}\}/g; +const EXACT_PLACEHOLDER = /^\{\{\s*payload\.([A-Za-z_][A-Za-z0-9_.]*)\s*\}\}$/; +const MAX_DEPTH = 10; + +export function renderMessageTemplate( + template: unknown, + payload: unknown +): unknown { + const serialized = JSON.stringify(template); + if (!serialized || serialized.length > MAX_TEMPLATE_BYTES) { + throw new Error('Message template is invalid or too large'); + } + return renderValue(template, payload, 0); +} + +function renderValue(value: unknown, payload: unknown, depth: number): unknown { + if (depth > MAX_DEPTH) { + throw new Error('Message template is nested too deeply'); + } + if (typeof value === 'string') { + return interpolateString(value, payload); + } + if (Array.isArray(value)) { + return value.map(item => renderValue(item, payload, depth + 1)); + } + if (value !== null && typeof value === 'object') { + const record = value as Record; + const output: Record = {}; + for (const key of Object.keys(record)) { + output[key] = renderValue(record[key], payload, depth + 1); + } + return output; + } + return value; +} + +function interpolateString(value: string, payload: unknown): unknown { + const exact = value.trim().match(EXACT_PLACEHOLDER); + if (exact) { + const resolved = lookupOwnPath(payload, exact[1]); + if (resolved === undefined) { + throw new Error(`Placeholder payload.${exact[1]} was not found`); + } + return resolved; + } + + return value.replace(PLACEHOLDER, (_match, path: string) => { + const resolved = lookupOwnPath(payload, path); + if (resolved === undefined) { + throw new Error(`Placeholder payload.${path} was not found`); + } + if (resolved === null || typeof resolved !== 'object') { + return String(resolved); + } + return JSON.stringify(resolved); + }); +} diff --git a/src/lib/models/Router.ts b/src/lib/models/Router.ts index 9360c8aa3..1aa55cfbb 100644 --- a/src/lib/models/Router.ts +++ b/src/lib/models/Router.ts @@ -86,3 +86,35 @@ export type UpdateSecurityClientRequest = { alias?: string; notes?: string; }; + +export type EventRelay = { + _id: string; + name: string; + notes?: string; + active: boolean; + busEvent: string; + socketEvent: string; + resourceType: string; + resourceIdPath: string; + permission: string; + messageTemplate: unknown; + createdAt: string; + updatedAt: string; +}; + +export type EventRelaysResponse = { + relays: EventRelay[]; + count: number; +}; + +export type EventRelayWriteRequest = { + name: string; + notes?: string; + active?: boolean; + busEvent: string; + socketEvent: string; + resourceType: string; + resourceIdPath: string; + permission: string; + messageTemplate: unknown; +}; From ba21f03c22f0e5b9c3a91b30a909c8862a295fcf Mon Sep 17 00:00:00 2001 From: Konstantinos Kopanidis Date: Sun, 13 Sep 2026 02:37:40 +0300 Subject: [PATCH 2/2] fix(router): event relays UI contract (stacks on #327) (#331) * fix(router): align event relays UI with relay contract Use auth.token client snippets with subscribe on connect, documentId defaults, server preview API, inline Active toggles, API 404 fallback, and docs for subscribe-only ReBAC delivery. * fix(router): match event relay preview API and follow-ups Align preview with Conduit #1604 (messageTemplate, samplePayload, { rendered }), debounce preview inputs by string, drop live-path unsubscribe from snippets, omit active on edit PATCH, and skip the sockets warning when Router settings fail to load. * fix(router): seed edit sample by resourceIdPath, unwrap 404 Derive preview sample JSON from the relay resourceIdPath on edit, walk error.cause in isAxiosNotFoundError for RSC-wrapped axios failures, and add contract tests for sample payload and 404 detection. * fix(router): rethrow session redirect from preview API Extract preview error coercion with isNextNavigationError so debounced preview does not swallow NEXT_REDIRECT as a field error. --------- --- package.json | 2 +- .../(modules)/router/event-relays/page.tsx | 48 ++++++- .../router/event-relays/event-relay-docs.tsx | 75 ++++++++--- .../router/event-relays/event-relay-form.tsx | 124 ++++++++++-------- .../router/event-relays/event-relay-list.tsx | 56 ++++++-- .../event-relays/use-event-relay-preview.ts | 91 +++++++++++++ src/components/router/event-relays/zod.ts | 54 +++++++- src/lib/api/router/index.ts | 24 ++++ src/lib/event-relays/client-snippet.test.ts | 26 ++++ src/lib/event-relays/client-snippet.ts | 15 +++ src/lib/event-relays/path.test.ts | 16 +++ src/lib/event-relays/preview-remote.test.ts | 51 +++++++ src/lib/event-relays/preview-remote.ts | 57 ++++++++ src/lib/event-relays/preview.test.ts | 31 ----- src/lib/event-relays/preview.ts | 33 ----- src/lib/event-relays/sample-payload.test.ts | 17 +++ src/lib/event-relays/sample-payload.ts | 10 ++ src/lib/event-relays/template.ts | 59 --------- src/lib/logic/api-error.test.ts | 17 +++ src/lib/logic/api-error.ts | 23 +++- src/lib/models/Router.ts | 10 ++ 21 files changed, 625 insertions(+), 214 deletions(-) create mode 100644 src/components/router/event-relays/use-event-relay-preview.ts create mode 100644 src/lib/event-relays/client-snippet.test.ts create mode 100644 src/lib/event-relays/client-snippet.ts create mode 100644 src/lib/event-relays/path.test.ts create mode 100644 src/lib/event-relays/preview-remote.test.ts create mode 100644 src/lib/event-relays/preview-remote.ts delete mode 100644 src/lib/event-relays/preview.test.ts delete mode 100644 src/lib/event-relays/preview.ts create mode 100644 src/lib/event-relays/sample-payload.test.ts create mode 100644 src/lib/event-relays/sample-payload.ts delete mode 100644 src/lib/event-relays/template.ts create mode 100644 src/lib/logic/api-error.test.ts diff --git a/package.json b/package.json index 62ec88456..4025177d9 100644 --- a/package.json +++ b/package.json @@ -9,7 +9,7 @@ "start": "next start -p 8080", "lint": "eslint .", "test:semantic-colors": "node --experimental-strip-types --test src/lib/semantic-colors.test.ts src/lib/reactflow-edge-colors.test.ts", - "test:event-relays": "node --experimental-strip-types --test src/lib/event-relays/preview.test.ts", + "test:event-relays": "node --experimental-strip-types --test src/lib/event-relays/path.test.ts src/lib/event-relays/client-snippet.test.ts src/lib/event-relays/preview-remote.test.ts src/lib/event-relays/sample-payload.test.ts src/lib/logic/api-error.test.ts", "test:schema-fields": "node --experimental-strip-types --test src/lib/database/schema-field-definition.test.ts src/lib/database/system-schema-fields.test.ts", "build:docker": "docker build --platform linux/amd64 -t ghcr.io/conduitplatform/conduit-ui:latest .", "prepare": "husky", diff --git a/src/app/(dashboard)/(modules)/router/event-relays/page.tsx b/src/app/(dashboard)/(modules)/router/event-relays/page.tsx index 509f7094a..132958497 100644 --- a/src/app/(dashboard)/(modules)/router/event-relays/page.tsx +++ b/src/app/(dashboard)/(modules)/router/event-relays/page.tsx @@ -1,5 +1,13 @@ import { getEventRelays, getRouterSettings } from '@/lib/api/router'; import { EventRelayList } from '@/components/router/event-relays/event-relay-list'; +import { + PageDescription, + PageHeader, + PageTitle, +} from '@/components/ui/page-header'; +import { EmptyState } from '@/components/ui/empty-state'; +import { Radio } from 'lucide-react'; +import { isAxiosNotFoundError } from '@/lib/logic/api-error'; export default async function EventRelaysPage(props: { searchParams: Promise<{ @@ -11,7 +19,8 @@ export default async function EventRelaysPage(props: { const searchParams = await props.searchParams; const skip = Number(searchParams.skip ?? 0); const limit = Number(searchParams.limit ?? 10); - const [{ relays, count }, { config }] = await Promise.all([ + + const [relaysResult, settingsResult] = await Promise.allSettled([ getEventRelays({ skip, limit, @@ -20,12 +29,47 @@ export default async function EventRelaysPage(props: { getRouterSettings(), ]); + if ( + relaysResult.status === 'rejected' && + isAxiosNotFoundError(relaysResult.reason) + ) { + return ( +
+ +
+ Event Relays + + Forward exact bus events to ReBAC-scoped socket subscribers. + +
+
+
+ +
+
+ ); + } + + if (relaysResult.status === 'rejected') { + throw relaysResult.reason; + } + + const { relays, count } = relaysResult.value; + const socketsEnabled = + settingsResult.status === 'fulfilled' + ? settingsResult.value.config.transports.sockets + : undefined; + return (
); diff --git a/src/components/router/event-relays/event-relay-docs.tsx b/src/components/router/event-relays/event-relay-docs.tsx index a1c3ed895..99623a699 100644 --- a/src/components/router/event-relays/event-relay-docs.tsx +++ b/src/components/router/event-relays/event-relay-docs.tsx @@ -8,14 +8,7 @@ import { CollapsibleTrigger, } from '@/components/ui/collapsible'; import { cn } from '@/lib/utils'; - -const CLIENT_SNIPPET = `const socket = io(\`\${SOCKET_URL}/events/\`, { - path: '/realtime', - extraHeaders: { authorization: \`Bearer \${accessToken}\` }, -}); -socket.emit('subscribe', relayId, resourceId); -socket.on('order-updated', payload => {}); -socket.emit('unsubscribe', relayId, resourceId);`; +import { EVENT_RELAY_DOCS_SNIPPET } from '@/lib/event-relays/client-snippet'; const STEPS = [ { @@ -73,7 +66,8 @@ export function EventRelayDocs({ open, onOpenChange }: EventRelayDocsProps) { Forward an exact bus event to permission-scoped socket - subscribers. Not a queue, and not a generic websocket broadcast. + subscribers. Subscribe-only, ephemeral, and not a generic + websocket broadcast.

- Configure a relay + Database realtime example

- Example: notify clients when an Order document changes. + Notify clients when an Order document changes via{' '} + database:change:Order. Database realtime payloads + expose documentId (not Mongo _id on + the wire).

+
+

+ CRUD bus channel (advanced) +

+

+ You can relay database:update:Order instead, but + the payload is the full document (including{' '} + _id). That duplicates what clients already get on{' '} + /database/ change — prefer the + database realtime channel unless you only consume{' '} + /events/. +

+
+ + + +
+
+

Subscribe from a client

Connect to {'/events/'} with{' '} - {'path: /realtime'} and a user bearer token. Then - subscribe with the relay id and resource id. + {'path: /realtime'} and{' '} + {'auth: { token: accessToken }'}. Re-subscribe + inside connect so reconnects re-join the room. + There is no replay — missed events are lost.

-                {CLIENT_SNIPPET}
+                {EVENT_RELAY_DOCS_SNIPPET}
               
@@ -182,13 +212,20 @@ export function EventRelayDocs({ open, onOpenChange }: EventRelayDocsProps) { Bus channels must match exactly. Patterns like{' '} {'database:change:*'} are not supported. +
  • + Subscribe-only: clients do not publish on{' '} + /events/. Modules write to the bus. +
  • Subscribe fails closed if Authorization is unavailable or the user lacks permission.
  • - Turn a relay off with Active to stop forwarding without - deleting it. Deleting drops current subscribers immediately. + No replay or ordering guarantee. Delivery is ephemeral. +
  • +
  • + Turn a relay off with Active to stop forwarding and evict + subscribers without deleting the relay.
  • diff --git a/src/components/router/event-relays/event-relay-form.tsx b/src/components/router/event-relays/event-relay-form.tsx index 92694b7e3..f2eb40d49 100644 --- a/src/components/router/event-relays/event-relay-form.tsx +++ b/src/components/router/event-relays/event-relay-form.tsx @@ -1,6 +1,5 @@ 'use client'; -import { useMemo } from 'react'; import { useForm, useWatch } from 'react-hook-form'; import { rhfZodResolver } from '@/lib/zod-form'; import { Form } from '@/components/ui/form'; @@ -12,14 +11,17 @@ import { Alert, AlertDescription, AlertTitle } from '@/components/ui/alert'; import { EventRelayFormSchema, EventRelayFormValues, - parseJsonField, + parseMessageTemplateField, } from '@/components/router/event-relays/zod'; -import { previewEventRelay } from '@/lib/event-relays/preview'; +import { buildEventRelayClientSnippet } from '@/lib/event-relays/client-snippet'; +import { useEventRelayPreview } from '@/components/router/event-relays/use-event-relay-preview'; import { EventRelay, EventRelayWriteRequest } from '@/lib/models/Router'; +import { + buildDefaultSamplePayload, + DEFAULT_CREATE_SAMPLE_PAYLOAD, +} from '@/lib/event-relays/sample-payload'; -const DEFAULT_TEMPLATE = '{\n "id": "{{payload._id}}"\n}'; -const DEFAULT_SAMPLE = - '{\n "_id": "64f1c0a2b4d0e1f2a3b4c5d6",\n "status": "paid"\n}'; +const DEFAULT_TEMPLATE = '{\n "id": "{{payload.documentId}}"\n}'; interface EventRelayFormProps { relay?: EventRelay | null; @@ -34,6 +36,8 @@ export function EventRelayForm({ onCancel, isSaving, }: EventRelayFormProps) { + const isEditing = Boolean(relay); + const form = useForm({ resolver: rhfZodResolver(EventRelayFormSchema), defaultValues: { @@ -43,63 +47,53 @@ export function EventRelayForm({ busEvent: relay?.busEvent ?? '', socketEvent: relay?.socketEvent ?? '', resourceType: relay?.resourceType ?? '', - resourceIdPath: relay?.resourceIdPath ?? '_id', + resourceIdPath: relay?.resourceIdPath ?? 'documentId', permission: relay?.permission ?? 'read', messageTemplate: relay ? JSON.stringify(relay.messageTemplate, null, 2) : DEFAULT_TEMPLATE, - samplePayload: DEFAULT_SAMPLE, + samplePayload: relay + ? buildDefaultSamplePayload(relay.resourceIdPath) + : DEFAULT_CREATE_SAMPLE_PAYLOAD, }, }); const watched = useWatch({ control: form.control }); - const preview = useMemo(() => { - try { - const template = parseJsonField( - watched.messageTemplate ?? '', - 'Message template' - ); - const sample = parseJsonField( - watched.samplePayload?.trim() || '{}', - 'Sample payload' - ); - return previewEventRelay({ - resourceIdPath: watched.resourceIdPath || '_id', - messageTemplate: template, - samplePayload: sample, - }); - } catch (err) { - return { - error: err instanceof Error ? err.message : String(err), - }; - } - }, [watched.messageTemplate, watched.samplePayload, watched.resourceIdPath]); + const preview = useEventRelayPreview({ + messageTemplate: watched.messageTemplate ?? '', + samplePayload: watched.samplePayload ?? '', + resourceIdPath: watched.resourceIdPath ?? 'documentId', + }); const handleSubmit = form.handleSubmit(async values => { - const messageTemplate = parseJsonField( - values.messageTemplate, - 'Message template' - ); + const messageTemplate = parseMessageTemplateField(values.messageTemplate); const notes = values.notes?.trim(); - await onSubmit({ + const payload: EventRelayWriteRequest = { name: values.name, notes: notes || undefined, - active: values.active, busEvent: values.busEvent, socketEvent: values.socketEvent, resourceType: values.resourceType, resourceIdPath: values.resourceIdPath, permission: values.permission, messageTemplate, - }); + }; + if (!isEditing) { + payload.active = values.active; + } + await onSubmit(payload); }); + const clientSnippet = buildEventRelayClientSnippet( + watched.socketEvent ?? 'your-event' + ); + return (
    - + {isEditing ? null : }

    Preview

    - Local only. Nothing is published to the bus. + Rendered by the Router Admin API. Nothing is published to the bus.

    - {preview.error ? ( -

    {preview.error}

    - ) : ( + {preview.kind === 'unavailable' ? ( +

    + Preview is not available on this Router build yet. Upgrade to a + version that includes{' '} + POST /router/event-relays/preview{' '} + (see{' '} + + Conduit #1604 + + ). +

    + ) : preview.kind === 'loading' ? ( +

    Rendering…

    + ) : preview.kind === 'error' ? ( +

    {preview.message}

    + ) : preview.kind === 'ready' ? (

    Resource{' '} - {watched.resourceType || 'Type'}:{preview.resourceId} + {watched.resourceType || 'Type'}: + {preview.resourceId ?? '…'}

                     {JSON.stringify(preview.payload, null, 2)}
                   
    + ) : ( +

    + Enter valid JSON in the message template and sample payload to + preview. +

    )}
    Client contract

    - Connect to /events/ with{' '} - path: /realtime and a bearer - token. Then emit{' '} - subscribe(relayId, resourceId). + Subscribe-only: clients listen on /events/{' '} + with path: /realtime and{' '} + auth.token (browsers ignore{' '} + extraHeaders). The Authorization + module must be available — subscribe fails closed without a + matching ReBAC grant. Events are ephemeral with no replay.

    -              {`const socket = io(\`\${SOCKET_URL}/events/\`, {
    -  path: '/realtime',
    -  extraHeaders: { authorization: \`Bearer \${accessToken}\` },
    -});
    -socket.emit('subscribe', relayId, resourceId);
    -socket.on('${watched.socketEvent || 'your-event'}', payload => {});
    -socket.emit('unsubscribe', relayId, resourceId);`}
    +              {clientSnippet}
                 
    diff --git a/src/components/router/event-relays/event-relay-list.tsx b/src/components/router/event-relays/event-relay-list.tsx index c64e5a456..adf48ab45 100644 --- a/src/components/router/event-relays/event-relay-list.tsx +++ b/src/components/router/event-relays/event-relay-list.tsx @@ -4,7 +4,8 @@ import { useCallback, useMemo, useState } from 'react'; import { ColumnDef } from '@tanstack/react-table'; import { BookOpen, Plus, Radio, Zap } from 'lucide-react'; import { Button } from '@/components/ui/button'; -import { Badge } from '@/components/ui/badge'; +import { Switch } from '@/components/ui/switch'; +import { Label } from '@/components/ui/label'; import { DataTable } from '@/components/ui/data-table'; import { EmptyState } from '@/components/ui/empty-state'; import { Alert, AlertDescription, AlertTitle } from '@/components/ui/alert'; @@ -43,7 +44,7 @@ import { useRouter } from 'next/navigation'; interface EventRelayListProps { relays: EventRelay[]; count: number; - socketsEnabled: boolean; + socketsEnabled?: boolean; } export function EventRelayList({ @@ -55,6 +56,7 @@ export function EventRelayList({ const [isCreateOpen, setIsCreateOpen] = useState(false); const [editing, setEditing] = useState(null); const [docsOpen, setDocsOpen] = useState(count === 0); + const [togglingId, setTogglingId] = useState(null); const { save, isSaving } = useSettingsSave('Event Relay'); const toggleDocs = useCallback(() => { @@ -88,6 +90,23 @@ export function EventRelayList({ } }; + const handleActiveToggle = useCallback( + async (relay: EventRelay, active: boolean) => { + setTogglingId(relay._id); + await save({ + action: async () => { + await patchEventRelay(relay._id, { active }); + await refresh(); + }, + successMessage: active + ? 'Event relay enabled' + : 'Event relay disabled', + }); + setTogglingId(null); + }, + [refresh, save] + ); + const handleUpdate = async (data: EventRelayWriteRequest) => { if (!editing) return; const result = await save({ @@ -147,12 +166,28 @@ export function EventRelayList({ }, { accessorKey: 'active', - header: 'Status', - cell: ({ row }) => ( - - {row.original.active ? 'Active' : 'Disabled'} - - ), + header: 'Active', + cell: ({ row }) => { + const relay = row.original; + const busy = togglingId === relay._id || isSaving; + return ( +
    + handleActiveToggle(relay, checked)} + aria-label={`${relay.active ? 'Disable' : 'Enable'} ${relay.name}`} + /> + +
    + ); + }, }, { id: 'actions', @@ -162,6 +197,7 @@ export function EventRelayList({ type="button" variant="outline" size="sm" + disabled={togglingId === row.original._id} onClick={() => setEditing(row.original)} > Edit @@ -183,7 +219,7 @@ export function EventRelayList({ ), }, ], - [save, refresh] + [handleActiveToggle, isSaving, save, refresh, togglingId] ); return ( @@ -225,7 +261,7 @@ export function EventRelayList({ - {!socketsEnabled ? ( + {socketsEnabled === false ? ( WebSockets are disabled diff --git a/src/components/router/event-relays/use-event-relay-preview.ts b/src/components/router/event-relays/use-event-relay-preview.ts new file mode 100644 index 000000000..61cf35e50 --- /dev/null +++ b/src/components/router/event-relays/use-event-relay-preview.ts @@ -0,0 +1,91 @@ +'use client'; + +import { useEffect, useState } from 'react'; +import { useDebounce } from '@uidotdev/usehooks'; +import { previewEventRelayRemote } from '@/lib/api/router'; +import { lookupOwnPath } from '@/lib/event-relays/path'; + +export type EventRelayPreviewState = + | { kind: 'idle' } + | { kind: 'loading' } + | { kind: 'unavailable' } + | { kind: 'error'; message: string } + | { kind: 'ready'; payload: unknown; resourceId?: string }; + +function tryParseJson(value: string): unknown | null { + try { + return JSON.parse(value); + } catch { + return null; + } +} + +export function useEventRelayPreview(options: { + messageTemplate: string; + samplePayload: string; + resourceIdPath: string; +}): EventRelayPreviewState { + const debouncedMessageTemplate = useDebounce(options.messageTemplate, 400); + const debouncedSamplePayload = useDebounce(options.samplePayload, 400); + const debouncedResourceIdPath = useDebounce(options.resourceIdPath, 400); + + const [state, setState] = useState({ kind: 'idle' }); + + useEffect(() => { + const template = tryParseJson(debouncedMessageTemplate); + if (template === null) { + setState({ kind: 'idle' }); + return; + } + + const sample = tryParseJson(debouncedSamplePayload?.trim() || '{}'); + if (sample === null) { + setState({ kind: 'idle' }); + return; + } + + let cancelled = false; + void (async () => { + setState({ kind: 'loading' }); + const result = await previewEventRelayRemote({ + messageTemplate: template, + samplePayload: sample, + }); + if (cancelled) return; + + if (result.status === 'unavailable') { + setState({ kind: 'unavailable' }); + return; + } + if (result.status === 'error') { + setState({ kind: 'error', message: result.message }); + return; + } + + let resourceId: string | undefined; + try { + const resolved = lookupOwnPath( + sample, + debouncedResourceIdPath.trim() || 'documentId' + ); + if (resolved !== undefined) { + resourceId = String(resolved); + } + } catch { + resourceId = undefined; + } + + setState({ + kind: 'ready', + payload: result.payload, + resourceId, + }); + })(); + + return () => { + cancelled = true; + }; + }, [debouncedMessageTemplate, debouncedSamplePayload, debouncedResourceIdPath]); + + return state; +} diff --git a/src/components/router/event-relays/zod.ts b/src/components/router/event-relays/zod.ts index bbffddbc4..8d4b34ae6 100644 --- a/src/components/router/event-relays/zod.ts +++ b/src/components/router/event-relays/zod.ts @@ -1,6 +1,44 @@ import { z } from 'zod'; import { parseDotPath, RESERVED_SOCKET_EVENTS } from '@/lib/event-relays/path'; +function jsonObjectString(label: string) { + return z + .string() + .trim() + .min(1, `${label} is required`) + .superRefine((value, ctx) => { + try { + const parsed = JSON.parse(value); + if (parsed === null || typeof parsed !== 'object' || Array.isArray(parsed)) { + ctx.addIssue({ + code: 'custom', + message: `${label} must be a JSON object`, + }); + } + } catch { + ctx.addIssue({ + code: 'custom', + message: `${label} must be valid JSON`, + }); + } + }); +} + +const optionalJsonObjectString = z + .string() + .optional() + .superRefine((value, ctx) => { + if (!value?.trim()) return; + try { + JSON.parse(value); + } catch { + ctx.addIssue({ + code: 'custom', + message: 'Sample payload must be valid JSON', + }); + } + }); + export const EventRelayFormSchema = z.object({ name: z .string() @@ -57,8 +95,8 @@ export const EventRelayFormSchema = z.object({ .trim() .min(1, 'Permission is required') .regex(/^[A-Za-z][A-Za-z0-9_]{0,63}$/, 'Permission is invalid'), - messageTemplate: z.string().trim().min(1, 'Message template is required'), - samplePayload: z.string().optional(), + messageTemplate: jsonObjectString('Message template'), + samplePayload: optionalJsonObjectString, }); export type EventRelayFormValues = z.infer; @@ -70,3 +108,15 @@ export function parseJsonField(value: string, label: string): unknown { throw new Error(`${label} must be valid JSON`); } } + +export function parseMessageTemplateField(value: string): Record { + const parsed = parseJsonField(value, 'Message template'); + if ( + parsed === null || + typeof parsed !== 'object' || + Array.isArray(parsed) + ) { + throw new Error('Message template must be a JSON object'); + } + return parsed as Record; +} diff --git a/src/lib/api/router/index.ts b/src/lib/api/router/index.ts index 5077e9b75..bc22dadc5 100644 --- a/src/lib/api/router/index.ts +++ b/src/lib/api/router/index.ts @@ -2,10 +2,17 @@ import { getApiClient } from '@/lib/api'; import { EventRelay, + EventRelayPreviewInput, + EventRelayPreviewRemoteResult, EventRelaysResponse, EventRelayWriteRequest, RouterSettings, } from '@/lib/models/Router'; +import { + buildEventRelayPreviewRequestBody, + coerceEventRelayPreviewRemoteError, + parseEventRelayPreviewResponse, +} from '@/lib/event-relays/preview-remote'; import { afterPatchServing } from '@/lib/api/modules/afterPatchServing'; import { PatchSettingsOptions } from '@/lib/api/modules/patch-settings-options'; @@ -131,3 +138,20 @@ export const patchEventRelay = async ( export const deleteEventRelay = async (id: string) => { await (await getApiClient()).delete(`/router/event-relays/${id}`); }; + +export const previewEventRelayRemote = async ( + body: EventRelayPreviewInput +): Promise => { + try { + const res = await ( + await getApiClient() + ).post( + '/router/event-relays/preview', + buildEventRelayPreviewRequestBody(body) + ); + const payload = parseEventRelayPreviewResponse(res.data); + return { status: 'ok', payload }; + } catch (err) { + return coerceEventRelayPreviewRemoteError(err); + } +}; diff --git a/src/lib/event-relays/client-snippet.test.ts b/src/lib/event-relays/client-snippet.test.ts new file mode 100644 index 000000000..f11bc4b03 --- /dev/null +++ b/src/lib/event-relays/client-snippet.test.ts @@ -0,0 +1,26 @@ +import assert from 'node:assert/strict'; +import { describe, it } from 'node:test'; +import { + buildEventRelayClientSnippet, + EVENT_RELAY_DOCS_SNIPPET, +} from './client-snippet.ts'; + +describe('event relay client snippet', () => { + it('uses auth.token and subscribe on connect', () => { + const snippet = buildEventRelayClientSnippet('order-updated'); + assert.match(snippet, /auth:\s*\{\s*token:\s*accessToken\s*\}/); + assert.match(snippet, /socket\.on\('connect',\s*\(\)\s*=>\s*\{/); + assert.match(snippet, /socket\.emit\('subscribe',\s*relayId,\s*resourceId\)/); + assert.match(snippet, /socket\.on\('order-updated',\s*payload\s*=>\s*\{\}\)/); + }); + + it('does not use extraHeaders or a live-path unsubscribe', () => { + for (const snippet of [ + buildEventRelayClientSnippet('x'), + EVENT_RELAY_DOCS_SNIPPET, + ]) { + assert.equal(snippet.includes('extraHeaders'), false); + assert.equal(snippet.includes("emit('unsubscribe'"), false); + } + }); +}); diff --git a/src/lib/event-relays/client-snippet.ts b/src/lib/event-relays/client-snippet.ts new file mode 100644 index 000000000..aa5eef912 --- /dev/null +++ b/src/lib/event-relays/client-snippet.ts @@ -0,0 +1,15 @@ +export function buildEventRelayClientSnippet(socketEvent: string): string { + const eventHandler = socketEvent.trim() || 'your-event'; + return `const socket = io(\`\${SOCKET_URL}/events/\`, { + path: '/realtime', + auth: { token: accessToken }, +}); +socket.on('connect', () => { + socket.emit('subscribe', relayId, resourceId); +}); +socket.on('${eventHandler}', payload => {});`; +} + +export const EVENT_RELAY_DOCS_SNIPPET = buildEventRelayClientSnippet( + 'order-updated' +); diff --git a/src/lib/event-relays/path.test.ts b/src/lib/event-relays/path.test.ts new file mode 100644 index 000000000..7bd0936bb --- /dev/null +++ b/src/lib/event-relays/path.test.ts @@ -0,0 +1,16 @@ +import assert from 'node:assert/strict'; +import { describe, it } from 'node:test'; +import { lookupOwnPath, parseDotPath } from './path.ts'; + +describe('event relay path helpers', () => { + it('resolves dot paths on own properties only', () => { + assert.deepEqual(parseDotPath('documentId'), ['documentId']); + assert.equal(lookupOwnPath({ documentId: 'abc' }, 'documentId'), 'abc'); + assert.equal(lookupOwnPath({ a: 1 }, 'b'), undefined); + }); + + it('rejects prototype paths', () => { + assert.throws(() => parseDotPath('__proto__')); + assert.throws(() => parseDotPath('constructor')); + }); +}); diff --git a/src/lib/event-relays/preview-remote.test.ts b/src/lib/event-relays/preview-remote.test.ts new file mode 100644 index 000000000..b6e5a365f --- /dev/null +++ b/src/lib/event-relays/preview-remote.test.ts @@ -0,0 +1,51 @@ +import assert from 'node:assert/strict'; +import { describe, it } from 'node:test'; +import { + buildEventRelayPreviewRequestBody, + coerceEventRelayPreviewRemoteError, + parseEventRelayPreviewResponse, +} from './preview-remote.ts'; + +describe('event relay preview remote contract', () => { + it('builds the Admin API request body', () => { + const body = buildEventRelayPreviewRequestBody({ + messageTemplate: { id: '{{payload.documentId}}' }, + samplePayload: { documentId: 'abc' }, + }); + assert.deepEqual(body, { + messageTemplate: { id: '{{payload.documentId}}' }, + samplePayload: { documentId: 'abc' }, + }); + assert.equal('template' in body, false); + assert.equal('sample' in body, false); + }); + + it('unwraps rendered from the Admin API response', () => { + const rendered = parseEventRelayPreviewResponse({ + rendered: { id: 'abc', status: 'paid' }, + }); + assert.deepEqual(rendered, { id: 'abc', status: 'paid' }); + }); + + it('rejects responses without rendered', () => { + assert.throws(() => parseEventRelayPreviewResponse({ payload: {} })); + assert.throws(() => parseEventRelayPreviewResponse(null)); + }); + + it('rethrows Next navigation errors for session redirect', () => { + const redirectErr = Object.assign(new Error('NEXT_REDIRECT'), { + digest: 'NEXT_REDIRECT;replace;/login?session-timeout=true', + }); + assert.throws( + () => coerceEventRelayPreviewRemoteError(redirectErr), + redirectErr + ); + }); + + it('maps axios 404 to unavailable', () => { + const result = coerceEventRelayPreviewRemoteError({ + response: { status: 404 }, + }); + assert.deepEqual(result, { status: 'unavailable' }); + }); +}); diff --git a/src/lib/event-relays/preview-remote.ts b/src/lib/event-relays/preview-remote.ts new file mode 100644 index 000000000..6ac141e9d --- /dev/null +++ b/src/lib/event-relays/preview-remote.ts @@ -0,0 +1,57 @@ +import { + formatAdminApiError, + isAxiosNotFoundError, + isNextNavigationError, +} from '../logic/api-error.ts'; +import type { EventRelayPreviewRemoteResult } from '../models/Router.ts'; + +export type EventRelayPreviewInput = { + messageTemplate: unknown; + samplePayload: unknown; +}; + +export type EventRelayPreviewRequestBody = { + messageTemplate: unknown; + samplePayload: unknown; +}; + +export type EventRelayPreviewResponseBody = { + rendered: unknown; +}; + +export function buildEventRelayPreviewRequestBody( + input: EventRelayPreviewInput +): EventRelayPreviewRequestBody { + return { + messageTemplate: input.messageTemplate, + samplePayload: input.samplePayload, + }; +} + +export function parseEventRelayPreviewResponse( + data: unknown +): EventRelayPreviewResponseBody['rendered'] { + if (data === null || typeof data !== 'object' || Array.isArray(data)) { + throw new Error('Preview response is invalid'); + } + const record = data as Record; + if (!('rendered' in record)) { + throw new Error('Preview response is missing rendered'); + } + return record.rendered; +} + +export function coerceEventRelayPreviewRemoteError( + err: unknown +): EventRelayPreviewRemoteResult { + if (isNextNavigationError(err)) { + throw err; + } + if (isAxiosNotFoundError(err)) { + return { status: 'unavailable' }; + } + if (err instanceof Error && err.message.startsWith('Preview response')) { + return { status: 'error', message: err.message }; + } + return { status: 'error', message: formatAdminApiError(err) }; +} diff --git a/src/lib/event-relays/preview.test.ts b/src/lib/event-relays/preview.test.ts deleted file mode 100644 index 50ecbd9df..000000000 --- a/src/lib/event-relays/preview.test.ts +++ /dev/null @@ -1,31 +0,0 @@ -import assert from 'node:assert/strict'; -import { describe, it } from 'node:test'; -import { lookupOwnPath } from './path.ts'; -import { renderMessageTemplate } from './template.ts'; -import { previewEventRelay } from './preview.ts'; - -describe('event relay preview helpers', () => { - it('resolves the resource id and rendered payload', () => { - const result = previewEventRelay({ - resourceIdPath: '_id', - messageTemplate: { id: '{{payload._id}}', status: '{{payload.status}}' }, - samplePayload: { _id: 'order-1', status: 'paid' }, - }); - assert.equal(result.error, undefined); - assert.equal(result.resourceId, 'order-1'); - assert.deepEqual(result.payload, { id: 'order-1', status: 'paid' }); - }); - - it('fails closed on missing fields and prototype paths', () => { - assert.equal(lookupOwnPath({ a: 1 }, 'b'), undefined); - assert.throws(() => - renderMessageTemplate({ x: '{{payload.missing}}' }, {}) - ); - const preview = previewEventRelay({ - resourceIdPath: '__proto__', - messageTemplate: { id: '{{payload._id}}' }, - samplePayload: { _id: '1' }, - }); - assert.equal(typeof preview.error, 'string'); - }); -}); diff --git a/src/lib/event-relays/preview.ts b/src/lib/event-relays/preview.ts deleted file mode 100644 index c085cbced..000000000 --- a/src/lib/event-relays/preview.ts +++ /dev/null @@ -1,33 +0,0 @@ -import { lookupOwnPath } from './path.ts'; -import { renderMessageTemplate } from './template.ts'; - -export type RelayPreview = { - resourceId?: string; - payload?: unknown; - error?: string; -}; - -export function previewEventRelay(options: { - resourceIdPath: string; - messageTemplate: unknown; - samplePayload: unknown; -}): RelayPreview { - try { - const resourceId = lookupOwnPath( - options.samplePayload, - options.resourceIdPath - ); - if (resourceId === undefined) { - return { - error: `Resource ID path '${options.resourceIdPath}' was not found`, - }; - } - const payload = renderMessageTemplate( - options.messageTemplate, - options.samplePayload - ); - return { resourceId: String(resourceId), payload }; - } catch (err) { - return { error: err instanceof Error ? err.message : String(err) }; - } -} diff --git a/src/lib/event-relays/sample-payload.test.ts b/src/lib/event-relays/sample-payload.test.ts new file mode 100644 index 000000000..476fa6cba --- /dev/null +++ b/src/lib/event-relays/sample-payload.test.ts @@ -0,0 +1,17 @@ +import assert from 'node:assert/strict'; +import { describe, it } from 'node:test'; +import { buildDefaultSamplePayload } from './sample-payload.ts'; + +describe('default event relay sample payload', () => { + it('uses documentId for database realtime paths', () => { + const json = buildDefaultSamplePayload('documentId'); + assert.match(json, /"documentId"/); + assert.doesNotMatch(json, /"_id"/); + }); + + it('uses _id for CRUD bus paths', () => { + const json = buildDefaultSamplePayload('_id'); + assert.match(json, /"_id"/); + assert.doesNotMatch(json, /"documentId"/); + }); +}); diff --git a/src/lib/event-relays/sample-payload.ts b/src/lib/event-relays/sample-payload.ts new file mode 100644 index 000000000..5ff16a7db --- /dev/null +++ b/src/lib/event-relays/sample-payload.ts @@ -0,0 +1,10 @@ +const SAMPLE_ID = '64f1c0a2b4d0e1f2a3b4c5d6'; + +export function buildDefaultSamplePayload(resourceIdPath: string): string { + const path = resourceIdPath.trim() || 'documentId'; + const key = path.split('.')[0]; + return JSON.stringify({ [key]: SAMPLE_ID, status: 'paid' }, null, 2); +} + +export const DEFAULT_CREATE_SAMPLE_PAYLOAD = + buildDefaultSamplePayload('documentId'); diff --git a/src/lib/event-relays/template.ts b/src/lib/event-relays/template.ts deleted file mode 100644 index 06b0808c7..000000000 --- a/src/lib/event-relays/template.ts +++ /dev/null @@ -1,59 +0,0 @@ -import { lookupOwnPath, MAX_TEMPLATE_BYTES } from './path.ts'; - -const PLACEHOLDER = /\{\{\s*payload\.([A-Za-z_][A-Za-z0-9_.]*)\s*\}\}/g; -const EXACT_PLACEHOLDER = /^\{\{\s*payload\.([A-Za-z_][A-Za-z0-9_.]*)\s*\}\}$/; -const MAX_DEPTH = 10; - -export function renderMessageTemplate( - template: unknown, - payload: unknown -): unknown { - const serialized = JSON.stringify(template); - if (!serialized || serialized.length > MAX_TEMPLATE_BYTES) { - throw new Error('Message template is invalid or too large'); - } - return renderValue(template, payload, 0); -} - -function renderValue(value: unknown, payload: unknown, depth: number): unknown { - if (depth > MAX_DEPTH) { - throw new Error('Message template is nested too deeply'); - } - if (typeof value === 'string') { - return interpolateString(value, payload); - } - if (Array.isArray(value)) { - return value.map(item => renderValue(item, payload, depth + 1)); - } - if (value !== null && typeof value === 'object') { - const record = value as Record; - const output: Record = {}; - for (const key of Object.keys(record)) { - output[key] = renderValue(record[key], payload, depth + 1); - } - return output; - } - return value; -} - -function interpolateString(value: string, payload: unknown): unknown { - const exact = value.trim().match(EXACT_PLACEHOLDER); - if (exact) { - const resolved = lookupOwnPath(payload, exact[1]); - if (resolved === undefined) { - throw new Error(`Placeholder payload.${exact[1]} was not found`); - } - return resolved; - } - - return value.replace(PLACEHOLDER, (_match, path: string) => { - const resolved = lookupOwnPath(payload, path); - if (resolved === undefined) { - throw new Error(`Placeholder payload.${path} was not found`); - } - if (resolved === null || typeof resolved !== 'object') { - return String(resolved); - } - return JSON.stringify(resolved); - }); -} diff --git a/src/lib/logic/api-error.test.ts b/src/lib/logic/api-error.test.ts new file mode 100644 index 000000000..a00a0645c --- /dev/null +++ b/src/lib/logic/api-error.test.ts @@ -0,0 +1,17 @@ +import assert from 'node:assert/strict'; +import { describe, it } from 'node:test'; +import { isAxiosNotFoundError } from './api-error.ts'; + +describe('isAxiosNotFoundError', () => { + it('detects 404 on the error or its cause chain', () => { + const axios404 = { response: { status: 404 } }; + assert.equal(isAxiosNotFoundError(axios404), true); + assert.equal( + isAxiosNotFoundError( + Object.assign(new Error('Server Components render'), { cause: axios404 }) + ), + true + ); + assert.equal(isAxiosNotFoundError({ response: { status: 500 } }), false); + }); +}); diff --git a/src/lib/logic/api-error.ts b/src/lib/logic/api-error.ts index c95fc1508..ce7877327 100644 --- a/src/lib/logic/api-error.ts +++ b/src/lib/logic/api-error.ts @@ -10,9 +10,28 @@ export function isAxiosLikeError(err: unknown): err is AxiosLikeError { return Boolean(err && typeof err === 'object' && 'response' in err); } +function collectErrorChain(err: unknown): unknown[] { + const chain: unknown[] = []; + const seen = new Set(); + let current: unknown = err; + while (current !== undefined && current !== null && !seen.has(current)) { + seen.add(current); + chain.push(current); + if (typeof current !== 'object' || !('cause' in current)) { + break; + } + current = (current as { cause?: unknown }).cause; + } + return chain; +} + export function getAxiosResponseStatus(err: unknown): number | undefined { - if (!isAxiosLikeError(err)) return undefined; - return err.response?.status; + for (const candidate of collectErrorChain(err)) { + if (!isAxiosLikeError(candidate)) continue; + const status = candidate.response?.status; + if (status !== undefined) return status; + } + return undefined; } export function isAxiosNotFoundError(err: unknown): boolean { diff --git a/src/lib/models/Router.ts b/src/lib/models/Router.ts index 1aa55cfbb..8cae38279 100644 --- a/src/lib/models/Router.ts +++ b/src/lib/models/Router.ts @@ -118,3 +118,13 @@ export type EventRelayWriteRequest = { permission: string; messageTemplate: unknown; }; + +export type EventRelayPreviewInput = { + messageTemplate: unknown; + samplePayload: unknown; +}; + +export type EventRelayPreviewRemoteResult = + | { status: 'ok'; payload: unknown } + | { status: 'error'; message: string } + | { status: 'unavailable' };