Skip to content

Commit bd54b5e

Browse files
committed
Avoid duplicate inbound persistence
1 parent 87ca2dd commit bd54b5e

2 files changed

Lines changed: 89 additions & 2 deletions

File tree

epaxos/node.go

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -557,7 +557,8 @@ func (n *RawNode) handlePreAccept(m Message) {
557557
attrs := n.computeAttrs(m.Command, m.Ref)
558558
attrs = mergeAttrs(attrs, m.Attributes())
559559
rec := InstanceRecord{Ref: m.Ref, Ballot: m.Ballot, Status: StatusPreAccepted, Seq: attrs.Seq, Deps: attrs.Deps, Command: inboundCommand(m.Command)}
560-
if old := n.instances[m.Ref]; old != nil {
560+
old := n.instances[m.Ref]
561+
if old != nil {
561562
if old.rec.Status >= StatusCommitted {
562563
n.sendCommitTo(m.From, old.rec)
563564
return
@@ -568,6 +569,11 @@ func (n *RawNode) handlePreAccept(m Message) {
568569
}
569570
}
570571
rec.Checksum = ChecksumRecord(rec)
572+
if old != nil && old.rec.Status == StatusPreAccepted && instanceRecordEqual(old.rec, rec) {
573+
resp := Message{Type: MsgPreAcceptResp, From: n.id, To: m.From, Ref: m.Ref, Ballot: rec.Ballot, Seq: rec.Seq, Deps: rec.Deps, RecordStatus: rec.Status}
574+
n.enqueueMessage(resp)
575+
return
576+
}
571577
n.instances[m.Ref] = &instance{rec: rec, phase: phasePreAccept}
572578
n.indexConflicts(rec)
573579
n.enqueueRecord(rec)
@@ -610,7 +616,8 @@ func (n *RawNode) handlePreAcceptResp(m Message) {
610616
}
611617

612618
func (n *RawNode) handleAccept(m Message) {
613-
if old := n.instances[m.Ref]; old != nil {
619+
old := n.instances[m.Ref]
620+
if old != nil {
614621
if old.rec.Status >= StatusCommitted {
615622
n.sendCommitTo(m.From, old.rec)
616623
return
@@ -622,6 +629,11 @@ func (n *RawNode) handleAccept(m Message) {
622629
}
623630
rec := InstanceRecord{Ref: m.Ref, Ballot: m.Ballot, Status: StatusAccepted, Seq: m.Seq, Deps: append([]InstanceNum(nil), m.Deps...), Command: inboundCommand(m.Command)}
624631
rec.Checksum = ChecksumRecord(rec)
632+
if old != nil && old.rec.Status == StatusAccepted && instanceRecordEqual(old.rec, rec) {
633+
resp := Message{Type: MsgAcceptResp, From: n.id, To: m.From, Ref: m.Ref, Ballot: rec.Ballot, Seq: rec.Seq, Deps: rec.Deps, RecordStatus: rec.Status}
634+
n.enqueueMessage(resp)
635+
return
636+
}
625637
n.instances[m.Ref] = &instance{rec: rec, phase: phaseAccept}
626638
n.indexConflicts(rec)
627639
n.enqueueRecord(rec)

epaxos/sim_test.go

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -158,6 +158,81 @@ func TestDuplicateMessagesAndMalformedInput(t *testing.T) {
158158
}
159159
}
160160

161+
func TestDuplicateInboundPreAcceptAndAcceptDoNotQueueDuplicateRecords(t *testing.T) {
162+
ref := InstanceRef{Replica: 1, Instance: 1, Conf: 1}
163+
command := Command{ID: CommandID{Client: 10, Sequence: 20}, Payload: []byte("duplicate-inbound"), ConflictKeys: [][]byte{[]byte("duplicate-key")}}
164+
tests := []struct {
165+
name string
166+
msg Message
167+
wantStatus Status
168+
}{
169+
{
170+
name: "preaccept",
171+
msg: Message{
172+
Type: MsgPreAccept,
173+
From: 1,
174+
To: 2,
175+
Ref: ref,
176+
Ballot: Ballot{Replica: 1},
177+
Seq: 1,
178+
Deps: []InstanceNum{0, 0, 0},
179+
Command: command,
180+
},
181+
wantStatus: StatusPreAccepted,
182+
},
183+
{
184+
name: "accept",
185+
msg: Message{
186+
Type: MsgAccept,
187+
From: 1,
188+
To: 2,
189+
Ref: ref,
190+
Ballot: Ballot{Replica: 1},
191+
Seq: 1,
192+
Deps: []InstanceNum{0, 0, 0},
193+
Command: command,
194+
},
195+
wantStatus: StatusAccepted,
196+
},
197+
}
198+
for _, tt := range tests {
199+
t.Run(tt.name, func(t *testing.T) {
200+
store := NewMemoryStorage()
201+
rn, err := NewRawNode(Config{ID: 2, Voters: makeIDs(3), Storage: store})
202+
if err != nil {
203+
t.Fatal(err)
204+
}
205+
tt.msg.Checksum = ChecksumMessage(tt.msg)
206+
207+
if err := rn.Step(tt.msg); err != nil {
208+
t.Fatalf("first Step(%s) err=%v", tt.msg.Type, err)
209+
}
210+
first := rn.Ready()
211+
if len(first.Records) != 1 || first.Records[0].Ref != ref || first.Records[0].Status != tt.wantStatus {
212+
t.Fatalf("first duplicate target ready records = %#v, want exactly one %s record for %s", first.Records, tt.wantStatus, ref)
213+
}
214+
if err := store.ApplyReady(first); err != nil {
215+
t.Fatal(err)
216+
}
217+
advanceOK(t, rn, first)
218+
219+
if err := rn.Step(tt.msg); err != nil {
220+
t.Fatalf("duplicate Step(%s) err=%v", tt.msg.Type, err)
221+
}
222+
duplicate := rn.Ready()
223+
if len(duplicate.Records) != 0 {
224+
t.Fatalf("duplicate %s ready records = %#v, want no durable records for unchanged %s", tt.msg.Type, duplicate.Records, ref)
225+
}
226+
if !duplicate.Empty() {
227+
if err := store.ApplyReady(duplicate); err != nil {
228+
t.Fatal(err)
229+
}
230+
advanceOK(t, rn, duplicate)
231+
}
232+
})
233+
}
234+
}
235+
161236
func TestCodecChecksumZeroCopy(t *testing.T) {
162237
m := Message{
163238
Type: MsgCommit,

0 commit comments

Comments
 (0)