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}