Skip to content
Merged
Show file tree
Hide file tree
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
4 changes: 4 additions & 0 deletions doc/rfc/stovepipe/steps/buildsignal.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,10 @@ For a delivery carrying build id `B`:
8. Else PublishAfter(B -> buildsignal, delayMs), partitioned by build id:
- delayMs = pollDelay(status): shorter while running, longer while accepted.
- a fresh message (retry_count resets to 0), not a nack — polling is not failure.
- the message id must be unique per tick. The queue dedups on (topic, partition_key, id)
and the delivery being processed is still un-acked, so its row is present: reusing B as
the message id makes every re-poll collide with the message that scheduled it and be
silently discarded, ending the poll loop after one tick.
- ack.
```

Expand Down
41 changes: 38 additions & 3 deletions stovepipe/controller/buildsignal/buildsignal.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ package buildsignal
import (
"context"
"fmt"
"strconv"
"strings"

"github.com/uber-go/tally"
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
Expand Down Expand Up @@ -151,7 +153,7 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
}

delayMs := pollDelay(effective)
if err := c.publishBuildSignal(ctx, build.ID, delayMs); err != nil {
if err := c.publishBuildSignal(ctx, build.ID, msg.ID, delayMs); err != nil {
return errs.NewRetryableError(fmt.Errorf("failed to reschedule poll for build %s: %w", build.ID, err))
}
c.logger.Debugw("rescheduled build status poll",
Expand Down Expand Up @@ -218,15 +220,48 @@ func (c *Controller) publishRecord(ctx context.Context, buildID, requestID strin
return c.publish(ctx, stovepipemq.TopicKeyRecord, msg, 0)
}

// pollIDInfix separates a build id from its poll generation in a re-poll
// message id: "<buildID>/poll/<generation>".
const pollIDInfix = "/poll/"

// nextPollMessageID returns the message id for the next poll of buildID, one
// generation past the delivery that scheduled it. currentMsgID is the id of the
// message being processed — either the initial publish from build (no
// generation, so the next is 1) or a previous re-poll.
//
// The generation has to advance because the queue dedups on
// (topic, partition_key, id) and the delivery being processed has not been acked
// yet, so its row is still present: a re-poll reusing the current id would
// collide with the message that scheduled it and be silently discarded, ending
// the poll loop after one tick.
//
// It advances deterministically rather than randomly so the id stays a pure
// function of the delivery. A redelivery racing the original computes the same
// next id, so dedup collapses the two into one message and the build keeps a
// single poll chain — where a random suffix would fork a second chain that
// doubles the poll rate and races the first one's status CAS.
func nextPollMessageID(buildID, currentMsgID string) string {
generation := 0
if rest, found := strings.CutPrefix(currentMsgID, buildID+pollIDInfix); found {
if n, err := strconv.Atoi(rest); err == nil && n > 0 {
generation = n
}
}
return fmt.Sprintf("%s%s%d", buildID, pollIDInfix, generation+1)
}

// publishBuildSignal re-publishes buildID to buildsignal after delayMs,
// partitioned by build id so each build's poll loop runs in its own
// partition. A fresh message, not a nack — polling is not failure.
func (c *Controller) publishBuildSignal(ctx context.Context, buildID string, delayMs int64) error {
//
// currentMsgID is the id of the delivery being processed; see nextPollMessageID
// for why the new id is derived from it rather than reused or randomized.
func (c *Controller) publishBuildSignal(ctx context.Context, buildID, currentMsgID string, delayMs int64) error {
payload, err := stovepipemq.Marshal(&stovepipemq.BuildSignal{Id: buildID})
if err != nil {
return fmt.Errorf("failed to serialize build signal: %w", err)
}
msg := entityqueue.NewMessage(buildID, payload, buildID, nil)
msg := entityqueue.NewMessage(nextPollMessageID(buildID, currentMsgID), payload, buildID, nil)
return c.publish(ctx, stovepipemq.TopicKeyBuildSignal, msg, delayMs)
}

Expand Down
74 changes: 74 additions & 0 deletions stovepipe/controller/buildsignal/buildsignal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -326,3 +326,77 @@ func TestProcess(t *testing.T) {
})
}
}

// TestPublishBuildSignalAdvancesPollGeneration is the regression test for the poll
// loop stalling. The queue dedups on (topic, partition_key, id) and the delivery that
// scheduled a re-poll is still un-acked when the re-poll is published, so a reused
// message id makes the reschedule a silent no-op and the build is never polled again.
// The generation therefore has to advance each tick. The partition key must stay the
// build id so each poll loop keeps its own partition.
func TestPublishBuildSignalAdvancesPollGeneration(t *testing.T) {
ctrl := gomock.NewController(t)
c, m := newController(t, ctrl)

var ids []string
m.publisher.EXPECT().PublishAfter(gomock.Any(), "buildsignal", gomock.Any(), PollDelayRunningMs).
DoAndReturn(func(_ context.Context, _ string, msg entityqueue.Message, _ int64) error {
ids = append(ids, msg.ID)
assert.Equal(t, testBuildID, msg.PartitionKey)
return nil
}).Times(3)

// Walk a chain: each publish is scheduled by the message the previous one minted.
current := testBuildID
for range 3 {
require.NoError(t, c.publishBuildSignal(context.Background(), testBuildID, current, PollDelayRunningMs))
current = ids[len(ids)-1]
}

assert.Equal(t, []string{
testBuildID + "/poll/1",
testBuildID + "/poll/2",
testBuildID + "/poll/3",
}, ids, "each tick must mint a fresh id, or the reschedule dedups against the message that scheduled it")
}

// TestPublishBuildSignalIsIdempotentPerDelivery pins the other half of the contract:
// the next id is a pure function of the delivery being processed. A redelivery racing
// the original computes the same id, so dedup collapses them and the build keeps a
// single poll chain rather than forking a second one that doubles the poll rate and
// races the first one's status CAS.
func TestPublishBuildSignalIsIdempotentPerDelivery(t *testing.T) {
ctrl := gomock.NewController(t)
c, m := newController(t, ctrl)

var ids []string
m.publisher.EXPECT().PublishAfter(gomock.Any(), "buildsignal", gomock.Any(), PollDelayRunningMs).
DoAndReturn(func(_ context.Context, _ string, msg entityqueue.Message, _ int64) error {
ids = append(ids, msg.ID)
return nil
}).Times(2)

scheduledBy := testBuildID + "/poll/7"
require.NoError(t, c.publishBuildSignal(context.Background(), testBuildID, scheduledBy, PollDelayRunningMs))
require.NoError(t, c.publishBuildSignal(context.Background(), testBuildID, scheduledBy, PollDelayRunningMs))

assert.Equal(t, ids[0], ids[1], "a redelivery of the same message must republish the same id so dedup collapses it")
assert.Equal(t, testBuildID+"/poll/8", ids[0])
}

func TestNextPollMessageID(t *testing.T) {
tests := []struct {
name string
current string
expected string
}{
{name: "initial publish from build starts at one", current: testBuildID, expected: testBuildID + "/poll/1"},
{name: "advances the generation", current: testBuildID + "/poll/4", expected: testBuildID + "/poll/5"},
{name: "unparsable generation restarts at one", current: testBuildID + "/poll/x", expected: testBuildID + "/poll/1"},
{name: "another build's id is not a prefix match", current: "other/poll/9", expected: testBuildID + "/poll/1"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
assert.Equal(t, tt.expected, nextPollMessageID(testBuildID, tt.current))
})
}
}
18 changes: 18 additions & 0 deletions test/e2e/stovepipe/harness_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -130,3 +130,21 @@ func (s *StovepipeE2ESuite) assertIngestPersisted(queue, id string) {
assert.Equal(t, id, s.uriMapping(queue), "URI mapping should point at the minted request id")
assert.Equal(t, 1, s.publishedMessageCount(id), "should have published one process message for %s", id)
}

// awaitBuildStatus blocks until the build row for a request reaches want. The
// pipeline runs inside the stovepipe-service container, so the build's own status
// column is the durable, black-box signal that the poll loop converged: buildsignal
// is its only writer after build creates the row.
func (s *StovepipeE2ESuite) awaitBuildStatus(requestID, want string) {
pollUntil(processPollInterval, func() bool {
var status string
err := s.db.QueryRow("SELECT status FROM build WHERE request_id = ?", requestID).Scan(&status)
if err != nil {
// sql.ErrNoRows means the build row is not created yet.
s.log.Logf("build for %s not readable yet: %v", requestID, err)
return false
}
s.log.Logf("build for %s status = %q (want %q)", requestID, status, want)
return status == want
})
}
35 changes: 31 additions & 4 deletions test/e2e/stovepipe/suite_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,10 +30,11 @@ package e2e_test
// bazel test //test/e2e/stovepipe:stovepipe_test
//
// The stack runs the Stovepipe gRPC service plus a storage MySQL (request,
// request_uri) and a queue MySQL (the process stage). Unlike the integration
// suite (test/integration/stovepipe), which asserts only that Ingest *publishes*
// a process message, this suite additionally drives the asynchronous process
// consumer to completion — proving the ingest→process pipeline runs end-to-end.
// request_uri, queue, build) and a queue MySQL (the pipeline stages). Unlike the
// integration suite (test/integration/stovepipe), which asserts only that Ingest
// *publishes* a process message, this suite additionally drives the asynchronous
// consumers to completion — proving the ingest→process→build→buildsignal pipeline
// runs end-to-end.

import (
"context"
Expand Down Expand Up @@ -157,3 +158,29 @@ func (s *StovepipeE2ESuite) TestIngest_Idempotent() {
id2 := s.ingest(queue)
assert.Equal(s.T(), id, id2, "re-ingest of the same head should dedup to the same id")
}

// TestIngest_SlowBuild_PollsToCompletion drives a build that is not terminal on its
// first poll, which is the only path that exercises buildsignal's reschedule.
//
// The queue name carries a fake-buildrunner marker: the fake SourceControl resolves a
// queue to "git://<queue>/HEAD", so the marker rides into the head URI and the fake
// BuildRunner reports running for a while before succeeding. Reaching a terminal build
// status therefore requires the poll loop to tick more than once.
//
// This is the regression test for the loop stalling: buildsignal re-publishes to its
// own topic to schedule the next poll, and the queue dedups on
// (topic, partition_key, id). While the delivery being processed is still un-acked its
// row is present, so a re-poll that reuses the build id as the message id is silently
// discarded and the build is never polled again — the build would sit at `running`
// forever.
func (s *StovepipeE2ESuite) TestIngest_SlowBuild_PollsToCompletion() {
const queue = "monorepo/slow?buildrunner-fake=build-slow"

id := s.ingest(queue)
s.log.Logf("Ingest succeeded: id=%s; waiting for the poll loop to reach a terminal build", id)

s.assertIngestPersisted(queue, id)

// Getting here at all means the reschedule produced a deliverable message.
s.awaitBuildStatus(id, "succeeded")
}
Loading