sync_client.go

  1package ygo
  2
  3import (
  4	"context"
  5	"encoding/json"
  6	"fmt"
  7	"sync"
  8	"time"
  9
 10	ygosync "git.kilimanjaro.io/ygo/sync"
 11)
 12
 13// SyncHook is the standard callback signature for sync events.
 14// It receives the document and current connection statistics.
 15type SyncHook func(doc *Doc, stats SyncStats) error
 16
 17// SyncHookCtx is the context-aware callback signature for sync events.
 18// It receives a context (for cancellation/deadlines), the document, and stats.
 19type SyncHookCtx func(ctx context.Context, doc *Doc, stats SyncStats) error
 20
 21// SyncStats holds connection statistics for a sync client.
 22type SyncStats struct {
 23	// PeerCount is the number of other clients currently connected (via awareness)
 24	PeerCount int
 25
 26	// LastUpdate is the time of the last received update from the server
 27	LastUpdate time.Time
 28
 29	// PendingUpdate is true if there are local changes not yet synced to the server
 30	PendingUpdate bool
 31
 32	// Status is the current connection state
 33	Status SyncStatus
 34}
 35
 36// SyncClient manages a WebSocket connection to a y-sweet server for document synchronization.
 37// It provides explicit lifecycle control, connection statistics, and event callbacks.
 38//
 39// The user must explicitly call doc.Destroy() when done with the document - this is never
 40// done automatically, even on disconnect, to allow for retry/reconnect logic.
 41type SyncClient struct {
 42	doc          *Doc
 43	inner        *ygosync.SyncClient
 44	config       *syncConfig
 45	onConnect    SyncHookCtx
 46	onDisconnect DisconnectHookCtx
 47	onUpdate     SyncHookCtx
 48	statsMu      sync.RWMutex
 49	cachedStats  SyncStats
 50}
 51
 52// NewSyncClient creates a new sync client for the document.
 53// The client is not connected yet - use Connect() to establish the connection.
 54// The returned client must be closed with Close() when done.
 55//
 56// IMPORTANT: The caller retains full ownership of the document. You must explicitly
 57// call doc.Destroy() when done. The sync client will NEVER destroy the document
 58// automatically, even on disconnect or errors, to allow for retry/reconnect logic.
 59func NewSyncClient(doc *Doc, opts ...SyncOption) (*SyncClient, error) {
 60	if doc == nil {
 61		return nil, ErrNilDocument
 62	}
 63	if doc.ptr == nil {
 64		return nil, ErrNilDocument
 65	}
 66
 67	cfg := &syncConfig{}
 68	for _, opt := range opts {
 69		opt(cfg)
 70	}
 71
 72	if cfg.endpoint == "" {
 73		return nil, ErrSyncEndpointRequired
 74	}
 75
 76	client := &SyncClient{
 77		doc:    doc,
 78		config: cfg,
 79	}
 80
 81	return client, nil
 82}
 83
 84// Connect establishes the WebSocket connection to the y-sweet server.
 85// This method blocks until the connection is established or fails.
 86// Once connected, the client runs in the background until Close() is called
 87// or the context is cancelled.
 88//
 89// If the context is cancelled, the client will gracefully shutdown,
 90// flush any pending updates, and close the connection.
 91func (c *SyncClient) Connect(ctx context.Context) error {
 92	// Capture callback values first, before acquiring the main lock
 93	c.statsMu.RLock()
 94	onUpdate := c.onUpdate
 95	onConnect := c.onConnect
 96	onDisconnect := c.onDisconnect
 97	c.statsMu.RUnlock()
 98
 99	c.statsMu.Lock()
100	defer c.statsMu.Unlock()
101
102	if c.inner != nil {
103		return fmt.Errorf("already connected")
104	}
105
106	// Set default disconnect timeout
107	disconnectAfter := 5 * time.Minute
108	if c.config.disconnectAfter > 0 {
109		disconnectAfter = c.config.disconnectAfter
110	}
111
112	// Create adapter to implement sync.DocInterface
113	adapter := &syncDocAdapter{doc: c.doc}
114
115	// Build internal sync options
116	opts := ygosync.Options{
117		Endpoint:         c.config.endpoint,
118		AwarenessTimeout: disconnectAfter,
119	}
120	if c.config.token != "" {
121		opts.AuthToken = c.config.token
122	}
123	if c.config.awarenessState != nil {
124		stateBytes, _ := json.Marshal(c.config.awarenessState)
125		opts.AwarenessState = stateBytes
126	}
127
128	// Create internal sync client
129	c.inner = ygosync.NewSyncClient(adapter,
130		ygosync.WithEndpoint(opts.Endpoint),
131		ygosync.WithAwarenessTimeout(opts.AwarenessTimeout))
132
133	if opts.AuthToken != "" {
134		c.inner.SetAuthToken(opts.AuthToken)
135	}
136
137	// Set up internal callbacks that wrap our SyncHookCtx callbacks
138	// Using captured callback values to avoid race conditions
139
140	if onUpdate != nil {
141		doc := c.doc
142		c.inner.SetOnUpdate(func(di ygosync.DocInterface) error {
143			stats := c.getStatsUnsafe()
144			return onUpdate(context.Background(), doc, stats)
145		})
146	}
147
148	// Set up disconnect callback with reason
149	disconnectCalled := false
150	c.inner.SetDisconnectCallback(func(reason ygosync.DisconnectReason) {
151		if !disconnectCalled && onDisconnect != nil {
152			disconnectCalled = true
153			stats := c.getStatsUnsafe()
154			// Convert internal reason to public reason
155			var publicReason DisconnectReason
156			switch reason {
157			case ygosync.DisconnectReasonClientInitiated:
158				publicReason = DisconnectReasonClientInitiated
159			case ygosync.DisconnectReasonYSweetInitiated:
160				publicReason = DisconnectReasonYSweetInitiated
161			case ygosync.DisconnectReasonDocumentIdle:
162				publicReason = DisconnectReasonDocumentIdle
163			}
164			onDisconnect(context.Background(), c.doc, stats, publicReason)
165		}
166	})
167
168	// Store reference in doc for awareness tracking
169	c.doc.sync = c.inner
170
171	// Start connection
172	err := c.inner.Start(ctx)
173	if err != nil {
174		c.inner = nil
175		return err
176	}
177
178	// Trigger OnConnect callback
179	if onConnect != nil {
180		go func() {
181			stats := c.Stats()
182			onConnect(context.Background(), c.doc, stats)
183		}()
184	}
185
186	// Monitor context for graceful shutdown
187	go func() {
188		<-ctx.Done()
189		// Context cancelled - graceful shutdown with flush
190		c.Close()
191	}()
192
193	return nil
194}
195
196// Close gracefully disconnects from the y-sweet server.
197// It flushes any pending updates before closing the connection.
198// Safe to call multiple times.
199//
200// Note: This does NOT destroy the document. You must call doc.Destroy() separately.
201func (c *SyncClient) Close() error {
202	c.statsMu.Lock()
203	defer c.statsMu.Unlock()
204
205	if c.inner == nil {
206		return nil
207	}
208
209	// Trigger OnDisconnect before closing
210	if c.onDisconnect != nil {
211		stats := c.getStatsUnsafe()
212		go func() {
213			c.onDisconnect(context.Background(), c.doc, stats, DisconnectReasonClientInitiated)
214		}()
215	}
216
217	// Send pending updates and close
218	err := c.inner.SendPendingAndClose()
219
220	// Clear doc reference
221	if c.doc != nil {
222		c.doc.sync = nil
223	}
224
225	c.inner = nil
226
227	return err
228}
229
230// Status returns the current connection status.
231func (c *SyncClient) Status() SyncStatus {
232	c.statsMu.RLock()
233	defer c.statsMu.RUnlock()
234
235	if c.inner == nil {
236		return SyncStatusDisconnected
237	}
238
239	return c.getStatsUnsafe().Status
240}
241
242// Stats returns current connection statistics.
243// Returns zero values if not connected.
244func (c *SyncClient) Stats() SyncStats {
245	c.statsMu.RLock()
246	defer c.statsMu.RUnlock()
247
248	return c.getStatsUnsafe()
249}
250
251// getStatsUnsafe returns stats without acquiring lock (must be called with lock held)
252func (c *SyncClient) getStatsUnsafe() SyncStats {
253	if c.inner == nil {
254		return SyncStats{
255			Status: SyncStatusDisconnected,
256		}
257	}
258
259	innerStats := c.inner.GetStats()
260	return SyncStats{
261		PeerCount:     innerStats.PeerCount,
262		LastUpdate:    innerStats.LastUpdate,
263		PendingUpdate: innerStats.PendingUpdate,
264		Status:        syncStatusFromInternal(innerStats.Status),
265	}
266}
267
268// Flush performs a best-effort sync of pending updates.
269// It sends any pending local changes to the server.
270// The context can be used to set a timeout or cancellation.
271// Useful before shutdown to ensure changes are transmitted.
272func (c *SyncClient) Flush(ctx context.Context) error {
273	c.statsMu.RLock()
274	inner := c.inner
275	c.statsMu.RUnlock()
276
277	if inner == nil {
278		return fmt.Errorf("not connected")
279	}
280
281	return inner.Flush(ctx)
282}
283
284// OnConnect sets a callback that is invoked when the client successfully connects.
285// The callback receives the document and current stats.
286func (c *SyncClient) OnConnect(fn SyncHook) {
287	// Wrap the non-context SyncHook to call as SyncHookCtx with background context
288	c.OnConnectCtx(func(ctx context.Context, doc *Doc, stats SyncStats) error {
289		return fn(doc, stats)
290	})
291}
292
293// OnConnectCtx sets a context-aware callback for connect events.
294func (c *SyncClient) OnConnectCtx(fn SyncHookCtx) {
295	c.statsMu.Lock()
296	c.onConnect = fn
297	c.statsMu.Unlock()
298}
299
300// OnDisconnect sets a callback that is invoked when the client disconnects.
301// The callback receives the document, stats, and the reason for disconnection.
302// Note: This does NOT mean you should destroy the document - reconnect is possible.
303func (c *SyncClient) OnDisconnect(fn DisconnectHook) {
304	// Wrap the non-context hook to call as DisconnectHookCtx
305	c.OnDisconnectCtx(func(ctx context.Context, doc *Doc, stats SyncStats, reason DisconnectReason) error {
306		return fn(doc, stats, reason)
307	})
308}
309
310// OnDisconnectCtx sets a context-aware callback for disconnect events.
311func (c *SyncClient) OnDisconnectCtx(fn DisconnectHookCtx) {
312	c.statsMu.Lock()
313	c.onDisconnect = fn
314	c.statsMu.Unlock()
315
316	// Also propagate to inner syncConn if already connected
317	if c.inner != nil {
318		c.inner.SetDisconnectCallback(func(reason ygosync.DisconnectReason) {
319			c.statsMu.RLock()
320			onDisconnect := c.onDisconnect
321			c.statsMu.RUnlock()
322
323			if onDisconnect != nil {
324				stats := c.getStatsUnsafe()
325				// Convert internal reason to public reason
326				var publicReason DisconnectReason
327				switch reason {
328				case ygosync.DisconnectReasonClientInitiated:
329					publicReason = DisconnectReasonClientInitiated
330				case ygosync.DisconnectReasonYSweetInitiated:
331					publicReason = DisconnectReasonYSweetInitiated
332				case ygosync.DisconnectReasonDocumentIdle:
333					publicReason = DisconnectReasonDocumentIdle
334				}
335				onDisconnect(context.Background(), c.doc, stats, publicReason)
336			}
337		})
338	}
339}
340
341// OnUpdate sets a callback that is invoked when the document receives updates from peers.
342// The callback receives the document and current stats.
343func (c *SyncClient) OnUpdate(fn SyncHook) {
344	// Wrap the non-context SyncHook to call as SyncHookCtx with background context
345	c.OnUpdateCtx(func(ctx context.Context, doc *Doc, stats SyncStats) error {
346		return fn(doc, stats)
347	})
348}
349
350// OnUpdateCtx sets a context-aware callback for update events.
351func (c *SyncClient) OnUpdateCtx(fn SyncHookCtx) {
352	c.statsMu.Lock()
353	c.onUpdate = fn
354	c.statsMu.Unlock()
355
356	// Also propagate to inner syncConn if already connected
357	if c.inner != nil {
358		c.inner.SetOnUpdate(func(di ygosync.DocInterface) error {
359			c.statsMu.RLock()
360			onUpdate := c.onUpdate
361			c.statsMu.RUnlock()
362
363			if onUpdate != nil {
364				return onUpdate(context.Background(), c.doc, c.getStatsUnsafe())
365			}
366			return nil
367		})
368	}
369}
370
371// syncStatusFromInternal converts internal ygosync.Status to public SyncStatus
372func syncStatusFromInternal(s ygosync.Status) SyncStatus {
373	switch s {
374	case ygosync.StatusConnecting:
375		return SyncStatusConnecting
376	case ygosync.StatusConnected:
377		return SyncStatusConnected
378	case ygosync.StatusDisconnecting:
379		return SyncStatusDisconnecting
380	default:
381		return SyncStatusDisconnected
382	}
383}