SPB Git

spb/vquant Public MIT

VibeQuant — AI-powered institutional-grade financial intelligence platform.

TypeScript 84.3% Python 11.7% JavaScript 1.6% CSS 1.5% HTML 0.7%
14.6 KB · 403 lines typescript
Raw Blame History
1/*2 * =============================================================================3 *  VibeQuant (vquant) — AI-Powered Financial Intelligence Platform4 * -----------------------------------------------------------------------------5 *  File:      server/storage.ts6 *7 *  Author:    Simon-Pierre Boucher8 *  Contact:   contact@spboucher.ai9 *  Website:   https://www.spboucher.ai10 *  Demo:      https://www.vquant.ai11 *  License:   MIT (see LICENSE)12 *13 *  Copyright © 2026 Simon-Pierre Boucher. All rights reserved.14 * =============================================================================15 */1617import { db, schema } from "./db";18import { eq, desc, and, lt, gte } from "drizzle-orm";19import { logger } from "./utils/logger";2021// Infer types from the active schema22type User = typeof schema.users.$inferSelect;23type InsertUser = Omit<User, "id" | "createdAt">;2425type CrawledPage = typeof schema.crawledPages.$inferSelect;26type InsertCrawledPage = Omit<CrawledPage, "id" | "crawledAt">;2728type Message = typeof schema.messages.$inferSelect;29type InsertMessage = Omit<Message, "id" | "createdAt">;3031type SharedReport = typeof schema.sharedReports.$inferSelect;32type InsertSharedReport = Omit<SharedReport, "id" | "createdAt">;3334type ConversationSession = typeof schema.conversationSessions.$inferSelect;35type InsertConversationSession = Omit<ConversationSession, "id" | "createdAt" | "updatedAt">;3637type ActiveUser = typeof schema.activeUsers.$inferSelect;38type InsertActiveUser = Omit<ActiveUser, "id" | "createdAt">;3940type AnalyticsMetric = typeof schema.analyticsMetrics.$inferSelect;41type InsertAnalyticsMetric = Omit<AnalyticsMetric, "id" | "timestamp">;4243type RequestLog = typeof schema.requestLogs.$inferSelect;44type InsertRequestLog = Omit<RequestLog, "id" | "timestamp">;4546const {47  users,48  crawledPages,49  messages,50  sharedReports,51  conversationSessions,52  activeUsers,53  analyticsMetrics,54  requestLogs,55} = schema;5657export interface IStorage {58  // Users59  getUser(id: string): Promise<User | undefined>;60  getUserByToken(token: string): Promise<User | undefined>;61  createUser(user: InsertUser): Promise<User>;62  getAllUsers(): Promise<User[]>;63  deleteUser(id: string): Promise<void>;6465  // Crawled Pages66  getCrawledPage(id: string): Promise<CrawledPage | undefined>;67  getCrawledPageByUrl(url: string): Promise<CrawledPage | undefined>;68  createCrawledPage(page: InsertCrawledPage): Promise<CrawledPage>;69  updateCrawledPage(id: string, page: Partial<InsertCrawledPage>): Promise<CrawledPage>;7071  // Messages72  getMessagesBySession(sessionId: string): Promise<Message[]>;73  createMessage(message: InsertMessage): Promise<Message>;7475  // Shared Reports76  getSharedReportByShareId(shareId: string): Promise<SharedReport | undefined>;77  createSharedReport(report: InsertSharedReport): Promise<SharedReport>;78  getAllSharedReports(): Promise<SharedReport[]>;79  deleteSharedReport(id: string): Promise<void>;8081  // Conversation Sessions82  getConversationSessions(userId?: string): Promise<ConversationSession[]>;83  getConversationSession(sessionId: string): Promise<ConversationSession | undefined>;84  getConversationSessionByUserAndSessionId(sessionId: string, userId: string): Promise<ConversationSession | undefined>;85  createConversationSession(session: InsertConversationSession): Promise<ConversationSession>;86  updateConversationSession(sessionId: string, updates: Partial<InsertConversationSession>): Promise<ConversationSession>;87  getAllConversationSessions(): Promise<ConversationSession[]>;88  deleteConversationSessionById(id: string): Promise<void>;8990  // Active Users91  getActiveUserBySessionId(sessionId: string): Promise<ActiveUser | undefined>;92  upsertActiveUser(user: InsertActiveUser): Promise<ActiveUser>;93  updateActiveUserHeartbeat(sessionId: string, status: string, currentQuery?: string): Promise<void>;94  getActiveUsers(): Promise<ActiveUser[]>;95  cleanupStaleUsers(timeoutMinutes?: number): Promise<void>;9697  // Analytics Metrics98  createAnalyticsMetric(metric: InsertAnalyticsMetric): Promise<AnalyticsMetric>;99  getLatestMetrics(periodType: string, limit?: number): Promise<AnalyticsMetric[]>;100  getMetricsInTimeRange(startTime: Date, endTime: Date): Promise<AnalyticsMetric[]>;101102  // Request Logs103  createRequestLog(log: InsertRequestLog): Promise<RequestLog>;104  getRecentRequestLogs(limit?: number): Promise<RequestLog[]>;105  getRequestLogsByStatus(status: string, limit?: number): Promise<RequestLog[]>;106}107108export class DatabaseStorage implements IStorage {109  // ─── Users ──────────────────────────────────────────────110111  async getUser(id: string): Promise<User | undefined> {112    const [user] = await (db as any).select().from(users).where(eq(users.id, id));113    return user ?? undefined;114  }115116  async getUserByToken(token: string): Promise<User | undefined> {117    const [user] = await (db as any).select().from(users).where(eq(users.token, token));118    return user ?? undefined;119  }120121  async createUser(insertUser: InsertUser): Promise<User> {122    const [user] = await (db as any)123      .insert(users)124      .values({ id: crypto.randomUUID(), ...insertUser, createdAt: new Date() })125      .returning();126    return user;127  }128129  async getAllUsers(): Promise<User[]> {130    return await (db as any).select().from(users).orderBy(desc(users.createdAt));131  }132133  async deleteUser(id: string): Promise<void> {134    await (db as any).delete(conversationSessions).where(eq(conversationSessions.userId, id));135    await (db as any).delete(users).where(eq(users.id, id));136  }137138  // ─── Crawled Pages ──────────────────────────────────────139140  async getCrawledPage(id: string): Promise<CrawledPage | undefined> {141    const [page] = await (db as any).select().from(crawledPages).where(eq(crawledPages.id, id));142    return page ?? undefined;143  }144145  async getCrawledPageByUrl(url: string): Promise<CrawledPage | undefined> {146    const [page] = await (db as any).select().from(crawledPages).where(eq(crawledPages.url, url));147    return page ?? undefined;148  }149150  async createCrawledPage(page: InsertCrawledPage): Promise<CrawledPage> {151    const [created] = await (db as any)152      .insert(crawledPages)153      .values({ id: crypto.randomUUID(), ...page, crawledAt: new Date() })154      .returning();155    return created;156  }157158  async updateCrawledPage(id: string, page: Partial<InsertCrawledPage>): Promise<CrawledPage> {159    const [updated] = await (db as any)160      .update(crawledPages)161      .set(page)162      .where(eq(crawledPages.id, id))163      .returning();164    return updated;165  }166167  // ─── Messages ───────────────────────────────────────────168169  async getMessagesBySession(sessionId: string): Promise<Message[]> {170    const msgs = await (db as any)171      .select()172      .from(messages)173      .where(eq(messages.sessionId, sessionId))174      .orderBy(desc(messages.createdAt));175176    return msgs.map((msg: Message) => ({177      ...msg,178      sources: typeof msg.sources === "string" ? JSON.parse(msg.sources) : msg.sources,179    }));180  }181182  async createMessage(message: InsertMessage): Promise<Message> {183    const [msg] = await (db as any)184      .insert(messages)185      .values({186        id: crypto.randomUUID(),187        ...message,188        sources: Array.isArray(message.sources) ? JSON.stringify(message.sources) : message.sources,189        createdAt: new Date(),190      })191      .returning();192193    return {194      ...msg,195      sources: typeof msg.sources === "string" ? JSON.parse(msg.sources) : msg.sources,196    };197  }198199  // ─── Shared Reports ─────────────────────────────────────200201  async getSharedReportByShareId(shareId: string): Promise<SharedReport | undefined> {202    const [report] = await (db as any)203      .select()204      .from(sharedReports)205      .where(eq(sharedReports.shareId, shareId));206    return report ?? undefined;207  }208209  async createSharedReport(report: InsertSharedReport): Promise<SharedReport> {210    const [created] = await (db as any)211      .insert(sharedReports)212      .values({ id: crypto.randomUUID(), ...report, createdAt: new Date() })213      .returning();214    return created;215  }216217  async getAllSharedReports(): Promise<SharedReport[]> {218    return await (db as any).select().from(sharedReports).orderBy(desc(sharedReports.createdAt));219  }220221  async deleteSharedReport(id: string): Promise<void> {222    await (db as any).delete(sharedReports).where(eq(sharedReports.id, id));223  }224225  // ─── Conversation Sessions ──────────────────────────────226227  async getConversationSessions(userId?: string): Promise<ConversationSession[]> {228    if (userId) {229      return await (db as any)230        .select()231        .from(conversationSessions)232        .where(eq(conversationSessions.userId, userId))233        .orderBy(desc(conversationSessions.updatedAt));234    }235    return await (db as any)236      .select()237      .from(conversationSessions)238      .orderBy(desc(conversationSessions.updatedAt));239  }240241  async getConversationSession(sessionId: string): Promise<ConversationSession | undefined> {242    const [session] = await (db as any)243      .select()244      .from(conversationSessions)245      .where(eq(conversationSessions.sessionId, sessionId));246    return session ?? undefined;247  }248249  async getConversationSessionByUserAndSessionId(250    sessionId: string,251    userId: string,252  ): Promise<ConversationSession | undefined> {253    const [session] = await (db as any)254      .select()255      .from(conversationSessions)256      .where(257        and(258          eq(conversationSessions.sessionId, sessionId),259          eq(conversationSessions.userId, userId),260        ),261      );262    return session ?? undefined;263  }264265  async createConversationSession(session: InsertConversationSession): Promise<ConversationSession> {266    const now = new Date();267    const [created] = await (db as any)268      .insert(conversationSessions)269      .values({ id: crypto.randomUUID(), ...session, createdAt: now, updatedAt: now })270      .returning();271    return created;272  }273274  async updateConversationSession(275    sessionId: string,276    updates: Partial<InsertConversationSession>,277  ): Promise<ConversationSession> {278    const [updated] = await (db as any)279      .update(conversationSessions)280      .set({ ...updates, updatedAt: new Date() })281      .where(eq(conversationSessions.sessionId, sessionId))282      .returning();283    return updated;284  }285286  async getAllConversationSessions(): Promise<ConversationSession[]> {287    return await (db as any)288      .select()289      .from(conversationSessions)290      .orderBy(desc(conversationSessions.updatedAt));291  }292293  async deleteConversationSessionById(id: string): Promise<void> {294    await (db as any).delete(conversationSessions).where(eq(conversationSessions.id, id));295  }296297  // ─── Active Users ───────────────────────────────────────298299  async getActiveUserBySessionId(sessionId: string): Promise<ActiveUser | undefined> {300    const [user] = await (db as any)301      .select()302      .from(activeUsers)303      .where(eq(activeUsers.sessionId, sessionId));304    return user ?? undefined;305  }306307  async upsertActiveUser(user: InsertActiveUser): Promise<ActiveUser> {308    const existing = await this.getActiveUserBySessionId(user.sessionId);309310    if (existing) {311      const [updated] = await (db as any)312        .update(activeUsers)313        .set({ ...user, lastHeartbeat: new Date() })314        .where(eq(activeUsers.sessionId, user.sessionId))315        .returning();316      return updated;317    }318319    const [created] = await (db as any)320      .insert(activeUsers)321      .values({ id: crypto.randomUUID(), ...user, lastHeartbeat: new Date(), createdAt: new Date() })322      .returning();323    return created;324  }325326  async updateActiveUserHeartbeat(sessionId: string, status: string, currentQuery?: string): Promise<void> {327    await (db as any)328      .update(activeUsers)329      .set({ lastHeartbeat: new Date(), status, currentQuery: currentQuery ?? null })330      .where(eq(activeUsers.sessionId, sessionId));331  }332333  async getActiveUsers(): Promise<ActiveUser[]> {334    const twoMinutesAgo = new Date(Date.now() - 2 * 60 * 1000);335    return await (db as any)336      .select()337      .from(activeUsers)338      .where(gte(activeUsers.lastHeartbeat, twoMinutesAgo))339      .orderBy(desc(activeUsers.lastHeartbeat));340  }341342  async cleanupStaleUsers(timeoutMinutes: number = 5): Promise<void> {343    const cutoffTime = new Date(Date.now() - timeoutMinutes * 60 * 1000);344    await (db as any).delete(activeUsers).where(lt(activeUsers.lastHeartbeat, cutoffTime));345  }346347  // ─── Analytics Metrics ──────────────────────────────────348349  async createAnalyticsMetric(metric: InsertAnalyticsMetric): Promise<AnalyticsMetric> {350    const [created] = await (db as any)351      .insert(analyticsMetrics)352      .values({ id: crypto.randomUUID(), ...metric, timestamp: new Date() })353      .returning();354    return created;355  }356357  async getLatestMetrics(periodType: string, limit: number = 60): Promise<AnalyticsMetric[]> {358    return await (db as any)359      .select()360      .from(analyticsMetrics)361      .where(eq(analyticsMetrics.periodType, periodType))362      .orderBy(desc(analyticsMetrics.timestamp))363      .limit(limit);364  }365366  async getMetricsInTimeRange(startTime: Date, endTime: Date): Promise<AnalyticsMetric[]> {367    return await (db as any)368      .select()369      .from(analyticsMetrics)370      .where(and(gte(analyticsMetrics.timestamp, startTime), lt(analyticsMetrics.timestamp, endTime)))371      .orderBy(desc(analyticsMetrics.timestamp));372  }373374  // ─── Request Logs ───────────────────────────────────────375376  async createRequestLog(log: InsertRequestLog): Promise<RequestLog> {377    const [created] = await (db as any)378      .insert(requestLogs)379      .values({ id: crypto.randomUUID(), ...log, timestamp: new Date() })380      .returning();381    return created;382  }383384  async getRecentRequestLogs(limit: number = 100): Promise<RequestLog[]> {385    return await (db as any)386      .select()387      .from(requestLogs)388      .orderBy(desc(requestLogs.timestamp))389      .limit(limit);390  }391392  async getRequestLogsByStatus(status: string, limit: number = 100): Promise<RequestLog[]> {393    return await (db as any)394      .select()395      .from(requestLogs)396      .where(eq(requestLogs.status, status))397      .orderBy(desc(requestLogs.timestamp))398      .limit(limit);399  }400}401402export const storage = new DatabaseStorage();403