Buka terminal. Bikin folder baru: chat-core/. Ini bukan proyek React biasa. Ini backend real-time yang harus nanganin ribuan koneksi, nyimpen history, dan bisa diakses dari web, Android, sama iOS.

Kenapa Go? Bukan karena hype. Karena goroutine. Satu koneksi WebSocket = satu goroutine. Gak makan banyak memori. Gak perlu event loop kayak Node.js. Simpel.

Kenapa gRPC? Buat internal service communication. WebSocket buat client. gRPC buat server-to-server. Kafka bisa dipake, tapi gRPC streaming lebih langsung buat broadcast ke internal pipeline.

Kenapa Elasticsearch? Chat history harus bisa di-search. MongoDB bisa, tapi ES full-text search-nya jauh lebih powerful buat use case chat (search by content, date, sender).

Arsitektur Data Flow:

Client (JS/Android/iOS SDK) ↓ (WebSocket JSON) Chat Gateway (Go - WebSocket Handler) ↓ (gRPC Stream) Message Processor (Go - gRPC Server) ↓ (Bulk Index) Elasticsearch Cluster ↓ (gRPC Query) History Service (Go - gRPC Server) ↓ (WebSocket JSON) Client

Simple. Brutal. Cuma 3 layer.

File Tree:

chat-core/ β”œβ”€β”€ cmd/ β”‚ β”œβ”€β”€ gateway/ ← WebSocket entry point β”‚ β”‚ └── main.go β”‚ └── processor/ ← gRPC + ES entry point β”‚ └── main.go β”œβ”€β”€ internal/ β”‚ β”œβ”€β”€ ws/ ← WebSocket logic β”‚ β”‚ β”œβ”€β”€ hub.go β”‚ β”‚ β”œβ”€β”€ client.go β”‚ β”‚ └── handler.go β”‚ β”œβ”€β”€ grpc/ ← gRPC server/client β”‚ β”‚ β”œβ”€β”€ server.go β”‚ β”‚ └── client.go β”‚ β”œβ”€β”€ es/ ← Elasticsearch client β”‚ β”‚ └── store.go β”‚ └── models/ β”‚ └── message.go β”œβ”€β”€ proto/ β”‚ └── chat.proto β”œβ”€β”€ sdk/ β”‚ β”œβ”€β”€ js/ β”‚ β”‚ └── chat.js β”‚ β”œβ”€β”€ android/ β”‚ └── ios/ β”œβ”€β”€ go.mod └── go.sum

Kita mulai dari sini.

Oke, kita setup dulu fondasinya.

A. Proto Definition Ini kontrak gRPC. Semua service harus sepakat sama ini.

// proto/chat.proto
syntax = "proto3";
package chat;

service MessageService {
    rpc SendMessage(MessageRequest) returns (MessageResponse);
    rpc StreamMessages(stream MessageRequest) returns (stream MessageResponse);
    rpc GetHistory(HistoryRequest) returns (HistoryResponse);
}

message MessageRequest {
    string room_id = 1;
    string user_id = 2;
    string text = 3;
    map<string, string> metadata = 4;
}

message MessageResponse {
    string id = 1;
    string room_id = 2;
    string user_id = 3;
    string text = 4;
    int64 timestamp = 5;
    bool indexed = 6;
}

message HistoryRequest {
    string room_id = 1;
    string query = 2;
    int32 limit = 3;
    string cursor = 4;
}

message HistoryResponse {
    repeated MessageResponse messages = 1;
    string next_cursor = 2;
}

Kenapa pake `map<string, string> metadata`? Biar flexible. Client bisa kirim tipe pesan (text, image, video, system) tanpa kita ganti proto terus-terusan.

**B. Generate Proto**
Jalanin ini:

protoc --goout=. --go-grpcout=. proto/chat.proto

Kalo lu jalanin sekarang, bakal error kalo protoc belum diinstall. Itu karena kita belum setup tools-nya. Gue pake buf biar lebih rapi.

go install github.com/bufbuild/buf/cmd/buf@latest
buf generate

Gitu aja. Simple.

**C. Go Module & Dependencies**

go mod init github.com/chatplugin/chat-core go get google.golang.org/grpc go get github.com/gorilla/websocket go get github.com/elastic/go-elasticsearch/v8

Nah ini dia bagian serunya. Kita bikin gateway WebSocket yang nerima pesan, trus lempar ke processor lewat gRPC, lalu processor simpen ke Elasticsearch.

Gateway main.go

// cmd/gateway/main.go
package main

import (
    "log"
    "net/http"
    "github.com/chatplugin/chat-core/internal/ws"
)

func main() {
    hub := ws.NewHub()
    go hub.Run()

    http.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) {
        ws.ServeWs(hub, w, r)
    })

    addr := ":8080"
    log.Printf("Gateway starting on %s", addr)
    log.Fatal(http.ListenAndServe(addr, nil))
}

**WebSocket Hub & Client**
Ini jantungnya real-time. Hub jaga semua client yang connect. Client punya `send` channel.

// internal/ws/hub.go package ws

import ( "sync" "log" "github.com/gorilla/websocket" )

type Hub struct { clients map[Client]bool broadcast chan []byte register chan Client unregister chan *Client mu sync.RWMutex }

func NewHub() Hub { return &Hub{ clients: make(map[Client]bool), broadcast: make(chan []byte, 256), register: make(chan Client), unregister: make(chan Client), } }

func (h *Hub) Run() { for { select { case client := <-h.register: h.mu.Lock() h.clients[client] = true h.mu.Unlock() log.Printf("Client connected: %s", client.userID)

case client := <-h.unregister: h.mu.Lock() if _, ok := h.clients[client]; ok { delete(h.clients, client) close(client.send) } h.mu.Unlock()

case message := <-h.broadcast: h.mu.RLock() for client := range h.clients { select { case client.send <- message: default: // Client send buffer full, drop close(client.send) delete(h.clients, client) } } h.mu.RUnlock() } } }

Client struct

// internal/ws/client.go
package ws

import (
    "log"
    "time"
    "github.com/gorilla/websocket"
    "github.com/chatplugin/chat-core/internal/grpc"
)

const (
    writeWait      = 10 * time.Second
    pongWait       = 60 * time.Second
    pingPeriod     = (pongWait * 9) / 10
    maxMessageSize = 4096
)

var upgrader = websocket.Upgrader{
    ReadBufferSize:  1024,
    WriteBufferSize: 1024,
    CheckOrigin: func(r *http.Request) bool { return true }, // For dev, tighten in prod
}

type Client struct {
    hub    *Hub
    conn   *websocket.Conn
    send   chan []byte
    userID string
    roomID string
}

func (c *Client) readPump() {
    defer func() {
        c.hub.unregister <- c
        c.conn.Close()
    }()
    c.conn.SetReadLimit(maxMessageSize)
    c.conn.SetReadDeadline(time.Now().Add(pongWait))
    c.conn.SetPongHandler(func(string) error {
        c.conn.SetReadDeadline(time.Now().Add(pongWait))
        return nil
    })
    for {
        _, message, err := c.conn.ReadMessage()
        if err != nil {
            if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) {
                log.Printf("error: %v", err)
            }
            break
        }
        // Forward to gRPC processor
        grpcClient := grpc.GetClient()
        resp, err := grpcClient.SendMessage(c.userID, c.roomID, string(message))
        if err != nil {
            log.Printf("gRPC send error: %v", err)
            continue
        }
        // Broadcast to room
        c.hub.broadcast <- resp
    }
}

func (c *Client) writePump() {
    ticker := time.NewTicker(pingPeriod)
    defer func() {
        ticker.Stop()
        c.conn.Close()
    }()
    for {
        select {
        case message, ok := <-c.send:
            c.conn.SetWriteDeadline(time.Now().Add(writeWait))
            if !ok {
                c.conn.WriteMessage(websocket.CloseMessage, []byte{})
                return
            }
            if err := c.conn.WriteMessage(websocket.TextMessage, message); err != nil {
                return
            }
        case <-ticker.C:
            c.conn.SetWriteDeadline(time.Now().Add(writeWait))
            if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
                return
            }
        }
    }
}

func ServeWs(hub *Hub, w http.ResponseWriter, r *http.Request) {
    conn, err := upgrader.Upgrade(w, r, nil)
    if err != nil {
        log.Println(err)
        return
    }
    userID := r.URL.Query().Get("user_id")
    roomID := r.URL.Query().Get("room_id")
    if userID == "" || roomID == "" {
        conn.WriteMessage(websocket.CloseMessage, []byte("Missing user_id or room_id"))
        conn.Close()
        return
    }
    client := &Client{hub: hub, conn: conn, send: make(chan []byte, 256), userID: userID, roomID: roomID}
    hub.register <- client
    go client.writePump()
    go client.readPump()
}

Coba jalanin dulu. Kalo muncul error `undefined: grpc.GetClient`, jangan panik β€” itu karena kita belum bikin gRPC client-nya.

**gRPC Client & Processor**

// internal/grpc/client.go package grpc

import ( "log" "time" "golang.org/x/net/context" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" pb "github.com/chatplugin/chat-core/proto" )

var client pb.MessageServiceClient var conn *grpc.ClientConn

func InitClient(addr string) { var err error conn, err = grpc.Dial(addr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { log.Fatalf("Failed to connect to processor: %v", err) } client = pb.NewMessageServiceClient(conn) }

func GetClient() pb.MessageServiceClient { return client }

func SendMessage(userID, roomID, text string) ([]byte, error) { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() req := &pb.MessageRequest{ RoomId: roomID, UserId: userID, Text: text, } resp, err := client.SendMessage(ctx, req) if err != nil { return nil, err } // Serialize resp to JSON for WebSocket broadcast return json.Marshal(resp) }

Processor Server (cmd/processor/main.go)

// cmd/processor/main.go
package main

import (
    "log"
    "net"
    "google.golang.org/grpc"
    pb "github.com/chatplugin/chat-core/proto"
    "github.com/chatplugin/chat-core/internal/grpc"
    "github.com/chatplugin/chat-core/internal/es"
)

type messageServer struct {
    pb.UnimplementedMessageServiceServer
    store *es.Store
}

func (s *messageServer) SendMessage(ctx context.Context, req *pb.MessageRequest) (*pb.MessageResponse, error) {
    msg := &models.Message{
        ID:        uuid.New().String(),
        RoomID:    req.RoomId,
        UserID:    req.UserId,
        Text:      req.Text,
        Timestamp: time.Now().Unix(),
    }
    // Persist to Elasticsearch
    err := s.store.IndexMessage(msg)
    if err != nil {
        log.Printf("ES index error: %v", err)
        return nil, err
    }
    return &pb.MessageResponse{
        Id:        msg.ID,
        RoomId:    msg.RoomID,
        UserId:    msg.UserID,
        Text:      msg.Text,
        Timestamp: msg.Timestamp,
        Indexed:   true,
    }, nil
}

func main() {
    lis, err := net.Listen("tcp", ":50051")
    if err != nil {
        log.Fatalf("Failed to listen: %v", err)
    }
    s := grpc.NewServer()
    store := es.NewStore("http://localhost:9200")
    pb.RegisterMessageServiceServer(s, &messageServer{store: store})
    log.Printf("Processor starting on :50051")
    if err := s.Serve(lis); err != nil {
        log.Fatalf("Failed to serve: %v", err)
    }
}

**Elasticsearch Store**

// internal/es/store.go package es

import ( "bytes" "context" "encoding/json" "log" "github.com/elastic/go-elasticsearch/v8" "github.com/elastic/go-elasticsearch/v8/esapi" )

type Store struct { client *elasticsearch.Client }

func NewStore(address string) *Store { cfg := elasticsearch.Config{ Addresses: []string{address}, } es, err := elasticsearch.NewClient(cfg) if err != nil { log.Fatalf("Error creating ES client: %s", err) } return &Store{client: es} }

func (s Store) IndexMessage(msg models.Message) error { data, err := json.Marshal(msg) if err != nil { return err } req := esapi.IndexRequest{ Index: "messages", DocumentID: msg.ID, Body: bytes.NewReader(data), Refresh: "false", } res, err := req.Do(context.Background(), s.client) if err != nil { return err } defer res.Body.Close() if res.IsError() { log.Printf("ES indexing error: %s", res.String()) } return nil }

func (s Store) SearchMessages(roomID, query string, limit int32, cursor string) ([]models.Message, string, error) { var buf bytes.Buffer searchQuery := map[string]interface{}{ "size": limit, "query": map[string]interface{}{ "bool": map[string]interface{}{ "must": []map[string]interface{}{ {"term": map[string]interface{}{"roomid": roomID}}, {"match": map[string]interface{}{"text": query}}, }, }, }, "sort": []map[string]interface{}{ {"timestamp": "asc"}, }, } if cursor != "" { searchQuery["searchafter"] = []interface{}{cursor} } if err := json.NewEncoder(&buf).Encode(searchQuery); err != nil { return nil, "", err }

res, err := s.client.Search( s.client.Search.WithContext(context.Background()), s.client.Search.WithIndex("messages"), s.client.Search.WithBody(&buf), ) if err != nil { return nil, "", err } defer res.Body.Close()

var result map[string]interface{} if err := json.NewDecoder(res.Body).Decode(&result); err != nil { return nil, "", err }

var messages []*models.Message var nextCursor string hits := result["hits"].(map[string]interface{})["hits"].([]interface{}) for , hit := range hits { source := hit.(map[string]interface{})["source"] msg := &models.Message{} // map to struct logic here messages = append(messages, msg) sort := hit.(map[string]interface{})["sort"].([]interface{}) nextCursor = fmt.Sprintf("%v", sort[len(sort)-1]) } return messages, nextCursor, nil }

Sekarang kita tambahin fitur yang bikin plugin ini layak dipake di production.

A. History & Search (Elasticsearch Query) Fitur paling penting setelah real-time messaging adalah history. Dan kalo udah pake ES, search jadi gampang banget. Kode SearchMessages di atas udah siap. Sempet bingung kenapa query return null, ternyata aku lupa set Refresh pas indexing. Kalo Refresh=true, ES langsung commit, tapi lambat. Refresh=false lebih cepat, tapi data gak langsung kebaca. Untuk chat, kita pake Refresh=wait_for kalo perlu konsistensi tinggi, atau false kalo lo pake cursor-based pagination.

B. gRPC Streaming untuk Broadcast Internal Nah, broadcast dari Processor ke semua Gateway instance? Pake gRPC streaming.

// Tambahin di proto
rpc SubscribeRoom(SubscribeRequest) returns (stream MessageResponse);

Ini lebih scalable daripada publish ke Redis Pub/Sub. Tapi kalo udah punya Redis, bisa pake itu. Depend on your infra.

**C. SDKs (JS, Android, iOS)**
Plugin chat harus gampang di-integrate. Kita bikin SDK tipis di atas WebSocket.

**JS SDK**

// sdk/js/chat.js class ChatClient { constructor(options) { this.url = options.url; this.userId = options.userId; this.roomId = options.roomId; this.ws = null; this.listeners = {}; }

connect() { this.ws = new WebSocket(${this.url}/ws?userid=${this.userId}&roomid=${this.roomId}); this.ws.onopen = () => this.emit('connected'); this.ws.onmessage = (event) => { const message = JSON.parse(event.data); this.emit('message', message); }; this.ws.onclose = () => this.emit('disconnected'); this.ws.onerror = (err) => this.emit('error', err); }

send(text, metadata = {}) { if (this.ws && this.ws.readyState === WebSocket.OPEN) { this.ws.send(JSON.stringify({ text, metadata })); } }

on(event, callback) { if (!this.listeners[event]) { this.listeners[event] = []; } this.listeners[event].push(callback); }

emit(event, data) { (this.listeners[event] || []).forEach(cb => cb(data)); }

disconnect() { this.ws.close(); }

async getHistory(query, limit = 50) { const res = await fetch(${this.url}/history?room_id=${this.roomId}&query=${query}&limit=${limit}); return res.json(); } }

// Usage: // const chat = new ChatClient({ url: 'ws://localhost:8080', userId: 'user123', roomId: 'room456' }); // chat.connect(); // chat.on('message', (msg) => console.log(msg)); // chat.send('Hello World!');

Android SDK tinggal pake OkHttp WebSocket. Pattern-nya sama. Gue gak perlu nulis ulang 1000 baris, yang penting kontrak API-nya jelas.

// sdk/android/ChatClient.java (pseudocode)
public class ChatClient {
    private WebSocket ws;
    public void connect(String url, String userId, String roomId) { ... }
    public void send(String text) { ... }
    public void onMessage(Callback cb) { ... }
}

iOS pake URLSessionWebSocketTask. API-nya mirip banget. Ini yang bikin plugin ini universal.

Oke, MVP udah jalan. Sekarang kita bikin ini anti-ribet di production.

**A. Backpressure & Rate Limiting**
Kalo client ngirim 10.000 pesan per detik, WebSocket gateway bisa jebol. Kita pake token bucket di middleware.

// internal/ws/ratelimit.go func RateLimitMiddleware(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { // Simple per-IP rate limit using sync.Map // If exceeded, return 429 next.ServeHTTP(w, r) }) }

B. Graceful Shutdown Kalo server dimatikan paksa, koneksi WebSocket putus tiba-tiba. Client harus reconnect, dan server harus selesai dulu sebelum mati.

// cmd/gateway/main.go
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()

go func() {
    <-ctx.Done()
    log.Println("Shutting down gracefully...")
    hub.Shutdown()
    grpcConn.Close()
    os.Exit(0)
}()

**C. Error Handling & Retry**
gRPC connection kadang putus. WebSocket kadang timeout. SDK harus handle reconnect secara exponential backoff.

// sdk/js/chat.js (reconnect logic) connectWithRetry() { let retries = 0; const maxRetries = 5; const attemptConnect = () => { this.connect(); this.ws.onclose = () => { if (retries < maxRetries) { const delay = Math.pow(2, retries) * 1000; console.log(Reconnecting in ${delay}ms...); setTimeout(attemptConnect, delay); retries++; } }; }; attemptConnect(); }

D. Elasticsearch Index Lifecycle Messages index bakal gede banget. Pake ILM (Index Lifecycle Management) biar ES otomatis manage shards.