feat(documents): dedupe intake channels on content, not just provenance (#1528)

* feat(documents): dedupe intake channels on content, not just provenance

Every ingestion path already computes and stores sha256_hash, but only
WhatsApp ever read it back: the manual upload, Resend inbound, and mail
hunt deduped on provenance keys alone (or not at all), so the same
receipt forwarded to two inboxes, re-hunted by a sweep, or uploaded
twice became a second archived document and a second inbox item. With
the hunt live and three channels feeding one inbox, that is an unbounded
duplicate generator (flows plan, prerequisite PR 1).

uploadDocument gains an opt-in dedupeByContent flag: before storing, it
looks for a current-version document in the same company with the same
SHA-256 and returns it (marked deduplicated) instead of archiving a
copy. Opt-in because archival callers must store what they produced even
when bytes repeat; the SELECT-then-insert race is accepted exactly as in
the WhatsApp intake precedent.

uploadAndExtract turns the flag on for every inbox channel. On a hit it
adopts the oldest inbox item for that document, so callers always
receive a real inbox_item_id, and only files a new item (against the
EXISTING document) when the content entered the archive outside the
inbox. The mail hunt skips outright: its provenance key catches the same
message re-hunted, the content check catches the same receipt arriving
through another inbox. WhatsApp keeps its own pre-check, which also
drives the duplicate reply to the sender.

No migration: the hash column and its index have existed since the
original archive schema.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(documents): review round: fail closed, adopt-or-file in the hunt, audit trail

CodeRabbit: both dedupe lookups failed OPEN, so a transient DB error
would silently archive the duplicate the feature exists to prevent; both
now throw before anything is stored, and a regression test locks it.
The ingest test also asserts the dedupeByContent flag in the production
call, so removing the flag fails the suite.

Swedish compliance review, both findings real: (1) the mail hunt's
unconditional skip could swallow a receipt whose content matches a
document that never passed the inbox (a manually attached copy), leaving
an affärshändelse without underlag routing (BFL 5 kap): the hunt now
mirrors the funnel's adopt-or-file semantics, skipping only when an
inbox item already carries the document and otherwise filing an item
against the EXISTING document. (2) The skip decision now lands in
behandlingshistorik as DocumentDuplicateSkipped (BFNAR 2013:2 kap 8),
not just the app log.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(receipt-hunt): keep the audit payload pseudonymous; lock the skip trail in tests

Review round 2. The DocumentDuplicateSkipped payload carried the mailbox
address, violating the processing-history contract (pseudonymous IDs
only, never emails); the digit-shaped PII validator would not have
caught it, which is exactly why the contract must hold at the call site.
Which mailbox first delivered the receipt is already on the existing
item's channel_context. Tests now assert the audit event lands with the
right identifiers and no address, and that a history outage still skips
rather than filing a duplicate.

Not changed: a duplicate-lookup error still soft-fails the attachment
(warn + continue). Aborting the candidate would contradict this
function's documented contract (one bad message never costs the night's
hunt); fail-closed holds either way, and the next sweep retries since
no item was filed.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Jakob Wennberg <311770904+jakobwennberg-oss@users.noreply.github.com>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Jakob Wennberg
2026-08-12 15:55:12 +02:00
committed by GitHub
co-authored by Claude Fable 5 Jakob Wennberg
parent 05b0c57ecd
commit 845add4573
6 changed files with 297 additions and 4 deletions
@@ -307,6 +307,72 @@ describe('uploadDocument', () => {
)
})
it('returns the existing document instead of storing a copy when dedupeByContent hits', async () => {
const existing = makeDocumentAttachment({ id: 'doc-orig', sha256_hash: 'same' })
results = [
{ data: [existing], error: null }, // dedupe lookup
]
const handler = vi.fn()
eventBus.on('document.uploaded', handler)
const upload = vi.fn().mockResolvedValue({ data: {}, error: null })
const supabase = makeClient({ upload })
const result = await uploadDocument(supabase as never, 'user-1', 'company-1', {
name: 'kvitto.pdf',
buffer: pdfBuffer(),
type: 'application/pdf',
}, { dedupeByContent: true })
expect(result.id).toBe('doc-orig')
expect(result.deduplicated).toBe(true)
// Nothing reaches storage and no document.uploaded fires: the archive
// already holds this content, so re-extraction must not run either.
expect(upload).not.toHaveBeenCalled()
expect(handler).not.toHaveBeenCalled()
})
it('rejects when the dedupe lookup itself fails, touching nothing', async () => {
// Fail closed: treating a broken lookup as "no match" would archive the
// duplicate the flag exists to prevent, silently, on transient DB errors.
results = [{ data: null, error: { message: 'connection reset' } }]
const handler = vi.fn()
eventBus.on('document.uploaded', handler)
const upload = vi.fn().mockResolvedValue({ data: {}, error: null })
const supabase = makeClient({ upload })
await expect(
uploadDocument(supabase as never, 'user-1', 'company-1', {
name: 'kvitto.pdf',
buffer: pdfBuffer(),
type: 'application/pdf',
}, { dedupeByContent: true }),
).rejects.toThrow(/dedupe lookup failed/i)
expect(upload).not.toHaveBeenCalled()
expect(handler).not.toHaveBeenCalled()
})
it('stores normally when dedupeByContent finds no match', async () => {
results = [
{ data: [], error: null }, // dedupe lookup: miss
{ data: makeDocumentAttachment({ id: 'doc-new' }), error: null }, // insert
]
const upload = vi.fn().mockResolvedValue({ data: {}, error: null })
const supabase = makeClient({ upload })
const result = await uploadDocument(supabase as never, 'user-1', 'company-1', {
name: 'kvitto.pdf',
buffer: pdfBuffer(),
type: 'application/pdf',
}, { dedupeByContent: true })
expect(result.id).toBe('doc-new')
expect(result.deduplicated).toBeUndefined()
expect(upload).toHaveBeenCalledOnce()
})
it('writes to the company-scoped key, not the legacy uploader-scoped key', async () => {
results = [{ data: makeDocumentAttachment({ id: 'doc-1' }), error: null }]
+33 -1
View File
@@ -597,8 +597,18 @@ export async function uploadDocument(
upload_source?: DocumentUploadSource
journal_entry_id?: string
journal_entry_line_id?: string
/**
* Content dedupe for intake channels: before storing, look for a
* current-version document in the same company with the same SHA-256 and
* return it (marked `deduplicated`) instead of archiving a copy. Opt-in,
* because archival callers (sent invoices, filings, bank exports) must
* store what they produced even when the bytes repeat. SELECT-then-insert
* leaves a small concurrent-upload race, accepted exactly as in the
* WhatsApp intake precedent: the loser stores a copy, nothing corrupts.
*/
dedupeByContent?: boolean
} = {}
): Promise<DocumentAttachment> {
): Promise<DocumentAttachment & { deduplicated?: boolean }> {
await ensureDocumentsBucket()
// Reject corrupt uploads at the boundary: see validateDocumentMagicBytes.
@@ -610,6 +620,28 @@ export async function uploadDocument(
// Compute SHA-256 hash
const sha256Hash = await computeSHA256(file.buffer)
if (metadata.dedupeByContent) {
// Oldest current-version match wins so repeated deliveries keep
// converging on the same archived original. Pre-dedupe data can hold
// several identical documents, hence limit(1) rather than maybeSingle.
const { data: existing, error: dedupeError } = await supabase
.from('document_attachments')
.select('*')
.eq('company_id', companyId)
.eq('sha256_hash', sha256Hash)
.eq('is_current_version', true)
.order('created_at', { ascending: true })
.limit(1)
if (dedupeError) {
// Fail closed: treating a broken lookup as "no match" would archive
// the duplicate this flag exists to prevent, silently, on every
// transient DB error. Intake callers (webhooks, sweeps) retry.
throw new Error(`Content dedupe lookup failed: ${dedupeError.message}`)
}
const hit = (existing as DocumentAttachment[] | null)?.[0]
if (hit) return { ...hit, deduplicated: true }
}
// Company-scoped storage key: the tenant id must be IN the key so the
// storage RLS policy can revoke access when a membership is removed.
const storagePath = buildDocumentStoragePath(companyId, userId, file.name)
+72 -2
View File
@@ -11,6 +11,11 @@ vi.mock('@/lib/core/documents/document-service', () => ({
uploadDocument: (...args: unknown[]) => mockUploadDocument(...args),
}))
const mockAppendHistory = vi.fn()
vi.mock('@/lib/processing-history/append', () => ({
appendProcessingHistory: (...args: unknown[]) => mockAppendHistory(...args),
}))
const mockFetchAttachment = vi.fn()
vi.mock('@/lib/mail-search/service', () => ({
getMailSearchService: () => ({
@@ -38,7 +43,12 @@ function candidate(overrides: Partial<MailCandidate> = {}): MailCandidate {
}
/** Table-dispatching Supabase stand-in with a settable existing-row answer. */
function mockSupabase(existing: { id: string } | null, insertResult: { data?: unknown; error?: unknown } = {}) {
function mockSupabase(
existing: { id: string } | null,
insertResult: { data?: unknown; error?: unknown } = {},
laterInboxAnswers: Array<{ id: string } | null> = [],
) {
const inboxQueue: Array<{ id: string } | null> = [existing, ...laterInboxAnswers]
const inserted: Array<Record<string, unknown>> = []
const client = {
from(table: string) {
@@ -48,7 +58,7 @@ function mockSupabase(existing: { id: string } | null, insertResult: { data?: un
Promise.resolve(
table === 'document_attachments'
? { data: { extracted_data: { total_amount: 425 } }, error: null }
: { data: existing, error: null },
: { data: inboxQueue.length ? inboxQueue.shift() ?? null : null, error: null },
),
)
chain.insert = vi.fn((row: Record<string, unknown>) => {
@@ -119,6 +129,66 @@ describe('ingestMailCandidate', () => {
await expect(ingestMailCandidate(client, 'co-1', 'user-1', candidate())).resolves.toBeNull()
})
it('skips filing when the archived duplicate already has an inbox item', async () => {
// The provenance key catches a re-hunted message; this catches the same
// receipt arriving through ANOTHER inbox: the receipt is already in the
// Underlag flow, so no second item is filed.
mockUploadDocument.mockResolvedValue({ id: 'doc-orig', deduplicated: true })
const { client, inserted } = mockSupabase(null, {}, [{ id: 'item-existing' }])
const result = await ingestMailCandidate(client, 'co-1', 'user-1', candidate())
expect(result).toBeNull()
expect(inserted).toHaveLength(0)
// Locks the flag itself: without dedupeByContent the mock still answers
// deduplicated, but production would silently archive copies again.
expect(mockUploadDocument).toHaveBeenCalledWith(
expect.anything(),
'user-1',
'co-1',
expect.anything(),
{ upload_source: 'mail_hunt', dedupeByContent: true },
)
// The skip is behandlingshistorik, not just an app log; and the payload
// stays pseudonymous (never a mailbox address).
expect(mockAppendHistory).toHaveBeenCalledTimes(1)
const event = mockAppendHistory.mock.calls[0]![0] as {
eventType: string
aggregateId: string
payload: Record<string, unknown>
}
expect(event.eventType).toBe('DocumentDuplicateSkipped')
expect(event.aggregateId).toBe('doc-orig')
expect(event.payload).toMatchObject({
channel: 'mail_hunt',
document_id: 'doc-orig',
inbox_item_id: 'item-existing',
reason: 'duplicate_content',
})
expect(JSON.stringify(event.payload)).not.toContain('@')
})
it('still skips the duplicate when the history append fails', async () => {
// The audit write is best-effort by design: a history outage must not
// turn a correct skip into a duplicate filing.
mockUploadDocument.mockResolvedValue({ id: 'doc-orig', deduplicated: true })
mockAppendHistory.mockRejectedValueOnce(new Error('history down'))
const { client, inserted } = mockSupabase(null, {}, [{ id: 'item-existing' }])
const result = await ingestMailCandidate(client, 'co-1', 'user-1', candidate())
expect(result).toBeNull()
expect(inserted).toHaveLength(0)
})
it('files an item against the existing document when the duplicate never passed the inbox', async () => {
// A content match against a document with no inbox item (a manually
// attached copy) must not swallow the receipt: the affärshändelse still
// needs routing to matching (BFL 5 kap), just without a second archive copy.
mockUploadDocument.mockResolvedValue({ id: 'doc-orig', deduplicated: true })
const { client, inserted } = mockSupabase(null, {}, [null])
const result = await ingestMailCandidate(client, 'co-1', 'user-1', candidate())
expect(result).toMatchObject({ documentId: 'doc-orig', inboxItemId: 'item-1' })
expect(inserted).toHaveLength(1)
expect(inserted[0].document_id).toBe('doc-orig')
})
it('ignores a body-only receipt, which has nothing to download', async () => {
const { client } = mockSupabase(null)
const result = await ingestMailCandidate(
+58 -1
View File
@@ -11,6 +11,7 @@
*/
import type { SupabaseClient } from '@supabase/supabase-js'
import { uploadDocument } from '@/lib/core/documents/document-service'
import { appendProcessingHistory } from '@/lib/processing-history/append'
import { getMailSearchService, type MailCandidate } from '@/lib/mail-search/service'
import { createLogger } from '@/lib/logger'
@@ -164,9 +165,65 @@ export async function ingestMailCandidate(
) as ArrayBuffer,
type: sniffMimeType(fetched.bytes, fetched.mimeType, fileName),
},
{ upload_source: 'mail_hunt' },
{ upload_source: 'mail_hunt', dedupeByContent: true },
)
// Content already archived for this company: the provenance key above
// only catches the SAME message re-hunted, while this catches the same
// receipt arriving through another inbox or channel (forwards are
// common). Skip ONLY when an inbox item already carries the document:
// then the receipt is in the Underlag flow (or handled). When the match
// is a document that never passed the inbox (a manually attached copy,
// an archival file), fall through and file an item for the EXISTING
// document: silently dropping the receipt could leave a real
// affärshändelse without underlag routing (BFL 5 kap).
if (document.deduplicated) {
const { data: dupItem, error: dupErr } = await supabase
.from('invoice_inbox_items')
.select('id')
.eq('company_id', companyId)
.eq('document_id', document.id)
.limit(1)
.maybeSingle()
// Fail closed: a broken lookup must not file a second item.
if (dupErr) throw new Error(dupErr.message)
if (dupItem) {
// Behandlingshistorik, not just an app log: the dedupe decision is
// part of the auditable trail (BFNAR 2013:2 kap 8).
try {
await appendProcessingHistory({
companyId,
correlationId: crypto.randomUUID(),
aggregateType: 'Document',
aggregateId: document.id,
eventType: 'DocumentDuplicateSkipped',
// No mailbox address here: the processing-history payload
// contract is pseudonymous IDs only (never emails). Which
// mailbox first delivered the receipt is on the existing
// item's channel_context.
payload: {
channel: 'mail_hunt',
document_id: document.id,
inbox_item_id: (dupItem as { id: string }).id,
mail_message_id: candidate.messageId,
reason: 'duplicate_content',
},
actor: { type: 'system', id: 'receipt-hunt' },
occurredAt: new Date(),
})
} catch (histErr) {
log.warn('could not append DocumentDuplicateSkipped', {
error: histErr instanceof Error ? histErr.message : String(histErr),
})
}
log.info('skipped duplicate attachment content', {
messageId: candidate.messageId,
documentId: document.id,
})
continue
}
}
// uploadDocument emits document.uploaded and awaits its handlers, so the
// extraction extension has already read the amount, date and vendor out
// of this file by the time we get here. Copying it onto the inbox item is