| 687 | } |
| 688 | |
| 689 | func TestMultipleUDPSends(t *testing.T) { |
| 690 | addr := "127.0.0.1:8126" |
| 691 | |
| 692 | address, _ := net.ResolveUDPAddr("udp", addr) |
| 693 | listener, err := net.ListenUDP("udp", address) |
| 694 | assert.Equal(t, nil, err) |
| 695 | |
| 696 | ch := make(chan *Packet, MAX_UNPROCESSED_PACKETS) |
| 697 | |
| 698 | wg := &sync.WaitGroup{} |
| 699 | wg.Add(1) |
| 700 | go func() { |
| 701 | parseTo(listener, false, ch) |
| 702 | wg.Done() |
| 703 | }() |
| 704 | |
| 705 | conn, err := net.DialTimeout("udp", addr, 50*time.Millisecond) |
| 706 | assert.Equal(t, nil, err) |
| 707 | |
| 708 | n, err := conn.Write([]byte("deploys.test.myservice:2|c")) |
| 709 | assert.Equal(t, nil, err) |
| 710 | assert.Equal(t, len("deploys.test.myservice:2|c"), n) |
| 711 | |
| 712 | n, err = conn.Write([]byte("deploys.test.my:service:2|c")) |
| 713 | |
| 714 | n, err = conn.Write([]byte("deploys.test.myservice:1|c")) |
| 715 | assert.Equal(t, nil, err) |
| 716 | assert.Equal(t, len("deploys.test.myservice:1|c"), n) |
| 717 | |
| 718 | select { |
| 719 | case pack := <-ch: |
| 720 | assert.Equal(t, "deploys.test.myservice", pack.Bucket) |
| 721 | assert.Equal(t, float64(2), pack.ValFlt) |
| 722 | assert.Equal(t, "c", pack.Modifier) |
| 723 | assert.Equal(t, float32(1), pack.Sampling) |
| 724 | case <-time.After(50 * time.Millisecond): |
| 725 | t.Fatal("packet receive timeout") |
| 726 | } |
| 727 | |
| 728 | select { |
| 729 | case pack := <-ch: |
| 730 | assert.Equal(t, "deploys.test.myservice", pack.Bucket) |
| 731 | assert.Equal(t, float64(1), pack.ValFlt) |
| 732 | assert.Equal(t, "c", pack.Modifier) |
| 733 | assert.Equal(t, float32(1), pack.Sampling) |
| 734 | case <-time.After(50 * time.Millisecond): |
| 735 | t.Fatal("packet receive timeout") |
| 736 | } |
| 737 | |
| 738 | listener.Close() |
| 739 | wg.Wait() |
| 740 | } |
| 741 | |
| 742 | func BenchmarkManyDifferentSensors(t *testing.B) { |
| 743 | r := rand.New(rand.NewSource(438)) |