feat: arcim inbox (Resend Inbound) + smart-match extension + commit metadata (#286)

* feat: multi-series SIE import, reusable FiscalYearSelector, library templates in picker

- SIE import preserves each voucher's source series (B/C/I/V/...), essential
  for Fortnox migrations where series carry semantic meaning (kundfakturor,
  inbetalningar, etc.). Target numbering still goes through next_voucher_number
  per series; source (series, number) is stored in the migration mapping for
  BFNAR 2013:2 audit trail.
- Execute route reads company_settings.default_voucher_series as the fallback
  for vouchers arriving without a series (SIE4I).
- Extract shared FiscalYearSelector component; adopt in /reports and
  /bookkeeping.
- Transaction TemplatePicker now surfaces user-created library templates
  (company + team scope) alongside the static registry, with a helper to
  convert simple library templates into the BookingTemplate shape.
- Exclude 8999 "Årets resultat" from income statement financial section and
  monthly breakdown so year-end closing entries don't cancel the net result.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* test: skip Bokio SIE regression when fixtures are absent

/dev_docs is gitignored (contains anonymised customer exports), so the
integration test can't find its input files in CI. Gate the suite on
fixture presence so it still runs locally.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* fix: address Greptile review feedback

- convertLibraryToBookingTemplate: default entity_applicability to 'all'
  when the source template has no entity_type, so TemplatePicker doesn't
  silently hide it for companies with a set entity type.
- FiscalYearSelector: fire onReady in the no-company early-return branch
  so consumers (e.g. ReportsPage) don't get stuck in a loading skeleton
  while the company context is still hydrating.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* feat: arcim inbox + smart-match extension + commit metadata

Three threads, all gated off in extensions.config.json (invoice-inbox and
inbox-smart-match are not in the enabled list for this PR).

invoice-inbox: Gmail OAuth -> Resend Inbound (v2.0.0)
- Remove gmail-scanner / gmail-helpers
- Add resend-inbound.ts (webhook verify, attachment fetch) and
  inbox-provisioning.ts (per-company @arcim.io address with rotation)
- Replace /gmail/* routes with /inbox/address and admin-only /inbox/rotate
- Workspace UI: card layout + MatchBlock surfacing AI transaction matches
- classify-document: tightened discount/total prompt; cap confidence at
  50% when line items do not reconcile with amount_incl_vat
- Manifest requires RESEND_API_KEY, RESEND_INBOUND_DOMAIN,
  RESEND_INBOUND_WEBHOOK_SECRET

inbox-smart-match (new extension)
- Event-driven AI matching of receipts to bank transactions
- Listens on inbox_item.classified (match now) and transaction.synced
  (retro-match receipts waiting for a transaction)
- Uses service-role client; processing_history append is scoped by
  company_id from the event payload

commit metadata + audit plumbing
- journal_entries gains commit_method and rubric_version columns
- commit_journal_entry RPC accepts both (BFNAR 2013:2 behandlingshistorik)
- processing-history PII detector strips UUID-shaped substrings before
  personnummer pattern matching (UUIDs were triggering false positives)
- New generic inbox_item.classified event

Migrations
- arcim_inbox: company_inboxes table, resend_email_id, email_body_text,
  auto-provision trigger, drops obsolete email_connections
- journal_entry_commit_metadata: new columns + updated RPC
- inbox_attachment_composite: resend_attachment_id + composite unique index
- inbox_smart_match: correlation_id, match_reasoning, expanded match_method
  CHECK, pending-match and correlation indexes

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
Jakob Wennberg
2026-04-20 21:50:07 +02:00
committed by GitHub
co-authored by Claude Opus 4.7
parent dfb953638c
commit 0076aa85f8
31 changed files with 2614 additions and 899 deletions
@@ -0,0 +1,96 @@
import { describe, it, expect } from 'vitest'
import { fetchCandidateTransactions } from '@/extensions/general/inbox-smart-match/lib/fetch-candidates'
import { createQueuedMockSupabase } from '@/tests/helpers'
import type { ReceiptExtractionResult } from '@/types'
function makeReceipt(overrides?: Partial<ReceiptExtractionResult>): ReceiptExtractionResult {
return {
merchant: { name: 'Willys Hemma', orgNumber: null, vatNumber: null, isForeign: false },
receipt: { date: '2026-04-15', time: null, currency: 'SEK' },
lineItems: [],
totals: { subtotal: 239, vatAmount: 60, total: 299 },
flags: { isRestaurant: false, isSystembolaget: false, isForeignMerchant: false },
confidence: 0.9,
...overrides,
} as ReceiptExtractionResult
}
describe('fetchCandidateTransactions', () => {
it('returns empty list when extracted data is missing', async () => {
const { supabase } = createQueuedMockSupabase()
const result = await fetchCandidateTransactions(supabase as never, 'company-1', null)
expect(result).toEqual([])
})
it('returns empty list when no anchors could be derived', async () => {
const { supabase } = createQueuedMockSupabase()
const result = await fetchCandidateTransactions(
supabase as never,
'company-1',
makeReceipt({ totals: { subtotal: 0, vatAmount: 0, total: 0 } } as never)
)
expect(result).toEqual([])
})
it('ranks candidates by amount proximity and returns top 5', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
// 1. already-matched lookup (none)
enqueue({ data: [] })
// 2. candidate transactions
enqueue({
data: [
{ id: 't1', date: '2026-04-15', description: 'A', amount: -400, amount_sek: null, currency: 'SEK', merchant_name: null },
{ id: 't2', date: '2026-04-14', description: 'B', amount: -299, amount_sek: null, currency: 'SEK', merchant_name: null }, // perfect match
{ id: 't3', date: '2026-04-16', description: 'C', amount: -305, amount_sek: null, currency: 'SEK', merchant_name: null }, // close
{ id: 't4', date: '2026-04-15', description: 'D', amount: -150, amount_sek: null, currency: 'SEK', merchant_name: null },
{ id: 't5', date: '2026-04-13', description: 'E', amount: -298, amount_sek: null, currency: 'SEK', merchant_name: null },
{ id: 't6', date: '2026-04-15', description: 'F', amount: -600, amount_sek: null, currency: 'SEK', merchant_name: null },
],
})
const result = await fetchCandidateTransactions(supabase as never, 'company-1', makeReceipt())
expect(result).toHaveLength(5)
expect(result[0].id).toBe('t2') // exact match sorted first
expect(result[1].id).toBe('t5') // ±1 next
})
it('excludes transactions already claimed by other inbox items', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
// 1. already-matched lookup — t2 is already taken
enqueue({ data: [{ matched_transaction_id: 't2' }] })
// 2. candidate transactions
enqueue({
data: [
{ id: 't1', date: '2026-04-15', description: 'A', amount: -400, amount_sek: null, currency: 'SEK', merchant_name: null },
{ id: 't2', date: '2026-04-14', description: 'B', amount: -299, amount_sek: null, currency: 'SEK', merchant_name: null },
{ id: 't3', date: '2026-04-16', description: 'C', amount: -305, amount_sek: null, currency: 'SEK', merchant_name: null },
],
})
const result = await fetchCandidateTransactions(supabase as never, 'company-1', makeReceipt())
const ids = result.map((c) => c.id)
expect(ids).not.toContain('t2')
expect(ids).toContain('t3')
})
it('uses amount_sek when receipt is foreign currency', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: [] }) // no already-matched
enqueue({
data: [
{ id: 't1', date: '2026-04-15', description: 'USD-denom', amount: -895, amount_sek: -895, currency: 'SEK', merchant_name: null },
],
})
const usdReceipt = makeReceipt({
receipt: { date: '2026-04-15', time: null, currency: 'USD' },
totals: { subtotal: 80, vatAmount: 0, total: 85 },
} as never)
const result = await fetchCandidateTransactions(supabase as never, 'company-1', usdReceipt)
// Amount doesn't perfectly match but the function should return it anyway
expect(result).toHaveLength(1)
expect(result[0].id).toBe('t1')
})
})
@@ -0,0 +1,134 @@
import { describe, it, expect, vi, beforeEach } from 'vitest'
// Mock Bedrock SDK before importing the module under test
const mockSend = vi.fn()
vi.mock('@aws-sdk/client-bedrock-runtime', () => {
class ConverseCommand {
public input: unknown
constructor(input: unknown) { this.input = input }
}
class BedrockRuntimeClient {
send(command: unknown) { return mockSend(command) }
}
return { BedrockRuntimeClient, ConverseCommand }
})
import { matchReceiptToCandidate } from '@/extensions/general/inbox-smart-match/lib/match-receipt'
import type { ReceiptExtractionResult } from '@/types'
import type { CandidateTransaction } from '@/extensions/general/inbox-smart-match/lib/fetch-candidates'
function makeExtracted(): ReceiptExtractionResult {
return {
merchant: { name: 'Willys Hemma', orgNumber: null, vatNumber: null, isForeign: false },
receipt: { date: '2026-04-15', time: null, currency: 'SEK' },
lineItems: [],
totals: { subtotal: 239, vatAmount: 60, total: 299 },
flags: { isRestaurant: false, isSystembolaget: false, isForeignMerchant: false },
confidence: 0.9,
} as ReceiptExtractionResult
}
function makeCandidates(): CandidateTransaction[] {
return [
{ id: 't1', date: '2026-04-15', description: 'WILLYS SÖDERM', amount: -299, amount_sek: null, currency: 'SEK', merchant_name: 'Willys' },
{ id: 't2', date: '2026-04-14', description: 'ICA MAXI', amount: -312, amount_sek: null, currency: 'SEK', merchant_name: 'ICA' },
]
}
function mockBedrockResponse(toolInput: Record<string, unknown>) {
mockSend.mockResolvedValue({
output: {
message: {
content: [
{
toolUse: {
toolUseId: 'id',
name: 'match_receipt',
input: toolInput,
},
},
],
},
},
usage: { inputTokens: 100, outputTokens: 20 },
})
}
describe('matchReceiptToCandidate', () => {
beforeEach(() => {
mockSend.mockReset()
process.env.AWS_ACCESS_KEY_ID = 'test'
process.env.AWS_SECRET_ACCESS_KEY = 'test'
process.env.AWS_REGION = 'eu-north-1'
})
it('returns a matched result when LLM picks a valid candidate', async () => {
mockBedrockResponse({
matched: true,
transaction_id: 't1',
confidence: 96,
reasoning: 'Exakt belopp och datum. Willys matchar WILLYS SÖDERM.',
})
const result = await matchReceiptToCandidate({
extracted: makeExtracted(),
candidates: makeCandidates(),
})
expect(result.matched).toBe(true)
expect(result.transactionId).toBe('t1')
expect(result.confidence).toBeCloseTo(0.96, 2)
expect(result.reasoning).toContain('Willys')
})
it('returns no match when LLM says matched=false', async () => {
mockBedrockResponse({
matched: false,
transaction_id: null,
confidence: 10,
reasoning: 'Ingen kandidat har rätt belopp eller handlare.',
})
const result = await matchReceiptToCandidate({
extracted: makeExtracted(),
candidates: makeCandidates(),
})
expect(result.matched).toBe(false)
expect(result.transactionId).toBeNull()
expect(result.confidence).toBeCloseTo(0.1, 2)
})
it('degrades to no-match when LLM returns unknown transaction_id', async () => {
mockBedrockResponse({
matched: true,
transaction_id: 'hallucinated-id-not-in-candidates',
confidence: 80,
reasoning: 'Detta är fel id',
})
const result = await matchReceiptToCandidate({
extracted: makeExtracted(),
candidates: makeCandidates(),
})
expect(result.matched).toBe(false)
expect(result.transactionId).toBeNull()
expect(result.confidence).toBe(0)
})
it('returns safe default when no tool use in response', async () => {
mockSend.mockResolvedValue({
output: { message: { content: [] } },
usage: { inputTokens: 0, outputTokens: 0 },
})
const result = await matchReceiptToCandidate({
extracted: makeExtracted(),
candidates: makeCandidates(),
})
expect(result.matched).toBe(false)
expect(result.transactionId).toBeNull()
})
})
@@ -0,0 +1,114 @@
import type { Extension } from '@/lib/extensions/types'
import type { EventPayload } from '@/lib/events/types'
import type { InvoiceInboxItem } from '@/types'
import type { SupabaseClient } from '@supabase/supabase-js'
import { createClient } from '@supabase/supabase-js'
import { processInboxItemMatch } from './lib/process-match'
const EXTENSION_ID = 'inbox-smart-match'
// The handler always uses a service-role client:
// - processing_history has no INSERT RLS policy (audit integrity) — only service-role can append
// - every query is scoped by company_id from the event payload
// - we write pseudonymous IDs only (PII validator enforces this at the append layer)
// Using ctx.supabase is unsafe because the registry wrapper may build an anon-key
// client when no user session exists (e.g. webhook path).
function getServiceSupabase(): SupabaseClient {
return createClient(
process.env.NEXT_PUBLIC_SUPABASE_URL!,
process.env.SUPABASE_SERVICE_ROLE_KEY!
)
}
export const inboxSmartMatchExtension: Extension = {
id: EXTENSION_ID,
name: 'Smart matchning',
version: '0.1.0',
eventHandlers: [
// When an inbox item is freshly classified, try to match it to a transaction
{
eventType: 'inbox_item.classified',
handler: async (payload: EventPayload<'inbox_item.classified'>) => {
// Only act on receipts for v1
if (payload.documentType !== 'receipt') return
const supabase = getServiceSupabase()
try {
await processInboxItemMatch(
{
supabase,
companyId: payload.companyId,
userId: payload.userId,
extensionId: EXTENSION_ID,
triggerReason: 'classified',
},
payload.inboxItem
)
} catch (err) {
console.error(
`[${EXTENSION_ID}] Failed to process classified inbox item ${payload.inboxItem.id}:`,
err
)
}
},
},
// When new transactions land, retry matching on any receipts still waiting
{
eventType: 'transaction.synced',
handler: async (payload: EventPayload<'transaction.synced'>) => {
const newTransactionIds = payload.transactions.map((t) => t.id).filter(Boolean)
if (newTransactionIds.length === 0) return
const supabase = getServiceSupabase()
// Find receipts in pending state for this company.
// Cap at 10 per sync so one big bank import doesn't time out the
// handler; leftover pending items pick up on the next sync.
const { data: pendingItems, error } = await supabase
.from('invoice_inbox_items')
.select('*')
.eq('company_id', payload.companyId)
.eq('document_type', 'receipt')
.eq('status', 'ready')
.eq('match_method', 'pending_transaction')
.order('created_at', { ascending: false })
.limit(10)
if (error) {
console.error(`[${EXTENSION_ID}] Failed to fetch pending receipts:`, error)
return
}
if (!pendingItems || pendingItems.length === 0) return
// Run LLM calls in parallel; one failing receipt shouldn't stop the others.
const results = await Promise.allSettled(
(pendingItems as InvoiceInboxItem[]).map((item) =>
processInboxItemMatch(
{
supabase,
companyId: payload.companyId,
userId: payload.userId,
extensionId: EXTENSION_ID,
triggerReason: 'transaction_synced',
},
item
)
)
)
results.forEach((r, i) => {
if (r.status === 'rejected') {
const itemId = (pendingItems[i] as InvoiceInboxItem).id
console.error(
`[${EXTENSION_ID}] Retroactive match failed for item ${itemId}:`,
r.reason
)
}
})
},
},
],
}
@@ -0,0 +1,125 @@
/**
* Candidate transaction fetcher — deterministic narrowing before the LLM call.
*
* Pulls unbooked expense transactions within ±7 days of the receipt date,
* ordered by how close their amount is to the receipt total. Limits to top 5
* so the LLM has a focused candidate set and the token cost stays bounded.
*/
import type { SupabaseClient } from '@supabase/supabase-js'
import type { ReceiptExtractionResult } from '@/types'
const DATE_WINDOW_DAYS = 7
const MAX_CANDIDATES = 5
export interface CandidateTransaction {
id: string
date: string
description: string
amount: number
amount_sek: number | null
currency: string
merchant_name: string | null
}
/**
* Extract the reference date and absolute amount from a classified receipt's
* extracted data. Returns null if required fields are missing.
*/
function getReceiptMatchAnchors(
extracted: ReceiptExtractionResult | null
): { date: string; amount: number; currency: string } | null {
if (!extracted) return null
const date = extracted.receipt?.date ?? null
const amount = extracted.totals?.total ?? null
const currency = extracted.receipt?.currency ?? 'SEK'
if (!date || amount == null || amount <= 0) return null
return { date, amount, currency }
}
/**
* Fetch up to MAX_CANDIDATES unbooked expense transactions near the receipt's
* date + amount. Ordering prefers exact amount matches first.
*/
export async function fetchCandidateTransactions(
supabase: SupabaseClient,
companyId: string,
extracted: ReceiptExtractionResult | null
): Promise<CandidateTransaction[]> {
const anchors = getReceiptMatchAnchors(extracted)
if (!anchors) return []
const receiptDate = new Date(anchors.date)
if (isNaN(receiptDate.getTime())) return []
const windowStart = new Date(receiptDate)
windowStart.setUTCDate(windowStart.getUTCDate() - DATE_WINDOW_DAYS)
const windowEnd = new Date(receiptDate)
windowEnd.setUTCDate(windowEnd.getUTCDate() + DATE_WINDOW_DAYS)
// Exclude transactions already claimed by any other inbox item in this
// company. The partial unique index on (company_id, matched_transaction_id)
// is the final guard against concurrent double-matches, but filtering up
// front saves an LLM roundtrip on the obvious cases.
const { data: claimed, error: claimedError } = await supabase
.from('invoice_inbox_items')
.select('matched_transaction_id')
.eq('company_id', companyId)
.not('matched_transaction_id', 'is', null)
if (claimedError) {
throw new Error(`Failed to load matched transactions: ${claimedError.message}`)
}
const excludedIds = new Set(
(claimed ?? [])
.map((row) => (row as { matched_transaction_id: string | null }).matched_transaction_id)
.filter((id): id is string => typeof id === 'string' && id.length > 0)
)
// Pull negative-amount (expense) transactions without a journal entry in the window
const { data, error } = await supabase
.from('transactions')
.select('id, date, description, amount, amount_sek, currency, merchant_name')
.eq('company_id', companyId)
.is('journal_entry_id', null)
.lt('amount', 0)
.gte('date', windowStart.toISOString().slice(0, 10))
.lte('date', windowEnd.toISOString().slice(0, 10))
.order('date', { ascending: false })
.limit(50)
if (error) {
throw new Error(`Failed to fetch candidate transactions: ${error.message}`)
}
if (!data || data.length === 0) return []
const filtered = excludedIds.size > 0
? data.filter((tx) => !excludedIds.has(tx.id as string))
: data
if (filtered.length === 0) return []
// Rank candidates by amount proximity. For SEK receipts we compare directly,
// for other currencies we prefer amount_sek if the receipt amount has been converted.
const receiptAbs = Math.abs(anchors.amount)
const scored = filtered.map((tx) => {
const txAmount = Math.abs(Number(tx.amount) || 0)
const txSek = tx.amount_sek == null ? null : Math.abs(Number(tx.amount_sek))
const primaryDiff = Math.abs(txAmount - receiptAbs)
const sekDiff = txSek == null ? Infinity : Math.abs(txSek - receiptAbs)
const bestDiff = Math.min(primaryDiff, sekDiff)
return { tx, diff: bestDiff }
})
scored.sort((a, b) => a.diff - b.diff)
return scored.slice(0, MAX_CANDIDATES).map(({ tx }) => ({
id: tx.id,
date: tx.date,
description: tx.description ?? '',
amount: Number(tx.amount),
amount_sek: tx.amount_sek == null ? null : Number(tx.amount_sek),
currency: tx.currency ?? 'SEK',
merchant_name: tx.merchant_name ?? null,
}))
}
@@ -0,0 +1,200 @@
/**
* Text-only LLM matcher — decides which candidate bank transaction (if any)
* corresponds to a classified receipt. Uses Bedrock Converse with structured
* tool output. No image input: the receipt is already represented by the
* extracted data, and candidates are pure text.
*/
import {
BedrockRuntimeClient,
ConverseCommand,
type ContentBlock,
type Message,
type ToolConfiguration,
} from '@aws-sdk/client-bedrock-runtime'
import type { ReceiptExtractionResult } from '@/types'
import type { CandidateTransaction } from './fetch-candidates'
export interface ReceiptMatchResult {
matched: boolean
transactionId: string | null
confidence: number // 0..1
reasoning: string
usage: { inputTokens: number; outputTokens: number }
}
let _client: BedrockRuntimeClient | null = null
function getClient(): BedrockRuntimeClient {
if (!_client) {
_client = new BedrockRuntimeClient({
region: process.env.AWS_REGION || 'eu-north-1',
credentials: {
accessKeyId: process.env.AWS_ACCESS_KEY_ID!,
secretAccessKey: process.env.AWS_SECRET_ACCESS_KEY!,
},
})
}
return _client
}
const SYSTEM_PROMPT = `Du är en expert på att matcha svenska kvitton mot banktransaktioner.
Du får:
- Kvittodata (handlare, belopp, valuta, datum) från AI-extraktion
- En lista med kandidat-banktransaktioner (id, beskrivning, belopp, valuta, datum, MCC)
Uppgift: identifiera vilken (om någon) banktransaktion som motsvarar kvittot.
Resonera utifrån:
- Belopp: bör vara identiskt eller mycket nära (ta hänsyn till valutaväxling om olika valutor)
- Datum: banktransaktion bokförs ofta 0-3 dagar efter kvittot
- Handlare: bankens beskrivning är ofta förkortad/versaler ("WILLYS SÖDERM" = "Willys Hemma Södermalm"). Matcha semantiskt, inte bokstavligt
- MCC-koder kan bekräfta branschtyp
Om inget förslag är trovärdigt — returnera matched=false.
Anropa ALLTID verktyget match_receipt med resultatet.
Motivering ska vara kort, på svenska och förklara varför transaktionen valdes.`
const MATCH_TOOL: ToolConfiguration = {
tools: [
{
toolSpec: {
name: 'match_receipt',
description: 'Returnera vilken kandidat-transaktion som matchar kvittot',
inputSchema: {
json: {
type: 'object',
required: ['matched', 'confidence', 'reasoning'],
properties: {
matched: {
type: 'boolean',
description: 'true om en kandidat matchar, false annars',
},
transaction_id: {
type: ['string', 'null'],
description: 'id för den matchande kandidaten (null om matched=false)',
},
confidence: {
type: 'integer',
minimum: 0,
maximum: 100,
description: 'Säkerhet 0-100. Sätt lågt när matched=false.',
},
reasoning: {
type: 'string',
description: '1-2 meningar på svenska som förklarar beslutet.',
},
},
},
},
},
},
],
toolChoice: { any: {} },
}
export interface MatchReceiptInput {
extracted: ReceiptExtractionResult
candidates: CandidateTransaction[]
}
/**
* Call Bedrock to choose the best matching transaction.
* Returns a neutral result (matched=false) if the model doesn't find a fit or
* the tool schema is missing from the response.
*/
export async function matchReceiptToCandidate(
input: MatchReceiptInput
): Promise<ReceiptMatchResult> {
const receiptBrief = {
merchant: input.extracted.merchant?.name ?? null,
amount: input.extracted.totals?.total ?? null,
currency: input.extracted.receipt?.currency ?? 'SEK',
date: input.extracted.receipt?.date ?? null,
vat_amount: input.extracted.totals?.vatAmount ?? null,
}
const candidateLines = input.candidates.map((c) => ({
id: c.id,
date: c.date,
description: c.description,
amount: c.amount,
amount_sek: c.amount_sek,
currency: c.currency,
merchant_name: c.merchant_name,
}))
const userPrompt = `Kvitto:
${JSON.stringify(receiptBrief, null, 2)}
Kandidat-transaktioner:
${JSON.stringify(candidateLines, null, 2)}
Vilken transaktion matchar kvittot? Om ingen matchar, returnera matched=false.`
const messages: Message[] = [
{
role: 'user',
content: [{ text: userPrompt }],
},
]
const modelId = process.env.BEDROCK_MODEL_ID || 'eu.anthropic.claude-sonnet-4-6'
const maxTokens = parseInt(process.env.BEDROCK_MAX_TOKENS || '1024', 10)
const command = new ConverseCommand({
modelId,
messages,
system: [{ text: SYSTEM_PROMPT }],
toolConfig: MATCH_TOOL,
inferenceConfig: { maxTokens, temperature: 0 },
})
const response = await getClient().send(command)
const usage = {
inputTokens: response.usage?.inputTokens ?? 0,
outputTokens: response.usage?.outputTokens ?? 0,
}
const outputMessage = response.output?.message
if (!outputMessage?.content) {
return { matched: false, transactionId: null, confidence: 0, reasoning: 'Inget LLM-svar', usage }
}
const toolUseBlock = outputMessage.content.find(
(block): block is ContentBlock.ToolUseMember => 'toolUse' in block && block.toolUse !== undefined
)
if (!toolUseBlock?.toolUse?.input) {
return { matched: false, transactionId: null, confidence: 0, reasoning: 'Inget verktygsanrop', usage }
}
const raw = toolUseBlock.toolUse.input as Record<string, unknown>
const matched = Boolean(raw.matched)
const rawId = typeof raw.transaction_id === 'string' ? raw.transaction_id : null
const transactionId = matched
? (input.candidates.find((c) => c.id === rawId)?.id ?? null)
: null
const confidenceRaw = Number(raw.confidence)
const confidence =
isFinite(confidenceRaw)
? Math.min(1, Math.max(0, confidenceRaw / 100))
: 0
const reasoning = typeof raw.reasoning === 'string' ? raw.reasoning.trim() : ''
// If LLM said matched but we can't resolve the transaction_id to a candidate,
// degrade gracefully to unmatched so downstream isn't left dangling.
if (matched && !transactionId) {
return {
matched: false,
transactionId: null,
confidence: 0,
reasoning: reasoning || 'LLM angav ogiltigt transaction_id',
usage,
}
}
return { matched, transactionId, confidence, reasoning, usage }
}
@@ -0,0 +1,212 @@
/**
* Core matching flow — deterministic narrowing + LLM call + persistence +
* processing_history audit events. Called from both the classify handler and
* the transaction-sync retroactive handler.
*/
import type { SupabaseClient } from '@supabase/supabase-js'
import type { InvoiceInboxItem, ReceiptExtractionResult } from '@/types'
import { appendProcessingHistory } from '@/lib/processing-history/append'
import { fetchCandidateTransactions } from './fetch-candidates'
import { matchReceiptToCandidate } from './match-receipt'
export interface MatchContext {
supabase: SupabaseClient
companyId: string
userId: string
extensionId: string
triggerReason: 'classified' | 'transaction_synced'
}
export interface MatchOutcome {
status: 'matched' | 'no_match' | 'pending_transaction' | 'skipped'
transactionId: string | null
confidence: number
reasoning: string
}
/**
* Process a single classified-receipt inbox item through the matcher pipeline.
* Writes match fields + appends processing_history. Swallows internal errors
* so one failing receipt doesn't break the whole event handler.
*/
export async function processInboxItemMatch(
ctx: MatchContext,
item: InvoiceInboxItem
): Promise<MatchOutcome> {
const tag = `[inbox-smart-match] item=${item.id} trigger=${ctx.triggerReason}`
// We only operate on receipts for v1
if (item.document_type !== 'receipt') {
return { status: 'skipped', transactionId: null, confidence: 0, reasoning: '' }
}
if (item.status !== 'ready') {
return { status: 'skipped', transactionId: null, confidence: 0, reasoning: '' }
}
if (!item.extracted_data) {
return { status: 'skipped', transactionId: null, confidence: 0, reasoning: '' }
}
const correlationId = item.correlation_id ?? crypto.randomUUID()
// If we just minted a fresh correlation_id (legacy row predating the column),
// persist it so retries reuse the same thread through processing_history.
if (!item.correlation_id) {
const { error: corrError } = await ctx.supabase
.from('invoice_inbox_items')
.update({ correlation_id: correlationId })
.eq('id', item.id)
if (corrError) {
console.error(`${tag} — failed to persist correlation_id:`, corrError)
// non-fatal — we still proceed with matching under the in-memory ID
}
}
const extracted = item.extracted_data as unknown as ReceiptExtractionResult
const candidates = await fetchCandidateTransactions(ctx.supabase, ctx.companyId, extracted)
// Append DeterministicMatch event — records that the narrowing ran
let deterministicEventId: string
try {
deterministicEventId = await appendProcessingHistory(ctx.supabase, {
companyId: ctx.companyId,
correlationId,
aggregateType: 'MatchProposal',
aggregateId: item.id,
eventType: 'MatchAttemptedDeterministic',
payload: {
inbox_item_id: item.id,
candidate_count: candidates.length,
candidate_ids: candidates.map((c) => c.id),
window_days: 7,
trigger: ctx.triggerReason,
},
actor: { type: 'system', id: ctx.extensionId },
occurredAt: new Date(),
})
} catch (err) {
console.error(`${tag} — failed to append MatchAttemptedDeterministic:`, err)
return { status: 'no_match', transactionId: null, confidence: 0, reasoning: '' }
}
// No candidates → mark pending, wait for bank sync
if (candidates.length === 0) {
await ctx.supabase
.from('invoice_inbox_items')
.update({
match_method: 'pending_transaction',
match_confidence: null,
matched_transaction_id: null,
match_reasoning: 'Inväntar matchande banktransaktion',
})
.eq('id', item.id)
return {
status: 'pending_transaction',
transactionId: null,
confidence: 0,
reasoning: 'Inväntar matchande banktransaktion',
}
}
// LLM chooses among candidates
let llm
try {
llm = await matchReceiptToCandidate({ extracted, candidates })
} catch (err) {
console.error(`${tag} — LLM matcher failed:`, err)
// Don't overwrite existing state on LLM failure; just log and exit
return { status: 'no_match', transactionId: null, confidence: 0, reasoning: '' }
}
// Record the LLM attempt in processing_history
try {
await appendProcessingHistory(ctx.supabase, {
companyId: ctx.companyId,
correlationId,
causationId: deterministicEventId,
aggregateType: 'MatchProposal',
aggregateId: item.id,
eventType: 'MatchAttemptedLlm',
payload: {
inbox_item_id: item.id,
matched: llm.matched,
chosen_transaction_id: llm.transactionId,
confidence: llm.confidence,
llm_input_tokens: llm.usage.inputTokens,
llm_output_tokens: llm.usage.outputTokens,
candidate_count: candidates.length,
},
actor: { type: 'llm', id: 'match_receipt' },
occurredAt: new Date(),
})
} catch (err) {
console.error(`${tag} — failed to append MatchAttemptedLlm:`, err)
}
// Persist match. The (company_id, matched_transaction_id) partial unique
// index means a concurrent second inbox item trying to claim the same
// transaction will get a 23505 — we catch that and downgrade this one
// to pending_transaction instead of overwriting the winner.
if (llm.matched && llm.transactionId) {
const { error: updateError } = await ctx.supabase
.from('invoice_inbox_items')
.update({
matched_transaction_id: llm.transactionId,
match_confidence: llm.confidence,
match_method: 'llm',
match_reasoning: llm.reasoning,
})
.eq('id', item.id)
if (updateError) {
const code = (updateError as { code?: string }).code
if (code === '23505') {
// Another receipt won the race for this transaction.
await ctx.supabase
.from('invoice_inbox_items')
.update({
matched_transaction_id: null,
match_method: 'pending_transaction',
match_confidence: null,
match_reasoning: 'Transaktionen matchades först till ett annat kvitto',
})
.eq('id', item.id)
return {
status: 'pending_transaction',
transactionId: null,
confidence: 0,
reasoning: 'Transaktionen matchades först till ett annat kvitto',
}
}
console.error(`${tag} — failed to persist match:`, updateError)
return { status: 'no_match', transactionId: null, confidence: 0, reasoning: '' }
}
return {
status: 'matched',
transactionId: llm.transactionId,
confidence: llm.confidence,
reasoning: llm.reasoning,
}
}
// LLM said no match among the candidates — record explanatory reasoning
await ctx.supabase
.from('invoice_inbox_items')
.update({
matched_transaction_id: null,
match_confidence: llm.confidence,
match_method: 'pending_transaction',
match_reasoning: llm.reasoning || 'AI kunde inte hitta matchande transaktion bland kandidaterna',
})
.eq('id', item.id)
return {
status: 'no_match',
transactionId: null,
confidence: llm.confidence,
reasoning: llm.reasoning,
}
}
@@ -0,0 +1,23 @@
{
"id": "inbox-smart-match",
"sector": "general",
"exportName": "inboxSmartMatchExtension",
"entryPoint": "@/extensions/general/inbox-smart-match",
"requiredEnvVars": [
"AWS_ACCESS_KEY_ID",
"AWS_SECRET_ACCESS_KEY",
"AWS_REGION"
],
"optionalEnvVars": ["BEDROCK_MODEL_ID", "BEDROCK_MAX_TOKENS"],
"npmDependencies": ["@aws-sdk/client-bedrock-runtime"],
"definition": {
"name": "Smart matchning",
"category": "operations",
"icon": "Sparkles",
"dataPattern": "core",
"hasOwnData": false,
"readsCoreTables": ["invoice_inbox_items", "transactions", "processing_history"],
"description": "AI-driven matchning av kvitton mot banktransaktioner",
"longDescription": "När ett kvitto klassificeras i inkorgen föreslår AI den mest sannolika matchande banktransaktionen, med motivering. Körs även retroaktivt när nya transaktioner synkas in. Kräver AWS Bedrock."
}
}
@@ -0,0 +1,213 @@
import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'
import { invoiceInboxExtension } from '@/extensions/general/invoice-inbox'
import { ResendSignatureError } from '@/extensions/general/invoice-inbox/lib/resend-inbound'
import { createQueuedMockSupabase, createMockRequest } from '@/tests/helpers'
vi.mock('@/extensions/general/invoice-inbox/lib/resend-inbound', async () => {
const actual = await vi.importActual<typeof import('@/extensions/general/invoice-inbox/lib/resend-inbound')>(
'@/extensions/general/invoice-inbox/lib/resend-inbound'
)
return {
...actual,
verifyInboundWebhook: vi.fn(),
fetchReceivingEmail: vi.fn(),
fetchInboundAttachment: vi.fn(),
}
})
vi.mock('@supabase/supabase-js', () => ({
createClient: vi.fn(),
}))
import { verifyInboundWebhook, fetchReceivingEmail, fetchInboundAttachment } from '@/extensions/general/invoice-inbox/lib/resend-inbound'
import { createClient } from '@supabase/supabase-js'
function findRoute(method: string, path: string) {
return invoiceInboxExtension.apiRoutes!.find((r) => r.method === method && r.path === path)!
}
const webhookRoute = findRoute('POST', '/inbound')
function mockReceivedEvent(overrides?: Record<string, unknown>) {
return {
type: 'email.received' as const,
created_at: '2026-04-20T10:00:00Z',
data: {
email_id: 'em_123',
created_at: '2026-04-20T10:00:00Z',
from: 'billing@supplier.com',
to: ['acme-ab-x7f2@arcim.io'],
cc: [],
bcc: [],
subject: 'Invoice #5678',
message_id: '<msg-id@supplier.com>',
attachments: [
{
id: 'att_1',
filename: 'invoice.pdf',
size: 12345,
content_type: 'application/pdf',
content_id: 'cid1',
content_disposition: 'attachment',
},
],
...overrides,
},
}
}
describe('POST /inbound', () => {
const originalEnv = { ...process.env }
beforeEach(() => {
vi.clearAllMocks()
process.env.RESEND_INBOUND_DOMAIN = 'arcim.io'
process.env.NEXT_PUBLIC_SUPABASE_URL = 'http://localhost'
process.env.SUPABASE_SERVICE_ROLE_KEY = 'test-service-key'
})
afterEach(() => {
process.env = { ...originalEnv }
})
it('returns 503 when RESEND_INBOUND_DOMAIN is not set', async () => {
delete process.env.RESEND_INBOUND_DOMAIN
const request = createMockRequest('/inbound', { method: 'POST', body: { type: 'email.received' } })
const res = await webhookRoute.handler(request)
expect(res.status).toBe(503)
})
it('returns 401 when signature verification fails', async () => {
vi.mocked(verifyInboundWebhook).mockImplementation(() => {
throw new ResendSignatureError('bad sig')
})
const request = createMockRequest('/inbound', { method: 'POST', body: { type: 'email.received' } })
const res = await webhookRoute.handler(request)
expect(res.status).toBe(401)
})
it('ignores non-received events with 200', async () => {
vi.mocked(verifyInboundWebhook).mockReturnValue({
type: 'email.sent',
created_at: '',
data: {},
} as never)
const request = createMockRequest('/inbound', { method: 'POST', body: {} })
const res = await webhookRoute.handler(request)
expect(res.status).toBe(200)
})
it('returns 404 when no recipient matches our domain', async () => {
vi.mocked(verifyInboundWebhook).mockReturnValue(
mockReceivedEvent({ to: ['random@contoso.com'] }) as never
)
const request = createMockRequest('/inbound', { method: 'POST', body: {} })
const res = await webhookRoute.handler(request)
expect(res.status).toBe(404)
})
it('returns 404 when the address is not in company_inboxes', async () => {
vi.mocked(verifyInboundWebhook).mockReturnValue(mockReceivedEvent() as never)
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: null }) // company_inboxes lookup returns nothing
vi.mocked(createClient).mockReturnValue(supabase as never)
const request = createMockRequest('/inbound', { method: 'POST', body: {} })
const res = await webhookRoute.handler(request)
expect(res.status).toBe(404)
})
it('returns 410 when the address is deprecated', async () => {
vi.mocked(verifyInboundWebhook).mockReturnValue(mockReceivedEvent() as never)
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: { id: 'inbox-1', company_id: 'company-1', status: 'deprecated' } })
vi.mocked(createClient).mockReturnValue(supabase as never)
const request = createMockRequest('/inbound', { method: 'POST', body: {} })
const res = await webhookRoute.handler(request)
expect(res.status).toBe(410)
})
it('skips already-processed attachments (per-attachment idempotency)', async () => {
vi.mocked(verifyInboundWebhook).mockReturnValue(mockReceivedEvent() as never)
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: { id: 'inbox-1', company_id: 'company-1', status: 'active' } }) // inbox lookup
enqueue({ data: { created_by: 'user-owner-1' } }) // company owner
enqueue({ data: { id: 'existing-item-1' } }) // per-attachment dup check finds existing row
vi.mocked(createClient).mockReturnValue(supabase as never)
vi.mocked(fetchReceivingEmail).mockResolvedValue({
object: 'email',
id: 'em_123',
to: ['acme-ab-x7f2@arcim.io'],
from: 'billing@supplier.com',
created_at: '2026-04-20T10:00:00Z',
subject: 'Invoice #5678',
bcc: null,
cc: null,
reply_to: null,
html: null,
text: 'Body',
headers: {},
message_id: '<msg@x>',
raw: null,
attachments: [
{ id: 'att_1', filename: 'invoice.pdf', size: 100, content_type: 'application/pdf', content_id: 'cid', content_disposition: 'attachment' },
],
} as never)
const request = createMockRequest('/inbound', { method: 'POST', body: {} })
const res = await webhookRoute.handler(request)
const body = await res.json()
expect(res.status).toBe(200)
expect(body.data.results[0].duplicate).toBe(true)
expect(body.data.results[0].inbox_item_id).toBe('existing-item-1')
expect(fetchInboundAttachment).not.toHaveBeenCalled()
})
it('returns 500 when the company has no created_by owner', async () => {
vi.mocked(verifyInboundWebhook).mockReturnValue(mockReceivedEvent() as never)
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: { id: 'inbox-1', company_id: 'company-1', status: 'active' } })
enqueue({ data: { created_by: null } }) // company with no owner
vi.mocked(createClient).mockReturnValue(supabase as never)
const request = createMockRequest('/inbound', { method: 'POST', body: {} })
const res = await webhookRoute.handler(request)
expect(res.status).toBe(500)
})
it('logs an error inbox item when email has no attachments', async () => {
vi.mocked(verifyInboundWebhook).mockReturnValue(
mockReceivedEvent({ attachments: [] }) as never
)
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: { id: 'inbox-1', company_id: 'company-1', status: 'active' } })
enqueue({ data: { created_by: 'user-owner-1' } })
enqueue({ data: null }) // insert succeeds
vi.mocked(createClient).mockReturnValue(supabase as never)
vi.mocked(fetchReceivingEmail).mockResolvedValue({
object: 'email',
id: 'em_123',
to: ['acme-ab-x7f2@arcim.io'],
from: 'billing@supplier.com',
created_at: '2026-04-20T10:00:00Z',
subject: 'No attachments here',
bcc: null,
cc: null,
reply_to: null,
html: null,
text: 'Body only',
headers: {},
message_id: '<msg@x>',
raw: null,
attachments: [],
} as never)
const request = createMockRequest('/inbound', { method: 'POST', body: {} })
const res = await webhookRoute.handler(request)
const body = await res.json()
expect(res.status).toBe(200)
expect(body.data.reason).toBe('no_attachments')
expect(fetchInboundAttachment).not.toHaveBeenCalled()
})
})
@@ -0,0 +1,98 @@
import { describe, it, expect } from 'vitest'
import {
composeInboxAddress,
getActiveInbox,
rotateCompanyInbox,
} from '@/extensions/general/invoice-inbox/lib/inbox-provisioning'
import { createQueuedMockSupabase } from '@/tests/helpers'
describe('composeInboxAddress', () => {
it('joins local_part and domain with @', () => {
expect(composeInboxAddress('acme-ab-x7f2', 'arcim.io')).toBe('acme-ab-x7f2@arcim.io')
})
})
describe('getActiveInbox', () => {
it('returns the active inbox row', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
const row = {
id: 'inbox-1',
company_id: 'company-1',
local_part: 'acme-x7f2',
status: 'active',
slug_seed: 'acme',
created_at: '2026-04-20T00:00:00Z',
updated_at: '2026-04-20T00:00:00Z',
deprecated_at: null,
}
enqueue({ data: row })
const result = await getActiveInbox(supabase as never, 'company-1')
expect(result).toEqual(row)
})
it('returns null when no active inbox exists', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: null })
const result = await getActiveInbox(supabase as never, 'company-1')
expect(result).toBeNull()
})
it('throws on database error', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: null, error: { message: 'db blew up' } })
await expect(getActiveInbox(supabase as never, 'company-1')).rejects.toThrow(/db blew up/)
})
})
describe('rotateCompanyInbox', () => {
it('delegates to the rotate_company_inbox RPC and returns the new row', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
const newRow = {
id: 'inbox-2',
company_id: 'company-1',
local_part: 'acme-new2',
status: 'active',
slug_seed: 'acme',
created_at: '2026-04-20T00:00:00Z',
updated_at: '2026-04-20T00:00:00Z',
deprecated_at: null,
}
enqueue({ data: newRow })
const result = await rotateCompanyInbox(supabase as never, 'company-1')
expect(result.local_part).toBe('acme-new2')
expect(result.status).toBe('active')
})
it('accepts the SETOF shape where the RPC wraps the row in an array', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
const newRow = {
id: 'inbox-3',
company_id: 'company-1',
local_part: 'new-co-abcd',
status: 'active',
slug_seed: 'new-co',
created_at: '2026-04-20T00:00:00Z',
updated_at: '2026-04-20T00:00:00Z',
deprecated_at: null,
}
enqueue({ data: [newRow] })
const result = await rotateCompanyInbox(supabase as never, 'company-1')
expect(result.local_part).toBe('new-co-abcd')
})
it('surfaces RPC errors', async () => {
const { supabase, enqueue } = createQueuedMockSupabase()
enqueue({ data: null, error: { message: 'Not authorized to rotate inbox for this company' } })
await expect(rotateCompanyInbox(supabase as never, 'nope'))
.rejects.toThrow(/Not authorized/)
})
})
@@ -0,0 +1,52 @@
import { describe, it, expect } from 'vitest'
import { extractLocalPartForDomain } from '@/extensions/general/invoice-inbox/lib/resend-inbound'
describe('extractLocalPartForDomain', () => {
it('returns the local part when a recipient matches the domain', () => {
const result = extractLocalPartForDomain(
['acme-ab-x7f2@arcim.io', 'billing@acme.se'],
'arcim.io'
)
expect(result).toBe('acme-ab-x7f2')
})
it('lowercases the local part and matches domain case-insensitively', () => {
const result = extractLocalPartForDomain(
['ACME-AB-X7F2@ARCIM.IO'],
'arcim.io'
)
expect(result).toBe('acme-ab-x7f2')
})
it('returns null when no recipient matches', () => {
const result = extractLocalPartForDomain(
['billing@acme.se', 'invoices@contoso.com'],
'arcim.io'
)
expect(result).toBeNull()
})
it('returns null for malformed addresses', () => {
const result = extractLocalPartForDomain(
['not-an-email', '@arcim.io', 'foo@'],
'arcim.io'
)
expect(result).toBeNull()
})
it('returns the first matching recipient when multiple match', () => {
const result = extractLocalPartForDomain(
['first-abcd@arcim.io', 'second-efgh@arcim.io'],
'arcim.io'
)
expect(result).toBe('first-abcd')
})
it('trims whitespace inside candidate addresses', () => {
const result = extractLocalPartForDomain(
[' acme-xxx@arcim.io '],
'arcim.io'
)
expect(result).toBe('acme-xxx')
})
})
+343 -224
View File
@@ -3,10 +3,23 @@ import { NextResponse } from 'next/server'
import { createClient } from '@supabase/supabase-js'
import { uploadDocument } from '@/lib/core/documents/document-service'
import { classifyDocument } from './lib/classify-document'
import { encryptState, decryptState, encryptToken, decryptToken } from './lib/gmail-helpers'
import { scanGmailConnection } from './lib/gmail-scanner'
import {
verifyInboundWebhook,
fetchReceivingEmail,
fetchInboundAttachment,
extractLocalPartForDomain,
isEmailReceivedEvent,
ResendSignatureError,
} from './lib/resend-inbound'
import {
rotateCompanyInbox,
getActiveInbox,
composeInboxAddress,
} from './lib/inbox-provisioning'
import { createSupplierInvoiceRegistrationEntry } from '@/lib/bookkeeping/supplier-invoice-entries'
import { CreateSupplierInvoiceSchema } from '@/lib/api/schemas'
import { appendProcessingHistory } from '@/lib/processing-history/append'
import { eventBus } from '@/lib/events/bus'
import type { InvoiceExtractionResult, InvoiceInboxItem, SupplierInvoice, SupplierInvoiceItem } from '@/types'
const MAX_FILE_SIZE = 10 * 1024 * 1024 // Match MAX_DOCUMENT_SIZE from document-service
@@ -20,7 +33,15 @@ const UPLOAD_ALLOWED_MIME_TYPES = new Set([
'image/webp',
])
const STATE_TTL_MS = 10 * 60 * 1000
interface EmailMeta {
from?: string | null
subject?: string | null
receivedAt?: string | null
messageId?: string | null
bodyText?: string | null
resendEmailId?: string | null
resendAttachmentId?: string | null
}
// ── Shared helper: upload + classify + create inbox item ─────
@@ -30,9 +51,12 @@ async function uploadAndClassify(
companyId: string,
file: { name: string; buffer: ArrayBuffer; type: string },
source: 'upload' | 'email',
emailMeta?: { from?: string | null; subject?: string | null; receivedAt?: string | null; messageId?: string },
emailMeta?: EmailMeta,
ctx?: ExtensionContext
) {
// Correlation ID threads through ingest → classify → match → book.
const correlationId = crypto.randomUUID()
// Store in WORM archive
const doc = await uploadDocument(supabase, userId, companyId, {
name: file.name,
@@ -42,6 +66,27 @@ async function uploadAndClassify(
upload_source: source === 'email' ? 'email' : 'file_upload',
})
// Audit: DocumentIngested
try {
await appendProcessingHistory(supabase, {
companyId,
correlationId,
aggregateType: 'Document',
aggregateId: doc.id,
eventType: 'DocumentIngested',
payload: {
channel: source,
document_id: doc.id,
mime_type: file.type,
size_bytes: file.buffer.byteLength,
},
actor: source === 'email' ? { type: 'system', id: 'resend-inbound' } : { type: 'user', id: userId },
occurredAt: new Date(),
})
} catch (err) {
console.error('[invoice-inbox] Failed to append DocumentIngested:', err)
}
// Classify with AI
let classificationResult
let classificationError: string | null = null
@@ -55,6 +100,30 @@ async function uploadAndClassify(
classificationError = err instanceof Error ? err.message : 'Classification failed'
}
// Audit: DocumentExtractionAttempted (fires whether classification succeeded or failed)
try {
await appendProcessingHistory(supabase, {
companyId,
correlationId,
aggregateType: 'Document',
aggregateId: doc.id,
eventType: 'DocumentExtractionAttempted',
payload: {
document_id: doc.id,
succeeded: !classificationError,
document_type: classificationResult?.documentType ?? null,
confidence: classificationResult?.confidence ? classificationResult.confidence / 100 : null,
llm_input_tokens: classificationResult?.usage?.inputTokens ?? 0,
llm_output_tokens: classificationResult?.usage?.outputTokens ?? 0,
error: classificationError,
},
actor: { type: 'llm', id: 'classify-document' },
occurredAt: new Date(),
})
} catch (err) {
console.error('[invoice-inbox] Failed to append DocumentExtractionAttempted:', err)
}
// Supplier matching
let matchedSupplierId: string | null = null
if (classificationResult?.documentType === 'supplier_invoice' && classificationResult.extractedData) {
@@ -104,20 +173,63 @@ async function uploadAndClassify(
email_from: emailMeta?.from || null,
email_subject: emailMeta?.subject || null,
email_received_at: emailMeta?.receivedAt || null,
email_body_text: emailMeta?.bodyText || null,
resend_email_id: emailMeta?.resendEmailId || null,
resend_attachment_id: emailMeta?.resendAttachmentId || null,
raw_email_payload: emailMeta?.messageId
? { messageId: emailMeta.messageId, filename: file.name }
: null,
error_message: classificationError,
correlation_id: correlationId,
})
.select('id, status, document_type, confidence, matched_supplier_id, error_message')
.select('*')
.single()
if (inboxError) throw new Error(`Failed to create inbox item: ${inboxError.message}`)
// Emit events for supplier invoices (non-blocking)
if (ctx && inbox.document_type === 'supplier_invoice') {
// Audit: DocumentClassified (only when classification succeeded)
if (!classificationError && classificationResult) {
try {
await ctx.emit({
await appendProcessingHistory(supabase, {
companyId,
correlationId,
aggregateType: 'Document',
aggregateId: doc.id,
eventType: 'DocumentClassified',
payload: {
document_id: doc.id,
inbox_item_id: inbox.id,
classification: classificationResult.documentType,
confidence: classificationResult.confidence / 100,
},
actor: { type: 'system', id: 'invoice-inbox' },
occurredAt: new Date(),
})
} catch (err) {
console.error('[invoice-inbox] Failed to append DocumentClassified:', err)
}
}
// Emit generic classified event for all document types.
// Always emit via eventBus directly so the webhook path (no ExtensionContext) still triggers handlers.
try {
await eventBus.emit({
type: 'inbox_item.classified',
payload: {
inboxItem: inbox as unknown as InvoiceInboxItem,
documentType: inbox.document_type,
confidence: inbox.confidence,
correlationId,
userId,
companyId,
},
})
} catch { /* non-blocking */ }
// Emit supplier-invoice-specific events (kept for backward compatibility)
if (inbox.document_type === 'supplier_invoice') {
try {
await eventBus.emit({
type: 'supplier_invoice.received',
payload: { inboxItem: inbox as unknown as InvoiceInboxItem, userId, companyId },
})
@@ -125,7 +237,7 @@ async function uploadAndClassify(
if (!classificationError && classificationResult?.confidence) {
try {
await ctx.emit({
await eventBus.emit({
type: 'supplier_invoice.extracted',
payload: { inboxItem: inbox as unknown as InvoiceInboxItem, confidence: classificationResult.confidence / 100, userId, companyId },
})
@@ -145,12 +257,28 @@ async function uploadAndClassify(
}
}
// ── Admin/owner check helper ──────────────────────────────────
async function isCompanyAdmin(
supabase: import('@supabase/supabase-js').SupabaseClient,
userId: string,
companyId: string
): Promise<boolean> {
const { data } = await supabase
.from('company_members')
.select('role')
.eq('company_id', companyId)
.eq('user_id', userId)
.maybeSingle()
return !!data && ['owner', 'admin'].includes(data.role)
}
// ── Extension definition ─────────────────────────────────────
export const invoiceInboxExtension: Extension = {
id: 'invoice-inbox',
name: 'Dokumentinkorg',
version: '1.0.0',
version: '2.0.0',
apiRoutes: [
// ── Upload ──────────────────────────────────────────────
@@ -212,7 +340,12 @@ export const invoiceInboxExtension: Extension = {
let query = ctx.supabase
.from('invoice_inbox_items')
.select('id, status, document_type, confidence, source, created_at, extracted_data, matched_supplier_id, document_id, email_from, email_subject, error_message')
.select(`
id, status, document_type, confidence, source, created_at, extracted_data,
matched_supplier_id, document_id, email_from, email_subject, error_message,
matched_transaction_id, match_confidence, match_method, match_reasoning,
matched_transaction:transactions!matched_transaction_id(id, description, amount, currency, date)
`)
.eq('company_id', ctx.companyId)
.order('created_at', { ascending: false })
.limit(limit)
@@ -252,251 +385,237 @@ export const invoiceInboxExtension: Extension = {
},
},
// ── Gmail OAuth: get auth URL ───────────────────────────
// ── Get this company's inbox address ────────────────────
{
method: 'GET',
path: '/gmail/auth',
path: '/inbox/address',
handler: async (_request: Request, ctx?: ExtensionContext) => {
if (!ctx) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
const clientId = process.env.GOOGLE_CLIENT_ID
if (!clientId) {
return NextResponse.json({ error: 'Google OAuth not configured' }, { status: 500 })
const domain = process.env.RESEND_INBOUND_DOMAIN
if (!domain) {
return NextResponse.json({ error: 'RESEND_INBOUND_DOMAIN not configured' }, { status: 503 })
}
if (!process.env.GMAIL_TOKEN_ENCRYPTION_KEY) {
return NextResponse.json({ error: 'GMAIL_TOKEN_ENCRYPTION_KEY is required' }, { status: 500 })
}
const appUrl = process.env.NEXT_PUBLIC_APP_URL || 'http://localhost:3000'
const redirectUri = `${appUrl}/api/extensions/ext/invoice-inbox/gmail/callback`
const state = encryptState({
companyId: ctx.companyId,
userId: ctx.userId,
exp: Date.now() + STATE_TTL_MS,
})
const params = new URLSearchParams({
client_id: clientId,
redirect_uri: redirectUri,
response_type: 'code',
scope: 'https://www.googleapis.com/auth/gmail.readonly https://www.googleapis.com/auth/gmail.labels',
access_type: 'offline',
prompt: 'consent',
state,
})
return NextResponse.json({ data: { authUrl: `https://accounts.google.com/o/oauth2/v2/auth?${params}` } })
},
},
// ── Gmail OAuth: callback (skipAuth — redirect from Google) ─
{
method: 'GET',
path: '/gmail/callback',
skipAuth: true,
handler: async (request: Request) => {
const url = new URL(request.url)
const code = url.searchParams.get('code')
const stateParam = url.searchParams.get('state')
const error = url.searchParams.get('error')
const appUrl = process.env.NEXT_PUBLIC_APP_URL || 'http://localhost:3000'
if (error) {
console.error('[gmail/callback] OAuth error:', error)
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_auth_denied`)
}
if (!code || !stateParam) {
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_missing_params`)
}
if (!process.env.GMAIL_TOKEN_ENCRYPTION_KEY) {
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_config_error`)
}
const state = decryptState(stateParam) as { companyId: string; userId: string; exp: number } | null
if (!state || Date.now() > state.exp) {
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_invalid_state`)
}
const { companyId, userId } = state
const redirectUri = `${appUrl}/api/extensions/ext/invoice-inbox/gmail/callback`
try {
// Exchange code for tokens
const tokenResponse = await fetch('https://oauth2.googleapis.com/token', {
method: 'POST',
headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
body: new URLSearchParams({
code,
client_id: process.env.GOOGLE_CLIENT_ID!,
client_secret: process.env.GOOGLE_CLIENT_SECRET!,
redirect_uri: redirectUri,
grant_type: 'authorization_code',
}),
const inbox = await getActiveInbox(ctx.supabase, ctx.companyId)
if (!inbox) {
return NextResponse.json({ error: 'No active inbox' }, { status: 404 })
}
return NextResponse.json({
data: {
address: composeInboxAddress(inbox.local_part, domain),
local_part: inbox.local_part,
status: inbox.status,
created_at: inbox.created_at,
},
})
if (!tokenResponse.ok) {
console.error('[gmail/callback] Token exchange failed:', await tokenResponse.text())
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_token_exchange`)
}
const tokens = await tokenResponse.json() as {
access_token: string; refresh_token?: string
}
if (!tokens.refresh_token) {
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_no_refresh_token`)
}
// Get user email
const profileResponse = await fetch('https://gmail.googleapis.com/gmail/v1/users/me/profile', {
headers: { Authorization: `Bearer ${tokens.access_token}` },
})
if (!profileResponse.ok) {
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_profile_error`)
}
const profile = await profileResponse.json() as { emailAddress: string }
// Create gnubok-processed label
let gmailLabelId: string | null = null
try {
const labelsResponse = await fetch('https://gmail.googleapis.com/gmail/v1/users/me/labels', {
headers: { Authorization: `Bearer ${tokens.access_token}` },
})
const labelsData = await labelsResponse.json() as { labels: { id: string; name: string }[] }
const existing = labelsData.labels?.find((l) => l.name === 'gnubok-processed')
if (existing) {
gmailLabelId = existing.id
} else {
const createLabelResponse = await fetch('https://gmail.googleapis.com/gmail/v1/users/me/labels', {
method: 'POST',
headers: {
Authorization: `Bearer ${tokens.access_token}`,
'Content-Type': 'application/json',
},
body: JSON.stringify({
name: 'gnubok-processed',
labelListVisibility: 'labelShow',
messageListVisibility: 'show',
}),
})
if (createLabelResponse.ok) {
const label = await createLabelResponse.json() as { id: string }
gmailLabelId = label.id
}
}
} catch (err) {
console.warn('[gmail/callback] Failed to create Gmail label:', err)
}
// Store connection (service-role — no auth cookie in callback)
const supabase = createClient(
process.env.NEXT_PUBLIC_SUPABASE_URL!,
process.env.SUPABASE_SERVICE_ROLE_KEY!
)
const { error: dbError } = await supabase
.from('email_connections')
.upsert(
{
company_id: companyId,
user_id: userId,
provider: 'gmail',
email_address: profile.emailAddress,
encrypted_token: encryptToken(tokens.refresh_token),
gmail_label_id: gmailLabelId,
status: 'active',
error_message: null,
},
{ onConflict: 'company_id,email_address' }
)
if (dbError) {
console.error('[gmail/callback] DB insert failed:', dbError)
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_db_error`)
}
console.log(`[gmail/callback] Gmail connected for ${profile.emailAddress} (company ${companyId})`)
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?gmail=connected`)
} catch (err) {
console.error('[gmail/callback] Unexpected error:', err)
return NextResponse.redirect(`${appUrl}/e/general/invoice-inbox?error=gmail_unexpected`)
return NextResponse.json(
{ error: err instanceof Error ? err.message : 'Failed to load inbox' },
{ status: 500 }
)
}
},
},
// ── Gmail: disconnect ───────────────────────────────────
// ── Rotate inbox address (admin/owner only) ─────────────
{
method: 'POST',
path: '/gmail/disconnect',
path: '/inbox/rotate',
handler: async (_request: Request, ctx?: ExtensionContext) => {
if (!ctx) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
const { error } = await ctx.supabase
.from('email_connections')
.delete()
.eq('company_id', ctx.companyId)
.eq('provider', 'gmail')
const domain = process.env.RESEND_INBOUND_DOMAIN
if (!domain) {
return NextResponse.json({ error: 'RESEND_INBOUND_DOMAIN not configured' }, { status: 503 })
}
if (error) return NextResponse.json({ error: error.message }, { status: 500 })
return NextResponse.json({ data: { disconnected: true } })
const isAdmin = await isCompanyAdmin(ctx.supabase, ctx.userId, ctx.companyId)
if (!isAdmin) return NextResponse.json({ error: 'Behörighet saknas.' }, { status: 403 })
try {
// rotate_company_inbox is SECURITY DEFINER and does its own
// auth.uid() role check. Call it through the user's JWT-bearing
// client — a service-role client has no session, so auth.uid()
// returns NULL and the in-RPC check always fails with 42501.
const newInbox = await rotateCompanyInbox(ctx.supabase, ctx.companyId)
return NextResponse.json({
data: {
address: composeInboxAddress(newInbox.local_part, domain),
local_part: newInbox.local_part,
status: newInbox.status,
},
})
} catch (err) {
console.error('[invoice-inbox/inbox/rotate] Failed:', err)
return NextResponse.json(
{ error: err instanceof Error ? err.message : 'Rotation failed' },
{ status: 500 }
)
}
},
},
// ── Gmail: connection status ────────────────────────────
{
method: 'GET',
path: '/gmail/status',
handler: async (_request: Request, ctx?: ExtensionContext) => {
if (!ctx) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
const { data, error } = await ctx.supabase
.from('email_connections')
.select('id, email_address, status, last_sync_at, error_message, created_at')
.eq('company_id', ctx.companyId)
.eq('provider', 'gmail')
if (error) return NextResponse.json({ error: error.message }, { status: 500 })
return NextResponse.json({ data: { connections: data || [] } })
},
},
// ── Gmail: manual scan trigger ──────────────────────────
// ── Resend Inbound webhook (Svix-signed, no user auth) ──
{
method: 'POST',
path: '/gmail/scan',
handler: async (_request: Request, ctx?: ExtensionContext) => {
if (!ctx) return NextResponse.json({ error: 'Unauthorized' }, { status: 401 })
if (!process.env.GMAIL_TOKEN_ENCRYPTION_KEY) {
return NextResponse.json({ error: 'GMAIL_TOKEN_ENCRYPTION_KEY is required' }, { status: 500 })
path: '/inbound',
skipAuth: true,
handler: async (request: Request) => {
const domain = process.env.RESEND_INBOUND_DOMAIN
if (!domain) {
console.error('[invoice-inbox/inbound] RESEND_INBOUND_DOMAIN not configured')
return NextResponse.json({ error: 'Inbound not configured' }, { status: 503 })
}
const { data: connections, error: connError } = await ctx.supabase
.from('email_connections')
.select('*')
.eq('company_id', ctx.companyId)
.eq('status', 'active')
const rawBody = await request.text()
if (connError) return NextResponse.json({ error: 'Failed to fetch connections' }, { status: 500 })
if (!connections?.length) {
return NextResponse.json({ data: { message: 'No active Gmail connections', scanned: 0 } })
// 1. Verify Svix signature
let event
try {
event = verifyInboundWebhook(rawBody, request.headers)
} catch (err) {
if (err instanceof ResendSignatureError) {
return NextResponse.json({ error: 'Invalid signature' }, { status: 401 })
}
console.error('[invoice-inbox/inbound] Verification error:', err)
return NextResponse.json({ error: 'Verification failed' }, { status: 500 })
}
let totalScanned = 0, totalClassified = 0, totalSkipped = 0, totalErrors = 0
for (const connection of connections) {
const result = await scanGmailConnection(ctx.supabase, connection, ctx.userId, ctx.companyId)
totalScanned += result.scanned
totalClassified += result.classified
totalSkipped += result.skipped
totalErrors += result.errors
// 2. Only process email.received events
if (!isEmailReceivedEvent(event)) {
return NextResponse.json({ data: { ignored: event.type } }, { status: 200 })
}
return NextResponse.json({
data: { scanned: totalScanned, classified: totalClassified, skipped: totalSkipped, errors: totalErrors },
})
const { email_id, to, from, subject, message_id, created_at } = event.data
// 3. Find the recipient that matches our domain
const localPart = extractLocalPartForDomain(to, domain)
if (!localPart) {
console.warn('[invoice-inbox/inbound] No recipient matched domain', { to, domain })
return NextResponse.json({ error: 'No matching recipient' }, { status: 404 })
}
// 4. Look up company_inbox (service role — webhook has no user session)
const serviceSupabase = createClient(
process.env.NEXT_PUBLIC_SUPABASE_URL!,
process.env.SUPABASE_SERVICE_ROLE_KEY!
)
const { data: inbox } = await serviceSupabase
.from('company_inboxes')
.select('id, company_id, status')
.eq('local_part', localPart)
.maybeSingle()
if (!inbox) {
return NextResponse.json({ error: 'Address not found' }, { status: 404 })
}
if (inbox.status !== 'active') {
// Deprecated or blocked → hard-bounce so Resend returns a 5xx to the sender
return NextResponse.json({ error: 'Address no longer active' }, { status: 410 })
}
// 5. Resolve a user_id for the inbox item (schema requires NOT NULL).
// Use company owner (companies.created_by).
const { data: company } = await serviceSupabase
.from('companies')
.select('created_by')
.eq('id', inbox.company_id)
.single()
if (!company?.created_by) {
console.error('[invoice-inbox/inbound] Company has no created_by', inbox.company_id)
return NextResponse.json({ error: 'Company owner missing' }, { status: 500 })
}
const userId = company.created_by
// 7. Fetch full email (body + attachment metadata)
let fullEmail
try {
fullEmail = await fetchReceivingEmail(email_id)
} catch (err) {
const message = err instanceof Error ? err.message : String(err)
console.error('[invoice-inbox/inbound] Failed to fetch received email:', err)
return NextResponse.json({ error: `Fetch failed: ${message}` }, { status: 500 })
}
const bodyText = fullEmail.text ?? null
const attachments = fullEmail.attachments ?? []
// 8. If no attachments, still log the email as an inbox item so the user sees it
if (attachments.length === 0) {
await serviceSupabase.from('invoice_inbox_items').insert({
company_id: inbox.company_id,
user_id: userId,
status: 'error',
source: 'email',
email_from: from,
email_subject: subject,
email_received_at: created_at,
email_body_text: bodyText,
resend_email_id: email_id,
document_type: 'unknown',
error_message: 'Email had no attachments',
raw_email_payload: { messageId: message_id },
})
return NextResponse.json({ data: { processed: 0, reason: 'no_attachments' } })
}
// 9. Download + classify each attachment (per-attachment idempotency)
const results: Array<{ attachment_id: string; inbox_item_id?: string; error?: string; duplicate?: boolean }> = []
for (const att of attachments) {
try {
// Skip if this (email_id, attachment_id) was already processed
const { data: existing } = await serviceSupabase
.from('invoice_inbox_items')
.select('id')
.eq('resend_email_id', email_id)
.eq('resend_attachment_id', att.id)
.maybeSingle()
if (existing) {
results.push({ attachment_id: att.id, inbox_item_id: existing.id, duplicate: true })
continue
}
const download = await fetchInboundAttachment(email_id, att.id)
if (!UPLOAD_ALLOWED_MIME_TYPES.has(download.contentType)) {
results.push({ attachment_id: att.id, error: `Unsupported type ${download.contentType}` })
continue
}
if (download.buffer.byteLength > MAX_FILE_SIZE) {
results.push({ attachment_id: att.id, error: 'Attachment too large' })
continue
}
const result = await uploadAndClassify(
serviceSupabase,
userId,
inbox.company_id,
{ name: download.filename, buffer: download.buffer, type: download.contentType },
'email',
{
from,
subject,
receivedAt: created_at,
messageId: message_id,
bodyText,
resendEmailId: email_id,
resendAttachmentId: att.id,
}
)
results.push({ attachment_id: att.id, inbox_item_id: result.inbox_item_id })
} catch (err) {
console.error('[invoice-inbox/inbound] Attachment processing failed:', err)
results.push({
attachment_id: att.id,
error: err instanceof Error ? err.message : 'Unknown error',
})
}
}
return NextResponse.json({ data: { processed: results.length, results } })
},
},
@@ -84,11 +84,17 @@ Kontext:
Instruktioner:
- Klassificera dokumenttypen
- Extrahera alla synliga fält — returnera null för fält som inte kan utläsas
- Belopp ska vara positiva tal med max 2 decimaler
- Belopp med max 2 decimaler. Totaler (amount_excl_vat, amount_incl_vat, vat_amount) är alltid positiva. line_items.amount är normalt positiva men KAN vara negativa för rabattrader.
- Datum i ISO-format (YYYY-MM-DD)
- Momssats som heltal (0, 6, 12, eller 25)
- Ge en confidence-poäng 0-100 för hur säker du är på klassificeringen och extraktionen
KRITISKT — Totaler och rabatter:
- amount_incl_vat MÅSTE motsvara fakturans slutbelopp ("Amount due", "Att betala", "Totalt", "Subtotal" när moms saknas). Summera ALDRIG delrader om fakturan anger ett explicit totalbelopp — använd det.
- Om en rad har en rabatt/discount under sig (t.ex. "Discount (-$10.00)", "Rabatt -100 kr"), extrahera radens NETTO-belopp (brutto − rabatt), INTE bruttobeloppet. Exempel: en rad "Compute Hours $50.68" följd av "Discount -$10.00" → line_items.amount = 40.68, inte 50.68.
- Summan av line_items.amount + moms MÅSTE bli lika med amount_incl_vat. Om det inte stämmer har du antingen missat en rabatt eller dubbelräknat en rad — kontrollera och justera.
- Om du är osäker på hur rabatter ska fördelas, lägg en enskild negativ "Rabatt"-rad i line_items så summan blir rätt.
Anropa ALLTID verktyget classify_document med resultatet.`
// ── Tool schema for structured output ────────────────────────
@@ -311,7 +317,12 @@ function mapToClassificationResult(
if (documentType === 'supplier_invoice') {
const extractedData = mapToInvoiceExtraction(raw)
if (!extractedData) return null
return { documentType, extractedData, confidence, rawResponse: raw, usage }
// mapToInvoiceExtraction caps its own confidence to 50% when line items
// don't reconcile with amount_incl_vat. Propagate that cap to the outer
// confidence so invoice_inbox_items.confidence (used by the UI badge)
// also reflects the reconciliation failure.
const effectiveConfidence = Math.min(confidence, Math.round(extractedData.confidence * 100))
return { documentType, extractedData, confidence: effectiveConfidence, rawResponse: raw, usage }
}
if (documentType === 'receipt') {
@@ -323,9 +334,32 @@ function mapToClassificationResult(
return null
}
// Sum of line items must approximately match the extracted subtotal.
// Allows 0.02 * max(|subtotal|, 1) tolerance for rounding. Returns true if the extraction
// is internally consistent; false signals the model skipped a discount or double-counted.
function invoiceTotalsAreConsistent(raw: Record<string, unknown>): boolean {
const subtotal = Number(raw.amount_excl_vat)
const total = Number(raw.amount_incl_vat)
const vat = Number(raw.vat_amount) || 0
const items = Array.isArray(raw.line_items) ? raw.line_items : []
if (!items.length) return true // nothing to compare against
if (!isFinite(subtotal) && !isFinite(total)) return true
const sumOfLines = items.reduce((acc, item) => {
if (typeof item !== 'object' || item === null) return acc
const amount = Number((item as Record<string, unknown>).amount)
return acc + (isFinite(amount) ? amount : 0)
}, 0)
const anchor = isFinite(subtotal) ? subtotal : total - vat
const tolerance = Math.max(0.02, Math.abs(anchor) * 0.02)
return Math.abs(sumOfLines - anchor) <= tolerance
}
function mapToInvoiceExtraction(raw: Record<string, unknown>): InvoiceExtractionResult | null {
const lineItems = mapInvoiceLineItems(raw.line_items)
const vatBreakdown = mapVatBreakdown(raw.vat_breakdown)
const totalsConsistent = invoiceTotalsAreConsistent(raw)
const result: InvoiceExtractionResult = {
supplier: {
@@ -350,7 +384,11 @@ function mapToInvoiceExtraction(raw: Record<string, unknown>): InvoiceExtraction
total: roundAmount(raw.amount_incl_vat),
},
vatBreakdown,
confidence: Math.min(1, Math.max(0, Number(raw.confidence) / 100 || 0)),
// If totals don't reconcile with line items, cap confidence at 50% so the UI flags it
confidence: (() => {
const raw_conf = Math.min(1, Math.max(0, Number(raw.confidence) / 100 || 0))
return totalsConsistent ? raw_conf : Math.min(raw_conf, 0.5)
})(),
}
return result
@@ -472,7 +510,7 @@ async function retryWithCorrection(
- document_type måste vara ett av: supplier_invoice, receipt, government_letter, unknown
- Datum i format YYYY-MM-DD
- Momssatser måste vara 0, 6, 12, eller 25
- Belopp ska vara positiva tal
- Totaler: amount_incl_vat ska vara fakturans slutbelopp. Summa av line_items.amount + vat_amount MÅSTE bli lika med amount_incl_vat. Om rader har rabatt under sig, använd NETTO-beloppet per rad, eller lägg till en separat negativ rabattrad så summan stämmer.
Försök igen med korrigerad data.`,
},
],
@@ -1,76 +0,0 @@
import crypto from 'crypto'
const ALGORITHM = 'aes-256-gcm'
export function getGmailEncryptionKey(): Buffer {
const secret = process.env.GMAIL_TOKEN_ENCRYPTION_KEY
if (!secret) throw new Error('GMAIL_TOKEN_ENCRYPTION_KEY is required')
return crypto.createHash('sha256').update(secret).digest()
}
export function encryptState(payload: Record<string, unknown>): string {
const key = getGmailEncryptionKey()
const iv = crypto.randomBytes(12)
const cipher = crypto.createCipheriv(ALGORITHM, key, iv)
const json = JSON.stringify(payload)
const encrypted = Buffer.concat([cipher.update(json, 'utf8'), cipher.final()])
const tag = cipher.getAuthTag()
return Buffer.concat([iv, tag, encrypted]).toString('base64url')
}
export function decryptState(encoded: string): Record<string, unknown> | null {
try {
const key = getGmailEncryptionKey()
const combined = Buffer.from(encoded, 'base64url')
const iv = combined.subarray(0, 12)
const tag = combined.subarray(12, 28)
const encrypted = combined.subarray(28)
const decipher = crypto.createDecipheriv(ALGORITHM, key, iv)
decipher.setAuthTag(tag)
const decrypted = Buffer.concat([decipher.update(encrypted), decipher.final()])
return JSON.parse(decrypted.toString('utf8'))
} catch {
return null
}
}
export function encryptToken(plaintext: string): string {
const key = getGmailEncryptionKey()
const iv = crypto.randomBytes(12)
const cipher = crypto.createCipheriv(ALGORITHM, key, iv)
const encrypted = Buffer.concat([cipher.update(plaintext, 'utf8'), cipher.final()])
const tag = cipher.getAuthTag()
return Buffer.concat([iv, tag, encrypted]).toString('base64url')
}
export function decryptToken(encoded: string): string | null {
try {
const key = getGmailEncryptionKey()
const combined = Buffer.from(encoded, 'base64url')
const iv = combined.subarray(0, 12)
const tag = combined.subarray(12, 28)
const encrypted = combined.subarray(28)
const decipher = crypto.createDecipheriv(ALGORITHM, key, iv)
decipher.setAuthTag(tag)
const decrypted = Buffer.concat([decipher.update(encrypted), decipher.final()])
return decrypted.toString('utf8')
} catch {
return null
}
}
export async function refreshAccessToken(refreshToken: string): Promise<string | null> {
const response = await fetch('https://oauth2.googleapis.com/token', {
method: 'POST',
headers: { 'Content-Type': 'application/x-www-form-urlencoded' },
body: new URLSearchParams({
refresh_token: refreshToken,
client_id: process.env.GOOGLE_CLIENT_ID!,
client_secret: process.env.GOOGLE_CLIENT_SECRET!,
grant_type: 'refresh_token',
}),
})
if (!response.ok) return null
const data = await response.json() as { access_token: string }
return data.access_token
}
@@ -1,336 +0,0 @@
import type { SupabaseClient } from '@supabase/supabase-js'
import { uploadDocument, computeSHA256 } from '@/lib/core/documents/document-service'
import { classifyDocument } from './classify-document'
import { decryptToken, refreshAccessToken } from './gmail-helpers'
import type { InvoiceExtractionResult } from '@/types'
const MIN_ATTACHMENT_SIZE = 3_000
const MAX_MESSAGES = 30
const ALLOWED_MIME_TYPES = new Set([
'application/pdf',
'application/octet-stream',
'image/jpeg',
'image/png',
'image/heic',
'image/heif',
'image/webp',
])
const SKIP_EXTENSIONS = new Set([
'ics', 'vcf', 'html', 'htm', 'zip', 'rar', 'gz',
'csv', 'json', 'xml', 'txt', 'eml', 'msg',
'mp3', 'mp4', 'mov', 'avi', 'wav',
])
interface GmailMessage {
id: string
payload: {
headers: { name: string; value: string }[]
parts?: GmailPart[]
mimeType: string
body?: { attachmentId?: string; size?: number; data?: string }
}
internalDate: string
}
interface GmailPart {
mimeType: string
filename: string
body: { attachmentId?: string; size?: number; data?: string }
parts?: GmailPart[]
}
function getHeader(message: GmailMessage, name: string): string | null {
return message.payload.headers.find(
(h) => h.name.toLowerCase() === name.toLowerCase()
)?.value ?? null
}
function collectAttachments(parts: GmailPart[] | undefined): GmailPart[] {
if (!parts) return []
const result: GmailPart[] = []
for (const part of parts) {
if (part.filename && part.body?.attachmentId) {
result.push(part)
}
if (part.parts) {
result.push(...collectAttachments(part.parts))
}
}
return result
}
function resolveActualMimeType(mimeType: string, filename: string): string | null {
const ext = filename.split('.').pop()?.toLowerCase()
if (ext && SKIP_EXTENSIONS.has(ext)) return null
if (mimeType === 'application/octet-stream') {
const extMap: Record<string, string> = {
pdf: 'application/pdf',
jpg: 'image/jpeg',
jpeg: 'image/jpeg',
png: 'image/png',
webp: 'image/webp',
}
return ext && extMap[ext] ? extMap[ext] : null
}
return mimeType
}
export interface ScanResult {
scanned: number
classified: number
skipped: number
errors: number
}
interface EmailConnection {
id: string
company_id: string
encrypted_token: string
last_sync_at: string | null
gmail_label_id: string | null
}
export async function scanGmailConnection(
supabase: SupabaseClient,
connection: EmailConnection,
userId: string,
companyId: string
): Promise<ScanResult> {
const result: ScanResult = { scanned: 0, classified: 0, skipped: 0, errors: 0 }
const seenFileHashes = new Set<string>()
const refreshToken = decryptToken(connection.encrypted_token)
if (!refreshToken) {
await supabase
.from('email_connections')
.update({ status: 'error', error_message: 'Failed to decrypt refresh token' })
.eq('id', connection.id)
result.errors++
return result
}
const accessToken = await refreshAccessToken(refreshToken)
if (!accessToken) {
await supabase
.from('email_connections')
.update({ status: 'revoked', error_message: 'Token refresh failed — user may have revoked access' })
.eq('id', connection.id)
result.errors++
return result
}
// Build Gmail search query
let afterDate: string
if (connection.last_sync_at) {
const d = new Date(connection.last_sync_at)
afterDate = `${d.getFullYear()}/${String(d.getMonth() + 1).padStart(2, '0')}/${String(d.getDate()).padStart(2, '0')}`
} else {
const d = new Date(Date.now() - 14 * 24 * 60 * 60 * 1000)
afterDate = `${d.getFullYear()}/${String(d.getMonth() + 1).padStart(2, '0')}/${String(d.getDate()).padStart(2, '0')}`
}
let query = `has:attachment after:${afterDate}`
if (connection.gmail_label_id) {
query += ' -label:gnubok-processed'
}
const listUrl = `https://gmail.googleapis.com/gmail/v1/users/me/messages?q=${encodeURIComponent(query)}&maxResults=${MAX_MESSAGES}`
const listResponse = await fetch(listUrl, {
headers: { Authorization: `Bearer ${accessToken}` },
})
if (!listResponse.ok) {
console.error('[gmail/scan] Failed to list messages:', await listResponse.text())
result.errors++
return result
}
const listData = await listResponse.json() as { messages?: { id: string }[] }
const messageIds = listData.messages || []
for (const { id: messageId } of messageIds) {
try {
const msgResponse = await fetch(
`https://gmail.googleapis.com/gmail/v1/users/me/messages/${messageId}?format=full`,
{ headers: { Authorization: `Bearer ${accessToken}` } }
)
if (!msgResponse.ok) continue
const message = await msgResponse.json() as GmailMessage
const emailFrom = getHeader(message, 'From')
const emailSubject = getHeader(message, 'Subject')
const emailDate = message.internalDate
? new Date(parseInt(message.internalDate)).toISOString()
: null
const attachments = collectAttachments(message.payload.parts)
for (const attachment of attachments) {
if (!attachment.body.attachmentId) continue
if ((attachment.body.size ?? 0) < MIN_ATTACHMENT_SIZE) continue
if (!ALLOWED_MIME_TYPES.has(attachment.mimeType)) continue
const resolvedMimeType = resolveActualMimeType(attachment.mimeType, attachment.filename)
if (!resolvedMimeType) continue
// Deduplicate by message ID + filename
const { data: existing } = await supabase
.from('invoice_inbox_items')
.select('id')
.eq('company_id', companyId)
.eq('source', 'email')
.filter('raw_email_payload->>messageId', 'eq', messageId)
.filter('raw_email_payload->>filename', 'eq', attachment.filename)
.limit(1)
.maybeSingle()
if (existing) {
result.skipped++
continue
}
// Download attachment
const attResponse = await fetch(
`https://gmail.googleapis.com/gmail/v1/users/me/messages/${messageId}/attachments/${attachment.body.attachmentId}`,
{ headers: { Authorization: `Bearer ${accessToken}` } }
)
if (!attResponse.ok) {
result.errors++
continue
}
const attData = await attResponse.json() as { data: string }
const fileBuffer = Buffer.from(attData.data, 'base64url')
// Deduplicate by file content hash
const fileHash = await computeSHA256(fileBuffer.buffer.slice(
fileBuffer.byteOffset,
fileBuffer.byteOffset + fileBuffer.byteLength
))
if (seenFileHashes.has(fileHash)) {
result.skipped++
continue
}
const { data: existingByHash } = await supabase
.from('document_attachments')
.select('id')
.eq('company_id', companyId)
.eq('sha256_hash', fileHash)
.limit(1)
.maybeSingle()
if (existingByHash) {
result.skipped++
seenFileHashes.add(fileHash)
continue
}
seenFileHashes.add(fileHash)
// Store in WORM archive
const doc = await uploadDocument(supabase, userId, companyId, {
name: attachment.filename,
buffer: fileBuffer.buffer.slice(
fileBuffer.byteOffset,
fileBuffer.byteOffset + fileBuffer.byteLength
),
type: resolvedMimeType,
}, {
upload_source: 'email',
})
// Classify
let classificationResult
let classificationError: string | null = null
try {
classificationResult = await classifyDocument({
fileBuffer,
mimeType: resolvedMimeType,
fileName: attachment.filename,
})
} catch (err) {
classificationError = err instanceof Error ? err.message : 'Classification failed'
console.error('[gmail/scan] Classification failed:', err)
}
// Find matching supplier
let matchedSupplierId: string | null = null
if (classificationResult?.documentType === 'supplier_invoice' && classificationResult.extractedData) {
const extractedData = classificationResult.extractedData as InvoiceExtractionResult
const orgNumber = extractedData.supplier?.orgNumber
if (orgNumber) {
const normalized = orgNumber.replace(/\D/g, '')
const { data: supplierByOrg } = await supabase
.from('suppliers')
.select('id')
.eq('company_id', companyId)
.eq('org_number', normalized)
.limit(1)
.maybeSingle()
if (supplierByOrg) matchedSupplierId = supplierByOrg.id
}
}
// Create inbox item
await supabase.from('invoice_inbox_items').insert({
company_id: companyId,
user_id: userId,
status: classificationError ? 'error' : 'ready',
source: 'email',
document_id: doc.id,
document_type: classificationResult?.documentType || 'unknown',
extracted_data: classificationResult?.extractedData || null,
raw_llm_response: classificationResult?.rawResponse || null,
confidence: classificationResult?.confidence
? classificationResult.confidence / 100
: null,
matched_supplier_id: matchedSupplierId,
email_from: emailFrom,
email_subject: emailSubject,
email_received_at: emailDate,
raw_email_payload: { messageId, filename: attachment.filename },
error_message: classificationError,
})
if (classificationError) {
result.errors++
} else {
result.classified++
}
result.scanned++
}
// Label message as processed
if (connection.gmail_label_id) {
try {
await fetch(
`https://gmail.googleapis.com/gmail/v1/users/me/messages/${messageId}/modify`,
{
method: 'POST',
headers: {
Authorization: `Bearer ${accessToken}`,
'Content-Type': 'application/json',
},
body: JSON.stringify({ addLabelIds: [connection.gmail_label_id] }),
}
)
} catch {
// Non-blocking
}
}
} catch (err) {
console.error('[gmail/scan] Error processing message:', err)
result.errors++
}
}
// Update last_sync_at
await supabase
.from('email_connections')
.update({ last_sync_at: new Date().toISOString(), error_message: null })
.eq('id', connection.id)
return result
}
@@ -0,0 +1,43 @@
import type { SupabaseClient } from '@supabase/supabase-js'
import type { CompanyInbox } from '@/types'
export function composeInboxAddress(localPart: string, domain: string): string {
return `${localPart}@${domain}`
}
export async function getActiveInbox(
supabase: SupabaseClient,
companyId: string
): Promise<CompanyInbox | null> {
const { data, error } = await supabase
.from('company_inboxes')
.select('*')
.eq('company_id', companyId)
.eq('status', 'active')
.maybeSingle()
if (error) throw new Error(`Failed to load inbox: ${error.message}`)
return (data as CompanyInbox | null) ?? null
}
// Rotate the company's inbox address. Delegates to the rotate_company_inbox
// RPC so the three steps (deprecate, generate, insert) run inside a single
// Postgres transaction; a failure on any step rolls the whole thing back
// and the company is never left without an active inbox.
export async function rotateCompanyInbox(
supabase: SupabaseClient,
companyId: string
): Promise<CompanyInbox> {
const { data, error } = await supabase
.rpc('rotate_company_inbox', { p_company_id: companyId })
if (error || !data) {
throw new Error(`Failed to rotate inbox: ${error?.message ?? 'no data'}`)
}
// The RPC returns a single row (SETOF company_inboxes).
const row = Array.isArray(data) ? data[0] : data
if (!row) throw new Error('Failed to rotate inbox: RPC returned no row')
return row as CompanyInbox
}
@@ -0,0 +1,99 @@
import { Resend } from 'resend'
import type { EmailReceivedEvent, GetReceivingEmailResponseSuccess, WebhookEventPayload } from 'resend'
export type ResendInboundEvent = EmailReceivedEvent
export type ResendReceivedEmail = GetReceivingEmailResponseSuccess
export interface ResendAttachmentDownload {
id: string
filename: string
contentType: string
buffer: ArrayBuffer
}
function getResend(): Resend {
const apiKey = process.env.RESEND_API_KEY
if (!apiKey) throw new Error('RESEND_API_KEY is required')
return new Resend(apiKey)
}
export class ResendSignatureError extends Error {
constructor(message: string) {
super(message)
this.name = 'ResendSignatureError'
}
}
// Verifies the Svix-signed webhook payload using the RESEND_INBOUND_WEBHOOK_SECRET.
// Throws ResendSignatureError on failure, returns the parsed event on success.
export function verifyInboundWebhook(rawBody: string, requestHeaders: Headers): WebhookEventPayload {
const secret = process.env.RESEND_INBOUND_WEBHOOK_SECRET
if (!secret) throw new Error('RESEND_INBOUND_WEBHOOK_SECRET is required')
// Resend's verify() expects Svix headers in a specific shape, not the raw Fetch Headers.
const svixHeaders = {
id: requestHeaders.get('svix-id') ?? '',
timestamp: requestHeaders.get('svix-timestamp') ?? '',
signature: requestHeaders.get('svix-signature') ?? '',
}
const resend = getResend()
try {
return resend.webhooks.verify({ payload: rawBody, headers: svixHeaders, webhookSecret: secret })
} catch (err) {
throw new ResendSignatureError(err instanceof Error ? err.message : 'Invalid signature')
}
}
// Fetches the full received email (body, headers, attachment metadata) by email_id.
export async function fetchReceivingEmail(emailId: string): Promise<ResendReceivedEmail> {
const resend = getResend()
const { data, error } = await resend.emails.receiving.get(emailId)
if (error || !data) {
throw new Error(`Failed to fetch received email ${emailId}: ${error?.message ?? 'no data'}`)
}
return data
}
// Fetches a single attachment's bytes via its short-lived download_url.
export async function fetchInboundAttachment(
emailId: string,
attachmentId: string
): Promise<ResendAttachmentDownload> {
const resend = getResend()
const { data, error } = await resend.emails.receiving.attachments.get({ emailId, id: attachmentId })
if (error || !data) {
throw new Error(`Failed to fetch attachment ${attachmentId}: ${error?.message ?? 'no data'}`)
}
const response = await fetch(data.download_url)
if (!response.ok) {
throw new Error(`Download URL returned ${response.status} for attachment ${attachmentId}`)
}
const buffer = await response.arrayBuffer()
return {
id: data.id,
filename: data.filename ?? `attachment-${data.id}`,
contentType: data.content_type,
buffer,
}
}
// Parses the first recipient whose domain matches our configured inbound domain,
// returning just the local_part. Returns null if no match.
export function extractLocalPartForDomain(recipients: string[], domain: string): string | null {
const normalized = domain.toLowerCase()
for (const addr of recipients) {
const match = addr.match(/^\s*([^@\s]+)@([^@\s]+?)\s*$/)
if (!match) continue
const [, localPart, addrDomain] = match
if (addrDomain.toLowerCase() === normalized) return localPart.toLowerCase()
}
return null
}
export function isEmailReceivedEvent(event: WebhookEventPayload): event is EmailReceivedEvent {
return event.type === 'email.received'
}
+10 -3
View File
@@ -4,8 +4,15 @@
"exportName": "invoiceInboxExtension",
"entryPoint": "@/extensions/general/invoice-inbox",
"workspace": "@/components/extensions/general/InvoiceInboxWorkspace",
"requiredEnvVars": ["AWS_ACCESS_KEY_ID", "AWS_SECRET_ACCESS_KEY", "AWS_REGION"],
"optionalEnvVars": ["BEDROCK_MODEL_ID", "BEDROCK_MAX_TOKENS", "GOOGLE_CLIENT_ID", "GOOGLE_CLIENT_SECRET", "GMAIL_TOKEN_ENCRYPTION_KEY"],
"requiredEnvVars": [
"AWS_ACCESS_KEY_ID",
"AWS_SECRET_ACCESS_KEY",
"AWS_REGION",
"RESEND_API_KEY",
"RESEND_INBOUND_DOMAIN",
"RESEND_INBOUND_WEBHOOK_SECRET"
],
"optionalEnvVars": ["BEDROCK_MODEL_ID", "BEDROCK_MAX_TOKENS"],
"npmDependencies": ["@aws-sdk/client-bedrock-runtime"],
"definition": {
"name": "Dokumentinkorg",
@@ -15,6 +22,6 @@
"hasOwnData": true,
"readsCoreTables": ["document_attachments", "suppliers", "transactions"],
"description": "AI-klassificering och extraktion av leverantörsfakturor och kvitton",
"longDescription": "Ladda upp eller ta emot dokument via Gmail. AI klassificerar dokumenttyp, extraherar strukturerad data (leverantör, belopp, moms) och matchar mot transaktioner. Kräver AWS Bedrock-åtkomst."
"longDescription": "Varje bolag får en unik fakturainkorg-adress. Fakturor som skickas dit fångas automatiskt, klassificeras med AI (leverantör, belopp, moms) och matchas mot transaktioner. Kräver AWS Bedrock och Resend."
}
}