Skip to content

Commit 8dc4fbb

Browse files
authored
feat(stovepipe): Announce validation results (#646)
## Summary **What**: - Publish an event whenever a commit's validation result becomes durable, carrying whether the commit is healthy or broken and what it was measured against. - Publish a distinct event when validation ends without a verdict, and fail the work when either event cannot be sent. **Why**: - Let subscriber systems act on trunk validation results as they land, instead of polling for them. - Keep the event vendor-neutral so each deployment wires its own subscriber to forward results onward, instead of building any single environment's delivery path into the pipeline. ## Test Plan - [x] Add unit tests. ## Revert Plan - Revert this PR. The events are additive and nothing consumes them yet. ## Issues - [CODEM-459](https://linear.app/uber/issue/CODEM-459/record-step-emit-events-internally-on-new-green)
1 parent a6a425d commit 8dc4fbb

13 files changed

Lines changed: 719 additions & 64 deletions

File tree

‎doc/rfc/stovepipe/steps/record.md‎

Lines changed: 34 additions & 26 deletions
Large diffs are not rendered by default.

‎doc/rfc/stovepipe/workflow.md‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -55,10 +55,10 @@ The ref is a *cache* of the last-green URI, not a second record of greenness. It
5555
|---|---|
5656
| **SourceControl** | Resolve a Queue name to its current head URI; answer ancestry/comparison questions between two URIs (is the new head a fast-forward descendant of the last green, or was history rewritten?); enumerate commits in a range; advance the Queue's **promotion ref** to a commit. The sole owner of URI semantics, including which refs a Queue name resolves to. |
5757
| **build-runner** | Build a scope at a URI (optionally relative to a baseline URI), returning pass/fail and the target graph. See [build-runner.md](../submitqueue/build-runner.md). |
58-
| **Hooks** | Deliver Stovepipe's greenness events to downstream systems — "this URI / this project is now green (or not green)". Fire-and-forget notification, decoupled so Stovepipe does not know or care who consumes the event. Not implemented yet; it will be the shared cross-domain hook seam rather than a Stovepipe-specific extension. See [hook-framework.md](../hook-framework.md). |
58+
| **Hooks** | Deliver Stovepipe's greenness events to downstream systems — "this URI / this project is now green (or not green)". Fire-and-forget notification, decoupled so Stovepipe does not know or care who consumes the event. The shared cross-domain hook seam rather than a Stovepipe-specific extension. See [hook-framework.md](../hook-framework.md). |
5959
| **Storage** | Persist Queues (incl. last-green URI), Requests, build records, and per-URI / per-project greenness. Key/value-shaped per the extension-design rules in [AGENTS.md](../../../AGENTS.md). |
6060

61-
Hooks are the notification boundary. When a validation fact is recorded — whole-repo green/not-green, or later a project green/not-green — the event reaches deployment systems, dashboards, and developer tooling without any of them polling Stovepipe's store, and each environment can route it to its own downstream (a deploy gate, a Slack notifier, an event bus) without changing the pipeline. The mechanism is the cross-domain hook framework rather than a call out of the recording stage: `record` publishes a `HookEvent` to Stovepipe's `hook` topic, and a dispatcher stage consumes it and invokes the wired hooks, so a slow or failing downstream cannot add latency to the pipeline. Neither half exists yet; see [record.md](steps/record.md#hooks) for the fact-to-event mapping and its open questions.
61+
Hooks are the notification boundary. When a validation fact is recorded — whole-repo green/not-green, or later a project green/not-green — the event reaches deployment systems, dashboards, and developer tooling without any of them polling Stovepipe's store, and each environment can route it to its own downstream (a deploy gate, a Slack notifier, an event bus) without changing the pipeline. The mechanism is the cross-domain hook framework rather than a call out of the recording stage: `record` publishes a `HookEvent` to Stovepipe's `hook` topic, and a dispatcher stage consumes it and invokes the wired hooks, so a slow or failing downstream cannot add latency to the pipeline. Both halves exist; what a deployment supplies is the hooks themselves, since the example server resolves every event to `noop`. See [record.md](steps/record.md#hooks) for the fact-to-event mapping.
6262

6363
## Workflow
6464

@@ -163,7 +163,7 @@ Per-stage design detail lives under `steps/` so this doc stays a pipeline overvi
163163
- [process.md](steps/process.md) — build-strategy decision, concurrency gate, backlog coalescing, [concurrency lifecycle](steps/process.md#concurrency-lifecycle), entity changes, [waiting for a slot](steps/process.md#waiting-for-a-slot)
164164
- [build.md](steps/build.md) — trigger-only stage: reads the decided scope off the Request, triggers the build-runner, hands off to buildsignal; the stovepipe `BuildRunner` contract and why it differs from SubmitQueue's
165165
- [buildsignal.md](steps/buildsignal.md) — the poll loop: hold-based re-poll cadence, target-graph return, per-build partitioning, and the fail-closed handoff to record
166-
- [record.md](steps/record.md) — turning a terminal build outcome into an immutable validation fact, monotonic last-green advancement and ref promotion, and the deferred hook and analyze handoffs
166+
- [record.md](steps/record.md) — turning a terminal build outcome into an immutable validation fact, monotonic last-green advancement and ref promotion, the hook event announcing the outcome, and the deferred analyze handoff
167167

168168
## Dedup, idempotency, and history rewrites
169169

‎platform/hook/BUILD.bazel‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ go_library(
55
srcs = [
66
"controller.go",
77
"dlq.go",
8+
"publisher.go",
89
],
910
importpath = "github.com/uber/submitqueue/platform/hook",
1011
visibility = ["//visibility:public"],
@@ -14,6 +15,7 @@ go_library(
1415
"//platform/errs:go_default_library",
1516
"//platform/extension/hook:go_default_library",
1617
"//platform/metrics:go_default_library",
18+
"//platform/publish:go_default_library",
1719
"@com_github_uber_go_tally//:go_default_library",
1820
"@org_uber_go_zap//:go_default_library",
1921
],
@@ -24,14 +26,17 @@ go_test(
2426
srcs = [
2527
"controller_test.go",
2628
"dlq_test.go",
29+
"publisher_test.go",
2730
],
2831
embed = [":go_default_library"],
2932
deps = [
3033
"//api/base/hook:go_default_library",
3134
"//platform/base/failure:go_default_library",
3235
"//platform/base/messagequeue:go_default_library",
36+
"//platform/consumer:go_default_library",
3337
"//platform/consumer/mock:go_default_library",
3438
"//platform/extension/hook:go_default_library",
39+
"//platform/extension/messagequeue/mock:go_default_library",
3540
"@com_github_stretchr_testify//assert:go_default_library",
3641
"@com_github_stretchr_testify//require:go_default_library",
3742
"@com_github_uber_go_tally//:go_default_library",

‎platform/hook/controller.go‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,9 +12,10 @@
1212
// See the License for the specific language governing permissions and
1313
// limitations under the License.
1414

15-
// Package hook holds the consumer side of the hooks framework: the controller
16-
// that turns hook events on a queue into hook.Hook calls, and the reconciler for
17-
// the events that never made it.
15+
// Package hook holds the domain-neutral mechanics of the hooks framework: the
16+
// controller that turns hook events on a queue into hook.Hook calls, the
17+
// reconciler for the events that never made it, and the helper a producer
18+
// publishes an event through.
1819
//
1920
// The controller is domain-neutral. Each domain runs its own hook topic and its
2021
// own instance of this stage — "per-domain" is about the topic and the wiring,

‎platform/hook/publisher.go‎

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
// Copyright (c) 2026 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package hook
16+
17+
import (
18+
"context"
19+
"fmt"
20+
21+
basehook "github.com/uber/submitqueue/api/base/hook"
22+
"github.com/uber/submitqueue/platform/consumer"
23+
"github.com/uber/submitqueue/platform/publish"
24+
)
25+
26+
// Publish sends one hook event to the domain's hook topic, partitioned by
27+
// partitionKey. The topic key is not a parameter: a domain runs a single hook
28+
// topic, and the caller's registry is what binds that key to a wire topic.
29+
//
30+
// The event id is the message id, so a redelivery republishing the same event
31+
// dedups into the original message instead of enqueuing a second one. Callers
32+
// pass the partition key their own topic partitions on, carrying that ordering
33+
// across the seam.
34+
//
35+
// Errors are returned unclassified, leaving retryability to the caller's
36+
// classifier: a malformed event is the caller's bug, not a transient fault.
37+
func Publish(
38+
ctx context.Context,
39+
registry consumer.TopicRegistry,
40+
event *basehook.HookEvent,
41+
partitionKey string,
42+
) error {
43+
if err := basehook.Validate(event); err != nil {
44+
return fmt.Errorf("refusing to publish a malformed hook event: %w", err)
45+
}
46+
47+
body, err := basehook.Marshal(event)
48+
if err != nil {
49+
return fmt.Errorf("failed to serialize hook event %s: %w", event.GetId(), err)
50+
}
51+
52+
if err := publish.Message(
53+
ctx, registry, basehook.TopicKeyHook, publish.IntentID(event.GetId()), body, partitionKey,
54+
); err != nil {
55+
return fmt.Errorf("failed to publish hook event %s: %w", event.GetId(), err)
56+
}
57+
return nil
58+
}

‎platform/hook/publisher_test.go‎

Lines changed: 125 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,125 @@
1+
// Copyright (c) 2026 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package hook
16+
17+
import (
18+
"context"
19+
"errors"
20+
"testing"
21+
22+
"github.com/stretchr/testify/assert"
23+
"github.com/stretchr/testify/require"
24+
basehook "github.com/uber/submitqueue/api/base/hook"
25+
entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
26+
"github.com/uber/submitqueue/platform/consumer"
27+
mqmock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
28+
"go.uber.org/mock/gomock"
29+
)
30+
31+
const (
32+
testEventID = "stovepipe/validation.repository.recorded/request/7/0"
33+
testPartitionKey = "request/7"
34+
)
35+
36+
func testEvent() *basehook.HookEvent {
37+
return &basehook.HookEvent{
38+
Id: testEventID,
39+
Source: "stovepipe",
40+
Type: "validation.repository.recorded",
41+
TimestampMs: 1756327200000,
42+
}
43+
}
44+
45+
// registryWithHookTopic returns a registry whose hook topic captures whatever is
46+
// published to it, and the slot the captured message lands in.
47+
func registryWithHookTopic(t *testing.T, ctrl *gomock.Controller, publishErr error) (consumer.TopicRegistry, *entityqueue.Message) {
48+
t.Helper()
49+
50+
var published entityqueue.Message
51+
publisher := mqmock.NewMockPublisher(ctrl)
52+
publisher.EXPECT().Publish(gomock.Any(), "domain-hook", gomock.Any()).
53+
DoAndReturn(func(_ context.Context, _ string, msg entityqueue.Message) error {
54+
published = msg
55+
return publishErr
56+
}).AnyTimes()
57+
58+
queue := mqmock.NewMockQueue(ctrl)
59+
queue.EXPECT().Publisher().Return(publisher).AnyTimes()
60+
61+
registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{
62+
{Key: basehook.TopicKeyHook, Name: "domain-hook", Queue: queue},
63+
})
64+
require.NoError(t, err)
65+
return registry, &published
66+
}
67+
68+
func TestPublish(t *testing.T) {
69+
ctrl := gomock.NewController(t)
70+
registry, published := registryWithHookTopic(t, ctrl, nil)
71+
72+
require.NoError(t, Publish(context.Background(), registry, testEvent(), testPartitionKey))
73+
74+
// The event id is the message id, so a redelivery republishing the same
75+
// event dedups into the original message.
76+
assert.Equal(t, testEventID, published.ID)
77+
assert.Equal(t, testPartitionKey, published.PartitionKey)
78+
79+
decoded := &basehook.HookEvent{}
80+
require.NoError(t, basehook.Unmarshal(published.Payload, decoded))
81+
assert.Equal(t, testEventID, decoded.GetId())
82+
assert.Equal(t, "validation.repository.recorded", decoded.GetType())
83+
}
84+
85+
func TestPublish_RejectsMalformedEvent(t *testing.T) {
86+
tests := []struct {
87+
name string
88+
event *basehook.HookEvent
89+
}{
90+
{name: "nil", event: nil},
91+
{name: "no id", event: &basehook.HookEvent{Source: "stovepipe", Type: "t"}},
92+
{name: "no source", event: &basehook.HookEvent{Id: testEventID, Type: "t"}},
93+
{name: "no type", event: &basehook.HookEvent{Id: testEventID, Source: "stovepipe"}},
94+
}
95+
96+
for _, tt := range tests {
97+
t.Run(tt.name, func(t *testing.T) {
98+
ctrl := gomock.NewController(t)
99+
publisher := mqmock.NewMockPublisher(ctrl)
100+
queue := mqmock.NewMockQueue(ctrl)
101+
queue.EXPECT().Publisher().Return(publisher).AnyTimes()
102+
registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{
103+
{Key: basehook.TopicKeyHook, Name: "domain-hook", Queue: queue},
104+
})
105+
require.NoError(t, err)
106+
107+
// No Publish expectation: a malformed event must not reach the queue.
108+
require.Error(t, Publish(context.Background(), registry, tt.event, testPartitionKey))
109+
})
110+
}
111+
}
112+
113+
func TestPublish_PropagatesPublishFailure(t *testing.T) {
114+
ctrl := gomock.NewController(t)
115+
registry, _ := registryWithHookTopic(t, ctrl, errors.New("boom"))
116+
117+
require.Error(t, Publish(context.Background(), registry, testEvent(), testPartitionKey))
118+
}
119+
120+
func TestPublish_FailsWhenHookTopicIsUnregistered(t *testing.T) {
121+
registry, err := consumer.NewTopicRegistry(nil)
122+
require.NoError(t, err)
123+
124+
require.Error(t, Publish(context.Background(), registry, testEvent(), testPartitionKey))
125+
}

‎service/stovepipe/server/main.go‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -447,7 +447,7 @@ func registerPrimaryControllers(
447447
}
448448
count++
449449

450-
recordController := record.NewController(logger, scope, store, sourceControl, stovepipemq.TopicKeyRecord, "stovepipe-record")
450+
recordController := record.NewController(logger, scope, store, sourceControl, registry, stovepipemq.TopicKeyRecord, "stovepipe-record")
451451
if err := c.Register(recordController); err != nil {
452452
return count, fmt.Errorf("failed to register record controller: %w", err)
453453
}
@@ -492,7 +492,7 @@ func registerDLQControllers(
492492
}
493493
count++
494494

495-
recordDLQController := record.NewController(logger, scope, store, sourceControl, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq")
495+
recordDLQController := record.NewController(logger, scope, store, sourceControl, registry, dlq.TopicKey(stovepipemq.TopicKeyRecord), "stovepipe-record-dlq")
496496
if err := c.Register(recordDLQController); err != nil {
497497
return count, fmt.Errorf("failed to register record dlq controller: %w", err)
498498
}

‎stovepipe/controller/record/BUILD.bazel‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,11 @@ go_library(
66
importpath = "github.com/uber/submitqueue/stovepipe/controller/record",
77
visibility = ["//visibility:public"],
88
deps = [
9+
"//api/base/hook:go_default_library",
910
"//platform/consumer:go_default_library",
11+
"//platform/hook:go_default_library",
1012
"//platform/metrics:go_default_library",
13+
"//stovepipe/core/hookevent:go_default_library",
1114
"//stovepipe/core/loader:go_default_library",
1215
"//stovepipe/core/messagequeue:go_default_library",
1316
"//stovepipe/entity:go_default_library",
@@ -23,10 +26,13 @@ go_test(
2326
srcs = ["record_test.go"],
2427
embed = [":go_default_library"],
2528
deps = [
29+
"//api/base/hook:go_default_library",
2630
"//platform/base/messagequeue:go_default_library",
2731
"//platform/consumer:go_default_library",
2832
"//platform/consumer/mock:go_default_library",
33+
"//platform/extension/messagequeue/mock:go_default_library",
2934
"//platform/metrics:go_default_library",
35+
"//stovepipe/core/hookevent:go_default_library",
3036
"//stovepipe/core/messagequeue:go_default_library",
3137
"//stovepipe/entity:go_default_library",
3238
"//stovepipe/extension/sourcecontrol:go_default_library",

0 commit comments

Comments
 (0)