Commit a9243f5
Eric Bower
·
2025-12-26 11:09:38 -0500 EST
parent f60c265
chore(pipe): more e2e tests for blocking and non-blocking behavior
1 files changed,
+352,
-0
+352,
-0
| ... | ... | @@ -659,3 +659,355 @@ func TestAccessControl_AllowedUserViaFullPath(t *testing.T) { | |
| 659 | 659 | t.Errorf("alice should receive bob's message on shared topic, got: %q", string(aliceReceived[:n])) | |
| 660 | 660 | } | |
| 661 | 661 | } | |
| 662 | + | ||
| 663 | + | func TestPubSub_BlockingWaitsForSubscriber(t *testing.T) { | |
| 664 | + | server := NewTestSSHServer(t) | |
| 665 | + | defer server.Shutdown() | |
| 666 | + | ||
| 667 | + | user := GenerateUser("alice") | |
| 668 | + | RegisterUserWithServer(server, user) | |
| 669 | + | ||
| 670 | + | pubClient, err := user.NewClient() | |
| 671 | + | if err != nil { | |
| 672 | + | t.Fatalf("failed to connect publisher: %v", err) | |
| 673 | + | } | |
| 674 | + | defer func() { _ = pubClient.Close() }() | |
| 675 | + | ||
| 676 | + | subClient, err := user.NewClient() | |
| 677 | + | if err != nil { | |
| 678 | + | t.Fatalf("failed to connect subscriber: %v", err) | |
| 679 | + | } | |
| 680 | + | defer func() { _ = subClient.Close() }() | |
| 681 | + | ||
| 682 | + | pubSession, err := pubClient.NewSession() | |
| 683 | + | if err != nil { | |
| 684 | + | t.Fatalf("failed to create pub session: %v", err) | |
| 685 | + | } | |
| 686 | + | defer func() { _ = pubSession.Close() }() | |
| 687 | + | ||
| 688 | + | pubStdin, err := pubSession.StdinPipe() | |
| 689 | + | if err != nil { | |
| 690 | + | t.Fatalf("failed to get pub stdin: %v", err) | |
| 691 | + | } | |
| 692 | + | ||
| 693 | + | pubStdout, err := pubSession.StdoutPipe() | |
| 694 | + | if err != nil { | |
| 695 | + | t.Fatalf("failed to get pub stdout: %v", err) | |
| 696 | + | } | |
| 697 | + | ||
| 698 | + | // Start publisher with blocking enabled (default -b=true) | |
| 699 | + | // Publisher should wait for subscriber | |
| 700 | + | if err := pubSession.Start("pub blockingtopic"); err != nil { | |
| 701 | + | t.Fatalf("failed to start pub: %v", err) | |
| 702 | + | } | |
| 703 | + | ||
| 704 | + | // Read initial output - should indicate waiting for subscribers | |
| 705 | + | initialOutput := make([]byte, 500) | |
| 706 | + | n, err := pubStdout.Read(initialOutput) | |
| 707 | + | if err != nil && err != io.EOF { | |
| 708 | + | t.Fatalf("failed to read initial output: %v", err) | |
| 709 | + | } | |
| 710 | + | output := string(initialOutput[:n]) | |
| 711 | + | if !strings.Contains(output, "waiting") { | |
| 712 | + | t.Errorf("expected 'waiting' message for blocking pub, got: %q", output) | |
| 713 | + | } | |
| 714 | + | ||
| 715 | + | // Now start subscriber - this should unblock the publisher | |
| 716 | + | subSession, err := subClient.NewSession() | |
| 717 | + | if err != nil { | |
| 718 | + | t.Fatalf("failed to create sub session: %v", err) | |
| 719 | + | } | |
| 720 | + | defer func() { _ = subSession.Close() }() | |
| 721 | + | ||
| 722 | + | subStdout, err := subSession.StdoutPipe() | |
| 723 | + | if err != nil { | |
| 724 | + | t.Fatalf("failed to get sub stdout: %v", err) | |
| 725 | + | } | |
| 726 | + | ||
| 727 | + | if err := subSession.Start("sub blockingtopic -c"); err != nil { | |
| 728 | + | t.Fatalf("failed to start sub: %v", err) | |
| 729 | + | } | |
| 730 | + | ||
| 731 | + | time.Sleep(100 * time.Millisecond) | |
| 732 | + | ||
| 733 | + | // Now send the message | |
| 734 | + | testMessage := "blocking message" | |
| 735 | + | _, err = pubStdin.Write([]byte(testMessage)) | |
| 736 | + | if err != nil { | |
| 737 | + | t.Fatalf("failed to write message: %v", err) | |
| 738 | + | } | |
| 739 | + | _ = pubStdin.Close() | |
| 740 | + | ||
| 741 | + | // Subscriber should receive the message | |
| 742 | + | received := make([]byte, 100) | |
| 743 | + | n, err = subStdout.Read(received) | |
| 744 | + | if err != nil && err != io.EOF { | |
| 745 | + | t.Logf("read error: %v", err) | |
| 746 | + | } | |
| 747 | + | ||
| 748 | + | if !strings.Contains(string(received[:n]), testMessage) { | |
| 749 | + | t.Errorf("subscriber did not receive blocking message, got: %q, want: %q", string(received[:n]), testMessage) | |
| 750 | + | } | |
| 751 | + | } | |
| 752 | + | ||
| 753 | + | func TestPubSub_NonBlockingDoesNotWait(t *testing.T) { | |
| 754 | + | server := NewTestSSHServer(t) | |
| 755 | + | defer server.Shutdown() | |
| 756 | + | ||
| 757 | + | user := GenerateUser("alice") | |
| 758 | + | RegisterUserWithServer(server, user) | |
| 759 | + | ||
| 760 | + | client, err := user.NewClient() | |
| 761 | + | if err != nil { | |
| 762 | + | t.Fatalf("failed to connect: %v", err) | |
| 763 | + | } | |
| 764 | + | defer func() { _ = client.Close() }() | |
| 765 | + | ||
| 766 | + | // Publish with -b=false (non-blocking) and no subscriber | |
| 767 | + | // Should complete immediately without waiting | |
| 768 | + | done := make(chan struct{}) | |
| 769 | + | var output string | |
| 770 | + | var cmdErr error | |
| 771 | + | ||
| 772 | + | go func() { | |
| 773 | + | output, cmdErr = user.RunCommandWithStdin(client, "pub nonblockingtopic -b=false -c", "non-blocking message") | |
| 774 | + | close(done) | |
| 775 | + | }() | |
| 776 | + | ||
| 777 | + | select { | |
| 778 | + | case <-done: | |
| 779 | + | // Command completed - this is expected for non-blocking | |
| 780 | + | if cmdErr != nil { | |
| 781 | + | t.Logf("non-blocking pub completed with: %v", cmdErr) | |
| 782 | + | } | |
| 783 | + | t.Logf("non-blocking pub output: %q", output) | |
| 784 | + | case <-time.After(2 * time.Second): | |
| 785 | + | t.Errorf("non-blocking pub should complete immediately, but it blocked") | |
| 786 | + | } | |
| 787 | + | } | |
| 788 | + | ||
| 789 | + | func TestPubSub_BlockingTimeout(t *testing.T) { | |
| 790 | + | server := NewTestSSHServer(t) | |
| 791 | + | defer server.Shutdown() | |
| 792 | + | ||
| 793 | + | user := GenerateUser("alice") | |
| 794 | + | RegisterUserWithServer(server, user) | |
| 795 | + | ||
| 796 | + | client, err := user.NewClient() | |
| 797 | + | if err != nil { | |
| 798 | + | t.Fatalf("failed to connect: %v", err) | |
| 799 | + | } | |
| 800 | + | defer func() { _ = client.Close() }() | |
| 801 | + | ||
| 802 | + | // Publish with blocking and short timeout, no subscriber | |
| 803 | + | // Should timeout after the specified duration | |
| 804 | + | done := make(chan struct{}) | |
| 805 | + | var output string | |
| 806 | + | ||
| 807 | + | go func() { | |
| 808 | + | output, _ = user.RunCommandWithStdin(client, "pub timeouttopic -b=true -t=500ms", "timeout message") | |
| 809 | + | close(done) | |
| 810 | + | }() | |
| 811 | + | ||
| 812 | + | select { | |
| 813 | + | case <-done: | |
| 814 | + | // Command completed due to timeout | |
| 815 | + | if !strings.Contains(output, "timeout") && !strings.Contains(output, "waiting") { | |
| 816 | + | t.Logf("blocking pub with timeout output: %q", output) | |
| 817 | + | } | |
| 818 | + | case <-time.After(3 * time.Second): | |
| 819 | + | t.Errorf("blocking pub with timeout should have timed out after 500ms") | |
| 820 | + | } | |
| 821 | + | } | |
| 822 | + | ||
| 823 | + | func TestSub_WaitsForPublisher(t *testing.T) { | |
| 824 | + | server := NewTestSSHServer(t) | |
| 825 | + | defer server.Shutdown() | |
| 826 | + | ||
| 827 | + | user := GenerateUser("alice") | |
| 828 | + | RegisterUserWithServer(server, user) | |
| 829 | + | ||
| 830 | + | subClient, err := user.NewClient() | |
| 831 | + | if err != nil { | |
| 832 | + | t.Fatalf("failed to connect subscriber: %v", err) | |
| 833 | + | } | |
| 834 | + | defer func() { _ = subClient.Close() }() | |
| 835 | + | ||
| 836 | + | pubClient, err := user.NewClient() | |
| 837 | + | if err != nil { | |
| 838 | + | t.Fatalf("failed to connect publisher: %v", err) | |
| 839 | + | } | |
| 840 | + | defer func() { _ = pubClient.Close() }() | |
| 841 | + | ||
| 842 | + | // Start subscriber first - it should wait for publisher | |
| 843 | + | subSession, err := subClient.NewSession() | |
| 844 | + | if err != nil { | |
| 845 | + | t.Fatalf("failed to create sub session: %v", err) | |
| 846 | + | } | |
| 847 | + | defer func() { _ = subSession.Close() }() | |
| 848 | + | ||
| 849 | + | subStdout, err := subSession.StdoutPipe() | |
| 850 | + | if err != nil { | |
| 851 | + | t.Fatalf("failed to get sub stdout: %v", err) | |
| 852 | + | } | |
| 853 | + | ||
| 854 | + | if err := subSession.Start("sub waitfortopic -c"); err != nil { | |
| 855 | + | t.Fatalf("failed to start sub: %v", err) | |
| 856 | + | } | |
| 857 | + | ||
| 858 | + | // Subscriber is now waiting - give it a moment | |
| 859 | + | time.Sleep(100 * time.Millisecond) | |
| 860 | + | ||
| 861 | + | // Now publish - subscriber should receive it | |
| 862 | + | testMessage := "delayed publish" | |
| 863 | + | _, err = user.RunCommandWithStdin(pubClient, "pub waitfortopic -c", testMessage) | |
| 864 | + | if err != nil { | |
| 865 | + | t.Logf("pub completed: %v", err) | |
| 866 | + | } | |
| 867 | + | ||
| 868 | + | received := make([]byte, 100) | |
| 869 | + | n, err := subStdout.Read(received) | |
| 870 | + | if err != nil && err != io.EOF { | |
| 871 | + | t.Logf("read error: %v", err) | |
| 872 | + | } | |
| 873 | + | ||
| 874 | + | if !strings.Contains(string(received[:n]), testMessage) { | |
| 875 | + | t.Errorf("subscriber waiting for publisher did not receive message, got: %q, want: %q", string(received[:n]), testMessage) | |
| 876 | + | } | |
| 877 | + | } | |
| 878 | + | ||
| 879 | + | func TestSub_KeepAliveReceivesMultipleMessages(t *testing.T) { | |
| 880 | + | server := NewTestSSHServer(t) | |
| 881 | + | defer server.Shutdown() | |
| 882 | + | ||
| 883 | + | user := GenerateUser("alice") | |
| 884 | + | RegisterUserWithServer(server, user) | |
| 885 | + | ||
| 886 | + | subClient, err := user.NewClient() | |
| 887 | + | if err != nil { | |
| 888 | + | t.Fatalf("failed to connect subscriber: %v", err) | |
| 889 | + | } | |
| 890 | + | defer func() { _ = subClient.Close() }() | |
| 891 | + | ||
| 892 | + | pubClient1, err := user.NewClient() | |
| 893 | + | if err != nil { | |
| 894 | + | t.Fatalf("failed to connect publisher 1: %v", err) | |
| 895 | + | } | |
| 896 | + | defer func() { _ = pubClient1.Close() }() | |
| 897 | + | ||
| 898 | + | pubClient2, err := user.NewClient() | |
| 899 | + | if err != nil { | |
| 900 | + | t.Fatalf("failed to connect publisher 2: %v", err) | |
| 901 | + | } | |
| 902 | + | defer func() { _ = pubClient2.Close() }() | |
| 903 | + | ||
| 904 | + | // Start subscriber with keepAlive (-k) flag | |
| 905 | + | subSession, err := subClient.NewSession() | |
| 906 | + | if err != nil { | |
| 907 | + | t.Fatalf("failed to create sub session: %v", err) | |
| 908 | + | } | |
| 909 | + | defer func() { _ = subSession.Close() }() | |
| 910 | + | ||
| 911 | + | subStdout, err := subSession.StdoutPipe() | |
| 912 | + | if err != nil { | |
| 913 | + | t.Fatalf("failed to get sub stdout: %v", err) | |
| 914 | + | } | |
| 915 | + | ||
| 916 | + | if err := subSession.Start("sub keepalivetopic -k -c"); err != nil { | |
| 917 | + | t.Fatalf("failed to start sub: %v", err) | |
| 918 | + | } | |
| 919 | + | ||
| 920 | + | time.Sleep(100 * time.Millisecond) | |
| 921 | + | ||
| 922 | + | // Send first message | |
| 923 | + | msg1 := "first message\n" | |
| 924 | + | _, err = user.RunCommandWithStdin(pubClient1, "pub keepalivetopic -c", msg1) | |
| 925 | + | if err != nil { | |
| 926 | + | t.Logf("pub 1 completed: %v", err) | |
| 927 | + | } | |
| 928 | + | ||
| 929 | + | received1 := make([]byte, 100) | |
| 930 | + | n1, _ := subStdout.Read(received1) | |
| 931 | + | if !strings.Contains(string(received1[:n1]), "first message") { | |
| 932 | + | t.Errorf("subscriber did not receive first message, got: %q", string(received1[:n1])) | |
| 933 | + | } | |
| 934 | + | ||
| 935 | + | // Send second message - subscriber with keepAlive should still receive it | |
| 936 | + | msg2 := "second message\n" | |
| 937 | + | _, err = user.RunCommandWithStdin(pubClient2, "pub keepalivetopic -c", msg2) | |
| 938 | + | if err != nil { | |
| 939 | + | t.Logf("pub 2 completed: %v", err) | |
| 940 | + | } | |
| 941 | + | ||
| 942 | + | received2 := make([]byte, 100) | |
| 943 | + | n2, _ := subStdout.Read(received2) | |
| 944 | + | if !strings.Contains(string(received2[:n2]), "second message") { | |
| 945 | + | t.Errorf("subscriber with keepAlive did not receive second message, got: %q", string(received2[:n2])) | |
| 946 | + | } | |
| 947 | + | } | |
| 948 | + | ||
| 949 | + | func TestSub_WithoutKeepAliveExitsAfterPublisher(t *testing.T) { | |
| 950 | + | server := NewTestSSHServer(t) | |
| 951 | + | defer server.Shutdown() | |
| 952 | + | ||
| 953 | + | user := GenerateUser("alice") | |
| 954 | + | RegisterUserWithServer(server, user) | |
| 955 | + | ||
| 956 | + | subClient, err := user.NewClient() | |
| 957 | + | if err != nil { | |
| 958 | + | t.Fatalf("failed to connect subscriber: %v", err) | |
| 959 | + | } | |
| 960 | + | defer func() { _ = subClient.Close() }() | |
| 961 | + | ||
| 962 | + | pubClient, err := user.NewClient() | |
| 963 | + | if err != nil { | |
| 964 | + | t.Fatalf("failed to connect publisher: %v", err) | |
| 965 | + | } | |
| 966 | + | defer func() { _ = pubClient.Close() }() | |
| 967 | + | ||
| 968 | + | // Start subscriber without keepAlive | |
| 969 | + | subSession, err := subClient.NewSession() | |
| 970 | + | if err != nil { | |
| 971 | + | t.Fatalf("failed to create sub session: %v", err) | |
| 972 | + | } | |
| 973 | + | ||
| 974 | + | subStdout, err := subSession.StdoutPipe() | |
| 975 | + | if err != nil { | |
| 976 | + | t.Fatalf("failed to get sub stdout: %v", err) | |
| 977 | + | } | |
| 978 | + | ||
| 979 | + | if err := subSession.Start("sub exitaftertopic -c"); err != nil { | |
| 980 | + | t.Fatalf("failed to start sub: %v", err) | |
| 981 | + | } | |
| 982 | + | ||
| 983 | + | time.Sleep(100 * time.Millisecond) | |
| 984 | + | ||
| 985 | + | // Publish a message | |
| 986 | + | testMessage := "single message" | |
| 987 | + | _, err = user.RunCommandWithStdin(pubClient, "pub exitaftertopic -c", testMessage) | |
| 988 | + | if err != nil { | |
| 989 | + | t.Logf("pub completed: %v", err) | |
| 990 | + | } | |
| 991 | + | ||
| 992 | + | // Read the message | |
| 993 | + | received := make([]byte, 100) | |
| 994 | + | n, _ := subStdout.Read(received) | |
| 995 | + | if !strings.Contains(string(received[:n]), testMessage) { | |
| 996 | + | t.Errorf("subscriber did not receive message, got: %q", string(received[:n])) | |
| 997 | + | } | |
| 998 | + | ||
| 999 | + | // Subscriber session should exit after publisher disconnects | |
| 1000 | + | done := make(chan error) | |
| 1001 | + | go func() { | |
| 1002 | + | done <- subSession.Wait() | |
| 1003 | + | }() | |
| 1004 | + | ||
| 1005 | + | select { | |
| 1006 | + | case err := <-done: | |
| 1007 | + | // Session ended as expected | |
| 1008 | + | t.Logf("subscriber session ended: %v", err) | |
| 1009 | + | case <-time.After(2 * time.Second): | |
| 1010 | + | t.Errorf("subscriber without keepAlive should have exited after publisher disconnected") | |
| 1011 | + | _ = subSession.Close() | |
| 1012 | + | } | |
| 1013 | + | } |