Revisited, September 2026. This was one of the first articles in the series, written while KickJS was in its early phases and Vibed (the Jira-like backend this series started with) was the test bed. The ideas held up — rooms, a stateless typing relay, presence in memory — but most of the code did not. It is updated here for
@forinda/kickjs8.4 and@forinda/kickjs-ws7.1. What changed is listed at the end for anyone upgrading code from the original. For the full realtime story, including brokers and multi-server setups, read the companion piece Realtime in KickJS: WebSockets, rooms, auth and scaling past one server.
TL;DR
-
@forinda/kickjs-wsruns the same decorator-driven controllers over plain WebSockets (WsAdapter) or Socket.IO (SocketIoAdapter). You pick the transport in one line. - Authenticate at the handshake with
auth.resolveUser. A socket that fails is refused before any handler runs. - Rooms (
ctx.join(),ctx.to().send()) are still the right primitive for channels. Join only after checking that the user belongs to the channel. -
ctx.to(room).send()includes the sender. The typing indicator filters out the user's own events on the client. - Presence is counted per user, not per socket, in an injectable service. The adapter's heartbeat removes dead sockets, so the cleanup cron is gone.
- To scale out, use Socket.IO's own adapter, or give
WsAdapterabroker. Presence stays per instance until you share it.
Choosing a Transport
Vibed has real-time chat built into its task management backend. Users join workspace channels and exchange messages. The socket handles delivery, typing indicators and presence. REST handles history, edits and deletes.
The first version of this article said KickJS "wraps Socket.IO". That was never quite true, and today it is a choice. The package ships two adapters behind one controller API:
// src/index.ts
import { bootstrap, getEnv } from '@forinda/kickjs'
import { SocketIoAdapter } from '@forinda/kickjs-ws/socket.io'
import { resolveSocketUser } from '@/realtime/socket-user'
import { modules } from './modules'
export const app = await bootstrap({
modules,
adapters: [
SocketIoAdapter({
path: '/ws',
cors: { origin: getEnv('WEB_ORIGIN'), credentials: true },
auth: { resolveUser: resolveSocketUser },
}),
],
})
Or, over plain WebSockets:
import { WsAdapter } from '@forinda/kickjs-ws'
WsAdapter({
path: '/ws',
heartbeatInterval: 30_000,
maxPayload: 1_048_576, // 1 MiB
auth: { resolveUser: resolveSocketUser },
})
Adapters are factories now: SocketIoAdapter({...}), not new WsAdapter({...}). (Adapters vs Plugins in KickJS explains why, and the adapters guide covers the lifecycle hooks.) SocketIoAdapter accepts any Socket.IO ServerOptions (cors, path, adapter, pingInterval…). Install socket.io alongside it, because it is an optional peer dependency.
WsAdapter (ws) |
SocketIoAdapter |
|
|---|---|---|
| Client | the browser's WebSocket
|
socket.io-client |
| Wire format |
{ "event": "…", "data": … } JSON frames |
Socket.IO packets |
| Rooms | shared across namespaces | scoped to a namespace |
| Reconnect, fallback to long-polling | yours to write | built in |
| Multi-instance |
broker option |
Socket.IO adapter (Redis, Postgres…) |
Reference docs for each: the WebSockets guide, the Socket.IO guide and the @forinda/kickjs-ws API.
The rest of this article uses Socket.IO, as the title says. Every controller below runs unchanged on WsAdapter. Where the two behave differently, the text says so.
Authenticating the Handshake
The original controller verified a JWT inside @OnConnect. When the token was bad, it sent an error event and left the socket open. An unauthenticated socket could then send channel:join and every other event, so each handler had to check ctx.get('userId') again.
auth.resolveUser moves that check to the door. (Where the token comes from is covered in Building JWT Auth with Refresh Token Rotation in KickJS and the authentication guide.) It runs once per socket, before any @OnConnect handler:
// src/realtime/socket-user.ts
import type { IncomingMessage } from 'node:http'
import { getEnv } from '@forinda/kickjs'
import type { WsAuthenticatedUser } from '@forinda/kickjs-ws'
import jwt from 'jsonwebtoken'
export async function resolveSocketUser(
request: IncomingMessage,
handshakeAuth?: Record<string, unknown>,
): Promise<WsAuthenticatedUser | null> {
// Socket.IO clients send `auth: { token }`. A browser WebSocket cannot set
// headers, so on `WsAdapter` read it from a cookie or subprotocol instead.
const token = typeof handshakeAuth?.token === 'string' ? handshakeAuth.token : null
if (!token) return null
try {
const payload = jwt.verify(token, getEnv('JWT_SECRET')) as { sub: string; name: string }
return { id: payload.sub, name: payload.name }
} catch {
return null
}
}
If it returns null or throws, Socket.IO clients get a connect_error with the message Unauthorized, and WsAdapter closes the socket with code 4401. If it succeeds, the user is already on the context as ctx.get('user') and ctx.get('userId'). The socket also joins user:<id>, which is what lets services push notifications to one person.
Messages the client sends while this runs are held and replayed once @OnConnect finishes, up to 64 messages or 1 MiB. A client that emits channel:join right after connecting no longer loses it to a race.
On the client:
import { io } from 'socket.io-client'
const socket = io('http://localhost:3000/chat', {
path: '/ws',
auth: (cb) => cb({ token: getAccessToken() }), // a function, so reconnects pick up a refreshed token
})
socket.on('connect_error', async (err) => {
if (err.message === 'Unauthorized') {
await refreshAccessToken()
socket.connect()
}
})
The WebSocket Controller
The decorators are the same as before (@WsController, @OnConnect, @OnDisconnect, @OnMessage), with @OnError added. Everything else in the core now imports from @forinda/kickjs, and logging goes through createLogger:
// src/modules/messages/chat.ws-controller.ts
import { Autowired, createLogger } from '@forinda/kickjs'
import { OnConnect, OnDisconnect, OnMessage, WsController } from '@forinda/kickjs-ws'
import type { WsContext } from '@forinda/kickjs-ws'
import { ChannelMembers } from '@/modules/channels/channel-members.service'
import { MessageRepository } from './message.repository'
import { PresenceService } from './presence.service'
const log = createLogger('ChatWsController')
const channelRoom = (channelId: string) => `channel:${channelId}`
@WsController('/chat')
export class ChatWsController {
@Autowired() private readonly messages!: MessageRepository
@Autowired() private readonly members!: ChannelMembers
@Autowired() private readonly presence!: PresenceService
@OnConnect()
connect(ctx: WsContext) {
const user = ctx.get<{ id: string; name: string }>('user')!
if (this.presence.connect(user.id, ctx.id)) {
ctx.broadcast('presence:online', { userId: user.id, userName: user.name })
}
ctx.send('welcome', { id: ctx.id, userId: user.id, online: this.presence.online() })
log.info(`${user.id} connected (${ctx.id})`)
}
@OnDisconnect()
disconnect(ctx: WsContext) {
const userId = ctx.get<string>('userId')
if (userId && this.presence.disconnect(userId, ctx.id)) {
ctx.broadcast('presence:offline', { userId })
}
}
// … join, send and typing handlers below
}
Three things changed from the original:
-
ctx.get('user')is always set. The handshake guarantees it, so handlers no longer need an "if not authenticated" branch. -
Presence lives in a service, not in a module-level
Mapexported from the controller file. Anything can inject it, and tests can build a fresh one. (See Dependency Injection in KickJS, and the@Service()gotchas if a class doesn't resolve.) -
ctx.broadcastinstead ofctx.broadcastAll. The user who just connected gets the full list inwelcomeand does not need their ownpresence:onlineevent.
The controller type is WsContext. Under SocketIoAdapter, handlers receive a SocketIoContext with the same methods. If you only ever run on Socket.IO, you can type the parameter as SocketIoContext and use ctx.socket for anything Socket.IO-specific.
@WsController('/chat') maps to the /chat namespace: io.of('/chat') on Socket.IO, and ws://host/ws/chat on WsAdapter.
Room-Based Broadcasting
Rooms are still the heart of it. A channel is a group, and a room is a group. The original join handler let any authenticated socket join any channel ID it named, including private channels it was not in. The fix is one check:
@OnMessage('channel:join')
async join(ctx: WsContext) {
const userId = ctx.get<string>('userId')!
const channelId = ctx.data?.channelId
if (typeof channelId !== 'string') return
if (!(await this.members.isMember(channelId, userId))) {
// Same answer as for a channel that does not exist — don't confirm it's there.
ctx.send('channel:join_denied', { channelId })
return
}
await ctx.join(channelRoom(channelId))
ctx.to(channelRoom(channelId)).send('channel:user_joined', { channelId, userId })
}
@OnMessage('channel:leave')
async leave(ctx: WsContext) {
const channelId = ctx.data?.channelId
if (typeof channelId !== 'string') return
await ctx.leave(channelRoom(channelId))
ctx.to(channelRoom(channelId)).send('channel:user_left', {
channelId,
userId: ctx.get('userId'),
})
}
join and leave are awaited. On Socket.IO they return a promise when the room adapter is asynchronous, as with Redis. On ws they return nothing, and awaiting that is harmless.
The methods:
-
ctx.join(room)/ctx.leave(room)add this socket to a room or remove it. A socket can be in many rooms. -
ctx.to(room).send(event, data)sends to every socket in the room, including this one. -
ctx.broadcast(event, data)sends to every socket in the namespace except this one.ctx.broadcastAllincludes it.
The first version said to().send() excluded the sender. It does not, on either adapter. That changes two handlers below.
Room names on ws are global. Under WsAdapter, a room called lobby joined from /chat is the same room as lobby joined from /admin. Prefix your room names (channel:, project:, user:) and you will not be surprised. Socket.IO rooms stay inside their namespace.
Sending Messages
@OnMessage('message:send')
async send(ctx: WsContext) {
const user = ctx.get<{ id: string; name: string }>('user')!
const { channelId, content, clientId } = ctx.data ?? {}
if (typeof channelId !== 'string' || typeof content !== 'string' || !content.trim()) return
if (!ctx.rooms().includes(channelRoom(channelId))) return // must have joined first
const message = await this.messages.create({ channelId, senderId: user.id, content })
ctx.to(channelRoom(channelId)).send('message:new', {
messageId: message.id,
clientId, // lets the sender match this to its optimistic copy
channelId,
senderId: user.id,
senderName: user.name,
content: message.content,
createdAt: message.createdAt,
})
}
The separate "echo back to sender" call is gone, because the room send already reaches the sender. The clientId the client generated comes back in the payload. An optimistic UI swaps its pending bubble for the saved one by matching on it, instead of guessing by content.
Checking ctx.rooms() before writing costs nothing and closes the same hole as the join check. A socket can post only to channels it was allowed to join.
Typing Indicators
Typing indicators need to be fast, cheap (no database writes) and scoped to the channel. The server is still a pure relay and keeps no typing state:
@OnMessage('channel:typing')
typing(ctx: WsContext) {
const channelId = ctx.data?.channelId
if (typeof channelId !== 'string' || !ctx.rooms().includes(channelRoom(channelId))) return
const user = ctx.get<{ id: string; name: string }>('user')!
ctx.to(channelRoom(channelId)).send('channel:typing', {
channelId,
userId: user.id,
userName: user.name,
})
}
@OnMessage('channel:stop_typing')
stopTyping(ctx: WsContext) {
const channelId = ctx.data?.channelId
if (typeof channelId !== 'string' || !ctx.rooms().includes(channelRoom(channelId))) return
ctx.to(channelRoom(channelId)).send('channel:stop_typing', {
channelId,
userId: ctx.get('userId'),
})
}
Client-Side Implementation
The first version had two gaps. Because the room send includes the sender, your own events come back and "Alice is typing…" would show on Alice's screen, including in her other tabs. And when a tab closed mid-sentence, stop_typing never arrived, so the indicator stayed forever.
Both are fixed on the client:
const TYPING_IDLE_MS = 2_000 // stop after 2s without a keystroke
const TYPING_REPEAT_MS = 3_000 // while typing, say so again every 3s
const TYPING_EXPIRES_MS = 5_000 // forget anyone we haven't heard from in 5s
let idleTimer: ReturnType<typeof setTimeout> | undefined
let lastSent = 0
export function onComposerInput(channelId: string) {
const now = Date.now()
if (now - lastSent > TYPING_REPEAT_MS) {
socket.emit('channel:typing', { channelId })
lastSent = now
}
clearTimeout(idleTimer)
idleTimer = setTimeout(() => {
socket.emit('channel:stop_typing', { channelId })
lastSent = 0
}, TYPING_IDLE_MS)
}
// Receiving
const typing = new Map<string, { name: string; expires: ReturnType<typeof setTimeout> }>()
socket.on('channel:typing', ({ userId, userName }) => {
if (userId === me.id) return // our own event, echoed by the room
clearTimeout(typing.get(userId)?.expires)
typing.set(userId, {
name: userName,
expires: setTimeout(() => forget(userId), TYPING_EXPIRES_MS),
})
render()
})
socket.on('channel:stop_typing', ({ userId }) => forget(userId))
function forget(userId: string) {
clearTimeout(typing.get(userId)?.expires)
if (typing.delete(userId)) render()
}
function render() {
const names = [...typing.values()].map((t) => t.name)
typingLabel.textContent =
names.length === 0 ? ''
: names.length === 1 ? `${names[0]} is typing…`
: names.length === 2 ? `${names[0]} and ${names[1]} are typing…`
: `${names[0]} and ${names.length - 1} others are typing…`
}
The sender repeats typing every 3 seconds while keys keep coming. The receiver forgets anyone it has not heard from in 5 seconds. A lost stop_typing now clears itself, and the server still keeps no state.
The filter uses the user ID, not the socket ID. Your own typing should not show in your other tabs either.
Why Rooms Still Make This Trivial
Without rooms, the typing handler would look up the channel's members, find their sockets, loop over them, and deal with users who have several tabs open. With rooms, it is one line. The adapter handles the fan-out, sends to every tab, and removes the socket from all its rooms when it disconnects.
Presence Tracking
Presence answers "who's online right now?". The original used a Map<socketId, user> exported from the controller file, and on every disconnect it scanned the whole map to see whether the user still had another tab open. Counting sockets per user makes that check constant-time and puts it where it can be tested:
// src/modules/messages/presence.service.ts
import { Service } from '@forinda/kickjs'
@Service()
export class PresenceService {
/** userId → the sockets that user has open on this instance */
private readonly sockets = new Map<string, Set<string>>()
/** Returns true when this is the user's first socket — they just came online. */
connect(userId: string, socketId: string): boolean {
const open = this.sockets.get(userId) ?? new Set<string>()
open.add(socketId)
this.sockets.set(userId, open)
return open.size === 1
}
/** Returns true when that was the user's last socket — they just went offline. */
disconnect(userId: string, socketId: string): boolean {
const open = this.sockets.get(userId)
if (!open?.delete(socketId)) return false
if (open.size > 0) return false
this.sockets.delete(userId)
return true
}
online(): string[] {
return [...this.sockets.keys()]
}
count(): number {
return this.sockets.size
}
}
Closing one of three tabs changes nothing. Closing the last one sends presence:offline.
Sharing Presence with Other Endpoints
Other code injects the service instead of importing a function from a controller file. The live stats endpoint pushes the online count over SSE (the SSE guide has the details, and WebSocket Chat + SSE Stats Streams in KickJS is the original walkthrough):
@Get('/:workspaceId/stats/live')
async live(ctx: RequestContext) {
const sse = ctx.sse<{ onlineUsers: number; timestamp: string }>()
const push = () =>
sse.send({ onlineUsers: this.presence.count(), timestamp: new Date().toISOString() }, 'stats:update')
push()
const interval = setInterval(push, 10_000)
sse.onClose(() => clearInterval(interval))
}
Workspace membership used to be checked with @Middleware(workspaceMembershipGuard). That check is now a context contributor applied at the module, so the same guard runs for HTTP, sockets and jobs. See Building Custom Context Decorators in KickJS and the context decorators guide.
Why the Cleanup Cron Is Gone
The first version registered a @Cron('*/5 * * * *') job to "clean up stale presence entries". The job body was a stub. It never ran against real data, and it turns out it was not needed:
-
Sockets that drop without saying goodbye (network loss, a killed browser) are what the heartbeat is for.
WsAdapterpings everyheartbeatIntervaland terminates sockets that don't answer. Socket.IO does the same withpingInterval/pingTimeout. Either way@OnDisconnectfires, and the service removes the socket. - A crashed server loses the map, but it also loses every socket in it. Clients reconnect to a live instance and connect again.
Stale entries cannot build up in memory, so a sweep has nothing to sweep. A timer is only worth adding for presence you store outside the process, which is the next section. When you do need one, Cron jobs in KickJS shows how to run it safely across restarts and several servers.
"Connected" and "actually here" are not the same thing, though. A tab left open in the background all weekend is connected. If that matters to you, have visible tabs send a small active event every minute or so. Record the time in the service and treat anyone silent for twice that long as away. That is a timestamp per user and a filter in online(), not a job.
Why Rooms Beat Individual Socket Tracking
Without rooms you maintain Map<channelId, Set<socketId>> yourself, loop over it to send, and prune sockets that vanished without leaving. Every edge case (disconnect without leave, multi-tab, crashes) is yours to handle.
With rooms:
await ctx.join(`channel:${channelId}`)
ctx.to(`channel:${channelId}`).send('channel:typing', { userId, userName })
await ctx.leave(`channel:${channelId}`)
// On disconnect the adapter removes the socket from every room.
Rooms also cover the case the original article handled by hand: per-user delivery. The handshake already joined each socket to user:<id>, so a service can reach every tab of one person without knowing about sockets:
import { Autowired, Service } from '@forinda/kickjs'
import { WS_USER_BROADCASTER, type WsUserBroadcaster } from '@forinda/kickjs-ws'
@Service()
export class MentionNotifier {
@Autowired(WS_USER_BROADCASTER) private readonly users!: WsUserBroadcaster
notify(userId: string, mention: { channelId: string; messageId: string }) {
this.users.toUser(userId).send('mention:new', mention)
}
}
Both adapters register WS_USER_BROADCASTER. For rooms other than user rooms, use WS_ROOM_MANAGER on WsAdapter or the SOCKET_IO server on Socket.IO.
The Complete Event Flow
Client connects with auth: { token }
→ resolveUser verifies it; failure = connect_error "Unauthorized"
→ socket joins user:A
→ @OnConnect: first socket for A → broadcast presence:online; send welcome { online }
User A opens #general
→ emit channel:join { channelId: 'general' }
→ server checks A is a member → join channel:general
→ room receives channel:user_joined { userId: 'A' }
User A types
→ emit channel:typing (at most every 3s)
→ room receives channel:typing { userId: 'A', userName: 'alice' }
→ A's own tabs ignore it; B shows "alice is typing…", expires after 5s of silence
User A sends
→ emit message:send { channelId, content: 'Hello!', clientId }
→ server checks A is in the room, persists, sends message:new to the room
→ B renders it; A replaces its optimistic bubble matched by clientId
User A stops typing (2s idle)
→ emit channel:stop_typing → B removes the indicator
User A closes the last tab
→ heartbeat or clean close → @OnDisconnect
→ last socket for A → broadcast presence:offline
REST still owns everything with history:
GET /channels/:channelId/messages → paginated history
PATCH /messages/:messageId → edit (author only)
DELETE /messages/:messageId → soft delete (author only)
Scaling Out
One instance is fine for a long time. With two or more, a room send on instance 1 has to reach sockets connected to instance 2. The realtime tutorial walks through both setups below end to end.
On Socket.IO, use a Socket.IO adapter. SocketIoAdapter accepts all ServerOptions, so it plugs straight in:
import { getEnv } from '@forinda/kickjs'
import { createAdapter } from '@socket.io/redis-adapter'
import { Redis } from 'ioredis'
const pub = new Redis(getEnv('REDIS_URL'))
SocketIoAdapter({
path: '/ws',
adapter: createAdapter(pub, pub.duplicate()),
auth: { resolveUser: resolveSocketUser },
})
If you run several instances behind a load balancer, keep sticky sessions on or restrict clients to transports: ['websocket']. Socket.IO's long-polling handshake has to land on the same instance.
On ws, give the adapter a broker. Room sends, namespace broadcasts and per-user sends are then relayed to every instance:
import { getEnv } from '@forinda/kickjs'
import { WsAdapter } from '@forinda/kickjs-ws'
import { redisBroker } from '@forinda/kickjs-ws/redis'
import { Redis } from 'ioredis'
const pub = new Redis(getEnv('REDIS_URL'))
WsAdapter({
path: '/ws',
broker: redisBroker({ publisher: pub, subscriber: pub.duplicate() }),
auth: { resolveUser: resolveSocketUser },
})
WsBroker is a small interface, so any pub/sub works. Postgres LISTEN/NOTIFY is enough if you'd rather not run Redis at all.
Presence does not scale by itself. Neither option shares PresenceService: each instance only knows its own sockets. When a user's two tabs land on different instances, closing one must not mark them offline. The next step is one of these:
- Keep the per-user socket counts in Redis (
HINCRBY presence:<userId> <instanceId> ±1) with a TTL that each instance refreshes. This is where a periodic job earns its place, clearing the counts of an instance that died. - Or have instances gossip their own "who I have" lists over the same broker and expire any list not refreshed within a couple of intervals.
Key Takeaways
-
Authenticate at the handshake.
auth.resolveUserrefuses a bad token before any handler runs. Don't verify inside@OnConnectand leave the socket open. -
Authorize every join. An authenticated socket is not a channel member. Check before
ctx.join, and check room membership before relaying. - Rooms include the sender. Filter your own typing events on the client and drop the manual echo on send.
-
Keep typing stateless on the server, and self-expiring on the client. Repeat while typing and expire on silence, so a lost
stop_typingcannot stick. - Count presence per user, in a service. First socket means online, last socket means offline. Inject it anywhere, including SSE endpoints.
-
Let the heartbeat clean up. Dead sockets fire
@OnDisconnect. A cleanup job belongs to shared presence, not in-memory presence. -
Pick the transport, keep the controllers. Socket.IO for its client, reconnects and ecosystem.
wsfor a lighter wire and the browser's nativeWebSocket. The handlers do not change.
What Changed Since the First Version
| First version | Now |
|---|---|
@forinda/kickjs-core |
@forinda/kickjs |
Logger.for('Name') |
createLogger('Name') |
new WsAdapter({...}) in config/adapters.ts
|
WsAdapter({...}) / SocketIoAdapter({...}) factories passed to bootstrap({ adapters })
|
| "KickJS wraps Socket.IO" |
ws by default, Socket.IO via @forinda/kickjs-ws/socket.io
|
JWT verified in @OnConnect, socket left open on failure |
auth.resolveUser: refused with connect_error or close code 4401
|
ctx.to(room).send() described as excluding the sender |
It includes the sender, on both adapters |
| Explicit echo to sender after a room send | Removed, and a clientId is returned for optimistic UI |
| Any channel ID could be joined | Membership checked before ctx.join
|
Exported module-level onlineUsers map with an O(n) multi-tab scan |
Injectable PresenceService with per-user socket sets |
@Cron presence cleanup (stub) |
Removed; the heartbeat fires @OnDisconnect
|
| Hand-rolled per-user rooms |
user:<id> joined automatically; WS_USER_BROADCASTER
|
@Middleware(workspaceMembershipGuard) |
Context contributor applied at the module |
| Redis adapter only | Socket.IO adapter, or broker (redisBroker or your own) on ws
|
Further Reading
From the same series:
- Building a Jira-like Backend with Decorator-Driven DDD in Node.js: where Vibed started
- Building a Jira-Like Task API with KickJS: Auth, WebSocket Chat, SSE, Queues & More
- Building a Complete Jira-like Task Management Backend with KickJS: From Scaffold to Production
- Realtime in KickJS: WebSockets, rooms, auth and scaling past one server
- The KickJS Request Lifecycle: Middleware, Contributors, and Typed Context
- From v3 to v5: Migrating a Production KickJS App, Slice by Slice
- I migrated a 593-route Express app to KickJS twice — here's what held up
KickJS:
- Docs: kickjs.app. Start with Getting started.
- Source: github.com/forinda/kick-js. Issues and stars welcome.
- npm:
@forinda/kickjs,@forinda/kickjs-ws
This is part of a series on building a Jira-like backend with KickJS.
Top comments (0)