fix(ordrestyring): finalize hour provenance and cadence (#38)
Co-authored-by: alexpolo1 <[email protected]>
This commit is contained in:
@@ -463,6 +463,96 @@ describe('OrdrestyringSyncService.getSyncStatus', () => {
|
||||
expect(mysql.createConnection).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
test('bounds hourly offer snapshot traffic below the shared GraphQL quota budget', () => {
|
||||
const previousPageSize = process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE;
|
||||
const previousMax = process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS;
|
||||
delete process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE;
|
||||
delete process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS;
|
||||
try {
|
||||
const service = new OrdrestyringSyncService();
|
||||
expect(service.getOfferSyncLimits(100)).toEqual({ pageSize: 100, maxOffers: 100 });
|
||||
process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE = '999';
|
||||
process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS = '9999';
|
||||
expect(service.getOfferSyncLimits(100)).toEqual({ pageSize: 200, maxOffers: 1000 });
|
||||
process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE = '1';
|
||||
process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS = '1';
|
||||
expect(service.getOfferSyncLimits(100)).toEqual({ pageSize: 20, maxOffers: 21 });
|
||||
} finally {
|
||||
if (previousPageSize === undefined) delete process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE;
|
||||
else process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE = previousPageSize;
|
||||
if (previousMax === undefined) delete process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS;
|
||||
else process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS = previousMax;
|
||||
}
|
||||
});
|
||||
|
||||
test('rotates a persisted offer cursor while always refreshing the newest offers', async () => {
|
||||
const previousMax = process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS;
|
||||
process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS = '21';
|
||||
try {
|
||||
const service = new OrdrestyringSyncService();
|
||||
const fresh = Array.from({ length: 20 }, (_, i) => ({ id: i + 1, number: `N${i + 1}` }));
|
||||
service.requestOfferSnapshotPage = jest.fn()
|
||||
.mockResolvedValueOnce({ items: fresh, hasMorePages: true, nextCursor: 'fresh-next' })
|
||||
.mockResolvedValueOnce({ items: [{ id: 21, number: 'OLD' }], hasMorePages: true, nextCursor: 'saved-next' });
|
||||
service.loadOfferSyncCursor = jest.fn().mockResolvedValue('saved-cursor');
|
||||
|
||||
await expect(service.fetchOfferSnapshots(100)).resolves.toHaveLength(21);
|
||||
expect(service.requestOfferSnapshotPage).toHaveBeenNthCalledWith(1, null, 20);
|
||||
expect(service.requestOfferSnapshotPage).toHaveBeenNthCalledWith(2, 'saved-cursor', 1);
|
||||
expect(service.pendingOfferSyncCursor).toBe('saved-next');
|
||||
} finally {
|
||||
if (previousMax === undefined) delete process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS;
|
||||
else process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS = previousMax;
|
||||
}
|
||||
});
|
||||
|
||||
test('resets offer backfill cursor after reaching the oldest page', async () => {
|
||||
const previousMax = process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS;
|
||||
process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS = '21';
|
||||
try {
|
||||
const service = new OrdrestyringSyncService();
|
||||
const fresh = Array.from({ length: 20 }, (_, i) => ({ id: i + 1 }));
|
||||
service.requestOfferSnapshotPage = jest.fn()
|
||||
.mockResolvedValueOnce({ items: fresh, hasMorePages: true, nextCursor: 'fresh-next' })
|
||||
.mockResolvedValueOnce({ items: [{ id: 99 }], hasMorePages: false, nextCursor: null });
|
||||
service.loadOfferSyncCursor = jest.fn().mockResolvedValue('last-cursor');
|
||||
|
||||
await service.fetchOfferSnapshots(100);
|
||||
|
||||
expect(service.requestOfferSnapshotPage).toHaveBeenNthCalledWith(2, 'last-cursor', 1);
|
||||
expect(service.pendingOfferSyncCursor).toBe(null);
|
||||
} finally {
|
||||
if (previousMax === undefined) delete process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS;
|
||||
else process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS = previousMax;
|
||||
}
|
||||
});
|
||||
|
||||
test('persists offer backfill cursor state in the local mirror database', async () => {
|
||||
const readConnection = {
|
||||
execute: jest.fn()
|
||||
.mockResolvedValueOnce([{ affectedRows: 0 }])
|
||||
.mockResolvedValueOnce([[{ next_cursor: 'cursor-2' }]]),
|
||||
end: jest.fn().mockResolvedValue()
|
||||
};
|
||||
const writeConnection = {
|
||||
execute: jest.fn().mockResolvedValue([{ affectedRows: 1 }]),
|
||||
end: jest.fn().mockResolvedValue()
|
||||
};
|
||||
mysql.createConnection
|
||||
.mockResolvedValueOnce(readConnection)
|
||||
.mockResolvedValueOnce(writeConnection);
|
||||
const service = new OrdrestyringSyncService();
|
||||
|
||||
await expect(service.loadOfferSyncCursor()).resolves.toBe('cursor-2');
|
||||
await service.saveOfferSyncCursor('cursor-3');
|
||||
|
||||
expect(readConnection.execute.mock.calls[0][0]).toMatch(/CREATE TABLE IF NOT EXISTS ordrestyring_offer_sync_state/);
|
||||
expect(readConnection.execute.mock.calls[1][0]).toMatch(/SELECT next_cursor/);
|
||||
expect(writeConnection.execute.mock.calls[1][1]).toEqual([
|
||||
'offer_backfill', 'cursor-3', expect.any(Number)
|
||||
]);
|
||||
});
|
||||
|
||||
test('offer snapshot sync contains connection-close failures after successful core work', async () => {
|
||||
const connection = {
|
||||
end: jest.fn().mockRejectedValue(new Error('private close failure'))
|
||||
@@ -471,6 +561,8 @@ describe('OrdrestyringSyncService.getSyncStatus', () => {
|
||||
const service = new OrdrestyringSyncService();
|
||||
service.fetchOfferSnapshots = jest.fn().mockResolvedValue([{ id: 1 }]);
|
||||
service.fetchOfferDetailsWithRetry = jest.fn().mockResolvedValue(null);
|
||||
service.pendingOfferSyncCursor = 'cursor-after-failed-offer';
|
||||
service.saveOfferSyncCursor = jest.fn().mockResolvedValue();
|
||||
|
||||
await expect(service.syncOfferSnapshots()).resolves.toEqual({
|
||||
hasChanges: false,
|
||||
@@ -478,6 +570,38 @@ describe('OrdrestyringSyncService.getSyncStatus', () => {
|
||||
degraded: true,
|
||||
reason: 'partial_offer_detail_failures'
|
||||
});
|
||||
expect(service.saveOfferSyncCursor).not.toHaveBeenCalled();
|
||||
expect(connection.end).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
test('checkpoints offer cursor only after all selected snapshots are stored', async () => {
|
||||
const connection = {
|
||||
execute: jest.fn().mockResolvedValue([{ affectedRows: 1 }]),
|
||||
end: jest.fn().mockResolvedValue()
|
||||
};
|
||||
mysql.createConnection.mockResolvedValueOnce(connection);
|
||||
const service = new OrdrestyringSyncService();
|
||||
service.fetchOfferSnapshots = jest.fn().mockResolvedValue([{ id: 1 }]);
|
||||
service.fetchOfferDetailsWithRetry = jest.fn().mockResolvedValue({
|
||||
id: 1,
|
||||
number: 'T-1',
|
||||
description: 'Tilbud',
|
||||
createdAt: 1758272500,
|
||||
customer: { name: 'Kunde', email: null },
|
||||
totals: { salesPrice: 10000, salesPriceWithVat: 12500 },
|
||||
status: { text: 'Åben' },
|
||||
tasks: []
|
||||
});
|
||||
service.pendingOfferSyncCursor = 'cursor-after-success';
|
||||
service.saveOfferSyncCursor = jest.fn().mockResolvedValue();
|
||||
|
||||
await expect(service.syncOfferSnapshots()).resolves.toEqual({
|
||||
hasChanges: true,
|
||||
changes: 1,
|
||||
degraded: false,
|
||||
reason: null
|
||||
});
|
||||
expect(service.saveOfferSyncCursor).toHaveBeenCalledWith('cursor-after-success');
|
||||
expect(connection.end).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
|
||||
@@ -29,6 +29,7 @@ class OrdrestyringSyncService {
|
||||
this.nextScheduledSyncAt = null;
|
||||
this.caseFeatureBatchSize = parseInt(process.env.ORDRESTYRING_CASE_FEATURE_BATCH_SIZE || '200', 10);
|
||||
this.caseLatestCacheReady = null;
|
||||
this.pendingOfferSyncCursor = undefined;
|
||||
|
||||
// Database connection for Ordrestyring data
|
||||
this.dbConfig = {
|
||||
@@ -1913,77 +1914,125 @@ class OrdrestyringSyncService {
|
||||
);
|
||||
}
|
||||
|
||||
async fetchOfferSnapshots(limit = 50) {
|
||||
const pageSize = Math.max(1, parseInt(process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE || `${limit || 100}`, 10) || 100);
|
||||
const maxOffers = Math.max(0, parseInt(process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS || '0', 10) || 0);
|
||||
getOfferSyncLimits(limit = 50) {
|
||||
const pageSize = Math.min(200, Math.max(20,
|
||||
parseInt(process.env.ORDRESTYRING_OFFER_SYNC_PAGE_SIZE || `${limit || 100}`, 10) || 100));
|
||||
const maxOffers = Math.min(1000, Math.max(21,
|
||||
parseInt(process.env.ORDRESTYRING_OFFER_SYNC_MAX_OFFERS || '100', 10) || 100));
|
||||
return { pageSize, maxOffers };
|
||||
}
|
||||
|
||||
async ensureOfferSyncStateTable(connection) {
|
||||
await connection.execute(`
|
||||
CREATE TABLE IF NOT EXISTS ordrestyring_offer_sync_state (
|
||||
state_key VARCHAR(64) NOT NULL PRIMARY KEY,
|
||||
next_cursor TEXT NULL,
|
||||
updated_at BIGINT NOT NULL
|
||||
) ENGINE=InnoDB
|
||||
`);
|
||||
}
|
||||
|
||||
async loadOfferSyncCursor() {
|
||||
let connection;
|
||||
try {
|
||||
connection = await mysql.createConnection(this.dbConfig);
|
||||
await this.ensureOfferSyncStateTable(connection);
|
||||
const [rows] = await connection.execute(
|
||||
'SELECT next_cursor FROM ordrestyring_offer_sync_state WHERE state_key = ?',
|
||||
['offer_backfill']
|
||||
);
|
||||
return rows[0]?.next_cursor || null;
|
||||
} finally {
|
||||
await this.closeConnectionSafely(connection, 'Offer sync cursor read');
|
||||
}
|
||||
}
|
||||
|
||||
async saveOfferSyncCursor(nextCursor) {
|
||||
let connection;
|
||||
try {
|
||||
connection = await mysql.createConnection(this.dbConfig);
|
||||
await this.ensureOfferSyncStateTable(connection);
|
||||
await connection.execute(`
|
||||
INSERT INTO ordrestyring_offer_sync_state (state_key, next_cursor, updated_at)
|
||||
VALUES (?, ?, ?)
|
||||
ON DUPLICATE KEY UPDATE
|
||||
next_cursor = VALUES(next_cursor),
|
||||
updated_at = VALUES(updated_at)
|
||||
`, ['offer_backfill', nextCursor || null, Math.floor(Date.now() / 1000)]);
|
||||
} finally {
|
||||
await this.closeConnectionSafely(connection, 'Offer sync cursor write');
|
||||
}
|
||||
}
|
||||
|
||||
async requestOfferSnapshotPage(cursor, limit) {
|
||||
const query = `
|
||||
query SyncOfferSnapshots($limit: Int!, $cursor: String) {
|
||||
offers(
|
||||
pagination: { cursor: $cursor, limit: $limit },
|
||||
orderBy: { field: "createdAt", direction: DESC }
|
||||
) {
|
||||
items {
|
||||
id
|
||||
number
|
||||
}
|
||||
items { id number }
|
||||
count
|
||||
hasMorePages
|
||||
nextCursor
|
||||
}
|
||||
}
|
||||
`;
|
||||
const response = await graphqlScheduler.schedule(signal => axios.post(
|
||||
'https://graphql.ordrestyring.dk/graphql',
|
||||
{ query, variables: { limit, cursor } },
|
||||
{
|
||||
headers: {
|
||||
Authorization: `Bearer ${this.apiToken}`,
|
||||
'Content-Type': 'application/json'
|
||||
},
|
||||
timeout: 30000,
|
||||
signal
|
||||
}
|
||||
));
|
||||
if (Array.isArray(response.data?.errors) && response.data.errors.length > 0) {
|
||||
throw new Error(response.data.errors[0]?.message || 'Ordrestyring GraphQL returned errors for offer list');
|
||||
}
|
||||
const page = response.data?.data?.offers;
|
||||
if (!page || !Array.isArray(page.items) || typeof page.hasMorePages !== 'boolean') {
|
||||
throw new Error('Invalid Ordrestyring GraphQL offer page');
|
||||
}
|
||||
const nextCursor = page.nextCursor || null;
|
||||
if (page.hasMorePages && !nextCursor) {
|
||||
throw new Error('Offer pagination reports more pages without a cursor');
|
||||
}
|
||||
return { items: page.items, hasMorePages: page.hasMorePages, nextCursor };
|
||||
}
|
||||
|
||||
async fetchOfferSnapshots(limit = 50) {
|
||||
const { pageSize, maxOffers } = this.getOfferSyncLimits(limit);
|
||||
this.pendingOfferSyncCursor = undefined;
|
||||
try {
|
||||
const offers = [];
|
||||
const seenOfferIds = new Set();
|
||||
let cursor = null;
|
||||
let hasMorePages = true;
|
||||
|
||||
while (hasMorePages) {
|
||||
const response = await graphqlScheduler.schedule(signal => axios.post(
|
||||
'https://graphql.ordrestyring.dk/graphql',
|
||||
{
|
||||
query,
|
||||
variables: { limit: pageSize, cursor }
|
||||
},
|
||||
{
|
||||
headers: {
|
||||
Authorization: `Bearer ${this.apiToken}`,
|
||||
'Content-Type': 'application/json'
|
||||
},
|
||||
timeout: 30000,
|
||||
signal
|
||||
}
|
||||
));
|
||||
|
||||
if (Array.isArray(response.data?.errors) && response.data.errors.length > 0) {
|
||||
throw new Error(response.data.errors[0]?.message || 'Ordrestyring GraphQL returned errors for offer list');
|
||||
}
|
||||
|
||||
const page = response.data?.data?.offers;
|
||||
const items = Array.isArray(page?.items) ? page.items : [];
|
||||
items.forEach(item => {
|
||||
if (!item?.id || seenOfferIds.has(item.id)) {
|
||||
return;
|
||||
}
|
||||
const appendItems = items => {
|
||||
for (const item of items) {
|
||||
if (!item?.id || seenOfferIds.has(item.id) || offers.length >= maxOffers) continue;
|
||||
seenOfferIds.add(item.id);
|
||||
offers.push(item);
|
||||
});
|
||||
|
||||
if (maxOffers > 0 && offers.length >= maxOffers) {
|
||||
return offers.slice(0, maxOffers);
|
||||
}
|
||||
};
|
||||
|
||||
hasMorePages = Boolean(page?.hasMorePages);
|
||||
cursor = page?.nextCursor || null;
|
||||
if (hasMorePages && !cursor) {
|
||||
logger.warn('Offer pagination reports more pages but no nextCursor; stopping early', {
|
||||
fetchedOffers: offers.length
|
||||
});
|
||||
break;
|
||||
}
|
||||
const freshLimit = Math.min(20, pageSize, maxOffers);
|
||||
const freshPage = await this.requestOfferSnapshotPage(null, freshLimit);
|
||||
appendItems(freshPage.items);
|
||||
|
||||
let cursor = await this.loadOfferSyncCursor();
|
||||
if (!cursor) cursor = freshPage.hasMorePages ? freshPage.nextCursor : null;
|
||||
|
||||
while (cursor && offers.length < maxOffers) {
|
||||
const remaining = maxOffers - offers.length;
|
||||
const page = await this.requestOfferSnapshotPage(cursor, Math.min(pageSize, remaining));
|
||||
appendItems(page.items);
|
||||
cursor = page.hasMorePages ? page.nextCursor : null;
|
||||
}
|
||||
|
||||
this.pendingOfferSyncCursor = cursor;
|
||||
return offers;
|
||||
} catch (error) {
|
||||
throw new Error(`Failed to load offer list from Ordrestyring GraphQL: ${error.message}`);
|
||||
@@ -2219,6 +2268,10 @@ class OrdrestyringSyncService {
|
||||
}
|
||||
}
|
||||
|
||||
if (failedDetails === 0 && this.pendingOfferSyncCursor !== undefined) {
|
||||
await this.saveOfferSyncCursor(this.pendingOfferSyncCursor);
|
||||
}
|
||||
|
||||
return {
|
||||
hasChanges: changes > 0,
|
||||
changes,
|
||||
|
||||
Reference in New Issue
Block a user