Skip to content

Commit 57b6570

Browse files
committed
fix(tunnel): expire incomplete reassemblies
1 parent c0e52d2 commit 57b6570

9 files changed

Lines changed: 95 additions & 58 deletions

File tree

networking/internal/tunnel/build_manager_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -834,7 +834,7 @@ func TestBuildManagerBuildsInboundAcrossTransitAndRejectsStaleRequests(t *testin
834834
t.Fatalf("inbound carrier message = %#v", carried)
835835
}
836836
blocks := make([]Block, 1)
837-
count, err := NewEndpoint(1, i2np.I2PDMaxPayload).Parse(carried[0].message.Payload, blocks)
837+
count, err := NewEndpoint(1, i2np.I2PDMaxPayload).Parse(carried[0].message.Payload, blocks, 0)
838838
if err != nil || count != 1 || blocks[0].Delivery != DeliveryRouter || blocks[0].Gateway != build.Hops[0].Router {
839839
t.Fatalf("carrier blocks = %d, %#v, %v", count, blocks[0], err)
840840
}
@@ -1291,7 +1291,7 @@ func TestInboundBuildGarlicWrapsAcrossDifferentCarrierEndpoint(t *testing.T) {
12911291
t.Fatalf("carrier output = %#v", carried)
12921292
}
12931293
blocks := make([]Block, 1)
1294-
count, err := NewEndpoint(1, i2np.I2PDMaxPayload).Parse(carried[0].message.Payload, blocks)
1294+
count, err := NewEndpoint(1, i2np.I2PDMaxPayload).Parse(carried[0].message.Payload, blocks, 0)
12951295
if err != nil || count != 1 {
12961296
t.Fatalf("carrier parse = %d, %v", count, err)
12971297
}

networking/internal/tunnel/gateway.go

Lines changed: 16 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -270,13 +270,12 @@ type Endpoint struct {
270270
mu sync.Mutex
271271
reasm *Reassembler
272272
maxMeta int
273-
clock uint64
274273
meta map[uint32]deliveryMeta
275274
}
276275

277276
type deliveryMeta struct {
278277
block Block
279-
touched uint64
278+
created uint64
280279
}
281280

282281
func NewEndpoint(maxEntries, maxMessage int) *Endpoint {
@@ -291,11 +290,11 @@ func NewEndpoint(maxEntries, maxMessage int) *Endpoint {
291290
}
292291

293292
// Parse reads a complete 1028-byte I2NP TunnelData payload and writes completed
294-
// deliveries to out. It returns ErrGatewayOutput if out cannot hold all
295-
// completed deliveries. Follow-on fragments may arrive before their initial
296-
// fragment; no delivery is emitted until both delivery metadata and every
297-
// fragment are present.
298-
func (e *Endpoint) Parse(payload []byte, out []Block) (int, error) {
293+
// deliveries to out. nowMillis timestamps incomplete reassemblies. It returns
294+
// ErrGatewayOutput if out cannot hold all completed deliveries. Follow-on
295+
// fragments may arrive before their initial fragment; no delivery is emitted
296+
// until both delivery metadata and every fragment are present.
297+
func (e *Endpoint) Parse(payload []byte, out []Block, nowMillis uint64) (int, error) {
299298
message, err := i2np.ParseTunnelData(payload)
300299
if err != nil {
301300
return 0, err
@@ -320,7 +319,6 @@ func (e *Endpoint) Parse(payload []byte, out []Block) (int, error) {
320319
}
321320
e.mu.Lock()
322321
defer e.mu.Unlock()
323-
e.clock++
324322
it := NewBlockIterator(data[start+1:])
325323
n := 0
326324
for {
@@ -335,7 +333,7 @@ func (e *Endpoint) Parse(payload []byte, out []Block) (int, error) {
335333
if block.Fragment == 0 {
336334
return n, ErrGatewayBlock
337335
}
338-
complete, done, err := e.reasm.Add(Fragment{MessageID: block.MessageID, Number: block.Fragment, Last: block.Last, Data: block.Data})
336+
complete, done, err := e.reasm.Add(Fragment{MessageID: block.MessageID, Number: block.Fragment, Last: block.Last, Data: block.Data}, nowMillis)
339337
if err != nil {
340338
return n, err
341339
}
@@ -363,10 +361,10 @@ func (e *Endpoint) Parse(payload []byte, out []Block) (int, error) {
363361
continue
364362
}
365363

366-
if err := e.remember(block); err != nil {
364+
if err := e.remember(block, nowMillis); err != nil {
367365
return n, err
368366
}
369-
complete, done, err := e.reasm.Add(Fragment{MessageID: block.MessageID, Number: 0, Data: block.Data})
367+
complete, done, err := e.reasm.Add(Fragment{MessageID: block.MessageID, Number: 0, Data: block.Data}, nowMillis)
370368
if err != nil {
371369
return n, err
372370
}
@@ -383,38 +381,37 @@ func (e *Endpoint) Parse(payload []byte, out []Block) (int, error) {
383381
}
384382
}
385383

386-
func (e *Endpoint) remember(block Block) error {
384+
func (e *Endpoint) remember(block Block, nowMillis uint64) error {
387385
block.Data = nil
388386
if prior, exists := e.meta[block.MessageID]; exists {
389387
if prior.block.Delivery != block.Delivery || prior.block.Gateway != block.Gateway || prior.block.TunnelID != block.TunnelID {
390388
return ErrFragment
391389
}
392-
prior.touched = e.clock
393390
e.meta[block.MessageID] = prior
394391
return nil
395392
}
396393
if len(e.meta) == e.maxMeta {
397394
var oldestID uint32
398395
var oldest uint64 = ^uint64(0)
399396
for id, item := range e.meta {
400-
if item.touched < oldest {
401-
oldestID, oldest = id, item.touched
397+
if item.created < oldest {
398+
oldestID, oldest = id, item.created
402399
}
403400
}
404401
delete(e.meta, oldestID)
405402
}
406-
e.meta[block.MessageID] = deliveryMeta{block: block, touched: e.clock}
403+
e.meta[block.MessageID] = deliveryMeta{block: block, created: nowMillis}
407404
return nil
408405
}
409406

410-
// Expire removes retained fragment metadata and incomplete reassemblies not
411-
// touched since the supplied endpoint clock tick.
407+
// Expire removes retained fragment metadata and incomplete reassemblies
408+
// created at or before cutoff.
412409
func (e *Endpoint) Expire(cutoff uint64) int {
413410
e.mu.Lock()
414411
defer e.mu.Unlock()
415412
removed := e.reasm.Expire(cutoff)
416413
for id, item := range e.meta {
417-
if item.touched <= cutoff {
414+
if item.created <= cutoff {
418415
delete(e.meta, id)
419416
removed++
420417
}

networking/internal/tunnel/gateway_test.go

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ func TestGatewayRoundTripsDeliveryInstructions(t *testing.T) {
3232
if !ok {
3333
t.Fatal("payload unavailable")
3434
}
35-
n, err := endpoint.Parse(payload, got)
35+
n, err := endpoint.Parse(payload, got, 0)
3636
if err != nil {
3737
t.Fatal(err)
3838
}
@@ -79,15 +79,15 @@ func TestGatewayFragmentsAndReassembles(t *testing.T) {
7979
if !ok {
8080
t.Fatalf("fragment %d payload unavailable", i)
8181
}
82-
if count, err := endpoint.Parse(payload, out); err != nil || count != 0 {
82+
if count, err := endpoint.Parse(payload, out, uint64(i)); err != nil || count != 0 {
8383
t.Fatalf("fragment %d = %d, %v", i, count, err)
8484
}
8585
}
8686
payload, ok := buffers[n-1].Payload()
8787
if !ok {
8888
t.Fatal("final fragment payload unavailable")
8989
}
90-
count, err := endpoint.Parse(payload, out)
90+
count, err := endpoint.Parse(payload, out, uint64(n))
9191
if err != nil || count != 1 {
9292
t.Fatalf("final fragment = %d, %v", count, err)
9393
}
@@ -104,7 +104,7 @@ func TestGatewayRejectsTruncatedAndOversizedBlocks(t *testing.T) {
104104
t.Fatal("oversized block accepted")
105105
}
106106
endpoint := NewEndpoint(1, 1)
107-
if _, err := endpoint.Parse(make([]byte, i2np.TunnelDataMessageLen-1), nil); err == nil {
107+
if _, err := endpoint.Parse(make([]byte, i2np.TunnelDataMessageLen-1), nil, 0); err == nil {
108108
t.Fatal("truncated TunnelData accepted")
109109
}
110110
}
@@ -122,7 +122,7 @@ func TestGatewayRejectsChecksumMismatch(t *testing.T) {
122122
t.Fatal("payload unavailable")
123123
}
124124
payload[len(payload)-1] ^= 1
125-
if _, err := NewEndpoint(1, 64).Parse(payload, make([]Block, 1)); err != ErrGatewayPayload {
125+
if _, err := NewEndpoint(1, 64).Parse(payload, make([]Block, 1), 0); err != ErrGatewayPayload {
126126
t.Fatalf("tampered payload error = %v, want %v", err, ErrGatewayPayload)
127127
}
128128
}

networking/internal/tunnel/java_compatibility_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -100,7 +100,7 @@ func testJavaUnknownI2NPRelay(t testing.TB, record int, fixture []byte) {
100100
t.Fatalf("record %d: %v", record, err)
101101
}
102102
blocks := make([]Block, 1)
103-
count, err := NewEndpoint(1, i2np.I2PDMaxPayload).Parse(payload, blocks)
103+
count, err := NewEndpoint(1, i2np.I2PDMaxPayload).Parse(payload, blocks, 0)
104104
if err != nil || count != 1 || blocks[0].Delivery != DeliveryLocal {
105105
t.Fatalf("record %d: gateway payload blocks = %d, %#v, %v", record, count, blocks[0], err)
106106
}

networking/internal/tunnel/reassembly.go

Lines changed: 7 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -29,15 +29,14 @@ type partial struct {
2929
last uint8
3030
hasLast bool
3131
size int
32-
touched uint64
32+
created uint64
3333
}
3434

3535
// Reassembler bounds retained incomplete data by entry count and message size.
3636
type Reassembler struct {
3737
mu sync.Mutex
3838
maxEntries int
3939
maxMessage int
40-
clock uint64
4140
entries map[uint32]*partial
4241
partialPool sync.Pool
4342
}
@@ -56,22 +55,21 @@ func NewReassembler(maxEntries, maxMessage int) *Reassembler {
5655

5756
// Add copies one fragment and returns a complete caller-owned message only
5857
// after every fragment from zero through the final number is present.
59-
func (r *Reassembler) Add(fragment Fragment) ([]byte, bool, error) {
58+
func (r *Reassembler) Add(fragment Fragment, nowMillis uint64) ([]byte, bool, error) {
6059
if len(fragment.Data) == 0 || fragment.Number > 127 {
6160
return nil, false, ErrFragment
6261
}
6362
r.mu.Lock()
6463
defer r.mu.Unlock()
65-
r.clock++
6664
entry := r.entries[fragment.MessageID]
6765
if entry == nil {
6866
if len(r.entries) == r.maxEntries {
6967
r.evictOldest()
7068
}
7169
entry = r.partialPool.Get().(*partial)
7270
r.entries[fragment.MessageID] = entry
71+
entry.created = nowMillis
7372
}
74-
entry.touched = r.clock
7573
if entry.hasLast && fragment.Number > entry.last {
7674
r.remove(fragment.MessageID)
7775
return nil, false, ErrFragment
@@ -125,13 +123,13 @@ func (r *Reassembler) Add(fragment Fragment) ([]byte, bool, error) {
125123
return message, true, nil
126124
}
127125

128-
// Expire removes incomplete assemblies not touched since cutoff clock ticks.
126+
// Expire removes incomplete assemblies created at or before cutoff.
129127
func (r *Reassembler) Expire(cutoff uint64) int {
130128
r.mu.Lock()
131129
defer r.mu.Unlock()
132130
removed := 0
133131
for id, entry := range r.entries {
134-
if entry.touched < cutoff {
132+
if entry.created <= cutoff {
135133
r.remove(id)
136134
removed++
137135
}
@@ -143,8 +141,8 @@ func (r *Reassembler) evictOldest() {
143141
var oldestID uint32
144142
var oldest uint64 = ^uint64(0)
145143
for id, entry := range r.entries {
146-
if entry.touched < oldest {
147-
oldestID, oldest = id, entry.touched
144+
if entry.created < oldest {
145+
oldestID, oldest = id, entry.created
148146
}
149147
}
150148
r.remove(oldestID)

networking/internal/tunnel/reassembly_conflict_test.go

Lines changed: 15 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -7,52 +7,52 @@ import (
77

88
func TestReassemblerRejectsAmbiguousLastAndDuplicates(t *testing.T) {
99
r := NewReassembler(2, 64)
10-
if _, _, err := r.Add(Fragment{MessageID: 1, Number: 5, Data: []byte("late")}); err != nil {
10+
if _, _, err := r.Add(Fragment{MessageID: 1, Number: 5, Data: []byte("late")}, 1); err != nil {
1111
t.Fatal(err)
1212
}
13-
if _, _, err := r.Add(Fragment{MessageID: 1, Number: 0, Last: true, Data: []byte("first")}); !errors.Is(err, ErrFragment) {
13+
if _, _, err := r.Add(Fragment{MessageID: 1, Number: 0, Last: true, Data: []byte("first")}, 1); !errors.Is(err, ErrFragment) {
1414
t.Fatalf("late fragment conflict = %v", err)
1515
}
16-
if _, _, err := r.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("a")}); err != nil {
16+
if _, _, err := r.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("a")}, 2); err != nil {
1717
t.Fatal(err)
1818
}
19-
if _, _, err := r.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("b")}); !errors.Is(err, ErrFragment) {
19+
if _, _, err := r.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("b")}, 2); !errors.Is(err, ErrFragment) {
2020
t.Fatalf("duplicate conflict = %v", err)
2121
}
2222
r = NewReassembler(2, 64)
23-
if _, _, err := r.Add(Fragment{MessageID: 3, Number: 0, Data: []byte("same")}); err != nil {
23+
if _, _, err := r.Add(Fragment{MessageID: 3, Number: 0, Data: []byte("same")}, 3); err != nil {
2424
t.Fatal(err)
2525
}
26-
if _, _, err := r.Add(Fragment{MessageID: 3, Number: 0, Last: true, Data: []byte("same")}); !errors.Is(err, ErrFragment) {
26+
if _, _, err := r.Add(Fragment{MessageID: 3, Number: 0, Last: true, Data: []byte("same")}, 3); !errors.Is(err, ErrFragment) {
2727
t.Fatalf("duplicate terminal conflict = %v", err)
2828
}
2929
}
3030

3131
func TestReassemblerEvictsAndExpiresIncompleteEntries(t *testing.T) {
3232
r := NewReassembler(1, 64)
33-
_, _, _ = r.Add(Fragment{MessageID: 1, Number: 0, Data: []byte("a")})
34-
_, _, _ = r.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("b")})
33+
_, _, _ = r.Add(Fragment{MessageID: 1, Number: 0, Data: []byte("a")}, 1)
34+
_, _, _ = r.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("b")}, 2)
3535
if len(r.entries) != 1 || r.entries[2] == nil {
3636
t.Fatalf("LRU entries = %#v", r.entries)
3737
}
38-
if removed := r.Expire(r.clock + 1); removed != 1 || len(r.entries) != 0 {
38+
if removed := r.Expire(2); removed != 1 || len(r.entries) != 0 {
3939
t.Fatalf("Expire = %d, entries=%d", removed, len(r.entries))
4040
}
4141
}
4242

4343
func TestReassemblerPooledEntryDoesNotRetainEvictedFragments(t *testing.T) {
4444
reassembler := NewReassembler(1, 64)
45-
if _, _, err := reassembler.Add(Fragment{MessageID: 1, Number: 0, Data: []byte("old")}); err != nil {
45+
if _, _, err := reassembler.Add(Fragment{MessageID: 1, Number: 0, Data: []byte("old")}, 1); err != nil {
4646
t.Fatal(err)
4747
}
48-
if _, _, err := reassembler.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("new")}); err != nil {
48+
if _, _, err := reassembler.Add(Fragment{MessageID: 2, Number: 0, Data: []byte("new")}, 2); err != nil {
4949
t.Fatal(err)
5050
}
51-
message, done, err := reassembler.Add(Fragment{MessageID: 2, Number: 1, Last: true, Data: []byte(" tail")})
51+
message, done, err := reassembler.Add(Fragment{MessageID: 2, Number: 1, Last: true, Data: []byte(" tail")}, 3)
5252
if err != nil || !done || string(message) != "new tail" {
5353
t.Fatalf("reused entry message = %q, %t, %v", message, done, err)
5454
}
55-
if message, done, err = reassembler.Add(Fragment{MessageID: 1, Number: 1, Last: true, Data: []byte(" tail")}); err != nil || done || message != nil {
55+
if message, done, err = reassembler.Add(Fragment{MessageID: 1, Number: 1, Last: true, Data: []byte(" tail")}, 4); err != nil || done || message != nil {
5656
t.Fatalf("evicted message resumed with pooled data = %q, %t, %v", message, done, err)
5757
}
5858
}
@@ -65,10 +65,10 @@ func BenchmarkReassemblerPooledFragments(b *testing.B) {
6565
for b.Loop() {
6666
first.MessageID++
6767
last.MessageID = first.MessageID
68-
if _, _, err := reassembler.Add(first); err != nil {
68+
if _, _, err := reassembler.Add(first, uint64(first.MessageID)); err != nil {
6969
b.Fatal(err)
7070
}
71-
if _, done, err := reassembler.Add(last); err != nil || !done {
71+
if _, done, err := reassembler.Add(last, uint64(last.MessageID)); err != nil || !done {
7272
b.Fatalf("complete = %t, %v", done, err)
7373
}
7474
}

networking/internal/tunnel/reassembly_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,10 +7,10 @@ import (
77

88
func TestReassemblerHandlesOutOfOrderFragments(t *testing.T) {
99
r := NewReassembler(2, 16)
10-
if _, done, err := r.Add(Fragment{MessageID: 1, Number: 1, Last: true, Data: []byte("world")}); err != nil || done {
10+
if _, done, err := r.Add(Fragment{MessageID: 1, Number: 1, Last: true, Data: []byte("world")}, 1); err != nil || done {
1111
t.Fatalf("last fragment = %t, %v", done, err)
1212
}
13-
message, done, err := r.Add(Fragment{MessageID: 1, Number: 0, Data: []byte("hello ")})
13+
message, done, err := r.Add(Fragment{MessageID: 1, Number: 0, Data: []byte("hello ")}, 2)
1414
if err != nil || !done || !bytes.Equal(message, []byte("hello world")) {
1515
t.Fatalf("message = %q, %t, %v", message, done, err)
1616
}

networking/internal/tunnel/runtime.go

Lines changed: 16 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,9 +16,11 @@ import (
1616
)
1717

1818
const (
19-
defaultTunnelTTL = time.Minute
20-
defaultDeliveryBlocks = TunnelPayloadLen / 3
21-
circuitShards = 64
19+
defaultTunnelTTL = time.Minute
20+
// Java I2P FragmentHandler removes incomplete messages after 45 seconds.
21+
fragmentReassemblyLifetime = 45 * time.Second
22+
defaultDeliveryBlocks = TunnelPayloadLen / 3
23+
circuitShards = 64
2224
)
2325

2426
var (
@@ -338,13 +340,23 @@ func (r *Runtime) Expire(nowMillis uint64) (removed int) {
338340
if r == nil {
339341
return 0
340342
}
343+
fragmentLifetimeMillis := uint64(fragmentReassemblyLifetime / time.Millisecond)
344+
expireFragments := nowMillis >= fragmentLifetimeMillis
345+
var fragmentCutoff uint64
346+
if expireFragments {
347+
fragmentCutoff = nowMillis - fragmentLifetimeMillis
348+
}
341349
for index := range r.shards {
342350
shard := &r.shards[index]
343351
shard.mu.Lock()
344352
for id, circuit := range shard.inbound {
345353
if circuit.expiresAt != 0 && circuit.expiresAt <= nowMillis {
346354
delete(shard.inbound, id)
347355
removed++
356+
continue
357+
}
358+
if expireFragments && circuit.endpoint != nil {
359+
circuit.endpoint.Expire(fragmentCutoff)
348360
}
349361
}
350362
for id, circuit := range shard.outbound {
@@ -570,7 +582,7 @@ func (r *Runtime) HandleContext(ctx context.Context, message i2np.Message) error
570582
}
571583

572584
blocks := r.blocks.Get().(*[]Block)
573-
count, err := circuit.endpoint.Parse(payload, *blocks)
585+
count, err := circuit.endpoint.Parse(payload, *blocks, r.now())
574586
if err == nil {
575587
sender := r.currentSender()
576588
for _, block := range (*blocks)[:count] {

0 commit comments

Comments
 (0)