diff --git a/package.json b/package.json index e5f90b63c..4025177d9 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/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 new file mode 100644 index 000000000..132958497 --- /dev/null +++ b/src/app/(dashboard)/(modules)/router/event-relays/page.tsx @@ -0,0 +1,76 @@ +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<{ + skip?: string; + limit?: string; + search?: string; + }>; +}) { + const searchParams = await props.searchParams; + const skip = Number(searchParams.skip ?? 0); + const limit = Number(searchParams.limit ?? 10); + + const [relaysResult, settingsResult] = await Promise.allSettled([ + getEventRelays({ + skip, + limit, + search: searchParams.search, + }), + 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/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..99623a699 --- /dev/null +++ b/src/components/router/event-relays/event-relay-docs.tsx @@ -0,0 +1,267 @@ +'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'; +import { EVENT_RELAY_DOCS_SNIPPET } from '@/lib/event-relays/client-snippet'; + +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. Subscribe-only, ephemeral, 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. + ))} +
+
+ +
+

+ Database realtime example +

+

+ 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{' '} + {'auth: { token: accessToken }'}. Re-subscribe + inside connect so reconnects re-join the room. + There is no replay — missed events are lost. +

+
+                {EVENT_RELAY_DOCS_SNIPPET}
+              
+
+ +
+

Limits

+
    +
  • + 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. +
  • +
  • + No replay or ordering guarantee. Delivery is ephemeral. +
  • +
  • + Turn a relay off with Active to stop forwarding and evict + subscribers without deleting the relay. +
  • +
+
+
+
+
+
+ ); +} + +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..f2eb40d49 --- /dev/null +++ b/src/components/router/event-relays/event-relay-form.tsx @@ -0,0 +1,227 @@ +'use client'; + +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, + parseMessageTemplateField, +} from '@/components/router/event-relays/zod'; +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.documentId}}"\n}'; + +interface EventRelayFormProps { + relay?: EventRelay | null; + onSubmit: (data: EventRelayWriteRequest) => Promise; + onCancel: () => void; + isSaving?: boolean; +} + +export function EventRelayForm({ + relay, + onSubmit, + onCancel, + isSaving, +}: EventRelayFormProps) { + const isEditing = Boolean(relay); + + 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 ?? 'documentId', + permission: relay?.permission ?? 'read', + messageTemplate: relay + ? JSON.stringify(relay.messageTemplate, null, 2) + : DEFAULT_TEMPLATE, + samplePayload: relay + ? buildDefaultSamplePayload(relay.resourceIdPath) + : DEFAULT_CREATE_SAMPLE_PAYLOAD, + }, + }); + + const watched = useWatch({ control: form.control }); + const preview = useEventRelayPreview({ + messageTemplate: watched.messageTemplate ?? '', + samplePayload: watched.samplePayload ?? '', + resourceIdPath: watched.resourceIdPath ?? 'documentId', + }); + + const handleSubmit = form.handleSubmit(async values => { + const messageTemplate = parseMessageTemplateField(values.messageTemplate); + const notes = values.notes?.trim(); + const payload: EventRelayWriteRequest = { + name: values.name, + notes: notes || undefined, + 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

+

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

+ {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 ?? '…'} + +

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

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

+ )} +
+ + Client contract + +

+ 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. +

+
+              {clientSnippet}
+            
+
+
+
+ + +
+ + + ); +} 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..adf48ab45 --- /dev/null +++ b/src/components/router/event-relays/event-relay-list.tsx @@ -0,0 +1,340 @@ +'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 { 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'; +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 [togglingId, setTogglingId] = useState(null); + 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 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({ + 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: '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', + cell: ({ row }) => ( +
+ + + save({ + action: async () => { + await deleteEventRelay(row.original._id); + await refresh(); + }, + successMessage: 'Event relay deleted', + }) + } + /> +
+ ), + }, + ], + [handleActiveToggle, isSaving, save, refresh, togglingId] + ); + + 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 === false ? ( + + + 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/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 new file mode 100644 index 000000000..8d4b34ae6 --- /dev/null +++ b/src/components/router/event-relays/zod.ts @@ -0,0 +1,122 @@ +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() + .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: jsonObjectString('Message template'), + samplePayload: optionalJsonObjectString, +}); + +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`); + } +} + +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 006227974..bc22dadc5 100644 --- a/src/lib/api/router/index.ts +++ b/src/lib/api/router/index.ts @@ -1,6 +1,18 @@ 'use server'; import { getApiClient } from '@/lib/api'; -import { RouterSettings } from '@/lib/models/Router'; +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'; @@ -94,3 +106,52 @@ 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}`); +}; + +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/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-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/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/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 9360c8aa3..8cae38279 100644 --- a/src/lib/models/Router.ts +++ b/src/lib/models/Router.ts @@ -86,3 +86,45 @@ 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; +}; + +export type EventRelayPreviewInput = { + messageTemplate: unknown; + samplePayload: unknown; +}; + +export type EventRelayPreviewRemoteResult = + | { status: 'ok'; payload: unknown } + | { status: 'error'; message: string } + | { status: 'unavailable' };