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.

Memuat komentar...