feat(zettle): sync paid purchases into webshop_orders (#2445)
Community PR #2416 by @olofpinzke, adopted and finished by maintainers (rebased so every commit is signed). Why the problem occurred: no Zettle integration; POS sales only reached the books as bank descriptors while Woo/Shopify already had order underlag via webshop_orders. The contributor's version also failed at the database (platform CHECKs listed only woocommerce/shopify), which the mocked unit tests never saw. What was simplified: reused the Orders/book/invoice path instead of a new inbox; Finance API payouts/fees deferred. Sales the one-account, revenue-per-rate model cannot book (split tender, gift cards, tips) import unbookable with a "bokför manuellt" title instead of guessing accounts. Reset parity uses the rename-and-wrap pattern instead of re-issuing the reset body. Why this solution: per-purchase rows give the radunderlag BFL verifikat need and the bulk-book path exists; daily kassarapport aggregation and Finance API fees/payouts are the follow-up (DECISIONS.md). Skeptic-refuted paths fixed before merge: concurrent refresh-token rotation (sync claim), cron offset paging (candidate snapshot), platform CHECKs, writer-role gate, migration-reset parity, white-label return origin re-validated at callback, VAT net from product rows. Not live until ZETTLE_CLIENT_ID / ZETTLE_CLIENT_SECRET / ZETTLE_CREDENTIALS_ENCRYPTION_KEY are set on Vercel and a Zettle developer app is registered with the callback redirect URI. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WtYqzKPoTSRHskYYdf7MwB
This commit is contained in:
co-authored by
Claude Fable 5.1
parent
500806f001
commit
6ea92f3152
@@ -295,6 +295,14 @@ describe('POST /api/company/[id]/delete', () => {
|
||||
client_secret_encrypted: null,
|
||||
})
|
||||
)
|
||||
expect(updateSpies.zettle_connections).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
status: 'revoked',
|
||||
disconnected_at: expect.any(String),
|
||||
refresh_token_encrypted: null,
|
||||
oauth_state: null,
|
||||
})
|
||||
)
|
||||
expect(updateSpies.stripe_connections).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
status: 'revoked',
|
||||
@@ -324,6 +332,7 @@ describe('POST /api/company/[id]/delete', () => {
|
||||
for (const table of [
|
||||
'woocommerce_connections',
|
||||
'shopify_connections',
|
||||
'zettle_connections',
|
||||
'stripe_connections',
|
||||
]) {
|
||||
expect(insertSpy).toHaveBeenCalledWith(
|
||||
@@ -369,5 +378,8 @@ describe('POST /api/company/[id]/delete', () => {
|
||||
expect(insertSpy).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ table_name: 'shopify_connections' }),
|
||||
)
|
||||
expect(insertSpy).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ table_name: 'zettle_connections' }),
|
||||
)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -177,6 +177,20 @@ export async function POST(
|
||||
.neq('status', 'revoked')
|
||||
.select('id'),
|
||||
},
|
||||
{
|
||||
table: 'zettle_connections',
|
||||
result: await service
|
||||
.from('zettle_connections')
|
||||
.update({
|
||||
status: 'revoked',
|
||||
disconnected_at: archivedAt,
|
||||
refresh_token_encrypted: null,
|
||||
oauth_state: null,
|
||||
})
|
||||
.eq('company_id', companyId)
|
||||
.neq('status', 'revoked')
|
||||
.select('id'),
|
||||
},
|
||||
{
|
||||
table: 'stripe_connections',
|
||||
result: await service
|
||||
|
||||
@@ -0,0 +1,238 @@
|
||||
import { describe, it, expect, vi, beforeEach } from 'vitest'
|
||||
import { eventBus } from '@/lib/events/bus'
|
||||
|
||||
const mockExchangeCodeForTokens = vi.fn()
|
||||
const mockFetchUserSelf = vi.fn()
|
||||
vi.mock('@/extensions/general/zettle/lib/oauth', () => ({
|
||||
exchangeCodeForTokens: (...args: unknown[]) => mockExchangeCodeForTokens(...args),
|
||||
fetchUserSelf: (...args: unknown[]) => mockFetchUserSelf(...args),
|
||||
}))
|
||||
|
||||
vi.mock('@/extensions/general/zettle/lib/credentials', () => ({
|
||||
encryptCredential: (value: string) => `enc:${value}`,
|
||||
}))
|
||||
|
||||
const { mockFrom, mockGetUser } = vi.hoisted(() => ({
|
||||
mockFrom: vi.fn(),
|
||||
mockGetUser: vi.fn(),
|
||||
}))
|
||||
|
||||
vi.mock('@/lib/supabase/server', () => ({
|
||||
createServiceClient: vi.fn().mockResolvedValue({ from: mockFrom }),
|
||||
createClient: vi.fn().mockResolvedValue({ auth: { getUser: mockGetUser } }),
|
||||
}))
|
||||
|
||||
vi.mock('@/lib/init', () => ({ ensureInitialized: vi.fn() }))
|
||||
|
||||
const registryGet = vi.fn((..._args: unknown[]) => ({ id: 'zettle' }) as unknown)
|
||||
vi.mock('@/lib/extensions/loader', () => ({ loadExtensions: vi.fn() }))
|
||||
vi.mock('@/lib/extensions/registry', () => ({
|
||||
extensionRegistry: { get: (...args: unknown[]) => registryGet(...args) },
|
||||
}))
|
||||
|
||||
// Only a host that resolves in the brands table is a valid return origin.
|
||||
vi.mock('@/lib/branding/resolve', () => ({
|
||||
resolveBrandByHost: vi.fn(async (host: string) =>
|
||||
host === 'brand.testbrand.example' ? { domain: 'brand.testbrand.example' } : null,
|
||||
),
|
||||
}))
|
||||
|
||||
vi.stubEnv('NEXT_PUBLIC_APP_URL', 'http://localhost:3000')
|
||||
|
||||
import { GET } from '../route'
|
||||
|
||||
const CONNECTION_ID = 'connection-1'
|
||||
const OAUTH_STATE = 'state-token-1'
|
||||
|
||||
function makeRequest(params: Record<string, string>) {
|
||||
const url = new URL('http://localhost:3000/api/extensions/zettle/callback')
|
||||
for (const [k, v] of Object.entries(params)) {
|
||||
url.searchParams.set(k, v)
|
||||
}
|
||||
return new Request(url.toString())
|
||||
}
|
||||
|
||||
function mockChain(result: { data?: unknown; error?: unknown }) {
|
||||
const chain: Record<string, unknown> = {}
|
||||
for (const m of ['select', 'eq', 'update', 'insert']) {
|
||||
chain[m] = vi.fn().mockReturnValue(chain)
|
||||
}
|
||||
chain.single = vi
|
||||
.fn()
|
||||
.mockResolvedValue({ data: result.data ?? null, error: result.error ?? null })
|
||||
chain.maybeSingle = vi
|
||||
.fn()
|
||||
.mockResolvedValue({ data: result.data ?? null, error: result.error ?? null })
|
||||
chain.then = (resolve: (v: unknown) => void) =>
|
||||
resolve({ data: result.data ?? null, error: result.error ?? null })
|
||||
return chain
|
||||
}
|
||||
|
||||
describe('GET /api/extensions/zettle/callback', () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
eventBus.clear()
|
||||
mockGetUser.mockResolvedValue({ data: { user: { id: 'user-1' } }, error: null })
|
||||
mockExchangeCodeForTokens.mockResolvedValue({
|
||||
access_token: 'access',
|
||||
refresh_token: 'refresh',
|
||||
expires_in: 7200,
|
||||
})
|
||||
mockFetchUserSelf.mockResolvedValue({ organizationUuid: 'org-uuid-1' })
|
||||
})
|
||||
|
||||
it('activates only while the pending oauth_state still matches', async () => {
|
||||
const findChain = mockChain({
|
||||
data: { id: CONNECTION_ID, user_id: 'user-1', company_id: 'company-1' },
|
||||
})
|
||||
const replayChain = mockChain({ error: null })
|
||||
const activateChain = mockChain({
|
||||
data: {
|
||||
id: CONNECTION_ID,
|
||||
company_id: 'company-1',
|
||||
user_id: 'user-1',
|
||||
organization_uuid: 'org-uuid-1',
|
||||
},
|
||||
})
|
||||
mockFrom
|
||||
.mockReturnValueOnce(findChain)
|
||||
.mockReturnValueOnce(replayChain)
|
||||
.mockReturnValueOnce(activateChain)
|
||||
|
||||
const response = await GET(makeRequest({ code: 'ac_123', state: OAUTH_STATE }))
|
||||
|
||||
expect(response.status).toBe(307)
|
||||
expect(response.headers.get('location')).toBe(
|
||||
'http://localhost:3000/import?mode=zettle&zettle_connected=true',
|
||||
)
|
||||
|
||||
const eqCalls = (activateChain.eq as ReturnType<typeof vi.fn>).mock.calls.map(
|
||||
(c) => c as [string, string],
|
||||
)
|
||||
expect(eqCalls).toEqual(
|
||||
expect.arrayContaining([
|
||||
['id', CONNECTION_ID],
|
||||
['status', 'pending'],
|
||||
['oauth_state', OAUTH_STATE],
|
||||
]),
|
||||
)
|
||||
expect(activateChain.maybeSingle).toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('refuses activation when /connect invalidated the pending row mid-callback', async () => {
|
||||
// Lookup still sees the original pending row (TOCTOU), then a concurrent
|
||||
// POST /connect flips it to error and clears oauth_state before activate.
|
||||
const findChain = mockChain({
|
||||
data: { id: CONNECTION_ID, user_id: 'user-1', company_id: 'company-1' },
|
||||
})
|
||||
const replayChain = mockChain({ error: null })
|
||||
const activateChain = mockChain({ data: null, error: null })
|
||||
mockFrom
|
||||
.mockReturnValueOnce(findChain)
|
||||
.mockReturnValueOnce(replayChain)
|
||||
.mockReturnValueOnce(activateChain)
|
||||
|
||||
const response = await GET(makeRequest({ code: 'ac_123', state: OAUTH_STATE }))
|
||||
|
||||
expect(response.headers.get('location')).toBe(
|
||||
'http://localhost:3000/import?mode=zettle&zettle_error=invalid_state',
|
||||
)
|
||||
const eqCalls = (activateChain.eq as ReturnType<typeof vi.fn>).mock.calls.map(
|
||||
(c) => c as [string, string],
|
||||
)
|
||||
expect(eqCalls).toEqual(
|
||||
expect.arrayContaining([
|
||||
['status', 'pending'],
|
||||
['oauth_state', OAUTH_STATE],
|
||||
]),
|
||||
)
|
||||
// Must not fall through into the conflict/error cleanup update.
|
||||
expect(mockFrom).toHaveBeenCalledTimes(3)
|
||||
})
|
||||
|
||||
it('returns the browser to the brand origin the connect flow started on', async () => {
|
||||
const findChain = mockChain({
|
||||
data: {
|
||||
id: CONNECTION_ID,
|
||||
user_id: 'user-1',
|
||||
company_id: 'company-1',
|
||||
return_origin: 'https://brand.testbrand.example',
|
||||
},
|
||||
})
|
||||
const replayChain = mockChain({ error: null })
|
||||
const activateChain = mockChain({
|
||||
data: {
|
||||
id: CONNECTION_ID,
|
||||
company_id: 'company-1',
|
||||
user_id: 'user-1',
|
||||
organization_uuid: 'org-uuid-1',
|
||||
},
|
||||
})
|
||||
mockFrom
|
||||
.mockReturnValueOnce(findChain)
|
||||
.mockReturnValueOnce(replayChain)
|
||||
.mockReturnValueOnce(activateChain)
|
||||
|
||||
const response = await GET(makeRequest({ code: 'ac_123', state: OAUTH_STATE }))
|
||||
|
||||
expect(response.headers.get('location')).toBe(
|
||||
'https://brand.testbrand.example/import?mode=zettle&zettle_connected=true',
|
||||
)
|
||||
})
|
||||
|
||||
it('never redirects to a tampered return_origin that is not a brand domain', async () => {
|
||||
// Members can UPDATE the row through RLS: the column is not trusted.
|
||||
const findChain = mockChain({
|
||||
data: {
|
||||
id: CONNECTION_ID,
|
||||
user_id: 'user-1',
|
||||
company_id: 'company-1',
|
||||
return_origin: 'https://evil.example',
|
||||
},
|
||||
})
|
||||
const replayChain = mockChain({ error: null })
|
||||
const activateChain = mockChain({
|
||||
data: { id: CONNECTION_ID, company_id: 'company-1', user_id: 'user-1', organization_uuid: 'org-uuid-1' },
|
||||
})
|
||||
mockFrom.mockReturnValueOnce(findChain).mockReturnValueOnce(replayChain).mockReturnValueOnce(activateChain)
|
||||
|
||||
const response = await GET(makeRequest({ code: 'ac_123', state: OAUTH_STATE }))
|
||||
|
||||
expect(response.headers.get('location')).toBe(
|
||||
'http://localhost:3000/import?mode=zettle&zettle_connected=true',
|
||||
)
|
||||
})
|
||||
|
||||
it('returns a denied authorization to the stored brand origin', async () => {
|
||||
mockFrom.mockReturnValueOnce(
|
||||
mockChain({ data: { return_origin: 'https://brand.testbrand.example' } }),
|
||||
)
|
||||
|
||||
const response = await GET(makeRequest({ error: 'access_denied', state: OAUTH_STATE }))
|
||||
|
||||
expect(response.headers.get('location')).toBe(
|
||||
'https://brand.testbrand.example/import?mode=zettle&zettle_error=access_denied',
|
||||
)
|
||||
expect(mockExchangeCodeForTokens).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('refuses with 503 when the zettle extension is not enabled', async () => {
|
||||
registryGet.mockReturnValueOnce(undefined)
|
||||
|
||||
const response = await GET(makeRequest({ code: 'ac_123', state: OAUTH_STATE }))
|
||||
|
||||
expect(response.status).toBe(503)
|
||||
expect(mockFrom).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('redirects with invalid_state when the oauth state is unknown', async () => {
|
||||
mockFrom.mockReturnValueOnce(mockChain({ data: null, error: { code: 'PGRST116' } }))
|
||||
|
||||
const response = await GET(makeRequest({ code: 'ac_123', state: 'unknown-state' }))
|
||||
|
||||
expect(response.headers.get('location')).toBe(
|
||||
'http://localhost:3000/import?mode=zettle&zettle_error=invalid_state',
|
||||
)
|
||||
expect(mockExchangeCodeForTokens).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,245 @@
|
||||
import { createServiceClient } from '@/lib/supabase/server'
|
||||
import { NextResponse } from 'next/server'
|
||||
import { ensureInitialized } from '@/lib/init'
|
||||
import { loadExtensions } from '@/lib/extensions/loader'
|
||||
import { extensionRegistry } from '@/lib/extensions/registry'
|
||||
import { eventBus } from '@/lib/events/bus'
|
||||
import { hashAuthCode } from '@/lib/auth/oauth-codes'
|
||||
import {
|
||||
requireFlowInitiator,
|
||||
FLOW_INITIATOR_MISMATCH_MESSAGE,
|
||||
} from '@/lib/auth/oauth-flow-binding'
|
||||
import { encryptCredential } from '@/extensions/general/zettle/lib/credentials'
|
||||
import { validateReturnOrigin } from '@/extensions/general/zettle/lib/return-origin'
|
||||
import {
|
||||
exchangeCodeForTokens,
|
||||
fetchUserSelf,
|
||||
} from '@/extensions/general/zettle/lib/oauth'
|
||||
|
||||
// This route emits zettle.connected (audit trail). ensureInitialized() must
|
||||
// run at module load so the event_log handler has subscribed before the first
|
||||
// emit on a cold instance.
|
||||
ensureInitialized()
|
||||
|
||||
/**
|
||||
* GET /api/extensions/zettle/callback
|
||||
*
|
||||
* OAuth callback for Zettle partner authorization. Must be a real Next.js
|
||||
* route (not an extension dispatcher handler) because Zettle redirects the
|
||||
* user's browser to this URL directly.
|
||||
*/
|
||||
export async function GET(request: Request) {
|
||||
// Physical route: refuse (503) when the extension is not enabled instead
|
||||
// of quietly activating connections for a feature the deployment turned off.
|
||||
loadExtensions()
|
||||
if (!extensionRegistry.get('zettle')) {
|
||||
return NextResponse.json(
|
||||
{ error: 'Zettle extension is not enabled', code: 'EXTENSION_DISABLED' },
|
||||
{ status: 503 },
|
||||
)
|
||||
}
|
||||
|
||||
const { searchParams } = new URL(request.url)
|
||||
|
||||
const code = searchParams.get('code')
|
||||
const state = searchParams.get('state')
|
||||
const error = searchParams.get('error')
|
||||
const errorDescription = searchParams.get('error_description')
|
||||
|
||||
const appBase = (process.env.NEXT_PUBLIC_APP_URL || 'http://localhost:3000').replace(/\/$/, '')
|
||||
// Return the browser to the origin the connect flow started on. Zettle
|
||||
// redirects to the one registered callback URL, so a white-label user
|
||||
// would otherwise land on the canonical app domain. The stored value is
|
||||
// re-validated here (members can update the row through RLS): the app
|
||||
// origin or a brand domain, never an arbitrary URL.
|
||||
const returnUrlFor = async (origin: string | null | undefined) =>
|
||||
`${await validateReturnOrigin(origin, appBase)}/import?mode=zettle`
|
||||
let returnUrl = `${appBase}/import?mode=zettle`
|
||||
|
||||
if (error) {
|
||||
const errorMessage = errorDescription || error
|
||||
const logDenied = error === 'access_denied' ? console.warn : console.error
|
||||
logDenied('[zettle] OAuth authorization denied', {
|
||||
error,
|
||||
error_description: errorDescription,
|
||||
has_state: !!state,
|
||||
})
|
||||
|
||||
if (state) {
|
||||
try {
|
||||
const supabase = await createServiceClient()
|
||||
const { data: denied } = await supabase
|
||||
.from('zettle_connections')
|
||||
.update({ status: 'error', error_message: errorMessage, oauth_state: null })
|
||||
.eq('oauth_state', state)
|
||||
.eq('status', 'pending')
|
||||
.select('return_origin')
|
||||
.maybeSingle()
|
||||
returnUrl = await returnUrlFor(denied?.return_origin)
|
||||
} catch (cleanupError) {
|
||||
console.error('[zettle] Failed to clean up pending connection:', cleanupError)
|
||||
}
|
||||
}
|
||||
|
||||
return NextResponse.redirect(
|
||||
`${returnUrl}&zettle_error=${encodeURIComponent(errorMessage)}`,
|
||||
)
|
||||
}
|
||||
|
||||
if (!code || !state) {
|
||||
return NextResponse.redirect(`${returnUrl}&zettle_error=missing_parameters`)
|
||||
}
|
||||
|
||||
const supabase = await createServiceClient()
|
||||
|
||||
try {
|
||||
const { data: pendingConnection, error: findError } = await supabase
|
||||
.from('zettle_connections')
|
||||
.select('id, user_id, company_id, return_origin')
|
||||
.eq('oauth_state', state)
|
||||
.eq('status', 'pending')
|
||||
.single()
|
||||
|
||||
if (pendingConnection) {
|
||||
returnUrl = await returnUrlFor(pendingConnection.return_origin)
|
||||
}
|
||||
|
||||
if (findError || !pendingConnection) {
|
||||
console.error('[zettle] No pending connection for oauth_state', {
|
||||
findError: findError
|
||||
? { message: findError.message, code: findError.code }
|
||||
: null,
|
||||
hasCode: !!code,
|
||||
})
|
||||
return NextResponse.redirect(
|
||||
`${returnUrl}&zettle_error=${encodeURIComponent('invalid_state')}`,
|
||||
)
|
||||
}
|
||||
|
||||
const initiator = await requireFlowInitiator(request, pendingConnection.user_id, {
|
||||
flow: 'zettle.callback',
|
||||
})
|
||||
if (!initiator.ok) {
|
||||
if (initiator.reason === 'no_session') {
|
||||
return initiator.response
|
||||
}
|
||||
return NextResponse.redirect(
|
||||
`${returnUrl}&zettle_error=${encodeURIComponent(FLOW_INITIATOR_MISMATCH_MESSAGE)}`,
|
||||
)
|
||||
}
|
||||
|
||||
const { error: replayError } = await supabase
|
||||
.from('oauth_used_codes')
|
||||
.insert({ code_hash: hashAuthCode(code) })
|
||||
if (replayError) {
|
||||
console.error('[zettle] Authorization code already used', {
|
||||
connectionId: pendingConnection.id,
|
||||
code: replayError.code,
|
||||
})
|
||||
return NextResponse.redirect(
|
||||
`${returnUrl}&zettle_error=${encodeURIComponent('invalid_state')}`,
|
||||
)
|
||||
}
|
||||
|
||||
const tokens = await exchangeCodeForTokens(code)
|
||||
const userSelf = await fetchUserSelf(tokens.access_token)
|
||||
|
||||
// Require the row still be the original pending state. POST /connect can
|
||||
// invalidate this row between lookup and activate; filtering only by id
|
||||
// would revive the abandoned flow and attach the wrong Zettle org.
|
||||
const { data: updatedConnection, error: updateError } = await supabase
|
||||
.from('zettle_connections')
|
||||
.update({
|
||||
organization_uuid: userSelf.organizationUuid,
|
||||
organization_name: null,
|
||||
refresh_token_encrypted: encryptCredential(tokens.refresh_token),
|
||||
status: 'active',
|
||||
connected_at: new Date().toISOString(),
|
||||
error_message: null,
|
||||
oauth_state: null,
|
||||
transaction_sync_enabled: true,
|
||||
})
|
||||
.eq('id', pendingConnection.id)
|
||||
.eq('status', 'pending')
|
||||
.eq('oauth_state', state)
|
||||
.select('id, company_id, user_id, organization_uuid')
|
||||
.maybeSingle()
|
||||
|
||||
if (updateError) {
|
||||
const isConflict = updateError.code === '23505'
|
||||
console.error('[zettle] Failed to activate connection', {
|
||||
connectionId: pendingConnection.id,
|
||||
error: { message: updateError.message, code: updateError.code },
|
||||
})
|
||||
await supabase
|
||||
.from('zettle_connections')
|
||||
.update({
|
||||
status: 'error',
|
||||
error_message: isConflict
|
||||
? 'Zettle-organisationen är redan ansluten till ett företag.'
|
||||
: 'Anslutningen kunde inte slutföras.',
|
||||
oauth_state: null,
|
||||
refresh_token_encrypted: null,
|
||||
})
|
||||
.eq('id', pendingConnection.id)
|
||||
.eq('status', 'pending')
|
||||
return NextResponse.redirect(
|
||||
`${returnUrl}&zettle_error=${encodeURIComponent(
|
||||
isConflict ? 'account_already_connected' : 'activation_failed',
|
||||
)}`,
|
||||
)
|
||||
}
|
||||
|
||||
if (!updatedConnection) {
|
||||
console.error('[zettle] Pending connection invalidated before activation', {
|
||||
connectionId: pendingConnection.id,
|
||||
})
|
||||
return NextResponse.redirect(
|
||||
`${returnUrl}&zettle_error=${encodeURIComponent('invalid_state')}`,
|
||||
)
|
||||
}
|
||||
|
||||
try {
|
||||
await eventBus.emit({
|
||||
type: 'zettle.connected',
|
||||
payload: {
|
||||
connectionId: updatedConnection.id,
|
||||
organizationUuid: updatedConnection.organization_uuid!,
|
||||
userId: updatedConnection.user_id,
|
||||
companyId: updatedConnection.company_id,
|
||||
},
|
||||
})
|
||||
} catch (emitError) {
|
||||
console.error('[zettle] Failed to emit zettle.connected event', {
|
||||
connectionId: updatedConnection.id,
|
||||
error: emitError instanceof Error ? emitError.message : String(emitError),
|
||||
})
|
||||
}
|
||||
|
||||
return NextResponse.redirect(`${returnUrl}&zettle_connected=true`)
|
||||
} catch (error) {
|
||||
console.error('[zettle] Callback error', {
|
||||
message: error instanceof Error ? error.message : String(error),
|
||||
name: error instanceof Error ? error.name : undefined,
|
||||
hasCode: !!code,
|
||||
})
|
||||
|
||||
try {
|
||||
await supabase
|
||||
.from('zettle_connections')
|
||||
.update({
|
||||
status: 'error',
|
||||
error_message: 'Anslutningen kunde inte slutföras.',
|
||||
oauth_state: null,
|
||||
})
|
||||
.eq('oauth_state', state)
|
||||
.eq('status', 'pending')
|
||||
} catch (cleanupError) {
|
||||
console.error('[zettle] Callback cleanup failed:', cleanupError)
|
||||
}
|
||||
|
||||
return NextResponse.redirect(
|
||||
`${returnUrl}&zettle_error=${encodeURIComponent('connection_failed')}`,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,117 @@
|
||||
import { describe, it, expect, vi, beforeEach } from 'vitest'
|
||||
|
||||
// Each test re-imports the route after vi.resetModules(); the cold import
|
||||
// exceeds the 5 s default under a loaded CI shard.
|
||||
vi.setConfig({ testTimeout: 30_000 })
|
||||
|
||||
const verifyCronSecret = vi.fn((..._args: unknown[]) => null as unknown)
|
||||
vi.mock('@/lib/auth/cron', () => ({
|
||||
verifyCronSecret: (...args: unknown[]) => verifyCronSecret(...args),
|
||||
}))
|
||||
|
||||
const registryGet = vi.fn()
|
||||
vi.mock('@/lib/extensions/loader', () => ({ loadExtensions: vi.fn() }))
|
||||
vi.mock('@/lib/extensions/registry', () => ({
|
||||
extensionRegistry: { get: (...args: unknown[]) => registryGet(...args) },
|
||||
}))
|
||||
|
||||
const rangeResult = vi.fn()
|
||||
vi.mock('@/lib/supabase/service-client', () => ({
|
||||
createServiceRoleClient: vi.fn(() => ({
|
||||
from: () => ({
|
||||
select: () => ({
|
||||
eq: () => ({
|
||||
eq: () => ({
|
||||
order: () => ({
|
||||
range: (...args: unknown[]) => rangeResult(...args),
|
||||
}),
|
||||
}),
|
||||
}),
|
||||
}),
|
||||
}),
|
||||
})),
|
||||
}))
|
||||
|
||||
const isZettleConfigured = vi.fn((..._args: unknown[]) => true)
|
||||
vi.mock('@/extensions/general/zettle/lib/credentials', () => ({
|
||||
isZettleConfigured: (...args: unknown[]) => isZettleConfigured(...args),
|
||||
}))
|
||||
|
||||
const syncZettlePurchases = vi.fn()
|
||||
vi.mock('@/extensions/general/zettle/lib/order-sync', () => ({
|
||||
syncZettlePurchases: (...args: unknown[]) => syncZettlePurchases(...args),
|
||||
}))
|
||||
|
||||
const hasCapability = vi.fn()
|
||||
vi.mock('@/lib/entitlements/has-capability', () => ({
|
||||
hasCapability: (...args: unknown[]) => hasCapability(...args),
|
||||
}))
|
||||
|
||||
const CONNECTION = { id: 'conn-1', company_id: 'company-1' }
|
||||
|
||||
async function callRoute() {
|
||||
const { GET } = await import('../route')
|
||||
return GET(new Request('https://example.test/api/extensions/zettle/orders/cron'))
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks()
|
||||
vi.resetModules()
|
||||
verifyCronSecret.mockReturnValue(null)
|
||||
registryGet.mockReturnValue({ id: 'zettle' })
|
||||
isZettleConfigured.mockReturnValue(true)
|
||||
hasCapability.mockResolvedValue(true)
|
||||
rangeResult.mockResolvedValue({ data: [CONNECTION], error: null })
|
||||
syncZettlePurchases.mockResolvedValue({
|
||||
fetched: 2,
|
||||
refundsFetched: 0,
|
||||
inserted: 2,
|
||||
updated: 0,
|
||||
unchanged: 0,
|
||||
frozenFlagged: 0,
|
||||
crossMarked: 0,
|
||||
errors: 0,
|
||||
})
|
||||
process.env.NEXT_PUBLIC_SUPABASE_URL = 'https://project.supabase.co'
|
||||
process.env.SUPABASE_SERVICE_ROLE_KEY = 'service-key'
|
||||
})
|
||||
|
||||
describe('GET /api/extensions/zettle/orders/cron', () => {
|
||||
it('returns 503 when the extension is disabled', async () => {
|
||||
registryGet.mockReturnValue(null)
|
||||
const res = await callRoute()
|
||||
expect(res.status).toBe(503)
|
||||
expect(syncZettlePurchases).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('syncs entitled active connections', async () => {
|
||||
const res = await callRoute()
|
||||
expect(res.status).toBe(200)
|
||||
expect(syncZettlePurchases).toHaveBeenCalled()
|
||||
const body = await res.json()
|
||||
expect(body.processed).toBe(1)
|
||||
expect(body.inserted).toBe(2)
|
||||
})
|
||||
|
||||
it('pages past a front of non-entitled connections so entitled ones still sync', async () => {
|
||||
// Old limit(50) before hasCapability starved eligible rows behind 50 skips.
|
||||
const notEntitled = Array.from({ length: 50 }, (_, i) => ({
|
||||
id: `skip-${i}`,
|
||||
company_id: `co-skip-${i}`,
|
||||
}))
|
||||
const entitled = { id: 'conn-entitled', company_id: 'company-entitled' }
|
||||
rangeResult.mockResolvedValue({ data: [...notEntitled, entitled], error: null })
|
||||
hasCapability.mockImplementation(async (_sb: unknown, companyId: string) => {
|
||||
return companyId === 'company-entitled'
|
||||
})
|
||||
|
||||
const res = await callRoute()
|
||||
expect(res.status).toBe(200)
|
||||
const body = await res.json()
|
||||
expect(body.processed).toBe(1)
|
||||
expect(syncZettlePurchases).toHaveBeenCalledTimes(1)
|
||||
expect(syncZettlePurchases.mock.calls[0][1]).toMatchObject(entitled)
|
||||
// Skips must not touch last_order_synced_at (purchase recovery cursor).
|
||||
expect(rangeResult).toHaveBeenCalledWith(0, 99)
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,156 @@
|
||||
import { createServiceRoleClient } from '@/lib/supabase/service-client'
|
||||
import { NextResponse } from 'next/server'
|
||||
import { withCronContext } from '@/lib/api/with-cron-context'
|
||||
import { errorResponse, errorResponseFromCode } from '@/lib/errors/get-structured-error'
|
||||
import { hasCapability } from '@/lib/entitlements/has-capability'
|
||||
import { CAPABILITY } from '@/lib/entitlements/keys'
|
||||
import { loadExtensions } from '@/lib/extensions/loader'
|
||||
import { extensionRegistry } from '@/lib/extensions/registry'
|
||||
import { isZettleConfigured } from '@/extensions/general/zettle/lib/credentials'
|
||||
import { syncZettlePurchases } from '@/extensions/general/zettle/lib/order-sync'
|
||||
import type { ZettleConnection } from '@/extensions/general/zettle/types'
|
||||
|
||||
export const maxDuration = 300
|
||||
|
||||
/** Cap of entitled connections synced per cron invocation. */
|
||||
const MAX_SYNCED = 50
|
||||
/** Page size when scanning candidates ordered by purchase cursor. */
|
||||
const CANDIDATE_PAGE_SIZE = 100
|
||||
/** Hard stop so a flood of non-entitled rows cannot burn the whole budget scanning. */
|
||||
const MAX_CANDIDATES_SCANNED = 2000
|
||||
|
||||
/**
|
||||
* GET /api/extensions/zettle/orders/cron
|
||||
* Nightly purchase sync for connections that opted in (transaction_sync_enabled):
|
||||
* upserts each connected org's paid purchases and refunds into webshop_orders.
|
||||
*
|
||||
* Candidates are paged by last_order_synced_at. Entitlement skips do not advance
|
||||
* that cursor (it controls purchase recovery) and do not consume the sync cap,
|
||||
* so a front of non-entitled rows cannot starve eligible connections behind them.
|
||||
*/
|
||||
export const GET = withCronContext('cron.zettle_order_sync', async (_request, ctx) => {
|
||||
loadExtensions()
|
||||
if (!extensionRegistry.get('zettle')) {
|
||||
ctx.log.warn('zettle extension is not enabled; cron refused')
|
||||
return NextResponse.json(
|
||||
{ error: 'Zettle extension is not enabled', code: 'EXTENSION_DISABLED' },
|
||||
{ status: 503 },
|
||||
)
|
||||
}
|
||||
|
||||
const supabaseUrl = process.env.NEXT_PUBLIC_SUPABASE_URL
|
||||
const supabaseServiceKey = process.env.SUPABASE_SERVICE_ROLE_KEY
|
||||
|
||||
if (!supabaseUrl || !supabaseServiceKey) {
|
||||
return errorResponseFromCode('INTERNAL_ERROR', ctx.log, {
|
||||
requestId: ctx.requestId,
|
||||
details: { reason: 'Missing Supabase configuration' },
|
||||
})
|
||||
}
|
||||
if (!isZettleConfigured()) {
|
||||
return NextResponse.json({ message: 'Zettle not configured', processed: 0 })
|
||||
}
|
||||
|
||||
const supabase = createServiceRoleClient(supabaseUrl, supabaseServiceKey)
|
||||
|
||||
const startTime = Date.now()
|
||||
const TIME_BUDGET_MS = 240_000
|
||||
const deadlineMs = startTime + TIME_BUDGET_MS
|
||||
|
||||
const results: Array<{
|
||||
connectionId: string
|
||||
inserted: number
|
||||
updated: number
|
||||
status: 'synced' | 'revoked' | 'locked' | 'error'
|
||||
}> = []
|
||||
|
||||
// Snapshot the candidate list BEFORE syncing any of it. Each sync moves its
|
||||
// row's last_order_synced_at to "now" (to the tail of this ordering), so
|
||||
// paging with a live offset would re-fetch already-synced rows on page 2
|
||||
// and never reach the eligible rows that slid into the gap.
|
||||
const candidates: ZettleConnection[] = []
|
||||
for (let offset = 0; offset < MAX_CANDIDATES_SCANNED; offset += CANDIDATE_PAGE_SIZE) {
|
||||
const { data: page, error: connError } = await supabase
|
||||
.from('zettle_connections')
|
||||
.select('*')
|
||||
.eq('status', 'active')
|
||||
.eq('transaction_sync_enabled', true)
|
||||
.order('last_order_synced_at', { ascending: true, nullsFirst: true })
|
||||
.range(offset, offset + CANDIDATE_PAGE_SIZE - 1)
|
||||
|
||||
if (connError) {
|
||||
ctx.log.error('failed to fetch zettle connections', connError, {
|
||||
message: connError.message,
|
||||
code: connError.code,
|
||||
})
|
||||
return errorResponse(connError, ctx.log, { requestId: ctx.requestId })
|
||||
}
|
||||
|
||||
if (!page || page.length === 0) break
|
||||
candidates.push(...(page as ZettleConnection[]))
|
||||
if (page.length < CANDIDATE_PAGE_SIZE) break
|
||||
}
|
||||
|
||||
let scanned = 0
|
||||
for (const connection of candidates) {
|
||||
if (results.length >= MAX_SYNCED) break
|
||||
if (Date.now() >= deadlineMs) {
|
||||
ctx.log.info('time budget reached', { processedSoFar: results.length, scanned })
|
||||
break
|
||||
}
|
||||
scanned += 1
|
||||
|
||||
if (!(await hasCapability(supabase, connection.company_id, CAPABILITY.zettle_sync))) {
|
||||
ctx.log.info('skip: capability not entitled', { companyId: connection.company_id })
|
||||
continue
|
||||
}
|
||||
|
||||
try {
|
||||
const summary = await syncZettlePurchases(supabase, connection, ctx.log, deadlineMs)
|
||||
if (summary.locked) {
|
||||
// A manual sync holds the claim; it will advance the cursor itself.
|
||||
results.push({ connectionId: connection.id, inserted: 0, updated: 0, status: 'locked' })
|
||||
continue
|
||||
}
|
||||
if (summary.deadlineReached) {
|
||||
ctx.log.info('connection stopped early on time budget; remaining rows resume next run', {
|
||||
connectionId: connection.id,
|
||||
})
|
||||
}
|
||||
results.push({
|
||||
connectionId: connection.id,
|
||||
inserted: summary.inserted,
|
||||
updated: summary.updated,
|
||||
status: summary.revoked ? 'revoked' : 'synced',
|
||||
})
|
||||
} catch (error) {
|
||||
ctx.log.error('zettle purchase sync failed for connection', error as Error, {
|
||||
connectionId: connection.id,
|
||||
companyId: connection.company_id,
|
||||
})
|
||||
results.push({
|
||||
connectionId: connection.id,
|
||||
inserted: 0,
|
||||
updated: 0,
|
||||
status: 'error',
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
if (candidates.length === 0) {
|
||||
return NextResponse.json({
|
||||
message: 'No connections with transaction sync enabled',
|
||||
processed: 0,
|
||||
})
|
||||
}
|
||||
|
||||
const totalInserted = results.reduce((acc, r) => acc + r.inserted, 0)
|
||||
ctx.log.info('zettle purchase sync summary', {
|
||||
processed: results.length,
|
||||
scanned,
|
||||
totalInserted,
|
||||
failed: results.filter((r) => r.status === 'error').length,
|
||||
})
|
||||
|
||||
return NextResponse.json({ processed: results.length, inserted: totalInserted, results })
|
||||
})
|
||||
Reference in New Issue
Block a user