sync.go

  1package ygo
  2
  3import (
  4	"context"
  5	"encoding/json"
  6	"time"
  7
  8	"git.kilimanjaro.io/ygo/sync"
  9)
 10
 11// syncDocAdapter adapts *Doc to sync.DocInterface
 12type syncDocAdapter struct {
 13	doc *Doc
 14}
 15
 16func (a *syncDocAdapter) WithReadTransaction(fn func(sync.Transaction) error) error {
 17	return a.doc.WithReadTransaction(func(txn *Transaction) error {
 18		adapter := &syncTransactionAdapter{txn: txn, doc: a.doc}
 19		return fn(adapter)
 20	})
 21}
 22
 23func (a *syncDocAdapter) WithWriteTransaction(fn func(sync.Transaction) error) error {
 24	return a.doc.WithWriteTransaction(func(txn *Transaction) error {
 25		adapter := &syncTransactionAdapter{txn: txn, doc: a.doc}
 26		return fn(adapter)
 27	})
 28}
 29
 30func (a *syncDocAdapter) GetStateDiff(stateVector []byte) []byte {
 31	var sv *StateVector
 32	if stateVector != nil {
 33		sv = &StateVector{data: stateVector}
 34	}
 35	var result []byte
 36	a.doc.WithReadTransaction(func(txn *Transaction) error {
 37		update := txn.GetStateDiff(sv)
 38		if update != nil {
 39			result = update.Data()
 40		}
 41		return nil
 42	})
 43	return result
 44}
 45
 46// syncTransactionAdapter adapts *Transaction to sync.Transaction
 47type syncTransactionAdapter struct {
 48	txn            *Transaction
 49	doc            *Doc
 50	postUpdateHook func()
 51}
 52
 53func (a *syncTransactionAdapter) ApplyUpdate(data []byte) error {
 54	update := UpdateFromBytes(data)
 55	if update == nil {
 56		return nil
 57	}
 58	return a.txn.ApplyUpdate(update)
 59}
 60
 61func (a *syncTransactionAdapter) GetStateVector() []byte {
 62	sv := a.txn.GetStateVector()
 63	if sv == nil {
 64		return nil
 65	}
 66	return sv.Data()
 67}
 68
 69// SyncStatus represents the connection state to the y-sweet server.
 70type SyncStatus int
 71
 72const (
 73	SyncStatusDisconnected SyncStatus = iota
 74	SyncStatusConnecting
 75	SyncStatusConnected
 76	SyncStatusDisconnecting
 77)
 78
 79// syncConfig holds configuration for the sync connection.
 80type syncConfig struct {
 81	endpoint        string
 82	token           string
 83	onUpdate        Hook
 84	disconnectAfter time.Duration
 85	awarenessState  *AwarenessState
 86	onDisconnect    DisconnectHookCtx
 87}
 88
 89// SyncOption configures the sync connection.
 90type SyncOption func(*syncConfig)
 91
 92// WithSyncEndpoint sets the WebSocket endpoint URL for the sync server.
 93func WithSyncEndpoint(url string) SyncOption {
 94	return func(c *syncConfig) {
 95		c.endpoint = url
 96	}
 97}
 98
 99// WithSyncAuthToken sets the authentication token for the sync server.
100func WithSyncAuthToken(token string) SyncOption {
101	return func(c *syncConfig) {
102		c.token = token
103	}
104}
105
106// WithOnUpdate sets a hook that is called after an update is received
107// from other clients via the y-sweet server.
108func WithOnUpdate(fn Hook) SyncOption {
109	return func(c *syncConfig) {
110		c.onUpdate = fn
111	}
112}
113
114// WithDisconnectOnNoClientsAfter sets the timeout for disconnecting
115// when no other clients are visible via awareness.
116// Default is 5 minutes. Use 0 to disable.
117func WithDisconnectOnNoClientsAfter(d time.Duration) SyncOption {
118	return func(c *syncConfig) {
119		c.disconnectAfter = d
120	}
121}
122
123// WithAwarenessState sets the local awareness state that is
124// broadcast to other clients.
125func WithAwarenessState(state *AwarenessState) SyncOption {
126	return func(c *syncConfig) {
127		c.awarenessState = state
128	}
129}
130
131// WithOnDisconnect sets a hook that is called when the sync client
132// disconnects, with the reason for the disconnection.
133func WithOnDisconnect(fn DisconnectHook) SyncOption {
134	return func(c *syncConfig) {
135		c.onDisconnect = func(ctx context.Context, doc *Doc, stats SyncStats, reason DisconnectReason) error {
136			return fn(doc, stats, reason)
137		}
138	}
139}
140
141// WithOnDisconnectCtx sets a context-aware hook for disconnect events.
142func WithOnDisconnectCtx(fn DisconnectHookCtx) SyncOption {
143	return func(c *syncConfig) {
144		c.onDisconnect = fn
145	}
146}
147
148// Sync establishes a sync connection to the y-sweet server.
149// The document will be synchronized with other clients.
150// When the context is cancelled or times out, the connection closes gracefully.
151// If already connected, this is a no-op.
152func (d *Doc) Sync(ctx context.Context, opts ...SyncOption) error {
153	// If already synced and connected, return no-op
154	if d.sync != nil && d.sync.Status() == sync.StatusConnected {
155		return nil
156	}
157
158	cfg := &syncConfig{}
159	for _, opt := range opts {
160		opt(cfg)
161	}
162
163	if cfg.endpoint == "" {
164		return ErrSyncEndpointRequired
165	}
166
167	// Set default disconnect timeout
168	disconnectAfter := 5 * time.Minute
169	if cfg.disconnectAfter > 0 {
170		disconnectAfter = cfg.disconnectAfter
171	}
172
173	// Create adapter to implement sync.DocInterface
174	adapter := &syncDocAdapter{doc: d}
175
176	// Build sync options
177	syncOpts := []sync.Option{
178		sync.WithEndpoint(cfg.endpoint),
179		sync.WithAwarenessTimeout(disconnectAfter),
180	}
181	if cfg.token != "" {
182		syncOpts = append(syncOpts, sync.WithAuthToken(cfg.token))
183	}
184	if cfg.awarenessState != nil {
185		stateBytes, _ := json.Marshal(cfg.awarenessState)
186		syncOpts = append(syncOpts, sync.WithAwarenessState(stateBytes))
187	}
188
189	syncClient := sync.NewSyncClient(adapter, syncOpts...)
190
191	doc := d
192	if cfg.onUpdate != nil {
193		syncClient.SetOnUpdate(func(di sync.DocInterface) error {
194			return cfg.onUpdate(doc)
195		})
196	}
197
198	// Set up disconnect callback if provided
199	if cfg.onDisconnect != nil {
200		syncClient.SetDisconnectCallback(func(reason sync.DisconnectReason) {
201			stats := syncClient.GetStats()
202			var publicReason DisconnectReason
203			switch reason {
204			case sync.DisconnectReasonClientInitiated:
205				publicReason = DisconnectReasonClientInitiated
206			case sync.DisconnectReasonYSweetInitiated:
207				publicReason = DisconnectReasonYSweetInitiated
208			case sync.DisconnectReasonDocumentIdle:
209				publicReason = DisconnectReasonDocumentIdle
210			}
211			publicStats := SyncStats{
212				PeerCount:     stats.PeerCount,
213				LastUpdate:    stats.LastUpdate,
214				PendingUpdate: stats.PendingUpdate,
215				Status:        syncStatusFromInternal(stats.Status),
216			}
217			cfg.onDisconnect(context.Background(), doc, publicStats, publicReason)
218		})
219	}
220
221	d.sync = syncClient
222
223	go func() {
224		<-ctx.Done()
225		syncClient.Close()
226	}()
227
228	return syncClient.Start(ctx)
229}