diff --git a/app/(auth)/verify/page.tsx b/app/(auth)/verify/page.tsx index edfd6050dc..7cca0ae033 100644 --- a/app/(auth)/verify/page.tsx +++ b/app/(auth)/verify/page.tsx @@ -13,8 +13,11 @@ import { } from '@/components/ui/card' import { InputOTP, InputOTPGroup, InputOTPSlot } from '@/components/ui/input-otp' import { client } from '@/lib/auth-client' +import { createLogger } from '@/lib/logs/console-logger' import { useNotificationStore } from '@/stores/notifications/store' +const logger = createLogger('VerifyPage') + // Extract the content into a separate component function VerifyContent() { const router = useRouter() @@ -46,7 +49,7 @@ function VerifyContent() { }) .then(() => {}) .catch((error) => { - console.error('Failed to send initial verification code:', error) + logger.error('Failed to send initial verification code:', error) addNotification?.( 'error', 'Failed to send verification code. Please use the resend button.', diff --git a/app/api/auth/oauth/connections/route.ts b/app/api/auth/oauth/connections/route.ts index df2686cf80..0b942033b0 100644 --- a/app/api/auth/oauth/connections/route.ts +++ b/app/api/auth/oauth/connections/route.ts @@ -2,10 +2,13 @@ import { NextRequest, NextResponse } from 'next/server' import { eq } from 'drizzle-orm' import { jwtDecode } from 'jwt-decode' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { OAuthService } from '@/lib/oauth' import { db } from '@/db' import { account } from '@/db/schema' +const logger = createLogger('OAuthConnectionsAPI') + interface GoogleIdToken { email?: string sub?: string @@ -18,12 +21,15 @@ const VALID_PROVIDERS = ['google', 'github', 'x'] * Get all OAuth connections for the current user */ export async function GET(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + try { // Get the session const session = await getSession() // Check if the user is authenticated if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthenticated request rejected`) return NextResponse.json({ error: 'User not authenticated' }, { status: 401 }) } @@ -47,7 +53,7 @@ export async function GET(request: NextRequest) { name = decoded.email } } catch (error) { - console.error('Error decoding ID token:', error) + logger.warn(`[${requestId}] Error decoding Google ID token`, { accountId: acc.id }) } } @@ -84,7 +90,7 @@ export async function GET(request: NextRequest) { return NextResponse.json({ connections }, { status: 200 }) } catch (error) { - console.error('Error fetching OAuth connections:', error) + logger.error(`[${requestId}] Error fetching OAuth connections`, error) return NextResponse.json({ error: 'Internal server error' }, { status: 500 }) } } diff --git a/app/api/auth/oauth/credentials/route.ts b/app/api/auth/oauth/credentials/route.ts index b60cd0a44e..23a411e2f4 100644 --- a/app/api/auth/oauth/credentials/route.ts +++ b/app/api/auth/oauth/credentials/route.ts @@ -2,11 +2,14 @@ import { NextRequest, NextResponse } from 'next/server' import { and, eq } from 'drizzle-orm' import { jwtDecode } from 'jwt-decode' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { parseProvider } from '@/lib/oauth' import { OAuthService } from '@/lib/oauth' import { db } from '@/db' import { account } from '@/db/schema' +const logger = createLogger('OAuthCredentialsAPI') + interface GoogleIdToken { email?: string sub?: string @@ -16,12 +19,15 @@ interface GoogleIdToken { * Get credentials for a specific provider */ export async function GET(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + try { // Get the session const session = await getSession() // Check if the user is authenticated if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthenticated credentials request rejected`) return NextResponse.json({ error: 'User not authenticated' }, { status: 401 }) } @@ -30,6 +36,7 @@ export async function GET(request: NextRequest) { const provider = searchParams.get('provider') as OAuthService | null if (!provider) { + logger.warn(`[${requestId}] Missing provider parameter`) return NextResponse.json({ error: 'Provider is required' }, { status: 400 }) } @@ -57,7 +64,7 @@ export async function GET(request: NextRequest) { name = decoded.email } } catch (error) { - console.error('Error decoding ID token:', error) + logger.warn(`[${requestId}] Error decoding Google ID token`, { accountId: acc.id }) } } @@ -73,7 +80,7 @@ export async function GET(request: NextRequest) { return NextResponse.json({ credentials }, { status: 200 }) } catch (error) { - console.error('Error fetching credentials:', error) + logger.error(`[${requestId}] Error fetching OAuth credentials`, error) return NextResponse.json({ error: 'Internal server error' }, { status: 500 }) } } diff --git a/app/api/auth/oauth/disconnect/route.ts b/app/api/auth/oauth/disconnect/route.ts index 0bc3f7f6d0..d87309c57b 100644 --- a/app/api/auth/oauth/disconnect/route.ts +++ b/app/api/auth/oauth/disconnect/route.ts @@ -1,19 +1,25 @@ import { NextRequest, NextResponse } from 'next/server' import { and, eq, like } from 'drizzle-orm' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { db } from '@/db' import { account } from '@/db/schema' +const logger = createLogger('OAuthDisconnectAPI') + /** * Disconnect an OAuth provider for the current user */ export async function POST(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + try { // Get the session const session = await getSession() // Check if the user is authenticated if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthenticated disconnect request rejected`) return NextResponse.json({ error: 'User not authenticated' }, { status: 401 }) } @@ -21,9 +27,15 @@ export async function POST(request: NextRequest) { const { provider, providerId } = await request.json() if (!provider) { + logger.warn(`[${requestId}] Missing provider in disconnect request`) return NextResponse.json({ error: 'Provider is required' }, { status: 400 }) } + logger.info(`[${requestId}] Processing OAuth disconnect request`, { + provider, + hasProviderId: !!providerId, + }) + // If a specific providerId is provided, delete only that account if (providerId) { await db @@ -39,7 +51,7 @@ export async function POST(request: NextRequest) { return NextResponse.json({ success: true }, { status: 200 }) } catch (error) { - console.error('Error disconnecting OAuth provider:', error) + logger.error(`[${requestId}] Error disconnecting OAuth provider`, error) return NextResponse.json({ error: 'Internal server error' }, { status: 500 }) } } diff --git a/app/api/auth/oauth/drive/file/route.ts b/app/api/auth/oauth/drive/file/route.ts index 8802726178..94150b305c 100644 --- a/app/api/auth/oauth/drive/file/route.ts +++ b/app/api/auth/oauth/drive/file/route.ts @@ -1,20 +1,26 @@ import { NextRequest, NextResponse } from 'next/server' import { eq } from 'drizzle-orm' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { refreshOAuthToken } from '@/lib/oauth' import { db } from '@/db' import { account } from '@/db/schema' +const logger = createLogger('GoogleDriveFileAPI') + /** * Get a single file from Google Drive by ID */ export async function GET(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) // Short request ID for correlation + try { // Get the session const session = await getSession() // Check if the user is authenticated if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthenticated request rejected`) return NextResponse.json({ error: 'User not authenticated' }, { status: 401 }) } @@ -23,18 +29,24 @@ export async function GET(request: NextRequest) { const credentialId = searchParams.get('credentialId') const fileId = searchParams.get('fileId') - if (!credentialId) { - return NextResponse.json({ error: 'Credential ID is required' }, { status: 400 }) - } - - if (!fileId) { - return NextResponse.json({ error: 'File ID is required' }, { status: 400 }) + if (!credentialId || !fileId) { + logger.warn(`[${requestId}] Missing required parameters`, { + credentialId: !!credentialId, + fileId: !!fileId, + }) + return NextResponse.json( + { + error: !credentialId ? 'Credential ID is required' : 'File ID is required', + }, + { status: 400 } + ) } // Get the credential from the database const credentials = await db.select().from(account).where(eq(account.id, credentialId)).limit(1) if (!credentials.length) { + logger.warn(`[${requestId}] Credential not found`, { credentialId }) return NextResponse.json({ error: 'Credential not found' }, { status: 404 }) } @@ -42,11 +54,13 @@ export async function GET(request: NextRequest) { // Check if the credential belongs to the user if (credential.userId !== session.user.id) { + logger.warn(`[${requestId}] Unauthorized credential access attempt`) return NextResponse.json({ error: 'Unauthorized' }, { status: 403 }) } // Check if the access token is valid if (!credential.accessToken) { + logger.warn(`[${requestId}] No access token available for credential`) return NextResponse.json({ error: 'No access token available' }, { status: 400 }) } @@ -68,7 +82,7 @@ export async function GET(request: NextRequest) { // If unauthorized, try to refresh the token if (response.status === 401 && credential.refreshToken) { - console.log('Access token expired, attempting to refresh...') + logger.info(`[${requestId}] Access token expired, attempting to refresh`) try { // Refresh the token using the centralized utility @@ -78,6 +92,8 @@ export async function GET(request: NextRequest) { ) if (refreshedToken) { + logger.info(`[${requestId}] Token refreshed successfully`) + // Update the token in the database await db .update(account) @@ -92,7 +108,7 @@ export async function GET(request: NextRequest) { response = await fetchFileWithToken(refreshedToken) } } catch (refreshError) { - console.error('Error refreshing token:', refreshError) + logger.error(`[${requestId}] Error refreshing token`, refreshError) return NextResponse.json({ error: 'Failed to refresh access token' }, { status: 401 }) } } @@ -100,6 +116,10 @@ export async function GET(request: NextRequest) { // Handle response if (!response.ok) { const error = await response.json().catch(() => ({ error: { message: 'Unknown error' } })) + logger.error(`[${requestId}] Google Drive API error`, { + status: response.status, + fileId, + }) return NextResponse.json( { error: error.error?.message || 'Failed to fetch file from Google Drive' }, { status: response.status } @@ -107,9 +127,10 @@ export async function GET(request: NextRequest) { } const file = await response.json() + logger.info(`[${requestId}] Successfully retrieved file from Google Drive`, { fileId }) return NextResponse.json({ file }, { status: 200 }) } catch (error) { - console.error('Error fetching file from Google Drive:', error) + logger.error(`[${requestId}] Error fetching file from Google Drive`, error) return NextResponse.json({ error: 'Internal server error' }, { status: 500 }) } } diff --git a/app/api/auth/oauth/drive/files/route.ts b/app/api/auth/oauth/drive/files/route.ts index 391088998d..6b5c9454f8 100644 --- a/app/api/auth/oauth/drive/files/route.ts +++ b/app/api/auth/oauth/drive/files/route.ts @@ -1,20 +1,27 @@ import { NextRequest, NextResponse } from 'next/server' import { eq } from 'drizzle-orm' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { refreshOAuthToken } from '@/lib/oauth' import { db } from '@/db' import { account } from '@/db/schema' +const logger = createLogger('GoogleDriveFilesAPI') + /** * Get files from Google Drive */ export async function GET(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) // Generate a short request ID for correlation + logger.info(`[${requestId}] Google Drive files request received`) + try { // Get the session const session = await getSession() // Check if the user is authenticated if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthenticated request rejected`) return NextResponse.json({ error: 'User not authenticated' }, { status: 401 }) } @@ -25,6 +32,7 @@ export async function GET(request: NextRequest) { const query = searchParams.get('query') || '' if (!credentialId) { + logger.warn(`[${requestId}] Missing credential ID`) return NextResponse.json({ error: 'Credential ID is required' }, { status: 400 }) } @@ -32,6 +40,7 @@ export async function GET(request: NextRequest) { const credentials = await db.select().from(account).where(eq(account.id, credentialId)).limit(1) if (!credentials.length) { + logger.warn(`[${requestId}] Credential not found`, { credentialId }) return NextResponse.json({ error: 'Credential not found' }, { status: 404 }) } @@ -39,11 +48,16 @@ export async function GET(request: NextRequest) { // Check if the credential belongs to the user if (credential.userId !== session.user.id) { + logger.warn(`[${requestId}] Unauthorized credential access attempt`, { + credentialUserId: credential.userId, + requestUserId: session.user.id, + }) return NextResponse.json({ error: 'Unauthorized' }, { status: 403 }) } // Check if the access token is valid if (!credential.accessToken) { + logger.warn(`[${requestId}] No access token available for credential`, { credentialId }) return NextResponse.json({ error: 'No access token available' }, { status: 400 }) } @@ -88,7 +102,7 @@ export async function GET(request: NextRequest) { // If unauthorized, try to refresh the token if (response.status === 401 && credential.refreshToken) { - console.log('Access token expired, attempting to refresh...') + logger.info(`[${requestId}] Access token expired, attempting to refresh`) try { // Refresh the token using the centralized utility @@ -98,6 +112,8 @@ export async function GET(request: NextRequest) { ) if (refreshedToken) { + logger.info(`[${requestId}] Token refreshed successfully`) + // Update the token in the database await db .update(account) @@ -112,13 +128,17 @@ export async function GET(request: NextRequest) { response = await fetchFilesWithToken(refreshedToken) } } catch (refreshError) { - console.error('Error refreshing token:', refreshError) + logger.error(`[${requestId}] Error refreshing token`, refreshError) return NextResponse.json({ error: 'Failed to refresh access token' }, { status: 401 }) } } if (!response.ok) { const error = await response.json().catch(() => ({ error: { message: 'Unknown error' } })) + logger.error(`[${requestId}] Google Drive API error`, { + status: response.status, + error: error.error?.message || 'Failed to fetch files from Google Drive', + }) return NextResponse.json( { error: error.error?.message || 'Failed to fetch files from Google Drive' }, { status: response.status } @@ -138,7 +158,7 @@ export async function GET(request: NextRequest) { return NextResponse.json({ files }, { status: 200 }) } catch (error) { - console.error('Error fetching files from Google Drive:', error) + logger.error(`[${requestId}] Error fetching files from Google Drive`, error) return NextResponse.json({ error: 'Internal server error' }, { status: 500 }) } } diff --git a/app/api/auth/oauth/token/route.ts b/app/api/auth/oauth/token/route.ts index 04cb40c5f1..7acf3741ad 100644 --- a/app/api/auth/oauth/token/route.ts +++ b/app/api/auth/oauth/token/route.ts @@ -1,22 +1,28 @@ import { NextRequest, NextResponse } from 'next/server' import { and, eq } from 'drizzle-orm' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { refreshOAuthToken } from '@/lib/oauth' import { db } from '@/db' import { account, workflow } from '@/db/schema' +const logger = createLogger('OAuthTokenAPI') + /** * Get an access token for a specific credential * Supports both session-based authentication (for client-side requests) * and workflow-based authentication (for server-side requests) */ export async function POST(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + try { // Parse request body const body = await request.json() const { credentialId, workflowId } = body if (!credentialId) { + logger.warn(`[${requestId}] Credential ID is required`) return NextResponse.json({ error: 'Credential ID is required' }, { status: 400 }) } @@ -33,6 +39,7 @@ export async function POST(request: NextRequest) { .limit(1) if (!workflows.length) { + logger.warn(`[${requestId}] Workflow not found`) return NextResponse.json({ error: 'Workflow not found' }, { status: 404 }) } @@ -43,6 +50,7 @@ export async function POST(request: NextRequest) { // Check if the user is authenticated if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthenticated token request rejected`) return NextResponse.json({ error: 'User not authenticated' }, { status: 401 }) } @@ -57,6 +65,7 @@ export async function POST(request: NextRequest) { .limit(1) if (!credentials.length) { + logger.warn(`[${requestId}] Credential not found`) return NextResponse.json({ error: 'Credential not found' }, { status: 404 }) } @@ -87,16 +96,18 @@ export async function POST(request: NextRequest) { }) .where(eq(account.id, credentialId)) + logger.info(`[${requestId}] Successfully refreshed access token`) return NextResponse.json({ accessToken: refreshedToken }, { status: 200 }) } catch (error) { - console.error('Error refreshing token:', error) + logger.error(`[${requestId}] Error refreshing token`, error) return NextResponse.json({ error: 'Failed to refresh access token' }, { status: 500 }) } } + logger.info(`[${requestId}] Access token is valid`) return NextResponse.json({ accessToken: credential.accessToken }, { status: 200 }) } catch (error) { - console.error('Error getting access token:', error) + logger.error(`[${requestId}] Error getting access token`, error) return NextResponse.json({ error: 'Internal server error' }, { status: 500 }) } } diff --git a/app/api/db/environment/route.ts b/app/api/db/environment/route.ts index 10bed29ffb..6e4914ff0c 100644 --- a/app/api/db/environment/route.ts +++ b/app/api/db/environment/route.ts @@ -2,71 +2,88 @@ import { NextRequest, NextResponse } from 'next/server' import { eq } from 'drizzle-orm' import { z } from 'zod' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { decryptSecret, encryptSecret } from '@/lib/utils' import { EnvironmentVariable } from '@/stores/settings/environment/types' import { db } from '@/db' import { environment } from '@/db/schema' +const logger = createLogger('EnvironmentAPI') + // Schema for environment variable updates const EnvVarSchema = z.object({ variables: z.record(z.string()), }) export async function POST(req: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + try { const session = await getSession() if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized environment variables update attempt`) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) } const body = await req.json() - const { variables } = EnvVarSchema.parse(body) - // Encrypt all variables - const encryptedVariables = await Object.entries(variables).reduce( - async (accPromise, [key, value]) => { - const acc = await accPromise - const { encrypted } = await encryptSecret(value) - return { ...acc, [key]: encrypted } - }, - Promise.resolve({}) - ) + try { + const { variables } = EnvVarSchema.parse(body) - // Replace all environment variables for user - await db - .insert(environment) - .values({ - id: crypto.randomUUID(), - userId: session.user.id, - variables: encryptedVariables, - updatedAt: new Date(), - }) - .onConflictDoUpdate({ - target: [environment.userId], - set: { + // Encrypt all variables + const encryptedVariables = await Object.entries(variables).reduce( + async (accPromise, [key, value]) => { + const acc = await accPromise + const { encrypted } = await encryptSecret(value) + return { ...acc, [key]: encrypted } + }, + Promise.resolve({}) + ) + + // Replace all environment variables for user + await db + .insert(environment) + .values({ + id: crypto.randomUUID(), + userId: session.user.id, variables: encryptedVariables, updatedAt: new Date(), - }, - }) + }) + .onConflictDoUpdate({ + target: [environment.userId], + set: { + variables: encryptedVariables, + updatedAt: new Date(), + }, + }) - return NextResponse.json({ success: true }) - } catch (error) { - console.error('Error updating environment variables:', error) - if (error instanceof z.ZodError) { - return NextResponse.json( - { error: 'Invalid request data', details: error.errors }, - { status: 400 } - ) + return NextResponse.json({ success: true }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid environment variables data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError } + } catch (error) { + logger.error(`[${requestId}] Error updating environment variables`, error) return NextResponse.json({ error: 'Failed to update environment variables' }, { status: 500 }) } } export async function GET(request: Request) { + const requestId = crypto.randomUUID().slice(0, 8) + try { // Get the session directly in the API route const session = await getSession() if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized environment variables access attempt`) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) } @@ -92,7 +109,7 @@ export async function GET(request: Request) { const { decrypted } = await decryptSecret(encryptedValue) decryptedVariables[key] = { key, value: decrypted } } catch (error) { - console.error(`Error decrypting variable ${key}:`, error) + logger.error(`[${requestId}] Error decrypting variable ${key}`, error) // If decryption fails, provide a placeholder decryptedVariables[key] = { key, value: '' } } @@ -100,7 +117,7 @@ export async function GET(request: Request) { return NextResponse.json({ data: decryptedVariables }, { status: 200 }) } catch (error: any) { - console.error('Environment fetch error:', error) + logger.error(`[${requestId}] Environment fetch error`, error) return NextResponse.json({ error: error.message }, { status: 500 }) } } diff --git a/app/api/db/settings/route.ts b/app/api/db/settings/route.ts index e793ea3f9b..f74796df38 100644 --- a/app/api/db/settings/route.ts +++ b/app/api/db/settings/route.ts @@ -2,49 +2,71 @@ import { NextResponse } from 'next/server' import { eq } from 'drizzle-orm' import { nanoid } from 'nanoid' import { z } from 'zod' +import { createLogger } from '@/lib/logs/console-logger' import { db } from '@/db' import { settings } from '@/db/schema' +const logger = createLogger('SettingsAPI') + const SettingsSchema = z.object({ userId: z.string(), isAutoConnectEnabled: z.boolean().default(true), }) export async function POST(request: Request) { + const requestId = crypto.randomUUID().slice(0, 8) + try { const body = await request.json() - const { userId, isAutoConnectEnabled } = SettingsSchema.parse(body) - // Store the settings - await db - .insert(settings) - .values({ - id: nanoid(), - userId, - general: { isAutoConnectEnabled }, - updatedAt: new Date(), - }) - .onConflictDoUpdate({ - target: [settings.userId], - set: { + try { + const { userId, isAutoConnectEnabled } = SettingsSchema.parse(body) + + // Store the settings + await db + .insert(settings) + .values({ + id: nanoid(), + userId, general: { isAutoConnectEnabled }, updatedAt: new Date(), - }, - }) + }) + .onConflictDoUpdate({ + target: [settings.userId], + set: { + general: { isAutoConnectEnabled }, + updatedAt: new Date(), + }, + }) - return NextResponse.json({ success: true }, { status: 200 }) + return NextResponse.json({ success: true }, { status: 200 }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid settings data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid settings data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError + } } catch (error: any) { - console.error('Settings update error:', error) + logger.error(`[${requestId}] Settings update error`, error) return NextResponse.json({ error: error.message }, { status: 500 }) } } export async function GET(request: Request) { + const requestId = crypto.randomUUID().slice(0, 8) + try { const { searchParams } = new URL(request.url) const userId = searchParams.get('userId') if (!userId) { + logger.warn(`[${requestId}] Missing userId parameter`) return NextResponse.json({ error: 'userId is required' }, { status: 400 }) } @@ -71,7 +93,7 @@ export async function GET(request: Request) { { status: 200 } ) } catch (error: any) { - console.error('Settings fetch error:', error) + logger.error(`[${requestId}] Settings fetch error`, error) return NextResponse.json({ error: error.message }, { status: 500 }) } } diff --git a/app/api/db/workflow-logs/route.ts b/app/api/db/workflow-logs/route.ts index 42f9b421ce..46eeb78a46 100644 --- a/app/api/db/workflow-logs/route.ts +++ b/app/api/db/workflow-logs/route.ts @@ -2,9 +2,13 @@ import { NextRequest, NextResponse } from 'next/server' import { SQL, and, eq, gte, lte, or, sql } from 'drizzle-orm' import { z } from 'zod' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { db } from '@/db' import { workflow, workflowLogs } from '@/db/schema' +// Create a logger for this module +const logger = createLogger('WorkflowLogsAPI') + // No cache export const dynamic = 'force-dynamic' export const revalidate = 0 @@ -22,113 +26,133 @@ const QueryParamsSchema = z.object({ }) export async function GET(request: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + try { // Get the session directly in the API route const session = await getSession() if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized workflow logs access attempt`) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) } const userId = session.user.id - // Parse query parameters - const { searchParams } = new URL(request.url) - const params = QueryParamsSchema.parse(Object.fromEntries(searchParams.entries())) + try { + // Parse query parameters + const { searchParams } = new URL(request.url) + const params = QueryParamsSchema.parse(Object.fromEntries(searchParams.entries())) - // Start building the query to get all workflows for the user - const userWorkflows = await db - .select({ id: workflow.id }) - .from(workflow) - .where(eq(workflow.userId, userId)) + // Start building the query to get all workflows for the user + const userWorkflows = await db + .select({ id: workflow.id }) + .from(workflow) + .where(eq(workflow.userId, userId)) - const workflowIds = userWorkflows.map((w) => w.id) + const workflowIds = userWorkflows.map((w) => w.id) - if (workflowIds.length === 0) { - return NextResponse.json({ data: [], total: 0 }, { status: 200 }) - } - - // Build the conditions for the query - let conditions: SQL | undefined - - // Start with the first workflowId - conditions = eq(workflowLogs.workflowId, workflowIds[0]) - - // Add additional workflowIds if there are more than one - if (workflowIds.length > 1) { - const workflowConditions = workflowIds.map((id) => eq(workflowLogs.workflowId, id)) - conditions = or(...workflowConditions) - } - - // Apply additional filters if provided - if (params.level) { - conditions = and(conditions, eq(workflowLogs.level, params.level)) - } - - if (params.workflowId) { - // Ensure the requested workflow belongs to the user - if (workflowIds.includes(params.workflowId)) { - conditions = and(conditions, eq(workflowLogs.workflowId, params.workflowId)) - } else { - return NextResponse.json({ error: 'Unauthorized access to workflow' }, { status: 403 }) - } - } - - if (params.startDate) { - const startDate = new Date(params.startDate) - conditions = and(conditions, gte(workflowLogs.createdAt, startDate)) - } - - if (params.endDate) { - const endDate = new Date(params.endDate) - conditions = and(conditions, lte(workflowLogs.createdAt, endDate)) - } - - // Execute the query with all conditions - const logs = await db - .select() - .from(workflowLogs) - .where(conditions) - .orderBy(sql`${workflowLogs.createdAt} DESC`) - .limit(params.limit) - .offset(params.offset) - - // Get total count for pagination - const countResult = await db - .select({ count: sql`count(*)` }) - .from(workflowLogs) - .where(conditions) - - const count = countResult[0]?.count || 0 - - // If includeWorkflow is true, fetch the associated workflow data - if (params.includeWorkflow === 'true' && logs.length > 0) { - // Get unique workflow IDs from logs - const uniqueWorkflowIds = [...new Set(logs.map((log) => log.workflowId))] - - // Create conditions for workflow query - let workflowConditions: SQL | undefined - - if (uniqueWorkflowIds.length === 1) { - workflowConditions = eq(workflow.id, uniqueWorkflowIds[0]) - } else { - workflowConditions = or(...uniqueWorkflowIds.map((id) => eq(workflow.id, id))) + if (workflowIds.length === 0) { + return NextResponse.json({ data: [], total: 0 }, { status: 200 }) } - // Fetch workflows - const workflowData = await db.select().from(workflow).where(workflowConditions) + // Build the conditions for the query + let conditions: SQL | undefined - // Create a map of workflow data for easy lookup - const workflowMap = new Map(workflowData.map((w) => [w.id, w])) + // Start with the first workflowId + conditions = eq(workflowLogs.workflowId, workflowIds[0]) - // Attach workflow data to each log - const logsWithWorkflow = logs.map((log) => ({ - ...log, - workflow: workflowMap.get(log.workflowId) || null, - })) + // Add additional workflowIds if there are more than one + if (workflowIds.length > 1) { + const workflowConditions = workflowIds.map((id) => eq(workflowLogs.workflowId, id)) + conditions = or(...workflowConditions) + } + // Apply additional filters if provided + if (params.level) { + conditions = and(conditions, eq(workflowLogs.level, params.level)) + } + + if (params.workflowId) { + // Ensure the requested workflow belongs to the user + if (workflowIds.includes(params.workflowId)) { + conditions = and(conditions, eq(workflowLogs.workflowId, params.workflowId)) + } else { + logger.warn(`[${requestId}] Unauthorized access to workflow logs`, { + requestedWorkflowId: params.workflowId, + }) + return NextResponse.json({ error: 'Unauthorized access to workflow' }, { status: 403 }) + } + } + + if (params.startDate) { + const startDate = new Date(params.startDate) + conditions = and(conditions, gte(workflowLogs.createdAt, startDate)) + } + + if (params.endDate) { + const endDate = new Date(params.endDate) + conditions = and(conditions, lte(workflowLogs.createdAt, endDate)) + } + + // Execute the query with all conditions + const logs = await db + .select() + .from(workflowLogs) + .where(conditions) + .orderBy(sql`${workflowLogs.createdAt} DESC`) + .limit(params.limit) + .offset(params.offset) + + // Get total count for pagination + const countResult = await db + .select({ count: sql`count(*)` }) + .from(workflowLogs) + .where(conditions) + + const count = countResult[0]?.count || 0 + + // If includeWorkflow is true, fetch the associated workflow data + if (params.includeWorkflow === 'true' && logs.length > 0) { + // Get unique workflow IDs from logs + const uniqueWorkflowIds = [...new Set(logs.map((log) => log.workflowId))] + + // Create conditions for workflow query + let workflowConditions: SQL | undefined + + if (uniqueWorkflowIds.length === 1) { + workflowConditions = eq(workflow.id, uniqueWorkflowIds[0]) + } else { + workflowConditions = or(...uniqueWorkflowIds.map((id) => eq(workflow.id, id))) + } + + // Fetch workflows + const workflowData = await db.select().from(workflow).where(workflowConditions) + + // Create a map of workflow data for easy lookup + const workflowMap = new Map(workflowData.map((w) => [w.id, w])) + + // Attach workflow data to each log + const logsWithWorkflow = logs.map((log) => ({ + ...log, + workflow: workflowMap.get(log.workflowId) || null, + })) + + return NextResponse.json( + { + data: logsWithWorkflow, + total: Number(count), + page: Math.floor(params.offset / params.limit) + 1, + pageSize: params.limit, + totalPages: Math.ceil(Number(count) / params.limit), + }, + { status: 200 } + ) + } + + // Return logs without workflow data return NextResponse.json( { - data: logsWithWorkflow, + data: logs, total: Number(count), page: Math.floor(params.offset / params.limit) + 1, pageSize: params.limit, @@ -136,27 +160,20 @@ export async function GET(request: NextRequest) { }, { status: 200 } ) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid workflow logs request parameters`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request parameters', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError } - - // Return logs without workflow data - return NextResponse.json( - { - data: logs, - total: Number(count), - page: Math.floor(params.offset / params.limit) + 1, - pageSize: params.limit, - totalPages: Math.ceil(Number(count) / params.limit), - }, - { status: 200 } - ) } catch (error: any) { - console.error('Workflow logs fetch error:', error) - if (error instanceof z.ZodError) { - return NextResponse.json( - { error: 'Invalid request parameters', details: error.errors }, - { status: 400 } - ) - } + logger.error(`[${requestId}] Workflow logs fetch error`, error) return NextResponse.json({ error: error.message }, { status: 500 }) } } diff --git a/app/api/db/workflow/route.ts b/app/api/db/workflow/route.ts index 845baf075f..f47f1f2a74 100644 --- a/app/api/db/workflow/route.ts +++ b/app/api/db/workflow/route.ts @@ -2,9 +2,12 @@ import { NextRequest, NextResponse } from 'next/server' import { eq } from 'drizzle-orm' import { z } from 'zod' import { getSession } from '@/lib/auth' +import { createLogger } from '@/lib/logs/console-logger' import { db } from '@/db' import { workflow } from '@/db/schema' +const logger = createLogger('WorkflowAPI') + // Schema for workflow data const WorkflowStateSchema = z.object({ blocks: z.record(z.any()), @@ -31,10 +34,13 @@ const SyncPayloadSchema = z.object({ }) export async function GET(request: Request) { + const requestId = crypto.randomUUID().slice(0, 8) + try { // Get the session directly in the API route const session = await getSession() if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized workflow access attempt`) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) } @@ -46,96 +52,110 @@ export async function GET(request: Request) { // Return the workflows return NextResponse.json({ data: workflows }, { status: 200 }) } catch (error: any) { - console.error('Workflow fetch error:', error) + logger.error(`[${requestId}] Workflow fetch error`, error) return NextResponse.json({ error: error.message }, { status: 500 }) } } export async function POST(req: NextRequest) { + const requestId = crypto.randomUUID().slice(0, 8) + try { const session = await getSession() if (!session?.user?.id) { + logger.warn(`[${requestId}] Unauthorized workflow sync attempt`) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) } const body = await req.json() - const { workflows: clientWorkflows } = SyncPayloadSchema.parse(body) - // Get all workflows for the user from the database - const dbWorkflows = await db.select().from(workflow).where(eq(workflow.userId, session.user.id)) + try { + const { workflows: clientWorkflows } = SyncPayloadSchema.parse(body) - const now = new Date() - const operations: Promise[] = [] + // Get all workflows for the user from the database + const dbWorkflows = await db + .select() + .from(workflow) + .where(eq(workflow.userId, session.user.id)) - // Create a map of DB workflows for easier lookup - const dbWorkflowMap = new Map(dbWorkflows.map((w) => [w.id, w])) - const processedIds = new Set() + const now = new Date() + const operations: Promise[] = [] - // Process client workflows - for (const [id, clientWorkflow] of Object.entries(clientWorkflows)) { - processedIds.add(id) - const dbWorkflow = dbWorkflowMap.get(id) + // Create a map of DB workflows for easier lookup + const dbWorkflowMap = new Map(dbWorkflows.map((w) => [w.id, w])) + const processedIds = new Set() - if (!dbWorkflow) { - // New workflow - create - operations.push( - db.insert(workflow).values({ - id: clientWorkflow.id, - userId: session.user.id, - name: clientWorkflow.name, - description: clientWorkflow.description, - color: clientWorkflow.color, - state: clientWorkflow.state, - lastSynced: now, - createdAt: now, - updatedAt: now, - }) - ) - } else { - // Existing workflow - update if needed - const needsUpdate = - JSON.stringify(dbWorkflow.state) !== JSON.stringify(clientWorkflow.state) || - dbWorkflow.name !== clientWorkflow.name || - dbWorkflow.description !== clientWorkflow.description || - dbWorkflow.color !== clientWorkflow.color + // Process client workflows + for (const [id, clientWorkflow] of Object.entries(clientWorkflows)) { + processedIds.add(id) + const dbWorkflow = dbWorkflowMap.get(id) - if (needsUpdate) { + if (!dbWorkflow) { + // New workflow - create operations.push( - db - .update(workflow) - .set({ - name: clientWorkflow.name, - description: clientWorkflow.description, - color: clientWorkflow.color, - state: clientWorkflow.state, - lastSynced: now, - updatedAt: now, - }) - .where(eq(workflow.id, id)) + db.insert(workflow).values({ + id: clientWorkflow.id, + userId: session.user.id, + name: clientWorkflow.name, + description: clientWorkflow.description, + color: clientWorkflow.color, + state: clientWorkflow.state, + lastSynced: now, + createdAt: now, + updatedAt: now, + }) ) + } else { + // Existing workflow - update if needed + const needsUpdate = + JSON.stringify(dbWorkflow.state) !== JSON.stringify(clientWorkflow.state) || + dbWorkflow.name !== clientWorkflow.name || + dbWorkflow.description !== clientWorkflow.description || + dbWorkflow.color !== clientWorkflow.color + + if (needsUpdate) { + operations.push( + db + .update(workflow) + .set({ + name: clientWorkflow.name, + description: clientWorkflow.description, + color: clientWorkflow.color, + state: clientWorkflow.state, + lastSynced: now, + updatedAt: now, + }) + .where(eq(workflow.id, id)) + ) + } } } - } - // Handle deletions - workflows in DB but not in client - for (const dbWorkflow of dbWorkflows) { - if (!processedIds.has(dbWorkflow.id)) { - operations.push(db.delete(workflow).where(eq(workflow.id, dbWorkflow.id))) + // Handle deletions - workflows in DB but not in client + for (const dbWorkflow of dbWorkflows) { + if (!processedIds.has(dbWorkflow.id)) { + operations.push(db.delete(workflow).where(eq(workflow.id, dbWorkflow.id))) + } } + + // Execute all operations in parallel + await Promise.all(operations) + + return NextResponse.json({ success: true }) + } catch (validationError) { + if (validationError instanceof z.ZodError) { + logger.warn(`[${requestId}] Invalid workflow data`, { + errors: validationError.errors, + }) + return NextResponse.json( + { error: 'Invalid request data', details: validationError.errors }, + { status: 400 } + ) + } + throw validationError } - - // Execute all operations in parallel - await Promise.all(operations) - - return NextResponse.json({ success: true }) } catch (error) { - console.error('Workflow sync error:', error) - if (error instanceof z.ZodError) { - return NextResponse.json( - { error: 'Invalid request data', details: error.errors }, - { status: 400 } - ) - } + logger.error(`[${requestId}] Workflow sync error`, error) return NextResponse.json({ error: 'Workflow sync failed' }, { status: 500 }) } } diff --git a/app/api/scheduled/execute/route.ts b/app/api/scheduled/execute/route.ts index a982e2b760..d4c0eb77ba 100644 --- a/app/api/scheduled/execute/route.ts +++ b/app/api/scheduled/execute/route.ts @@ -3,7 +3,7 @@ import { Cron } from 'croner' import { eq, lte } from 'drizzle-orm' import { v4 as uuidv4 } from 'uuid' import { z } from 'zod' -import { persistExecutionError, persistExecutionLogs } from '@/lib/logging' +import { persistExecutionError, persistExecutionLogs } from '@/lib/logs/execution-logger' import { decryptSecret } from '@/lib/utils' import { mergeSubblockState } from '@/stores/workflows/utils' import { BlockState, WorkflowState } from '@/stores/workflows/workflow/types' diff --git a/app/api/webhooks/trigger/[path]/route.ts b/app/api/webhooks/trigger/[path]/route.ts index c7f8a36200..e8bfdecfe3 100644 --- a/app/api/webhooks/trigger/[path]/route.ts +++ b/app/api/webhooks/trigger/[path]/route.ts @@ -1,7 +1,7 @@ import { NextRequest, NextResponse } from 'next/server' import { and, eq } from 'drizzle-orm' import { v4 as uuidv4 } from 'uuid' -import { persistExecutionError, persistExecutionLogs } from '@/lib/logging' +import { persistExecutionError, persistExecutionLogs } from '@/lib/logs/execution-logger' import { closeRedisConnection, hasProcessedMessage, markMessageAsProcessed } from '@/lib/redis' import { decryptSecret } from '@/lib/utils' import { mergeSubblockState, mergeSubblockStateAsync } from '@/stores/workflows/utils' diff --git a/app/api/workflow/[id]/execute/route.ts b/app/api/workflow/[id]/execute/route.ts index 9f77f1b817..d4934d2808 100644 --- a/app/api/workflow/[id]/execute/route.ts +++ b/app/api/workflow/[id]/execute/route.ts @@ -2,7 +2,7 @@ import { NextRequest } from 'next/server' import { eq } from 'drizzle-orm' import { v4 as uuidv4 } from 'uuid' import { z } from 'zod' -import { persistExecutionError, persistExecutionLogs } from '@/lib/logging' +import { persistExecutionError, persistExecutionLogs } from '@/lib/logs/execution-logger' import { decryptSecret } from '@/lib/utils' import { mergeSubblockState } from '@/stores/workflows/utils' import { WorkflowState } from '@/stores/workflows/workflow/types' diff --git a/app/api/workflow/[id]/log/route.ts b/app/api/workflow/[id]/log/route.ts index fdaa721c60..bd18f60d92 100644 --- a/app/api/workflow/[id]/log/route.ts +++ b/app/api/workflow/[id]/log/route.ts @@ -1,6 +1,6 @@ import { NextRequest } from 'next/server' import { v4 as uuidv4 } from 'uuid' -import { persistLog } from '@/lib/logging' +import { persistLog } from '@/lib/logs/execution-logger' import { validateWorkflowAccess } from '../../middleware' import { createErrorResponse, createSuccessResponse } from '../../utils' diff --git a/app/layout.tsx b/app/layout.tsx index 6bb1723603..276f2341ac 100644 --- a/app/layout.tsx +++ b/app/layout.tsx @@ -1,7 +1,10 @@ import type { Metadata, Viewport } from 'next' +import { createLogger } from '@/lib/logs/console-logger' import './globals.css' import { ZoomPrevention } from './zoom-prevention' +const logger = createLogger('RootLayout') + // Add browser extension attributes that we want to ignore const BROWSER_EXTENSION_ATTRIBUTES = [ 'data-new-gr-c-s-check-loaded', @@ -24,14 +27,12 @@ if (typeof window !== 'undefined') { ) if (!isExtensionError) { - // Log non-extension hydration errors with more detail - console.group('Hydration Error') - console.log('Error Details:', args) - console.log( - 'Component Stack:', - args.find((arg) => typeof arg === 'string' && arg.includes('component stack')) - ) - console.groupEnd() + logger.error('Hydration Error', { + details: args, + componentStack: args.find( + (arg) => typeof arg === 'string' && arg.includes('component stack') + ), + }) } } originalError.apply(console, args)