Files
etalon f87b954f27 feat: atomization pipeline with job queue, atomizer client, WS hub, and board UI (phase 7)
- internal/extsvc: circuit breaker (§11.16) + retrying bearer JSON client
- atomizer client: §5.1 contract incl. defensive coefficient normalization
  and extension coefficient range (0,2]
- persisted jobs collection: atomic claim, panic recovery, exponential
  backoff (max 3), stale-running requeue, shutdown-safe bookkeeping
- subdivide (async 202, atomizing status, re-run replaces unpublished
  children after confirm) and extend (sibling task, source budget)
- failure path reverts status and notifies the consultant on final attempt
- PATCH task editing with bounty recompute, budget cascade to descendants
- publish single + bulk; task detail endpoint role-scoped
- coder/websocket hub: multiplexed channels, origin check, 30s heartbeats;
  client ws.js with reconnect + 15s polling fallback
- consultant atomization board: tree by root, coefficient sliders with live
  per-parent sum indicator, bounty preview, modals, shimmer, atomizer-down
  gating via /api/v1/service-health
- fix: tasks external-ref unique index is partial, not sparse (compound
  sparse indexed every task through customerId and broke child inserts)

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-06-12 19:17:09 +02:00

191 lines
4.4 KiB
Go

// Package ws implements the single multiplexed live-update WebSocket (§2.2):
// one connection per logged-in session carrying `chat`, `board`, and
// `notifications` channels, with 30 s heartbeats so idle proxies keep the
// connection alive (§12).
package ws
import (
"context"
"encoding/json"
"log/slog"
"net/http"
"sync"
"time"
"github.com/coder/websocket"
)
// Message is the wire format in both directions.
type Message struct {
Channel string `json:"channel"`
Event string `json:"event"`
Data json.RawMessage `json:"data,omitempty"`
}
type client struct {
userID string
send chan []byte
}
// Inbound handles client→server messages (e.g. typing indicators).
type Inbound func(userID string, msg Message)
type Hub struct {
log *slog.Logger
originOK func(r *http.Request) bool
inbound Inbound
mu sync.RWMutex
clients map[*client]struct{}
byUser map[string]map[*client]struct{}
}
func NewHub(log *slog.Logger, originOK func(*http.Request) bool) *Hub {
return &Hub{
log: log,
originOK: originOK,
clients: map[*client]struct{}{},
byUser: map[string]map[*client]struct{}{},
}
}
// SetInbound installs the client-message handler (chat typing etc.).
func (h *Hub) SetInbound(f Inbound) { h.inbound = f }
// ConnectedUserIDs lists users with at least one live socket.
func (h *Hub) ConnectedUserIDs() []string {
h.mu.RLock()
defer h.mu.RUnlock()
out := make([]string, 0, len(h.byUser))
for id := range h.byUser {
out = append(out, id)
}
return out
}
func (h *Hub) register(c *client) {
h.mu.Lock()
defer h.mu.Unlock()
h.clients[c] = struct{}{}
if h.byUser[c.userID] == nil {
h.byUser[c.userID] = map[*client]struct{}{}
}
h.byUser[c.userID][c] = struct{}{}
}
func (h *Hub) unregister(c *client) {
h.mu.Lock()
defer h.mu.Unlock()
delete(h.clients, c)
if set := h.byUser[c.userID]; set != nil {
delete(set, c)
if len(set) == 0 {
delete(h.byUser, c.userID)
}
}
}
func encode(channel, event string, payload any) []byte {
data, err := json.Marshal(payload)
if err != nil {
data = []byte("null")
}
b, _ := json.Marshal(Message{Channel: channel, Event: event, Data: data})
return b
}
// Broadcast sends to every connected client.
func (h *Hub) Broadcast(channel, event string, payload any) {
b := encode(channel, event, payload)
h.mu.RLock()
defer h.mu.RUnlock()
for c := range h.clients {
select {
case c.send <- b:
default: // slow client: drop rather than block the publisher
}
}
}
// SendTo delivers to all sockets of the given users.
func (h *Hub) SendTo(userIDs []string, channel, event string, payload any) {
b := encode(channel, event, payload)
h.mu.RLock()
defer h.mu.RUnlock()
for _, id := range userIDs {
for c := range h.byUser[id] {
select {
case c.send <- b:
default:
}
}
}
}
// Handle upgrades the request (the caller has already authenticated userID).
func (h *Hub) Handle(w http.ResponseWriter, r *http.Request, userID string) {
if h.originOK != nil && !h.originOK(r) {
http.Error(w, "origin not allowed", http.StatusForbidden)
return
}
conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{
// Origin verified above against APP_BASE_URL (§12).
InsecureSkipVerify: true,
})
if err != nil {
h.log.Warn("ws accept", "err", err)
return
}
c := &client{userID: userID, send: make(chan []byte, 64)}
h.register(c)
defer h.unregister(c)
ctx, cancel := context.WithCancel(r.Context())
defer cancel()
go h.writeLoop(ctx, cancel, conn, c)
h.readLoop(ctx, conn, c)
conn.Close(websocket.StatusNormalClosure, "bye")
}
func (h *Hub) writeLoop(ctx context.Context, cancel func(), conn *websocket.Conn, c *client) {
defer cancel()
heartbeat := time.NewTicker(30 * time.Second)
defer heartbeat.Stop()
for {
select {
case <-ctx.Done():
return
case msg := <-c.send:
wctx, wcancel := context.WithTimeout(ctx, 10*time.Second)
err := conn.Write(wctx, websocket.MessageText, msg)
wcancel()
if err != nil {
return
}
case <-heartbeat.C:
pctx, pcancel := context.WithTimeout(ctx, 10*time.Second)
err := conn.Ping(pctx)
pcancel()
if err != nil {
return
}
}
}
}
func (h *Hub) readLoop(ctx context.Context, conn *websocket.Conn, c *client) {
for {
_, data, err := conn.Read(ctx)
if err != nil {
return
}
if h.inbound == nil {
continue
}
var msg Message
if err := json.Unmarshal(data, &msg); err != nil {
continue
}
h.inbound(c.userID, msg)
}
}