|
51 | 51 | import io.grpc.StringMarshaller; |
52 | 52 | import io.grpc.SynchronizationContext; |
53 | 53 | import io.grpc.internal.ClientStreamListener.RpcProgress; |
| 54 | +import java.util.ArrayList; |
| 55 | +import java.util.Arrays; |
| 56 | +import java.util.Collections; |
| 57 | +import java.util.List; |
54 | 58 | import java.util.concurrent.CyclicBarrier; |
55 | 59 | import java.util.concurrent.TimeUnit; |
56 | 60 | import java.util.concurrent.atomic.AtomicBoolean; |
@@ -774,152 +778,165 @@ public void pendingStream_appendTimeoutInsight_waitForReady_withLastPickFailure( |
774 | 778 |
|
775 | 779 | @Test |
776 | 780 | public void streamDelayMetrics() { |
777 | | - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); |
778 | | - ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer }; |
779 | | - |
780 | | - SubchannelPicker connectingPicker = mock(SubchannelPicker.class); |
781 | | - when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class))) |
782 | | - .thenReturn(PickResult.withNoResult("connecting", "pick_first: attempting to connect")); |
783 | | - |
784 | | - delayedTransport.reprocess(connectingPicker); |
| 781 | + FakeStreamTracer fakeTracer = new FakeStreamTracer(); |
| 782 | + ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer }; |
| 783 | + |
| 784 | + delayedTransport.reprocess(fakePicker( |
| 785 | + PickResult.withNoResult("connecting", "pick_first: attempting to connect"))); |
785 | 786 | delayedTransport.newStream(method, headers, callOptions, customTracers); |
786 | | - |
787 | | - InOrder inOrder = inOrder(mockTracer); |
788 | | - inOrder.verify(mockTracer).recordAttemptDelayStart( |
789 | | - "connecting", "pick_first: attempting to connect"); |
790 | | - |
791 | | - SubchannelPicker customDelayPicker = mock(SubchannelPicker.class); |
792 | | - when(customDelayPicker.pickSubchannel(any(PickSubchannelArgs.class))) |
793 | | - .thenReturn(PickResult.withNoResult("rls_lookup_pending", "RLS request pending.")); |
794 | | - |
795 | | - delayedTransport.reprocess(customDelayPicker); |
796 | | - |
797 | | - inOrder.verify(mockTracer).recordAttemptDelayEnd(); |
798 | | - inOrder.verify(mockTracer).recordAttemptDelayStart( |
799 | | - "rls_lookup_pending", "RLS request pending."); |
800 | | - |
| 787 | + |
| 788 | + assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes); |
| 789 | + assertEquals(Collections.singletonList("pick_first: attempting to connect"), |
| 790 | + fakeTracer.startedDelayReasons); |
| 791 | + |
| 792 | + delayedTransport.reprocess(fakePicker( |
| 793 | + PickResult.withNoResult("rls_lookup_pending", "RLS request pending."))); |
| 794 | + |
| 795 | + assertEquals(1, fakeTracer.delayEndedCount); |
| 796 | + assertEquals(Arrays.asList("connecting", "rls_lookup_pending"), |
| 797 | + fakeTracer.startedDelayTypes); |
| 798 | + assertEquals(Arrays.asList("pick_first: attempting to connect", "RLS request pending."), |
| 799 | + fakeTracer.startedDelayReasons); |
| 800 | + |
801 | 801 | delayedTransport.reprocess(mockPicker); |
802 | | - |
803 | | - inOrder.verify(mockTracer).recordAttemptDelayEnd(); |
| 802 | + |
| 803 | + assertEquals(2, fakeTracer.delayEndedCount); |
804 | 804 | } |
805 | 805 |
|
806 | 806 | @Test |
807 | 807 | public void streamDelayMetrics_cancelled() { |
808 | | - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); |
809 | | - ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer }; |
810 | | - |
811 | | - SubchannelPicker connectingPicker = mock(SubchannelPicker.class); |
812 | | - when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class))) |
813 | | - .thenReturn(PickResult.withNoResult("connecting", "pick_first: attempting to connect")); |
814 | | - |
815 | | - delayedTransport.reprocess(connectingPicker); |
| 808 | + FakeStreamTracer fakeTracer = new FakeStreamTracer(); |
| 809 | + ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer }; |
| 810 | + |
| 811 | + delayedTransport.reprocess(fakePicker( |
| 812 | + PickResult.withNoResult("connecting", "pick_first: attempting to connect"))); |
816 | 813 | ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers); |
817 | 814 | stream.start(streamListener); |
818 | | - |
819 | | - verify(mockTracer).recordAttemptDelayStart( |
820 | | - "connecting", "pick_first: attempting to connect"); |
821 | | - |
| 815 | + |
| 816 | + assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes); |
| 817 | + |
822 | 818 | stream.cancel(Status.CANCELLED); |
823 | | - |
824 | | - verify(mockTracer).recordAttemptDelayEnd(); |
| 819 | + |
| 820 | + assertEquals(1, fakeTracer.delayEndedCount); |
825 | 821 | } |
826 | 822 |
|
827 | 823 | @Test |
828 | 824 | public void streamDelayMetrics_shutdownNow() { |
829 | | - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); |
830 | | - ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer }; |
| 825 | + FakeStreamTracer fakeTracer = new FakeStreamTracer(); |
| 826 | + ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer }; |
831 | 827 |
|
832 | | - SubchannelPicker connectingPicker = mock(SubchannelPicker.class); |
833 | | - when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class))) |
834 | | - .thenReturn(PickResult.withNoResult("connecting", "pick_first: attempting to connect")); |
835 | | - |
836 | | - delayedTransport.reprocess(connectingPicker); |
| 828 | + delayedTransport.reprocess(fakePicker( |
| 829 | + PickResult.withNoResult("connecting", "pick_first: attempting to connect"))); |
837 | 830 | ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers); |
838 | 831 | stream.start(streamListener); |
839 | 832 |
|
840 | | - verify(mockTracer).recordAttemptDelayStart( |
841 | | - "connecting", "pick_first: attempting to connect"); |
| 833 | + assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes); |
842 | 834 |
|
843 | 835 | delayedTransport.shutdownNow(Status.UNAVAILABLE); |
844 | 836 |
|
845 | | - verify(mockTracer).recordAttemptDelayEnd(); |
| 837 | + assertEquals(1, fakeTracer.delayEndedCount); |
846 | 838 | } |
847 | 839 |
|
848 | 840 | @Test |
849 | 841 | public void streamDelayMetrics_cadenceReasonUpdate_doesNotStartNewTypeSegment() { |
850 | | - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); |
851 | | - ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer }; |
852 | | - |
853 | | - SubchannelPicker picker1 = mock(SubchannelPicker.class); |
854 | | - when(picker1.pickSubchannel(any(PickSubchannelArgs.class))) |
855 | | - .thenReturn(PickResult.withNoResult("connecting", "attempt 1")); |
| 842 | + FakeStreamTracer fakeTracer = new FakeStreamTracer(); |
| 843 | + ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer }; |
856 | 844 |
|
857 | | - delayedTransport.reprocess(picker1); |
| 845 | + delayedTransport.reprocess(fakePicker( |
| 846 | + PickResult.withNoResult("connecting", "attempt 1"))); |
858 | 847 | delayedTransport.newStream(method, headers, callOptions, customTracers); |
859 | 848 |
|
860 | | - verify(mockTracer, times(1)).recordAttemptDelayStart("connecting", "attempt 1"); |
| 849 | + assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes); |
| 850 | + assertEquals(Collections.singletonList("attempt 1"), fakeTracer.startedDelayReasons); |
861 | 851 |
|
862 | | - SubchannelPicker picker2 = mock(SubchannelPicker.class); |
863 | | - when(picker2.pickSubchannel(any(PickSubchannelArgs.class))) |
864 | | - .thenReturn(PickResult.withNoResult("connecting", "attempt 2")); |
865 | | - |
866 | | - delayedTransport.reprocess(picker2); |
| 852 | + delayedTransport.reprocess(fakePicker( |
| 853 | + PickResult.withNoResult("connecting", "attempt 2"))); |
867 | 854 |
|
868 | | - verify(mockTracer, times(1)).recordAttemptDelayStart("connecting", "attempt 1"); |
869 | | - verify(mockTracer).recordAttemptDelayReasonChanged("attempt 2"); |
870 | | - verify(mockTracer, never()).recordAttemptDelayEnd(); |
| 855 | + assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes); |
| 856 | + assertEquals(Collections.singletonList("attempt 1"), fakeTracer.startedDelayReasons); |
| 857 | + assertEquals(Collections.singletonList("attempt 2"), fakeTracer.changedDelayReasons); |
| 858 | + assertEquals(0, fakeTracer.delayEndedCount); |
871 | 859 | } |
872 | 860 |
|
873 | 861 | @Test |
874 | 862 | public void streamDelayMetrics_channelFallback_clientChannelInit() { |
875 | | - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); |
876 | | - ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer }; |
| 863 | + FakeStreamTracer fakeTracer = new FakeStreamTracer(); |
| 864 | + ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer }; |
877 | 865 |
|
878 | 866 | // No picker reprocessed yet (lastPicker == null) |
879 | 867 | delayedTransport.newStream(method, headers, callOptions, customTracers); |
880 | 868 |
|
881 | | - verify(mockTracer).recordAttemptDelayStart( |
882 | | - "connecting", "client channel: waiting for picker"); |
| 869 | + assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes); |
| 870 | + assertEquals(Collections.singletonList("client channel: waiting for picker"), |
| 871 | + fakeTracer.startedDelayReasons); |
883 | 872 | } |
884 | 873 |
|
885 | 874 | @Test |
886 | 875 | public void streamDelayMetrics_channelFallback_subchannelStateMismatch() { |
887 | | - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); |
888 | | - ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer }; |
| 876 | + FakeStreamTracer fakeTracer = new FakeStreamTracer(); |
| 877 | + ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer }; |
889 | 878 |
|
890 | 879 | io.grpc.LoadBalancer.Subchannel disconnectedSubchannel = |
891 | 880 | mock(io.grpc.LoadBalancer.Subchannel.class); |
892 | 881 | when(disconnectedSubchannel.getInternalSubchannel()) |
893 | 882 | .thenReturn(newTransportProvider(null)); |
894 | 883 |
|
895 | | - SubchannelPicker stalePicker = mock(SubchannelPicker.class); |
896 | | - when(stalePicker.pickSubchannel(any(PickSubchannelArgs.class))) |
897 | | - .thenReturn(PickResult.withSubchannel(disconnectedSubchannel)); |
898 | | - |
899 | | - delayedTransport.reprocess(stalePicker); |
| 884 | + delayedTransport.reprocess(fakePicker(PickResult.withSubchannel(disconnectedSubchannel))); |
900 | 885 | delayedTransport.newStream(method, headers, callOptions, customTracers); |
901 | 886 |
|
902 | | - verify(mockTracer).recordAttemptDelayStart( |
903 | | - "subchannel_state_mismatch", |
904 | | - "subchannel returned by LB picker has no connected subchannel"); |
| 887 | + assertEquals(Collections.singletonList("subchannel_state_mismatch"), |
| 888 | + fakeTracer.startedDelayTypes); |
| 889 | + assertEquals(Collections.singletonList( |
| 890 | + "subchannel returned by LB picker has no connected subchannel"), |
| 891 | + fakeTracer.startedDelayReasons); |
905 | 892 | } |
906 | 893 |
|
907 | 894 | @Test |
908 | 895 | public void streamDelayMetrics_channelFallback_waitForReadyFailed() { |
909 | | - ClientStreamTracer mockTracer = mock(ClientStreamTracer.class); |
910 | | - ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer }; |
911 | | - |
912 | | - SubchannelPicker failPicker = mock(SubchannelPicker.class); |
913 | | - when(failPicker.pickSubchannel(any(PickSubchannelArgs.class))) |
914 | | - .thenReturn(PickResult.withError(Status.UNAVAILABLE)); |
| 896 | + FakeStreamTracer fakeTracer = new FakeStreamTracer(); |
| 897 | + ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer }; |
915 | 898 |
|
916 | | - delayedTransport.reprocess(failPicker); |
| 899 | + delayedTransport.reprocess(fakePicker(PickResult.withError(Status.UNAVAILABLE))); |
917 | 900 | CallOptions wfrOptions = callOptions.withWaitForReady(); |
918 | 901 | delayedTransport.newStream(method, headers, wfrOptions, customTracers); |
919 | 902 |
|
920 | | - verify(mockTracer).recordAttemptDelayStart( |
921 | | - "picker_failing_with_wait_for_ready", |
922 | | - "wait_for_ready RPC failed with status: " + Status.UNAVAILABLE); |
| 903 | + assertEquals(Collections.singletonList("picker_failing_with_wait_for_ready"), |
| 904 | + fakeTracer.startedDelayTypes); |
| 905 | + assertEquals(Collections.singletonList( |
| 906 | + "wait_for_ready RPC failed with status: " + Status.UNAVAILABLE), |
| 907 | + fakeTracer.startedDelayReasons); |
| 908 | + } |
| 909 | + |
| 910 | + private static final class FakeStreamTracer extends ClientStreamTracer { |
| 911 | + final List<String> startedDelayTypes = new ArrayList<>(); |
| 912 | + final List<String> startedDelayReasons = new ArrayList<>(); |
| 913 | + final List<String> changedDelayReasons = new ArrayList<>(); |
| 914 | + int delayEndedCount = 0; |
| 915 | + |
| 916 | + @Override |
| 917 | + public void recordAttemptDelayStart(String delayType, String delayReason) { |
| 918 | + startedDelayTypes.add(delayType); |
| 919 | + startedDelayReasons.add(delayReason); |
| 920 | + } |
| 921 | + |
| 922 | + @Override |
| 923 | + public void recordAttemptDelayReasonChanged(String delayReason) { |
| 924 | + changedDelayReasons.add(delayReason); |
| 925 | + } |
| 926 | + |
| 927 | + @Override |
| 928 | + public void recordAttemptDelayEnd() { |
| 929 | + delayEndedCount++; |
| 930 | + } |
| 931 | + } |
| 932 | + |
| 933 | + private static SubchannelPicker fakePicker(final PickResult result) { |
| 934 | + return new SubchannelPicker() { |
| 935 | + @Override |
| 936 | + public PickResult pickSubchannel(PickSubchannelArgs args) { |
| 937 | + return result; |
| 938 | + } |
| 939 | + }; |
923 | 940 | } |
924 | 941 |
|
925 | 942 | private static TransportProvider newTransportProvider(final ClientTransport transport) { |
|
0 commit comments