Skip to content

Commit 6022978

Browse files
Terbaubrage-andreas
authored andcommitted
feat: realtime change notifications
1 parent bbb9e7c commit 6022978

7 files changed

Lines changed: 100 additions & 48 deletions

File tree

apps/rpc/src/bin/server.ts

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -59,10 +59,24 @@ controller.signal.addEventListener("abort", () => {
5959
})
6060

6161
export async function createFastifyContext({ req }: CreateFastifyContextOptions) {
62+
let rawToken: string | undefined
63+
6264
const bearer = req.headers.authorization
6365
if (bearer !== undefined) {
64-
const token = bearer.substring("Bearer ".length)
65-
const principal = await dependencies.rpcJwtService.verify(token)
66+
rawToken = bearer.substring("Bearer ".length)
67+
} else if (typeof req.query === "object" && req.query !== null && "connectionParams" in req.query) {
68+
try {
69+
const params = JSON.parse(req.query.connectionParams as string)
70+
if (typeof params.token === "string") {
71+
rawToken = params.token
72+
}
73+
} catch {
74+
// malformed connectionParams — treat as unauthenticated
75+
}
76+
}
77+
78+
if (rawToken !== undefined) {
79+
const principal = await dependencies.rpcJwtService.verify(rawToken)
6680
const subject = principal.payload.sub
6781
if (subject === undefined) {
6882
return createTrpcContext(null, serviceLayer)

apps/rpc/src/modules/core.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -189,7 +189,7 @@ export async function createServiceLayer(
189189
const feedbackFormAnswerRepository = getFeedbackFormAnswerRepository()
190190
const notificationRepository = getNotificationRepository()
191191

192-
const notificationService = getNotificationService(notificationRepository, userRepository, attendanceRepository)
192+
const notificationService = getNotificationService(notificationRepository, userRepository, attendanceRepository, eventEmitter)
193193
const membershipService = getMembershipService()
194194
const emailService = isAmazonSesEmailFeatureEnabled(configuration)
195195
? getEmailService(clients.sesClient, clients.sqsClient, configuration)

apps/rpc/src/modules/notification/notification-router.ts

Lines changed: 35 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
1+
import { on } from "node:events"
12
import type { inferProcedureInput, inferProcedureOutput } from "@trpc/server"
23
import { z } from "zod"
3-
import { isEditor } from "../../authorization"
44
import { withAuditLogEntry, withAuthentication, withAuthorization, withDatabaseTransaction } from "../../middlewares"
55
import { procedure, t } from "../../trpc"
66
import { BasePaginateInputSchema, PaginateInputSchema } from "@dotkomonline/utils"
@@ -11,6 +11,8 @@ import {
1111
UserNotificationDTOSchema,
1212
UserNotificationSchema,
1313
} from "./notification-types"
14+
import type { Notification } from "./notification-types"
15+
import { isCommitteeMember } from "src/authorization"
1416

1517
export type GetNotificationInput = inferProcedureInput<typeof getNotificationProcedure>
1618
export type GetNotificationOutput = inferProcedureOutput<typeof getNotificationProcedure>
@@ -28,11 +30,20 @@ const createNotificationProcedure = procedure
2830
.input(NotificationWriteSchema)
2931
.output(NotificationDTOSchema)
3032
.use(withAuthentication())
31-
.use(withAuthorization(isEditor()))
33+
.use(withAuthorization(isCommitteeMember()))
3234
.use(withDatabaseTransaction())
3335
.use(withAuditLogEntry())
3436
.mutation(({ input, ctx }) => {
35-
return ctx.notificationService.createWithRecipients(ctx.handle, input)
37+
return ctx.notificationService.create(
38+
ctx.handle,
39+
input.recipientIds,
40+
input.type,
41+
input.title,
42+
input.shortDescription ?? input.title,
43+
input.actorGroupId,
44+
input.payloadType,
45+
input.payload
46+
)
3647
})
3748

3849
export type EditNotificationInput = inferProcedureInput<typeof editNotificationProcedure>
@@ -46,7 +57,8 @@ const editNotificationProcedure = procedure
4657
)
4758
.output(NotificationDTOSchema)
4859
.use(withAuthentication())
49-
.use(withAuthorization(isEditor()))
60+
.use(withAuthorization(isCommitteeMember()))
61+
5062
.use(withDatabaseTransaction())
5163
.use(withAuditLogEntry())
5264
.mutation(async ({ input: changes, ctx }) => {
@@ -59,7 +71,8 @@ const deleteNotificationProcedure = procedure
5971
.input(NotificationSchema.shape.id)
6072
.output(z.boolean())
6173
.use(withAuthentication())
62-
.use(withAuthorization(isEditor()))
74+
.use(withAuthorization(isCommitteeMember()))
75+
6376
.use(withDatabaseTransaction())
6477
.use(withAuditLogEntry())
6578
.mutation(async ({ input, ctx }) => {
@@ -118,7 +131,8 @@ const findNotificationsProcedure = procedure
118131
.input(PaginateInputSchema)
119132
.output(z.object({ items: z.array(NotificationDTOSchema), nextCursor: NotificationSchema.shape.id.optional() }))
120133
.use(withAuthentication())
121-
.use(withAuthorization(isEditor()))
134+
.use(withAuthorization(isCommitteeMember()))
135+
122136
.use(withDatabaseTransaction())
123137
.query(async ({ input, ctx }) => {
124138
const items = await ctx.notificationService.findMany(ctx.handle, input)
@@ -129,6 +143,20 @@ const findNotificationsProcedure = procedure
129143
}
130144
})
131145

146+
export type OnNewNotificationInput = inferProcedureInput<typeof onNewNotificationProcedure>
147+
export type OnNewNotificationOutput = inferProcedureOutput<typeof onNewNotificationProcedure>
148+
const onNewNotificationProcedure = procedure
149+
.use(withAuthentication())
150+
.subscription(async function* ({ ctx, signal }) {
151+
for await (const [data] of on(ctx.eventEmitter, "notification:new", { signal })) {
152+
const { userId, notification } = data as { userId: string; notification: Notification }
153+
if (userId !== ctx.principal.subject) {
154+
continue
155+
}
156+
yield notification
157+
}
158+
})
159+
132160
export const notificationRouter = t.router({
133161
get: getNotificationProcedure,
134162
create: createNotificationProcedure,
@@ -139,4 +167,5 @@ export const notificationRouter = t.router({
139167
markAsRead: markAsReadProcedure,
140168
markAllAsRead: markAllAsReadProcedure,
141169
findMany: findNotificationsProcedure,
170+
onNewNotification: onNewNotificationProcedure,
142171
})

apps/rpc/src/modules/notification/notification-service.ts

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
import type EventEmitter from "node:events"
12
import type { DBHandle } from "@dotkomonline/db"
23
import type { UserId } from "@dotkomonline/types"
34
import type { NotificationRepository } from "./notification-repository"
@@ -57,7 +58,8 @@ export interface NotificationService {
5758
export function getNotificationService(
5859
notificationRepository: NotificationRepository,
5960
userRepository: UserRepository,
60-
attendanceRepository: AttendanceRepository
61+
attendanceRepository: AttendanceRepository,
62+
eventEmitter: EventEmitter
6163
): NotificationService {
6264
return {
6365
async findById(handle, notificationId) {
@@ -69,7 +71,7 @@ export function getNotificationService(
6971
},
7072

7173
async create(handle, recipientIds, notificationType, title, shortDescription, actorGroupId, payloadType, payload) {
72-
return await notificationRepository.createWithRecipients(handle, {
74+
const notification = await notificationRepository.createWithRecipients(handle, {
7375
title,
7476
shortDescription,
7577
content: shortDescription ?? title,
@@ -80,6 +82,10 @@ export function getNotificationService(
8082
taskId: null,
8183
recipientIds,
8284
})
85+
for (const userId of recipientIds) {
86+
eventEmitter.emit("notification:new", { userId, notification })
87+
}
88+
return notification
8389
},
8490

8591
async retrieveIntendedRecipientIds(handle, notificationType, eventId) {

apps/web/src/components/Navbar/Notifications/NotificationDropdown.tsx

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import { cn } from "@dotkomonline/ui"
66
import { useTRPC } from "@/utils/trpc/client"
77
import { useInfiniteQuery, useMutation, useQueryClient } from "@tanstack/react-query"
88
import type { UserNotificationDTO } from "@dotkomonline/rpc"
9+
import { useSubscription } from "@trpc/tanstack-react-query"
910

1011
interface NotificationDropdownProps extends ComponentProps<typeof DropdownMenu.Root> {
1112
open?: boolean
@@ -16,6 +17,17 @@ export const NotificationDropdown = ({ open, amountUnread, ...props }: Notificat
1617
const trpc = useTRPC()
1718
const queryClient = useQueryClient()
1819

20+
useSubscription(
21+
trpc.notification.onNewNotification.subscriptionOptions(undefined, {
22+
onData: () => {
23+
queryClient.invalidateQueries(trpc.notification.getUnreadCount.queryOptions())
24+
queryClient.invalidateQueries({
25+
queryKey: trpc.notification.getMyNotifications.infiniteQueryOptions({ take: 10 }).queryKey,
26+
})
27+
},
28+
})
29+
)
30+
1931
const { data, fetchNextPage, hasNextPage, isFetchingNextPage, isLoading } = useInfiniteQuery({
2032
...trpc.notification.getMyNotifications.infiniteQueryOptions({ take: 10 }),
2133
enabled: open,
@@ -109,7 +121,7 @@ export const NotificationDropdown = ({ open, amountUnread, ...props }: Notificat
109121
"mx-4 xs:ml-4 xs:w-102 xs:-mr-16 lg:-mr-4",
110122
"rounded-3xl shadow-sm",
111123
"bg-blue-50 border border-blue-100",
112-
"dark:border-white/10 dark:bg-stone-700"
124+
"dark:border-white/10 dark:bg-stone-800"
113125
)}
114126
>
115127
<div className="p-5 border-b border-black/10 dark:border-white/10">

apps/web/src/utils/trpc/QueryProvider.tsx

Lines changed: 27 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
"use client"
22

33
import { env } from "@/env"
4-
import { getAccessToken } from "@auth0/nextjs-auth0"
4+
import { getAccessToken, useUser } from "@auth0/nextjs-auth0"
55
import type { AppRouter } from "@dotkomonline/rpc"
66
import { QueryClient, QueryClientProvider } from "@tanstack/react-query"
77
import {
@@ -45,11 +45,16 @@ export const useTRPCSSERegisterChangeConnectionState = () => {
4545
}
4646

4747
export const QueryProvider = ({ children }: PropsWithChildren) => {
48+
const { user } = useUser()
49+
const userId = user?.sub
50+
4851
const [trpcSSERegisterChangeConnectionState, setTRPCSSERegisterChangeConnectionState] =
4952
useState<TRPCSSEConnectionState>("connecting")
5053

51-
const trpcConfig: CreateTRPCClientOptions<AppRouter> = useMemo(
52-
() => ({
54+
const trpcConfig: CreateTRPCClientOptions<AppRouter> = useMemo(() => {
55+
const hasAuthenticatedUser = typeof userId === "string" && userId !== ""
56+
57+
return {
5358
links: [
5459
loggerLink({
5560
enabled: (opts) => opts.direction === "down" && opts.result instanceof Error,
@@ -59,6 +64,23 @@ export const QueryProvider = ({ children }: PropsWithChildren) => {
5964
true: httpSubscriptionLink({
6065
transformer: superjson,
6166
url: `${env.NEXT_PUBLIC_RPC_HOST}/api/trpc`,
67+
async connectionParams() {
68+
if (!hasAuthenticatedUser) {
69+
return {}
70+
}
71+
72+
try {
73+
const token = await getAccessToken()
74+
75+
if (typeof token === "string" && token !== "") {
76+
return { token }
77+
}
78+
} catch {
79+
// not authenticated
80+
}
81+
82+
return {}
83+
},
6284
}),
6385
false: httpBatchLink({
6486
transformer: superjson,
@@ -85,9 +107,8 @@ export const QueryProvider = ({ children }: PropsWithChildren) => {
85107
}),
86108
}),
87109
],
88-
}),
89-
[]
90-
)
110+
}
111+
}, [userId])
91112

92113
const trpcClient = useMemo(() => createTRPCClient(trpcConfig), [trpcConfig])
93114

packages/db/prisma/migrations/20260325214642_fix_notification_recipient_table_name/migration.sql

Lines changed: 0 additions & 30 deletions
This file was deleted.

0 commit comments

Comments
 (0)