summaryrefslogtreecommitdiff
path: root/src/server/api/stream/index.ts
diff options
context:
space:
mode:
Diffstat (limited to 'src/server/api/stream/index.ts')
-rw-r--r--src/server/api/stream/index.ts35
1 files changed, 17 insertions, 18 deletions
diff --git a/src/server/api/stream/index.ts b/src/server/api/stream/index.ts
index ccd555e149..da4ea5ec99 100644
--- a/src/server/api/stream/index.ts
+++ b/src/server/api/stream/index.ts
@@ -14,6 +14,7 @@ import { AccessToken } from '@/models/entities/access-token';
import { UserProfile } from '@/models/entities/user-profile';
import { publishChannelStream, publishGroupMessagingStream, publishMessagingStream } from '@/services/stream';
import { UserGroup } from '@/models/entities/user-group';
+import { StreamEventEmitter, StreamMessages } from './types';
import { Packed } from '@/misc/schema';
/**
@@ -28,7 +29,7 @@ export default class Connection {
public followingChannels: Set<ChannelModel['id']> = new Set();
public token?: AccessToken;
private wsConnection: websocket.connection;
- public subscriber: EventEmitter;
+ public subscriber: StreamEventEmitter;
private channels: Channel[] = [];
private subscribingNotes: any = {};
private cachedNotes: Packed<'Note'>[] = [];
@@ -46,8 +47,8 @@ export default class Connection {
this.wsConnection.on('message', this.onWsConnectionMessage);
- this.subscriber.on('broadcast', async ({ type, body }) => {
- this.onBroadcastMessage(type, body);
+ this.subscriber.on('broadcast', data => {
+ this.onBroadcastMessage(data);
});
if (this.user) {
@@ -57,43 +58,41 @@ export default class Connection {
this.updateFollowingChannels();
this.updateUserProfile();
- this.subscriber.on(`user:${this.user.id}`, ({ type, body }) => {
- this.onUserEvent(type, body);
- });
+ this.subscriber.on(`user:${this.user.id}`, this.onUserEvent);
}
}
@autobind
- private onUserEvent(type: string, body: any) {
- switch (type) {
+ private onUserEvent(data: StreamMessages['user']['payload']) { // { type, body }と展開するとそれぞれ型が分離してしまう
+ switch (data.type) {
case 'follow':
- this.following.add(body.id);
+ this.following.add(data.body.id);
break;
case 'unfollow':
- this.following.delete(body.id);
+ this.following.delete(data.body.id);
break;
case 'mute':
- this.muting.add(body.id);
+ this.muting.add(data.body.id);
break;
case 'unmute':
- this.muting.delete(body.id);
+ this.muting.delete(data.body.id);
break;
// TODO: block events
case 'followChannel':
- this.followingChannels.add(body.id);
+ this.followingChannels.add(data.body.id);
break;
case 'unfollowChannel':
- this.followingChannels.delete(body.id);
+ this.followingChannels.delete(data.body.id);
break;
case 'updateUserProfile':
- this.userProfile = body;
+ this.userProfile = data.body;
break;
case 'terminate':
@@ -145,8 +144,8 @@ export default class Connection {
}
@autobind
- private onBroadcastMessage(type: string, body: any) {
- this.sendMessageToWs(type, body);
+ private onBroadcastMessage(data: StreamMessages['broadcast']['payload']) {
+ this.sendMessageToWs(data.type, data.body);
}
@autobind
@@ -249,7 +248,7 @@ export default class Connection {
}
@autobind
- private async onNoteStreamMessage(data: any) {
+ private async onNoteStreamMessage(data: StreamMessages['note']['payload']) {
this.sendMessageToWs('noteUpdated', {
id: data.body.id,
type: data.type,