Skip to content

Commit e4d589e

Browse files
fubhytim-smart
andauthored
Fix clearAddress leaves primary-key deduplication state (#7038)
Co-authored-by: Tim Smart <hello@timsmart.co>
1 parent dce8219 commit e4d589e

3 files changed

Lines changed: 39 additions & 0 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"effect": patch
3+
---
4+
5+
Clear in-memory message primary-key indexes when clearing an entity address.

packages/effect/src/unstable/cluster/MessageStorage.ts

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -991,6 +991,14 @@ export class MemoryDriver extends Context.Service<MemoryDriver>()("effect/cluste
991991
resetAddress: () => Effect.void,
992992
clearAddress: (address) =>
993993
Effect.sync(() => {
994+
for (const [primaryKey, entry] of requestsByPrimaryKey) {
995+
const envelope = entry.envelope
996+
const sameAddress = address.entityType === envelope.address.entityType &&
997+
address.entityId === envelope.address.entityId
998+
if (sameAddress) {
999+
requestsByPrimaryKey.delete(primaryKey)
1000+
}
1001+
}
9941002
for (let i = journal.length - 1; i >= 0; i--) {
9951003
const envelope = journal[i]
9961004
const sameAddress = address.entityType === envelope.address.entityType &&

packages/effect/test/cluster/MessageStorage.test.ts

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,32 @@ const MemoryLive = MessageStorage.layerMemory.pipe(
2323

2424
describe("MessageStorage", () => {
2525
describe("memory", () => {
26+
it.effect("removes the primary-key index when clearing an address", () =>
27+
Effect.gen(function*() {
28+
const driver = yield* MessageStorage.MemoryDriver
29+
const address = EntityAddress.make({
30+
shardId: ShardId.make("default", 1),
31+
entityType: EntityType.make("Repro"),
32+
entityId: EntityId.make("one")
33+
})
34+
const envelope: Envelope.PartialRequestEncoded = {
35+
_tag: "Request",
36+
requestId: "1",
37+
address: { shardId: { group: "default", id: 1 }, entityType: "Repro", entityId: "one" },
38+
tag: "Repro",
39+
payload: {},
40+
headers: {}
41+
}
42+
yield* driver.encoded.saveEnvelope({ envelope, primaryKey: "dedup-key", deliverAt: null })
43+
yield* driver.encoded.clearAddress(address)
44+
const result = yield* driver.encoded.saveEnvelope({
45+
envelope: { ...envelope, requestId: "2" },
46+
primaryKey: "dedup-key",
47+
deliverAt: null
48+
})
49+
expect(result._tag).toEqual("Success")
50+
}).pipe(Effect.provide(MessageStorage.MemoryDriver.layer)))
51+
2652
it.effect("saves a request", () =>
2753
Effect.gen(function*() {
2854
const storage = yield* MessageStorage.MessageStorage

0 commit comments

Comments
 (0)