From a66c1e3b1604c8a19c2c090caf8d0c0d136893a0 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Thu, 15 May 2025 19:38:40 -0700 Subject: [PATCH] feat(logs): added log archive --- apps/sim/app/api/logs/cleanup/route.ts | 128 ++++++++++++++++++++----- 1 file changed, 104 insertions(+), 24 deletions(-) diff --git a/apps/sim/app/api/logs/cleanup/route.ts b/apps/sim/app/api/logs/cleanup/route.ts index ddcea99e46..058e316b4e 100644 --- a/apps/sim/app/api/logs/cleanup/route.ts +++ b/apps/sim/app/api/logs/cleanup/route.ts @@ -1,11 +1,21 @@ import { NextResponse } from 'next/server' -import { sql } from 'drizzle-orm' +import { PutObjectCommand } from '@aws-sdk/client-s3' +import { and, eq, inArray, lt, sql } from 'drizzle-orm' import { createLogger } from '@/lib/logs/console-logger' +import { s3Client } from '@/lib/uploads/s3-client' import { db } from '@/db' -import { subscription, user, workflowLogs } from '@/db/schema' +import { subscription, user, workflow, workflowLogs } from '@/db/schema' + +export const dynamic = 'force-dynamic' const logger = createLogger('LogsCleanup') +const BATCH_SIZE = 500 +const S3_CONFIG = { + bucket: process.env.S3_LOGS_BUCKET_NAME || '', + region: process.env.AWS_REGION || '', +} + export async function POST(request: Request) { try { const authHeader = request.headers.get('authorization') @@ -19,9 +29,13 @@ export async function POST(request: Request) { return new NextResponse('Unauthorized', { status: 401 }) } + if (!S3_CONFIG.bucket || !S3_CONFIG.region) { + return new NextResponse('Configuration error: S3 bucket or region not set', { status: 500 }) + } + const retentionDate = new Date() retentionDate.setDate( - retentionDate.getDate() - Number(process.env.FREE_PLAN_LOG_RETENTION_DAYS) + retentionDate.getDate() - Number(process.env.FREE_PLAN_LOG_RETENTION_DAYS || '7') ) const freeUsers = await db @@ -39,39 +53,105 @@ export async function POST(request: Request) { } const freeUserIds = freeUsers.map((u) => u.userId) - logger.info(`Found ${freeUserIds.length} free users for log cleanup`) - const freeUserWorkflows = await db - .select({ workflowId: workflowLogs.workflowId }) + const workflowsQuery = await db + .select({ id: workflow.id }) + .from(workflow) + .where(inArray(workflow.userId, freeUserIds)) + + if (workflowsQuery.length === 0) { + logger.info('No workflows found for free users') + return NextResponse.json({ message: 'No workflows found for cleanup' }) + } + + const workflowIds = workflowsQuery.map((w) => w.id) + + const oldLogs = await db + .select({ + id: workflowLogs.id, + workflowId: workflowLogs.workflowId, + executionId: workflowLogs.executionId, + level: workflowLogs.level, + message: workflowLogs.message, + duration: workflowLogs.duration, + trigger: workflowLogs.trigger, + createdAt: workflowLogs.createdAt, + metadata: workflowLogs.metadata, + }) .from(workflowLogs) - .innerJoin( - sql`workflow`, - sql`${workflowLogs.workflowId} = workflow.id AND workflow.user_id IN (${sql.join(freeUserIds)})` + .where( + and( + inArray(workflowLogs.workflowId, workflowIds), + lt(workflowLogs.createdAt, retentionDate) + ) ) - .groupBy(workflowLogs.workflowId) + .limit(BATCH_SIZE) - if (freeUserWorkflows.length === 0) { - logger.info('No free user workflows found for log cleanup') + logger.info(`Found ${oldLogs.length} logs older than ${retentionDate.toISOString()} to archive`) + + if (oldLogs.length === 0) { return NextResponse.json({ message: 'No logs to clean up' }) } - const workflowIds = freeUserWorkflows.map((w) => w.workflowId) + const results = { + total: oldLogs.length, + archived: 0, + archiveFailed: 0, + deleted: 0, + deleteFailed: 0, + } - const result = await db - .delete(workflowLogs) - .where( - sql`${workflowLogs.workflowId} IN (${sql.join(workflowIds)}) AND ${workflowLogs.createdAt} < ${retentionDate}` - ) - .returning({ id: workflowLogs.id }) + for (const log of oldLogs) { + const today = new Date().toISOString().split('T')[0] - logger.info(`Successfully cleaned up ${result.length} logs for free users`) + const logKey = `archived-logs/${today}/${log.id}.json` + const logData = JSON.stringify(log) + + try { + await s3Client.send( + new PutObjectCommand({ + Bucket: S3_CONFIG.bucket, + Key: logKey, + Body: logData, + ContentType: 'application/json', + Metadata: { + logId: String(log.id), + workflowId: String(log.workflowId), + archivedAt: new Date().toISOString(), + }, + }) + ) + + results.archived++ + + try { + const deleteResult = await db + .delete(workflowLogs) + .where(eq(workflowLogs.id, log.id)) + .returning({ id: workflowLogs.id }) + + if (deleteResult.length > 0) { + results.deleted++ + } else { + results.deleteFailed++ + logger.warn(`Failed to delete log ${log.id} after archiving: No rows deleted`) + } + } catch (deleteError) { + results.deleteFailed++ + logger.error(`Error deleting log ${log.id} after archiving:`, { deleteError }) + } + } catch (archiveError) { + results.archiveFailed++ + logger.error(`Failed to archive log ${log.id}:`, { archiveError }) + } + } return NextResponse.json({ - message: `Successfully cleaned up ${result.length} logs for free users`, - deletedCount: result.length, + message: `Successfully processed ${results.total} logs: archived ${results.archived}, deleted ${results.deleted}`, + results, }) } catch (error) { - logger.error('Error cleaning up logs:', { error }) - return NextResponse.json({ error: 'Failed to clean up logs' }, { status: 500 }) + logger.error('Error in log cleanup process:', { error }) + return NextResponse.json({ error: 'Failed to process log cleanup' }, { status: 500 }) } }