package proto_socket import ( "sync" "testing" "time" "google.golang.org/protobuf/proto" "git.toki-labs.com/toki/proto-socket/go/packets" ) // TestFrameReorderBufferReleasesInOrder: out-of-order results are buffered and // released only once the contiguous prefix is complete. func TestFrameReorderBufferReleasesInOrder(t *testing.T) { b := newFrameReorderBuffer() if got := b.release(workerResult{seq: 3}); len(got) != 0 { t.Fatalf("seq 3 should buffer until 1,2 arrive, released %d", len(got)) } if got := b.release(workerResult{seq: 2}); len(got) != 0 { t.Fatalf("seq 2 should still buffer until 1 arrives, released %d", len(got)) } got := b.release(workerResult{seq: 1}) if len(got) != 3 { t.Fatalf("expected 1,2,3 to flush together, got %d", len(got)) } for i, r := range got { if r.seq != int64(i+1) { t.Fatalf("reorder violated: position %d has seq %d", i, r.seq) } } // A seq below nextSeq is a duplicate/late arrival and is ignored. if got := b.release(workerResult{seq: 1}); len(got) != 0 { t.Fatalf("late duplicate seq should be ignored, released %d", len(got)) } } // TestWorkerGatewayReordersOutOfOrderCompletion: workers that finish out of // order still reach the sink in input seq order. Earlier seqs are made to decode // slower so the pool genuinely completes them last. func TestWorkerGatewayReordersOutOfOrderCompletion(t *testing.T) { const total = 8 var mu sync.Mutex var order []int32 parse := func(raw []byte) (*packets.PacketBase, error) { base := &packets.PacketBase{} if err := proto.Unmarshal(raw, base); err != nil { return nil, err } // Lower nonce (earlier seq) sleeps longer to force reversed completion. delay := time.Duration(total-base.GetNonce()) * 5 * time.Millisecond time.Sleep(delay) return base, nil } done := make(chan struct{}) sink := func(f DecodedFrame) { mu.Lock() order = append(order, f.IncomingNonce) n := len(order) mu.Unlock() if n == total { close(done) } } g := newWorkerGateway(4, 16, parse, sink, nil) defer g.Close() for seq := int64(1); seq <= total; seq++ { data, err := proto.Marshal(&packets.PacketBase{TypeName: "x", Nonce: int32(seq)}) if err != nil { t.Fatal(err) } g.Submit(InboundFrame{Seq: seq, Bytes: data}) } select { case <-done: case <-time.After(5 * time.Second): t.Fatal("timed out waiting for gateway results") } mu.Lock() defer mu.Unlock() for i, n := range order { if n != int32(i+1) { t.Fatalf("seq ordering not preserved at position %d: %v", i, order) } } } // TestWorkerGatewayReportsDecodeErrorAndAdvances: a frame that fails to decode is // reported to onError (not the sink) and does not stall later frames; subsequent // good frames still reach the sink in order. func TestWorkerGatewayReportsDecodeErrorAndAdvances(t *testing.T) { var mu sync.Mutex var seen []int32 var errCount int done := make(chan struct{}) sink := func(f DecodedFrame) { mu.Lock() seen = append(seen, f.IncomingNonce) n := len(seen) mu.Unlock() if n == 2 { close(done) } } onError := func(error) { mu.Lock() errCount++ mu.Unlock() } g := NewWorkerGateway(2, 8, sink, onError) defer g.Close() good1, _ := proto.Marshal(&packets.PacketBase{TypeName: "x", Nonce: 1}) good2, _ := proto.Marshal(&packets.PacketBase{TypeName: "x", Nonce: 3}) g.Submit(InboundFrame{Seq: 1, Bytes: good1}) g.Submit(InboundFrame{Seq: 2, Bytes: []byte{0xff, 0xff, 0xff}}) // undecodable g.Submit(InboundFrame{Seq: 3, Bytes: good2}) select { case <-done: case <-time.After(2 * time.Second): t.Fatal("timed out; a decode error stalled the reorder window") } mu.Lock() defer mu.Unlock() if len(seen) != 2 || seen[0] != 1 || seen[1] != 3 { t.Fatalf("expected sink to see [1 3] skipping the bad frame, got %v", seen) } if errCount != 1 { t.Fatalf("expected exactly one decode error reported, got %d", errCount) } }