Files
docs.a2v.space/content/leadera/leadera-websocket-gantt.md
T

747 lines
23 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
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<WSEvent>();
private reconnectTimer: any;
private chartId: string | null = null;
/** Observable для подписки на события */
get events$(): Observable<WSEvent> {
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<string> = 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
<!-- Онлайн-индикатор -->
<div class="online-indicator" *ngIf="onlineUsers.size > 0">
<span class="pulse-dot"></span>
{{ onlineUsers.size }} онлайн
</div>
```
---
## 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 в комментариях к задачам