Skip to content

Commit 998114c

Browse files
committed
Refactor grpc pkg to support general grpc service.
Signed-off-by: xuezhaojun <zxue@redhat.com>
1 parent 68bb7fc commit 998114c

18 files changed

Lines changed: 845 additions & 38 deletions

File tree

grpc-interceptor-flow.md

Lines changed: 185 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,185 @@
1+
# gRPC 拦截器和认证授权流程
2+
3+
本文档描述了 gRPC 服务器中拦截器链的工作流程,特别是认证(Authentication)和授权(Authorization)的处理过程。
4+
5+
## 系统架构概览
6+
7+
```mermaid
8+
graph TB
9+
Client[gRPC Client] --> Server[gRPC Server]
10+
Server --> Chain[Interceptor Chain]
11+
Chain --> AuthN[Authentication Interceptor]
12+
Chain --> AuthZ[Authorization Interceptor]
13+
Chain --> Handler[Business Handler]
14+
15+
AuthN --> AuthN1[Authenticator 1<br/>Token Auth]
16+
AuthN --> AuthN2[Authenticator 2<br/>mTLS Auth]
17+
AuthN --> AuthN3[Authenticator N<br/>Custom Auth]
18+
19+
AuthZ --> AuthZ1[Authorizer 1<br/>SAR Check]
20+
AuthZ --> AuthZ2[Authorizer 2<br/>RBAC Check]
21+
AuthZ --> AuthZ3[Authorizer N<br/>Custom Check]
22+
```
23+
24+
## 详细请求处理流程
25+
26+
```mermaid
27+
sequenceDiagram
28+
participant C as gRPC Client
29+
participant S as gRPC Server
30+
participant IC as Interceptor Chain
31+
participant AI as Auth Interceptor
32+
participant A1 as TokenAuthenticator
33+
participant A2 as mTLSAuthenticator
34+
participant AZ as Authz Interceptor
35+
participant SAR as SARAuthorizer
36+
participant H as Business Handler
37+
participant K8S as Kubernetes API
38+
39+
C->>S: gRPC Request
40+
S->>IC: Process Request
41+
IC->>AI: Authentication Phase
42+
43+
Note over AI: OR Logic - Any authenticator success
44+
AI->>A1: Try Token Authentication
45+
A1->>K8S: TokenReview API Call
46+
K8S-->>A1: Authentication Failed
47+
A1-->>AI: Failed
48+
49+
AI->>A2: Try mTLS Authentication
50+
A2->>A2: Verify Client Certificate
51+
A2-->>AI: Success + User Context
52+
53+
Note over AI: First success, skip remaining authenticators
54+
AI->>AZ: Pass with User Context
55+
56+
Note over AZ: AND Logic - All authorizers must pass
57+
AZ->>SAR: Check Permissions
58+
SAR->>K8S: SubjectAccessReview API Call
59+
K8S-->>SAR: Permission Granted
60+
SAR-->>AZ: Authorized
61+
62+
Note over AZ: All authorizers passed
63+
AZ->>H: Execute Business Logic
64+
H-->>AZ: Response
65+
AZ-->>AI: Response
66+
AI-->>IC: Response
67+
IC-->>S: Response
68+
S-->>C: gRPC Response
69+
```
70+
71+
## 拦截器链配置
72+
73+
```mermaid
74+
graph LR
75+
subgraph "Server Configuration"
76+
SO[Server Options]
77+
SO --> CUI[ChainUnaryInterceptor]
78+
SO --> CSI[ChainStreamInterceptor]
79+
end
80+
81+
subgraph "Unary Interceptor Chain"
82+
CUI --> AI1[newAuthnUnaryInterceptor]
83+
AI1 --> AZ1[newAuthzUnaryInterceptor]
84+
AZ1 --> BH1[Business Handler]
85+
end
86+
87+
subgraph "Stream Interceptor Chain"
88+
CSI --> AI2[newAuthnStreamInterceptor]
89+
AI2 --> AZ2[newAuthzStreamInterceptor]
90+
AZ2 --> BH2[Business Handler]
91+
end
92+
93+
subgraph "Authenticators"
94+
AI1 --> TA[TokenAuthenticator]
95+
AI1 --> MA[mTLSAuthenticator]
96+
AI2 --> TA
97+
AI2 --> MA
98+
end
99+
100+
subgraph "Authorizers"
101+
AZ1 --> SARA[SARAuthorizer]
102+
AZ2 --> SARA
103+
end
104+
```
105+
106+
## 认证器逻辑(OR 关系)
107+
108+
```mermaid
109+
flowchart TD
110+
Start([Request Arrives]) --> AuthLoop{For each Authenticator}
111+
AuthLoop --> TryAuth[Try Authentication]
112+
TryAuth --> AuthSuccess{Success?}
113+
AuthSuccess -->|Yes| SetContext[Set User Context]
114+
SetContext --> CallHandler[Call Next Handler]
115+
AuthSuccess -->|No| NextAuth{More Authenticators?}
116+
NextAuth -->|Yes| AuthLoop
117+
NextAuth -->|No| AuthFailed[Return Auth Error]
118+
CallHandler --> End([Continue to Authorization])
119+
AuthFailed --> ErrorEnd([Return Error])
120+
121+
style SetContext fill:#90EE90
122+
style AuthFailed fill:#FFB6C1
123+
style CallHandler fill:#87CEEB
124+
```
125+
126+
## 授权器逻辑(AND 关系)
127+
128+
```mermaid
129+
flowchart TD
130+
Start([From Authentication]) --> AuthzLoop{For each Authorizer}
131+
AuthzLoop --> TryAuthz[Try Authorization]
132+
TryAuthz --> AuthzSuccess{Success?}
133+
AuthzSuccess -->|No| AuthzFailed[Return Authz Error]
134+
AuthzSuccess -->|Yes| NextAuthz{More Authorizers?}
135+
NextAuthz -->|Yes| AuthzLoop
136+
NextAuthz -->|No| AllPassed[All Authorizers Passed]
137+
AllPassed --> CallHandler[Call Business Handler]
138+
CallHandler --> End([Return Response])
139+
AuthzFailed --> ErrorEnd([Return Error])
140+
141+
style AllPassed fill:#90EE90
142+
style AuthzFailed fill:#FFB6C1
143+
style CallHandler fill:#87CEEB
144+
```
145+
146+
## 性能影响分析
147+
148+
```mermaid
149+
graph TB
150+
subgraph "Performance Impact"
151+
PI[Performance Impact]
152+
PI --> FO[Framework Overhead<br/>~1μs]
153+
PI --> AO[Auth Overhead<br/>1-10ms]
154+
PI --> NO[Network Overhead<br/>K8s API Calls]
155+
end
156+
157+
subgraph "Optimization Strategies"
158+
OS[Optimization]
159+
OS --> Cache[Token/SAR Caching]
160+
OS --> Pool[Connection Pooling]
161+
OS --> Batch[Batch Processing]
162+
OS --> Async[Async Prewarming]
163+
end
164+
165+
FO -.->|Minimal| Cache
166+
AO -.->|Significant| Cache
167+
NO -.->|Major| Pool
168+
```
169+
170+
## 关键特性总结
171+
172+
### 认证(Authentication)
173+
- **逻辑关系**: OR(任意一个成功即可)
174+
- **早期退出**: 第一个成功的认证器会立即返回
175+
- **上下文传递**: 成功的认证器会在 context 中设置用户信息
176+
177+
### 授权(Authorization)
178+
- **逻辑关系**: AND(所有授权器都必须通过)
179+
- **严格检查**: 任意一个授权器失败都会导致请求被拒绝
180+
- **依赖认证**: 需要认证阶段提供的用户上下文
181+
182+
### 性能考虑
183+
- **拦截器框架开销**: 微秒级,可忽略
184+
- **主要瓶颈**: Kubernetes API 调用(TokenReview, SubjectAccessReview)
185+
- **优化方案**: 缓存、连接池、批处理等策略

pkg/cloudevents/server/grpc/authz/kube/sar.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,8 @@ import (
1616
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/lease"
1717
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/payload"
1818
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
19-
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/authn"
2019
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/authz"
20+
"open-cluster-management.io/sdk-go/pkg/grpc/authn"
2121
)
2222

2323
type SARAuthorizer struct {

pkg/cloudevents/server/grpc/authz/kube/sar_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@ import (
1717
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/lease"
1818
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/payload"
1919
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
20-
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/authn"
20+
"open-cluster-management.io/sdk-go/pkg/grpc/authn"
2121
)
2222

2323
func TestSARAuthorize(t *testing.T) {
Lines changed: 5 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,10 @@ import (
44
"context"
55
"encoding/json"
66
"fmt"
7-
"net"
87
"sync"
98
"time"
109

1110
"k8s.io/apimachinery/pkg/api/errors"
12-
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
1311

1412
cloudevents "github.com/cloudevents/sdk-go/v2"
1513
"github.com/cloudevents/sdk-go/v2/binding"
@@ -44,23 +42,24 @@ var _ server.AgentEventServer = &GRPCBroker{}
4442
// It broadcasts resource spec to agents and listens for resource status updates from them.
4543
type GRPCBroker struct {
4644
pbv1.UnimplementedCloudEventServiceServer
47-
grpcServer *grpc.Server
4845
services map[types.CloudEventsDataType]server.Service
4946
subscribers map[string]*subscriber // registered subscribers
5047
mu sync.RWMutex
5148
}
5249

5350
// NewGRPCBroker creates a new gRPC broker with the given gRPC server.
54-
func NewGRPCBroker(srv *grpc.Server) server.AgentEventServer {
51+
func NewGRPCBroker() *GRPCBroker {
5552
broker := &GRPCBroker{
56-
grpcServer: srv,
5753
subscribers: make(map[string]*subscriber),
5854
services: make(map[types.CloudEventsDataType]server.Service),
5955
}
60-
pbv1.RegisterCloudEventServiceServer(broker.grpcServer, broker)
6156
return broker
6257
}
6358

59+
func (bkr *GRPCBroker) Register(grpcServer *grpc.Server) {
60+
pbv1.RegisterCloudEventServiceServer(grpcServer, bkr)
61+
}
62+
6463
func (bkr *GRPCBroker) RegisterService(t types.CloudEventsDataType, service server.Service) {
6564
bkr.services[t] = service
6665
service.RegisterHandler(bkr)
@@ -78,26 +77,6 @@ func (bkr *GRPCBroker) Subscribers() sets.Set[string] {
7877
return subscribers
7978
}
8079

81-
// Start starts the gRPC broker at the given address
82-
func (bkr *GRPCBroker) Start(ctx context.Context, addr string) {
83-
logger := klog.FromContext(ctx)
84-
logger.Info("Starting gRPC broker at addr", "addr", addr)
85-
lis, err := net.Listen("tcp", addr)
86-
if err != nil {
87-
utilruntime.Must(fmt.Errorf("failed to listen: %v", err))
88-
}
89-
go func() {
90-
if err := bkr.grpcServer.Serve(lis); err != nil {
91-
utilruntime.Must(fmt.Errorf("failed to serve gRPC broker: %v", err))
92-
}
93-
}()
94-
95-
// wait until context is canceled
96-
<-ctx.Done()
97-
klog.Infof("Shutting down gRPC broker")
98-
bkr.grpcServer.GracefulStop()
99-
}
100-
10180
// Publish in stub implementation for agent publish resource status.
10281
func (bkr *GRPCBroker) Publish(ctx context.Context, pubReq *pbv1.PublishRequest) (*emptypb.Empty, error) {
10382
logger := klog.FromContext(ctx)

pkg/cloudevents/server/grpc/server_test.go renamed to pkg/cloudevents/server/grpc/broker_test.go

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ package grpc
33
import (
44
"context"
55
"errors"
6+
"net"
67
"testing"
78

89
cloudevents "github.com/cloudevents/sdk-go/v2"
@@ -59,14 +60,24 @@ func (s *testService) create(evt *cloudevents.Event) error {
5960
func TestServer(t *testing.T) {
6061
grpcServerOptions := []grpc.ServerOption{}
6162
grpcServer := grpc.NewServer(grpcServerOptions...)
62-
grpcEventServer := NewGRPCBroker(grpcServer)
63+
grpcEventServer := NewGRPCBroker()
64+
grpcEventServer.Register(grpcServer)
6365

6466
svc := &testService{evts: make(map[string]*cloudevents.Event)}
6567
grpcEventServer.RegisterService(dataType, svc)
6668

6769
ctx, cancel := context.WithCancel(context.Background())
6870
defer cancel()
69-
go grpcEventServer.Start(ctx, ":8888")
71+
72+
lis, err := net.Listen("tcp", ":8888")
73+
if err != nil {
74+
t.Fatalf("failed to listen: %v", err)
75+
}
76+
go func() {
77+
if err := grpcServer.Serve(lis); err != nil {
78+
t.Errorf("failed to serve: %v", err)
79+
}
80+
}()
7081

7182
grpcClientOptions := grpccli.NewGRPCOptions()
7283
grpcClientOptions.Dialer = &grpccli.GRPCDialer{URL: "localhost:8888"}

pkg/cloudevents/server/grpc/options/server.go

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -18,8 +18,8 @@ import (
1818
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
1919
"open-cluster-management.io/sdk-go/pkg/cloudevents/server"
2020
grpcserver "open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc"
21-
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/authn"
2221
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/authz"
22+
"open-cluster-management.io/sdk-go/pkg/grpc/authn"
2323
)
2424

2525
// PreStartHook is an interface to start hook before grpc server is started.
@@ -116,7 +116,8 @@ func (s *Server) Run(ctx context.Context) error {
116116
newAuthzStreamInterceptor(s.authorizers...)))
117117

118118
grpcServer := grpc.NewServer(grpcServerOptions...)
119-
grpcEventServer := grpcserver.NewGRPCBroker(grpcServer)
119+
grpcEventServer := grpcserver.NewGRPCBroker()
120+
grpcEventServer.Register(grpcServer)
120121

121122
for t, service := range s.services {
122123
grpcEventServer.RegisterService(t, service)
@@ -127,7 +128,7 @@ func (s *Server) Run(ctx context.Context) error {
127128
hook.Run(ctx)
128129
}
129130

130-
go grpcEventServer.Start(ctx, ":"+s.options.ServerBindPort)
131+
// go grpcEventServer.Start(ctx, ":"+s.options.ServerBindPort)
131132
<-ctx.Done()
132133
return nil
133134
}

pkg/cloudevents/server/grpc/options/server_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ import (
77
"testing"
88

99
certutil "k8s.io/client-go/util/cert"
10-
"open-cluster-management.io/sdk-go/pkg/cloudevents/server/grpc/authn"
10+
"open-cluster-management.io/sdk-go/pkg/grpc/authn"
1111
)
1212

1313
type testHook struct{}

pkg/cloudevents/server/interface.go

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,9 +17,6 @@ type AgentEventServer interface {
1717
// RegisterService registers a backend service with a certain data type.
1818
RegisterService(t types.CloudEventsDataType, service Service)
1919

20-
// Start initiates the EventServer to listen to agents.
21-
Start(ctx context.Context, addr string)
22-
2320
// Subscribers returns all current subscribers who subscribe to this server.
2421
Subscribers() sets.Set[string]
2522
}
File renamed without changes.

0 commit comments

Comments
 (0)