diff --git a/crdt.go b/crdt.go index 4559dfd..77360e3 100644 --- a/crdt.go +++ b/crdt.go @@ -4,16 +4,39 @@ import ( "bytes" "context" "encoding/gob" + "encoding/json" "fmt" "maps" "strconv" "strings" "sync" + "time" "github.com/ipfs/go-datastore" "github.com/ipfs/go-datastore/query" ) +// PeerMetadata represents metadata about a peer in the cluster +type PeerMetadata struct { + // BestBefore is the Unix timestamp after which this peer + // should be considered offline if not refreshed + BestBefore uint64 `json:"best_before"` + + // Metadata contains user-defined key-value pairs for this peer + // Examples: {"name": "node-1", "role": "write"} + Metadata map[string]string `json:"metadata"` + + // Sequence is the last known sequence number for this peer + Sequence uint64 `json:"sequence"` +} + +// MembershipState represents the current cluster membership view +type MembershipState struct { + // Members maps peer ID to their metadata + // Key format: libp2p peer ID as string + Members map[string]*PeerMetadata `json:"members"` +} + type CRDT struct { PeerID string PeerSeq uint64 @@ -22,6 +45,12 @@ type CRDT struct { mu sync.Mutex hooksMu sync.RWMutex + // Namespace is an optional prefix for all keys to allow multiple + // independent CRDT instances to coexist in the same datastore. + // If empty, keys are stored at the root (e.g., "/data/key"). + // If set (e.g., "ipservice"), keys are namespaced (e.g., "/ipservice/data/key"). + namespace string + // Transactional hooks (called inside transaction, can abort) insertTxnHook func(ctx context.Context, txn datastore.Write, key string, val []byte, meta CRDTKeyMeta) error updateTxnHook func(ctx context.Context, txn datastore.Write, key string, oldVal []byte, oldMeta CRDTKeyMeta, newVal []byte, newMeta CRDTKeyMeta) error @@ -31,6 +60,73 @@ type CRDT struct { insertHooks []func(key string, val []byte, meta CRDTKeyMeta) updateHooks []func(key string, oldVal []byte, oldMeta CRDTKeyMeta, newVal []byte, newMeta CRDTKeyMeta) deleteHooks []func(key string, oldVal []byte, oldMeta CRDTKeyMeta) + + // Metadata support + localMetadata map[string]string + peerMetadata map[string]*PeerMetadata // peerID -> metadata + metadataMu sync.RWMutex + membershipHooks []func(members map[string]*PeerMetadata) + peerTTL time.Duration +} + +// Key construction helpers - these ensure namespace is consistently applied +// across all internal storage. If namespace is empty, keys are at root level. + +func (c *CRDT) dataKey(userKey string) datastore.Key { + if c.namespace == "" { + return datastore.NewKey("data").ChildString(userKey) + } + return datastore.NewKey(c.namespace).ChildString("data").ChildString(userKey) +} + +func (c *CRDT) changeKey(peerID string, seq uint64) datastore.Key { + seqStr := fmt.Sprintf("%010d", seq) + if c.namespace == "" { + return datastore.NewKey("change").ChildString(peerID).ChildString(seqStr) + } + return datastore.NewKey(c.namespace).ChildString("change").ChildString(peerID).ChildString(seqStr) +} + +func (c *CRDT) trackedKey(peerID string) datastore.Key { + if c.namespace == "" { + return datastore.NewKey("tracked").ChildString(peerID) + } + return datastore.NewKey(c.namespace).ChildString("tracked").ChildString(peerID) +} + +func (c *CRDT) metadataKey(peerID string) datastore.Key { + if c.namespace == "" { + return datastore.NewKey("metadata").ChildString(peerID) + } + return datastore.NewKey(c.namespace).ChildString("metadata").ChildString(peerID) +} + +func (c *CRDT) dataPrefix() string { + if c.namespace == "" { + return "/data/" + } + return datastore.NewKey(c.namespace).ChildString("data").String() + "/" +} + +func (c *CRDT) changePrefix(peerID string) string { + if c.namespace == "" { + return "/change/" + peerID + "/" + } + return datastore.NewKey(c.namespace).ChildString("change").ChildString(peerID).String() + "/" +} + +func (c *CRDT) trackedPrefix() string { + if c.namespace == "" { + return "/tracked/" + } + return datastore.NewKey(c.namespace).ChildString("tracked").String() + "/" +} + +func (c *CRDT) metadataPrefix() string { + if c.namespace == "" { + return "/metadata/" + } + return datastore.NewKey(c.namespace).ChildString("metadata").String() + "/" } func (c *CRDT) AddInsertHook(h func(key string, val []byte, meta CRDTKeyMeta)) { @@ -92,9 +188,12 @@ func (c *CRDT) runDeleteHooks(key string, oldVal []byte, oldMeta CRDTKeyMeta, ho func New(peerID string, ds datastore.Datastore, opts ...func(*CRDT)) *CRDT { crdt := &CRDT{ - PeerID: peerID, - trackedPeers: make(map[string]uint64), - ds: ds, + PeerID: peerID, + trackedPeers: make(map[string]uint64), + ds: ds, + localMetadata: make(map[string]string), + peerMetadata: make(map[string]*PeerMetadata), + peerTTL: 7 * 24 * time.Hour, // default 7 days, matching go-ds-crdt } for _, opt := range opts { @@ -102,6 +201,7 @@ func New(peerID string, ds datastore.Datastore, opts ...func(*CRDT)) *CRDT { } _ = crdt.loadTrackedPeers() + _ = crdt.loadPeerMetadata(context.Background()) crdt.PeerSeq = crdt.trackedPeers[peerID] return crdt } @@ -161,16 +261,60 @@ func WithDeleteTxnHook(h func(ctx context.Context, txn datastore.Write, key stri return func(c *CRDT) { c.deleteTxnHook = h } } +// WithMembershipHook sets a callback that is invoked whenever +// cluster membership changes. This includes: +// - New peer joining the cluster +// - Existing peer updating their metadata +// - Peer TTL expiring (peer removed from membership) +// +// The hook receives a map of all current members with their metadata. +// The hook is called asynchronously (in a goroutine) and should not block. +func WithMembershipHook(hook func(members map[string]*PeerMetadata)) func(*CRDT) { + return func(c *CRDT) { + c.hooksMu.Lock() + defer c.hooksMu.Unlock() + c.membershipHooks = append(c.membershipHooks, hook) + } +} + +// WithPeerTTL sets the duration that a peer's membership remains valid +// without refresh. Default is 7 days (matching go-ds-crdt default). +// +// The TTL is automatically refreshed on each rebroadcast interval. +func WithPeerTTL(ttl time.Duration) func(*CRDT) { + return func(c *CRDT) { + c.peerTTL = ttl + } +} + +// WithNamespace sets a namespace prefix for all keys in the datastore. +// This allows multiple independent CRDT instances to coexist in the same +// underlying datastore without conflicts. +// +// Example: +// +// crdt := clset.New(peerID, ds, clset.WithNamespace("ipservice")) +// // User keys stored as: /ipservice/data/key +// // Internal keys stored as: /ipservice/change/..., /ipservice/tracked/..., etc. +// +// If not set, keys are stored at the root level (e.g., /data/key) for +// backward compatibility with existing datastores. +func WithNamespace(namespace string) func(*CRDT) { + return func(c *CRDT) { + c.namespace = namespace + } +} + func (c *CRDT) loadTrackedPeers() error { ctx := context.Background() // Use QueryIter for efficient prefix iteration - q := query.Query{Prefix: "/tracked/"} + q := query.Query{Prefix: c.trackedPrefix()} for entry, err := range datastore.QueryIter(ctx, c.ds, q) { if err != nil { return err } - peerID := strings.TrimPrefix(entry.Key, "/tracked/") + peerID := strings.TrimPrefix(entry.Key, c.trackedPrefix()) seq, err := strconv.ParseUint(string(entry.Value), 10, 64) if err != nil { return err @@ -181,7 +325,7 @@ func (c *CRDT) loadTrackedPeers() error { } func (c *CRDT) saveTrackedPeer(ctx context.Context, writer datastore.Write, peerID string, seq uint64) error { - key := datastore.NewKey("/tracked/" + peerID) + key := c.trackedKey(peerID) val := []byte(strconv.FormatUint(seq, 10)) return writer.Put(ctx, key, val) } @@ -197,7 +341,7 @@ func decodeEntry(data []byte, entry *KeyEntry) error { } func (c *CRDT) getEntry(ctx context.Context, reader datastore.Read, key string) (KeyEntry, error) { - dsKey := datastore.NewKey("/data/" + key) + dsKey := c.dataKey(key) val, err := reader.Get(ctx, dsKey) if err != nil { return KeyEntry{}, err @@ -245,14 +389,14 @@ func (c *CRDT) applyEntry( ) error { // clean up old index if oldEntry.Meta.PeerID != "" { - oldIndexKey := datastore.NewKey(fmt.Sprintf("/change/%s/%010d", oldEntry.Meta.PeerID, oldEntry.Meta.PeerSeq)) + oldIndexKey := c.changeKey(oldEntry.Meta.PeerID, oldEntry.Meta.PeerSeq) _ = writer.Delete(ctx, oldIndexKey) } encoded, _ := encodeEntry(newEntry) - if err := writer.Put(ctx, datastore.NewKey("/data/"+newEntry.Key), encoded); err != nil { + if err := writer.Put(ctx, c.dataKey(newEntry.Key), encoded); err != nil { return err } - indexKey := datastore.NewKey(fmt.Sprintf("/change/%s/%010d", newEntry.Meta.PeerID, newEntry.Meta.PeerSeq)) + indexKey := c.changeKey(newEntry.Meta.PeerID, newEntry.Meta.PeerSeq) return writer.Put(ctx, indexKey, encoded) } @@ -611,7 +755,7 @@ func (c *CRDT) KeyCount() (uint64, error) { ctx := context.Background() // Use QueryIter for efficient prefix iteration - q := query.Query{Prefix: "/data/"} + q := query.Query{Prefix: c.dataPrefix()} for entry, err := range datastore.QueryIter(ctx, c.ds, q) { if err != nil { return 0, err @@ -628,6 +772,58 @@ func (c *CRDT) KeyCount() (uint64, error) { return count, nil } +// Query executes a query against the datastore. +// This method exposes the underlying datastore's Query functionality, +// allowing callers to perform operations like prefix-based key filtering. +// Query executes a query against the CRDT's data. +// The query prefix is automatically adjusted to include the namespace and /data/ prefix. +// For example, if the namespace is "ipservice" and the query prefix is "/pool/", +// the actual query will use "/ipservice/data/pool/". +// Results have the namespace and /data/ prefix stripped from keys. +func (c *CRDT) Query(ctx context.Context, q query.Query) (query.Results, error) { + // Adjust the query prefix to include namespace and /data/ + if q.Prefix != "" { + // User provides prefixes like "/pool/" or "pool/" + // We need to convert to "/namespace/data/pool/" or "/data/pool/" + q.Prefix = c.dataPrefix() + strings.TrimPrefix(q.Prefix, "/") + } + + results, err := c.ds.Query(ctx, q) + if err != nil { + return nil, err + } + + // Wrap results to strip namespace+data prefix from keys + return &namespacedQueryResults{ + Results: results, + prefix: c.dataPrefix(), + }, nil +} + +// namespacedQueryResults wraps query.Results to strip the namespace+data prefix from keys +type namespacedQueryResults struct { + query.Results + prefix string +} + +func (n *namespacedQueryResults) Next() <-chan query.Result { + originalCh := n.Results.Next() + strippedCh := make(chan query.Result) + + go func() { + defer close(strippedCh) + for result := range originalCh { + // Strip the namespace+data prefix from the key + if strings.HasPrefix(result.Key, n.prefix) { + result.Key = result.Key[len(n.prefix):] + } + strippedCh <- result + } + }() + + return strippedCh +} + // GetTrackedPeers returns a copy of the tracked peers map func (c *CRDT) GetTrackedPeers() map[string]uint64 { c.mu.Lock() @@ -649,7 +845,7 @@ func (c *CRDT) GetLatestChanges(requestorID string, requestorTracked map[string] count := 0 for peerID := range c.trackedPeers { startSeq := requestorTracked[peerID] + 1 - prefix := fmt.Sprintf("/change/%s/", peerID) + prefix := c.changePrefix(peerID) // Create query for this peer's changes q := query.Query{Prefix: prefix} @@ -719,7 +915,6 @@ func (c *CRDT) mergeChangeBatch(ctx context.Context, writer datastore.Write, cha } if isRemoteWinner(incoming, existing) { - err := c.applyEntry(ctx, writer, existing, incoming) if err != nil { return nil, err @@ -869,3 +1064,159 @@ func (c *CRDT) MergeChanges(fromPeerID string, changes []KeyEntry, tracked map[s return nil } + +// UpdateMeta sets or updates the local peer's metadata. +// This metadata will be broadcast to all other peers via gossip. +// +// The TTL is automatically managed - each rebroadcast updates the +// BestBefore timestamp to Now + TTL duration. +// +// Returns error if metadata cannot be persisted. +func (c *CRDT) UpdateMeta(ctx context.Context, metadata map[string]string) error { + c.metadataMu.Lock() + c.localMetadata = make(map[string]string) + for k, v := range metadata { + c.localMetadata[k] = v + } + c.metadataMu.Unlock() + + // Persist to datastore + key := c.metadataKey(c.PeerID) + data, err := json.Marshal(metadata) + if err != nil { + return fmt.Errorf("failed to marshal metadata: %w", err) + } + + if err := c.ds.Put(ctx, key, data); err != nil { + return fmt.Errorf("failed to persist metadata: %w", err) + } + + return nil +} + +// GetState returns the current membership state of the cluster. +// This includes all known peers with their metadata and TTL status. +// +// The returned state is a snapshot at the time of the call. +// Expired peers are filtered out based on their BestBefore timestamp. +func (c *CRDT) GetState(ctx context.Context) *MembershipState { + c.metadataMu.RLock() + defer c.metadataMu.RUnlock() + + state := &MembershipState{ + Members: make(map[string]*PeerMetadata), + } + + now := uint64(time.Now().Unix()) + + for peerID, meta := range c.peerMetadata { + // Filter out expired peers + if meta.BestBefore >= now { + // Deep copy to prevent mutation + state.Members[peerID] = &PeerMetadata{ + BestBefore: meta.BestBefore, + Metadata: copyStringMap(meta.Metadata), + Sequence: c.trackedPeers[peerID], + } + } + } + + return state +} + +// AddSelfToPeerMetadata manually adds the local peer to its own peer metadata. +// This is useful for single-node scenarios or testing where the node needs to +// see itself in the membership immediately without waiting for P2P gossip. +// +// In multi-node scenarios, peers learn about each other through P2P gossip, but +// a node ignores its own gossip messages. This method ensures the local peer +// can see itself in the membership state. +func (c *CRDT) AddSelfToPeerMetadata() { + c.metadataMu.Lock() + + // Get local metadata + metadata := make(map[string]string) + for k, v := range c.localMetadata { + metadata[k] = v + } + + // Calculate TTL for local peer + bestBefore := uint64(time.Now().Add(c.peerTTL).Unix()) + + // Add self to peer metadata + c.peerMetadata[c.PeerID] = &PeerMetadata{ + BestBefore: bestBefore, + Metadata: metadata, + Sequence: c.PeerSeq, + } + c.metadataMu.Unlock() + + // Trigger membership hooks + c.runMembershipHooks() +} + +// loadPeerMetadata loads peer metadata from the datastore on startup +func (c *CRDT) loadPeerMetadata(ctx context.Context) error { + q := query.Query{Prefix: c.metadataPrefix()} + results, err := c.ds.Query(ctx, q) + if err != nil { + return err + } + defer results.Close() + + c.metadataMu.Lock() + defer c.metadataMu.Unlock() + + for result := range results.Next() { + if result.Error != nil { + return result.Error + } + + peerID := strings.TrimPrefix(result.Key, c.metadataPrefix()) + + var metadata map[string]string + if err := json.Unmarshal(result.Value, &metadata); err != nil { + return err + } + + // Initialize with expired TTL - will be updated by gossip + if c.peerMetadata[peerID] == nil { + c.peerMetadata[peerID] = &PeerMetadata{ + BestBefore: 0, + Metadata: metadata, + } + } else { + c.peerMetadata[peerID].Metadata = metadata + } + } + + return nil +} + +// runMembershipHooks calls all registered membership hooks with current state +func (c *CRDT) runMembershipHooks() { + c.hooksMu.RLock() + hooks := c.membershipHooks + c.hooksMu.RUnlock() + + if len(hooks) == 0 { + return + } + + state := c.GetState(context.Background()) + + for _, hook := range hooks { + go hook(state.Members) + } +} + +func copyStringMap(m map[string]string) map[string]string { + if m == nil { + return nil + } + result := make(map[string]string, len(m)) + for k, v := range m { + result[k] = v + } + return result +} diff --git a/crdt_test.go b/crdt_test.go index d1a90e0..e060708 100644 --- a/crdt_test.go +++ b/crdt_test.go @@ -5,6 +5,7 @@ import ( "testing" "github.com/ipfs/go-datastore" + "github.com/ipfs/go-datastore/query" badgerds "github.com/ipfs/go-ds-badger4" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -489,3 +490,151 @@ func TestCRDT_RemoteHooks_OnMerge(t *testing.T) { require.NoError(t, crdt2.MergeChanges("p1", changes, tracked)) assert.Contains(t, inserts, "foo:v3") } + +func TestCRDT_Query(t *testing.T) { + ds := createTestDatastore(t) + crdt := clset.New("peer1", ds) + + // Setup: Create test data with different prefixes + testData := map[string][]byte{ + "pool/pool1/subscriber/sub1": []byte("allocation1"), + "pool/pool1/subscriber/sub2": []byte("allocation2"), + "pool/pool2/subscriber/sub3": []byte("allocation3"), + "subscriber/sub1": []byte("metadata1"), + "subscriber/sub2": []byte("metadata2"), + "other/data": []byte("other"), + } + + for key, value := range testData { + require.NoError(t, crdt.Set(key, value)) + } + + t.Run("query with prefix - pool1 allocations", func(t *testing.T) { + ctx := t.Context() + q := query.Query{Prefix: "/data/pool/pool1/"} + + results, err := crdt.Query(ctx, q) + require.NoError(t, err) + require.NotNil(t, results) + defer results.Close() + + // Collect results + var keys []string + for result := range results.Next() { + require.NoError(t, result.Error) + keys = append(keys, result.Key) + } + + // Should find 2 allocations under pool1 + assert.Len(t, keys, 2) + assert.Contains(t, keys, "/data/pool/pool1/subscriber/sub1") + assert.Contains(t, keys, "/data/pool/pool1/subscriber/sub2") + }) + + t.Run("query with prefix - subscriber metadata", func(t *testing.T) { + ctx := t.Context() + q := query.Query{Prefix: "/data/subscriber/"} + + results, err := crdt.Query(ctx, q) + require.NoError(t, err) + require.NotNil(t, results) + defer results.Close() + + // Collect results + var keys []string + for result := range results.Next() { + require.NoError(t, result.Error) + keys = append(keys, result.Key) + } + + // Should find 2 subscriber records + assert.Len(t, keys, 2) + assert.Contains(t, keys, "/data/subscriber/sub1") + assert.Contains(t, keys, "/data/subscriber/sub2") + }) + + t.Run("query with prefix - all pools", func(t *testing.T) { + ctx := t.Context() + q := query.Query{Prefix: "/data/pool/"} + + results, err := crdt.Query(ctx, q) + require.NoError(t, err) + require.NotNil(t, results) + defer results.Close() + + // Collect results + var keys []string + for result := range results.Next() { + require.NoError(t, result.Error) + keys = append(keys, result.Key) + } + + // Should find 3 pool allocations (2 from pool1, 1 from pool2) + assert.Len(t, keys, 3) + }) + + t.Run("query with non-matching prefix", func(t *testing.T) { + ctx := t.Context() + q := query.Query{Prefix: "/data/nonexistent/"} + + results, err := crdt.Query(ctx, q) + require.NoError(t, err) + require.NotNil(t, results) + defer results.Close() + + // Collect results + var keys []string + for result := range results.Next() { + require.NoError(t, result.Error) + keys = append(keys, result.Key) + } + + // Should find no results + assert.Len(t, keys, 0) + }) + + t.Run("query all data keys", func(t *testing.T) { + ctx := t.Context() + q := query.Query{Prefix: "/data/"} + + results, err := crdt.Query(ctx, q) + require.NoError(t, err) + require.NotNil(t, results) + defer results.Close() + + // Collect results + var keys []string + for result := range results.Next() { + require.NoError(t, result.Error) + keys = append(keys, result.Key) + } + + // Should find all 6 test keys + assert.Len(t, keys, 6) + }) + + t.Run("query returns values", func(t *testing.T) { + ctx := t.Context() + q := query.Query{Prefix: "/data/pool/pool1/"} + + results, err := crdt.Query(ctx, q) + require.NoError(t, err) + require.NotNil(t, results) + defer results.Close() + + // Collect at least one result + foundResults := false + for result := range results.Next() { + require.NoError(t, result.Error) + + // Result should have key and value from the datastore + // Note: The value contains CRDT metadata, not the raw user value + assert.NotEmpty(t, result.Key) + assert.NotNil(t, result.Value) + foundResults = true + } + + // Should have found at least one result + assert.True(t, foundResults, "Expected to find at least one result") + }) +} diff --git a/metadata_test.go b/metadata_test.go new file mode 100644 index 0000000..248bb1f --- /dev/null +++ b/metadata_test.go @@ -0,0 +1,167 @@ +package clset + +import ( + "context" + "testing" + "time" + + "github.com/ipfs/go-datastore" + dssync "github.com/ipfs/go-datastore/sync" +) + +func TestUpdateMeta(t *testing.T) { + ds := dssync.MutexWrap(datastore.NewMapDatastore()) + crdt := New("peer1", ds) + + metadata := map[string]string{ + "name": "TestApp", + "role": "write", + } + + err := crdt.UpdateMeta(context.Background(), metadata) + if err != nil { + t.Fatalf("UpdateMeta failed: %v", err) + } + + // Verify metadata was persisted + crdt2 := New("peer1", ds) + crdt2.metadataMu.RLock() + if len(crdt2.peerMetadata) == 0 { + t.Fatal("Expected metadata to be loaded from datastore") + } + crdt2.metadataMu.RUnlock() +} + +func TestGetState(t *testing.T) { + ds := dssync.MutexWrap(datastore.NewMapDatastore()) + crdt := New("peer1", ds) + + // Add some peer metadata manually + crdt.metadataMu.Lock() + crdt.peerMetadata["peer2"] = &PeerMetadata{ + BestBefore: uint64(time.Now().Add(1 * time.Hour).Unix()), + Metadata: map[string]string{ + "name": "TestApp", + "role": "write", + }, + Sequence: 42, + } + // Add an expired peer + crdt.peerMetadata["peer3"] = &PeerMetadata{ + BestBefore: uint64(time.Now().Add(-1 * time.Hour).Unix()), // expired + Metadata: map[string]string{ + "name": "TestApp", + "role": "read", + }, + Sequence: 10, + } + crdt.metadataMu.Unlock() + + state := crdt.GetState(context.Background()) + + // Should only include non-expired peers + if len(state.Members) != 1 { + t.Fatalf("Expected 1 member (peer3 should be filtered as expired), got %d", len(state.Members)) + } + + peer2, exists := state.Members["peer2"] + if !exists { + t.Fatal("Expected peer2 to be in state") + } + + if peer2.Metadata["name"] != "TestApp" { + t.Errorf("Expected name=TestApp, got %s", peer2.Metadata["name"]) + } + + if peer2.Metadata["role"] != "write" { + t.Errorf("Expected role=write, got %s", peer2.Metadata["role"]) + } +} + +func TestMembershipHook(t *testing.T) { + ds := dssync.MutexWrap(datastore.NewMapDatastore()) + + hookCalled := make(chan bool, 1) + var receivedMembers map[string]*PeerMetadata + + crdt := New("peer1", ds, WithMembershipHook(func(members map[string]*PeerMetadata) { + receivedMembers = members + select { + case hookCalled <- true: + default: + } + })) + + // Add a peer + crdt.metadataMu.Lock() + crdt.peerMetadata["peer2"] = &PeerMetadata{ + BestBefore: uint64(time.Now().Add(1 * time.Hour).Unix()), + Metadata: map[string]string{ + "name": "TestApp", + }, + } + crdt.metadataMu.Unlock() + + // Trigger hook + crdt.runMembershipHooks() + + // Wait for hook to be called (with timeout) + select { + case <-hookCalled: + // Success + case <-time.After(1 * time.Second): + t.Fatal("Membership hook was not called within timeout") + } + + // Verify received data + if len(receivedMembers) != 1 { + t.Fatalf("Expected 1 member in hook callback, got %d", len(receivedMembers)) + } + + if receivedMembers["peer2"].Metadata["name"] != "TestApp" { + t.Error("Hook received incorrect metadata") + } +} + +func TestWithPeerTTL(t *testing.T) { + ds := dssync.MutexWrap(datastore.NewMapDatastore()) + customTTL := 1 * time.Hour + + crdt := New("peer1", ds, WithPeerTTL(customTTL)) + + if crdt.peerTTL != customTTL { + t.Errorf("Expected TTL to be %v, got %v", customTTL, crdt.peerTTL) + } +} + +func TestCopyStringMap(t *testing.T) { + original := map[string]string{ + "key1": "value1", + "key2": "value2", + } + + copied := copyStringMap(original) + + // Verify contents match + if len(copied) != len(original) { + t.Fatalf("Expected %d entries, got %d", len(original), len(copied)) + } + + for k, v := range original { + if copied[k] != v { + t.Errorf("Expected %s=%s, got %s", k, v, copied[k]) + } + } + + // Verify it's a deep copy (modifying copied doesn't affect original) + copied["key3"] = "value3" + if _, exists := original["key3"]; exists { + t.Error("Modifying copied map affected original map") + } + + // Test nil map + nilCopy := copyStringMap(nil) + if nilCopy != nil { + t.Error("copyStringMap(nil) should return nil") + } +} diff --git a/p2p_sync.go b/p2p_sync.go index 873f169..0dff4cc 100644 --- a/p2p_sync.go +++ b/p2p_sync.go @@ -60,8 +60,10 @@ type Peer struct { } type SummaryMessage struct { - PeerID string `json:"peer_id"` - Tracked map[string]uint64 `json:"tracked"` + PeerID string `json:"peer_id"` + Tracked map[string]uint64 `json:"tracked"` + Metadata map[string]string `json:"metadata"` // NEW: Peer metadata + TTL uint64 `json:"ttl"` // NEW: BestBefore timestamp } type cachedSummary struct { @@ -176,6 +178,9 @@ func NewPeer(crdt *CRDT, ctx context.Context, host host.Host, opts ...PeerOption go p2p.broadcastSummaries() go p2p.scheduledSync() + // Start TTL cleanup for expired peers + go p2p.cleanupExpiredPeers() + log.Printf("P2P CRDT node started") log.Printf("Peer ID: %s", host.ID()) log.Printf("Listening on: %v", host.Addrs()) @@ -583,6 +588,18 @@ func (p *Peer) broadcastSummaries() { return case <-p.summaryCh: tracked := p.crdt.GetTrackedPeers() + + // Get local metadata + p.crdt.metadataMu.RLock() + metadata := make(map[string]string) + for k, v := range p.crdt.localMetadata { + metadata[k] = v + } + p.crdt.metadataMu.RUnlock() + + // Calculate TTL (BestBefore timestamp) + ttl := uint64(time.Now().Add(p.crdt.peerTTL).Unix()) + // Avoid re-sending identical state too often if maps.Equal(tracked, p.lastSummary) && time.Since(p.lastPub) < p.config.MinSummaryInterval { continue @@ -591,12 +608,15 @@ func (p *Peer) broadcastSummaries() { p.lastPub = time.Now() msg := SummaryMessage{ - PeerID: p.Host.ID().String(), - Tracked: tracked, + PeerID: p.Host.ID().String(), + Tracked: tracked, + Metadata: metadata, + TTL: ttl, } data, err := json.Marshal(msg) if err != nil { log.Printf("Failed to marshal summary: %v", err) + continue } if err := p.topic.Publish(p.ctx, data); err != nil { log.Printf("Failed to publish summary: %v", err) @@ -640,6 +660,34 @@ func (p *Peer) readSummaries() { continue } + // Update peer metadata + p.crdt.metadataMu.Lock() + changed := false + + existing := p.crdt.peerMetadata[msg.PeerID] + if existing == nil { + p.crdt.peerMetadata[msg.PeerID] = &PeerMetadata{ + BestBefore: msg.TTL, + Metadata: msg.Metadata, + } + changed = true + } else { + // Update if changed + if existing.BestBefore != msg.TTL || + !mapsEqual(existing.Metadata, msg.Metadata) { + existing.BestBefore = msg.TTL + existing.Metadata = msg.Metadata + changed = true + } + } + p.crdt.metadataMu.Unlock() + + // Trigger membership hook if metadata changed + if changed { + p.crdt.runMembershipHooks() + } + + // Cache summary for sync scheduling p.cacheMu.Lock() p.summaryCache[msg.PeerID] = cachedSummary{ Msg: msg, @@ -649,6 +697,19 @@ func (p *Peer) readSummaries() { } } +// mapsEqual compares two string maps for equality +func mapsEqual(a, b map[string]string) bool { + if len(a) != len(b) { + return false + } + for k, v := range a { + if b[k] != v { + return false + } + } + return true +} + // scheduledSync periodically evaluates cached peer summaries and chooses a peer to synchronize with. // It selects the peer with the largest sequence number gap in tracked states. func (p *Peer) scheduledSync() { @@ -719,3 +780,35 @@ func (p *Peer) cleanupSummaryCache() { } } } + +// cleanupExpiredPeers periodically removes peers whose TTL has expired. +// This runs in a background goroutine and triggers the membership hook when peers are removed. +func (p *Peer) cleanupExpiredPeers() { + ticker := time.NewTicker(1 * time.Minute) + defer ticker.Stop() + + for { + select { + case <-p.ctx.Done(): + return + case <-ticker.C: + p.crdt.metadataMu.Lock() + now := uint64(time.Now().Unix()) + removed := false + + for peerID, meta := range p.crdt.peerMetadata { + if meta.BestBefore < now { + delete(p.crdt.peerMetadata, peerID) + removed = true + log.Printf("Removed expired peer: %s (TTL expired)", peerID) + } + } + p.crdt.metadataMu.Unlock() + + // Trigger membership hook if any peers were removed + if removed { + p.crdt.runMembershipHooks() + } + } + } +}