Skip to content

Commit 86696bf

Browse files
committed
Apply code review fixes
1 parent 48eaa1a commit 86696bf

3 files changed

Lines changed: 24 additions & 7 deletions

File tree

‎packages/typescript-client/src/client.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1571,6 +1571,13 @@ export class ShapeStream<T extends Row<unknown> = Row>
15711571
// Filter messages using snapshot tracker
15721572
const messagesToProcess = batch.filter((message) => {
15731573
if (isChangeMessage(message)) {
1574+
const changeLsn = message.headers.lsn
1575+
if (typeof changeLsn === `string` && changeLsn) {
1576+
// A quiet SSE stream may deliver its first later change before the
1577+
// next up-to-date boundary. Retire snapshots against that change's
1578+
// WAL position before resolving its wrapped transaction ID.
1579+
this.#snapshotTracker.lastSeenUpdate(BigInt(changeLsn))
1580+
}
15741581
return !this.#snapshotTracker.shouldRejectMessage(message)
15751582
}
15761583

‎packages/typescript-client/test/pbt-micro.test.ts‎

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -546,12 +546,16 @@ describe(`SnapshotTracker database LSN retirement`, () => {
546546
{
547547
key: `k1`,
548548
value: { id: 1, version: `duplicate` },
549-
headers: { operation: `update`, txids: [679_865_407] },
549+
headers: {
550+
operation: `update`,
551+
txids: [679_865_407],
552+
lsn: `122`,
553+
},
550554
},
551555
{
552556
headers: {
553557
control: `up-to-date`,
554-
global_last_seen_lsn: `123`,
558+
global_last_seen_lsn: `100`,
555559
},
556560
},
557561
]),
@@ -565,7 +569,11 @@ describe(`SnapshotTracker database LSN retirement`, () => {
565569
{
566570
key: `k1`,
567571
value: { id: 1, version: `future` },
568-
headers: { operation: `update`, txids: [2_827_350_156] },
572+
headers: {
573+
operation: `update`,
574+
txids: [2_827_350_156],
575+
lsn: `124`,
576+
},
569577
},
570578
{
571579
headers: {

‎packages/typescript-client/test/snapshot-tracker.test.ts‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -116,7 +116,7 @@ describe(`SnapshotTracker`, () => {
116116
}
117117
tracker.addSnapshot(metadata, new Set([`user:1`]))
118118

119-
const message: ChangeMessage<Row<unknown>> = {
119+
const newerMessage: ChangeMessage<Row<unknown>> = {
120120
key: `user:1`,
121121
value: { id: 1, name: `Alice` },
122122
headers: {
@@ -125,9 +125,11 @@ describe(`SnapshotTracker`, () => {
125125
},
126126
}
127127

128-
expect(tracker.shouldRejectMessage(message)).toBe(false)
128+
expect(tracker.shouldRejectMessage(newerMessage)).toBe(false)
129129

130-
const olderMessage: ChangeMessage<Row<unknown>> = {
130+
// This would be rejected if the newer message had not retired the
131+
// snapshot after resolving its wrapped xid against xmax.
132+
const messageAfterRetirement: ChangeMessage<Row<unknown>> = {
131133
key: `user:1`,
132134
value: { id: 1, name: `Alice` },
133135
headers: {
@@ -136,7 +138,7 @@ describe(`SnapshotTracker`, () => {
136138
},
137139
}
138140

139-
expect(tracker.shouldRejectMessage(olderMessage)).toBe(false)
141+
expect(tracker.shouldRejectMessage(messageAfterRetirement)).toBe(false)
140142
})
141143
})
142144

0 commit comments

Comments
 (0)