Skip to content
Santekno.com | Level Up Your Engineering Skills
ID
📖 0%
11 Jul 2025 · 5 mnt baca ·Artikel 33 / 110
Go

33 Menulis Stream Interceptor Sendiri

IH
Ihsan Arif
Penulis di Santekno · Backend Engineer

Jika Anda pernah membangun aplikasi berbasis gRPC, pasti sudah akrab dengan konsep interceptor. Di ekosistem gRPC, interceptor menempati peran layaknya middleware di framework web seperti Express (Node.js) atau Gin (Go). Namun, ada perbedaan fundamental: selain interceptor untuk unary RPC, gRPC punya mekanisme khusus untuk menangani streaming RPC, yaitu Stream Interceptor. Kali ini, kita akan membedah bagaimana menulis stream interceptor sendiri—mulai dari konsep, perancangan, hingga implementasi dan simulasi penggunaannya.

Kenapa Stream Interceptor Penting?

gRPC mendukung beberapa jenis RPC:

  • Unary: Satu request, satu response.
  • Server streaming: Satu request, banyak response.
  • Client streaming: Banyak request, satu response.
  • Bidirectional streaming: Banyak request, banyak response secara simultan.

Saat kita menghadapi streaming—baik dari sisi client, server, ataupun keduanya—middleware standar tidak cukup. Di sinilah stream interceptor menjadi penting karena:

  • Dapat melakukan validasi atau manipulasi pada every message (bukan hanya di initial request/response).
  • Bisa logging, monitoring, pembatasan rate per message, manipulasi konteks, hingga error handling pada aliran data.
  • Membuka pintu ke berbagai fitur seperti metrik per stream, tracing, dsb.

Anatomy: Bagaimana Stream Interceptor Bekerja di gRPC Go?

Mari kita lihat signature–nya (versi Go):

go
1type StreamServerInterceptor func(
2    srv interface{},
3    ss grpc.ServerStream,
4    info *grpc.StreamServerInfo,
5    handler grpc.StreamHandler,
6) error

Argumen penting di sini:

  • srv: implementasi service.
  • ss: objek stream (untuk kirim/terima message).
  • info: informasi stream (nama RPC, metode, dsb).
  • handler: fungsi “next”, mirip middleware handler di web.

Gambaran alur interceptor:

MERMAID
flowchart LR
    Client--Stream Msgs--> Server[gRPC Server]
    Server-->|Incoming Context| Interceptor
    Interceptor-->|Manipulate or Pass| Handler
    Handler-->|Responses| Interceptor
    Interceptor-->|Manipulate or Pass| Client

Studi Kasus: Membuat Logging Stream Interceptor

Mari implementasikan logging di setiap pesan pada stream server-side. Kita ingin tahu kapan client mengirim/menerima pesan, serta waktu proses setiap pesan.

1. Membungkus ServerStream

Pada gRPC Go, grpc.ServerStream adalah interface untuk stream. Untuk intercept pesan, kita perlu membungkus metode RecvMsg dan SendMsg.

go
 1type wrappedStream struct {
 2    grpc.ServerStream
 3}
 4
 5func (w *wrappedStream) RecvMsg(m interface{}) error {
 6    err := w.ServerStream.RecvMsg(m)
 7    log.Printf("[Interceptor] Received message: %#v, Error: %v", m, err)
 8    return err
 9}
10
11func (w *wrappedStream) SendMsg(m interface{}) error {
12    log.Printf("[Interceptor] Sending message: %#v", m)
13    return w.ServerStream.SendMsg(m)
14}

2. Menulis Interceptor

Lalu, kita siapkan interceptor yang menggunakan wrappedStream:

go
 1func LoggingStreamInterceptor(
 2    srv interface{},
 3    ss grpc.ServerStream,
 4    info *grpc.StreamServerInfo,
 5    handler grpc.StreamHandler,
 6) error {
 7    log.Printf("[Interceptor] Stream started: %s", info.FullMethod)
 8    err := handler(srv, &wrappedStream{ss})
 9    if err != nil {
10        log.Printf("[Interceptor] Stream error: %v", err)
11    } else {
12        log.Printf("[Interceptor] Stream completed successfully")
13    }
14    return err
15}

3. Register Interceptor di Server

Ketika membangun server:

go
1grpcServer := grpc.NewServer(
2    grpc.StreamInterceptor(LoggingStreamInterceptor),
3)

4. Simulasi: Service Sederhana

Misal, service streaming upload:

proto
1service FileService {
2    rpc Upload(stream FileChunk) returns (UploadStatus);
3}
4message FileChunk {
5    bytes data = 1;
6}
7message UploadStatus {
8    string message = 1;
9}

Handler:

go
 1func (s *fileServiceServer) Upload(stream pb.FileService_UploadServer) error {
 2    var total int
 3    for {
 4        chunk, err := stream.Recv()
 5        if err == io.EOF {
 6            return stream.SendAndClose(&pb.UploadStatus{Message: "Received"})
 7        }
 8        if err != nil {
 9            return err
10        }
11        total += len(chunk.Data)
12        // Processing chunk...
13    }
14}

Ketika client streaming file, setiap chunk akan “di-log” oleh interceptor.

5. Output Simulasi

ActionMessageData
Stream started/FileService/Upload-
Received messageFileChunk{data:…}len: 1024 bytes
Received messageFileChunk{data:…}len: 483 bytes
Sending messageUploadStatus{…}message: “Received”
Stream completed--

Log real-time akan terlihat di konsol server.

Memodifikasi Interceptor: Menambah Fitur Rate Limiting

Ayo kembangkan interceptor kita: implementasi rate limiter berbasis pesan per detik.

1. Rate Limiter Sederhana

go
 1type rateLimitStream struct {
 2    grpc.ServerStream
 3    tokens chan struct{}
 4}
 5
 6func NewRateLimitStream(ss grpc.ServerStream, rate int) *rateLimitStream {
 7    rl := &rateLimitStream{
 8        ServerStream: ss,
 9        tokens:       make(chan struct{}, rate),
10    }
11    // refill tokens setiap detik
12    go func() {
13        ticker := time.NewTicker(time.Second)
14        defer ticker.Stop()
15        for {
16            <-ticker.C
17            for i := 0; i < rate; i++ {
18                select {
19                case rl.tokens <- struct{}{}:
20                default:
21                }
22            }
23        }
24    }()
25    return rl
26}
27
28func (w *rateLimitStream) RecvMsg(m interface{}) error {
29    <-w.tokens // akan blocking jika quota habis
30    return w.ServerStream.RecvMsg(m)
31}

2. Integrasi ke Interceptor

go
 1func RateLimitStreamInterceptor(rate int) grpc.StreamServerInterceptor {
 2    return func(
 3        srv interface{},
 4        ss grpc.ServerStream,
 5        info *grpc.StreamServerInfo,
 6        handler grpc.StreamHandler,
 7    ) error {
 8        limited := NewRateLimitStream(ss, rate)
 9        return handler(srv, limited)
10    }
11}

3. Penggunaan di Server

go
1grpcServer := grpc.NewServer(
2    grpc.StreamInterceptor(RateLimitStreamInterceptor(5)), // max 5 pesan/detik
3)

Ringkasan: Stream Interceptor adalah “Middleware on Steroids”

Membuat stream interceptor sendiri di gRPC sangat powerful untuk kebutuhan observabilitas, security, atau traffic shaping. Anda bisa membangun fitur logging granular, tracing, rate limiting, custom authentication, quota enforcement, dan berbagai kasus lain—semua dengan single point of change tanpa mengutak-atik handler service.

👇 Tips produksi:

  • Jangan lupa handle context cancellation/error.
  • Jika ingin multi-interceptor, gunakan chaining atau integrasi dengan third party (misal, go-grpc-middleware ).
  • Logging hati-hati, jangan log data sensitif.
  • Ukur performa, overhead interceptor bisa memengaruhi throughput stream.

Kapan Anda akan menulis stream interceptor sendiri? Jika ingin kontrol penuh atas semua data yang streaming masuk/keluar server, inilah caranya. Selamat bereksperimen!


Referensi:


Jika artikel ini bermanfaat, silakan bookmark dan share ke rekan developer Anda. Sampai jumpa di eksperimen berikutnya! 🚀

Artikel Terkait

💬 Komentar