File size: 2,320 Bytes
020ed68
 
 
 
461678c
020ed68
 
 
 
 
 
 
 
 
 
 
461678c
020ed68
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
461678c
 
 
 
 
 
 
020ed68
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
import { NextRequest } from 'next/server'
import { Db } from 'mongodb'
import { writeShard, invalidateUsage } from '@/lib/mongoPool'
import { checkAuth } from '@/lib/auth'
import { record } from '@/lib/ingestLog'

// Atlas raises this once a cluster is over its storage quota.
function isQuotaError(err: unknown): boolean {
  return /over your space quota|quota exceeded|you are over/i.test(String(err))
}

export async function POST(req: NextRequest) {
  try {
    const rawBody = Buffer.from(await req.arrayBuffer())
    const auth = checkAuth(rawBody, req.headers)
    if (!auth.ok) {
      record({ route: 'lambda-executions', ok: false, count: 0, detail: auth.error ?? 'auth failed' })
      return Response.json({ error: auth.error }, { status: 401 })
    }

    const body = JSON.parse(rawBody.toString('utf-8'))
    const events = body?.events
    if (!Array.isArray(events) || events.length === 0) {
      return Response.json({ error: 'events array required' }, { status: 400 })
    }

    const now = new Date()
    const docs = events.map((ev) => ({
      event_id: ev.event_id ?? crypto.randomUUID(),
      project_id: ev.project_id ?? 'unknown',
      func_name: ev.func_name ?? 'unknown',
      stage: ev.stage ?? 'prod',
      created_at: now,
    }))

    const write = async (db: Db) => db.collection('lambda_events').insertMany(docs, { ordered: false })

    let db = await writeShard('lambda')
    try {
      await write(db)
    } catch (err) {
      if (!isQuotaError(err)) throw err
      // The cached size estimate was stale and this shard filled between checks.
      // Re-probe and roll to the next one rather than dropping the batch.
      console.warn('lambda shard full, rolling over')
      invalidateUsage('lambda')
      db = await writeShard('lambda')
      await write(db)
    }

    const funcs = [...new Set(docs.map((d) => d.func_name))]
    record({
      route: 'lambda-executions',
      ok: true,
      count: docs.length,
      detail: `${docs[0]?.stage ?? '?'} · ${funcs.slice(0, 3).join(', ')}${funcs.length > 3 ? ` +${funcs.length - 3}` : ''}`,
    })
    return Response.json({ success: true, inserted: docs.length })
  } catch (err) {
    console.error('lambda-executions error:', err)
    return Response.json({ error: 'Internal server error' }, { status: 500 })
  }
}