Skip to content

Commit dbc312e

Browse files
committed
realtime updates using supabase
1 parent 44dd412 commit dbc312e

10 files changed

Lines changed: 928 additions & 4 deletions

File tree

‎apps/codebility/app/home/kanban/[projectId]/[id]/_components/KanbanBoardColumnContainer.tsx‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,12 +20,15 @@ import {
2020
} from "@dnd-kit/sortable";
2121
import toast from "react-hot-toast";
2222

23+
import { broadcastColumnsReorder } from "@/lib/kanban/board-broadcast";
2324
import {
2425
commitTaskMove,
2526
filterColumnsByMember,
2627
} from "@/lib/kanban/board-mutations";
28+
import { registerColumnReorderBlock } from "@/lib/kanban/own-writes";
2729
import {
2830
useKanbanBoardActions,
31+
useKanbanBoardMeta,
2932
useKanbanColumns,
3033
} from "@/store/kanban-board/KanbanBoardProvider";
3134

@@ -71,6 +74,7 @@ export default function KanbanBoardColumnContainer({
7174
}: Props) {
7275
const router = useRouter();
7376
const columns = useKanbanColumns();
77+
const { boardId } = useKanbanBoardMeta();
7478
const { setColumns, removeTaskLocal } = useKanbanBoardActions();
7579

7680
const boardData = useMemo(
@@ -100,7 +104,14 @@ export default function KanbanBoardColumnContainer({
100104
const newCols = arrayMove(orderedColumns, oldIndex, newIndex).map(
101105
(column, index) => ({ ...column, position: index }),
102106
);
107+
registerColumnReorderBlock();
103108
setColumns(newCols);
109+
broadcastColumnsReorder(boardId, {
110+
columns: newCols.map((column) => ({
111+
id: column.id,
112+
position: column.position ?? 0,
113+
})),
114+
});
104115

105116
try {
106117
await Promise.all(
@@ -160,7 +171,7 @@ export default function KanbanBoardColumnContainer({
160171
}
161172
}
162173
},
163-
[boardData, columns, router, setColumns],
174+
[boardData, boardId, columns, router, setColumns],
164175
);
165176

166177
const { sensors } = useDragAndDrop({

‎apps/codebility/hooks/use-kanban-board-sync.ts‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,12 @@
22

33
import { useRouter } from "next/navigation";
44

5+
import {
6+
broadcastTaskPatch,
7+
broadcastTaskRemove,
8+
} from "@/lib/kanban/board-broadcast";
59
import { clearTaskDetailCache } from "@/lib/kanban/task-detail-cache";
10+
import { registerOwnTaskWrite } from "@/lib/kanban/own-writes";
611
import { Task } from "@/types/home/codev";
712
import { tryGetKanbanBoardStore } from "@/store/kanban-board/registry";
813

@@ -15,11 +20,25 @@ export function useKanbanBoardSync() {
1520
};
1621

1722
const removeTask = (taskId: string) => {
18-
tryGetKanbanBoardStore()?.getState().removeTaskLocal(taskId);
23+
const store = tryGetKanbanBoardStore();
24+
if (!store) {
25+
return;
26+
}
27+
28+
registerOwnTaskWrite(taskId);
29+
store.getState().removeTaskLocal(taskId);
30+
broadcastTaskRemove(store.getState().boardId, taskId);
1931
};
2032

2133
const patchTask = (taskId: string, patch: Partial<Task>) => {
22-
tryGetKanbanBoardStore()?.getState().patchTaskLocal(taskId, patch);
34+
const store = tryGetKanbanBoardStore();
35+
if (!store) {
36+
return;
37+
}
38+
39+
registerOwnTaskWrite(taskId);
40+
store.getState().patchTaskLocal(taskId, patch);
41+
broadcastTaskPatch(store.getState().boardId, { taskId, patch });
2342
};
2443

2544
return { refreshBoard, removeTask, patchTask };
Lines changed: 309 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,309 @@
1+
"use client";
2+
3+
import { useEffect, useRef } from "react";
4+
import type {
5+
RealtimeChannel,
6+
RealtimePostgresChangesPayload,
7+
} from "@supabase/supabase-js";
8+
9+
import {
10+
registerKanbanBroadcastChannel,
11+
unregisterKanbanBroadcastChannel,
12+
} from "@/lib/kanban/board-broadcast";
13+
import {
14+
bindPostgresSyncQueue,
15+
enqueuePostgresColumnChange,
16+
enqueuePostgresTaskChange,
17+
unbindPostgresSyncQueue,
18+
} from "@/lib/kanban/postgres-board-sync";
19+
import {
20+
registerBroadcastApplied,
21+
shouldIgnoreRemoteColumnPatches,
22+
shouldIgnoreRemoteTask,
23+
} from "@/lib/kanban/own-writes";
24+
import type { KanbanBoardStoreApi } from "@/store/kanban-board/create-kanban-board-store";
25+
import { Task } from "@/types/home/codev";
26+
import { createClientClientComponent } from "@/utils/supabase/client";
27+
28+
type TaskRow = Record<string, unknown>;
29+
type ColumnRow = Record<string, unknown>;
30+
31+
type TaskMovePayload = {
32+
taskId: string;
33+
columnId: string;
34+
position: number;
35+
};
36+
37+
type TaskPatchPayload = {
38+
taskId: string;
39+
patch: Partial<Task>;
40+
};
41+
42+
type ColumnsReorderPayload = {
43+
columns: Array<{ id: string; position: number }>;
44+
};
45+
46+
async function syncRealtimeAuth(
47+
supabase: NonNullable<ReturnType<typeof createClientClientComponent>>,
48+
) {
49+
const {
50+
data: { session },
51+
} = await supabase.auth.getSession();
52+
53+
if (session?.access_token) {
54+
await supabase.realtime.setAuth(session.access_token);
55+
return true;
56+
}
57+
58+
return false;
59+
}
60+
61+
function applyBroadcastTaskMove(
62+
store: KanbanBoardStoreApi,
63+
payload: TaskMovePayload,
64+
) {
65+
if (shouldIgnoreRemoteTask(payload.taskId)) {
66+
return;
67+
}
68+
69+
const state = store.getState();
70+
const sourceColumn = state.columns.find((column) =>
71+
(column.tasks ?? []).some((task) => task.id === payload.taskId),
72+
);
73+
const targetColumn = state.columns.find(
74+
(column) => column.id === payload.columnId,
75+
);
76+
77+
for (const column of [sourceColumn, targetColumn]) {
78+
if (!column) {
79+
continue;
80+
}
81+
82+
for (const task of column.tasks ?? []) {
83+
registerBroadcastApplied(task.id);
84+
}
85+
}
86+
87+
store
88+
.getState()
89+
.moveTaskLocal(payload.taskId, payload.columnId, payload.position);
90+
}
91+
92+
function applyBroadcastTaskRemove(
93+
store: KanbanBoardStoreApi,
94+
payload: { taskId: string },
95+
) {
96+
if (shouldIgnoreRemoteTask(payload.taskId)) {
97+
return;
98+
}
99+
100+
registerBroadcastApplied(payload.taskId);
101+
store.getState().removeTaskLocal(payload.taskId);
102+
}
103+
104+
function applyBroadcastTaskPatch(
105+
store: KanbanBoardStoreApi,
106+
payload: TaskPatchPayload,
107+
) {
108+
if (shouldIgnoreRemoteTask(payload.taskId)) {
109+
return;
110+
}
111+
112+
registerBroadcastApplied(payload.taskId);
113+
store.getState().patchTaskLocal(payload.taskId, payload.patch);
114+
}
115+
116+
function applyBroadcastColumnsReorder(
117+
store: KanbanBoardStoreApi,
118+
payload: ColumnsReorderPayload,
119+
) {
120+
if (shouldIgnoreRemoteColumnPatches()) {
121+
return;
122+
}
123+
124+
const positions = new Map(
125+
payload.columns.map((column) => [column.id, column.position]),
126+
);
127+
128+
const nextColumns = store
129+
.getState()
130+
.columns.map((column) => {
131+
const position = positions.get(column.id);
132+
return position === undefined ? column : { ...column, position };
133+
})
134+
.sort((a, b) => (a.position ?? 0) - (b.position ?? 0));
135+
136+
store.getState().setColumns(nextColumns);
137+
}
138+
139+
type UseKanbanRealtimeOptions = {
140+
boardId: string;
141+
store: KanbanBoardStoreApi;
142+
onReconnect: () => void;
143+
};
144+
145+
export function useKanbanRealtime({
146+
boardId,
147+
store,
148+
onReconnect,
149+
}: UseKanbanRealtimeOptions) {
150+
const onReconnectRef = useRef(onReconnect);
151+
onReconnectRef.current = onReconnect;
152+
153+
useEffect(() => {
154+
const supabase = createClientClientComponent();
155+
if (!supabase) {
156+
return;
157+
}
158+
159+
let channel: RealtimeChannel | null = null;
160+
let needsReconnectRefresh = false;
161+
let disposed = false;
162+
163+
const getColumnIds = () =>
164+
new Set(store.getState().columns.map((column) => column.id));
165+
166+
bindPostgresSyncQueue(boardId, store, getColumnIds);
167+
168+
const handleChannelStatus = (status: string, err?: Error) => {
169+
if (err && process.env.NODE_ENV === "development") {
170+
console.error("[kanban-realtime] subscription error:", err);
171+
}
172+
173+
if (status === "SUBSCRIBED") {
174+
if (channel) {
175+
registerKanbanBroadcastChannel(channel, boardId);
176+
}
177+
178+
if (needsReconnectRefresh) {
179+
needsReconnectRefresh = false;
180+
onReconnectRef.current();
181+
}
182+
return;
183+
}
184+
185+
if (
186+
status === "CLOSED" ||
187+
status === "CHANNEL_ERROR" ||
188+
status === "TIMED_OUT"
189+
) {
190+
needsReconnectRefresh = true;
191+
if (channel) {
192+
unregisterKanbanBroadcastChannel(channel);
193+
}
194+
}
195+
};
196+
197+
const removeChannel = () => {
198+
if (channel) {
199+
unregisterKanbanBroadcastChannel(channel);
200+
supabase.removeChannel(channel);
201+
channel = null;
202+
}
203+
};
204+
205+
const subscribe = async () => {
206+
if (disposed) {
207+
return;
208+
}
209+
210+
const authed = await syncRealtimeAuth(supabase);
211+
if (!authed || disposed) {
212+
return;
213+
}
214+
215+
removeChannel();
216+
217+
const columnIds = getColumnIds();
218+
219+
channel = supabase.channel(`kanban-board:${boardId}`, {
220+
config: { broadcast: { self: false } },
221+
});
222+
223+
channel
224+
.on("broadcast", { event: "task_move" }, ({ payload }) => {
225+
applyBroadcastTaskMove(store, payload as TaskMovePayload);
226+
})
227+
.on("broadcast", { event: "task_remove" }, ({ payload }) => {
228+
applyBroadcastTaskRemove(store, payload as { taskId: string });
229+
})
230+
.on("broadcast", { event: "task_patch" }, ({ payload }) => {
231+
applyBroadcastTaskPatch(store, payload as TaskPatchPayload);
232+
})
233+
.on("broadcast", { event: "columns_reorder" }, ({ payload }) => {
234+
applyBroadcastColumnsReorder(store, payload as ColumnsReorderPayload);
235+
})
236+
.on(
237+
"postgres_changes",
238+
{
239+
event: "*",
240+
schema: "public",
241+
table: "kanban_columns",
242+
filter: `board_id=eq.${boardId}`,
243+
},
244+
(payload) => {
245+
enqueuePostgresColumnChange(
246+
boardId,
247+
payload as RealtimePostgresChangesPayload<ColumnRow>,
248+
);
249+
},
250+
);
251+
252+
for (const columnId of columnIds) {
253+
channel.on(
254+
"postgres_changes",
255+
{
256+
event: "*",
257+
schema: "public",
258+
table: "tasks",
259+
filter: `kanban_column_id=eq.${columnId}`,
260+
},
261+
(payload) => {
262+
enqueuePostgresTaskChange(
263+
boardId,
264+
payload as RealtimePostgresChangesPayload<TaskRow>,
265+
);
266+
},
267+
);
268+
}
269+
270+
channel.subscribe(handleChannelStatus);
271+
};
272+
273+
void subscribe();
274+
275+
let previousColumnIds = new Set(
276+
store.getState().columns.map((column) => column.id),
277+
);
278+
279+
const unsubscribeColumns = store.subscribe((state) => {
280+
const nextColumnIds = new Set(state.columns.map((column) => column.id));
281+
const structureChanged =
282+
nextColumnIds.size !== previousColumnIds.size ||
283+
[...nextColumnIds].some((id) => !previousColumnIds.has(id));
284+
285+
if (structureChanged) {
286+
previousColumnIds = nextColumnIds;
287+
void subscribe();
288+
}
289+
});
290+
291+
const {
292+
data: { subscription: authSubscription },
293+
} = supabase.auth.onAuthStateChange((_event, session) => {
294+
if (disposed || !session?.access_token) {
295+
return;
296+
}
297+
298+
void supabase.realtime.setAuth(session.access_token);
299+
});
300+
301+
return () => {
302+
disposed = true;
303+
authSubscription.unsubscribe();
304+
unsubscribeColumns();
305+
removeChannel();
306+
unbindPostgresSyncQueue(boardId);
307+
};
308+
}, [boardId, store]);
309+
}

0 commit comments

Comments
 (0)