feat(logger): added new logger to half of the api routes

This commit is contained in:
Waleed Latif
2025-03-12 00:13:47 -07:00
parent f64715aa16
commit 3e02c81e2a
16 changed files with 412 additions and 255 deletions
+4 -1
View File
@@ -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.',
+8 -2
View File
@@ -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 })
}
}
+9 -2
View File
@@ -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 })
}
}
+13 -1
View File
@@ -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 })
}
}
+30 -9
View File
@@ -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 })
}
}
+23 -3
View File
@@ -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 })
}
}
+13 -2
View File
@@ -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 })
}
}
+51 -34
View File
@@ -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 })
}
}
+40 -18
View File
@@ -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 })
}
}
+124 -107
View File
@@ -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<unknown> | 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<number>`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<unknown> | 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<unknown> | 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<number>`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<unknown> | 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 })
}
}
+84 -64
View File
@@ -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<any>[] = []
// 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<string>()
const now = new Date()
const operations: Promise<any>[] = []
// 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<string>()
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 })
}
}
+1 -1
View File
@@ -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'
+1 -1
View File
@@ -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'
+1 -1
View File
@@ -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'
+1 -1
View File
@@ -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'
+9 -8
View File
@@ -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)