Skip to content

Commit d0a19d9

Browse files
markturanskyAmbient Code Botclaude
authored
fix(api-server): authorize runner OIDC service account in WatchSessionMessages (#1183)
## Summary - Runner's BOT_TOKEN is an OIDC JWT with `preferred_username=service-account-ocm-ams-service` - `AMBIENT_API_TOKEN` is not set on the api-server deployment → `IsServiceCaller` is always false - Runner's JWT was parsed as a regular user → ownership check fired → `PERMISSION_DENIED: not authorized to watch this session` - Fix: read `GRPC_SERVICE_ACCOUNT` env var at startup into a package-level var; bypass ownership enforcement when the authenticated username matches that value ## Test plan - [ ] Deploy updated api-server image - [ ] Run `demo-github.sh` — runner pod should connect to gRPC stream and execute the task without PERMISSION_DENIED - [ ] Verify session reaches `Completed` phase 🤖 Generated with [Claude Code](https://claude.ai/code) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **New Features** * Added service account authentication support for gRPC operations * Introduced automatic token refresh for long-running sessions * **Bug Fixes** * Fixed unnecessary credential patch emissions when values haven't changed * Improved event handler retry mechanisms for better reliability <!-- end of auto-generated comment: release notes by coderabbit.ai --> --------- Co-authored-by: Ambient Code Bot <bot@ambient-code.local> Co-authored-by: Claude <noreply@anthropic.com>
1 parent 281a643 commit d0a19d9

8 files changed

Lines changed: 227 additions & 62 deletions

File tree

components/ambient-api-server/plugins/sessions/grpc_handler.go

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -21,18 +21,20 @@ import (
2121

2222
type sessionGRPCHandler struct {
2323
pb.UnimplementedSessionServiceServer
24-
service SessionService
25-
generic services.GenericService
26-
brokerFunc func() *server.EventBroker
27-
msgService MessageService
24+
service SessionService
25+
generic services.GenericService
26+
brokerFunc func() *server.EventBroker
27+
msgService MessageService
28+
grpcServiceAccount string
2829
}
2930

30-
func NewSessionGRPCHandler(service SessionService, generic services.GenericService, brokerFunc func() *server.EventBroker, msgService MessageService) pb.SessionServiceServer {
31+
func NewSessionGRPCHandler(service SessionService, generic services.GenericService, brokerFunc func() *server.EventBroker, msgService MessageService, grpcServiceAccount string) pb.SessionServiceServer {
3132
return &sessionGRPCHandler{
32-
service: service,
33-
generic: generic,
34-
brokerFunc: brokerFunc,
35-
msgService: msgService,
33+
service: service,
34+
generic: generic,
35+
brokerFunc: brokerFunc,
36+
msgService: msgService,
37+
grpcServiceAccount: grpcServiceAccount,
3638
}
3739
}
3840

@@ -286,7 +288,7 @@ func (h *sessionGRPCHandler) WatchSessionMessages(req *pb.WatchSessionMessagesRe
286288

287289
if !middleware.IsServiceCaller(ctx) {
288290
username := auth.GetUsernameFromContext(ctx)
289-
if username != "" {
291+
if username != "" && (h.grpcServiceAccount == "" || username != h.grpcServiceAccount) {
290292
session, svcErr := h.service.Get(ctx, req.GetSessionId())
291293
if svcErr != nil {
292294
return grpcutil.ServiceErrorToGRPC(svcErr)

components/ambient-api-server/plugins/sessions/plugin.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,10 @@ package sessions
22

33
import (
44
"net/http"
5+
"os"
56
"sync"
67

8+
79
pb "github.com/ambient-code/platform/components/ambient-api-server/pkg/api/grpc/ambient/v1"
810
pkgrbac "github.com/ambient-code/platform/components/ambient-api-server/plugins/rbac"
911
"github.com/gorilla/mux"
@@ -22,6 +24,8 @@ import (
2224

2325
const EventSource = "Sessions"
2426

27+
var grpcServiceAccount = os.Getenv("GRPC_SERVICE_ACCOUNT")
28+
2529
type ServiceLocator func() SessionService
2630

2731
func NewServiceLocator(env *environments.Env) ServiceLocator {
@@ -135,7 +139,7 @@ func init() {
135139
}
136140
return nil
137141
}
138-
pb.RegisterSessionServiceServer(grpcServer, NewSessionGRPCHandler(sessionService, genericService, brokerFunc, msgService))
142+
pb.RegisterSessionServiceServer(grpcServer, NewSessionGRPCHandler(sessionService, genericService, brokerFunc, msgService, grpcServiceAccount))
139143
})
140144

141145
db.RegisterMigration(migration())

components/ambient-cli/cmd/acpctl/apply/cmd.go

Lines changed: 17 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -294,24 +294,24 @@ func buildCredentialPatch(existing *sdktypes.Credential, doc resource) (map[stri
294294
patch = patch.Description(doc.Description)
295295
changed = true
296296
}
297-
if doc.URL != "" {
297+
if doc.URL != "" && doc.URL != existing.Url {
298298
patch = patch.Url(doc.URL)
299299
changed = true
300300
}
301-
if doc.Email != "" {
301+
if doc.Email != "" && doc.Email != existing.Email {
302302
patch = patch.Email(doc.Email)
303303
changed = true
304304
}
305305
token := os.ExpandEnv(doc.Token)
306-
if token != "" {
306+
if token != "" && token != existing.Token {
307307
patch = patch.Token(token)
308308
changed = true
309309
}
310-
if len(doc.Labels) > 0 {
310+
if len(doc.Labels) > 0 && marshalStringMap(doc.Labels) != existing.Labels {
311311
patch = patch.Labels(marshalStringMap(doc.Labels))
312312
changed = true
313313
}
314-
if len(doc.Annotations) > 0 {
314+
if len(doc.Annotations) > 0 && marshalStringMap(doc.Annotations) != existing.Annotations {
315315
patch = patch.Annotations(marshalStringMap(doc.Annotations))
316316
changed = true
317317
}
@@ -649,6 +649,18 @@ func strategicMerge(base, patch resource) resource {
649649
if patch.Prompt != "" {
650650
base.Prompt = patch.Prompt
651651
}
652+
if patch.Provider != "" {
653+
base.Provider = patch.Provider
654+
}
655+
if patch.Token != "" {
656+
base.Token = patch.Token
657+
}
658+
if patch.URL != "" {
659+
base.URL = patch.URL
660+
}
661+
if patch.Email != "" {
662+
base.Email = patch.Email
663+
}
652664
for k, v := range patch.Labels {
653665
if base.Labels == nil {
654666
base.Labels = make(map[string]string)

components/ambient-cli/demo-github.sh

Lines changed: 38 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -143,6 +143,14 @@ CREATED_PROJECT=""
143143
CREATED_CREDENTIAL_ID=""
144144

145145
cleanup() {
146+
if [[ -n "${NO_CLEANUP:-}" ]]; then
147+
echo
148+
yellow " NO_CLEANUP set — skipping cleanup"
149+
dim " session: ${CREATED_SESSION_ID}"
150+
dim " credential: ${CREATED_CREDENTIAL_ID}"
151+
dim " project: ${CREATED_PROJECT}"
152+
return
153+
fi
146154
echo
147155
announce "Cleanup"
148156
if [[ -n "${CREATED_SESSION_ID}" ]]; then
@@ -270,13 +278,32 @@ echo
270278
announce "4 · Create GitHub credential"
271279

272280
sep; bold "▶ Create credential: ${CRED_NAME}"; sleep "$PAUSE"
281+
_CRED_MANIFEST=$(mktemp --suffix=.yaml)
282+
cat > "${_CRED_MANIFEST}" <<'CRED_EOF'
283+
kind: Credential
284+
name: CRED_NAME_PLACEHOLDER
285+
provider: github
286+
token: $DEMO_GITHUB_PAT
287+
description: CRED_DESC_PLACEHOLDER
288+
CRED_EOF
289+
sed -i \
290+
-e "s/CRED_NAME_PLACEHOLDER/${CRED_NAME}/" \
291+
-e "s/CRED_DESC_PLACEHOLDER/GitHub PAT for demo ${RUN_ID}/" \
292+
"${_CRED_MANIFEST}"
293+
DEMO_GITHUB_PAT="${GITHUB_TOKEN_VALUE}" \
294+
"$ACPCTL" apply -f "${_CRED_MANIFEST}" 2>/dev/null
295+
rm -f "${_CRED_MANIFEST}"
273296
CRED_JSON=$(
274-
"$ACPCTL" credential create \
275-
--name "${CRED_NAME}" \
276-
--provider github \
277-
--token "${GITHUB_TOKEN_VALUE}" \
278-
--description "GitHub PAT for demo ${RUN_ID}" \
279-
-o json 2>/dev/null
297+
"$ACPCTL" get credentials -o json 2>/dev/null \
298+
| python3 -c "
299+
import sys, json
300+
data = json.load(sys.stdin)
301+
items = data.get('items', []) if isinstance(data, dict) else data
302+
for c in items:
303+
if c.get('name') == '${CRED_NAME}':
304+
print(json.dumps(c))
305+
break
306+
" 2>/dev/null
280307
)
281308
CREDENTIAL_ID=$(json_field "$CRED_JSON" "id")
282309
[[ -z "${CREDENTIAL_ID}" ]] && die "Failed to parse credential ID"
@@ -291,15 +318,15 @@ step "Verify credential visible" \
291318

292319
announce "5 · Bind credential to agent"
293320

294-
sep; bold "▶ Look up credential:reader role ID"; sleep "$PAUSE"
321+
sep; bold "▶ Look up credential:token-reader role ID"; sleep "$PAUSE"
295322
ROLES_JSON=$("$ACPCTL" get roles -o json 2>/dev/null)
296323
READER_ROLE_ID=$(
297324
echo "$ROLES_JSON" | python3 -c "
298325
import sys, json
299326
data = json.load(sys.stdin)
300327
items = data.get('items', []) if isinstance(data, dict) else data
301328
for r in items:
302-
if r.get('name') == 'credential:reader':
329+
if r.get('name') == 'credential:token-reader':
303330
print(r['id'])
304331
break
305332
" 2>/dev/null
@@ -311,13 +338,13 @@ MY_USER_ID=$(
311338
)
312339

313340
if [[ -z "${READER_ROLE_ID}" ]]; then
314-
yellow " credential:reader role not in this deployment — skipping role binding"
341+
yellow " credential:token-reader role not in this deployment — skipping role binding"
315342
dim " (credential roles are seeded by the api-server migration; redeploy may be needed)"
316343
else
317-
dim " credential:reader role ID: ${READER_ROLE_ID}"
344+
dim " credential:token-reader role ID: ${READER_ROLE_ID}"
318345
dim " my user ID: ${MY_USER_ID}"
319346

320-
sep; bold "▶ Create role-binding: credential:reader scope=agent"; sleep "$PAUSE"
347+
sep; bold "▶ Create role-binding: credential:token-reader scope=agent"; sleep "$PAUSE"
321348
"$ACPCTL" create role-binding \
322349
--user-id "${MY_USER_ID}" \
323350
--role-id "${READER_ROLE_ID}" \

components/ambient-control-plane/cmd/ambient-control-plane/main.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -176,6 +176,9 @@ func runKubeMode(ctx context.Context, cfg *config.ControlPlaneConfig) error {
176176
sessionReconcilers := createSessionReconcilers(cfg.Reconcilers, factory, kube, projectKube, provisioner, kubeReconcilerCfg, log.Logger)
177177
for _, sessionRec := range sessionReconcilers {
178178
inf.RegisterHandler("sessions", sessionRec.Reconcile)
179+
if kr, ok := sessionRec.(*reconciler.SimpleKubeReconciler); ok {
180+
kr.StartTokenRefreshLoop(ctx)
181+
}
179182
}
180183

181184
return inf.Run(ctx)

components/ambient-control-plane/internal/informer/informer.go

Lines changed: 53 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -86,9 +86,10 @@ const (
8686
)
8787

8888
type retryEvent struct {
89-
event ResourceEvent
90-
attempt int
91-
fireAt time.Time
89+
event ResourceEvent
90+
handlerIndex int
91+
attempt int
92+
fireAt time.Time
9293
}
9394

9495
type Informer struct {
@@ -162,13 +163,15 @@ func (inf *Informer) retryLoop(ctx context.Context) {
162163
case re := <-inf.retryCh:
163164
wait := time.Until(re.fireAt)
164165
if wait > 0 {
166+
timer := time.NewTimer(wait)
165167
select {
166-
case <-time.After(wait):
168+
case <-timer.C:
167169
case <-ctx.Done():
170+
timer.Stop()
168171
return
169172
}
170173
}
171-
inf.dispatchEvent(ctx, re.event, re.attempt)
174+
inf.dispatchHandler(ctx, re.event, re.handlerIndex, re.attempt)
172175
}
173176
}
174177
}
@@ -178,34 +181,53 @@ func (inf *Informer) dispatchEvent(ctx context.Context, event ResourceEvent, att
178181
handlers := inf.handlers[event.Resource]
179182
inf.mu.RUnlock()
180183

181-
for _, handler := range handlers {
184+
for i, handler := range handlers {
182185
if err := handler(ctx, event); err != nil {
183-
if attempt < retryMaxAttempts {
184-
delay := retryBaseDelay * (1 << attempt)
185-
if delay > retryMaxDelay {
186-
delay = retryMaxDelay
187-
}
188-
inf.logger.Warn().
189-
Err(err).
190-
Str("resource", event.Resource).
191-
Str("event_type", string(event.Type)).
192-
Int("attempt", attempt+1).
193-
Int("max_attempts", retryMaxAttempts).
194-
Dur("retry_in", delay).
195-
Msg("handler failed, will retry")
196-
select {
197-
case inf.retryCh <- retryEvent{event: event, attempt: attempt + 1, fireAt: time.Now().Add(delay)}:
198-
case <-ctx.Done():
199-
}
200-
} else {
201-
inf.logger.Error().
202-
Err(err).
203-
Str("resource", event.Resource).
204-
Str("event_type", string(event.Type)).
205-
Int("attempts", attempt+1).
206-
Msg("handler failed after max retries")
207-
}
186+
inf.scheduleRetry(ctx, event, i, attempt, err)
187+
}
188+
}
189+
}
190+
191+
func (inf *Informer) dispatchHandler(ctx context.Context, event ResourceEvent, handlerIndex, attempt int) {
192+
inf.mu.RLock()
193+
handlers := inf.handlers[event.Resource]
194+
inf.mu.RUnlock()
195+
196+
if handlerIndex >= len(handlers) {
197+
return
198+
}
199+
if err := handlers[handlerIndex](ctx, event); err != nil {
200+
inf.scheduleRetry(ctx, event, handlerIndex, attempt, err)
201+
}
202+
}
203+
204+
func (inf *Informer) scheduleRetry(ctx context.Context, event ResourceEvent, handlerIndex, attempt int, err error) {
205+
if attempt < retryMaxAttempts {
206+
delay := retryBaseDelay * (1 << attempt)
207+
if delay > retryMaxDelay {
208+
delay = retryMaxDelay
209+
}
210+
inf.logger.Warn().
211+
Err(err).
212+
Str("resource", event.Resource).
213+
Str("event_type", string(event.Type)).
214+
Int("handler", handlerIndex).
215+
Int("attempt", attempt+1).
216+
Int("max_attempts", retryMaxAttempts).
217+
Dur("retry_in", delay).
218+
Msg("handler failed, will retry")
219+
select {
220+
case inf.retryCh <- retryEvent{event: event, handlerIndex: handlerIndex, attempt: attempt + 1, fireAt: time.Now().Add(delay)}:
221+
case <-ctx.Done():
208222
}
223+
} else {
224+
inf.logger.Error().
225+
Err(err).
226+
Str("resource", event.Resource).
227+
Str("event_type", string(event.Type)).
228+
Int("handler", handlerIndex).
229+
Int("attempts", attempt+1).
230+
Msg("handler failed after max retries")
209231
}
210232
}
211233

components/ambient-control-plane/internal/kubeclient/kubeclient.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -220,6 +220,10 @@ func (kc *KubeClient) CreateSecret(ctx context.Context, obj *unstructured.Unstru
220220
return kc.dynamic.Resource(SecretGVR).Namespace(obj.GetNamespace()).Create(ctx, obj, metav1.CreateOptions{})
221221
}
222222

223+
func (kc *KubeClient) UpdateSecret(ctx context.Context, obj *unstructured.Unstructured) (*unstructured.Unstructured, error) {
224+
return kc.dynamic.Resource(SecretGVR).Namespace(obj.GetNamespace()).Update(ctx, obj, metav1.UpdateOptions{})
225+
}
226+
223227
func (kc *KubeClient) DeleteSecretsByLabel(ctx context.Context, namespace, labelSelector string) error {
224228
return kc.dynamic.Resource(SecretGVR).Namespace(namespace).DeleteCollection(ctx, metav1.DeleteOptions{}, metav1.ListOptions{LabelSelector: labelSelector})
225229
}

0 commit comments

Comments
 (0)