Spaces:
Sleeping
Sleeping
Sasha commited on
Commit ·
0252fea
1
Parent(s): adbcce3
feat: add automatic self-healing for empty stream viewers
Browse files- server/db.js +130 -2
server/db.js
CHANGED
|
@@ -1091,6 +1091,102 @@ export async function getStreamsList() {
|
|
| 1091 |
}
|
| 1092 |
}
|
| 1093 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1094 |
/**
|
| 1095 |
* Get top active chatters
|
| 1096 |
*/
|
|
@@ -1149,11 +1245,33 @@ export async function getTopChatters(streamId = null, limit = 50) {
|
|
| 1149 |
.order('message_count', { ascending: false })
|
| 1150 |
.limit(limit);
|
| 1151 |
}
|
| 1152 |
-
|
| 1153 |
if (error) {
|
| 1154 |
console.error('[Supabase] Error getting top chatters:', error);
|
| 1155 |
return [];
|
| 1156 |
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1157 |
return data || [];
|
| 1158 |
}
|
| 1159 |
}
|
|
@@ -1201,10 +1319,20 @@ export async function getStatsSummary(streamId = null) {
|
|
| 1201 |
// Unique chatters count
|
| 1202 |
let uniqueChatters = 0;
|
| 1203 |
if (streamId) {
|
| 1204 |
-
|
| 1205 |
.from('stream_chatter_stats')
|
| 1206 |
.select('*', { count: 'exact', head: true })
|
| 1207 |
.eq('stream_id', streamId);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1208 |
uniqueChatters = count || 0;
|
| 1209 |
} else {
|
| 1210 |
const { count } = await supabase
|
|
|
|
| 1091 |
}
|
| 1092 |
}
|
| 1093 |
|
| 1094 |
+
/**
|
| 1095 |
+
* Self-healing: Rebuild stream_viewers and chat_users for a stream from its messages
|
| 1096 |
+
*/
|
| 1097 |
+
export async function rebuildViewersForStream(streamId) {
|
| 1098 |
+
if (!streamId || dbMode !== 'supabase') return;
|
| 1099 |
+
console.log(`[Self-Healing] Rebuilding viewers and chat users for stream ${streamId} from messages...`);
|
| 1100 |
+
|
| 1101 |
+
// Fetch messages from Supabase (only needed columns to reduce egress)
|
| 1102 |
+
const { data: messages, error } = await supabase
|
| 1103 |
+
.from('messages')
|
| 1104 |
+
.select('username, display_name, timestamp, is_mod, is_sub')
|
| 1105 |
+
.eq('stream_id', streamId);
|
| 1106 |
+
|
| 1107 |
+
if (error) {
|
| 1108 |
+
console.error(`[Self-Healing] Error fetching messages for stream ${streamId}:`, error);
|
| 1109 |
+
return;
|
| 1110 |
+
}
|
| 1111 |
+
|
| 1112 |
+
if (!messages || messages.length === 0) {
|
| 1113 |
+
console.log(`[Self-Healing] No messages found for stream ${streamId}.`);
|
| 1114 |
+
return;
|
| 1115 |
+
}
|
| 1116 |
+
|
| 1117 |
+
const viewersMap = new Map();
|
| 1118 |
+
const chatUsersMap = new Map();
|
| 1119 |
+
|
| 1120 |
+
for (const msg of messages) {
|
| 1121 |
+
const username = msg.username.toLowerCase();
|
| 1122 |
+
const displayName = msg.display_name || msg.username;
|
| 1123 |
+
const ts = msg.timestamp;
|
| 1124 |
+
|
| 1125 |
+
// Viewer key: stream_id-username
|
| 1126 |
+
const viewerKey = `${streamId}-${username}`;
|
| 1127 |
+
if (!viewersMap.has(viewerKey)) {
|
| 1128 |
+
viewersMap.set(viewerKey, {
|
| 1129 |
+
stream_id: streamId,
|
| 1130 |
+
username: username,
|
| 1131 |
+
display_name: displayName,
|
| 1132 |
+
has_chatted: true,
|
| 1133 |
+
is_mod: msg.is_mod || false,
|
| 1134 |
+
is_sub: msg.is_sub || false,
|
| 1135 |
+
first_seen: ts
|
| 1136 |
+
});
|
| 1137 |
+
} else {
|
| 1138 |
+
const viewer = viewersMap.get(viewerKey);
|
| 1139 |
+
if (new Date(ts) < new Date(viewer.first_seen)) {
|
| 1140 |
+
viewer.first_seen = ts;
|
| 1141 |
+
}
|
| 1142 |
+
if (msg.is_mod) viewer.is_mod = true;
|
| 1143 |
+
if (msg.is_sub) viewer.is_sub = true;
|
| 1144 |
+
}
|
| 1145 |
+
|
| 1146 |
+
// Chat user profile
|
| 1147 |
+
const userKey = username;
|
| 1148 |
+
if (!chatUsersMap.has(userKey)) {
|
| 1149 |
+
chatUsersMap.set(userKey, {
|
| 1150 |
+
username: username,
|
| 1151 |
+
display_name: displayName,
|
| 1152 |
+
is_mod: msg.is_mod || false,
|
| 1153 |
+
is_sub: msg.is_sub || false,
|
| 1154 |
+
is_vip: false,
|
| 1155 |
+
last_seen: ts
|
| 1156 |
+
});
|
| 1157 |
+
} else {
|
| 1158 |
+
const user = chatUsersMap.get(userKey);
|
| 1159 |
+
if (new Date(ts) > new Date(user.last_seen)) {
|
| 1160 |
+
user.last_seen = ts;
|
| 1161 |
+
user.display_name = displayName;
|
| 1162 |
+
}
|
| 1163 |
+
if (msg.is_mod) user.is_mod = true;
|
| 1164 |
+
if (msg.is_sub) user.is_sub = true;
|
| 1165 |
+
}
|
| 1166 |
+
}
|
| 1167 |
+
|
| 1168 |
+
const viewers = Array.from(viewersMap.values());
|
| 1169 |
+
const chatUsers = Array.from(chatUsersMap.values());
|
| 1170 |
+
|
| 1171 |
+
// Upsert chat_users
|
| 1172 |
+
const { error: userErr } = await supabase
|
| 1173 |
+
.from('chat_users')
|
| 1174 |
+
.upsert(chatUsers, { onConflict: 'username' });
|
| 1175 |
+
if (userErr) {
|
| 1176 |
+
console.error(`[Self-Healing] Error upserting chat_users:`, userErr);
|
| 1177 |
+
}
|
| 1178 |
+
|
| 1179 |
+
// Upsert stream_viewers
|
| 1180 |
+
const { error: viewErr } = await supabase
|
| 1181 |
+
.from('stream_viewers')
|
| 1182 |
+
.upsert(viewers, { onConflict: 'stream_id, username' });
|
| 1183 |
+
if (viewErr) {
|
| 1184 |
+
console.error(`[Self-Healing] Error upserting stream_viewers:`, viewErr);
|
| 1185 |
+
}
|
| 1186 |
+
|
| 1187 |
+
console.log(`[Self-Healing] Rebuilt and saved ${viewers.length} viewers for stream ${streamId}`);
|
| 1188 |
+
}
|
| 1189 |
+
|
| 1190 |
/**
|
| 1191 |
* Get top active chatters
|
| 1192 |
*/
|
|
|
|
| 1245 |
.order('message_count', { ascending: false })
|
| 1246 |
.limit(limit);
|
| 1247 |
}
|
| 1248 |
+
let { data, error } = await query;
|
| 1249 |
if (error) {
|
| 1250 |
console.error('[Supabase] Error getting top chatters:', error);
|
| 1251 |
return [];
|
| 1252 |
}
|
| 1253 |
+
|
| 1254 |
+
// Trigger self-healing if we expected chatters but got 0
|
| 1255 |
+
if (streamId && (!data || data.length === 0)) {
|
| 1256 |
+
const { count: msgCount } = await supabase
|
| 1257 |
+
.from('messages')
|
| 1258 |
+
.select('*', { count: 'exact', head: true })
|
| 1259 |
+
.eq('stream_id', streamId);
|
| 1260 |
+
|
| 1261 |
+
if (msgCount && msgCount > 0) {
|
| 1262 |
+
await rebuildViewersForStream(streamId);
|
| 1263 |
+
const requery = await supabase
|
| 1264 |
+
.from('stream_chatter_stats')
|
| 1265 |
+
.select('*')
|
| 1266 |
+
.eq('stream_id', streamId)
|
| 1267 |
+
.order('message_count', { ascending: false })
|
| 1268 |
+
.limit(limit);
|
| 1269 |
+
if (!requery.error && requery.data) {
|
| 1270 |
+
data = requery.data;
|
| 1271 |
+
}
|
| 1272 |
+
}
|
| 1273 |
+
}
|
| 1274 |
+
|
| 1275 |
return data || [];
|
| 1276 |
}
|
| 1277 |
}
|
|
|
|
| 1319 |
// Unique chatters count
|
| 1320 |
let uniqueChatters = 0;
|
| 1321 |
if (streamId) {
|
| 1322 |
+
let { count } = await supabase
|
| 1323 |
.from('stream_chatter_stats')
|
| 1324 |
.select('*', { count: 'exact', head: true })
|
| 1325 |
.eq('stream_id', streamId);
|
| 1326 |
+
|
| 1327 |
+
// Self-healing check if summary shows 0 viewers but we have messages
|
| 1328 |
+
if ((!count || count === 0) && msgCount > 0) {
|
| 1329 |
+
await rebuildViewersForStream(streamId);
|
| 1330 |
+
const requery = await supabase
|
| 1331 |
+
.from('stream_chatter_stats')
|
| 1332 |
+
.select('*', { count: 'exact', head: true })
|
| 1333 |
+
.eq('stream_id', streamId);
|
| 1334 |
+
count = requery.count || 0;
|
| 1335 |
+
}
|
| 1336 |
uniqueChatters = count || 0;
|
| 1337 |
} else {
|
| 1338 |
const { count } = await supabase
|