sync.go

  1package waypoint
  2
  3import (
  4	"fmt"
  5
  6	"git.kilimanjaro.io/ygo"
  7)
  8
  9// SyncToArray synchronizes waypoints to a shared Y.Array at the specified key.
 10// It creates the array if it doesn't exist, and only writes changed values
 11// to minimize Yjs update operations.
 12//
 13// The waypoints are stored as JSON objects for interoperability with JS clients.
 14// Each waypoint is marshaled individually and inserted as a JSON value.
 15//
 16// Waypoints are compared position-by-position using deep equality of all fields.
 17// Only positions where waypoint data has actually changed are updated.
 18func SyncToArray(doc *ygo.Doc, key string, waypoints []Waypoint) error {
 19	// Get or create the array
 20	arr, err := doc.GetArray(key)
 21	if err != nil {
 22		return fmt.Errorf("failed to get array %q: %w", key, err)
 23	}
 24	defer arr.Destroy()
 25
 26	return doc.WithWriteTransaction(func(txn *ygo.Transaction) error {
 27		// Read existing waypoints from array
 28		existing, err := readWaypointsFromArray(arr, txn)
 29		if err != nil {
 30			return fmt.Errorf("failed to read existing waypoints: %w", err)
 31		}
 32
 33		// Compute diff and apply changes
 34		if err := applyDiff(arr, txn, existing, waypoints); err != nil {
 35			return fmt.Errorf("failed to apply diff: %w", err)
 36		}
 37
 38		_ = arr.Len()
 39
 40		return nil
 41	})
 42}
 43
 44// readWaypointsFromArray reads all waypoints from the Y.Array as []Waypoint.
 45// Returns empty slice if array is empty or contains non-waypoint data.
 46func readWaypointsFromArray(arr *ygo.Array, txn *ygo.Transaction) ([]Waypoint, error) {
 47	length := arr.Len()
 48	if length == 0 {
 49		return nil, nil
 50	}
 51
 52	waypoints := make([]Waypoint, 0, length)
 53
 54	err := arr.ForEach(txn, func(index uint32, value *ygo.Output) error {
 55		// Try to unmarshal as Waypoint
 56		var wp Waypoint
 57		if err := value.JSON(&wp); err != nil {
 58			// Skip non-waypoint entries
 59			return nil
 60		}
 61		waypoints = append(waypoints, wp)
 62		return nil
 63	})
 64
 65	if err != nil {
 66		return nil, err
 67	}
 68
 69	return waypoints, nil
 70}
 71
 72// applyDiff computes and applies minimal changes to synchronize arrays.
 73// Compares waypoints position-by-position using deep equality.
 74// Only updates positions where waypoint data has actually changed.
 75func applyDiff(arr *ygo.Array, txn *ygo.Transaction, existing, desired []Waypoint) error {
 76	// Process each position
 77	i := 0
 78	for {
 79		arrLen := int(arr.Len())
 80		hasExisting := i < arrLen
 81		hasDesired := i < len(desired)
 82
 83		// Exit when no more elements to process
 84		if !hasExisting && !hasDesired {
 85			break
 86		}
 87
 88		switch {
 89		case hasExisting && hasDesired:
 90			// Both exist - compare all fields using original slice for comparison
 91			if !waypointsEqual(existing[i], desired[i]) {
 92				if err := updateWaypointAtIndex(arr, txn, uint32(i), desired[i]); err != nil {
 93					return fmt.Errorf("failed to update waypoint at %d: %w", i, err)
 94				}
 95			}
 96			// If equal, do nothing (no Yjs update)
 97			i++
 98
 99		case !hasExisting && hasDesired:
100			// Only desired exists - append new waypoint
101			if err := arr.Push(txn, mustJSON(desired[i])); err != nil {
102				return fmt.Errorf("failed to append waypoint at %d: %w", i, err)
103			}
104			i++
105
106		case hasExisting && !hasDesired:
107			// Only existing exists - remove it
108			if err := arr.RemoveRange(txn, uint32(i), 1); err != nil {
109				return fmt.Errorf("failed to remove waypoint at %d: %w", i, err)
110			}
111			// Don't increment i - we need to check same index again since array shrank
112		}
113	}
114
115	return nil
116}
117
118// updateWaypointAtIndex updates a waypoint at a specific index.
119// Removes the old value and inserts the new one at the same position.
120func updateWaypointAtIndex(arr *ygo.Array, txn *ygo.Transaction, index uint32, wp Waypoint) error {
121	// Remove old value
122	if err := arr.RemoveRange(txn, index, 1); err != nil {
123		return err
124	}
125	// Insert new value at same position
126	return insertWaypointAtIndex(arr, txn, index, wp)
127}
128
129// insertWaypointAtIndex inserts a waypoint at a specific index.
130func insertWaypointAtIndex(arr *ygo.Array, txn *ygo.Transaction, index uint32, wp Waypoint) error {
131	input, err := ygo.JSON(wp)
132	if err != nil {
133		return fmt.Errorf("failed to marshal waypoint: %w", err)
134	}
135
136	// For index at end, use Push
137	if index >= arr.Len() {
138		return arr.Push(txn, input)
139	}
140
141	// Otherwise use InsertRange
142	return arr.InsertRange(txn, index, []ygo.Input{input})
143}
144
145// waypointsEqual compares two waypoints for equality across all fields.
146func waypointsEqual(a, b Waypoint) bool {
147	if a.Type != b.Type || a.ID != b.ID || a.Label != b.Label || a.GID != b.GID || a.Country != b.Country || a.Nonroutable != b.Nonroutable {
148		return false
149	}
150	if len(a.Point) != len(b.Point) {
151		return false
152	}
153	for i := range a.Point {
154		if a.Point[i] != b.Point[i] {
155			return false
156		}
157	}
158	// Compare From field - map[string]Route comparison
159	if len(a.From) != len(b.From) {
160		return false
161	}
162	for key, aRoute := range a.From {
163		bRoute, ok := b.From[key]
164		if !ok {
165			return false
166		}
167		if aRoute.ID != bRoute.ID || aRoute.Distance != bRoute.Distance || aRoute.Time != bRoute.Time {
168			return false
169		}
170		if len(aRoute.Bbox) != len(bRoute.Bbox) {
171			return false
172		}
173		for i := range aRoute.Bbox {
174			if aRoute.Bbox[i] != bRoute.Bbox[i] {
175				return false
176			}
177		}
178	}
179	// Compare Cons field - []string comparison
180	if len(a.Cons) != len(b.Cons) {
181		return false
182	}
183	for i := range a.Cons {
184		if a.Cons[i] != b.Cons[i] {
185			return false
186		}
187	}
188	return true
189}
190
191// mustJSON marshals a waypoint to ygo.Input, panicking on error (should never happen)
192func mustJSON(wp Waypoint) ygo.Input {
193	input, err := ygo.JSON(wp)
194	if err != nil {
195		panic(fmt.Sprintf("failed to marshal waypoint: %v", err))
196	}
197	return input
198}