Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
66 changes: 54 additions & 12 deletions internal/lambda-managed-instances/aws-lambda-rie/test/rie_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package test
import (
"io"
"net/http"
"net/http/httptrace"
"os"
"strings"
"sync"
Expand Down Expand Up @@ -210,12 +211,31 @@ func TestRie_InvokeWaitingForInitError(t *testing.T) {
server, rieHandler, _, err := internal.Run(supv, args, mockFileUtil, sigCh)
require.NoError(t, err)

var wg sync.WaitGroup
wg.Add(1)
req, err := http.NewRequest(http.MethodPost, "http://"+server.Addr.String()+"/2015-03-31/functions/function/invocations", strings.NewReader("{}"))
require.NoError(t, err)
req.Header.Set("Content-Type", "application/json")

// Closed once the invoke request has been written out, so that the test never has to
// guess how long the invoke goroutine takes to get scheduled.
requestSent := make(chan struct{})
var requestSentOnce sync.Once
req = req.WithContext(httptrace.WithClientTrace(req.Context(), &httptrace.ClientTrace{
WroteRequest: func(httptrace.WroteRequestInfo) {
requestSentOnce.Do(func() { close(requestSent) })
},
}))

invokeDone := make(chan struct{})
go func() {
resp, err := http.Post("http://"+server.Addr.String()+"/2015-03-31/functions/function/invocations", "application/json", strings.NewReader("{}"))
require.NoError(t, err)
defer func() { require.NoError(t, resp.Body.Close()) }()
// This is not the test goroutine, so it must only use assert: a failed require would
// end the goroutine via runtime.Goexit and leave the test blocked on invokeDone.
defer close(invokeDone)

resp, err := http.DefaultClient.Do(req)
if !assert.NoError(t, err) {
return
}
defer func() { assert.NoError(t, resp.Body.Close()) }()

assert.Equal(t, http.StatusOK, resp.StatusCode)

Expand All @@ -228,22 +248,37 @@ func TestRie_InvokeWaitingForInitError(t *testing.T) {
}

body, err := io.ReadAll(resp.Body)
require.NoError(t, err)
if !assert.NoError(t, err) {
return
}
assert.JSONEq(t, `{"errorType":"Runtime.ExitError"}`, string(body))

wg.Done()
}()

time.Sleep(200 * time.Millisecond)
// The invoke request is what triggers initialization here, so waiting for it to be sent
// before touching the handler keeps the invoke in flight while init runs and fails.
waitForClose(t, requestSent, "invoke request was never sent")
waitForClose(t, server.Done(), "server did not shut down after the init error")

initErr := rieHandler.Init()
require.Error(t, initErr)

<-server.Done()
serverErr := server.Err()
assert.Error(t, serverErr)
assert.Equal(t, initErr, serverErr)

wg.Wait()
waitForClose(t, invokeDone, "invoke request did not complete")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[CONCURRENCY] The invoke goroutine is only joined on the happy path. Every t.Fatal/require failure before this line exits the test goroutine through runtime.Goexit while the goroutine is still inside http.DefaultClient.Do:

  • waitForClose(t, requestSent, ...) failing means the request is still in flight by definition
  • waitForClose(t, server.Done(), ...) and require.Error(t, initErr) can both fail with the request outstanding

When Do later returns, the goroutine calls assert.NoError(t, err) on a *testing.T that is already done, and the testing package panics with Log in goroutine after TestRie_InvokeWaitingForInitError has completed — taking the whole test binary down, which is the same collateral damage this PR is removing.

Registering the join as a cleanup makes it run on every exit path, including Goexit:

invokeDone := make(chan struct{})
t.Cleanup(func() {
select {
case <-invokeDone:
case <-time.After(30  time.Second):
t.Error("invoke request did not complete")
}
})

go func() {
defer close(invokeDone)
...
}()

and then the trailing waitForClose(t, invokeDone, ...) can be dropped. t.Error is used rather than t.Fatal because the cleanup runs outside the test body.

}

// waitForClose blocks until ch is closed and fails the test instead of letting the whole
// package hit the go test timeout when it never is.
func waitForClose(t *testing.T, ch <-chan struct{}, msg string) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[CONCURRENCY] TestRie_SigtermDuringInvoke (lines 118-147 of this file) still has the exact deadlock this helper was added to remove, and it is the sibling of the test being fixed:

go func() {
resp, err := http.Post(...)
require.NoError(t, err)   // off the test goroutine
...
wg.Done()                 // not deferred
}()

time.Sleep(200  time.Millisecond)
sigCh <- syscall.SIGTERM

If the goroutine is not scheduled within 200ms, SIGTERM shuts the server down first and http.Post hits a closed listener. require.NoError then ends the goroutine via runtime.Goexit, the un-deferred wg.Done() is skipped, and wg.Wait() blocks until the go test timeout panics the whole package — the failure mode described in the PR description.

The ordering there is load-bearing (the test wants the invoke to land after SIGTERM but while the listener still accepts), so it cannot take the requestSent trick verbatim. At minimum, switching the goroutine to assert and deferring the WaitGroup release converts the hang into a reported failure:

go func() {
defer wg.Done()

resp, err := http.Post(...)
if !assert.NoError(t, err) {
return
}
defer func() { assert.NoError(t, resp.Body.Close()) }()
...
}()

t.Helper()

select {
case <-ch:
case <-time.After(30 * time.Second):
t.Fatal(msg)
}
}

func TestRie_InvokeFatalError(t *testing.T) {
Expand Down Expand Up @@ -475,11 +510,18 @@ func TestRIE_TelemetryAPI(t *testing.T) {
},
}

const expectedLogLines = 6

for _, mock := range []*functional.InMemoryEventsApi{httpEventsApi, tcpEventsApi} {
// Log lines are relayed to subscribers asynchronously and are not guaranteed to
// have all been delivered by the time the server reports itself shut down.
require.Eventually(t, func() bool { return len(mock.LogLines()) == expectedLogLines },

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[GENERAL] Polling for the log lines fixes the flake, but folding the count assertion into the Eventually condition with == loses the upper bound that assert.Len(t, mock.LogLines(), 6) provided, and makes the outcome timing-dependent if a 7th line ever shows up: whether the poll happens to observe the transient 6 decides between pass and a 10-second timeout with a message that never reports the count actually seen.

Poll for the lower bound and keep the exact assertion afterwards, so the wait is unchanged but an unexpected extra line fails deterministically and the failure prints the real length:

require.Eventually(t, func() bool { return len(mock.LogLines()) >= expectedLogLines },
10time.Second, 10*time.Millisecond,
"expected at least %d log lines to be delivered over the Telemetry API", expectedLogLines)

mock.CheckSimpleInitExpectations(initStartTime, initFinishTime, expectedInitEvents, initPayload)
mock.CheckSimpleExtensionExpectations(expectedExtensionEvents)
mock.CheckSimpleInvokeExpectations(invokeStartTime, invokeFinishTime, invokeID, expectedInvokeEvents, initPayload)
assert.Len(t, mock.LogLines(), expectedLogLines)

10*time.Second, 10*time.Millisecond,
"expected %d log lines to be delivered over the Telemetry API", expectedLogLines)

mock.CheckSimpleInitExpectations(initStartTime, initFinishTime, expectedInitEvents, initPayload)
mock.CheckSimpleExtensionExpectations(expectedExtensionEvents)
mock.CheckSimpleInvokeExpectations(invokeStartTime, invokeFinishTime, invokeID, expectedInvokeEvents, initPayload)
assert.Len(t, mock.LogLines(), 6)
}
})
}
Expand Down