feat: enhance moderation functionality with type improvements and global broadcaster integration
This commit is contained in:
@@ -9,10 +9,22 @@ import {
|
||||
getPendingMessagesByConversation,
|
||||
updateMessageAIAnalysis,
|
||||
} from "./messageStore";
|
||||
import type { AnalysisQueueStatus, MessageRecord } from "./types";
|
||||
import type {
|
||||
AnalysisQueueStatus,
|
||||
MessageRecord,
|
||||
ModerationBroadcaster,
|
||||
} from "./types";
|
||||
|
||||
const logger = createChildLogger("ai-analyzer");
|
||||
|
||||
type ModerationGlobal = typeof globalThis & {
|
||||
moderationBroadcaster?: ModerationBroadcaster;
|
||||
};
|
||||
|
||||
function getModerationBroadcaster(): ModerationBroadcaster | undefined {
|
||||
return (globalThis as ModerationGlobal).moderationBroadcaster;
|
||||
}
|
||||
|
||||
// Debounce state per conversation key
|
||||
const conversationDebounceTimers = new Map<string, NodeJS.Timeout>();
|
||||
// Track conversations currently being processed
|
||||
@@ -117,7 +129,7 @@ async function processBatch(
|
||||
|
||||
// Broadcast analyzed messages
|
||||
for (const row of analyzedRows) {
|
||||
(globalThis as any).moderationBroadcaster?.messageAnalyzed(row);
|
||||
getModerationBroadcaster()?.messageAnalyzed(row);
|
||||
}
|
||||
|
||||
// Clear error cooldown on success
|
||||
@@ -147,7 +159,7 @@ async function processBatch(
|
||||
error: lastError,
|
||||
});
|
||||
if (row) {
|
||||
(globalThis as any).moderationBroadcaster?.messageAnalyzed(row);
|
||||
getModerationBroadcaster()?.messageAnalyzed(row);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,12 +1,25 @@
|
||||
import type { Client, Message } from "discord.js-selfbot-v13";
|
||||
import type { Channel, Client, Message } from "discord.js-selfbot-v13";
|
||||
import { config } from "../config";
|
||||
import { createChildLogger } from "../logger";
|
||||
import { captureMessage } from "./messageCapture";
|
||||
|
||||
const logger = createChildLogger("backlog-sync");
|
||||
|
||||
type BacklogChannel = Channel & {
|
||||
messages: {
|
||||
fetch(options: { limit: number; before?: string }): Promise<{
|
||||
size: number;
|
||||
values(): IterableIterator<Message>;
|
||||
}>;
|
||||
};
|
||||
};
|
||||
|
||||
function hasMessageBacklog(channel: Channel): channel is BacklogChannel {
|
||||
return "messages" in channel;
|
||||
}
|
||||
|
||||
async function syncChannelMessages(
|
||||
channel: any,
|
||||
channel: BacklogChannel,
|
||||
cutoffTime: number,
|
||||
): Promise<number> {
|
||||
let before: string | undefined;
|
||||
@@ -77,6 +90,10 @@ export async function syncSelectedChannelBacklog(
|
||||
logger.warn({ guildId, channelId }, "Channel not found for backlog sync");
|
||||
return 0;
|
||||
}
|
||||
if (!hasMessageBacklog(channel)) {
|
||||
logger.warn({ guildId, channelId }, "Channel cannot fetch message backlog");
|
||||
return 0;
|
||||
}
|
||||
|
||||
const cutoffTime = Date.now() - config.BACKLOG_SYNC_HOURS * 60 * 60 * 1000;
|
||||
logger.info(
|
||||
@@ -85,7 +102,7 @@ export async function syncSelectedChannelBacklog(
|
||||
);
|
||||
|
||||
try {
|
||||
const count = await syncChannelMessages(channel as any, cutoffTime);
|
||||
const count = await syncChannelMessages(channel, cutoffTime);
|
||||
logger.info(
|
||||
{ channelId, count },
|
||||
"Backlog sync completed for selected channel",
|
||||
|
||||
@@ -7,11 +7,14 @@ import type {
|
||||
ModerationWsEvent,
|
||||
} from "./types";
|
||||
|
||||
type ClientLike = Pick<WebSocket, "readyState" | "send">;
|
||||
export type BroadcasterClient = Pick<WebSocket, "readyState" | "send">;
|
||||
|
||||
const log = createChildLogger("broadcaster");
|
||||
|
||||
function sendJson(clients: Set<ClientLike>, event: ModerationWsEvent): void {
|
||||
function sendJson(
|
||||
clients: Set<BroadcasterClient>,
|
||||
event: ModerationWsEvent,
|
||||
): void {
|
||||
const payload = JSON.stringify({ ...event, timestamp: Date.now() });
|
||||
for (const client of clients) {
|
||||
if (client.readyState === 1) {
|
||||
@@ -28,14 +31,14 @@ function sendJson(clients: Set<ClientLike>, event: ModerationWsEvent): void {
|
||||
}
|
||||
|
||||
export function createBroadcaster() {
|
||||
const clients = new Set<ClientLike>();
|
||||
const clients = new Set<BroadcasterClient>();
|
||||
|
||||
return {
|
||||
addClient(client: ClientLike) {
|
||||
addClient(client: BroadcasterClient) {
|
||||
clients.add(client);
|
||||
log.debug({ clientCount: clients.size }, "Client added");
|
||||
},
|
||||
removeClient(client: ClientLike) {
|
||||
removeClient(client: BroadcasterClient) {
|
||||
clients.delete(client);
|
||||
log.debug({ clientCount: clients.size }, "Client removed");
|
||||
},
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { and, asc, desc, eq, isNull, or, sql } from "drizzle-orm";
|
||||
import { and, asc, desc, eq, isNull, or, type SQL, sql } from "drizzle-orm";
|
||||
import { getDatabase } from "../database/drizzle";
|
||||
import { attachmentsTable, messagesTable } from "../database/schema";
|
||||
import { createChildLogger } from "../logger";
|
||||
@@ -11,6 +11,29 @@ import type {
|
||||
|
||||
const logger = createChildLogger("message-store");
|
||||
|
||||
interface QueryBuilder<T = unknown> extends PromiseLike<T> {
|
||||
from(...args: unknown[]): QueryBuilder<T>;
|
||||
where(...args: unknown[]): QueryBuilder<T>;
|
||||
orderBy(...args: unknown[]): QueryBuilder<T>;
|
||||
limit(...args: unknown[]): QueryBuilder<T>;
|
||||
offset(...args: unknown[]): QueryBuilder<T>;
|
||||
values(...args: unknown[]): QueryBuilder<T>;
|
||||
onConflictDoNothing(...args: unknown[]): QueryBuilder<T>;
|
||||
returning(...args: unknown[]): QueryBuilder<T>;
|
||||
set(...args: unknown[]): QueryBuilder<T>;
|
||||
}
|
||||
|
||||
interface MessageDatabase {
|
||||
select<T = unknown[]>(...args: unknown[]): QueryBuilder<T>;
|
||||
selectDistinct<T = unknown[]>(...args: unknown[]): QueryBuilder<T>;
|
||||
insert<T = unknown>(...args: unknown[]): QueryBuilder<T>;
|
||||
update(...args: unknown[]): QueryBuilder<unknown>;
|
||||
}
|
||||
|
||||
function db(): MessageDatabase {
|
||||
return getDatabase() as unknown as MessageDatabase;
|
||||
}
|
||||
|
||||
// Cursor helpers for pagination
|
||||
interface CursorData {
|
||||
created_at: number;
|
||||
@@ -36,8 +59,8 @@ export function decodeCursor(cursor?: string): CursorData | null {
|
||||
|
||||
export async function insertMessage(message: MessageRecord): Promise<void> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
await db.insert(messagesTable).values(message).onConflictDoNothing();
|
||||
const database = db();
|
||||
await database.insert(messagesTable).values(message).onConflictDoNothing();
|
||||
|
||||
logger.debug(
|
||||
{ messageId: message.id, channelId: message.channel_id },
|
||||
@@ -59,14 +82,14 @@ export async function upsertMessageForCapture(
|
||||
message: MessageRecord,
|
||||
): Promise<boolean> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const database = db();
|
||||
const messageWithAIStatus = {
|
||||
...message,
|
||||
ai_status: "pending" as const,
|
||||
};
|
||||
|
||||
const rows = await db
|
||||
.insert(messagesTable)
|
||||
const rows = await database
|
||||
.insert<Array<{ id: string }>>(messagesTable)
|
||||
.values(messageWithAIStatus)
|
||||
.onConflictDoNothing()
|
||||
.returning({ id: messagesTable.id });
|
||||
@@ -95,8 +118,8 @@ export async function updateMessageAsEdited(
|
||||
editedAt: number,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
await db
|
||||
const database = db();
|
||||
await database
|
||||
.update(messagesTable)
|
||||
.set({
|
||||
edited_content: editedContent,
|
||||
@@ -130,8 +153,8 @@ export async function updateMessageAsDeleted(
|
||||
deletedAt: number,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
await db
|
||||
const database = db();
|
||||
await database
|
||||
.update(messagesTable)
|
||||
.set({
|
||||
deleted_at: deletedAt,
|
||||
@@ -158,8 +181,8 @@ export async function getMessagesByChannel(
|
||||
offset: number = 0,
|
||||
): Promise<MessageRecord[]> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const rows = await db
|
||||
const database = db();
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(messagesTable)
|
||||
.where(
|
||||
@@ -189,8 +212,11 @@ export async function insertAttachment(
|
||||
attachment: AttachmentRecord,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
await db.insert(attachmentsTable).values(attachment).onConflictDoNothing();
|
||||
const database = db();
|
||||
await database
|
||||
.insert(attachmentsTable)
|
||||
.values(attachment)
|
||||
.onConflictDoNothing();
|
||||
|
||||
logger.debug(
|
||||
{ attachmentId: attachment.id, messageId: attachment.message_id },
|
||||
@@ -214,8 +240,8 @@ export async function getAttachmentsByChannel(
|
||||
offset: number = 0,
|
||||
): Promise<AttachmentRecord[]> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const rows = await db
|
||||
const database = db();
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(attachmentsTable)
|
||||
.where(
|
||||
@@ -247,8 +273,8 @@ export async function updateAttachmentAsUploaded(
|
||||
uploadedAt: number,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
await db
|
||||
const database = db();
|
||||
await database
|
||||
.update(attachmentsTable)
|
||||
.set({
|
||||
uploaded_url: uploadedUrl,
|
||||
@@ -278,8 +304,8 @@ export async function updateAttachmentAsFailedUpload(
|
||||
error: string,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
await db
|
||||
const database = db();
|
||||
await database
|
||||
.update(attachmentsTable)
|
||||
.set({
|
||||
upload_status: "failed",
|
||||
@@ -315,8 +341,8 @@ export async function updateMessageAIAnalysis(
|
||||
result: AIAnalysisUpdate,
|
||||
): Promise<MessageRecord | null> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
await db
|
||||
const database = db();
|
||||
await database
|
||||
.update(messagesTable)
|
||||
.set({
|
||||
ai_status: result.status,
|
||||
@@ -329,7 +355,7 @@ export async function updateMessageAIAnalysis(
|
||||
})
|
||||
.where(eq(messagesTable.id, messageId));
|
||||
|
||||
const rows = await db
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(messagesTable)
|
||||
.where(eq(messagesTable.id, messageId));
|
||||
@@ -351,8 +377,8 @@ export async function getPendingAIAnalysisMessages(
|
||||
limit: number = 25,
|
||||
): Promise<MessageRecord[]> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const rows = await db
|
||||
const database = db();
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(messagesTable)
|
||||
.where(
|
||||
@@ -378,8 +404,8 @@ export async function getMessageById(
|
||||
messageId: string,
|
||||
): Promise<MessageRecord | null> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const rows = await db
|
||||
const database = db();
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(messagesTable)
|
||||
.where(eq(messagesTable.id, messageId));
|
||||
@@ -401,8 +427,8 @@ export async function listMessages(
|
||||
query: MessageQuery,
|
||||
): Promise<PageResult<MessageRecord>> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const conditions: any[] = [];
|
||||
const database = db();
|
||||
const conditions: SQL[] = [];
|
||||
|
||||
// Apply filters
|
||||
if (query.guildId) {
|
||||
@@ -411,10 +437,7 @@ export async function listMessages(
|
||||
|
||||
if (query.channelId) {
|
||||
conditions.push(
|
||||
or(
|
||||
eq(messagesTable.channel_id, query.channelId),
|
||||
eq(messagesTable.thread_id, query.channelId),
|
||||
),
|
||||
sql`(${messagesTable.channel_id} = ${query.channelId} or ${messagesTable.thread_id} = ${query.channelId})`,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -427,11 +450,7 @@ export async function listMessages(
|
||||
}
|
||||
|
||||
if (query.status && query.status.length > 0) {
|
||||
conditions.push(
|
||||
or(
|
||||
...query.status.map((status) => eq(messagesTable.ai_status, status)),
|
||||
),
|
||||
);
|
||||
conditions.push(sql`${messagesTable.ai_status} in ${query.status}`);
|
||||
}
|
||||
|
||||
// Text search
|
||||
@@ -445,20 +464,14 @@ export async function listMessages(
|
||||
const cursorData = decodeCursor(query.cursor);
|
||||
if (cursorData) {
|
||||
conditions.push(
|
||||
or(
|
||||
sql`${messagesTable.created_at} < ${cursorData.created_at}`,
|
||||
and(
|
||||
eq(messagesTable.created_at, cursorData.created_at),
|
||||
sql`${messagesTable.id} < ${cursorData.id}`,
|
||||
),
|
||||
),
|
||||
sql`(${messagesTable.created_at} < ${cursorData.created_at} or (${messagesTable.created_at} = ${cursorData.created_at} and ${messagesTable.id} < ${cursorData.id}))`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch limit + 1 to determine if there's a next page
|
||||
const fetchLimit = query.limit + 1;
|
||||
const rows = await db
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(messagesTable)
|
||||
.where(conditions.length > 0 ? and(...conditions) : undefined)
|
||||
@@ -506,7 +519,7 @@ export async function getConversationContextBefore(input: {
|
||||
limit: number;
|
||||
}): Promise<MessageRecord[]> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const database = db();
|
||||
const { channelId, threadId, beforeCreatedAt, limit } = input;
|
||||
|
||||
// Query same thread if threadId exists, otherwise channelId
|
||||
@@ -514,7 +527,7 @@ export async function getConversationContextBefore(input: {
|
||||
? eq(messagesTable.thread_id, threadId)
|
||||
: eq(messagesTable.channel_id, channelId);
|
||||
|
||||
const rows = await db
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(messagesTable)
|
||||
.where(
|
||||
@@ -547,11 +560,11 @@ export async function getPendingMessagesByConversation(
|
||||
limit: number = 25,
|
||||
): Promise<MessageRecord[]> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
// conversationKey is either thread_id or channel_id
|
||||
// Query both to safely handle the key
|
||||
const rows = await db
|
||||
const rows = await database
|
||||
.select()
|
||||
.from(messagesTable)
|
||||
.where(
|
||||
@@ -584,11 +597,11 @@ export async function getPendingConversationKeys(
|
||||
limit: number = 100,
|
||||
): Promise<string[]> {
|
||||
try {
|
||||
const db = getDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
// Get distinct conversation keys (thread_id or channel_id) for pending messages
|
||||
const rows = await db
|
||||
.selectDistinct({
|
||||
const rows = await database
|
||||
.selectDistinct<Array<{ thread_id: string | null; channel_id: string }>>({
|
||||
thread_id: messagesTable.thread_id,
|
||||
channel_id: messagesTable.channel_id,
|
||||
})
|
||||
@@ -602,7 +615,7 @@ export async function getPendingConversationKeys(
|
||||
.limit(limit);
|
||||
|
||||
const keys: string[] = [];
|
||||
for (const row of rows as any[]) {
|
||||
for (const row of rows) {
|
||||
const key = row.thread_id || row.channel_id;
|
||||
if (key && !keys.includes(key)) {
|
||||
keys.push(key);
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
import type { ModerationBroadcaster } from "./broadcaster";
|
||||
import type { BroadcasterClient, ModerationBroadcaster } from "./broadcaster";
|
||||
|
||||
export type AIStatus = "pending" | "clean" | "warn" | "flagged" | "error";
|
||||
|
||||
export type { ModerationBroadcaster };
|
||||
export type { BroadcasterClient, ModerationBroadcaster };
|
||||
|
||||
export interface MessageRecord {
|
||||
id: string;
|
||||
|
||||
+65
-29
@@ -8,8 +8,28 @@ import { createChildLogger } from "./logger";
|
||||
|
||||
const logger = createChildLogger("muxer-queue");
|
||||
|
||||
// Type alias for backward compatibility
|
||||
export type SqliteDatabase = any;
|
||||
interface QueryBuilder<T = unknown> extends PromiseLike<T> {
|
||||
from(...args: unknown[]): QueryBuilder<T>;
|
||||
where(...args: unknown[]): QueryBuilder<T>;
|
||||
orderBy(...args: unknown[]): QueryBuilder<T>;
|
||||
limit(...args: unknown[]): QueryBuilder<T>;
|
||||
values(...args: unknown[]): QueryBuilder<T>;
|
||||
onConflictDoNothing(...args: unknown[]): QueryBuilder<T>;
|
||||
onConflictDoUpdate(...args: unknown[]): QueryBuilder<T>;
|
||||
set(...args: unknown[]): QueryBuilder<T>;
|
||||
groupBy(...args: unknown[]): QueryBuilder<T>;
|
||||
}
|
||||
|
||||
export interface SqliteDatabase {
|
||||
select<T = unknown[]>(...args: unknown[]): QueryBuilder<T>;
|
||||
insert(...args: unknown[]): QueryBuilder<unknown>;
|
||||
update(...args: unknown[]): QueryBuilder<unknown>;
|
||||
delete(...args: unknown[]): QueryBuilder<unknown>;
|
||||
}
|
||||
|
||||
function db(): SqliteDatabase {
|
||||
return getDrizzleDatabase() as unknown as SqliteDatabase;
|
||||
}
|
||||
|
||||
export interface MuxerJobData {
|
||||
userId: string;
|
||||
@@ -18,6 +38,22 @@ export interface MuxerJobData {
|
||||
outputDir: string;
|
||||
}
|
||||
|
||||
interface StoredJobRow {
|
||||
id: string;
|
||||
data: string;
|
||||
status: "pending" | "processing" | "completed" | "failed";
|
||||
attempts: number;
|
||||
maxAttempts: number;
|
||||
createdAt: number;
|
||||
updatedAt: number;
|
||||
error: string | null;
|
||||
}
|
||||
|
||||
interface JobStatsRow {
|
||||
status: "pending" | "processing" | "completed" | "failed";
|
||||
count: number | string | { count: number | string };
|
||||
}
|
||||
|
||||
interface StoredJob {
|
||||
id: string;
|
||||
data: string;
|
||||
@@ -31,7 +67,7 @@ interface StoredJob {
|
||||
|
||||
// Export getDatabase for backward compatibility with webserver.ts
|
||||
export function getDatabase(): SqliteDatabase {
|
||||
return getDrizzleDatabase() as any;
|
||||
return db();
|
||||
}
|
||||
|
||||
export async function getPersistedValue<T>(
|
||||
@@ -39,10 +75,10 @@ export async function getPersistedValue<T>(
|
||||
fallback: T,
|
||||
): Promise<T> {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
const row = await db
|
||||
.select()
|
||||
const row = await database
|
||||
.select<Array<{ value: string }>>()
|
||||
.from(uiStateTable)
|
||||
.where(eq(uiStateTable.key, key))
|
||||
.limit(1);
|
||||
@@ -61,9 +97,9 @@ export async function setPersistedValue(
|
||||
value: unknown,
|
||||
): Promise<void> {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
await db
|
||||
await database
|
||||
.insert(uiStateTable)
|
||||
.values({
|
||||
key,
|
||||
@@ -82,12 +118,12 @@ export async function setPersistedValue(
|
||||
export async function enqueueMuxerJob(data: MuxerJobData): Promise<string> {
|
||||
try {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
const jobId = `${data.userId}-${data.sessionId}`;
|
||||
const now = Date.now();
|
||||
|
||||
await db
|
||||
await database
|
||||
.insert(muxerJobsTable)
|
||||
.values({
|
||||
id: jobId,
|
||||
@@ -120,16 +156,16 @@ export async function enqueueMuxerJob(data: MuxerJobData): Promise<string> {
|
||||
|
||||
export async function getPendingJobs(): Promise<StoredJob[]> {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
const rows = await db
|
||||
.select()
|
||||
const rows = await database
|
||||
.select<StoredJobRow[]>()
|
||||
.from(muxerJobsTable)
|
||||
.where(eq(muxerJobsTable.status, "pending"))
|
||||
.orderBy(asc(muxerJobsTable.createdAt))
|
||||
.limit(10);
|
||||
|
||||
return rows.map((row: any) => ({
|
||||
return rows.map((row) => ({
|
||||
id: row.id,
|
||||
data: row.data,
|
||||
status: row.status as "pending" | "processing" | "completed" | "failed",
|
||||
@@ -147,11 +183,11 @@ export async function updateJobStatus(
|
||||
error?: string,
|
||||
): Promise<void> {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
const now = Date.now();
|
||||
|
||||
if (status === "failed") {
|
||||
await db
|
||||
await database
|
||||
.update(muxerJobsTable)
|
||||
.set({
|
||||
status,
|
||||
@@ -161,7 +197,7 @@ export async function updateJobStatus(
|
||||
})
|
||||
.where(eq(muxerJobsTable.id, jobId));
|
||||
} else {
|
||||
await db
|
||||
await database
|
||||
.update(muxerJobsTable)
|
||||
.set({
|
||||
status,
|
||||
@@ -175,10 +211,10 @@ export async function updateJobStatus(
|
||||
|
||||
export async function retryFailedJob(jobId: string): Promise<boolean> {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
const jobs = await db
|
||||
.select()
|
||||
const jobs = await database
|
||||
.select<StoredJobRow[]>()
|
||||
.from(muxerJobsTable)
|
||||
.where(eq(muxerJobsTable.id, jobId))
|
||||
.limit(1);
|
||||
@@ -198,7 +234,7 @@ export async function retryFailedJob(jobId: string): Promise<boolean> {
|
||||
return false;
|
||||
}
|
||||
|
||||
await db
|
||||
await database
|
||||
.update(muxerJobsTable)
|
||||
.set({
|
||||
status: "pending",
|
||||
@@ -215,10 +251,10 @@ export async function cleanupCompletedJobs(
|
||||
olderThanMs: number = 24 * 60 * 60 * 1000,
|
||||
): Promise<number> {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
const cutoffTime = Date.now() - olderThanMs;
|
||||
|
||||
const result = await db
|
||||
const result = await database
|
||||
.delete(muxerJobsTable)
|
||||
.where(
|
||||
and(
|
||||
@@ -228,8 +264,8 @@ export async function cleanupCompletedJobs(
|
||||
);
|
||||
|
||||
const deletedCount =
|
||||
typeof result === "object" && "rowsAffected" in result
|
||||
? result.rowsAffected
|
||||
typeof result === "object" && result !== null && "rowsAffected" in result
|
||||
? Number(result.rowsAffected)
|
||||
: 0;
|
||||
|
||||
logger.info({ deletedCount }, "Cleaned up completed jobs");
|
||||
@@ -244,10 +280,10 @@ export async function getJobStats(): Promise<{
|
||||
failed: number;
|
||||
}> {
|
||||
await initializeDatabase();
|
||||
const db = getDrizzleDatabase() as any;
|
||||
const database = db();
|
||||
|
||||
const rows = await db
|
||||
.select({
|
||||
const rows = await database
|
||||
.select<JobStatsRow[]>({
|
||||
status: muxerJobsTable.status,
|
||||
count: sql<number>`COUNT(*)`,
|
||||
})
|
||||
@@ -264,7 +300,7 @@ export async function getJobStats(): Promise<{
|
||||
for (const row of rows) {
|
||||
const count =
|
||||
typeof row.count === "object" && "count" in row.count
|
||||
? (row.count as any).count
|
||||
? Number((row.count as { count: number | string }).count)
|
||||
: Number(row.count);
|
||||
if (row.status === "pending") stats.pending = count;
|
||||
else if (row.status === "processing") stats.processing = count;
|
||||
|
||||
+3
-1
@@ -1,6 +1,7 @@
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import {
|
||||
type DiscordGatewayAdapterCreator,
|
||||
EndBehaviorType,
|
||||
entersState,
|
||||
getVoiceConnection,
|
||||
@@ -41,7 +42,8 @@ export async function startRecording(
|
||||
const connection = joinVoiceChannel({
|
||||
channelId: channel.id,
|
||||
guildId: channel.guild.id,
|
||||
adapterCreator: channel.guild.voiceAdapterCreator as any,
|
||||
adapterCreator: channel.guild
|
||||
.voiceAdapterCreator as DiscordGatewayAdapterCreator,
|
||||
selfDeaf: false,
|
||||
selfMute: false,
|
||||
debug: true,
|
||||
|
||||
@@ -88,13 +88,16 @@ export class VoiceController {
|
||||
await guild.channels.fetch().catch(() => null);
|
||||
|
||||
const threads: ChannelSummary[] = [];
|
||||
type ThreadFetchResult = {
|
||||
threads: Map<string, { id: string; name: string; type: string }>;
|
||||
};
|
||||
for (const channel of guild.channels.cache.values()) {
|
||||
const threadParent = channel as typeof channel & {
|
||||
threads?: {
|
||||
fetch: (options: {
|
||||
archived: boolean;
|
||||
limit: number;
|
||||
}) => Promise<any>;
|
||||
}) => Promise<ThreadFetchResult>;
|
||||
};
|
||||
};
|
||||
if (!threadParent.threads?.fetch) continue;
|
||||
|
||||
+17
-4
@@ -10,6 +10,7 @@ import { AppError } from "./errors";
|
||||
import { createChildLogger, logger } from "./logger";
|
||||
import { getMetrics, uptimeGauge } from "./metrics";
|
||||
import { createBroadcaster } from "./moderation/broadcaster";
|
||||
import type { ModerationBroadcaster } from "./moderation/types";
|
||||
import { getPersistedValue, setPersistedValue } from "./muxer-queue";
|
||||
import { discordPlayer } from "./player";
|
||||
import { createAnalysisRoutes } from "./routes/analysisRoutes";
|
||||
@@ -26,6 +27,15 @@ const activeUsers = new Map<
|
||||
{ username: string; avatar: string; speaking: boolean }
|
||||
>();
|
||||
|
||||
type VoiceGlobals = typeof globalThis & {
|
||||
moderationBroadcaster?: ModerationBroadcaster;
|
||||
broadcastPcmToWeb?: (chunk: Buffer, userId: string) => void;
|
||||
updateActiveUser?: (
|
||||
userId: string,
|
||||
data: { username: string; avatar: string; speaking: boolean },
|
||||
) => void;
|
||||
};
|
||||
|
||||
interface SharedUIState {
|
||||
selectedGuild: string;
|
||||
selectedVoiceChannel: string;
|
||||
@@ -118,7 +128,7 @@ export async function startWebserver(
|
||||
|
||||
// Create broadcaster instance
|
||||
const broadcaster = createBroadcaster();
|
||||
(globalThis as any).moderationBroadcaster = broadcaster;
|
||||
(globalThis as VoiceGlobals).moderationBroadcaster = broadcaster;
|
||||
|
||||
// Security headers. CSP disabled because the current static UI uses inline scripts/styles.
|
||||
app.use(
|
||||
@@ -196,7 +206,10 @@ export async function startWebserver(
|
||||
app.use("/api", createSyncRoutes(_client));
|
||||
|
||||
// Inbound: Discord PCM → tagged chunks → browser
|
||||
(global as any).broadcastPcmToWeb = (chunk: Buffer, userId: string) => {
|
||||
(globalThis as VoiceGlobals).broadcastPcmToWeb = (
|
||||
chunk: Buffer,
|
||||
userId: string,
|
||||
) => {
|
||||
let hash = 0;
|
||||
for (let i = 0; i < userId.length; i++) {
|
||||
hash = (hash << 5) - hash + userId.charCodeAt(i);
|
||||
@@ -210,7 +223,7 @@ export async function startWebserver(
|
||||
}
|
||||
};
|
||||
|
||||
(global as any).updateActiveUser = (
|
||||
(globalThis as VoiceGlobals).updateActiveUser = (
|
||||
userId: string,
|
||||
data: { username: string; avatar: string; speaking: boolean },
|
||||
) => {
|
||||
@@ -327,7 +340,7 @@ export async function startWebserver(
|
||||
);
|
||||
ws.send(JSON.stringify({ type: "ui_state", state: getSharedUIState() }));
|
||||
|
||||
ws.on("message", (data: any) => {
|
||||
ws.on("message", (data: Buffer | ArrayBuffer | Buffer[]) => {
|
||||
if (!Buffer.isBuffer(data)) return;
|
||||
lastBrowserAudioTime = Date.now();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user