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}