File
Blob: overlay/attachment_test.go
| 1 | package overlay |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "testing" |
| 6 | "time" |
| 7 | |
| 8 | "go.miragespace.co/specter/spec/protocol" |
| 9 | "go.miragespace.co/specter/spec/rpc" |
| 10 | "go.miragespace.co/specter/spec/transport" |
| 11 | |
| 12 | "github.com/stretchr/testify/require" |
| 13 | ) |
| 14 | |
| 15 | func TestAttachmentBoundToPhysicalConn(t *testing.T) { |
| 16 | ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) |
| 17 | defer cancel() |
| 18 | dial := newReaperTestListener(t, ctx) |
| 19 | client, server := dial() |
| 20 | stream, err := openStream(client, &protocol.Stream{Type: protocol.Stream_RPC}) |
| 21 | require.NoError(t, err) |
| 22 | defer stream.Close() |
| 23 | incoming, err := server.AcceptStream(ctx) |
| 24 | require.NoError(t, err) |
| 25 | conn := WrapQuicConnection(incoming, server) |
| 26 | defer conn.Close() |
| 27 | var header protocol.Stream |
| 28 | require.NoError(t, rpc.BoundedReceive(conn, &header, 1024)) |
| 29 | pc := conn.(transport.PhysicalConnProvider).PhysicalConn() |
| 30 | direct, err := pc.OpenStream(protocol.Stream_DIRECT) |
| 31 | require.NoError(t, err) |
| 32 | defer direct.Close() |
| 33 | received, err := client.AcceptStream(ctx) |
| 34 | require.NoError(t, err) |
| 35 | require.NoError(t, rpc.BoundedReceive(received, &header, 1024)) |
| 36 | require.Equal(t, protocol.Stream_DIRECT, header.Type) |
| 37 | require.NoError(t, client.CloseWithError(0, "test disconnect")) |
| 38 | select { |
| 39 | case <-pc.Done(): |
| 40 | case <-ctx.Done(): |
| 41 | t.Fatal("attachment did not close") |
| 42 | } |
| 43 | _, replacement := dial() |
| 44 | require.NotEqual(t, pc, transport.PhysicalConn(attachment{replacement})) |
| 45 | _, err = pc.OpenStream(protocol.Stream_DIRECT) |
| 46 | require.Error(t, err) |
| 47 | require.Error(t, pc.Err()) |
| 48 | } |