Skip to content

Commit baaed28

Browse files
committed
[ACTP] keep executor alive during key sync
1 parent 406b628 commit baaed28

2 files changed

Lines changed: 68 additions & 13 deletions

File tree

pkg/privateactionrunner/executor/server.go

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,7 @@ type Server struct {
5151

5252
ready atomic.Bool
5353
active atomic.Int32
54+
busy atomic.Int32
5455

5556
lastActivity atomic.Int64
5657
clock clock.Clock
@@ -72,8 +73,17 @@ func (s *Server) touch() {
7273
s.lastActivity.Store(s.clock.Now().UnixNano())
7374
}
7475

76+
func (s *Server) beginActivity() func() {
77+
s.touch()
78+
s.busy.Add(1)
79+
return func() {
80+
s.touch()
81+
s.busy.Add(-1)
82+
}
83+
}
84+
7585
func (s *Server) idleFor() time.Duration {
76-
if s.active.Load() > 0 {
86+
if s.busy.Load() > 0 {
7787
return 0
7888
}
7989
return s.clock.Since(time.Unix(0, s.lastActivity.Load()))
@@ -109,7 +119,8 @@ func (s *Server) SyncKeys(ctx context.Context, req *pb.SyncKeysRequest) (*pb.Syn
109119
Key: append([]byte(nil), key.GetKey()...),
110120
})
111121
}
112-
s.touch()
122+
finishActivity := s.beginActivity()
123+
defer finishActivity()
113124
if err := s.keysManager.Seed(seed); err != nil {
114125
return nil, status.Errorf(codes.InvalidArgument, "invalid signing-key seed: %v", err)
115126
}
@@ -143,11 +154,11 @@ func (s *Server) RunAction(req *pb.RunActionRequest, stream pb.Executor_RunActio
143154
))
144155
}
145156

146-
s.touch()
157+
finishActivity := s.beginActivity()
147158
s.active.Add(1)
148159
defer func() {
149-
s.touch()
150160
s.active.Add(-1)
161+
finishActivity()
151162
}()
152163

153164
// Raw bytes must stay unmodified for signature verification.

pkg/privateactionrunner/executor/server_test.go

Lines changed: 53 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -71,13 +71,27 @@ func (f *fakeExecutor) RunPrepared(ctx context.Context, _ *runners.PreparedWorkf
7171
}
7272

7373
type fakeKeysManager struct {
74-
seed []taskverifier.SigningKey
75-
snapshot []taskverifier.SigningKey
74+
seed []taskverifier.SigningKey
75+
snapshot []taskverifier.SigningKey
76+
readyGate <-chan struct{}
77+
waitStarted chan<- struct{}
7678
}
7779

78-
func (f *fakeKeysManager) Start(context.Context) {}
79-
func (f *fakeKeysManager) GetKey(string) types.DecodedKey { return nil }
80-
func (f *fakeKeysManager) WaitForReady(context.Context) error { return nil }
80+
func (f *fakeKeysManager) Start(context.Context) {}
81+
func (f *fakeKeysManager) GetKey(string) types.DecodedKey { return nil }
82+
func (f *fakeKeysManager) WaitForReady(ctx context.Context) error {
83+
if f.waitStarted != nil {
84+
f.waitStarted <- struct{}{}
85+
}
86+
if f.readyGate != nil {
87+
select {
88+
case <-f.readyGate:
89+
case <-ctx.Done():
90+
return ctx.Err()
91+
}
92+
}
93+
return nil
94+
}
8195
func (f *fakeKeysManager) Seed(keys []taskverifier.SigningKey) error {
8296
f.seed = keys
8397
return nil
@@ -262,6 +276,34 @@ func TestServeSyncKeysSeedsAndReturnsCurrentSnapshot(t *testing.T) {
262276
require.True(t, health.GetReady())
263277
}
264278

279+
func TestSyncKeysSuppressesIdleExitWhileWaitingForKeys(t *testing.T) {
280+
readyGate := make(chan struct{})
281+
waitStarted := make(chan struct{}, 1)
282+
keysManager := &fakeKeysManager{readyGate: readyGate, waitStarted: waitStarted}
283+
srv := NewServer(&fakeExecutor{}, "test-version", keysManager)
284+
mockClock := clock.NewMock()
285+
srv.clock = mockClock
286+
srv.touch()
287+
288+
done := make(chan error, 1)
289+
go func() {
290+
_, err := srv.SyncKeys(context.Background(), &pb.SyncKeysRequest{})
291+
done <- err
292+
}()
293+
<-waitStarted
294+
295+
mockClock.Add(2 * time.Minute)
296+
assert.Zero(t, srv.idleFor(), "key synchronization is executor activity")
297+
health, err := srv.Health(context.Background(), &pb.HealthRequest{})
298+
require.NoError(t, err)
299+
assert.Zero(t, health.ActiveActions, "key synchronization is not an action")
300+
301+
close(readyGate)
302+
require.NoError(t, <-done)
303+
mockClock.Add(time.Minute)
304+
assert.Equal(t, time.Minute, srv.idleFor())
305+
}
306+
265307
func TestServeRunActionStreamsOutputAndForwardsRawTask(t *testing.T) {
266308
fake := &fakeExecutor{
267309
prepared: &runners.PreparedWorkflowTask{Task: &types.Task{}},
@@ -595,11 +637,13 @@ func TestIdleTracking(t *testing.T) {
595637
mockClock.Add(time.Second)
596638
assert.Equal(t, time.Second, srv.idleFor(), "health checks should not count as activity")
597639

598-
srv.active.Add(1)
640+
finishActivity := srv.beginActivity()
599641
mockClock.Add(2 * timeout)
600-
assert.Zero(t, srv.idleFor(), "an active action should suppress the idle state")
601-
srv.active.Add(-1)
602-
srv.touch()
642+
assert.Zero(t, srv.idleFor(), "activity should suppress the idle state")
643+
health, err := srv.Health(context.Background(), &pb.HealthRequest{})
644+
require.NoError(t, err)
645+
assert.Zero(t, health.ActiveActions, "non-action activity must not affect action accounting")
646+
finishActivity()
603647

604648
mockClock.Add(timeout)
605649
assert.Equal(t, timeout, srv.idleFor())

0 commit comments

Comments
 (0)