Skip to content

fix: don't discard final message on server-close - #262

Open
mkulke wants to merge 1 commit into
containerd:mainfrom
mkulke:mkulke/fix-missing-response-on-close
Open

mkulke wants to merge 1 commit into
containerd:mainfrom
mkulke:mkulke/fix-missing-response-on-close

Conversation

@mkulke

@mkulke mkulke commented Oct 5, 2026

Copy link
Copy Markdown

A ttrpc server can send valid response and then close the connection. In the client there is a race between the handling the cancellation and receiving the final response. In the case of pending msg the client should retrieve this first instead of returning ErrClosed.

(this surfaced in kata-container runtime/agent comm that terminates a session with DestroySandbox. In some cases the client will not receive an acknowledgment of proper termination but an error instead).

Reproduction

on main, we run 64 RPCs, a certain amount of replies will be discarded (~23 of 64 in local tests).

$ cat <<EOF> response_close_race_test.go
package ttrpc

import (
	"context"
	"errors"
	"net"
	"testing"
	"time"

	"google.golang.org/protobuf/proto"
	"google.golang.org/protobuf/types/known/wrapperspb"
)

type responseBeforeCloseConn struct {
	net.Conn
	beforeWriteReturns func()
}

func (c *responseBeforeCloseConn) Write(p []byte) (int, error) {
	n, err := c.Conn.Write(p)
	if err == nil {
		c.beforeWriteReturns()
	}
	return n, err
}

func TestResponseBeforePeerClose(t *testing.T) {
	const attempts = 64
	closedReplies := 0
	for attempt := 0; attempt < attempts; attempt++ {
		clientConn, serverConn := net.Pipe()
		conn := &responseBeforeCloseConn{Conn: clientConn}
		client := NewClient(conn)
		// Hold dispatch before its select until the reply is buffered and EOF
		// has closed the client, making both competing outcomes ready.
		conn.beforeWriteReturns = func() { <-client.userCloseWaitCh }

		serverResult := make(chan error, 1)
		go func() {
			defer serverConn.Close()
			channel := newChannel(serverConn)
			header, request, err := channel.recv()
			if err != nil {
				serverResult <- err
				return
			}
			channel.putmbuf(request)
			payload, err := proto.Marshal(wrapperspb.String("response before EOF"))
			if err != nil {
				serverResult <- err
				return
			}
			response, err := proto.Marshal(&Response{Payload: payload})
			if err != nil {
				serverResult <- err
				return
			}
			serverResult <- channel.send(header.StreamID, messageTypeResponse, 0, response)
		}()

		ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
		timeout := time.AfterFunc(2*time.Second, func() { _ = client.Close() })
		var response wrapperspb.StringValue
		err := client.Call(ctx, "test.Service", "ReplyAndClose", wrapperspb.String("request"), &response)
		timeout.Stop()
		cancel()
		_ = client.Close()
		if serverErr := <-serverResult; serverErr != nil {
			t.Fatalf("server failed to send the response: %v", serverErr)
		}
		switch {
		case errors.Is(err, ErrClosed):
			closedReplies++
		case err != nil:
			t.Fatalf("unexpected RPC error: %v", err)
		case response.Value != "response before EOF":
			t.Fatalf("unexpected response: %q", response.Value)
		}
	}
	if closedReplies != 0 {
		t.Fatalf("discarded %d/%d responses already received before peer EOF", closedReplies, attempts)
	}
}
EOF
$ go test -run '^TestResponseBeforePeerClose$' -count=1 -timeout=20s .
--- FAIL: TestResponseBeforePeerClose (0.00s)
    response_close_race_test.go:81: discarded 22/64 responses already received before peer EOF
FAIL
FAIL    github.com/containerd/ttrpc     0.008s
FAIL

w/ the fix applied:

$ git checkout mkulke/fix-missing-response-on-close
$ go test -run '^TestResponseBeforePeerClose$' -count=1 -timeout=20s .
ok      github.com/containerd/ttrpc     0.008s
$ go test -race -run '^TestResponseBeforePeerClose$' -count=1 -timeout=20s .
ok      github.com/containerd/ttrpc     1.031s

A ttrpc server can send valid response and then close the connection.
In the client there is a race between the handling the cancellation
and receiving the final response. In the case of pending msg the client
should retrieve this first instead of returning `ErrClosed`.

(this surfaced in kata-container runtime/agent comm that terminates
a session with `DestroySandbox`. In some cases the client will not
receive an acknowledgment of proper termination but an error instead)

Signed-off-by: Magnus Kulke <magnuskulke@microsoft.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant