--- title: "Leadera — WebSocket для диаграмм Ганта" date: 2026-04-07 lastmod: 2026-04-07 tags: ["leadera", "websocket", "realtime", "gantt", "architecture"] weight: 110 --- # WebSocket для диаграмм Ганта — realtime-совместное редактирование **Дата:** 07.04.2026 **Статус:** Планирование **Проект:** Leadera (dev — sand.a2v.space) --- ## 1. Постановка задачи Диаграммы Ганта должны обновляться в реальном времени при совместной работе: если два или более пользователей редактируют один Гант, изменения каждого мгновенно отображаются у остальных — без перезагрузки страницы и ручного обновления. --- ## 2. Архитектура ``` Angular frontend ←── WebSocket ──→ Go backend (Hub) ──→ PostgreSQL │ │ │ HTTP (мутации) │ Broadcast └──→ REST API ──→ Service ──→ DB ──→ Hub ──→ все клиенты комнаты ``` **Ключевой принцип:** мутации данных идут через HTTP (REST), WebSocket используется **только для доставки уведомлений**. Сервер — источник истины, клиент не может обойти проверки через WS. ### Комнаты Каждый открытый Гант — это «комната» с идентификатором `chartId`. При подключении клиент подписывается на комнату. Hub рассылает события всем подписчикам комнаты. --- ## 3. Backend (Go) ### 3.1. Зависимости ```bash go get github.com/gorilla/websocket ``` ### 3.2. Новые файлы — `internal/ws/` #### `message.go` — типы сообщений ```go package ws import "time" // EventType — тип события type EventType string const ( // Данные EventSectionCreated EventType = "section_created" EventSectionUpdated EventType = "section_updated" EventSectionDeleted EventType = "section_deleted" EventSectionsReordered EventType = "sections_reordered" EventTaskCreated EventType = "task_created" EventTaskUpdated EventType = "task_updated" EventTaskDeleted EventType = "task_deleted" // Присутствие EventUserJoined EventType = "user_joined" EventUserLeft EventType = "user_left" // Безопасность / доступ EventAccessRevoked EventType = "access_revoked" EventChartArchived EventType = "chart_archived" EventRoleChanged EventType = "role_changed" EventTokenExpired EventType = "token_expired" ) // Event — WS-событие, отправляемое клиентам type Event struct { Type EventType `json:"type"` ChartID string `json:"chartId"` UserID string `json:"userId,omitempty"` Payload interface{} `json:"payload,omitempty"` Timestamp time.Time `json:"timestamp"` } // IncomingMessage — сообщение от клиента (для ping/pong) type IncomingMessage struct { Type string `json:"type"` } ``` #### `hub.go` — центральный брокер ```go package ws import ( "sync" ) // Hub управляет комнатами и рассылкой событий type Hub struct { mu sync.RWMutex rooms map[string]map[*Client]bool // chartId → клиенты } func NewHub() *Hub { return &Hub{ rooms: make(map[string]map[*Client]bool), } } // Subscribe добавляет клиента в комнату func (h *Hub) Subscribe(chartID string, client *Client) { h.mu.Lock() defer h.mu.Unlock() if h.rooms[chartID] == nil { h.rooms[chartID] = make(map[*Client]bool) } h.rooms[chartID][client] = true // Уведомить остальных: user_joined h.broadcastExcept(chartID, client, Event{ Type: EventUserJoined, ChartID: chartID, UserID: client.UserID, Payload: map[string]string{"username": client.Username}, Timestamp: time.Now(), }) } // Unsubscribe удаляет клиента из комнаты func (h *Hub) Unsubscribe(chartID string, client *Client) { h.mu.Lock() defer h.mu.Unlock() if clients, ok := h.rooms[chartID]; ok { delete(clients, client) if len(clients) == 0 { delete(h.rooms, chartID) } } // Уведомить остальных: user_left h.broadcastExcept(chartID, client, Event{ Type: EventUserLeft, ChartID: chartID, UserID: client.UserID, Timestamp: time.Now(), }) } // Broadcast рассылает событие всем клиентам в комнате func (h *Hub) Broadcast(chartID string, event Event) { h.mu.RLock() defer h.mu.RUnlock() if clients, ok := h.rooms[chartID]; ok { for client := range clients { client.Send(event) } } } // SendToUser отправляет событие конкретному пользователю во всех его комнатах func (h *Hub) SendToUser(userID string, event Event) { h.mu.RLock() defer h.mu.RUnlock() for _, clients := range h.rooms { for client := range clients { if client.UserID == userID { client.Send(event) } } } } // broadcastExcept рассылает всем, кроме отправителя func (h *Hub) broadcastExcept(chartID string, sender *Client, event Event) { if clients, ok := h.rooms[chartID]; ok { for client := range clients { if client != sender { client.Send(event) } } } } ``` #### `client.go` — обёртка WS-соединения ```go package ws import ( "encoding/json" "sync" "time" "github.com/gorilla/websocket" ) const ( writeWait = 10 * time.Second pongWait = 60 * time.Second pingPeriod = (pongWait * 9) / 10 maxMessageSize = 1024 ) type Client struct { Hub *Hub Conn *websocket.Conn UserID string Username string ChartID string SendCh chan Event mu sync.Mutex } func NewClient(hub *Hub, conn *websocket.Conn, userID, username, chartID string) *Client { return &Client{ Hub: hub, Conn: conn, UserID: userID, Username: username, ChartID: chartID, SendCh: make(chan Event, 64), } } // ReadPump читает входящие сообщения (ping/pong) func (c *Client) ReadPump() { defer func() { c.Hub.Unsubscribe(c.ChartID, 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 { break } // Клиент может слать только ping, данные игнорируем _ = message } } // WritePump отправляет события клиенту func (c *Client) WritePump() { ticker := time.NewTicker(pingPeriod) defer func() { ticker.Stop() c.Conn.Close() }() for { select { case event, ok := <-c.SendCh: c.Conn.SetWriteDeadline(time.Now().Add(writeWait)) if !ok { c.Conn.WriteMessage(websocket.CloseMessage, []byte{}) return } data, _ := json.Marshal(event) c.mu.Lock() c.Conn.WriteMessage(websocket.TextMessage, data) c.mu.Unlock() case <-ticker.C: c.Conn.SetWriteDeadline(time.Now().Add(writeWait)) if err := c.Conn.WriteMessage(websocket.PingMessage, nil); err != nil { return } } } } // Send ставит событие в очередь отправки func (c *Client) Send(event Event) { select { case c.SendCh <- event: default: // Очередь полна — отключаем медленного клиента close(c.SendCh) } } ``` ### 3.3. WS endpoint — `internal/handlers/ws_handler.go` ```go package handlers import ( "net/http" "time" "app-leadera-api/internal/middleware" "app-leadera-api/internal/ws" "github.com/gin-gonic/gin" "github.com/google/uuid" "github.com/gorilla/websocket" ) var upgrader = websocket.Upgrader{ ReadBufferSize: 1024, WriteBufferSize: 1024, CheckOrigin: func(r *http.Request) bool { return true }, // nginx проксирует } type WSHandler struct { hub *ws.Hub ganttMW *middleware.GanttMiddleware } func NewWSHandler(hub *ws.Hub, ganttMW *middleware.GanttMiddleware) *WSHandler { return &WSHandler{hub: hub, ganttMW: ganttMW} } func (h *WSHandler) HandleGanttWS(c *gin.Context) { // 1. Получить пользователя из контекста (auth middleware уже отработал) userID, exists := c.Get("userID") if !exists { c.AbortWithStatusJSON(401, gin.H{"error": "unauthorized"}) return } // 2. Проверить доступ к chartId (gantt middleware) chartID := c.Param("chartId") if _, err := uuid.Parse(chartID); err != nil { c.AbortWithStatusJSON(400, gin.H{"error": "invalid chart id"}) return } // TODO: проверить через ganttMW что userID имеет доступ к этому chart // 3. Upgrade до WebSocket conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) if err != nil { return } // 4. Создать клиент и подписать на комнату client := ws.NewClient( h.hub, conn, userID.(string), c.GetString("username"), chartID, ) h.hub.Subscribe(chartID, client) // 5. Запустить pump'ы go client.WritePump() go client.ReadPump() } ``` ### 3.4. Регистрация маршрута В `cmd/server/main.go` или где регистрируются routes: ```go // После создания hub wsHub := ws.NewHub() // WS endpoint wsHandler := handlers.NewWSHandler(wsHub, ganttMW) api.GET("/ws/gantt/:chartId", authMW.RequireAuth(), wsHandler.HandleGanttWS) // Передать hub в ganttHandler для broadcast'ов ganttHandler := handlers.NewGanttHandler(ganttService, sectionService, taskService, blockService, wsHub) ``` ### 3.5. Интеграция broadcast в handlers В каждый write-handler — одна строка после успешной операции: ```go // Пример: после успешного UpdateTask h.hub.Broadcast(chartID, ws.Event{ Type: ws.EventTaskUpdated, ChartID: chartID, UserID: userID, Payload: updatedTask, Timestamp: time.Now(), }) ``` Точки интеграции в `gantt_handler.go`: - `CreateGanttChart` → не нужен (пользователь только что создал, он один) - `UpdateGanttChart` / `ArchiveGanttChart` → `chart_archived` - `CreateSection` → `section_created` - `UpdateSection` → `section_updated` - `DeleteSection` → `section_deleted` - `ReorderSections` → `sections_reordered` - `CreateTask` → `task_created` - `UpdateTask` → `task_updated` - `DeleteTask` → `task_deleted` ### 3.6. Интеграция событий безопасности В соответствующих service-методах: ```go // При исключении из space (SpaceService.RemoveMember) wsHub.SendToUser(removedUserID, ws.Event{ Type: ws.EventAccessRevoked, Payload: map[string]string{"spaceId": spaceID}, }) // При архивации ганта wsHub.Broadcast(chartID, ws.Event{ Type: ws.EventChartArchived, ChartID: chartID, }) // При изменении роли wsHub.SendToUser(userID, ws.Event{ Type: ws.EventRoleChanged, Payload: map[string]string{"role": newRole}, }) ``` --- ## 4. Frontend (Angular) ### 4.1. WebSocket сервис — `gantt-ws.service.ts` ```typescript // src/app/pages/gantt/services/gantt-ws.service.ts import { Injectable, inject, OnDestroy } from '@angular/core'; import { Subject, Observable } from 'rxjs'; import { environment } from '../../../environments/environment'; import { SpaceStorageService } from '../../../core/services/space-storage.service'; import { AuthService } from '../../../core/auth/services/auth.service'; export interface WSEvent { type: string; chartId?: string; userId?: string; payload?: any; timestamp?: string; } @Injectable({ providedIn: 'root' }) export class GanttWSService implements OnDestroy { private spaceStorage = inject(SpaceStorageService); private auth = inject(AuthService); private socket: WebSocket | null = null; private eventSubject = new Subject(); private reconnectTimer: any; private chartId: string | null = null; /** Observable для подписки на события */ get events$(): Observable { return this.eventSubject.asObservable(); } /** Подключиться к комнате ганта */ connect(chartId: string): void { this.disconnect(); this.chartId = chartId; const token = this.auth.getAccessToken(); const wsUrl = this.buildWSUrl(chartId, token); this.socket = new WebSocket(wsUrl); this.socket.onopen = () => { console.log(`[WS] Connected to gantt ${chartId}`); }; this.socket.onmessage = (msg) => { try { const event: WSEvent = JSON.parse(msg.data); this.eventSubject.next(event); } catch (e) { console.error('[WS] Parse error', e); } }; this.socket.onclose = (event) => { console.log(`[WS] Closed: ${event.code}`); if (event.code !== 1000) { this.scheduleReconnect(); } }; this.socket.onerror = () => { this.socket?.close(); }; } /** Отключиться */ disconnect(): void { clearTimeout(this.reconnectTimer); if (this.socket) { this.socket.close(1000, 'manual'); this.socket = null; } this.chartId = null; } ngOnDestroy(): void { this.disconnect(); } private buildWSUrl(chartId: string, token: string): string { const base = environment.wsUrl || location.protocol.replace('http', 'ws') + '//' + location.host; return `${base}/api/v1/ws/gantt/${chartId}?token=${token}`; } private scheduleReconnect(): void { if (!this.chartId) return; this.reconnectTimer = setTimeout(() => { console.log('[WS] Reconnecting...'); if (this.chartId) this.connect(this.chartId); }, 3000); } } ``` ### 4.2. Интеграция в gantt-chart-canvas компонент ```typescript // В gantt-chart-canvas.component.ts export class GanttChartCanvasComponent implements OnInit, OnDestroy { private ws = inject(GanttWSService); private onlineUsers: Set = new Set(); private wsSubscription: Subscription | null = null; ngOnInit(): void { // Подключиться к WS при загрузке ганта this.ws.connect(this.chartId); this.wsSubscription = this.ws.events$.subscribe(event => { this.handleWSEvent(event); }); } ngOnDestroy(): void { this.wsSubscription?.unsubscribe(); this.ws.disconnect(); } private handleWSEvent(event: WSEvent): void { // Если это событие от текущего пользователя — пропускаем // (мы уже обновили UI оптимистично или после HTTP-ответа) if (event.userId === this.currentUserId) return; switch (event.type) { case 'section_created': this.addSectionLocal(event.payload); this.redrawCanvas(); break; case 'section_updated': this.updateSectionLocal(event.payload); this.redrawCanvas(); break; case 'section_deleted': this.removeSectionLocal(event.payload.id); this.redrawCanvas(); break; case 'sections_reordered': this.reorderSectionsLocal(event.payload.ordered_ids); this.redrawCanvas(); break; case 'task_created': this.addTaskLocal(event.payload); this.redrawCanvas(); break; case 'task_updated': this.updateTaskLocal(event.payload); this.redrawCanvas(); break; case 'task_deleted': this.removeTaskLocal(event.payload.id); this.redrawCanvas(); break; case 'user_joined': this.onlineUsers.add(event.userId!); this.showOnlineIndicator(); break; case 'user_left': this.onlineUsers.delete(event.userId!); this.showOnlineIndicator(); break; // === Безопасность === case 'access_revoked': this.router.navigate(['/gantt'], { state: { alert: 'Доступ к пространству был отозван' } }); break; case 'chart_archived': this.router.navigate(['/gantt'], { state: { alert: 'Диаграмма была заархивирована' } }); break; case 'role_changed': this.currentRole = event.payload.role; this.updateEditableState(); break; case 'token_expired': this.router.navigate(['/auth/login']); break; } } private redrawCanvas(): void { // Перерисовать Konva canvas с обновлёнными данными } } ``` ### 4.3. Индикатор онлайн-пользователей Добавить в шаблон canvas-компонента: ```html
{{ onlineUsers.size }} онлайн
``` --- ## 5. Безопасность ### 5.1. Авторизация при подключении | Этап | Проверка | |------|----------| | WS Handshake | JWT в query-параметре `?token=...` | | Подписка на комнату | Проверка membership в space | | Каждое broadcast-событие | Hub знает userId каждого клиента | ### 5.2. Сценарии отзыва доступа | Событие | Триггер | Действие на клиенте | |---------|---------|---------------------| | `access_revoked` | Исключение из space | Редирект на список гантов | | `chart_archived` | Архивация ганта | Редирект на список гантов | | `role_changed` | Смена роли (editor→viewer) | Отключить редактирование | | `token_expired` | Токен протух | Редирект на логин | ### 5.3. Защита от злоупотреблений | Угроза | Мера | |--------|------| | Спам WS-сообщений | Rate limit: 30 msg/min на клиента, `ReadLimit` в gorilla/websocket | | Подключение без прав | Проверка space membership при handshake | | Истёкший токен | Ping/pong + периодическая проверка JWT (раз в 5 мин) | | XSS через WS | Payload — JSON, sanitize имён на frontend | | Мутации через WS | WS **только для чтения** (уведомления). Мутации — через HTTP API с полной валидацией | ### 5.4. Nginx конфигурация В конфиг `sand.a2v.space` добавить: ```nginx # WebSocket проксирование location /api/v1/ws/ { proxy_pass http://127.0.0.1:8080; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; proxy_set_header X-Forwarded-Proto $scheme; proxy_read_timeout 86400s; # 24h для long-lived WS proxy_send_timeout 86400s; } ``` --- ## 6. Оценка трудозатрат | Компонент | Строки кода | Время | |-----------|-------------|-------| | `internal/ws/` (hub, client, message) | ~200 | 2-3 часа | | `ws_handler.go` (endpoint) | ~60 | 30 мин | | Интеграция broadcast в handlers | ~20 (по строке на endpoint) | 30 мин | | События безопасности в services | ~30 | 30 мин | | `gantt-ws.service.ts` (Angular) | ~80 | 1 час | | Интеграция в canvas-компонент | ~80 | 1-2 часа | | Индикатор онлайн | ~30 | 30 мин | | Nginx конфигурация | ~10 | 10 мин | | Тестирование | — | 2-3 часа | | **Итого** | **~510 строк** | **~1-2 дня** | --- ## 7. Порядок реализации 1. **Backend: `internal/ws/`** — hub, client, message types 2. **Backend: WS endpoint** — handler + route 3. **Backend: Broadcast в handlers** — интеграция в существующие CRUD-эндпоинты 4. **Backend: События безопасности** — access_revoked, chart_archived, role_changed 5. **Nginx: WS location** — проксирование WebSocket 6. **Frontend: WS service** — подключение, реконнект, парсинг событий 7. **Frontend: Интеграция в canvas** — обработка событий, обновление Konva 8. **Frontend: Индикатор онлайн** — визуальная обратная связь 9. **Тестирование** — многопользовательский сценарий, отзыв прав, обрывы связи --- ## 8. Дальнейшие улучшения (после MVP) - **Cursor других пользователей** — показывать позицию курсора на таймлайне - **Блокировка элемента** — «кто-то редактирует эту задачу» (optimistic locking через WS) - **Чат в ганте** — комментарии к задачам прямо на канвасе - **History/undo** — откат изменений с визуальной историей - **Уведомления** — push-уведомления при @mention в комментариях к задачам