Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
94 changes: 65 additions & 29 deletions Realism/app/api/jobs/webhook/route.ts
Original file line number Diff line number Diff line change
@@ -1,29 +1,65 @@
import { NextRequest, NextResponse } from 'next/server'
import { inngest } from '@/lib/inngest'

export async function POST(req: NextRequest) {
let jobId: string | undefined

try {
const body = await req.json()
jobId = body.jobId
} catch {
return NextResponse.json({ error: 'Invalid body' }, { status: 400 })
}

if (!jobId) {
return NextResponse.json({ error: 'jobId required' }, { status: 400 })
}

try {
await inngest.send({
name: 'job/webhook',
data: { jobId },
})
} catch (err) {
console.error('[webhook] Failed to send Inngest event:', err)
}

// Always 200 to prevent QStash retries
return NextResponse.json({ ok: true })
}
import { NextRequest, NextResponse } from 'next/server'
import { inngest } from '@/lib/inngest'
import { Receiver } from '@upstash/qstash'
import { checkAndLockMessage, unlockMessage } from '@/lib/jobs'

// SECURITY FIX: Initialize QStash receiver for signature verification
const receiver = new Receiver({
currentSigningKey: process.env.QSTASH_CURRENT_SIGNING_KEY || '',
nextSigningKey: process.env.QSTASH_NEXT_SIGNING_KEY || '',
})

export async function POST(req: NextRequest) {
// SECURITY FIX: Verify provider signatures to reject unsigned/unauthenticated requests.
const signature = req.headers.get('Upstash-Signature')
const messageId = req.headers.get('Upstash-Message-Id')

if (!signature || !messageId) {
return NextResponse.json({ error: 'Missing QStash signature or message ID' }, { status: 401 })
}

const bodyText = await req.text()

try {
const isValid = await receiver.verify({ signature, body: bodyText })
if (!isValid) throw new Error("Invalid signature")
} catch (err) {
return NextResponse.json({ error: 'Unauthorized: Invalid signature' }, { status: 401 })
}

let jobId: string | undefined

try {
const body = JSON.parse(bodyText)
jobId = body.jobId
} catch {
return NextResponse.json({ error: 'Invalid JSON body' }, { status: 400 })
}

if (!jobId) {
return NextResponse.json({ error: 'jobId required' }, { status: 400 })
}

// SECURITY FIX: Persist a one-time delivery ID (Idempotency Key) with SETNX.
// This prevents the same webhook request from being processed multiple times.
const isDuplicate = await checkAndLockMessage(messageId)
if (isDuplicate) {
return NextResponse.json({ ok: true, note: 'Message already processed' })
}

try {
await inngest.send({
name: 'job/webhook',
data: { jobId },
})
} catch (err) {
console.error('[webhook] Failed to send Inngest event:', err)

// SECURITY FIX: Do NOT swallow errors. If delivery fails, we release the idempotency lock
// and return a 500 status code so the upstream scheduler can properly retry the failed delivery.
await unlockMessage(messageId)
return NextResponse.json({ error: 'Event delivery failed' }, { status: 500 })
}

return NextResponse.json({ ok: true })
}
208 changes: 123 additions & 85 deletions Realism/lib/jobs.ts
Original file line number Diff line number Diff line change
@@ -1,85 +1,123 @@
import { randomUUID } from 'crypto'
import { setJob, getJob, updateJob } from '@/lib/redis'
import type { Job, JobType, JobCadence, Artifact } from '@/types'

export async function createJob(params: {
id?: string
userId: string
goal: string
budget: number
type: JobType
cadence?: JobCadence
spendRuleId?: string
qstashScheduleId?: string
}): Promise<Job> {
const job: Job = {
id: params.id ?? randomUUID(),
userId: params.userId,
goal: params.goal,
budget: params.budget,
type: params.type,
status: 'pending',
spendRuleId: params.spendRuleId,
spendTotal: 0,
cadence: params.cadence,
qstashScheduleId: params.qstashScheduleId,
createdAt: new Date().toISOString(),
}

await setJob(job)
return job
}

export async function startJob(id: string): Promise<Job | null> {
return updateJob(id, { status: 'running' })
}

export async function completeJob(id: string, artifact: Artifact): Promise<Job | null> {
return updateJob(id, {
status: 'complete',
artifact,
completedAt: new Date().toISOString(),
})
}

export async function failJob(id: string, reason: string): Promise<Job | null> {
return updateJob(id, {
status: 'failed',
failureReason: reason,
completedAt: new Date().toISOString(),
})
}

export async function addSpend(id: string, amount: number): Promise<Job | null> {
const job = await getJob(id)
if (!job) return null
return updateJob(id, { spendTotal: job.spendTotal + amount })
}

export async function pauseJob(id: string): Promise<Job | null> {
return updateJob(id, { status: 'paused' })
}

export async function recordRun(
id: string,
artifact: Artifact,
spentThisRun: number
): Promise<Job | null> {
const job = await getJob(id)
if (!job) return null

const now = new Date()
const nextRun = job.cadence === 'daily'
? new Date(now.getTime() + 24 * 60 * 60 * 1000)
: new Date(now.getTime() + 7 * 24 * 60 * 60 * 1000)

return updateJob(id, {
artifact,
lastRunAt: now.toISOString(),
nextRunAt: nextRun.toISOString(),
spendTotal: (job.spendTotal ?? 0) + spentThisRun,
status: 'running',
})
}

export { getJob }
import { randomUUID } from 'crypto'
// Added 'redis' client export to utilize SETNX commands directly
import { setJob, getJob, updateJob, redis } from '@/lib/redis'
import type { Job, JobType, JobCadence, Artifact } from '@/types'

// SECURITY FIX: Webhook Idempotency locking via SETNX
export async function checkAndLockMessage(messageId: string): Promise<boolean> {
try {
const key = `webhook:idempotency:${messageId}`
// SETNX returns 1 if the key was set successfully, 0 if it already existed
const acquired = await redis.setnx(key, '1')
if (acquired === 1) {
// Set expiration to clean up old locks (e.g., 7 days)
await redis.expire(key, 86400 * 7)
return false // Not a duplicate
}
return true // Duplicate execution rejected
} catch (err) {
console.warn('Idempotency check failed, defaulting to block:', err)
return true
}
}

// Releases the lock if the downstream action fails, allowing for a retry
export async function unlockMessage(messageId: string): Promise<void> {
try {
await redis.del(`webhook:idempotency:${messageId}`)
} catch (err) {
console.error('Failed to unlock message:', err)
}
}

export async function createJob(params: {
id?: string
userId: string
goal: string
budget: number
type: JobType
cadence?: JobCadence
spendRuleId?: string
qstashScheduleId?: string
}): Promise<Job> {
const job: Job = {
id: params.id ?? randomUUID(),
userId: params.userId,
goal: params.goal,
budget: params.budget,
type: params.type,
status: 'pending',
spendRuleId: params.spendRuleId,
spendTotal: 0,
cadence: params.cadence,
qstashScheduleId: params.qstashScheduleId,
createdAt: new Date().toISOString(),
}

await setJob(job)
return job
}

export async function startJob(id: string): Promise<Job | null> {
const job = await getJob(id)
if (!job) return null

// SECURITY FIX: Enforce an atomic state transition check.
// Reject duplicate/concurrent executions by ensuring a job can only transition to
// 'running' if it is currently 'pending' or 'paused'.
if (job.status !== 'pending' && job.status !== 'paused') {
throw new Error(`Cannot start job. Current status is '${job.status}'. Duplicate execution rejected.`)
}

return updateJob(id, { status: 'running' })
}

export async function completeJob(id: string, artifact: Artifact): Promise<Job | null> {
return updateJob(id, {
status: 'complete',
artifact,
completedAt: new Date().toISOString(),
})
}

export async function failJob(id: string, reason: string): Promise<Job | null> {
return updateJob(id, {
status: 'failed',
failureReason: reason,
completedAt: new Date().toISOString(),
})
}

export async function addSpend(id: string, amount: number): Promise<Job | null> {
const job = await getJob(id)
if (!job) return null
return updateJob(id, { spendTotal: job.spendTotal + amount })
}

export async function pauseJob(id: string): Promise<Job | null> {
return updateJob(id, { status: 'paused' })
}

export async function recordRun(
id: string,
artifact: Artifact,
spentThisRun: number
): Promise<Job | null> {
const job = await getJob(id)
if (!job) return null

const now = new Date()
const nextRun = job.cadence === 'daily'
? new Date(now.getTime() + 24 * 60 * 60 * 1000)
: new Date(now.getTime() + 7 * 24 * 60 * 60 * 1000)

return updateJob(id, {
artifact,
lastRunAt: now.toISOString(),
nextRunAt: nextRun.toISOString(),
spendTotal: (job.spendTotal ?? 0) + spentThisRun,
status: 'running',
})
}

export { getJob }