Skip to content

Commit 06d9971

Browse files
Merge commit from fork
* fix(controlplane): contain panics in request-spawned goroutines Goroutines spawned from gRPC handlers ran under a bare `go`: the attestation CAS upload, the integration dispatcher, and the per-plugin fan-out dispatch. The gRPC recovery middleware only wraps the handler's own stack, so an unrecovered panic in one of these detached goroutines terminated the whole control plane process. Add an internal panicguard helper that runs background work with a recovery frame, logging and swallowing any panic, and route the three detached goroutines through it. Signed-off-by: Matías Insaurralde <matias@chainloop.dev> * fix(controlplane): build fan-out dispatch metadata without nil derefs The fan-out dispatcher built its workflow-run metadata by unconditionally dereferencing CreatedAt, FinishedAt and Attestation. These are nil until a run is finalized, so dispatching for an unfinished or attestation-less run panicked. Extract the metadata construction into newWorkflowMetadata and guard each pointer, mapping a nil to the zero value. Add a regression test that builds metadata for an unfinished run and asserts it does not panic. Signed-off-by: Matías Insaurralde <matias@chainloop.dev> --------- Signed-off-by: Matías Insaurralde <matias@chainloop.dev>
1 parent 28d66cd commit 06d9971

5 files changed

Lines changed: 284 additions & 25 deletions

File tree

‎app/controlplane/internal/dispatcher/dispatcher.go‎

Lines changed: 38 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ import (
2828

2929
"github.com/cenkalti/backoff/v4"
3030

31+
"github.com/chainloop-dev/chainloop/app/controlplane/internal/panicguard"
3132
"github.com/chainloop-dev/chainloop/app/controlplane/pkg/biz"
3233
"github.com/chainloop-dev/chainloop/app/controlplane/plugins/sdk/v1"
3334
"github.com/chainloop-dev/chainloop/pkg/attestation/renderer/chainloop"
@@ -122,30 +123,15 @@ func (d *FanOutDispatcher) Run(ctx context.Context, opts *RunOpts) error {
122123
return fmt.Errorf("workflowRun not found")
123124
}
124125

125-
workflowMetadata := &sdk.ChainloopMetadata{
126-
Workflow: &sdk.ChainloopMetadataWorkflow{
127-
ID: opts.WorkflowID,
128-
Name: wf.Name,
129-
Project: wf.Project,
130-
Team: wf.Team,
131-
},
132-
WorkflowRun: &sdk.ChainloopMetadataWorkflowRun{
133-
ID: opts.WorkflowRunID,
134-
State: wfRun.State,
135-
StartedAt: *wfRun.CreatedAt,
136-
FinishedAt: *wfRun.FinishedAt,
137-
RunnerType: wfRun.RunnerType,
138-
RunURL: wfRun.RunURL,
139-
AttestationDigest: wfRun.Attestation.Digest,
140-
},
141-
}
126+
workflowMetadata := newWorkflowMetadata(opts, wf, wfRun)
142127

143128
// Dispatch the integrations
144129
for _, item := range queue {
145130
req := generateRequest(item, workflowMetadata)
146-
go func(p sdk.FanOut, r *sdk.ExecutionRequest) {
147-
_ = dispatch(ctx, p, req, d.log)
148-
}(item.plugin, req)
131+
plugin := item.plugin
132+
panicguard.Go(d.log, "fanout-dispatch", func() {
133+
_ = dispatch(ctx, plugin, req, d.log)
134+
})
149135
}
150136

151137
return nil
@@ -326,6 +312,38 @@ func dispatch(ctx context.Context, plugin sdk.FanOut, opts *sdk.ExecutionRequest
326312
)
327313
}
328314

315+
// newWorkflowMetadata builds the dispatch metadata for a run. CreatedAt,
316+
// FinishedAt and Attestation are nil until the run is finalized, so a dispatch
317+
// for an unfinished or attestation-less run must not dereference them.
318+
func newWorkflowMetadata(opts *RunOpts, wf *biz.Workflow, wfRun *biz.WorkflowRun) *sdk.ChainloopMetadata {
319+
runMetadata := &sdk.ChainloopMetadataWorkflowRun{
320+
ID: opts.WorkflowRunID,
321+
State: wfRun.State,
322+
RunnerType: wfRun.RunnerType,
323+
RunURL: wfRun.RunURL,
324+
}
325+
326+
if wfRun.CreatedAt != nil {
327+
runMetadata.StartedAt = *wfRun.CreatedAt
328+
}
329+
if wfRun.FinishedAt != nil {
330+
runMetadata.FinishedAt = *wfRun.FinishedAt
331+
}
332+
if wfRun.Attestation != nil {
333+
runMetadata.AttestationDigest = wfRun.Attestation.Digest
334+
}
335+
336+
return &sdk.ChainloopMetadata{
337+
Workflow: &sdk.ChainloopMetadataWorkflow{
338+
ID: opts.WorkflowID,
339+
Name: wf.Name,
340+
Project: wf.Project,
341+
Team: wf.Team,
342+
},
343+
WorkflowRun: runMetadata,
344+
}
345+
}
346+
329347
func generateRequest(in *dispatchItem, metadata *sdk.ChainloopMetadata) *sdk.ExecutionRequest {
330348
return &sdk.ExecutionRequest{
331349
ChainloopMetadata: metadata,
Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
1+
//
2+
// Copyright 2026 The Chainloop Authors.
3+
//
4+
// Licensed under the Apache License, Version 2.0 (the "License");
5+
// you may not use this file except in compliance with the License.
6+
// You may obtain a copy of the License at
7+
//
8+
// http://www.apache.org/licenses/LICENSE-2.0
9+
//
10+
// Unless required by applicable law or agreed to in writing, software
11+
// distributed under the License is distributed on an "AS IS" BASIS,
12+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
// See the License for the specific language governing permissions and
14+
// limitations under the License.
15+
16+
package dispatcher
17+
18+
import (
19+
"testing"
20+
"time"
21+
22+
"github.com/chainloop-dev/chainloop/app/controlplane/pkg/biz"
23+
"github.com/google/uuid"
24+
"github.com/stretchr/testify/assert"
25+
"github.com/stretchr/testify/require"
26+
)
27+
28+
// TestNewWorkflowMetadata_UnfinishedRun pins the crash condition: a run that is
29+
// not yet finished has a nil FinishedAt (and may have a nil Attestation). The
30+
// dispatcher must build its metadata without dereferencing those nil pointers.
31+
func TestNewWorkflowMetadata_UnfinishedRun(t *testing.T) {
32+
createdAt := time.Now()
33+
wfRun := &biz.WorkflowRun{
34+
ID: uuid.New(),
35+
State: "initialized",
36+
CreatedAt: &createdAt,
37+
FinishedAt: nil, // unfinished run
38+
RunnerType: "GENERIC",
39+
RunURL: "https://ci.example/run/1",
40+
// Attestation is nil until the bundle is stored.
41+
}
42+
wf := &biz.Workflow{Name: "wf", Project: "proj", Team: "team"}
43+
opts := &RunOpts{WorkflowID: uuid.NewString(), WorkflowRunID: wfRun.ID.String()}
44+
45+
require.NotPanics(t, func() {
46+
got := newWorkflowMetadata(opts, wf, wfRun)
47+
require.NotNil(t, got)
48+
require.NotNil(t, got.WorkflowRun)
49+
50+
assert.True(t, got.WorkflowRun.FinishedAt.IsZero(), "nil FinishedAt must map to the zero time")
51+
assert.Equal(t, createdAt, got.WorkflowRun.StartedAt)
52+
assert.Empty(t, got.WorkflowRun.AttestationDigest, "nil Attestation must map to an empty digest")
53+
assert.Equal(t, "wf", got.Workflow.Name)
54+
assert.Equal(t, "proj", got.Workflow.Project)
55+
})
56+
}
57+
58+
// TestNewWorkflowMetadata_FinishedRun confirms the populated fields still flow
59+
// through unchanged once the run is finalized.
60+
func TestNewWorkflowMetadata_FinishedRun(t *testing.T) {
61+
createdAt := time.Now().Add(-time.Hour)
62+
finishedAt := time.Now()
63+
wfRun := &biz.WorkflowRun{
64+
ID: uuid.New(),
65+
State: "success",
66+
CreatedAt: &createdAt,
67+
FinishedAt: &finishedAt,
68+
RunnerType: "GENERIC",
69+
RunURL: "https://ci.example/run/1",
70+
Attestation: &biz.Attestation{Digest: "sha256:abc"},
71+
}
72+
wf := &biz.Workflow{Name: "wf", Project: "proj", Team: "team"}
73+
opts := &RunOpts{WorkflowID: uuid.NewString(), WorkflowRunID: wfRun.ID.String()}
74+
75+
got := newWorkflowMetadata(opts, wf, wfRun)
76+
77+
assert.Equal(t, createdAt, got.WorkflowRun.StartedAt)
78+
assert.Equal(t, finishedAt, got.WorkflowRun.FinishedAt)
79+
assert.Equal(t, "sha256:abc", got.WorkflowRun.AttestationDigest)
80+
assert.Equal(t, "success", got.WorkflowRun.State)
81+
}
Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
//
2+
// Copyright 2026 The Chainloop Authors.
3+
//
4+
// Licensed under the Apache License, Version 2.0 (the "License");
5+
// you may not use this file except in compliance with the License.
6+
// You may obtain a copy of the License at
7+
//
8+
// http://www.apache.org/licenses/LICENSE-2.0
9+
//
10+
// Unless required by applicable law or agreed to in writing, software
11+
// distributed under the License is distributed on an "AS IS" BASIS,
12+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
// See the License for the specific language governing permissions and
14+
// limitations under the License.
15+
16+
// Package panicguard runs background work in a goroutine with a recovery frame.
17+
//
18+
// The gRPC recovery middleware only wraps the request handler's own stack. A
19+
// goroutine spawned from a handler escapes it, so an unrecovered panic there
20+
// terminates the whole process. Any request-spawned goroutine must therefore
21+
// run through this package instead of a bare `go`.
22+
package panicguard
23+
24+
import (
25+
"runtime/debug"
26+
27+
"github.com/go-kratos/kratos/v2/log"
28+
)
29+
30+
// Go runs fn in a new goroutine with a recovery frame. A recovered panic is
31+
// logged with its stack and swallowed so it cannot crash the process.
32+
func Go(logger *log.Helper, task string, fn func()) {
33+
go Recover(logger, task, fn)
34+
}
35+
36+
// Recover runs fn in the current goroutine with a recovery frame. It is the
37+
// synchronous core of Go, exposed for callers that already own a goroutine and
38+
// for tests.
39+
func Recover(logger *log.Helper, task string, fn func()) {
40+
defer func() {
41+
if r := recover(); r != nil {
42+
if logger != nil {
43+
logger.Errorw(
44+
"msg", "recovered from panic in background task",
45+
"task", task,
46+
"panic", r,
47+
"stack", string(debug.Stack()),
48+
)
49+
}
50+
}
51+
}()
52+
53+
fn()
54+
}
Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,104 @@
1+
//
2+
// Copyright 2026 The Chainloop Authors.
3+
//
4+
// Licensed under the Apache License, Version 2.0 (the "License");
5+
// you may not use this file except in compliance with the License.
6+
// You may obtain a copy of the License at
7+
//
8+
// http://www.apache.org/licenses/LICENSE-2.0
9+
//
10+
// Unless required by applicable law or agreed to in writing, software
11+
// distributed under the License is distributed on an "AS IS" BASIS,
12+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
// See the License for the specific language governing permissions and
14+
// limitations under the License.
15+
16+
package panicguard_test
17+
18+
import (
19+
"sync"
20+
"testing"
21+
"time"
22+
23+
"github.com/chainloop-dev/chainloop/app/controlplane/internal/panicguard"
24+
"github.com/go-kratos/kratos/v2/log"
25+
"github.com/stretchr/testify/assert"
26+
"github.com/stretchr/testify/require"
27+
)
28+
29+
// captureLogger counts Log calls so a test can assert a recovery was logged.
30+
type captureLogger struct {
31+
mu sync.Mutex
32+
n int
33+
}
34+
35+
func (c *captureLogger) Log(_ log.Level, _ ...interface{}) error {
36+
c.mu.Lock()
37+
c.n++
38+
c.mu.Unlock()
39+
return nil
40+
}
41+
42+
func (c *captureLogger) count() int {
43+
c.mu.Lock()
44+
defer c.mu.Unlock()
45+
return c.n
46+
}
47+
48+
// TestRecover_ContainsPanic is the core assertion: a panicking fn returns
49+
// normally instead of crashing the process. Without the recovery frame this
50+
// panic would abort the test binary.
51+
func TestRecover_ContainsPanic(t *testing.T) {
52+
cl := &captureLogger{}
53+
logger := log.NewHelper(cl)
54+
55+
require.NotPanics(t, func() {
56+
panicguard.Recover(logger, "unit", func() {
57+
panic("boom")
58+
})
59+
})
60+
61+
assert.Equal(t, 1, cl.count(), "the recovered panic must be logged once")
62+
}
63+
64+
// TestRecover_RunsFnAndReturns confirms the happy path is untouched.
65+
func TestRecover_RunsFnAndReturns(t *testing.T) {
66+
cl := &captureLogger{}
67+
logger := log.NewHelper(cl)
68+
69+
ran := false
70+
panicguard.Recover(logger, "unit", func() { ran = true })
71+
72+
assert.True(t, ran, "fn must run")
73+
assert.Equal(t, 0, cl.count(), "no panic means nothing is logged")
74+
}
75+
76+
// TestGo_ContainsConcurrentPanics confirms the async wrapper contains panics in
77+
// detached goroutines. The process surviving to the assertion is the proof.
78+
func TestGo_ContainsConcurrentPanics(t *testing.T) {
79+
cl := &captureLogger{}
80+
logger := log.NewHelper(cl)
81+
82+
const n = 50
83+
for i := 0; i < n; i++ {
84+
panicguard.Go(logger, "unit", func() { panic("boom") })
85+
}
86+
87+
require.Eventually(t, func() bool { return cl.count() == n }, time.Second, time.Millisecond,
88+
"every detached panic must be recovered and logged")
89+
}
90+
91+
// TestGo_NilLoggerDoesNotPanic guards the defensive nil-logger branch.
92+
func TestGo_NilLoggerDoesNotPanic(t *testing.T) {
93+
done := make(chan struct{})
94+
panicguard.Go(nil, "unit", func() {
95+
defer close(done)
96+
panic("boom")
97+
})
98+
99+
select {
100+
case <-done:
101+
case <-time.After(time.Second):
102+
t.Fatal("goroutine did not run")
103+
}
104+
}

‎app/controlplane/internal/service/attestation.go‎

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ import (
2828
cpAPI "github.com/chainloop-dev/chainloop/app/controlplane/api/controlplane/v1"
2929
conf "github.com/chainloop-dev/chainloop/app/controlplane/internal/conf/controlplane/config/v1"
3030
"github.com/chainloop-dev/chainloop/app/controlplane/internal/dispatcher"
31+
"github.com/chainloop-dev/chainloop/app/controlplane/internal/panicguard"
3132
"github.com/chainloop-dev/chainloop/app/controlplane/internal/usercontext"
3233
"github.com/chainloop-dev/chainloop/app/controlplane/internal/usercontext/attjwtmiddleware"
3334
"github.com/chainloop-dev/chainloop/app/controlplane/pkg/authz"
@@ -360,11 +361,12 @@ func (s *AttestationService) storeAttestation(ctx context.Context, bundle []byte
360361

361362
if !casBackend.Inline {
362363
// Detach from the request context so the upload survives request completion.
363-
go func(digest v1.Hash) {
364-
if err := s.uploadAttestationToCASWithRetry(context.Background(), bundle, casBackend, workflowRunID, digest); err != nil {
364+
dgst := *digest
365+
panicguard.Go(s.log, "attestation-cas-upload", func() {
366+
if err := s.uploadAttestationToCASWithRetry(context.Background(), bundle, casBackend, workflowRunID, dgst); err != nil {
365367
_ = handleUseCaseErr(err, s.log)
366368
}
367-
}(*digest)
369+
})
368370
}
369371
}
370372

@@ -394,7 +396,7 @@ func (s *AttestationService) storeAttestation(ctx context.Context, bundle []byte
394396
secretName := casBackend.SecretName
395397

396398
// Run integrations dispatcher
397-
go func() {
399+
panicguard.Go(s.log, "integration-dispatcher", func() {
398400
if err := s.integrationDispatcher.Run(context.TODO(), &dispatcher.RunOpts{
399401
Envelope: dsseEnv, OrgID: robotAccount.OrgID, WorkflowID: wf.ID.String(),
400402
DownloadBackendType: string(casBackend.Provider),
@@ -403,7 +405,7 @@ func (s *AttestationService) storeAttestation(ctx context.Context, bundle []byte
403405
}); err != nil {
404406
_ = handleUseCaseErr(err, s.log)
405407
}
406-
}()
408+
})
407409

408410
// promote release if the workflowRun is successful
409411
if markAsReleased != nil && *markAsReleased {

0 commit comments

Comments
 (0)